Skip to main content

risingwave_pb/
connector_service.rs

1// This file is @generated by prost-build.
2#[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    /// to be deprecated
29    #[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    /// Backward-compatible replacement for `SINK_TYPE_FORCE_APPEND_ONLY`.
40    ///
41    /// Should not directly access this field. Use method `ignore_delete()` instead.
42    #[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}
51/// Nested message and enum types in `SinkWriterStreamRequest`.
52pub 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    /// Nested message and enum types in `WriteBatch`.
72    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            /// This is a reference pointer to a StreamChunk. The StreamChunk is owned
85            /// by the JniSinkWriterStreamRequest, which should handle the release of StreamChunk.
86            /// Index set to 5 because 3 and 4 have been occupied by `batch_id` and `epoch`
87            #[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}
116/// Nested message and enum types in `SinkWriterStreamResponse`.
117pub 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    /// On validation failure, we return the error.
158    #[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}
167/// Nested message and enum types in `SinkMetadata`.
168pub 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}
188/// Nested message and enum types in `SinkCoordinatorStreamRequest`.
189pub 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}
219/// Nested message and enum types in `SinkCoordinatorStreamResponse`.
220pub 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    /// The value of the Debezium message
243    #[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    /// The key of the Debezium message, which only used by `mongodb-cdc` connector.
256    #[prost(string, tag = "7")]
257    pub key: ::prost::alloc::string::String,
258    /// The CDC source type of this message. This is required to disambiguate the parsing logic for
259    /// `full_table_name` (e.g. `db.table` in MySQL vs `schema.table` in
260    /// Postgres/SQL Server/Oracle).
261    #[prost(enumeration = "SourceType", tag = "8")]
262    pub source_type: i32,
263}
264/// Nested message and enum types in `CdcMessage`.
265pub 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        /// String value of the enum field names used in the ProtoBuf definition.
288        ///
289        /// The values are not transformed in any way and thus are considered stable
290        /// (if the ProtoBuf definition does not change) and safe for programmatic use.
291        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        /// Creates an enum from field names used in the ProtoBuf definition.
301        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}
342/// Nested message and enum types in `GetEventStreamResponse`.
343pub 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    /// On validation failure, we return the error.
374    #[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}
383/// Nested message and enum types in `CoordinateRequest`.
384pub mod coordinate_request {
385    /// The first request that starts a coordination between sink writer and coordinator.
386    /// The service will respond after sink writers of all vnodes have sent the request.
387    #[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}
434/// Nested message and enum types in `CoordinateResponse`.
435pub 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    /// String value of the enum field names used in the ProtoBuf definition.
483    ///
484    /// The values are not transformed in any way and thus are considered stable
485    /// (if the ProtoBuf definition does not change) and safe for programmatic use.
486    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    /// Creates an enum from field names used in the ProtoBuf definition.
498    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}
511/// Generated client implementations.
512pub 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        /// Attempt to create a new client by connecting to a given endpoint.
528        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        /// Compress requests with the given encoding.
572        ///
573        /// This requires the server to support it otherwise it might respond with an
574        /// error.
575        #[must_use]
576        pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
577            self.inner = self.inner.send_compressed(encoding);
578            self
579        }
580        /// Enable decompressing responses.
581        #[must_use]
582        pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
583            self.inner = self.inner.accept_compressed(encoding);
584            self
585        }
586        /// Limits the maximum size of a decoded message.
587        ///
588        /// Default: `4MB`
589        #[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        /// Limits the maximum size of an encoded message.
595        ///
596        /// Default: `usize::MAX`
597        #[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}
752/// Generated server implementations.
753pub 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    /// Generated trait containing gRPC methods that should be implemented for use with ConnectorServiceServer.
763    #[async_trait]
764    pub trait ConnectorService: std::marker::Send + std::marker::Sync + 'static {
765        /// Server streaming response type for the SinkWriterStream method.
766        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        /// Server streaming response type for the SinkCoordinatorStream method.
782        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        /// Server streaming response type for the GetEventStream method.
807        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        /// Enable decompressing requests with the given encoding.
858        #[must_use]
859        pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
860            self.accept_compression_encodings.enable(encoding);
861            self
862        }
863        /// Compress responses with the given encoding, if the client supports it.
864        #[must_use]
865        pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
866            self.send_compression_encodings.enable(encoding);
867            self
868        }
869        /// Limits the maximum size of a decoded message.
870        ///
871        /// Default: `4MB`
872        #[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        /// Limits the maximum size of an encoded message.
878        ///
879        /// Default: `usize::MAX`
880        #[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    /// Generated gRPC service name
1179    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}
1184/// Generated client implementations.
1185pub 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        /// Attempt to create a new client by connecting to a given endpoint.
1201        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        /// Compress requests with the given encoding.
1247        ///
1248        /// This requires the server to support it otherwise it might respond with an
1249        /// error.
1250        #[must_use]
1251        pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
1252            self.inner = self.inner.send_compressed(encoding);
1253            self
1254        }
1255        /// Enable decompressing responses.
1256        #[must_use]
1257        pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
1258            self.inner = self.inner.accept_compressed(encoding);
1259            self
1260        }
1261        /// Limits the maximum size of a decoded message.
1262        ///
1263        /// Default: `4MB`
1264        #[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        /// Limits the maximum size of an encoded message.
1270        ///
1271        /// Default: `usize::MAX`
1272        #[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}
1308/// Generated server implementations.
1309pub 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    /// Generated trait containing gRPC methods that should be implemented for use with SinkCoordinationServiceServer.
1319    #[async_trait]
1320    pub trait SinkCoordinationService: std::marker::Send + std::marker::Sync + 'static {
1321        /// Server streaming response type for the Coordinate method.
1322        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        /// Enable decompressing requests with the given encoding.
1363        #[must_use]
1364        pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
1365            self.accept_compression_encodings.enable(encoding);
1366            self
1367        }
1368        /// Compress responses with the given encoding, if the client supports it.
1369        #[must_use]
1370        pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
1371            self.send_compression_encodings.enable(encoding);
1372            self
1373        }
1374        /// Limits the maximum size of a decoded message.
1375        ///
1376        /// Default: `4MB`
1377        #[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        /// Limits the maximum size of an encoded message.
1383        ///
1384        /// Default: `usize::MAX`
1385        #[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    /// Generated gRPC service name
1493    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}