Skip to main content

nominal_api_conjure/conjure/clients/ingest/api/
streaming_session_service.rs

1use conjure_http::endpoint;
2#[conjure_http::conjure_client(name = "StreamingSessionService")]
3pub trait StreamingSessionService<
4    #[response_body]
5    I: Iterator<
6            Item = Result<conjure_http::private::Bytes, conjure_http::private::Error>,
7        >,
8> {
9    #[endpoint(
10        method = POST,
11        path = "/ingest/v1/internal/streaming-session/dataset/{datasetRid}/resolve",
12        name = "resolve",
13        accept = conjure_http::client::StdResponseDeserializer
14    )]
15    fn resolve(
16        &self,
17        #[auth]
18        auth_: &conjure_object::BearerToken,
19        #[path(
20            name = "datasetRid",
21            encoder = conjure_http::client::conjure::PlainEncoder
22        )]
23        dataset_rid: &super::super::super::super::objects::api::rids::DatasetRid,
24        #[body(serializer = conjure_http::client::StdRequestSerializer)]
25        request: &super::super::super::super::objects::ingest::api::ResolveStreamingSessionRequest,
26    ) -> Result<
27        super::super::super::super::objects::ingest::api::ResolveStreamingSessionResponse,
28        conjure_http::private::Error,
29    >;
30    #[endpoint(
31        method = PUT,
32        path = "/ingest/v1/internal/streaming-session/{sessionRid}/heartbeat",
33        name = "heartbeat",
34        accept = conjure_http::client::conjure::EmptyResponseDeserializer
35    )]
36    fn heartbeat(
37        &self,
38        #[auth]
39        auth_: &conjure_object::BearerToken,
40        #[path(
41            name = "sessionRid",
42            encoder = conjure_http::client::conjure::PlainEncoder
43        )]
44        session_rid: &super::super::super::super::objects::ingest::api::StreamingSessionRid,
45        #[body(serializer = conjure_http::client::StdRequestSerializer)]
46        request: &super::super::super::super::objects::ingest::api::HeartbeatStreamingSessionRequest,
47    ) -> Result<(), conjure_http::private::Error>;
48    /// Returns a paginated list of streaming sessions, optionally filtered by dataset.
49    #[endpoint(
50        method = POST,
51        path = "/ingest/v1/streaming-sessions/search",
52        name = "searchStreamingSessions",
53        accept = conjure_http::client::StdResponseDeserializer
54    )]
55    fn search_streaming_sessions(
56        &self,
57        #[auth]
58        auth_: &conjure_object::BearerToken,
59        #[body(serializer = conjure_http::client::StdRequestSerializer)]
60        request: &super::super::super::super::objects::ingest::api::SearchStreamingSessionsRequest,
61    ) -> Result<
62        super::super::super::super::objects::ingest::api::SearchStreamingSessionsResponse,
63        conjure_http::private::Error,
64    >;
65    /// Returns a paginated list of streaming sessions aggregated by (source, dataset) group,
66    /// ordered by latest activity descending.
67    #[endpoint(
68        method = POST,
69        path = "/ingest/v1/streaming-sources/search",
70        name = "searchStreamingSources",
71        accept = conjure_http::client::StdResponseDeserializer
72    )]
73    fn search_streaming_sources(
74        &self,
75        #[auth]
76        auth_: &conjure_object::BearerToken,
77        #[body(serializer = conjure_http::client::StdRequestSerializer)]
78        request: &super::super::super::super::objects::ingest::api::SearchStreamingSourcesRequest,
79    ) -> Result<
80        super::super::super::super::objects::ingest::api::SearchStreamingSourcesResponse,
81        conjure_http::private::Error,
82    >;
83}
84#[conjure_http::conjure_client(name = "StreamingSessionService")]
85pub trait AsyncStreamingSessionService<
86    #[response_body]
87    I: conjure_http::private::Stream<
88            Item = Result<conjure_http::private::Bytes, conjure_http::private::Error>,
89        >,
90> {
91    #[endpoint(
92        method = POST,
93        path = "/ingest/v1/internal/streaming-session/dataset/{datasetRid}/resolve",
94        name = "resolve",
95        accept = conjure_http::client::StdResponseDeserializer
96    )]
97    async fn resolve(
98        &self,
99        #[auth]
100        auth_: &conjure_object::BearerToken,
101        #[path(
102            name = "datasetRid",
103            encoder = conjure_http::client::conjure::PlainEncoder
104        )]
105        dataset_rid: &super::super::super::super::objects::api::rids::DatasetRid,
106        #[body(serializer = conjure_http::client::StdRequestSerializer)]
107        request: &super::super::super::super::objects::ingest::api::ResolveStreamingSessionRequest,
108    ) -> Result<
109        super::super::super::super::objects::ingest::api::ResolveStreamingSessionResponse,
110        conjure_http::private::Error,
111    >;
112    #[endpoint(
113        method = PUT,
114        path = "/ingest/v1/internal/streaming-session/{sessionRid}/heartbeat",
115        name = "heartbeat",
116        accept = conjure_http::client::conjure::EmptyResponseDeserializer
117    )]
118    async fn heartbeat(
119        &self,
120        #[auth]
121        auth_: &conjure_object::BearerToken,
122        #[path(
123            name = "sessionRid",
124            encoder = conjure_http::client::conjure::PlainEncoder
125        )]
126        session_rid: &super::super::super::super::objects::ingest::api::StreamingSessionRid,
127        #[body(serializer = conjure_http::client::StdRequestSerializer)]
128        request: &super::super::super::super::objects::ingest::api::HeartbeatStreamingSessionRequest,
129    ) -> Result<(), conjure_http::private::Error>;
130    /// Returns a paginated list of streaming sessions, optionally filtered by dataset.
131    #[endpoint(
132        method = POST,
133        path = "/ingest/v1/streaming-sessions/search",
134        name = "searchStreamingSessions",
135        accept = conjure_http::client::StdResponseDeserializer
136    )]
137    async fn search_streaming_sessions(
138        &self,
139        #[auth]
140        auth_: &conjure_object::BearerToken,
141        #[body(serializer = conjure_http::client::StdRequestSerializer)]
142        request: &super::super::super::super::objects::ingest::api::SearchStreamingSessionsRequest,
143    ) -> Result<
144        super::super::super::super::objects::ingest::api::SearchStreamingSessionsResponse,
145        conjure_http::private::Error,
146    >;
147    /// Returns a paginated list of streaming sessions aggregated by (source, dataset) group,
148    /// ordered by latest activity descending.
149    #[endpoint(
150        method = POST,
151        path = "/ingest/v1/streaming-sources/search",
152        name = "searchStreamingSources",
153        accept = conjure_http::client::StdResponseDeserializer
154    )]
155    async fn search_streaming_sources(
156        &self,
157        #[auth]
158        auth_: &conjure_object::BearerToken,
159        #[body(serializer = conjure_http::client::StdRequestSerializer)]
160        request: &super::super::super::super::objects::ingest::api::SearchStreamingSourcesRequest,
161    ) -> Result<
162        super::super::super::super::objects::ingest::api::SearchStreamingSourcesResponse,
163        conjure_http::private::Error,
164    >;
165}
166#[conjure_http::conjure_client(name = "StreamingSessionService", local)]
167pub trait LocalAsyncStreamingSessionService<
168    #[response_body]
169    I: conjure_http::private::Stream<
170            Item = Result<conjure_http::private::Bytes, conjure_http::private::Error>,
171        >,
172> {
173    #[endpoint(
174        method = POST,
175        path = "/ingest/v1/internal/streaming-session/dataset/{datasetRid}/resolve",
176        name = "resolve",
177        accept = conjure_http::client::StdResponseDeserializer
178    )]
179    async fn resolve(
180        &self,
181        #[auth]
182        auth_: &conjure_object::BearerToken,
183        #[path(
184            name = "datasetRid",
185            encoder = conjure_http::client::conjure::PlainEncoder
186        )]
187        dataset_rid: &super::super::super::super::objects::api::rids::DatasetRid,
188        #[body(serializer = conjure_http::client::StdRequestSerializer)]
189        request: &super::super::super::super::objects::ingest::api::ResolveStreamingSessionRequest,
190    ) -> Result<
191        super::super::super::super::objects::ingest::api::ResolveStreamingSessionResponse,
192        conjure_http::private::Error,
193    >;
194    #[endpoint(
195        method = PUT,
196        path = "/ingest/v1/internal/streaming-session/{sessionRid}/heartbeat",
197        name = "heartbeat",
198        accept = conjure_http::client::conjure::EmptyResponseDeserializer
199    )]
200    async fn heartbeat(
201        &self,
202        #[auth]
203        auth_: &conjure_object::BearerToken,
204        #[path(
205            name = "sessionRid",
206            encoder = conjure_http::client::conjure::PlainEncoder
207        )]
208        session_rid: &super::super::super::super::objects::ingest::api::StreamingSessionRid,
209        #[body(serializer = conjure_http::client::StdRequestSerializer)]
210        request: &super::super::super::super::objects::ingest::api::HeartbeatStreamingSessionRequest,
211    ) -> Result<(), conjure_http::private::Error>;
212    /// Returns a paginated list of streaming sessions, optionally filtered by dataset.
213    #[endpoint(
214        method = POST,
215        path = "/ingest/v1/streaming-sessions/search",
216        name = "searchStreamingSessions",
217        accept = conjure_http::client::StdResponseDeserializer
218    )]
219    async fn search_streaming_sessions(
220        &self,
221        #[auth]
222        auth_: &conjure_object::BearerToken,
223        #[body(serializer = conjure_http::client::StdRequestSerializer)]
224        request: &super::super::super::super::objects::ingest::api::SearchStreamingSessionsRequest,
225    ) -> Result<
226        super::super::super::super::objects::ingest::api::SearchStreamingSessionsResponse,
227        conjure_http::private::Error,
228    >;
229    /// Returns a paginated list of streaming sessions aggregated by (source, dataset) group,
230    /// ordered by latest activity descending.
231    #[endpoint(
232        method = POST,
233        path = "/ingest/v1/streaming-sources/search",
234        name = "searchStreamingSources",
235        accept = conjure_http::client::StdResponseDeserializer
236    )]
237    async fn search_streaming_sources(
238        &self,
239        #[auth]
240        auth_: &conjure_object::BearerToken,
241        #[body(serializer = conjure_http::client::StdRequestSerializer)]
242        request: &super::super::super::super::objects::ingest::api::SearchStreamingSourcesRequest,
243    ) -> Result<
244        super::super::super::super::objects::ingest::api::SearchStreamingSourcesResponse,
245        conjure_http::private::Error,
246    >;
247}