Skip to main content

temporalio_client/
grpc.rs

1//! gRPC service traits for direct access to Temporal services.
2//!
3//! Most users should use the higher-level methods on [`Client`] or [`Connection`] instead.
4//! These traits are useful for advanced scenarios like custom interceptors, testing with mocks,
5//! or making raw gRPC calls not covered by the higher-level API.
6
7use crate::{
8    Client, Connection, LONG_POLL_TIMEOUT, PayloadErrorLimits, RequestExt, SharedReplaceableClient,
9    TEMPORAL_NAMESPACE_HEADER_KEY, TemporalServiceClient,
10    metrics::namespace_kv,
11    retry::make_future_retry,
12    worker::{ClientWorkerSet, Slot},
13};
14use dyn_clone::DynClone;
15use futures_util::{FutureExt, TryFutureExt, future::BoxFuture};
16use parking_lot::RwLock;
17use std::{any::Any, marker::PhantomData, sync::Arc};
18use temporalio_common::{
19    payload_limits::{PayloadLimits, validate_known_payload_limits},
20    protos::{
21        grpc::health::v1::{health_client::HealthClient, *},
22        temporal::api::{
23            cloud::cloudservice::{v1 as cloudreq, v1::cloud_service_client::CloudServiceClient},
24            operatorservice::v1::{operator_service_client::OperatorServiceClient, *},
25            taskqueue::v1::TaskQueue,
26            testservice::v1::{test_service_client::TestServiceClient, *},
27            workflowservice::v1::{workflow_service_client::WorkflowServiceClient, *},
28        },
29    },
30    telemetry::metrics::MetricKeyValue,
31};
32use tonic::{
33    Request, Response, Status,
34    body::Body,
35    client::GrpcService,
36    metadata::{AsciiMetadataValue, KeyAndValueRef},
37};
38
39/// Something that has access to the raw grpc services
40pub(crate) trait RawClientProducer {
41    /// Returns information about workers associated with this client. Implementers outside of
42    /// core can safely return `None`.
43    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>>;
44
45    /// Return a workflow service client instance
46    fn workflow_client(&mut self) -> Box<dyn WorkflowService>;
47
48    /// Return a mutable ref to the operator service client instance
49    fn operator_client(&mut self) -> Box<dyn OperatorService>;
50
51    /// Return a mutable ref to the cloud service client instance
52    fn cloud_client(&mut self) -> Box<dyn CloudService>;
53
54    /// Return a mutable ref to the test service client instance
55    fn test_client(&mut self) -> Box<dyn TestService>;
56
57    /// Return a mutable ref to the health service client instance
58    fn health_client(&mut self) -> Box<dyn HealthService>;
59}
60
61/// Any client that can make gRPC calls. The default implementation simply invokes the passed-in
62/// function. Implementers may override this to provide things like retry behavior.
63#[async_trait::async_trait]
64pub(crate) trait RawGrpcCaller: Send + Sync + 'static {
65    /// Make a gRPC call. The default implementation simply invokes the provided function.
66    async fn call<F, Req, Resp>(
67        &mut self,
68        _call_name: &'static str,
69        mut callfn: F,
70        req: Request<Req>,
71    ) -> Result<Response<Resp>, Status>
72    where
73        Req: Clone + Unpin + Send + Sync + 'static,
74        Resp: Send + 'static,
75        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
76        F: Send + Sync + Unpin + 'static,
77    {
78        callfn(req).await
79    }
80}
81
82trait ErasedRawClient: Send + Sync + 'static {
83    fn erased_call(
84        &mut self,
85        call_name: &'static str,
86        op: &mut dyn ErasedCallOp,
87    ) -> BoxFuture<'static, Result<Response<Box<dyn Any + Send>>, Status>>;
88}
89
90trait ErasedCallOp: Send {
91    fn invoke(
92        &mut self,
93        raw: &mut dyn ErasedRawClient,
94        call_name: &'static str,
95    ) -> BoxFuture<'static, Result<Response<Box<dyn Any + Send>>, Status>>;
96}
97
98struct CallShim<F, Req, Resp> {
99    callfn: F,
100    seed_req: Option<Request<Req>>,
101    _resp: PhantomData<Resp>,
102}
103
104impl<F, Req, Resp> CallShim<F, Req, Resp> {
105    fn new(callfn: F, seed_req: Request<Req>) -> Self {
106        Self {
107            callfn,
108            seed_req: Some(seed_req),
109            _resp: PhantomData,
110        }
111    }
112}
113impl<F, Req, Resp> ErasedCallOp for CallShim<F, Req, Resp>
114where
115    Req: Clone + Unpin + Send + Sync + 'static,
116    Resp: Send + 'static,
117    F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
118    F: Send + Sync + Unpin + 'static,
119{
120    fn invoke(
121        &mut self,
122        _raw: &mut dyn ErasedRawClient,
123        _call_name: &'static str,
124    ) -> BoxFuture<'static, Result<Response<Box<dyn Any + Send>>, Status>> {
125        (self.callfn)(
126            self.seed_req
127                .take()
128                .expect("CallShim must have request populated"),
129        )
130        .map(|res| res.map(|payload| payload.map(|t| Box::new(t) as Box<dyn Any + Send>)))
131        .boxed()
132    }
133}
134
135#[async_trait::async_trait]
136impl RawGrpcCaller for Connection {
137    async fn call<F, Req, Resp>(
138        &mut self,
139        call_name: &'static str,
140        mut callfn: F,
141        mut req: Request<Req>,
142    ) -> Result<Response<Resp>, Status>
143    where
144        Req: Clone + Unpin + Send + Sync + 'static,
145        Resp: Send + 'static,
146        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
147        F: Send + Sync + Unpin + 'static,
148    {
149        // Validate payload sizes after any request mutation but before encoding/metrics.
150        validate_request_payload_limits(
151            &req,
152            self.inner.payloads_warn_size,
153            self.inner.memo_warn_size,
154        )?;
155
156        let info = self
157            .inner
158            .retry_options
159            .get_call_info(call_name, Some(&req));
160        req.extensions_mut().insert(info.call_type);
161        if info.call_type.is_long() {
162            req.set_default_timeout(LONG_POLL_TIMEOUT);
163        }
164
165        let fact = || {
166            let req_clone = req_cloner(&req);
167            callfn(req_clone)
168        };
169
170        let res = make_future_retry(info, fact);
171        res.map_err(|(e, _attempt)| e).map_ok(|x| x.0).await
172    }
173}
174
175/// Helper for cloning a tonic request as long as the inner message may be cloned.
176fn req_cloner<T: Clone>(cloneme: &Request<T>) -> Request<T> {
177    let msg = cloneme.get_ref().clone();
178    let mut new_req = Request::new(msg);
179    let new_met = new_req.metadata_mut();
180    for kv in cloneme.metadata().iter() {
181        match kv {
182            KeyAndValueRef::Ascii(k, v) => {
183                new_met.insert(k, v.clone());
184            }
185            KeyAndValueRef::Binary(k, v) => {
186                new_met.insert_bin(k, v.clone());
187            }
188        }
189    }
190    *new_req.extensions_mut() = cloneme.extensions().clone();
191    new_req
192}
193
194/// `*_warn` are the connection's configured warn thresholds; per-call error limits ride a
195/// [`PayloadErrorLimits`] extension. On an error-level violation, returns a [`Status`] carrying
196/// the [`PayloadLimitViolation`] as its source (extract via [crate::payload_limit_violation_from]).
197fn validate_request_payload_limits<Req: Any>(
198    req: &Request<Req>,
199    blob_warn: usize,
200    memo_warn: usize,
201) -> Result<(), Status> {
202    let mut limits = PayloadLimits {
203        blob_warn,
204        memo_warn,
205        blob_error: 0,
206        memo_error: 0,
207    };
208    if let Some(error_limits) = req.extensions().get::<PayloadErrorLimits>() {
209        limits.blob_error = error_limits.blob;
210        limits.memo_error = error_limits.memo;
211    }
212    if let Some(violation) = validate_known_payload_limits(req.get_ref(), &limits) {
213        let mut status = Status::invalid_argument(violation.to_string());
214        status.set_source(Arc::new(violation));
215        return Err(status);
216    }
217    Ok(())
218}
219
220#[async_trait::async_trait]
221impl RawGrpcCaller for dyn ErasedRawClient {
222    async fn call<F, Req, Resp>(
223        &mut self,
224        call_name: &'static str,
225        callfn: F,
226        req: Request<Req>,
227    ) -> Result<Response<Resp>, Status>
228    where
229        Req: Clone + Unpin + Send + Sync + 'static,
230        Resp: Send + 'static,
231        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
232        F: Send + Sync + Unpin + 'static,
233    {
234        let mut shim = CallShim::new(callfn, req);
235        let erased_resp = ErasedRawClient::erased_call(self, call_name, &mut shim).await?;
236        Ok(erased_resp.map(|boxed| {
237            *boxed
238                .downcast()
239                .expect("RawGrpcCaller erased response type mismatch")
240        }))
241    }
242}
243
244impl<T> ErasedRawClient for T
245where
246    T: RawGrpcCaller + 'static,
247{
248    fn erased_call(
249        &mut self,
250        call_name: &'static str,
251        op: &mut dyn ErasedCallOp,
252    ) -> BoxFuture<'static, Result<Response<Box<dyn Any + Send>>, Status>> {
253        let raw: &mut dyn ErasedRawClient = self;
254        op.invoke(raw, call_name)
255    }
256}
257
258impl RawClientProducer for Connection {
259    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
260        Some(self.workers())
261    }
262
263    fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
264        self.inner.service.workflow_service()
265    }
266
267    fn operator_client(&mut self) -> Box<dyn OperatorService> {
268        self.inner.service.operator_service()
269    }
270
271    fn cloud_client(&mut self) -> Box<dyn CloudService> {
272        self.inner.service.cloud_service()
273    }
274
275    fn test_client(&mut self) -> Box<dyn TestService> {
276        self.inner.service.test_service()
277    }
278
279    fn health_client(&mut self) -> Box<dyn HealthService> {
280        self.inner.service.health_service()
281    }
282}
283
284impl RawClientProducer for TemporalServiceClient {
285    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
286        None
287    }
288
289    fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
290        self.workflow_service()
291    }
292
293    fn operator_client(&mut self) -> Box<dyn OperatorService> {
294        self.operator_service()
295    }
296
297    fn cloud_client(&mut self) -> Box<dyn CloudService> {
298        self.cloud_service()
299    }
300
301    fn test_client(&mut self) -> Box<dyn TestService> {
302        self.test_service()
303    }
304
305    fn health_client(&mut self) -> Box<dyn HealthService> {
306        self.health_service()
307    }
308}
309
310impl RawGrpcCaller for TemporalServiceClient {}
311
312impl RawClientProducer for Client {
313    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
314        Some(self.connection.workers())
315    }
316
317    fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
318        self.connection.workflow_client()
319    }
320
321    fn operator_client(&mut self) -> Box<dyn OperatorService> {
322        self.connection.operator_client()
323    }
324
325    fn cloud_client(&mut self) -> Box<dyn CloudService> {
326        self.connection.cloud_client()
327    }
328
329    fn test_client(&mut self) -> Box<dyn TestService> {
330        self.connection.test_client()
331    }
332
333    fn health_client(&mut self) -> Box<dyn HealthService> {
334        self.connection.health_client()
335    }
336}
337
338#[async_trait::async_trait]
339impl RawGrpcCaller for Client {
340    async fn call<F, Req, Resp>(
341        &mut self,
342        call_name: &'static str,
343        callfn: F,
344        req: Request<Req>,
345    ) -> Result<Response<Resp>, Status>
346    where
347        Req: Clone + Unpin + Send + Sync + 'static,
348        Resp: Send + 'static,
349        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
350        F: Send + Sync + Unpin + 'static,
351    {
352        self.connection.call(call_name, callfn, req).await
353    }
354}
355
356impl<RC> RawClientProducer for SharedReplaceableClient<RC>
357where
358    RC: RawClientProducer + Clone + Send + Sync + 'static,
359{
360    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
361        self.inner_cow().get_workers_info()
362    }
363    fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
364        self.inner_mut_refreshed().workflow_client()
365    }
366
367    fn operator_client(&mut self) -> Box<dyn OperatorService> {
368        self.inner_mut_refreshed().operator_client()
369    }
370
371    fn cloud_client(&mut self) -> Box<dyn CloudService> {
372        self.inner_mut_refreshed().cloud_client()
373    }
374
375    fn test_client(&mut self) -> Box<dyn TestService> {
376        self.inner_mut_refreshed().test_client()
377    }
378
379    fn health_client(&mut self) -> Box<dyn HealthService> {
380        self.inner_mut_refreshed().health_client()
381    }
382}
383
384#[async_trait::async_trait]
385impl<RC> RawGrpcCaller for SharedReplaceableClient<RC>
386where
387    RC: RawGrpcCaller + Clone + Sync + 'static,
388{
389    async fn call<F, Req, Resp>(
390        &mut self,
391        call_name: &'static str,
392        callfn: F,
393        req: Request<Req>,
394    ) -> Result<Response<Resp>, Status>
395    where
396        Req: Clone + Unpin + Send + Sync + 'static,
397        Resp: Send + 'static,
398        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
399        F: Send + Sync + Unpin + 'static,
400    {
401        self.inner_mut_refreshed()
402            .call(call_name, callfn, req)
403            .await
404    }
405}
406
407/// Wraps a client and injects a worker's configured [`PayloadErrorLimits`] (when set) as a request
408/// extension on every gRPC call, so the gRPC layer enforces error limits uniformly across all
409/// outbound requests — current and future — without each call site having to opt in.
410///
411/// The limits are shared and may be updated after construction (e.g. once a worker learns the
412/// namespace limits) via [`set_error_limits`](Self::set_error_limits); clones observe the update.
413#[doc(hidden)]
414#[derive(Clone)]
415pub struct PayloadLimitsClient<C> {
416    inner: C,
417    error_limits: Arc<RwLock<Option<PayloadErrorLimits>>>,
418}
419
420impl<C> PayloadLimitsClient<C> {
421    /// Wrap `inner`; no limits are enforced until [`set_error_limits`](Self::set_error_limits).
422    pub fn new(inner: C) -> Self {
423        Self {
424            inner,
425            error_limits: Arc::new(RwLock::new(None)),
426        }
427    }
428
429    /// Set (or clear) the error limits injected on every call. Shared across clones.
430    pub fn set_error_limits(&self, limits: Option<PayloadErrorLimits>) {
431        *self.error_limits.write() = limits;
432    }
433
434    /// Return the error limits currently injected into requests.
435    pub fn error_limits(&self) -> Option<PayloadErrorLimits> {
436        *self.error_limits.read()
437    }
438}
439
440impl<C: RawClientProducer> RawClientProducer for PayloadLimitsClient<C> {
441    fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
442        self.inner.get_workers_info()
443    }
444    fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
445        self.inner.workflow_client()
446    }
447    fn operator_client(&mut self) -> Box<dyn OperatorService> {
448        self.inner.operator_client()
449    }
450    fn cloud_client(&mut self) -> Box<dyn CloudService> {
451        self.inner.cloud_client()
452    }
453    fn test_client(&mut self) -> Box<dyn TestService> {
454        self.inner.test_client()
455    }
456    fn health_client(&mut self) -> Box<dyn HealthService> {
457        self.inner.health_client()
458    }
459}
460
461#[async_trait::async_trait]
462impl<C> RawGrpcCaller for PayloadLimitsClient<C>
463where
464    C: RawGrpcCaller + Clone + Sync + 'static,
465{
466    async fn call<F, Req, Resp>(
467        &mut self,
468        call_name: &'static str,
469        callfn: F,
470        mut req: Request<Req>,
471    ) -> Result<Response<Resp>, Status>
472    where
473        Req: Clone + Unpin + Send + Sync + 'static,
474        Resp: Send + 'static,
475        F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
476        F: Send + Sync + Unpin + 'static,
477    {
478        let limits = *self.error_limits.read();
479        if let Some(limits) = limits {
480            req.extensions_mut().insert(limits);
481        }
482        self.inner.call(call_name, callfn, req).await
483    }
484}
485
486#[derive(Clone, Debug)]
487pub(super) struct AttachMetricLabels {
488    pub(super) labels: Vec<MetricKeyValue>,
489    pub(super) normal_task_queue: Option<String>,
490    pub(super) sticky_task_queue: Option<String>,
491}
492impl AttachMetricLabels {
493    pub(super) fn new(kvs: impl Into<Vec<MetricKeyValue>>) -> Self {
494        Self {
495            labels: kvs.into(),
496            normal_task_queue: None,
497            sticky_task_queue: None,
498        }
499    }
500    pub(super) fn namespace(ns: impl Into<String>) -> Self {
501        AttachMetricLabels::new(vec![namespace_kv(ns.into())])
502    }
503    pub(super) fn task_q(&mut self, tq: Option<TaskQueue>) -> &mut Self {
504        if let Some(tq) = tq {
505            if !tq.normal_name.is_empty() {
506                self.sticky_task_queue = Some(tq.name);
507                self.normal_task_queue = Some(tq.normal_name);
508            } else {
509                self.normal_task_queue = Some(tq.name);
510            }
511        }
512        self
513    }
514    pub(super) fn task_q_str(&mut self, tq: impl Into<String>) -> &mut Self {
515        self.normal_task_queue = Some(tq.into());
516        self
517    }
518}
519
520/// A request extension that, when set, should make the [RetryClient] consider this call to be a
521/// [super::retry::CallType::UserLongPoll]
522#[derive(Copy, Clone, Debug)]
523pub(super) struct IsUserLongPoll;
524
525macro_rules! proxy_def {
526    ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty, defaults) => {
527        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
528        fn $method(
529            &mut self,
530            _request: tonic::Request<$req>,
531        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
532            async { Ok(tonic::Response::new(<$resp>::default())) }.boxed()
533        }
534    };
535    ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty) => {
536        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
537        fn $method(
538            &mut self,
539            _request: tonic::Request<$req>,
540        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>>;
541    };
542}
543
544/// Helps re-declare gRPC client methods
545///
546/// There are four forms:
547///
548/// * The first takes a closure that can modify the request. This is only called once, before the
549///   actual rpc call is made, and before determinations are made about the kind of call (long poll
550///   or not) and retry policy.
551/// * The second takes three closures. The first can modify the request like in the first form.
552///   The second can modify the request and return a value, and is called right before every call
553///   (including on retries). The third is called with the response to the call after it resolves.
554/// * The third and fourth are equivalents of the above that skip calling through the `call` method
555///   and are implemented directly on the generated gRPC clients (IE: the bottom of the stack).
556macro_rules! proxy_impl {
557    ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty $(, $closure:expr)?) => {
558        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
559        fn $method(
560            &mut self,
561            #[allow(unused_mut)]
562            mut request: tonic::Request<$req>,
563        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
564            $( type_closure_arg(&mut request, $closure); )*
565            let mut self_clone = self.clone();
566            #[allow(unused_mut)]
567            let fact = move |mut req: tonic::Request<$req>| {
568                let mut c = self_clone.$client_meth();
569                async move { c.$method(req).await }.boxed()
570            };
571            self.call(stringify!($method), fact, request)
572        }
573    };
574    ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty,
575     $closure_request:expr, $closure_before:expr, $closure_after:expr) => {
576        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
577        fn $method(
578            &mut self,
579            mut request: tonic::Request<$req>,
580        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
581            type_closure_arg(&mut request, $closure_request);
582            let workers_info = self.get_workers_info();
583            let mut self_clone = self.clone();
584            #[allow(unused_mut)]
585            let fact = move |mut req: tonic::Request<$req>| {
586                let data = type_closure_two_arg(&mut req, workers_info.clone(), $closure_before);
587                let mut c = self_clone.$client_meth();
588                async move {
589                    type_closure_two_arg(c.$method(req).await, data, $closure_after)
590                }.boxed()
591            };
592            self.call(stringify!($method), fact, request)
593        }
594    };
595    ($client_type:tt, $method:ident, $req:ty, $resp:ty $(, $closure:expr)?) => {
596        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
597        fn $method(
598            &mut self,
599            #[allow(unused_mut)]
600            mut request: tonic::Request<$req>,
601        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
602            $( type_closure_arg(&mut request, $closure); )*
603            async move { <$client_type<_>>::$method(self, request).await }.boxed()
604        }
605    };
606    ($client_type:tt, $method:ident, $req:ty, $resp:ty,
607     $closure_request:expr, $closure_before:expr, $closure_after:expr) => {
608        #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
609        fn $method(
610            &mut self,
611            mut request: tonic::Request<$req>,
612        ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
613            type_closure_arg(&mut request, $closure_request);
614            let data = type_closure_two_arg(&mut request, Option::<Arc<ClientWorkerSet>>::None,
615                                            $closure_before);
616            async move {
617                type_closure_two_arg(<$client_type<_>>::$method(self, request).await,
618                                     data, $closure_after)
619            }.boxed()
620        }
621    };
622}
623macro_rules! proxier_impl {
624    ($trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident;
625     [$( proxy_def!($($def_args:tt)*); )*];
626     $(($method:ident, $req:ty, $resp:ty
627       $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
628        #[cfg(test)]
629        const $impl_list_name: &'static [&'static str] = &[$(stringify!($method)),*];
630
631        #[doc = concat!("Trait version of [", stringify!($client_type), "]")]
632        pub trait $trait_name: Send + Sync + DynClone
633        {
634            $( proxy_def!($($def_args)*); )*
635        }
636        dyn_clone::clone_trait_object!($trait_name);
637
638        // `allow(deprecated)`: upstream proto deprecations (e.g. cloud-api
639        // AddNamespaceRegion) generate calls to deprecated tonic client
640        // methods inside the impls; we still wire them for back-compat.
641        #[allow(deprecated)]
642        impl<RC> $trait_name for RC
643        where
644            RC: RawGrpcCaller + RawClientProducer + Clone + Unpin,
645        {
646            $(
647                proxy_impl!($client_type, $client_meth, $method, $req, $resp
648                            $(,$closure $(,$closure_before, $closure_after)*)*);
649            )*
650        }
651
652        impl<T: Send + Sync + 'static> RawGrpcCaller for $client_type<T> {}
653
654        #[allow(deprecated)]
655        impl<T> $trait_name for $client_type<T>
656        where
657            T: GrpcService<Body> + Clone + Send + Sync + 'static,
658            T::ResponseBody: tonic::codegen::Body<Data = tonic::codegen::Bytes> + Send + 'static,
659            T::Error: Into<tonic::codegen::StdError>,
660            <T::ResponseBody as tonic::codegen::Body>::Error: Into<tonic::codegen::StdError> + Send,
661            <T as tonic::client::GrpcService<Body>>::Future: Send
662        {
663            $(
664                proxy_impl!($client_type, $method, $req, $resp
665                            $(,$closure $(,$closure_before, $closure_after)*)*);
666            )*
667        }
668    };
669}
670
671macro_rules! proxier {
672    ( $trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident;
673      $(($method:ident, $req:ty, $resp:ty
674         $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
675        proxier_impl!($trait_name; $impl_list_name; $client_type; $client_meth;
676                      [$(proxy_def!($client_type, $client_meth, $method, $req, $resp);)*];
677                      $(($method, $req, $resp $(, $closure $(, $closure_before, $closure_after)?)?);)*);
678    };
679    ( $trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident; defaults;
680      $(($method:ident, $req:ty, $resp:ty
681         $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
682        proxier_impl!($trait_name; $impl_list_name; $client_type; $client_meth;
683                      [$(proxy_def!($client_type, $client_meth, $method, $req, $resp, defaults);)*];
684                      $(($method, $req, $resp $(, $closure $(, $closure_before, $closure_after)?)?);)*);
685    };
686}
687
688macro_rules! namespaced_request {
689    ($req:ident) => {{
690        let ns_str = $req.get_ref().namespace.clone();
691        // Attach namespace header
692        $req.metadata_mut().insert(
693            TEMPORAL_NAMESPACE_HEADER_KEY,
694            ns_str.parse().unwrap_or_else(|e| {
695                warn!("Unable to parse namespace for header: {e:?}");
696                AsciiMetadataValue::from_static("")
697            }),
698        );
699        // Init metric labels
700        AttachMetricLabels::namespace(ns_str)
701    }};
702}
703
704// Nice little trick to avoid the callsite asking to type the closure parameter
705fn type_closure_arg<T, R>(arg: T, f: impl FnOnce(T) -> R) -> R {
706    f(arg)
707}
708
709fn type_closure_two_arg<T, R, S>(arg1: R, arg2: T, f: impl FnOnce(R, T) -> S) -> S {
710    f(arg1, arg2)
711}
712
713proxier! {
714    WorkflowService; ALL_IMPLEMENTED_WORKFLOW_SERVICE_RPCS; WorkflowServiceClient; workflow_client; defaults;
715    (
716        register_namespace,
717        RegisterNamespaceRequest,
718        RegisterNamespaceResponse,
719        |r| {
720            let labels = namespaced_request!(r);
721            r.extensions_mut().insert(labels);
722        }
723    );
724    (
725        describe_namespace,
726        DescribeNamespaceRequest,
727        DescribeNamespaceResponse,
728        |r| {
729            let labels = namespaced_request!(r);
730            r.extensions_mut().insert(labels);
731        }
732    );
733    (
734        list_namespaces,
735        ListNamespacesRequest,
736        ListNamespacesResponse
737    );
738    (
739        update_namespace,
740        UpdateNamespaceRequest,
741        UpdateNamespaceResponse,
742        |r| {
743            let labels = namespaced_request!(r);
744            r.extensions_mut().insert(labels);
745        }
746    );
747    (
748        deprecate_namespace,
749        DeprecateNamespaceRequest,
750        DeprecateNamespaceResponse,
751        |r| {
752            let labels = namespaced_request!(r);
753            r.extensions_mut().insert(labels);
754        }
755    );
756    (
757        start_workflow_execution,
758        StartWorkflowExecutionRequest,
759        StartWorkflowExecutionResponse,
760        |r| {
761            let mut labels = namespaced_request!(r);
762            labels.task_q(r.get_ref().task_queue.clone());
763            r.extensions_mut().insert(labels);
764        },
765        |r, workers| {
766            if let Some(workers) = workers {
767                let mut slot: Option<Box<dyn Slot + Send>> = None;
768                let req_mut = r.get_mut();
769                if req_mut.request_eager_execution {
770                    let namespace = req_mut.namespace.clone();
771                    let task_queue = req_mut.task_queue.as_ref()
772                                        .map(|tq| tq.name.clone()).unwrap_or_default();
773                    match workers.try_reserve_wft_slot(namespace, task_queue) {
774                        Some(reservation) => {
775                            // Populate eager_worker_deployment_options from the slot reservation
776                            if let Some(opts) = reservation.deployment_options {
777                                req_mut.eager_worker_deployment_options = Some(temporalio_common::protos::temporal::api::deployment::v1::WorkerDeploymentOptions {
778                                    deployment_name: opts.version.deployment_name,
779                                    build_id: opts.version.build_id,
780                                    worker_versioning_mode: if opts.use_worker_versioning {
781                                        temporalio_common::protos::temporal::api::enums::v1::WorkerVersioningMode::Versioned.into()
782                                    } else {
783                                        temporalio_common::protos::temporal::api::enums::v1::WorkerVersioningMode::Unversioned.into()
784                                    },
785                                });                            }
786                            slot = Some(reservation.slot);
787                        }
788                        None => req_mut.request_eager_execution = false
789                    }
790                }
791                slot
792            } else {
793                None
794            }
795        },
796        |resp, slot| {
797            if let Some(s) = slot
798                && let Ok(response) = resp.as_ref()
799                    && let Some(task) = response.get_ref().clone().eager_workflow_task
800                        && let Err(e) = s.schedule_wft(task) {
801                            // This is a latency issue, i.e., the client does not need to handle
802                            //  this error, because the WFT will be retried after a timeout.
803                            warn!(details = ?e, "Eager workflow task rejected by worker.");
804                        }
805            resp
806        }
807    );
808    (
809        get_workflow_execution_history,
810        GetWorkflowExecutionHistoryRequest,
811        GetWorkflowExecutionHistoryResponse,
812        |r| {
813            let labels = namespaced_request!(r);
814            r.extensions_mut().insert(labels);
815            if r.get_ref().wait_new_event {
816                r.extensions_mut().insert(IsUserLongPoll);
817            }
818        }
819    );
820    (
821        get_workflow_execution_history_reverse,
822        GetWorkflowExecutionHistoryReverseRequest,
823        GetWorkflowExecutionHistoryReverseResponse,
824        |r| {
825            let labels = namespaced_request!(r);
826            r.extensions_mut().insert(labels);
827        }
828    );
829    (
830        poll_workflow_task_queue,
831        PollWorkflowTaskQueueRequest,
832        PollWorkflowTaskQueueResponse,
833        |r| {
834            let mut labels = namespaced_request!(r);
835            labels.task_q(r.get_ref().task_queue.clone());
836            r.extensions_mut().insert(labels);
837        }
838    );
839    (
840        respond_workflow_task_completed,
841        RespondWorkflowTaskCompletedRequest,
842        RespondWorkflowTaskCompletedResponse,
843        |r| {
844            let labels = namespaced_request!(r);
845            r.extensions_mut().insert(labels);
846        }
847    );
848    (
849        respond_workflow_task_failed,
850        RespondWorkflowTaskFailedRequest,
851        RespondWorkflowTaskFailedResponse,
852        |r| {
853            let labels = namespaced_request!(r);
854            r.extensions_mut().insert(labels);
855        }
856    );
857    (
858        poll_activity_task_queue,
859        PollActivityTaskQueueRequest,
860        PollActivityTaskQueueResponse,
861        |r| {
862            let mut labels = namespaced_request!(r);
863            labels.task_q(r.get_ref().task_queue.clone());
864            r.extensions_mut().insert(labels);
865        }
866    );
867    (
868        record_activity_task_heartbeat,
869        RecordActivityTaskHeartbeatRequest,
870        RecordActivityTaskHeartbeatResponse,
871        |r| {
872            let labels = namespaced_request!(r);
873            r.extensions_mut().insert(labels);
874        }
875    );
876    (
877        record_activity_task_heartbeat_by_id,
878        RecordActivityTaskHeartbeatByIdRequest,
879        RecordActivityTaskHeartbeatByIdResponse,
880        |r| {
881            let labels = namespaced_request!(r);
882            r.extensions_mut().insert(labels);
883        }
884    );
885    (
886        respond_activity_task_completed,
887        RespondActivityTaskCompletedRequest,
888        RespondActivityTaskCompletedResponse,
889        |r| {
890            let labels = namespaced_request!(r);
891            r.extensions_mut().insert(labels);
892        }
893    );
894    (
895        respond_activity_task_completed_by_id,
896        RespondActivityTaskCompletedByIdRequest,
897        RespondActivityTaskCompletedByIdResponse,
898        |r| {
899            let labels = namespaced_request!(r);
900            r.extensions_mut().insert(labels);
901        }
902    );
903
904    (
905        respond_activity_task_failed,
906        RespondActivityTaskFailedRequest,
907        RespondActivityTaskFailedResponse,
908        |r| {
909            let labels = namespaced_request!(r);
910            r.extensions_mut().insert(labels);
911        }
912    );
913    (
914        respond_activity_task_failed_by_id,
915        RespondActivityTaskFailedByIdRequest,
916        RespondActivityTaskFailedByIdResponse,
917        |r| {
918            let labels = namespaced_request!(r);
919            r.extensions_mut().insert(labels);
920        }
921    );
922    (
923        respond_activity_task_canceled,
924        RespondActivityTaskCanceledRequest,
925        RespondActivityTaskCanceledResponse,
926        |r| {
927            let labels = namespaced_request!(r);
928            r.extensions_mut().insert(labels);
929        }
930    );
931    (
932        respond_activity_task_canceled_by_id,
933        RespondActivityTaskCanceledByIdRequest,
934        RespondActivityTaskCanceledByIdResponse,
935        |r| {
936            let labels = namespaced_request!(r);
937            r.extensions_mut().insert(labels);
938        }
939    );
940    (
941        request_cancel_workflow_execution,
942        RequestCancelWorkflowExecutionRequest,
943        RequestCancelWorkflowExecutionResponse,
944        |r| {
945            let labels = namespaced_request!(r);
946            r.extensions_mut().insert(labels);
947        }
948    );
949    (
950        signal_workflow_execution,
951        SignalWorkflowExecutionRequest,
952        SignalWorkflowExecutionResponse,
953        |r| {
954            let labels = namespaced_request!(r);
955            r.extensions_mut().insert(labels);
956        }
957    );
958    (
959        signal_with_start_workflow_execution,
960        SignalWithStartWorkflowExecutionRequest,
961        SignalWithStartWorkflowExecutionResponse,
962        |r| {
963            let mut labels = namespaced_request!(r);
964            labels.task_q(r.get_ref().task_queue.clone());
965            r.extensions_mut().insert(labels);
966        }
967    );
968    (
969        reset_workflow_execution,
970        ResetWorkflowExecutionRequest,
971        ResetWorkflowExecutionResponse,
972        |r| {
973            let labels = namespaced_request!(r);
974            r.extensions_mut().insert(labels);
975        }
976    );
977    (
978        terminate_workflow_execution,
979        TerminateWorkflowExecutionRequest,
980        TerminateWorkflowExecutionResponse,
981        |r| {
982            let labels = namespaced_request!(r);
983            r.extensions_mut().insert(labels);
984        }
985    );
986    (
987        delete_workflow_execution,
988        DeleteWorkflowExecutionRequest,
989        DeleteWorkflowExecutionResponse,
990        |r| {
991            let labels = namespaced_request!(r);
992            r.extensions_mut().insert(labels);
993        }
994    );
995    (
996        list_open_workflow_executions,
997        ListOpenWorkflowExecutionsRequest,
998        ListOpenWorkflowExecutionsResponse,
999        |r| {
1000            let labels = namespaced_request!(r);
1001            r.extensions_mut().insert(labels);
1002        }
1003    );
1004    (
1005        list_closed_workflow_executions,
1006        ListClosedWorkflowExecutionsRequest,
1007        ListClosedWorkflowExecutionsResponse,
1008        |r| {
1009            let labels = namespaced_request!(r);
1010            r.extensions_mut().insert(labels);
1011        }
1012    );
1013    (
1014        list_workflow_executions,
1015        ListWorkflowExecutionsRequest,
1016        ListWorkflowExecutionsResponse,
1017        |r| {
1018            let labels = namespaced_request!(r);
1019            r.extensions_mut().insert(labels);
1020        }
1021    );
1022    (
1023        list_archived_workflow_executions,
1024        ListArchivedWorkflowExecutionsRequest,
1025        ListArchivedWorkflowExecutionsResponse,
1026        |r| {
1027            let labels = namespaced_request!(r);
1028            r.extensions_mut().insert(labels);
1029        }
1030    );
1031    (
1032        scan_workflow_executions,
1033        ScanWorkflowExecutionsRequest,
1034        ScanWorkflowExecutionsResponse,
1035        |r| {
1036            let labels = namespaced_request!(r);
1037            r.extensions_mut().insert(labels);
1038        }
1039    );
1040    (
1041        count_workflow_executions,
1042        CountWorkflowExecutionsRequest,
1043        CountWorkflowExecutionsResponse,
1044        |r| {
1045            let labels = namespaced_request!(r);
1046            r.extensions_mut().insert(labels);
1047        }
1048    );
1049    (
1050        create_workflow_rule,
1051        CreateWorkflowRuleRequest,
1052        CreateWorkflowRuleResponse,
1053        |r| {
1054            let labels = namespaced_request!(r);
1055            r.extensions_mut().insert(labels);
1056        }
1057    );
1058    (
1059        describe_workflow_rule,
1060        DescribeWorkflowRuleRequest,
1061        DescribeWorkflowRuleResponse,
1062        |r| {
1063            let labels = namespaced_request!(r);
1064            r.extensions_mut().insert(labels);
1065        }
1066    );
1067    (
1068        delete_workflow_rule,
1069        DeleteWorkflowRuleRequest,
1070        DeleteWorkflowRuleResponse,
1071        |r| {
1072            let labels = namespaced_request!(r);
1073            r.extensions_mut().insert(labels);
1074        }
1075    );
1076    (
1077        list_workflow_rules,
1078        ListWorkflowRulesRequest,
1079        ListWorkflowRulesResponse,
1080        |r| {
1081            let labels = namespaced_request!(r);
1082            r.extensions_mut().insert(labels);
1083        }
1084    );
1085    (
1086        trigger_workflow_rule,
1087        TriggerWorkflowRuleRequest,
1088        TriggerWorkflowRuleResponse,
1089        |r| {
1090            let labels = namespaced_request!(r);
1091            r.extensions_mut().insert(labels);
1092        }
1093    );
1094    (
1095        get_search_attributes,
1096        GetSearchAttributesRequest,
1097        GetSearchAttributesResponse
1098    );
1099    (
1100        respond_query_task_completed,
1101        RespondQueryTaskCompletedRequest,
1102        RespondQueryTaskCompletedResponse,
1103        |r| {
1104            let labels = namespaced_request!(r);
1105            r.extensions_mut().insert(labels);
1106        }
1107    );
1108    (
1109        reset_sticky_task_queue,
1110        ResetStickyTaskQueueRequest,
1111        ResetStickyTaskQueueResponse,
1112        |r| {
1113            let labels = namespaced_request!(r);
1114            r.extensions_mut().insert(labels);
1115        }
1116    );
1117    (
1118        query_workflow,
1119        QueryWorkflowRequest,
1120        QueryWorkflowResponse,
1121        |r| {
1122            let labels = namespaced_request!(r);
1123            r.extensions_mut().insert(labels);
1124        }
1125    );
1126    (
1127        describe_workflow_execution,
1128        DescribeWorkflowExecutionRequest,
1129        DescribeWorkflowExecutionResponse,
1130        |r| {
1131            let labels = namespaced_request!(r);
1132            r.extensions_mut().insert(labels);
1133        }
1134    );
1135    (
1136        describe_task_queue,
1137        DescribeTaskQueueRequest,
1138        DescribeTaskQueueResponse,
1139        |r| {
1140            let mut labels = namespaced_request!(r);
1141            labels.task_q(r.get_ref().task_queue.clone());
1142            r.extensions_mut().insert(labels);
1143        }
1144    );
1145    (
1146        get_cluster_info,
1147        GetClusterInfoRequest,
1148        GetClusterInfoResponse
1149    );
1150    (
1151        get_system_info,
1152        GetSystemInfoRequest,
1153        GetSystemInfoResponse
1154    );
1155    (
1156        list_task_queue_partitions,
1157        ListTaskQueuePartitionsRequest,
1158        ListTaskQueuePartitionsResponse,
1159        |r| {
1160            let mut labels = namespaced_request!(r);
1161            labels.task_q(r.get_ref().task_queue.clone());
1162            r.extensions_mut().insert(labels);
1163        }
1164    );
1165    (
1166        create_schedule,
1167        CreateScheduleRequest,
1168        CreateScheduleResponse,
1169        |r| {
1170            let labels = namespaced_request!(r);
1171            r.extensions_mut().insert(labels);
1172        }
1173    );
1174    (
1175        describe_schedule,
1176        DescribeScheduleRequest,
1177        DescribeScheduleResponse,
1178        |r| {
1179            let labels = namespaced_request!(r);
1180            r.extensions_mut().insert(labels);
1181        }
1182    );
1183    (
1184        update_schedule,
1185        UpdateScheduleRequest,
1186        UpdateScheduleResponse,
1187        |r| {
1188            let labels = namespaced_request!(r);
1189            r.extensions_mut().insert(labels);
1190        }
1191    );
1192    (
1193        patch_schedule,
1194        PatchScheduleRequest,
1195        PatchScheduleResponse,
1196        |r| {
1197            let labels = namespaced_request!(r);
1198            r.extensions_mut().insert(labels);
1199        }
1200    );
1201    (
1202        list_schedule_matching_times,
1203        ListScheduleMatchingTimesRequest,
1204        ListScheduleMatchingTimesResponse,
1205        |r| {
1206            let labels = namespaced_request!(r);
1207            r.extensions_mut().insert(labels);
1208        }
1209    );
1210    (
1211        delete_schedule,
1212        DeleteScheduleRequest,
1213        DeleteScheduleResponse,
1214        |r| {
1215            let labels = namespaced_request!(r);
1216            r.extensions_mut().insert(labels);
1217        }
1218    );
1219    (
1220        list_schedules,
1221        ListSchedulesRequest,
1222        ListSchedulesResponse,
1223        |r| {
1224            let labels = namespaced_request!(r);
1225            r.extensions_mut().insert(labels);
1226        }
1227    );
1228    (
1229        count_schedules,
1230        CountSchedulesRequest,
1231        CountSchedulesResponse,
1232        |r| {
1233            let labels = namespaced_request!(r);
1234            r.extensions_mut().insert(labels);
1235        }
1236    );
1237    (
1238        update_worker_build_id_compatibility,
1239        UpdateWorkerBuildIdCompatibilityRequest,
1240        UpdateWorkerBuildIdCompatibilityResponse,
1241        |r| {
1242            let mut labels = namespaced_request!(r);
1243            labels.task_q_str(r.get_ref().task_queue.clone());
1244            r.extensions_mut().insert(labels);
1245        }
1246    );
1247    (
1248        get_worker_build_id_compatibility,
1249        GetWorkerBuildIdCompatibilityRequest,
1250        GetWorkerBuildIdCompatibilityResponse,
1251        |r| {
1252            let mut labels = namespaced_request!(r);
1253            labels.task_q_str(r.get_ref().task_queue.clone());
1254            r.extensions_mut().insert(labels);
1255        }
1256    );
1257    (
1258        get_worker_task_reachability,
1259        GetWorkerTaskReachabilityRequest,
1260        GetWorkerTaskReachabilityResponse,
1261        |r| {
1262            let labels = namespaced_request!(r);
1263            r.extensions_mut().insert(labels);
1264        }
1265    );
1266    (
1267        update_workflow_execution,
1268        UpdateWorkflowExecutionRequest,
1269        UpdateWorkflowExecutionResponse,
1270        |r| {
1271            let labels = namespaced_request!(r);
1272            let exts = r.extensions_mut();
1273            exts.insert(labels);
1274            exts.insert(IsUserLongPoll);
1275        }
1276    );
1277    (
1278        poll_workflow_execution_update,
1279        PollWorkflowExecutionUpdateRequest,
1280        PollWorkflowExecutionUpdateResponse,
1281        |r| {
1282            let labels = namespaced_request!(r);
1283            r.extensions_mut().insert(labels);
1284        }
1285    );
1286    (
1287        start_batch_operation,
1288        StartBatchOperationRequest,
1289        StartBatchOperationResponse,
1290        |r| {
1291            let labels = namespaced_request!(r);
1292            r.extensions_mut().insert(labels);
1293        }
1294    );
1295    (
1296        stop_batch_operation,
1297        StopBatchOperationRequest,
1298        StopBatchOperationResponse,
1299        |r| {
1300            let labels = namespaced_request!(r);
1301            r.extensions_mut().insert(labels);
1302        }
1303    );
1304    (
1305        describe_batch_operation,
1306        DescribeBatchOperationRequest,
1307        DescribeBatchOperationResponse,
1308        |r| {
1309            let labels = namespaced_request!(r);
1310            r.extensions_mut().insert(labels);
1311        }
1312    );
1313    (
1314        describe_deployment,
1315        DescribeDeploymentRequest,
1316        DescribeDeploymentResponse,
1317        |r| {
1318            let labels = namespaced_request!(r);
1319            r.extensions_mut().insert(labels);
1320        }
1321    );
1322    (
1323        list_batch_operations,
1324        ListBatchOperationsRequest,
1325        ListBatchOperationsResponse,
1326        |r| {
1327            let labels = namespaced_request!(r);
1328            r.extensions_mut().insert(labels);
1329        }
1330    );
1331    (
1332        list_deployments,
1333        ListDeploymentsRequest,
1334        ListDeploymentsResponse,
1335        |r| {
1336            let labels = namespaced_request!(r);
1337            r.extensions_mut().insert(labels);
1338        }
1339    );
1340    (
1341        execute_multi_operation,
1342        ExecuteMultiOperationRequest,
1343        ExecuteMultiOperationResponse,
1344        |r| {
1345            let mut labels = namespaced_request!(r);
1346            if let Some(execute_multi_operation_request::operation::Operation::StartWorkflow(
1347                start_req,
1348            )) = r
1349                .get_ref()
1350                .operations
1351                .first()
1352                .and_then(|op| op.operation.as_ref())
1353            {
1354                labels.task_q(start_req.task_queue.clone());
1355            }
1356            let exts = r.extensions_mut();
1357            exts.insert(labels);
1358            // Update-with-start blocks until the update reaches the requested wait stage, so it
1359            // must be retried/timed out like other user long-polls.
1360            exts.insert(IsUserLongPoll);
1361        }
1362    );
1363    (
1364        get_current_deployment,
1365        GetCurrentDeploymentRequest,
1366        GetCurrentDeploymentResponse,
1367        |r| {
1368            let labels = namespaced_request!(r);
1369            r.extensions_mut().insert(labels);
1370        }
1371    );
1372    (
1373        get_deployment_reachability,
1374        GetDeploymentReachabilityRequest,
1375        GetDeploymentReachabilityResponse,
1376        |r| {
1377            let labels = namespaced_request!(r);
1378            r.extensions_mut().insert(labels);
1379        }
1380    );
1381    (
1382        get_worker_versioning_rules,
1383        GetWorkerVersioningRulesRequest,
1384        GetWorkerVersioningRulesResponse,
1385        |r| {
1386            let mut labels = namespaced_request!(r);
1387            labels.task_q_str(&r.get_ref().task_queue);
1388            r.extensions_mut().insert(labels);
1389        }
1390    );
1391    (
1392        update_worker_versioning_rules,
1393        UpdateWorkerVersioningRulesRequest,
1394        UpdateWorkerVersioningRulesResponse,
1395        |r| {
1396            let mut labels = namespaced_request!(r);
1397            labels.task_q_str(&r.get_ref().task_queue);
1398            r.extensions_mut().insert(labels);
1399        }
1400    );
1401    (
1402        poll_nexus_task_queue,
1403        PollNexusTaskQueueRequest,
1404        PollNexusTaskQueueResponse,
1405        |r| {
1406            let mut labels = namespaced_request!(r);
1407            labels.task_q(r.get_ref().task_queue.clone());
1408            r.extensions_mut().insert(labels);
1409        }
1410    );
1411    (
1412        respond_nexus_task_completed,
1413        RespondNexusTaskCompletedRequest,
1414        RespondNexusTaskCompletedResponse,
1415        |r| {
1416            let labels = namespaced_request!(r);
1417            r.extensions_mut().insert(labels);
1418        }
1419    );
1420    (
1421        respond_nexus_task_failed,
1422        RespondNexusTaskFailedRequest,
1423        RespondNexusTaskFailedResponse,
1424        |r| {
1425            let labels = namespaced_request!(r);
1426            r.extensions_mut().insert(labels);
1427        }
1428    );
1429    (
1430        set_current_deployment,
1431        SetCurrentDeploymentRequest,
1432        SetCurrentDeploymentResponse,
1433        |r| {
1434            let labels = namespaced_request!(r);
1435            r.extensions_mut().insert(labels);
1436        }
1437    );
1438    (
1439        shutdown_worker,
1440        ShutdownWorkerRequest,
1441        ShutdownWorkerResponse,
1442        |r| {
1443            let labels = namespaced_request!(r);
1444            r.extensions_mut().insert(labels);
1445        }
1446    );
1447    (
1448        update_activity_options,
1449        UpdateActivityOptionsRequest,
1450        UpdateActivityOptionsResponse,
1451        |r| {
1452            let labels = namespaced_request!(r);
1453            r.extensions_mut().insert(labels);
1454        }
1455    );
1456    (
1457        pause_activity,
1458        PauseActivityRequest,
1459        PauseActivityResponse,
1460        |r| {
1461            let labels = namespaced_request!(r);
1462            r.extensions_mut().insert(labels);
1463        }
1464    );
1465    (
1466        unpause_activity,
1467        UnpauseActivityRequest,
1468        UnpauseActivityResponse,
1469        |r| {
1470            let labels = namespaced_request!(r);
1471            r.extensions_mut().insert(labels);
1472        }
1473    );
1474    (
1475        update_workflow_execution_options,
1476        UpdateWorkflowExecutionOptionsRequest,
1477        UpdateWorkflowExecutionOptionsResponse,
1478        |r| {
1479            let labels = namespaced_request!(r);
1480            r.extensions_mut().insert(labels);
1481        }
1482    );
1483    (
1484        reset_activity,
1485        ResetActivityRequest,
1486        ResetActivityResponse,
1487        |r| {
1488            let labels = namespaced_request!(r);
1489            r.extensions_mut().insert(labels);
1490        }
1491    );
1492    (
1493        delete_worker_deployment,
1494        DeleteWorkerDeploymentRequest,
1495        DeleteWorkerDeploymentResponse,
1496        |r| {
1497            let labels = namespaced_request!(r);
1498            r.extensions_mut().insert(labels);
1499        }
1500    );
1501    (
1502        delete_worker_deployment_version,
1503        DeleteWorkerDeploymentVersionRequest,
1504        DeleteWorkerDeploymentVersionResponse,
1505        |r| {
1506            let labels = namespaced_request!(r);
1507            r.extensions_mut().insert(labels);
1508        }
1509    );
1510    (
1511        describe_worker_deployment,
1512        DescribeWorkerDeploymentRequest,
1513        DescribeWorkerDeploymentResponse,
1514        |r| {
1515            let labels = namespaced_request!(r);
1516            r.extensions_mut().insert(labels);
1517        }
1518    );
1519    (
1520        describe_worker_deployment_version,
1521        DescribeWorkerDeploymentVersionRequest,
1522        DescribeWorkerDeploymentVersionResponse,
1523        |r| {
1524            let labels = namespaced_request!(r);
1525            r.extensions_mut().insert(labels);
1526        }
1527    );
1528    (
1529        list_worker_deployments,
1530        ListWorkerDeploymentsRequest,
1531        ListWorkerDeploymentsResponse,
1532        |r| {
1533            let labels = namespaced_request!(r);
1534            r.extensions_mut().insert(labels);
1535        }
1536    );
1537    (
1538        set_worker_deployment_current_version,
1539        SetWorkerDeploymentCurrentVersionRequest,
1540        SetWorkerDeploymentCurrentVersionResponse,
1541        |r| {
1542            let labels = namespaced_request!(r);
1543            r.extensions_mut().insert(labels);
1544        }
1545    );
1546    (
1547        set_worker_deployment_ramping_version,
1548        SetWorkerDeploymentRampingVersionRequest,
1549        SetWorkerDeploymentRampingVersionResponse,
1550        |r| {
1551            let labels = namespaced_request!(r);
1552            r.extensions_mut().insert(labels);
1553        }
1554    );
1555    (
1556        update_worker_deployment_version_metadata,
1557        UpdateWorkerDeploymentVersionMetadataRequest,
1558        UpdateWorkerDeploymentVersionMetadataResponse,
1559        |r| {
1560            let labels = namespaced_request!(r);
1561            r.extensions_mut().insert(labels);
1562        }
1563    );
1564    (
1565        list_workers,
1566        ListWorkersRequest,
1567        ListWorkersResponse,
1568        |r| {
1569            let labels = namespaced_request!(r);
1570            r.extensions_mut().insert(labels);
1571        }
1572    );
1573    (
1574        count_workers,
1575        CountWorkersRequest,
1576        CountWorkersResponse,
1577        |r| {
1578            let labels = namespaced_request!(r);
1579            r.extensions_mut().insert(labels);
1580        }
1581    );
1582    (
1583        record_worker_heartbeat,
1584        RecordWorkerHeartbeatRequest,
1585        RecordWorkerHeartbeatResponse,
1586        |r| {
1587            let labels = namespaced_request!(r);
1588            r.extensions_mut().insert(labels);
1589        }
1590    );
1591    (
1592        update_task_queue_config,
1593        UpdateTaskQueueConfigRequest,
1594        UpdateTaskQueueConfigResponse,
1595        |r| {
1596            let mut labels = namespaced_request!(r);
1597            labels.task_q_str(r.get_ref().task_queue.clone());
1598            r.extensions_mut().insert(labels);
1599        }
1600    );
1601    (
1602        fetch_worker_config,
1603        FetchWorkerConfigRequest,
1604        FetchWorkerConfigResponse,
1605        |r| {
1606            let labels = namespaced_request!(r);
1607            r.extensions_mut().insert(labels);
1608        }
1609    );
1610    (
1611        update_worker_config,
1612        UpdateWorkerConfigRequest,
1613        UpdateWorkerConfigResponse,
1614        |r| {
1615            let labels = namespaced_request!(r);
1616            r.extensions_mut().insert(labels);
1617        }
1618    );
1619    (
1620        describe_worker,
1621        DescribeWorkerRequest,
1622        DescribeWorkerResponse,
1623        |r| {
1624            let labels = namespaced_request!(r);
1625            r.extensions_mut().insert(labels);
1626        }
1627    );
1628    (
1629        set_worker_deployment_manager,
1630        SetWorkerDeploymentManagerRequest,
1631        SetWorkerDeploymentManagerResponse,
1632        |r| {
1633            let labels = namespaced_request!(r);
1634            r.extensions_mut().insert(labels);
1635        }
1636    );
1637    (
1638        pause_workflow_execution,
1639        PauseWorkflowExecutionRequest,
1640        PauseWorkflowExecutionResponse,
1641        |r| {
1642            let labels = namespaced_request!(r);
1643            r.extensions_mut().insert(labels);
1644        }
1645    );
1646    (
1647        unpause_workflow_execution,
1648        UnpauseWorkflowExecutionRequest,
1649        UnpauseWorkflowExecutionResponse,
1650        |r| {
1651            let labels = namespaced_request!(r);
1652            r.extensions_mut().insert(labels);
1653        }
1654    );
1655    (
1656        start_activity_execution,
1657        StartActivityExecutionRequest,
1658        StartActivityExecutionResponse,
1659        |r| {
1660            let labels = namespaced_request!(r);
1661            r.extensions_mut().insert(labels);
1662        }
1663    );
1664    (
1665        describe_activity_execution,
1666        DescribeActivityExecutionRequest,
1667        DescribeActivityExecutionResponse,
1668        |r| {
1669            let labels = namespaced_request!(r);
1670            r.extensions_mut().insert(labels);
1671        }
1672    );
1673    (
1674        poll_activity_execution,
1675        PollActivityExecutionRequest,
1676        PollActivityExecutionResponse,
1677        |r| {
1678            let labels = namespaced_request!(r);
1679            r.extensions_mut().insert(labels);
1680        }
1681    );
1682    (
1683        list_activity_executions,
1684        ListActivityExecutionsRequest,
1685        ListActivityExecutionsResponse,
1686        |r| {
1687            let labels = namespaced_request!(r);
1688            r.extensions_mut().insert(labels);
1689        }
1690    );
1691    (
1692        count_activity_executions,
1693        CountActivityExecutionsRequest,
1694        CountActivityExecutionsResponse,
1695        |r| {
1696            let labels = namespaced_request!(r);
1697            r.extensions_mut().insert(labels);
1698        }
1699    );
1700    (
1701        request_cancel_activity_execution,
1702        RequestCancelActivityExecutionRequest,
1703        RequestCancelActivityExecutionResponse,
1704        |r| {
1705            let labels = namespaced_request!(r);
1706            r.extensions_mut().insert(labels);
1707        }
1708    );
1709    (
1710        terminate_activity_execution,
1711        TerminateActivityExecutionRequest,
1712        TerminateActivityExecutionResponse,
1713        |r| {
1714            let labels = namespaced_request!(r);
1715            r.extensions_mut().insert(labels);
1716        }
1717    );
1718    (
1719        delete_activity_execution,
1720        DeleteActivityExecutionRequest,
1721        DeleteActivityExecutionResponse,
1722        |r| {
1723            let labels = namespaced_request!(r);
1724            r.extensions_mut().insert(labels);
1725        }
1726    );
1727    (
1728        pause_activity_execution,
1729        PauseActivityExecutionRequest,
1730        PauseActivityExecutionResponse,
1731        |r| {
1732            let labels = namespaced_request!(r);
1733            r.extensions_mut().insert(labels);
1734        }
1735    );
1736    (
1737        unpause_activity_execution,
1738        UnpauseActivityExecutionRequest,
1739        UnpauseActivityExecutionResponse,
1740        |r| {
1741            let labels = namespaced_request!(r);
1742            r.extensions_mut().insert(labels);
1743        }
1744    );
1745    (
1746        reset_activity_execution,
1747        ResetActivityExecutionRequest,
1748        ResetActivityExecutionResponse,
1749        |r| {
1750            let labels = namespaced_request!(r);
1751            r.extensions_mut().insert(labels);
1752        }
1753    );
1754    (
1755        update_activity_execution_options,
1756        UpdateActivityExecutionOptionsRequest,
1757        UpdateActivityExecutionOptionsResponse,
1758        |r| {
1759            let labels = namespaced_request!(r);
1760            r.extensions_mut().insert(labels);
1761        }
1762    );
1763    (
1764        count_nexus_operation_executions,
1765        CountNexusOperationExecutionsRequest,
1766        CountNexusOperationExecutionsResponse,
1767        |r| {
1768            let labels = namespaced_request!(r);
1769            r.extensions_mut().insert(labels);
1770        }
1771    );
1772    (
1773        create_worker_deployment,
1774        CreateWorkerDeploymentRequest,
1775        CreateWorkerDeploymentResponse,
1776        |r| {
1777            let labels = namespaced_request!(r);
1778            r.extensions_mut().insert(labels);
1779        }
1780    );
1781    (
1782        create_worker_deployment_version,
1783        CreateWorkerDeploymentVersionRequest,
1784        CreateWorkerDeploymentVersionResponse,
1785        |r| {
1786            let labels = namespaced_request!(r);
1787            r.extensions_mut().insert(labels);
1788        }
1789    );
1790    (
1791        delete_nexus_operation_execution,
1792        DeleteNexusOperationExecutionRequest,
1793        DeleteNexusOperationExecutionResponse,
1794        |r| {
1795            let labels = namespaced_request!(r);
1796            r.extensions_mut().insert(labels);
1797        }
1798    );
1799    (
1800        describe_nexus_operation_execution,
1801        DescribeNexusOperationExecutionRequest,
1802        DescribeNexusOperationExecutionResponse,
1803        |r| {
1804            let labels = namespaced_request!(r);
1805            r.extensions_mut().insert(labels);
1806        }
1807    );
1808    (
1809        list_nexus_operation_executions,
1810        ListNexusOperationExecutionsRequest,
1811        ListNexusOperationExecutionsResponse,
1812        |r| {
1813            let labels = namespaced_request!(r);
1814            r.extensions_mut().insert(labels);
1815        }
1816    );
1817    (
1818        poll_nexus_operation_execution,
1819        PollNexusOperationExecutionRequest,
1820        PollNexusOperationExecutionResponse,
1821        |r| {
1822            let labels = namespaced_request!(r);
1823            r.extensions_mut().insert(labels);
1824        }
1825    );
1826    (
1827        poll_workflow_execution_time_skipping,
1828        PollWorkflowExecutionTimeSkippingRequest,
1829        PollWorkflowExecutionTimeSkippingResponse,
1830        |r| {
1831            let labels = namespaced_request!(r);
1832            r.extensions_mut().insert(labels);
1833            r.extensions_mut().insert(IsUserLongPoll);
1834        }
1835    );
1836    (
1837        request_cancel_nexus_operation_execution,
1838        RequestCancelNexusOperationExecutionRequest,
1839        RequestCancelNexusOperationExecutionResponse,
1840        |r| {
1841            let labels = namespaced_request!(r);
1842            r.extensions_mut().insert(labels);
1843        }
1844    );
1845    (
1846        start_nexus_operation_execution,
1847        StartNexusOperationExecutionRequest,
1848        StartNexusOperationExecutionResponse,
1849        |r| {
1850            let labels = namespaced_request!(r);
1851            r.extensions_mut().insert(labels);
1852        }
1853    );
1854    (
1855        terminate_nexus_operation_execution,
1856        TerminateNexusOperationExecutionRequest,
1857        TerminateNexusOperationExecutionResponse,
1858        |r| {
1859            let labels = namespaced_request!(r);
1860            r.extensions_mut().insert(labels);
1861        }
1862    );
1863    (
1864        update_worker_deployment_version_compute_config,
1865        UpdateWorkerDeploymentVersionComputeConfigRequest,
1866        UpdateWorkerDeploymentVersionComputeConfigResponse,
1867        |r| {
1868            let labels = namespaced_request!(r);
1869            r.extensions_mut().insert(labels);
1870        }
1871    );
1872    (
1873        validate_worker_deployment_version_compute_config,
1874        ValidateWorkerDeploymentVersionComputeConfigRequest,
1875        ValidateWorkerDeploymentVersionComputeConfigResponse,
1876        |r| {
1877            let labels = namespaced_request!(r);
1878            r.extensions_mut().insert(labels);
1879        }
1880    );
1881}
1882
1883proxier! {
1884    OperatorService; ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS; OperatorServiceClient; operator_client; defaults;
1885    (add_search_attributes, AddSearchAttributesRequest, AddSearchAttributesResponse);
1886    (remove_search_attributes, RemoveSearchAttributesRequest, RemoveSearchAttributesResponse);
1887    (list_search_attributes, ListSearchAttributesRequest, ListSearchAttributesResponse);
1888    (delete_namespace, DeleteNamespaceRequest, DeleteNamespaceResponse,
1889        |r| {
1890            let labels = namespaced_request!(r);
1891            r.extensions_mut().insert(labels);
1892        }
1893    );
1894    (add_or_update_remote_cluster, AddOrUpdateRemoteClusterRequest, AddOrUpdateRemoteClusterResponse);
1895    (remove_remote_cluster, RemoveRemoteClusterRequest, RemoveRemoteClusterResponse);
1896    (list_clusters, ListClustersRequest, ListClustersResponse);
1897    (get_nexus_endpoint, GetNexusEndpointRequest, GetNexusEndpointResponse);
1898    (create_nexus_endpoint, CreateNexusEndpointRequest, CreateNexusEndpointResponse);
1899    (update_nexus_endpoint, UpdateNexusEndpointRequest, UpdateNexusEndpointResponse);
1900    (delete_nexus_endpoint, DeleteNexusEndpointRequest, DeleteNexusEndpointResponse);
1901    (list_nexus_endpoints, ListNexusEndpointsRequest, ListNexusEndpointsResponse);
1902}
1903
1904proxier! {
1905    CloudService; ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS; CloudServiceClient; cloud_client; defaults;
1906    (get_users, cloudreq::GetUsersRequest, cloudreq::GetUsersResponse);
1907    (get_user, cloudreq::GetUserRequest, cloudreq::GetUserResponse);
1908    (create_user, cloudreq::CreateUserRequest, cloudreq::CreateUserResponse);
1909    (update_user, cloudreq::UpdateUserRequest, cloudreq::UpdateUserResponse);
1910    (delete_user, cloudreq::DeleteUserRequest, cloudreq::DeleteUserResponse);
1911    (set_user_namespace_access, cloudreq::SetUserNamespaceAccessRequest, cloudreq::SetUserNamespaceAccessResponse);
1912    (get_async_operation, cloudreq::GetAsyncOperationRequest, cloudreq::GetAsyncOperationResponse);
1913    (create_namespace, cloudreq::CreateNamespaceRequest, cloudreq::CreateNamespaceResponse);
1914    (get_namespaces, cloudreq::GetNamespacesRequest, cloudreq::GetNamespacesResponse);
1915    (get_namespace, cloudreq::GetNamespaceRequest, cloudreq::GetNamespaceResponse,
1916        |r| {
1917            let labels = namespaced_request!(r);
1918            r.extensions_mut().insert(labels);
1919        }
1920    );
1921    (update_namespace, cloudreq::UpdateNamespaceRequest, cloudreq::UpdateNamespaceResponse,
1922        |r| {
1923            let labels = namespaced_request!(r);
1924            r.extensions_mut().insert(labels);
1925        }
1926    );
1927    (rename_custom_search_attribute, cloudreq::RenameCustomSearchAttributeRequest, cloudreq::RenameCustomSearchAttributeResponse);
1928    (delete_namespace, cloudreq::DeleteNamespaceRequest, cloudreq::DeleteNamespaceResponse,
1929        |r| {
1930            let labels = namespaced_request!(r);
1931            r.extensions_mut().insert(labels);
1932        }
1933    );
1934    (failover_namespace_region, cloudreq::FailoverNamespaceRegionRequest, cloudreq::FailoverNamespaceRegionResponse);
1935    (add_namespace_region, cloudreq::AddNamespaceRegionRequest, cloudreq::AddNamespaceRegionResponse);
1936    (delete_namespace_region, cloudreq::DeleteNamespaceRegionRequest, cloudreq::DeleteNamespaceRegionResponse);
1937    (get_regions, cloudreq::GetRegionsRequest, cloudreq::GetRegionsResponse);
1938    (get_region, cloudreq::GetRegionRequest, cloudreq::GetRegionResponse);
1939    (get_api_keys, cloudreq::GetApiKeysRequest, cloudreq::GetApiKeysResponse);
1940    (get_api_key, cloudreq::GetApiKeyRequest, cloudreq::GetApiKeyResponse);
1941    (create_api_key, cloudreq::CreateApiKeyRequest, cloudreq::CreateApiKeyResponse);
1942    (update_api_key, cloudreq::UpdateApiKeyRequest, cloudreq::UpdateApiKeyResponse);
1943    (delete_api_key, cloudreq::DeleteApiKeyRequest, cloudreq::DeleteApiKeyResponse);
1944    (get_nexus_endpoints, cloudreq::GetNexusEndpointsRequest, cloudreq::GetNexusEndpointsResponse);
1945    (get_nexus_endpoint, cloudreq::GetNexusEndpointRequest, cloudreq::GetNexusEndpointResponse);
1946    (create_nexus_endpoint, cloudreq::CreateNexusEndpointRequest, cloudreq::CreateNexusEndpointResponse);
1947    (update_nexus_endpoint, cloudreq::UpdateNexusEndpointRequest, cloudreq::UpdateNexusEndpointResponse);
1948    (delete_nexus_endpoint, cloudreq::DeleteNexusEndpointRequest, cloudreq::DeleteNexusEndpointResponse);
1949    (get_user_groups, cloudreq::GetUserGroupsRequest, cloudreq::GetUserGroupsResponse);
1950    (get_user_group, cloudreq::GetUserGroupRequest, cloudreq::GetUserGroupResponse);
1951    (create_user_group, cloudreq::CreateUserGroupRequest, cloudreq::CreateUserGroupResponse);
1952    (update_user_group, cloudreq::UpdateUserGroupRequest, cloudreq::UpdateUserGroupResponse);
1953    (delete_user_group, cloudreq::DeleteUserGroupRequest, cloudreq::DeleteUserGroupResponse);
1954    (add_user_group_member, cloudreq::AddUserGroupMemberRequest, cloudreq::AddUserGroupMemberResponse);
1955    (remove_user_group_member, cloudreq::RemoveUserGroupMemberRequest, cloudreq::RemoveUserGroupMemberResponse);
1956    (get_user_group_members, cloudreq::GetUserGroupMembersRequest, cloudreq::GetUserGroupMembersResponse);
1957    (set_user_group_namespace_access, cloudreq::SetUserGroupNamespaceAccessRequest, cloudreq::SetUserGroupNamespaceAccessResponse);
1958    (create_service_account, cloudreq::CreateServiceAccountRequest, cloudreq::CreateServiceAccountResponse);
1959    (get_service_account, cloudreq::GetServiceAccountRequest, cloudreq::GetServiceAccountResponse);
1960    (get_service_accounts, cloudreq::GetServiceAccountsRequest, cloudreq::GetServiceAccountsResponse);
1961    (update_service_account, cloudreq::UpdateServiceAccountRequest, cloudreq::UpdateServiceAccountResponse);
1962    (delete_service_account, cloudreq::DeleteServiceAccountRequest, cloudreq::DeleteServiceAccountResponse);
1963    (get_usage, cloudreq::GetUsageRequest, cloudreq::GetUsageResponse);
1964    (get_account, cloudreq::GetAccountRequest, cloudreq::GetAccountResponse);
1965    (update_account, cloudreq::UpdateAccountRequest, cloudreq::UpdateAccountResponse);
1966    (create_namespace_export_sink, cloudreq::CreateNamespaceExportSinkRequest, cloudreq::CreateNamespaceExportSinkResponse);
1967    (get_namespace_export_sink, cloudreq::GetNamespaceExportSinkRequest, cloudreq::GetNamespaceExportSinkResponse);
1968    (get_namespace_export_sinks, cloudreq::GetNamespaceExportSinksRequest, cloudreq::GetNamespaceExportSinksResponse);
1969    (update_namespace_export_sink, cloudreq::UpdateNamespaceExportSinkRequest, cloudreq::UpdateNamespaceExportSinkResponse);
1970    (delete_namespace_export_sink, cloudreq::DeleteNamespaceExportSinkRequest, cloudreq::DeleteNamespaceExportSinkResponse);
1971    (validate_namespace_export_sink, cloudreq::ValidateNamespaceExportSinkRequest, cloudreq::ValidateNamespaceExportSinkResponse);
1972    (update_namespace_tags, cloudreq::UpdateNamespaceTagsRequest, cloudreq::UpdateNamespaceTagsResponse);
1973    (create_connectivity_rule, cloudreq::CreateConnectivityRuleRequest, cloudreq::CreateConnectivityRuleResponse);
1974    (get_connectivity_rule, cloudreq::GetConnectivityRuleRequest, cloudreq::GetConnectivityRuleResponse);
1975    (get_connectivity_rules, cloudreq::GetConnectivityRulesRequest, cloudreq::GetConnectivityRulesResponse);
1976    (delete_connectivity_rule, cloudreq::DeleteConnectivityRuleRequest, cloudreq::DeleteConnectivityRuleResponse);
1977    (set_service_account_namespace_access, cloudreq::SetServiceAccountNamespaceAccessRequest, cloudreq::SetServiceAccountNamespaceAccessResponse);
1978    (validate_account_audit_log_sink, cloudreq::ValidateAccountAuditLogSinkRequest, cloudreq::ValidateAccountAuditLogSinkResponse);
1979    (get_current_identity, cloudreq::GetCurrentIdentityRequest, cloudreq::GetCurrentIdentityResponse);
1980    (get_audit_logs, cloudreq::GetAuditLogsRequest, cloudreq::GetAuditLogsResponse);
1981    (create_account_audit_log_sink, cloudreq::CreateAccountAuditLogSinkRequest, cloudreq::CreateAccountAuditLogSinkResponse);
1982    (get_account_audit_log_sink, cloudreq::GetAccountAuditLogSinkRequest, cloudreq::GetAccountAuditLogSinkResponse);
1983    (get_account_audit_log_sinks, cloudreq::GetAccountAuditLogSinksRequest, cloudreq::GetAccountAuditLogSinksResponse);
1984    (update_account_audit_log_sink, cloudreq::UpdateAccountAuditLogSinkRequest, cloudreq::UpdateAccountAuditLogSinkResponse);
1985    (delete_account_audit_log_sink, cloudreq::DeleteAccountAuditLogSinkRequest, cloudreq::DeleteAccountAuditLogSinkResponse);
1986    (get_namespace_capacity_info, cloudreq::GetNamespaceCapacityInfoRequest, cloudreq::GetNamespaceCapacityInfoResponse);
1987    (create_billing_report, cloudreq::CreateBillingReportRequest, cloudreq::CreateBillingReportResponse);
1988    (get_billing_report, cloudreq::GetBillingReportRequest, cloudreq::GetBillingReportResponse);
1989    (get_custom_roles, cloudreq::GetCustomRolesRequest, cloudreq::GetCustomRolesResponse);
1990    (get_custom_role, cloudreq::GetCustomRoleRequest, cloudreq::GetCustomRoleResponse);
1991    (create_custom_role, cloudreq::CreateCustomRoleRequest, cloudreq::CreateCustomRoleResponse);
1992    (update_custom_role, cloudreq::UpdateCustomRoleRequest, cloudreq::UpdateCustomRoleResponse);
1993    (delete_custom_role, cloudreq::DeleteCustomRoleRequest, cloudreq::DeleteCustomRoleResponse);
1994    (get_user_namespace_assignments, cloudreq::GetUserNamespaceAssignmentsRequest, cloudreq::GetUserNamespaceAssignmentsResponse);
1995    (get_service_account_namespace_assignments, cloudreq::GetServiceAccountNamespaceAssignmentsRequest, cloudreq::GetServiceAccountNamespaceAssignmentsResponse);
1996    (get_user_group_namespace_assignments, cloudreq::GetUserGroupNamespaceAssignmentsRequest, cloudreq::GetUserGroupNamespaceAssignmentsResponse);
1997}
1998
1999proxier! {
2000    TestService; ALL_IMPLEMENTED_TEST_SERVICE_RPCS; TestServiceClient; test_client; defaults;
2001    (lock_time_skipping, LockTimeSkippingRequest, LockTimeSkippingResponse);
2002    (unlock_time_skipping, UnlockTimeSkippingRequest, UnlockTimeSkippingResponse);
2003    (sleep, SleepRequest, SleepResponse);
2004    (sleep_until, SleepUntilRequest, SleepResponse);
2005    (unlock_time_skipping_with_sleep, SleepRequest, SleepResponse);
2006    (get_current_time, (), GetCurrentTimeResponse);
2007}
2008
2009proxier! {
2010    HealthService; ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS; HealthClient; health_client;
2011    (check, HealthCheckRequest, HealthCheckResponse);
2012    (watch, HealthCheckRequest, tonic::codec::Streaming<HealthCheckResponse>);
2013}
2014
2015#[cfg(test)]
2016mod tests {
2017    use super::*;
2018    use crate::{ClientOptions, ConnectionOptions};
2019    use std::collections::HashSet;
2020    use temporalio_common::{
2021        protos::temporal::api::{
2022            operatorservice::v1::DeleteNamespaceRequest, workflowservice::v1::ListNamespacesRequest,
2023        },
2024        worker::WorkerTaskTypes,
2025    };
2026    use tonic::IntoRequest;
2027    use url::Url;
2028    use uuid::Uuid;
2029
2030    #[test]
2031    fn payload_limits_warn_only_vs_error() {
2032        use temporalio_common::protos::temporal::api::{
2033            common::v1::{Payload, Payloads},
2034            workflowservice::v1::StartWorkflowExecutionRequest,
2035        };
2036        let big = Payloads {
2037            payloads: vec![Payload {
2038                data: vec![0u8; 1000],
2039                ..Default::default()
2040            }],
2041        };
2042        let new_req = || {
2043            StartWorkflowExecutionRequest {
2044                input: Some(big.clone()),
2045                ..Default::default()
2046            }
2047            .into_request()
2048        };
2049
2050        // warn thresholds = 1 byte. No per-call error limits: over-warn is allowed (warn-only).
2051        assert!(validate_request_payload_limits(&new_req(), 1, 1).is_ok());
2052
2053        // With per-call error limits below the payload size: rejected, carrying the typed violation.
2054        let mut req = new_req();
2055        req.extensions_mut()
2056            .insert(PayloadErrorLimits { blob: 10, memo: 10 });
2057        let err = validate_request_payload_limits(&req, 1, 1).unwrap_err();
2058        let violation =
2059            crate::payload_limit_violation_from(&err).expect("violation carried on status");
2060        assert_eq!(violation.path, "input");
2061        assert_eq!(
2062            violation.class,
2063            temporalio_common::payload_limits::LimitClass::Blob
2064        );
2065        assert!(violation.size > violation.limit);
2066
2067        // A zero error threshold means "no limit" for that class, so it does not reject.
2068        let mut req = new_req();
2069        req.extensions_mut()
2070            .insert(PayloadErrorLimits { blob: 0, memo: 0 });
2071        assert!(validate_request_payload_limits(&req, 1, 1).is_ok());
2072
2073        // Zero warn thresholds disable warnings (and there are no error limits): always ok.
2074        assert!(validate_request_payload_limits(&new_req(), 0, 0).is_ok());
2075    }
2076
2077    // Just to help make sure some stuff compiles. Not run.
2078    #[allow(dead_code)]
2079    async fn raw_client_retry_compiles() {
2080        let opts = ConnectionOptions::new(Url::parse("http://localhost:7233").unwrap())
2081            .client_name("test")
2082            .client_version("0.0.0")
2083            .build();
2084        let connection = Connection::connect(opts).await.unwrap();
2085        let mut client = Client::new(connection, ClientOptions::new("default").build()).unwrap();
2086
2087        let list_ns_req = ListNamespacesRequest::default();
2088        let wf_client = client.workflow_client();
2089        let fact = move |req| {
2090            let mut c = wf_client.clone();
2091            async move { c.list_namespaces(req).await }.boxed()
2092        };
2093        client
2094            .call("whatever", fact, Request::new(list_ns_req.clone()))
2095            .await
2096            .unwrap();
2097
2098        // Operator svc method
2099        let op_del_ns_req = DeleteNamespaceRequest::default();
2100        let op_client = client.operator_client();
2101        let fact = move |req| {
2102            let mut c = op_client.clone();
2103            async move { c.delete_namespace(req).await }.boxed()
2104        };
2105        client
2106            .call("whatever", fact, Request::new(op_del_ns_req.clone()))
2107            .await
2108            .unwrap();
2109
2110        // Cloud svc method
2111        let cloud_del_ns_req = cloudreq::DeleteNamespaceRequest::default();
2112        let cloud_client = client.cloud_client();
2113        let fact = move |req| {
2114            let mut c = cloud_client.clone();
2115            async move { c.delete_namespace(req).await }.boxed()
2116        };
2117        client
2118            .call("whatever", fact, Request::new(cloud_del_ns_req.clone()))
2119            .await
2120            .unwrap();
2121
2122        // Verify calling through traits works
2123        client
2124            .list_namespaces(list_ns_req.into_request())
2125            .await
2126            .unwrap();
2127        // Have to disambiguate operator and cloud service
2128        OperatorService::delete_namespace(&mut client, op_del_ns_req.into_request())
2129            .await
2130            .unwrap();
2131        CloudService::delete_namespace(&mut client, cloud_del_ns_req.into_request())
2132            .await
2133            .unwrap();
2134        client.get_current_time(().into_request()).await.unwrap();
2135        client
2136            .check(HealthCheckRequest::default().into_request())
2137            .await
2138            .unwrap();
2139    }
2140
2141    fn verify_methods(proto_def_str: &str, impl_list: &[&str]) {
2142        let methods: Vec<_> = proto_def_str
2143            .lines()
2144            .map(|l| l.trim())
2145            .filter(|l| l.starts_with("rpc"))
2146            .map(|l| {
2147                let stripped = l.strip_prefix("rpc ").unwrap();
2148                stripped[..stripped.find('(').unwrap()].trim()
2149            })
2150            .collect();
2151        let no_underscores: HashSet<_> = impl_list.iter().map(|x| x.replace('_', "")).collect();
2152        let mut not_implemented = vec![];
2153        for method in methods {
2154            if !no_underscores.contains(&method.to_lowercase()) {
2155                not_implemented.push(method);
2156            }
2157        }
2158        if !not_implemented.is_empty() {
2159            panic!(
2160                "The following RPC methods are not implemented by raw client: {not_implemented:?}"
2161            );
2162        }
2163    }
2164    #[test]
2165    fn verify_all_workflow_service_methods_implemented() {
2166        // This is less work than trying to hook into the codegen process
2167        let proto_def = include_str!(
2168            "../../protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto"
2169        );
2170        verify_methods(proto_def, ALL_IMPLEMENTED_WORKFLOW_SERVICE_RPCS);
2171    }
2172
2173    #[test]
2174    fn verify_all_operator_service_methods_implemented() {
2175        let proto_def = include_str!(
2176            "../../protos/protos/api_upstream/temporal/api/operatorservice/v1/service.proto"
2177        );
2178        verify_methods(proto_def, ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS);
2179    }
2180
2181    #[test]
2182    fn verify_all_cloud_service_methods_implemented() {
2183        let proto_def = include_str!(
2184            "../../protos/protos/api_cloud_upstream/temporal/api/cloud/cloudservice/v1/service.proto"
2185        );
2186        verify_methods(proto_def, ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS);
2187    }
2188
2189    #[test]
2190    fn verify_all_test_service_methods_implemented() {
2191        let proto_def = include_str!(
2192            "../../protos/protos/testsrv_upstream/temporal/api/testservice/v1/service.proto"
2193        );
2194        verify_methods(proto_def, ALL_IMPLEMENTED_TEST_SERVICE_RPCS);
2195    }
2196
2197    #[test]
2198    fn verify_all_health_service_methods_implemented() {
2199        let proto_def = include_str!("../../protos/protos/grpc/health/v1/health.proto");
2200        verify_methods(proto_def, ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS);
2201    }
2202
2203    #[tokio::test]
2204    async fn can_mock_services() {
2205        #[derive(Clone)]
2206        struct MyFakeServices {}
2207        impl RawGrpcCaller for MyFakeServices {}
2208        impl WorkflowService for MyFakeServices {
2209            fn list_namespaces(
2210                &mut self,
2211                _request: Request<ListNamespacesRequest>,
2212            ) -> BoxFuture<'_, Result<Response<ListNamespacesResponse>, Status>> {
2213                async {
2214                    Ok(Response::new(ListNamespacesResponse {
2215                        namespaces: vec![DescribeNamespaceResponse {
2216                            failover_version: 12345,
2217                            ..Default::default()
2218                        }],
2219                        ..Default::default()
2220                    }))
2221                }
2222                .boxed()
2223            }
2224        }
2225        impl OperatorService for MyFakeServices {}
2226        impl CloudService for MyFakeServices {}
2227        impl TestService for MyFakeServices {}
2228        // Health service isn't possible to create a default impl for.
2229        impl HealthService for MyFakeServices {
2230            fn check(
2231                &mut self,
2232                _request: tonic::Request<HealthCheckRequest>,
2233            ) -> BoxFuture<'_, Result<tonic::Response<HealthCheckResponse>, tonic::Status>>
2234            {
2235                todo!()
2236            }
2237            fn watch(
2238                &mut self,
2239                _request: tonic::Request<HealthCheckRequest>,
2240            ) -> BoxFuture<
2241                '_,
2242                Result<
2243                    tonic::Response<tonic::codec::Streaming<HealthCheckResponse>>,
2244                    tonic::Status,
2245                >,
2246            > {
2247                todo!()
2248            }
2249        }
2250        let mut mocked_client = TemporalServiceClient::from_services(
2251            Box::new(MyFakeServices {}),
2252            Box::new(MyFakeServices {}),
2253            Box::new(MyFakeServices {}),
2254            Box::new(MyFakeServices {}),
2255            Box::new(MyFakeServices {}),
2256        );
2257        let r = mocked_client
2258            .list_namespaces(ListNamespacesRequest::default().into_request())
2259            .await
2260            .unwrap();
2261        assert_eq!(r.into_inner().namespaces[0].failover_version, 12345);
2262    }
2263
2264    #[rstest::rstest]
2265    #[case::with_versioning(true)]
2266    #[case::without_versioning(false)]
2267    #[tokio::test]
2268    async fn eager_reservations_attach_deployment_options(#[case] use_worker_versioning: bool) {
2269        use crate::worker::{MockClientWorker, MockSlot};
2270        use temporalio_common::{
2271            protos::temporal::api::enums::v1::WorkerVersioningMode,
2272            worker::{WorkerDeploymentOptions, WorkerDeploymentVersion},
2273        };
2274
2275        let expected_mode = if use_worker_versioning {
2276            WorkerVersioningMode::Versioned
2277        } else {
2278            WorkerVersioningMode::Unversioned
2279        };
2280
2281        #[derive(Clone)]
2282        struct MyFakeServices {
2283            client_worker_set: Arc<ClientWorkerSet>,
2284            expected_mode: WorkerVersioningMode,
2285        }
2286        impl RawGrpcCaller for MyFakeServices {}
2287        impl RawClientProducer for MyFakeServices {
2288            fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
2289                Some(self.client_worker_set.clone())
2290            }
2291            fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
2292                Box::new(MyFakeWfClient {
2293                    expected_mode: self.expected_mode,
2294                })
2295            }
2296            fn operator_client(&mut self) -> Box<dyn OperatorService> {
2297                unimplemented!()
2298            }
2299            fn cloud_client(&mut self) -> Box<dyn CloudService> {
2300                unimplemented!()
2301            }
2302            fn test_client(&mut self) -> Box<dyn TestService> {
2303                unimplemented!()
2304            }
2305            fn health_client(&mut self) -> Box<dyn HealthService> {
2306                unimplemented!()
2307            }
2308        }
2309
2310        let deployment_opts = WorkerDeploymentOptions::new(
2311            WorkerDeploymentVersion::builder()
2312                .deployment_name("test-deployment".to_string())
2313                .build_id("test-build-123".to_string())
2314                .build(),
2315        )
2316        .use_worker_versioning(use_worker_versioning)
2317        .build();
2318
2319        let mut mock_provider = MockClientWorker::new();
2320        mock_provider
2321            .expect_namespace()
2322            .return_const("test-namespace".to_string());
2323        mock_provider
2324            .expect_task_queue()
2325            .return_const("test-task-queue".to_string());
2326        let mut mock_slot = MockSlot::new();
2327        mock_slot.expect_schedule_wft().returning(|_| Ok(()));
2328        mock_provider
2329            .expect_try_reserve_wft_slot()
2330            .return_once(|| Some(Box::new(mock_slot)));
2331        mock_provider
2332            .expect_deployment_options()
2333            .return_const(Some(deployment_opts.clone()));
2334        mock_provider.expect_heartbeat_enabled().return_const(false);
2335        let uuid = Uuid::new_v4();
2336        mock_provider
2337            .expect_worker_instance_key()
2338            .return_const(uuid);
2339        mock_provider
2340            .expect_worker_task_types()
2341            .return_const(WorkerTaskTypes {
2342                enable_workflows: true,
2343                enable_local_activities: true,
2344                enable_remote_activities: true,
2345                enable_nexus: true,
2346            });
2347
2348        let client_worker_set = Arc::new(ClientWorkerSet::new());
2349        client_worker_set
2350            .register_worker(Arc::new(mock_provider), true)
2351            .unwrap();
2352
2353        #[derive(Clone)]
2354        struct MyFakeWfClient {
2355            expected_mode: WorkerVersioningMode,
2356        }
2357        impl WorkflowService for MyFakeWfClient {
2358            fn start_workflow_execution(
2359                &mut self,
2360                request: tonic::Request<StartWorkflowExecutionRequest>,
2361            ) -> BoxFuture<'_, Result<tonic::Response<StartWorkflowExecutionResponse>, tonic::Status>>
2362            {
2363                let req = request.into_inner();
2364                let expected_mode = self.expected_mode;
2365
2366                assert!(
2367                    req.eager_worker_deployment_options.is_some(),
2368                    "eager_worker_deployment_options should be populated"
2369                );
2370
2371                let opts = req.eager_worker_deployment_options.as_ref().unwrap();
2372                assert_eq!(opts.deployment_name, "test-deployment");
2373                assert_eq!(opts.build_id, "test-build-123");
2374                assert_eq!(opts.worker_versioning_mode, expected_mode as i32);
2375
2376                async { Ok(Response::new(StartWorkflowExecutionResponse::default())) }.boxed()
2377            }
2378        }
2379
2380        let mut mfs = MyFakeServices {
2381            client_worker_set,
2382            expected_mode,
2383        };
2384
2385        // Create a request with eager execution enabled
2386        let req = StartWorkflowExecutionRequest {
2387            namespace: "test-namespace".to_string(),
2388            workflow_id: "test-wf-id".to_string(),
2389            workflow_type: Some(
2390                temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2391                    name: "test-workflow".to_string(),
2392                },
2393            ),
2394            task_queue: Some(TaskQueue {
2395                name: "test-task-queue".to_string(),
2396                kind: 0,
2397                normal_name: String::new(),
2398            }),
2399            request_eager_execution: true,
2400            ..Default::default()
2401        };
2402
2403        mfs.start_workflow_execution(req.into_request())
2404            .await
2405            .unwrap();
2406    }
2407
2408    /// Tests that Connection's RawClientProducer impl correctly provides worker info
2409    /// so that eager workflow start can reserve a slot and dispatch the WFT.
2410    #[tokio::test]
2411    async fn connection_eager_start_dispatches_wft() {
2412        use crate::{
2413            ConnectionOptions,
2414            callback_based::{CallbackBasedGrpcService, GrpcSuccessResponse},
2415            worker::{MockClientWorker, MockSlot},
2416        };
2417        use prost::Message;
2418        use std::sync::atomic::{AtomicBool, Ordering};
2419        use temporalio_common::protos::temporal::api::workflowservice::v1::PollWorkflowTaskQueueResponse;
2420
2421        let dispatched = Arc::new(AtomicBool::new(false));
2422        let dispatched_clone = dispatched.clone();
2423
2424        // Create a callback-based service that returns an eager_workflow_task in the response
2425        let service_override = CallbackBasedGrpcService {
2426            callback: Arc::new(|_req| {
2427                Box::pin(async {
2428                    let resp = StartWorkflowExecutionResponse {
2429                        run_id: "test-run-id".to_string(),
2430                        eager_workflow_task: Some(PollWorkflowTaskQueueResponse {
2431                            task_token: vec![1, 2, 3],
2432                            ..Default::default()
2433                        }),
2434                        ..Default::default()
2435                    };
2436                    let proto = resp.encode_to_vec();
2437                    Ok(GrpcSuccessResponse {
2438                        headers: Default::default(),
2439                        proto,
2440                    })
2441                })
2442            }),
2443        };
2444
2445        let opts = ConnectionOptions::new(url::Url::parse("http://localhost:7233").unwrap())
2446            .skip_get_system_info(true)
2447            .service_override(service_override)
2448            .dns_load_balancing(None)
2449            .build();
2450        let mut connection = crate::Connection::connect(opts).await.unwrap();
2451
2452        // Register a mock worker on the connection's worker set
2453        let mut mock_worker = MockClientWorker::new();
2454        mock_worker
2455            .expect_namespace()
2456            .return_const("default".to_string());
2457        mock_worker
2458            .expect_task_queue()
2459            .return_const("test-tq".to_string());
2460        mock_worker
2461            .expect_deployment_options()
2462            .return_const(None::<temporalio_common::worker::WorkerDeploymentOptions>);
2463        mock_worker.expect_heartbeat_enabled().return_const(false);
2464        let uuid = Uuid::new_v4();
2465        mock_worker.expect_worker_instance_key().return_const(uuid);
2466        mock_worker
2467            .expect_worker_task_types()
2468            .return_const(WorkerTaskTypes {
2469                enable_workflows: true,
2470                enable_local_activities: false,
2471                enable_remote_activities: false,
2472                enable_nexus: false,
2473            });
2474
2475        let mut mock_slot = MockSlot::new();
2476        mock_slot.expect_schedule_wft().returning(move |_| {
2477            dispatched_clone.store(true, Ordering::SeqCst);
2478            Ok(())
2479        });
2480        mock_worker
2481            .expect_try_reserve_wft_slot()
2482            .return_once(|| Some(Box::new(mock_slot)));
2483
2484        connection
2485            .workers()
2486            .register_worker(Arc::new(mock_worker), true)
2487            .unwrap();
2488
2489        // Make an eager start_workflow_execution call through Connection
2490        let req = StartWorkflowExecutionRequest {
2491            namespace: "default".to_string(),
2492            workflow_id: "test-wf".to_string(),
2493            workflow_type: Some(
2494                temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2495                    name: "test-workflow".to_string(),
2496                },
2497            ),
2498            task_queue: Some(TaskQueue {
2499                name: "test-tq".to_string(),
2500                kind: 0,
2501                normal_name: String::new(),
2502            }),
2503            request_eager_execution: true,
2504            ..Default::default()
2505        };
2506
2507        connection
2508            .start_workflow_execution(req.into_request())
2509            .await
2510            .unwrap();
2511
2512        assert!(
2513            dispatched.load(Ordering::SeqCst),
2514            "Eager workflow task should have been dispatched to the worker"
2515        );
2516    }
2517}