Skip to main content

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

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