Skip to main content

iscp/
lib.rs

1//! iscpクレートは、iSCP version 2を用いたリアルタイムAPIにアクセスするためのクライアントライブラリです。
2//!
3//! iscpクレートを使用することで、iSCPを利用するクライアントアプリケーションを実装することができます。
4//!
5//! iSCPで通信を行うためには、コネクション[`iscp::Conn`]を確立した後、ストリームやE2Eコールを使用してデータを送受信してください。
6//!
7//! # Examples
8//!
9//! ## Connect To intdash API
10//!
11//! このサンプルではintdash APIとのコネクションを確立します。
12//!
13//! ```no_run
14//! #[tokio::main]
15//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
16//!     #[derive(Clone)]
17//!     struct TokenSource {
18//!         access_token: String,
19//!     }
20//!
21//!     impl iscp::TokenSource for TokenSource {
22//!         async fn token(&mut self) -> Result<iscp::AccessToken, iscp::TokenSourceError> {
23//!             Ok(iscp::AccessToken::new(&self.access_token))
24//!         }
25//!     }
26//!
27//!     let url = std::env::var("EXAMPLE_URL").unwrap_or_else(|_| "wss://xxx.xxx.jp".to_string());
28//!     let api_token = std::env::var("EXAMPLE_TOKEN").unwrap_or_else(|_| {
29//!         "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx".to_string()
30//!     });
31//!     let node_id = std::env::var("EXAMPLE_NODE_ID")
32//!         .unwrap_or_else(|_| "11111111-1111-1111-1111-111111111111".to_string());
33//!
34//!     let token_source = TokenSource {
35//!         // Get an access token by oauth2 from intdash API in actual usage
36//!         access_token: api_token,
37//!     };
38//!
39//!     let connector = iscp::transport::websocket::WebSocketConnector::new(url);
40//!     let conn = iscp::ConnBuilder::new(connector)
41//!         .node_id(node_id)
42//!         .compression(iscp::transport::Compression::new().enable(9))
43//!         .token_source(token_source)
44//!         .build()
45//!         .await?;
46//!
47//!     conn.close().await?;
48//!     Ok(())
49//! }
50//! ```
51//!
52//! ## Start Upstream
53//!
54//! アップストリームの送信サンプルです。このサンプルでは、コネクションからアップストリームを開き、基準時刻のメタデータと、文字列型のデータポイントをiSCPサーバーへ送信しています。
55//!
56//! ```no_run
57//! #[tokio::main]
58//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
59//!     #[derive(Clone)]
60//!     struct TokenSource {
61//!         access_token: String,
62//!     }
63//!
64//!     impl iscp::TokenSource for TokenSource {
65//!         async fn token(&mut self) -> Result<iscp::AccessToken, iscp::TokenSourceError> {
66//!             Ok(iscp::AccessToken::new(&self.access_token))
67//!         }
68//!     }
69//!
70//!     let url = std::env::var("EXAMPLE_URL").unwrap_or_else(|_| "wss://xxx.xxx.jp".to_string());
71//!     let api_token = std::env::var("EXAMPLE_TOKEN").unwrap_or_else(|_| {
72//!         "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx".to_string()
73//!     });
74//!     let node_id = std::env::var("EXAMPLE_NODE_ID")
75//!         .unwrap_or_else(|_| "11111111-1111-1111-1111-111111111111".to_string());
76//!
77//!     let token_source = TokenSource {
78//!         // Get an access token by oauth2 from intdash API in actual usage
79//!         access_token: api_token,
80//!     };
81//!
82//!     let connector = iscp::transport::websocket::WebSocketConnector::new(url);
83//!     let conn = iscp::ConnBuilder::new(connector)
84//!         .node_id(node_id)
85//!         .compression(iscp::transport::Compression::new().enable(9))
86//!         .token_source(token_source)
87//!         .build()
88//!         .await?;
89//!
90//!     // Open an upstream
91//!     let session_id = uuid::Uuid::new_v4();
92//!     let mut up = conn
93//!         .upstream_builder(session_id.to_string())
94//!         .flush_policy(iscp::FlushPolicy::Immediately)
95//!         .expiry_interval(std::time::Duration::from_secs(60))
96//!         .persist(true)
97//!         .build()
98//!         .await?;
99//!
100//!     // Send base time
101//!     let start_clock = std::time::Instant::now();
102//!     let base_time = iscp::metadata::BaseTime {
103//!         base_time: std::time::SystemTime::now(),
104//!         elapsed_time: std::time::Duration::from_secs(0),
105//!         name: "edge_rtc".into(),
106//!         priority: 0,
107//!         session_id: session_id.to_string(),
108//!     };
109//!     conn.metadata_sender(base_time).persist(true).send().await?;
110//!
111//!     // Send a data point
112//!     let data_point = iscp::DataPoint {
113//!         payload: b"hello, world".to_vec().into(),
114//!         elapsed_time: (std::time::Instant::now() - start_clock).as_nanos() as _,
115//!     };
116//!     up.write_data_points(
117//!         iscp::DataId::new("greeting", "string"),
118//!         vec![data_point],
119//!     )
120//!     .await?;
121//!
122//!     conn.close().await?;
123//!     Ok(())
124//! }
125//! ```
126//!
127//! ## Start Downstream
128//!
129//! 前述のアップストリームで送信されたデータをダウンストリームで受信するサンプルです。 このサンプルでは、アップストリーム開始のメタデータ、基準時刻のメタデータ、 文字列型のデータポイントを受信しています。
130//!
131//! ```no_run
132//! #[tokio::main]
133//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
134//!     #[derive(Clone)]
135//!     struct TokenSource {
136//!         access_token: String,
137//!     }
138//!
139//!     impl iscp::TokenSource for TokenSource {
140//!         async fn token(&mut self) -> Result<iscp::AccessToken, iscp::TokenSourceError> {
141//!             Ok(iscp::AccessToken::new(&self.access_token))
142//!         }
143//!     }
144//!
145//!     let url = std::env::var("EXAMPLE_URL").unwrap_or_else(|_| "wss://xxx.xxx.jp".to_string());
146//!     let api_token = std::env::var("EXAMPLE_TOKEN").unwrap_or_else(|_| {
147//!         "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx".to_string()
148//!     });
149//!     let node_id = std::env::var("EXAMPLE_NODE_ID")
150//!         .unwrap_or_else(|_| "11111111-1111-1111-1111-111111111111".to_string());
151//!     let src_node_id = std::env::var("EXAMPLE_SRC_NODE_ID")
152//!         .unwrap_or_else(|_| "22222222-2222-2222-2222-222222222222".to_string());
153//!
154//!     let token_source = TokenSource {
155//!         // Get an access token by oauth2 from intdash API in actual usage
156//!         access_token: api_token,
157//!     };
158//!
159//!     let connector = iscp::transport::websocket::WebSocketConnector::new(url);
160//!     let conn = iscp::ConnBuilder::new(connector)
161//!         .node_id(node_id)
162//!         .compression(iscp::transport::Compression::new().enable(9))
163//!         .token_source(token_source)
164//!         .build()
165//!         .await?;
166//!
167//!     // Open a downstream
168//!     let mut config = iscp::DownstreamConfig::default();
169//!     config.filters = vec![iscp::DownstreamFilter {
170//!         source_node_id: src_node_id.into(),
171//!         data_filters: vec![iscp::DataFilter {
172//!             name: "greeting".into(),
173//!             type_: "string".into(),
174//!         }],
175//!     }];
176//!     let (mut down, mut metadata_reader) = conn.open_downstream_with_config(config).await?;
177//!
178//!     // Read received chunks and metadata
179//!     for _ in 0..4 {
180//!         tokio::select! {
181//!             Ok(chunk) = down.read_chunk() => {
182//!                 println!("received {}: {:?}", chunk.upstream.stream_id, chunk);
183//!             }
184//!             Ok(metadata) = metadata_reader.read() => {
185//!                 println!("received metadata {:?}", metadata);
186//!             }
187//!             else => break,
188//!         }
189//!     }
190//!
191//!     conn.close().await?;
192//!     Ok(())
193//! }
194//! ```
195
196#![cfg_attr(docsrs, feature(doc_cfg))]
197
198#[macro_use]
199mod internal;
200
201pub mod encoding;
202mod error;
203mod iscp;
204pub mod message;
205mod token_source;
206pub mod transport;
207pub mod wire;
208
209pub use crate::error::*;
210pub use crate::iscp::*;
211pub use crate::token_source::*;
212
213/// iSCPプロトコルバージョン
214pub const ISCP_VERSION: &str = "3.0.0";