1#[derive(prost_helpers::AnyPB)]
3#[derive(Clone, PartialEq, ::prost::Message)]
4pub struct TableSchema {
5 #[prost(message, repeated, tag = "1")]
6 pub columns: ::prost::alloc::vec::Vec<super::plan_common::ColumnDesc>,
7 #[prost(uint32, repeated, tag = "2")]
8 pub pk_indices: ::prost::alloc::vec::Vec<u32>,
9}
10#[derive(prost_helpers::AnyPB)]
11#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
12pub struct ValidationError {
13 #[prost(string, tag = "1")]
14 pub error_message: ::prost::alloc::string::String,
15}
16#[derive(prost_helpers::AnyPB)]
17#[derive(Clone, PartialEq, ::prost::Message)]
18pub struct SinkParam {
19 #[prost(uint32, tag = "1", wrapper = "crate::id::SinkId")]
20 pub sink_id: crate::id::SinkId,
21 #[prost(btree_map = "string, string", tag = "2")]
22 pub properties: ::prost::alloc::collections::BTreeMap<
23 ::prost::alloc::string::String,
24 ::prost::alloc::string::String,
25 >,
26 #[prost(message, optional, tag = "3")]
27 pub table_schema: ::core::option::Option<TableSchema>,
28 #[prost(enumeration = "super::catalog::SinkType", tag = "4")]
30 pub sink_type: i32,
31 #[prost(string, tag = "5")]
32 pub db_name: ::prost::alloc::string::String,
33 #[prost(string, tag = "6")]
34 pub sink_from_name: ::prost::alloc::string::String,
35 #[prost(message, optional, tag = "7")]
36 pub format_desc: ::core::option::Option<super::catalog::SinkFormatDesc>,
37 #[prost(string, tag = "8")]
38 pub sink_name: ::prost::alloc::string::String,
39 #[prost(bool, tag = "9")]
43 pub raw_ignore_delete: bool,
44}
45#[derive(prost_helpers::AnyPB)]
46#[derive(Clone, PartialEq, ::prost::Message)]
47pub struct SinkWriterStreamRequest {
48 #[prost(oneof = "sink_writer_stream_request::Request", tags = "1, 3, 4")]
49 pub request: ::core::option::Option<sink_writer_stream_request::Request>,
50}
51pub mod sink_writer_stream_request {
53 #[derive(prost_helpers::AnyPB)]
54 #[derive(Clone, PartialEq, ::prost::Message)]
55 pub struct StartSink {
56 #[prost(message, optional, tag = "1")]
57 pub sink_param: ::core::option::Option<super::SinkParam>,
58 #[prost(message, optional, tag = "3")]
59 pub payload_schema: ::core::option::Option<super::TableSchema>,
60 }
61 #[derive(prost_helpers::AnyPB)]
62 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
63 pub struct WriteBatch {
64 #[prost(uint64, tag = "3")]
65 pub batch_id: u64,
66 #[prost(uint64, tag = "4")]
67 pub epoch: u64,
68 #[prost(oneof = "write_batch::Payload", tags = "2, 5")]
69 pub payload: ::core::option::Option<write_batch::Payload>,
70 }
71 pub mod write_batch {
73 #[derive(prost_helpers::AnyPB)]
74 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
75 pub struct StreamChunkPayload {
76 #[prost(bytes = "vec", tag = "1")]
77 pub binary_data: ::prost::alloc::vec::Vec<u8>,
78 }
79 #[derive(prost_helpers::AnyPB)]
80 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
81 pub enum Payload {
82 #[prost(message, tag = "2")]
83 StreamChunkPayload(StreamChunkPayload),
84 #[prost(int64, tag = "5")]
88 StreamChunkRefPointer(i64),
89 }
90 }
91 #[derive(prost_helpers::AnyPB)]
92 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
93 pub struct Barrier {
94 #[prost(uint64, tag = "1")]
95 pub epoch: u64,
96 #[prost(bool, tag = "2")]
97 pub is_checkpoint: bool,
98 }
99 #[derive(prost_helpers::AnyPB)]
100 #[derive(Clone, PartialEq, ::prost::Oneof)]
101 pub enum Request {
102 #[prost(message, tag = "1")]
103 Start(StartSink),
104 #[prost(message, tag = "3")]
105 WriteBatch(WriteBatch),
106 #[prost(message, tag = "4")]
107 Barrier(Barrier),
108 }
109}
110#[derive(prost_helpers::AnyPB)]
111#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
112pub struct SinkWriterStreamResponse {
113 #[prost(oneof = "sink_writer_stream_response::Response", tags = "1, 2, 3")]
114 pub response: ::core::option::Option<sink_writer_stream_response::Response>,
115}
116pub mod sink_writer_stream_response {
118 #[derive(prost_helpers::AnyPB)]
119 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
120 pub struct StartResponse {}
121 #[derive(prost_helpers::AnyPB)]
122 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
123 pub struct CommitResponse {
124 #[prost(uint64, tag = "1")]
125 pub epoch: u64,
126 #[prost(message, optional, tag = "2")]
127 pub metadata: ::core::option::Option<super::SinkMetadata>,
128 }
129 #[derive(prost_helpers::AnyPB)]
130 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
131 pub struct BatchWrittenResponse {
132 #[prost(uint64, tag = "1")]
133 pub epoch: u64,
134 #[prost(uint64, tag = "2")]
135 pub batch_id: u64,
136 }
137 #[derive(prost_helpers::AnyPB)]
138 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
139 pub enum Response {
140 #[prost(message, tag = "1")]
141 Start(StartResponse),
142 #[prost(message, tag = "2")]
143 Commit(CommitResponse),
144 #[prost(message, tag = "3")]
145 Batch(BatchWrittenResponse),
146 }
147}
148#[derive(prost_helpers::AnyPB)]
149#[derive(Clone, PartialEq, ::prost::Message)]
150pub struct ValidateSinkRequest {
151 #[prost(message, optional, tag = "1")]
152 pub sink_param: ::core::option::Option<SinkParam>,
153}
154#[derive(prost_helpers::AnyPB)]
155#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
156pub struct ValidateSinkResponse {
157 #[prost(message, optional, tag = "1")]
159 pub error: ::core::option::Option<ValidationError>,
160}
161#[derive(prost_helpers::AnyPB)]
162#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
163pub struct SinkMetadata {
164 #[prost(oneof = "sink_metadata::Metadata", tags = "1")]
165 pub metadata: ::core::option::Option<sink_metadata::Metadata>,
166}
167pub mod sink_metadata {
169 #[derive(prost_helpers::AnyPB)]
170 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
171 pub struct SerializedMetadata {
172 #[prost(bytes = "vec", tag = "1")]
173 pub metadata: ::prost::alloc::vec::Vec<u8>,
174 }
175 #[derive(prost_helpers::AnyPB)]
176 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)]
177 pub enum Metadata {
178 #[prost(message, tag = "1")]
179 Serialized(SerializedMetadata),
180 }
181}
182#[derive(prost_helpers::AnyPB)]
183#[derive(Clone, PartialEq, ::prost::Message)]
184pub struct SinkCoordinatorStreamRequest {
185 #[prost(oneof = "sink_coordinator_stream_request::Request", tags = "1, 2")]
186 pub request: ::core::option::Option<sink_coordinator_stream_request::Request>,
187}
188pub mod sink_coordinator_stream_request {
190 #[derive(prost_helpers::AnyPB)]
191 #[derive(Clone, PartialEq, ::prost::Message)]
192 pub struct StartCoordinator {
193 #[prost(message, optional, tag = "1")]
194 pub param: ::core::option::Option<super::SinkParam>,
195 }
196 #[derive(prost_helpers::AnyPB)]
197 #[derive(Clone, PartialEq, ::prost::Message)]
198 pub struct CommitMetadata {
199 #[prost(uint64, tag = "1")]
200 pub epoch: u64,
201 #[prost(message, repeated, tag = "2")]
202 pub metadata: ::prost::alloc::vec::Vec<super::SinkMetadata>,
203 }
204 #[derive(prost_helpers::AnyPB)]
205 #[derive(Clone, PartialEq, ::prost::Oneof)]
206 pub enum Request {
207 #[prost(message, tag = "1")]
208 Start(StartCoordinator),
209 #[prost(message, tag = "2")]
210 Commit(CommitMetadata),
211 }
212}
213#[derive(prost_helpers::AnyPB)]
214#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
215pub struct SinkCoordinatorStreamResponse {
216 #[prost(oneof = "sink_coordinator_stream_response::Response", tags = "1, 2")]
217 pub response: ::core::option::Option<sink_coordinator_stream_response::Response>,
218}
219pub mod sink_coordinator_stream_response {
221 #[derive(prost_helpers::AnyPB)]
222 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
223 pub struct StartResponse {}
224 #[derive(prost_helpers::AnyPB)]
225 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
226 pub struct CommitResponse {
227 #[prost(uint64, tag = "1")]
228 pub epoch: u64,
229 }
230 #[derive(prost_helpers::AnyPB)]
231 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
232 pub enum Response {
233 #[prost(message, tag = "1")]
234 Start(StartResponse),
235 #[prost(message, tag = "2")]
236 Commit(CommitResponse),
237 }
238}
239#[derive(prost_helpers::AnyPB)]
240#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
241pub struct CdcMessage {
242 #[prost(string, tag = "1")]
244 pub payload: ::prost::alloc::string::String,
245 #[prost(string, tag = "2")]
246 pub partition: ::prost::alloc::string::String,
247 #[prost(string, tag = "3")]
248 pub offset: ::prost::alloc::string::String,
249 #[prost(string, tag = "4")]
250 pub full_table_name: ::prost::alloc::string::String,
251 #[prost(int64, tag = "5")]
252 pub source_ts_ms: i64,
253 #[prost(enumeration = "cdc_message::CdcMessageType", tag = "6")]
254 pub msg_type: i32,
255 #[prost(string, tag = "7")]
257 pub key: ::prost::alloc::string::String,
258 #[prost(enumeration = "SourceType", tag = "8")]
262 pub source_type: i32,
263}
264pub mod cdc_message {
266 #[derive(prost_helpers::AnyPB)]
267 #[derive(
268 Clone,
269 Copy,
270 Debug,
271 PartialEq,
272 Eq,
273 Hash,
274 PartialOrd,
275 Ord,
276 ::prost::Enumeration
277 )]
278 #[repr(i32)]
279 pub enum CdcMessageType {
280 Unspecified = 0,
281 Heartbeat = 1,
282 Data = 2,
283 TransactionMeta = 3,
284 SchemaChange = 4,
285 }
286 impl CdcMessageType {
287 pub fn as_str_name(&self) -> &'static str {
292 match self {
293 Self::Unspecified => "UNSPECIFIED",
294 Self::Heartbeat => "HEARTBEAT",
295 Self::Data => "DATA",
296 Self::TransactionMeta => "TRANSACTION_META",
297 Self::SchemaChange => "SCHEMA_CHANGE",
298 }
299 }
300 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
302 match value {
303 "UNSPECIFIED" => Some(Self::Unspecified),
304 "HEARTBEAT" => Some(Self::Heartbeat),
305 "DATA" => Some(Self::Data),
306 "TRANSACTION_META" => Some(Self::TransactionMeta),
307 "SCHEMA_CHANGE" => Some(Self::SchemaChange),
308 _ => None,
309 }
310 }
311 }
312}
313#[derive(prost_helpers::AnyPB)]
314#[derive(Clone, PartialEq, ::prost::Message)]
315pub struct GetEventStreamRequest {
316 #[prost(uint64, tag = "1")]
317 pub source_id: u64,
318 #[prost(enumeration = "SourceType", tag = "2")]
319 pub source_type: i32,
320 #[prost(string, tag = "3")]
321 pub start_offset: ::prost::alloc::string::String,
322 #[prost(btree_map = "string, string", tag = "4")]
323 pub properties: ::prost::alloc::collections::BTreeMap<
324 ::prost::alloc::string::String,
325 ::prost::alloc::string::String,
326 >,
327 #[prost(bool, tag = "5")]
328 pub snapshot_done: bool,
329 #[prost(bool, tag = "6")]
330 pub is_source_job: bool,
331}
332#[derive(prost_helpers::AnyPB)]
333#[derive(Clone, PartialEq, ::prost::Message)]
334pub struct GetEventStreamResponse {
335 #[prost(uint64, tag = "1")]
336 pub source_id: u64,
337 #[prost(message, repeated, tag = "2")]
338 pub events: ::prost::alloc::vec::Vec<CdcMessage>,
339 #[prost(message, optional, tag = "3")]
340 pub control: ::core::option::Option<get_event_stream_response::ControlInfo>,
341}
342pub mod get_event_stream_response {
344 #[derive(prost_helpers::AnyPB)]
345 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
346 pub struct ControlInfo {
347 #[prost(bool, tag = "1")]
348 pub handshake_ok: bool,
349 }
350}
351#[derive(prost_helpers::AnyPB)]
352#[derive(Clone, PartialEq, ::prost::Message)]
353pub struct ValidateSourceRequest {
354 #[prost(uint64, tag = "1")]
355 pub source_id: u64,
356 #[prost(enumeration = "SourceType", tag = "2")]
357 pub source_type: i32,
358 #[prost(btree_map = "string, string", tag = "3")]
359 pub properties: ::prost::alloc::collections::BTreeMap<
360 ::prost::alloc::string::String,
361 ::prost::alloc::string::String,
362 >,
363 #[prost(message, optional, tag = "4")]
364 pub table_schema: ::core::option::Option<TableSchema>,
365 #[prost(bool, tag = "5")]
366 pub is_source_job: bool,
367 #[prost(bool, tag = "6")]
368 pub is_backfill_table: bool,
369}
370#[derive(prost_helpers::AnyPB)]
371#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
372pub struct ValidateSourceResponse {
373 #[prost(message, optional, tag = "1")]
375 pub error: ::core::option::Option<ValidationError>,
376}
377#[derive(prost_helpers::AnyPB)]
378#[derive(Clone, PartialEq, ::prost::Message)]
379pub struct CoordinateRequest {
380 #[prost(oneof = "coordinate_request::Msg", tags = "1, 2, 3, 4, 5")]
381 pub msg: ::core::option::Option<coordinate_request::Msg>,
382}
383pub mod coordinate_request {
385 #[derive(prost_helpers::AnyPB)]
388 #[derive(Clone, PartialEq, ::prost::Message)]
389 pub struct StartCoordinationRequest {
390 #[prost(message, optional, tag = "1")]
391 pub vnode_bitmap: ::core::option::Option<super::super::common::Buffer>,
392 #[prost(message, optional, tag = "2")]
393 pub param: ::core::option::Option<super::SinkParam>,
394 }
395 #[derive(prost_helpers::AnyPB)]
396 #[derive(Clone, PartialEq, ::prost::Message)]
397 pub struct CommitRequest {
398 #[prost(uint64, tag = "1")]
399 pub epoch: u64,
400 #[prost(message, optional, tag = "2")]
401 pub metadata: ::core::option::Option<super::SinkMetadata>,
402 #[prost(message, optional, tag = "3")]
403 pub schema_change: ::core::option::Option<
404 super::super::stream_plan::SinkSchemaChange,
405 >,
406 }
407 #[derive(prost_helpers::AnyPB)]
408 #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
409 pub struct UpdateVnodeBitmapRequest {
410 #[prost(message, optional, tag = "1")]
411 pub vnode_bitmap: ::core::option::Option<super::super::common::Buffer>,
412 }
413 #[derive(prost_helpers::AnyPB)]
414 #[derive(Clone, PartialEq, ::prost::Oneof)]
415 pub enum Msg {
416 #[prost(message, tag = "1")]
417 StartRequest(StartCoordinationRequest),
418 #[prost(message, tag = "2")]
419 CommitRequest(CommitRequest),
420 #[prost(message, tag = "3")]
421 UpdateVnodeRequest(UpdateVnodeBitmapRequest),
422 #[prost(bool, tag = "4")]
423 Stop(bool),
424 #[prost(uint64, tag = "5")]
425 AlignInitialEpochRequest(u64),
426 }
427}
428#[derive(prost_helpers::AnyPB)]
429#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
430pub struct CoordinateResponse {
431 #[prost(oneof = "coordinate_response::Msg", tags = "1, 2, 3, 4")]
432 pub msg: ::core::option::Option<coordinate_response::Msg>,
433}
434pub mod coordinate_response {
436 #[derive(prost_helpers::AnyPB)]
437 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
438 pub struct StartCoordinationResponse {
439 #[prost(uint64, optional, tag = "1")]
440 pub log_store_rewind_start_epoch: ::core::option::Option<u64>,
441 }
442 #[derive(prost_helpers::AnyPB)]
443 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
444 pub struct CommitResponse {
445 #[prost(uint64, tag = "1")]
446 pub epoch: u64,
447 }
448 #[derive(prost_helpers::AnyPB)]
449 #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Oneof)]
450 pub enum Msg {
451 #[prost(message, tag = "1")]
452 StartResponse(StartCoordinationResponse),
453 #[prost(message, tag = "2")]
454 CommitResponse(CommitResponse),
455 #[prost(bool, tag = "3")]
456 Stopped(bool),
457 #[prost(uint64, tag = "4")]
458 AlignInitialEpochResponse(u64),
459 }
460}
461#[derive(prost_helpers::AnyPB)]
462#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
463pub struct IcebergPkIndexPreCommitState {
464 #[prost(bytes = "vec", tag = "1")]
465 pub agg_result: ::prost::alloc::vec::Vec<u8>,
466 #[prost(int64, tag = "2")]
467 pub snapshot_id: i64,
468}
469#[derive(prost_helpers::AnyPB)]
470#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
471#[repr(i32)]
472pub enum SourceType {
473 Unspecified = 0,
474 Mysql = 1,
475 Postgres = 2,
476 Citus = 3,
477 Mongodb = 4,
478 SqlServer = 5,
479 Oracle = 6,
480}
481impl SourceType {
482 pub fn as_str_name(&self) -> &'static str {
487 match self {
488 Self::Unspecified => "UNSPECIFIED",
489 Self::Mysql => "MYSQL",
490 Self::Postgres => "POSTGRES",
491 Self::Citus => "CITUS",
492 Self::Mongodb => "MONGODB",
493 Self::SqlServer => "SQL_SERVER",
494 Self::Oracle => "ORACLE",
495 }
496 }
497 pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
499 match value {
500 "UNSPECIFIED" => Some(Self::Unspecified),
501 "MYSQL" => Some(Self::Mysql),
502 "POSTGRES" => Some(Self::Postgres),
503 "CITUS" => Some(Self::Citus),
504 "MONGODB" => Some(Self::Mongodb),
505 "SQL_SERVER" => Some(Self::SqlServer),
506 "ORACLE" => Some(Self::Oracle),
507 _ => None,
508 }
509 }
510}
511pub mod connector_service_client {
513 #![allow(
514 unused_variables,
515 dead_code,
516 missing_docs,
517 clippy::wildcard_imports,
518 clippy::let_unit_value,
519 )]
520 use tonic::codegen::*;
521 use tonic::codegen::http::Uri;
522 #[derive(Debug, Clone)]
523 pub struct ConnectorServiceClient<T> {
524 inner: tonic::client::Grpc<T>,
525 }
526 impl ConnectorServiceClient<tonic::transport::Channel> {
527 pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
529 where
530 D: TryInto<tonic::transport::Endpoint>,
531 D::Error: Into<StdError>,
532 {
533 let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
534 Ok(Self::new(conn))
535 }
536 }
537 impl<T> ConnectorServiceClient<T>
538 where
539 T: tonic::client::GrpcService<tonic::body::Body>,
540 T::Error: Into<StdError>,
541 T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
542 <T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
543 {
544 pub fn new(inner: T) -> Self {
545 let inner = tonic::client::Grpc::new(inner);
546 Self { inner }
547 }
548 pub fn with_origin(inner: T, origin: Uri) -> Self {
549 let inner = tonic::client::Grpc::with_origin(inner, origin);
550 Self { inner }
551 }
552 pub fn with_interceptor<F>(
553 inner: T,
554 interceptor: F,
555 ) -> ConnectorServiceClient<InterceptedService<T, F>>
556 where
557 F: tonic::service::Interceptor,
558 T::ResponseBody: Default,
559 T: tonic::codegen::Service<
560 http::Request<tonic::body::Body>,
561 Response = http::Response<
562 <T as tonic::client::GrpcService<tonic::body::Body>>::ResponseBody,
563 >,
564 >,
565 <T as tonic::codegen::Service<
566 http::Request<tonic::body::Body>,
567 >>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
568 {
569 ConnectorServiceClient::new(InterceptedService::new(inner, interceptor))
570 }
571 #[must_use]
576 pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
577 self.inner = self.inner.send_compressed(encoding);
578 self
579 }
580 #[must_use]
582 pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
583 self.inner = self.inner.accept_compressed(encoding);
584 self
585 }
586 #[must_use]
590 pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
591 self.inner = self.inner.max_decoding_message_size(limit);
592 self
593 }
594 #[must_use]
598 pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
599 self.inner = self.inner.max_encoding_message_size(limit);
600 self
601 }
602 pub async fn sink_writer_stream(
603 &mut self,
604 request: impl tonic::IntoStreamingRequest<
605 Message = super::SinkWriterStreamRequest,
606 >,
607 ) -> std::result::Result<
608 tonic::Response<tonic::codec::Streaming<super::SinkWriterStreamResponse>>,
609 tonic::Status,
610 > {
611 self.inner
612 .ready()
613 .await
614 .map_err(|e| {
615 tonic::Status::unknown(
616 format!("Service was not ready: {}", e.into()),
617 )
618 })?;
619 let codec = tonic_prost::ProstCodec::default();
620 let path = http::uri::PathAndQuery::from_static(
621 "/connector_service.ConnectorService/SinkWriterStream",
622 );
623 let mut req = request.into_streaming_request();
624 req.extensions_mut()
625 .insert(
626 GrpcMethod::new(
627 "connector_service.ConnectorService",
628 "SinkWriterStream",
629 ),
630 );
631 self.inner.streaming(req, path, codec).await
632 }
633 pub async fn sink_coordinator_stream(
634 &mut self,
635 request: impl tonic::IntoStreamingRequest<
636 Message = super::SinkCoordinatorStreamRequest,
637 >,
638 ) -> std::result::Result<
639 tonic::Response<
640 tonic::codec::Streaming<super::SinkCoordinatorStreamResponse>,
641 >,
642 tonic::Status,
643 > {
644 self.inner
645 .ready()
646 .await
647 .map_err(|e| {
648 tonic::Status::unknown(
649 format!("Service was not ready: {}", e.into()),
650 )
651 })?;
652 let codec = tonic_prost::ProstCodec::default();
653 let path = http::uri::PathAndQuery::from_static(
654 "/connector_service.ConnectorService/SinkCoordinatorStream",
655 );
656 let mut req = request.into_streaming_request();
657 req.extensions_mut()
658 .insert(
659 GrpcMethod::new(
660 "connector_service.ConnectorService",
661 "SinkCoordinatorStream",
662 ),
663 );
664 self.inner.streaming(req, path, codec).await
665 }
666 pub async fn validate_sink(
667 &mut self,
668 request: impl tonic::IntoRequest<super::ValidateSinkRequest>,
669 ) -> std::result::Result<
670 tonic::Response<super::ValidateSinkResponse>,
671 tonic::Status,
672 > {
673 self.inner
674 .ready()
675 .await
676 .map_err(|e| {
677 tonic::Status::unknown(
678 format!("Service was not ready: {}", e.into()),
679 )
680 })?;
681 let codec = tonic_prost::ProstCodec::default();
682 let path = http::uri::PathAndQuery::from_static(
683 "/connector_service.ConnectorService/ValidateSink",
684 );
685 let mut req = request.into_request();
686 req.extensions_mut()
687 .insert(
688 GrpcMethod::new("connector_service.ConnectorService", "ValidateSink"),
689 );
690 self.inner.unary(req, path, codec).await
691 }
692 pub async fn get_event_stream(
693 &mut self,
694 request: impl tonic::IntoRequest<super::GetEventStreamRequest>,
695 ) -> std::result::Result<
696 tonic::Response<tonic::codec::Streaming<super::GetEventStreamResponse>>,
697 tonic::Status,
698 > {
699 self.inner
700 .ready()
701 .await
702 .map_err(|e| {
703 tonic::Status::unknown(
704 format!("Service was not ready: {}", e.into()),
705 )
706 })?;
707 let codec = tonic_prost::ProstCodec::default();
708 let path = http::uri::PathAndQuery::from_static(
709 "/connector_service.ConnectorService/GetEventStream",
710 );
711 let mut req = request.into_request();
712 req.extensions_mut()
713 .insert(
714 GrpcMethod::new(
715 "connector_service.ConnectorService",
716 "GetEventStream",
717 ),
718 );
719 self.inner.server_streaming(req, path, codec).await
720 }
721 pub async fn validate_source(
722 &mut self,
723 request: impl tonic::IntoRequest<super::ValidateSourceRequest>,
724 ) -> std::result::Result<
725 tonic::Response<super::ValidateSourceResponse>,
726 tonic::Status,
727 > {
728 self.inner
729 .ready()
730 .await
731 .map_err(|e| {
732 tonic::Status::unknown(
733 format!("Service was not ready: {}", e.into()),
734 )
735 })?;
736 let codec = tonic_prost::ProstCodec::default();
737 let path = http::uri::PathAndQuery::from_static(
738 "/connector_service.ConnectorService/ValidateSource",
739 );
740 let mut req = request.into_request();
741 req.extensions_mut()
742 .insert(
743 GrpcMethod::new(
744 "connector_service.ConnectorService",
745 "ValidateSource",
746 ),
747 );
748 self.inner.unary(req, path, codec).await
749 }
750 }
751}
752pub mod connector_service_server {
754 #![allow(
755 unused_variables,
756 dead_code,
757 missing_docs,
758 clippy::wildcard_imports,
759 clippy::let_unit_value,
760 )]
761 use tonic::codegen::*;
762 #[async_trait]
764 pub trait ConnectorService: std::marker::Send + std::marker::Sync + 'static {
765 type SinkWriterStreamStream: tonic::codegen::tokio_stream::Stream<
767 Item = std::result::Result<
768 super::SinkWriterStreamResponse,
769 tonic::Status,
770 >,
771 >
772 + std::marker::Send
773 + 'static;
774 async fn sink_writer_stream(
775 &self,
776 request: tonic::Request<tonic::Streaming<super::SinkWriterStreamRequest>>,
777 ) -> std::result::Result<
778 tonic::Response<Self::SinkWriterStreamStream>,
779 tonic::Status,
780 >;
781 type SinkCoordinatorStreamStream: tonic::codegen::tokio_stream::Stream<
783 Item = std::result::Result<
784 super::SinkCoordinatorStreamResponse,
785 tonic::Status,
786 >,
787 >
788 + std::marker::Send
789 + 'static;
790 async fn sink_coordinator_stream(
791 &self,
792 request: tonic::Request<
793 tonic::Streaming<super::SinkCoordinatorStreamRequest>,
794 >,
795 ) -> std::result::Result<
796 tonic::Response<Self::SinkCoordinatorStreamStream>,
797 tonic::Status,
798 >;
799 async fn validate_sink(
800 &self,
801 request: tonic::Request<super::ValidateSinkRequest>,
802 ) -> std::result::Result<
803 tonic::Response<super::ValidateSinkResponse>,
804 tonic::Status,
805 >;
806 type GetEventStreamStream: tonic::codegen::tokio_stream::Stream<
808 Item = std::result::Result<super::GetEventStreamResponse, tonic::Status>,
809 >
810 + std::marker::Send
811 + 'static;
812 async fn get_event_stream(
813 &self,
814 request: tonic::Request<super::GetEventStreamRequest>,
815 ) -> std::result::Result<
816 tonic::Response<Self::GetEventStreamStream>,
817 tonic::Status,
818 >;
819 async fn validate_source(
820 &self,
821 request: tonic::Request<super::ValidateSourceRequest>,
822 ) -> std::result::Result<
823 tonic::Response<super::ValidateSourceResponse>,
824 tonic::Status,
825 >;
826 }
827 #[derive(Debug)]
828 pub struct ConnectorServiceServer<T> {
829 inner: Arc<T>,
830 accept_compression_encodings: EnabledCompressionEncodings,
831 send_compression_encodings: EnabledCompressionEncodings,
832 max_decoding_message_size: Option<usize>,
833 max_encoding_message_size: Option<usize>,
834 }
835 impl<T> ConnectorServiceServer<T> {
836 pub fn new(inner: T) -> Self {
837 Self::from_arc(Arc::new(inner))
838 }
839 pub fn from_arc(inner: Arc<T>) -> Self {
840 Self {
841 inner,
842 accept_compression_encodings: Default::default(),
843 send_compression_encodings: Default::default(),
844 max_decoding_message_size: None,
845 max_encoding_message_size: None,
846 }
847 }
848 pub fn with_interceptor<F>(
849 inner: T,
850 interceptor: F,
851 ) -> InterceptedService<Self, F>
852 where
853 F: tonic::service::Interceptor,
854 {
855 InterceptedService::new(Self::new(inner), interceptor)
856 }
857 #[must_use]
859 pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
860 self.accept_compression_encodings.enable(encoding);
861 self
862 }
863 #[must_use]
865 pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
866 self.send_compression_encodings.enable(encoding);
867 self
868 }
869 #[must_use]
873 pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
874 self.max_decoding_message_size = Some(limit);
875 self
876 }
877 #[must_use]
881 pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
882 self.max_encoding_message_size = Some(limit);
883 self
884 }
885 }
886 impl<T, B> tonic::codegen::Service<http::Request<B>> for ConnectorServiceServer<T>
887 where
888 T: ConnectorService,
889 B: Body + std::marker::Send + 'static,
890 B::Error: Into<StdError> + std::marker::Send + 'static,
891 {
892 type Response = http::Response<tonic::body::Body>;
893 type Error = std::convert::Infallible;
894 type Future = BoxFuture<Self::Response, Self::Error>;
895 fn poll_ready(
896 &mut self,
897 _cx: &mut Context<'_>,
898 ) -> Poll<std::result::Result<(), Self::Error>> {
899 Poll::Ready(Ok(()))
900 }
901 fn call(&mut self, req: http::Request<B>) -> Self::Future {
902 match req.uri().path() {
903 "/connector_service.ConnectorService/SinkWriterStream" => {
904 #[allow(non_camel_case_types)]
905 struct SinkWriterStreamSvc<T: ConnectorService>(pub Arc<T>);
906 impl<
907 T: ConnectorService,
908 > tonic::server::StreamingService<super::SinkWriterStreamRequest>
909 for SinkWriterStreamSvc<T> {
910 type Response = super::SinkWriterStreamResponse;
911 type ResponseStream = T::SinkWriterStreamStream;
912 type Future = BoxFuture<
913 tonic::Response<Self::ResponseStream>,
914 tonic::Status,
915 >;
916 fn call(
917 &mut self,
918 request: tonic::Request<
919 tonic::Streaming<super::SinkWriterStreamRequest>,
920 >,
921 ) -> Self::Future {
922 let inner = Arc::clone(&self.0);
923 let fut = async move {
924 <T as ConnectorService>::sink_writer_stream(&inner, request)
925 .await
926 };
927 Box::pin(fut)
928 }
929 }
930 let accept_compression_encodings = self.accept_compression_encodings;
931 let send_compression_encodings = self.send_compression_encodings;
932 let max_decoding_message_size = self.max_decoding_message_size;
933 let max_encoding_message_size = self.max_encoding_message_size;
934 let inner = self.inner.clone();
935 let fut = async move {
936 let method = SinkWriterStreamSvc(inner);
937 let codec = tonic_prost::ProstCodec::default();
938 let mut grpc = tonic::server::Grpc::new(codec)
939 .apply_compression_config(
940 accept_compression_encodings,
941 send_compression_encodings,
942 )
943 .apply_max_message_size_config(
944 max_decoding_message_size,
945 max_encoding_message_size,
946 );
947 let res = grpc.streaming(method, req).await;
948 Ok(res)
949 };
950 Box::pin(fut)
951 }
952 "/connector_service.ConnectorService/SinkCoordinatorStream" => {
953 #[allow(non_camel_case_types)]
954 struct SinkCoordinatorStreamSvc<T: ConnectorService>(pub Arc<T>);
955 impl<
956 T: ConnectorService,
957 > tonic::server::StreamingService<
958 super::SinkCoordinatorStreamRequest,
959 > for SinkCoordinatorStreamSvc<T> {
960 type Response = super::SinkCoordinatorStreamResponse;
961 type ResponseStream = T::SinkCoordinatorStreamStream;
962 type Future = BoxFuture<
963 tonic::Response<Self::ResponseStream>,
964 tonic::Status,
965 >;
966 fn call(
967 &mut self,
968 request: tonic::Request<
969 tonic::Streaming<super::SinkCoordinatorStreamRequest>,
970 >,
971 ) -> Self::Future {
972 let inner = Arc::clone(&self.0);
973 let fut = async move {
974 <T as ConnectorService>::sink_coordinator_stream(
975 &inner,
976 request,
977 )
978 .await
979 };
980 Box::pin(fut)
981 }
982 }
983 let accept_compression_encodings = self.accept_compression_encodings;
984 let send_compression_encodings = self.send_compression_encodings;
985 let max_decoding_message_size = self.max_decoding_message_size;
986 let max_encoding_message_size = self.max_encoding_message_size;
987 let inner = self.inner.clone();
988 let fut = async move {
989 let method = SinkCoordinatorStreamSvc(inner);
990 let codec = tonic_prost::ProstCodec::default();
991 let mut grpc = tonic::server::Grpc::new(codec)
992 .apply_compression_config(
993 accept_compression_encodings,
994 send_compression_encodings,
995 )
996 .apply_max_message_size_config(
997 max_decoding_message_size,
998 max_encoding_message_size,
999 );
1000 let res = grpc.streaming(method, req).await;
1001 Ok(res)
1002 };
1003 Box::pin(fut)
1004 }
1005 "/connector_service.ConnectorService/ValidateSink" => {
1006 #[allow(non_camel_case_types)]
1007 struct ValidateSinkSvc<T: ConnectorService>(pub Arc<T>);
1008 impl<
1009 T: ConnectorService,
1010 > tonic::server::UnaryService<super::ValidateSinkRequest>
1011 for ValidateSinkSvc<T> {
1012 type Response = super::ValidateSinkResponse;
1013 type Future = BoxFuture<
1014 tonic::Response<Self::Response>,
1015 tonic::Status,
1016 >;
1017 fn call(
1018 &mut self,
1019 request: tonic::Request<super::ValidateSinkRequest>,
1020 ) -> Self::Future {
1021 let inner = Arc::clone(&self.0);
1022 let fut = async move {
1023 <T as ConnectorService>::validate_sink(&inner, request)
1024 .await
1025 };
1026 Box::pin(fut)
1027 }
1028 }
1029 let accept_compression_encodings = self.accept_compression_encodings;
1030 let send_compression_encodings = self.send_compression_encodings;
1031 let max_decoding_message_size = self.max_decoding_message_size;
1032 let max_encoding_message_size = self.max_encoding_message_size;
1033 let inner = self.inner.clone();
1034 let fut = async move {
1035 let method = ValidateSinkSvc(inner);
1036 let codec = tonic_prost::ProstCodec::default();
1037 let mut grpc = tonic::server::Grpc::new(codec)
1038 .apply_compression_config(
1039 accept_compression_encodings,
1040 send_compression_encodings,
1041 )
1042 .apply_max_message_size_config(
1043 max_decoding_message_size,
1044 max_encoding_message_size,
1045 );
1046 let res = grpc.unary(method, req).await;
1047 Ok(res)
1048 };
1049 Box::pin(fut)
1050 }
1051 "/connector_service.ConnectorService/GetEventStream" => {
1052 #[allow(non_camel_case_types)]
1053 struct GetEventStreamSvc<T: ConnectorService>(pub Arc<T>);
1054 impl<
1055 T: ConnectorService,
1056 > tonic::server::ServerStreamingService<super::GetEventStreamRequest>
1057 for GetEventStreamSvc<T> {
1058 type Response = super::GetEventStreamResponse;
1059 type ResponseStream = T::GetEventStreamStream;
1060 type Future = BoxFuture<
1061 tonic::Response<Self::ResponseStream>,
1062 tonic::Status,
1063 >;
1064 fn call(
1065 &mut self,
1066 request: tonic::Request<super::GetEventStreamRequest>,
1067 ) -> Self::Future {
1068 let inner = Arc::clone(&self.0);
1069 let fut = async move {
1070 <T as ConnectorService>::get_event_stream(&inner, request)
1071 .await
1072 };
1073 Box::pin(fut)
1074 }
1075 }
1076 let accept_compression_encodings = self.accept_compression_encodings;
1077 let send_compression_encodings = self.send_compression_encodings;
1078 let max_decoding_message_size = self.max_decoding_message_size;
1079 let max_encoding_message_size = self.max_encoding_message_size;
1080 let inner = self.inner.clone();
1081 let fut = async move {
1082 let method = GetEventStreamSvc(inner);
1083 let codec = tonic_prost::ProstCodec::default();
1084 let mut grpc = tonic::server::Grpc::new(codec)
1085 .apply_compression_config(
1086 accept_compression_encodings,
1087 send_compression_encodings,
1088 )
1089 .apply_max_message_size_config(
1090 max_decoding_message_size,
1091 max_encoding_message_size,
1092 );
1093 let res = grpc.server_streaming(method, req).await;
1094 Ok(res)
1095 };
1096 Box::pin(fut)
1097 }
1098 "/connector_service.ConnectorService/ValidateSource" => {
1099 #[allow(non_camel_case_types)]
1100 struct ValidateSourceSvc<T: ConnectorService>(pub Arc<T>);
1101 impl<
1102 T: ConnectorService,
1103 > tonic::server::UnaryService<super::ValidateSourceRequest>
1104 for ValidateSourceSvc<T> {
1105 type Response = super::ValidateSourceResponse;
1106 type Future = BoxFuture<
1107 tonic::Response<Self::Response>,
1108 tonic::Status,
1109 >;
1110 fn call(
1111 &mut self,
1112 request: tonic::Request<super::ValidateSourceRequest>,
1113 ) -> Self::Future {
1114 let inner = Arc::clone(&self.0);
1115 let fut = async move {
1116 <T as ConnectorService>::validate_source(&inner, request)
1117 .await
1118 };
1119 Box::pin(fut)
1120 }
1121 }
1122 let accept_compression_encodings = self.accept_compression_encodings;
1123 let send_compression_encodings = self.send_compression_encodings;
1124 let max_decoding_message_size = self.max_decoding_message_size;
1125 let max_encoding_message_size = self.max_encoding_message_size;
1126 let inner = self.inner.clone();
1127 let fut = async move {
1128 let method = ValidateSourceSvc(inner);
1129 let codec = tonic_prost::ProstCodec::default();
1130 let mut grpc = tonic::server::Grpc::new(codec)
1131 .apply_compression_config(
1132 accept_compression_encodings,
1133 send_compression_encodings,
1134 )
1135 .apply_max_message_size_config(
1136 max_decoding_message_size,
1137 max_encoding_message_size,
1138 );
1139 let res = grpc.unary(method, req).await;
1140 Ok(res)
1141 };
1142 Box::pin(fut)
1143 }
1144 _ => {
1145 Box::pin(async move {
1146 let mut response = http::Response::new(
1147 tonic::body::Body::default(),
1148 );
1149 let headers = response.headers_mut();
1150 headers
1151 .insert(
1152 tonic::Status::GRPC_STATUS,
1153 (tonic::Code::Unimplemented as i32).into(),
1154 );
1155 headers
1156 .insert(
1157 http::header::CONTENT_TYPE,
1158 tonic::metadata::GRPC_CONTENT_TYPE,
1159 );
1160 Ok(response)
1161 })
1162 }
1163 }
1164 }
1165 }
1166 impl<T> Clone for ConnectorServiceServer<T> {
1167 fn clone(&self) -> Self {
1168 let inner = self.inner.clone();
1169 Self {
1170 inner,
1171 accept_compression_encodings: self.accept_compression_encodings,
1172 send_compression_encodings: self.send_compression_encodings,
1173 max_decoding_message_size: self.max_decoding_message_size,
1174 max_encoding_message_size: self.max_encoding_message_size,
1175 }
1176 }
1177 }
1178 pub const SERVICE_NAME: &str = "connector_service.ConnectorService";
1180 impl<T> tonic::server::NamedService for ConnectorServiceServer<T> {
1181 const NAME: &'static str = SERVICE_NAME;
1182 }
1183}
1184pub mod sink_coordination_service_client {
1186 #![allow(
1187 unused_variables,
1188 dead_code,
1189 missing_docs,
1190 clippy::wildcard_imports,
1191 clippy::let_unit_value,
1192 )]
1193 use tonic::codegen::*;
1194 use tonic::codegen::http::Uri;
1195 #[derive(Debug, Clone)]
1196 pub struct SinkCoordinationServiceClient<T> {
1197 inner: tonic::client::Grpc<T>,
1198 }
1199 impl SinkCoordinationServiceClient<tonic::transport::Channel> {
1200 pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
1202 where
1203 D: TryInto<tonic::transport::Endpoint>,
1204 D::Error: Into<StdError>,
1205 {
1206 let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
1207 Ok(Self::new(conn))
1208 }
1209 }
1210 impl<T> SinkCoordinationServiceClient<T>
1211 where
1212 T: tonic::client::GrpcService<tonic::body::Body>,
1213 T::Error: Into<StdError>,
1214 T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
1215 <T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
1216 {
1217 pub fn new(inner: T) -> Self {
1218 let inner = tonic::client::Grpc::new(inner);
1219 Self { inner }
1220 }
1221 pub fn with_origin(inner: T, origin: Uri) -> Self {
1222 let inner = tonic::client::Grpc::with_origin(inner, origin);
1223 Self { inner }
1224 }
1225 pub fn with_interceptor<F>(
1226 inner: T,
1227 interceptor: F,
1228 ) -> SinkCoordinationServiceClient<InterceptedService<T, F>>
1229 where
1230 F: tonic::service::Interceptor,
1231 T::ResponseBody: Default,
1232 T: tonic::codegen::Service<
1233 http::Request<tonic::body::Body>,
1234 Response = http::Response<
1235 <T as tonic::client::GrpcService<tonic::body::Body>>::ResponseBody,
1236 >,
1237 >,
1238 <T as tonic::codegen::Service<
1239 http::Request<tonic::body::Body>,
1240 >>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
1241 {
1242 SinkCoordinationServiceClient::new(
1243 InterceptedService::new(inner, interceptor),
1244 )
1245 }
1246 #[must_use]
1251 pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
1252 self.inner = self.inner.send_compressed(encoding);
1253 self
1254 }
1255 #[must_use]
1257 pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
1258 self.inner = self.inner.accept_compressed(encoding);
1259 self
1260 }
1261 #[must_use]
1265 pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
1266 self.inner = self.inner.max_decoding_message_size(limit);
1267 self
1268 }
1269 #[must_use]
1273 pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
1274 self.inner = self.inner.max_encoding_message_size(limit);
1275 self
1276 }
1277 pub async fn coordinate(
1278 &mut self,
1279 request: impl tonic::IntoStreamingRequest<Message = super::CoordinateRequest>,
1280 ) -> std::result::Result<
1281 tonic::Response<tonic::codec::Streaming<super::CoordinateResponse>>,
1282 tonic::Status,
1283 > {
1284 self.inner
1285 .ready()
1286 .await
1287 .map_err(|e| {
1288 tonic::Status::unknown(
1289 format!("Service was not ready: {}", e.into()),
1290 )
1291 })?;
1292 let codec = tonic_prost::ProstCodec::default();
1293 let path = http::uri::PathAndQuery::from_static(
1294 "/connector_service.SinkCoordinationService/Coordinate",
1295 );
1296 let mut req = request.into_streaming_request();
1297 req.extensions_mut()
1298 .insert(
1299 GrpcMethod::new(
1300 "connector_service.SinkCoordinationService",
1301 "Coordinate",
1302 ),
1303 );
1304 self.inner.streaming(req, path, codec).await
1305 }
1306 }
1307}
1308pub mod sink_coordination_service_server {
1310 #![allow(
1311 unused_variables,
1312 dead_code,
1313 missing_docs,
1314 clippy::wildcard_imports,
1315 clippy::let_unit_value,
1316 )]
1317 use tonic::codegen::*;
1318 #[async_trait]
1320 pub trait SinkCoordinationService: std::marker::Send + std::marker::Sync + 'static {
1321 type CoordinateStream: tonic::codegen::tokio_stream::Stream<
1323 Item = std::result::Result<super::CoordinateResponse, tonic::Status>,
1324 >
1325 + std::marker::Send
1326 + 'static;
1327 async fn coordinate(
1328 &self,
1329 request: tonic::Request<tonic::Streaming<super::CoordinateRequest>>,
1330 ) -> std::result::Result<tonic::Response<Self::CoordinateStream>, tonic::Status>;
1331 }
1332 #[derive(Debug)]
1333 pub struct SinkCoordinationServiceServer<T> {
1334 inner: Arc<T>,
1335 accept_compression_encodings: EnabledCompressionEncodings,
1336 send_compression_encodings: EnabledCompressionEncodings,
1337 max_decoding_message_size: Option<usize>,
1338 max_encoding_message_size: Option<usize>,
1339 }
1340 impl<T> SinkCoordinationServiceServer<T> {
1341 pub fn new(inner: T) -> Self {
1342 Self::from_arc(Arc::new(inner))
1343 }
1344 pub fn from_arc(inner: Arc<T>) -> Self {
1345 Self {
1346 inner,
1347 accept_compression_encodings: Default::default(),
1348 send_compression_encodings: Default::default(),
1349 max_decoding_message_size: None,
1350 max_encoding_message_size: None,
1351 }
1352 }
1353 pub fn with_interceptor<F>(
1354 inner: T,
1355 interceptor: F,
1356 ) -> InterceptedService<Self, F>
1357 where
1358 F: tonic::service::Interceptor,
1359 {
1360 InterceptedService::new(Self::new(inner), interceptor)
1361 }
1362 #[must_use]
1364 pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
1365 self.accept_compression_encodings.enable(encoding);
1366 self
1367 }
1368 #[must_use]
1370 pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
1371 self.send_compression_encodings.enable(encoding);
1372 self
1373 }
1374 #[must_use]
1378 pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
1379 self.max_decoding_message_size = Some(limit);
1380 self
1381 }
1382 #[must_use]
1386 pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
1387 self.max_encoding_message_size = Some(limit);
1388 self
1389 }
1390 }
1391 impl<T, B> tonic::codegen::Service<http::Request<B>>
1392 for SinkCoordinationServiceServer<T>
1393 where
1394 T: SinkCoordinationService,
1395 B: Body + std::marker::Send + 'static,
1396 B::Error: Into<StdError> + std::marker::Send + 'static,
1397 {
1398 type Response = http::Response<tonic::body::Body>;
1399 type Error = std::convert::Infallible;
1400 type Future = BoxFuture<Self::Response, Self::Error>;
1401 fn poll_ready(
1402 &mut self,
1403 _cx: &mut Context<'_>,
1404 ) -> Poll<std::result::Result<(), Self::Error>> {
1405 Poll::Ready(Ok(()))
1406 }
1407 fn call(&mut self, req: http::Request<B>) -> Self::Future {
1408 match req.uri().path() {
1409 "/connector_service.SinkCoordinationService/Coordinate" => {
1410 #[allow(non_camel_case_types)]
1411 struct CoordinateSvc<T: SinkCoordinationService>(pub Arc<T>);
1412 impl<
1413 T: SinkCoordinationService,
1414 > tonic::server::StreamingService<super::CoordinateRequest>
1415 for CoordinateSvc<T> {
1416 type Response = super::CoordinateResponse;
1417 type ResponseStream = T::CoordinateStream;
1418 type Future = BoxFuture<
1419 tonic::Response<Self::ResponseStream>,
1420 tonic::Status,
1421 >;
1422 fn call(
1423 &mut self,
1424 request: tonic::Request<
1425 tonic::Streaming<super::CoordinateRequest>,
1426 >,
1427 ) -> Self::Future {
1428 let inner = Arc::clone(&self.0);
1429 let fut = async move {
1430 <T as SinkCoordinationService>::coordinate(&inner, request)
1431 .await
1432 };
1433 Box::pin(fut)
1434 }
1435 }
1436 let accept_compression_encodings = self.accept_compression_encodings;
1437 let send_compression_encodings = self.send_compression_encodings;
1438 let max_decoding_message_size = self.max_decoding_message_size;
1439 let max_encoding_message_size = self.max_encoding_message_size;
1440 let inner = self.inner.clone();
1441 let fut = async move {
1442 let method = CoordinateSvc(inner);
1443 let codec = tonic_prost::ProstCodec::default();
1444 let mut grpc = tonic::server::Grpc::new(codec)
1445 .apply_compression_config(
1446 accept_compression_encodings,
1447 send_compression_encodings,
1448 )
1449 .apply_max_message_size_config(
1450 max_decoding_message_size,
1451 max_encoding_message_size,
1452 );
1453 let res = grpc.streaming(method, req).await;
1454 Ok(res)
1455 };
1456 Box::pin(fut)
1457 }
1458 _ => {
1459 Box::pin(async move {
1460 let mut response = http::Response::new(
1461 tonic::body::Body::default(),
1462 );
1463 let headers = response.headers_mut();
1464 headers
1465 .insert(
1466 tonic::Status::GRPC_STATUS,
1467 (tonic::Code::Unimplemented as i32).into(),
1468 );
1469 headers
1470 .insert(
1471 http::header::CONTENT_TYPE,
1472 tonic::metadata::GRPC_CONTENT_TYPE,
1473 );
1474 Ok(response)
1475 })
1476 }
1477 }
1478 }
1479 }
1480 impl<T> Clone for SinkCoordinationServiceServer<T> {
1481 fn clone(&self) -> Self {
1482 let inner = self.inner.clone();
1483 Self {
1484 inner,
1485 accept_compression_encodings: self.accept_compression_encodings,
1486 send_compression_encodings: self.send_compression_encodings,
1487 max_decoding_message_size: self.max_decoding_message_size,
1488 max_encoding_message_size: self.max_encoding_message_size,
1489 }
1490 }
1491 }
1492 pub const SERVICE_NAME: &str = "connector_service.SinkCoordinationService";
1494 impl<T> tonic::server::NamedService for SinkCoordinationServiceServer<T> {
1495 const NAME: &'static str = SERVICE_NAME;
1496 }
1497}