1use 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
39pub(crate) trait RawClientProducer {
41 fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>>;
44
45 fn workflow_client(&mut self) -> Box<dyn WorkflowService>;
47
48 fn operator_client(&mut self) -> Box<dyn OperatorService>;
50
51 fn cloud_client(&mut self) -> Box<dyn CloudService>;
53
54 fn test_client(&mut self) -> Box<dyn TestService>;
56
57 fn health_client(&mut self) -> Box<dyn HealthService>;
59}
60
61#[async_trait::async_trait]
64pub(crate) trait RawGrpcCaller: Send + Sync + 'static {
65 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_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
175fn 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
194fn 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#[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 pub fn new(inner: C) -> Self {
423 Self {
424 inner,
425 error_limits: Arc::new(RwLock::new(None)),
426 }
427 }
428
429 pub fn set_error_limits(&self, limits: Option<PayloadErrorLimits>) {
431 *self.error_limits.write() = limits;
432 }
433}
434
435impl<C: RawClientProducer> RawClientProducer for PayloadLimitsClient<C> {
436 fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
437 self.inner.get_workers_info()
438 }
439 fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
440 self.inner.workflow_client()
441 }
442 fn operator_client(&mut self) -> Box<dyn OperatorService> {
443 self.inner.operator_client()
444 }
445 fn cloud_client(&mut self) -> Box<dyn CloudService> {
446 self.inner.cloud_client()
447 }
448 fn test_client(&mut self) -> Box<dyn TestService> {
449 self.inner.test_client()
450 }
451 fn health_client(&mut self) -> Box<dyn HealthService> {
452 self.inner.health_client()
453 }
454}
455
456#[async_trait::async_trait]
457impl<C> RawGrpcCaller for PayloadLimitsClient<C>
458where
459 C: RawGrpcCaller + Clone + Sync + 'static,
460{
461 async fn call<F, Req, Resp>(
462 &mut self,
463 call_name: &'static str,
464 callfn: F,
465 mut req: Request<Req>,
466 ) -> Result<Response<Resp>, Status>
467 where
468 Req: Clone + Unpin + Send + Sync + 'static,
469 Resp: Send + 'static,
470 F: FnMut(Request<Req>) -> BoxFuture<'static, Result<Response<Resp>, Status>>,
471 F: Send + Sync + Unpin + 'static,
472 {
473 let limits = *self.error_limits.read();
474 if let Some(limits) = limits {
475 req.extensions_mut().insert(limits);
476 }
477 self.inner.call(call_name, callfn, req).await
478 }
479}
480
481#[derive(Clone, Debug)]
482pub(super) struct AttachMetricLabels {
483 pub(super) labels: Vec<MetricKeyValue>,
484 pub(super) normal_task_queue: Option<String>,
485 pub(super) sticky_task_queue: Option<String>,
486}
487impl AttachMetricLabels {
488 pub(super) fn new(kvs: impl Into<Vec<MetricKeyValue>>) -> Self {
489 Self {
490 labels: kvs.into(),
491 normal_task_queue: None,
492 sticky_task_queue: None,
493 }
494 }
495 pub(super) fn namespace(ns: impl Into<String>) -> Self {
496 AttachMetricLabels::new(vec![namespace_kv(ns.into())])
497 }
498 pub(super) fn task_q(&mut self, tq: Option<TaskQueue>) -> &mut Self {
499 if let Some(tq) = tq {
500 if !tq.normal_name.is_empty() {
501 self.sticky_task_queue = Some(tq.name);
502 self.normal_task_queue = Some(tq.normal_name);
503 } else {
504 self.normal_task_queue = Some(tq.name);
505 }
506 }
507 self
508 }
509 pub(super) fn task_q_str(&mut self, tq: impl Into<String>) -> &mut Self {
510 self.normal_task_queue = Some(tq.into());
511 self
512 }
513}
514
515#[derive(Copy, Clone, Debug)]
518pub(super) struct IsUserLongPoll;
519
520macro_rules! proxy_def {
521 ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty, defaults) => {
522 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
523 fn $method(
524 &mut self,
525 _request: tonic::Request<$req>,
526 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
527 async { Ok(tonic::Response::new(<$resp>::default())) }.boxed()
528 }
529 };
530 ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty) => {
531 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
532 fn $method(
533 &mut self,
534 _request: tonic::Request<$req>,
535 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>>;
536 };
537}
538
539macro_rules! proxy_impl {
552 ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty $(, $closure:expr)?) => {
553 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
554 fn $method(
555 &mut self,
556 #[allow(unused_mut)]
557 mut request: tonic::Request<$req>,
558 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
559 $( type_closure_arg(&mut request, $closure); )*
560 let mut self_clone = self.clone();
561 #[allow(unused_mut)]
562 let fact = move |mut req: tonic::Request<$req>| {
563 let mut c = self_clone.$client_meth();
564 async move { c.$method(req).await }.boxed()
565 };
566 self.call(stringify!($method), fact, request)
567 }
568 };
569 ($client_type:tt, $client_meth:ident, $method:ident, $req:ty, $resp:ty,
570 $closure_request:expr, $closure_before:expr, $closure_after:expr) => {
571 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
572 fn $method(
573 &mut self,
574 mut request: tonic::Request<$req>,
575 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
576 type_closure_arg(&mut request, $closure_request);
577 let workers_info = self.get_workers_info();
578 let mut self_clone = self.clone();
579 #[allow(unused_mut)]
580 let fact = move |mut req: tonic::Request<$req>| {
581 let data = type_closure_two_arg(&mut req, workers_info.clone(), $closure_before);
582 let mut c = self_clone.$client_meth();
583 async move {
584 type_closure_two_arg(c.$method(req).await, data, $closure_after)
585 }.boxed()
586 };
587 self.call(stringify!($method), fact, request)
588 }
589 };
590 ($client_type:tt, $method:ident, $req:ty, $resp:ty $(, $closure:expr)?) => {
591 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
592 fn $method(
593 &mut self,
594 #[allow(unused_mut)]
595 mut request: tonic::Request<$req>,
596 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
597 $( type_closure_arg(&mut request, $closure); )*
598 async move { <$client_type<_>>::$method(self, request).await }.boxed()
599 }
600 };
601 ($client_type:tt, $method:ident, $req:ty, $resp:ty,
602 $closure_request:expr, $closure_before:expr, $closure_after:expr) => {
603 #[doc = concat!("See [", stringify!($client_type), "::", stringify!($method), "]")]
604 fn $method(
605 &mut self,
606 mut request: tonic::Request<$req>,
607 ) -> BoxFuture<'_, Result<tonic::Response<$resp>, tonic::Status>> {
608 type_closure_arg(&mut request, $closure_request);
609 let data = type_closure_two_arg(&mut request, Option::<Arc<ClientWorkerSet>>::None,
610 $closure_before);
611 async move {
612 type_closure_two_arg(<$client_type<_>>::$method(self, request).await,
613 data, $closure_after)
614 }.boxed()
615 }
616 };
617}
618macro_rules! proxier_impl {
619 ($trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident;
620 [$( proxy_def!($($def_args:tt)*); )*];
621 $(($method:ident, $req:ty, $resp:ty
622 $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
623 #[cfg(test)]
624 const $impl_list_name: &'static [&'static str] = &[$(stringify!($method)),*];
625
626 #[doc = concat!("Trait version of [", stringify!($client_type), "]")]
627 pub trait $trait_name: Send + Sync + DynClone
628 {
629 $( proxy_def!($($def_args)*); )*
630 }
631 dyn_clone::clone_trait_object!($trait_name);
632
633 #[allow(deprecated)]
637 impl<RC> $trait_name for RC
638 where
639 RC: RawGrpcCaller + RawClientProducer + Clone + Unpin,
640 {
641 $(
642 proxy_impl!($client_type, $client_meth, $method, $req, $resp
643 $(,$closure $(,$closure_before, $closure_after)*)*);
644 )*
645 }
646
647 impl<T: Send + Sync + 'static> RawGrpcCaller for $client_type<T> {}
648
649 #[allow(deprecated)]
650 impl<T> $trait_name for $client_type<T>
651 where
652 T: GrpcService<Body> + Clone + Send + Sync + 'static,
653 T::ResponseBody: tonic::codegen::Body<Data = tonic::codegen::Bytes> + Send + 'static,
654 T::Error: Into<tonic::codegen::StdError>,
655 <T::ResponseBody as tonic::codegen::Body>::Error: Into<tonic::codegen::StdError> + Send,
656 <T as tonic::client::GrpcService<Body>>::Future: Send
657 {
658 $(
659 proxy_impl!($client_type, $method, $req, $resp
660 $(,$closure $(,$closure_before, $closure_after)*)*);
661 )*
662 }
663 };
664}
665
666macro_rules! proxier {
667 ( $trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident;
668 $(($method:ident, $req:ty, $resp:ty
669 $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
670 proxier_impl!($trait_name; $impl_list_name; $client_type; $client_meth;
671 [$(proxy_def!($client_type, $client_meth, $method, $req, $resp);)*];
672 $(($method, $req, $resp $(, $closure $(, $closure_before, $closure_after)?)?);)*);
673 };
674 ( $trait_name:ident; $impl_list_name:ident; $client_type:tt; $client_meth:ident; defaults;
675 $(($method:ident, $req:ty, $resp:ty
676 $(, $closure:expr $(, $closure_before:expr, $closure_after:expr)?)? );)* ) => {
677 proxier_impl!($trait_name; $impl_list_name; $client_type; $client_meth;
678 [$(proxy_def!($client_type, $client_meth, $method, $req, $resp, defaults);)*];
679 $(($method, $req, $resp $(, $closure $(, $closure_before, $closure_after)?)?);)*);
680 };
681}
682
683macro_rules! namespaced_request {
684 ($req:ident) => {{
685 let ns_str = $req.get_ref().namespace.clone();
686 $req.metadata_mut().insert(
688 TEMPORAL_NAMESPACE_HEADER_KEY,
689 ns_str.parse().unwrap_or_else(|e| {
690 warn!("Unable to parse namespace for header: {e:?}");
691 AsciiMetadataValue::from_static("")
692 }),
693 );
694 AttachMetricLabels::namespace(ns_str)
696 }};
697}
698
699fn type_closure_arg<T, R>(arg: T, f: impl FnOnce(T) -> R) -> R {
701 f(arg)
702}
703
704fn type_closure_two_arg<T, R, S>(arg1: R, arg2: T, f: impl FnOnce(R, T) -> S) -> S {
705 f(arg1, arg2)
706}
707
708proxier! {
709 WorkflowService; ALL_IMPLEMENTED_WORKFLOW_SERVICE_RPCS; WorkflowServiceClient; workflow_client; defaults;
710 (
711 register_namespace,
712 RegisterNamespaceRequest,
713 RegisterNamespaceResponse,
714 |r| {
715 let labels = namespaced_request!(r);
716 r.extensions_mut().insert(labels);
717 }
718 );
719 (
720 describe_namespace,
721 DescribeNamespaceRequest,
722 DescribeNamespaceResponse,
723 |r| {
724 let labels = namespaced_request!(r);
725 r.extensions_mut().insert(labels);
726 }
727 );
728 (
729 list_namespaces,
730 ListNamespacesRequest,
731 ListNamespacesResponse
732 );
733 (
734 update_namespace,
735 UpdateNamespaceRequest,
736 UpdateNamespaceResponse,
737 |r| {
738 let labels = namespaced_request!(r);
739 r.extensions_mut().insert(labels);
740 }
741 );
742 (
743 deprecate_namespace,
744 DeprecateNamespaceRequest,
745 DeprecateNamespaceResponse,
746 |r| {
747 let labels = namespaced_request!(r);
748 r.extensions_mut().insert(labels);
749 }
750 );
751 (
752 start_workflow_execution,
753 StartWorkflowExecutionRequest,
754 StartWorkflowExecutionResponse,
755 |r| {
756 let mut labels = namespaced_request!(r);
757 labels.task_q(r.get_ref().task_queue.clone());
758 r.extensions_mut().insert(labels);
759 },
760 |r, workers| {
761 if let Some(workers) = workers {
762 let mut slot: Option<Box<dyn Slot + Send>> = None;
763 let req_mut = r.get_mut();
764 if req_mut.request_eager_execution {
765 let namespace = req_mut.namespace.clone();
766 let task_queue = req_mut.task_queue.as_ref()
767 .map(|tq| tq.name.clone()).unwrap_or_default();
768 match workers.try_reserve_wft_slot(namespace, task_queue) {
769 Some(reservation) => {
770 if let Some(opts) = reservation.deployment_options {
772 req_mut.eager_worker_deployment_options = Some(temporalio_common::protos::temporal::api::deployment::v1::WorkerDeploymentOptions {
773 deployment_name: opts.version.deployment_name,
774 build_id: opts.version.build_id,
775 worker_versioning_mode: if opts.use_worker_versioning {
776 temporalio_common::protos::temporal::api::enums::v1::WorkerVersioningMode::Versioned.into()
777 } else {
778 temporalio_common::protos::temporal::api::enums::v1::WorkerVersioningMode::Unversioned.into()
779 },
780 }); }
781 slot = Some(reservation.slot);
782 }
783 None => req_mut.request_eager_execution = false
784 }
785 }
786 slot
787 } else {
788 None
789 }
790 },
791 |resp, slot| {
792 if let Some(s) = slot
793 && let Ok(response) = resp.as_ref()
794 && let Some(task) = response.get_ref().clone().eager_workflow_task
795 && let Err(e) = s.schedule_wft(task) {
796 warn!(details = ?e, "Eager workflow task rejected by worker.");
799 }
800 resp
801 }
802 );
803 (
804 get_workflow_execution_history,
805 GetWorkflowExecutionHistoryRequest,
806 GetWorkflowExecutionHistoryResponse,
807 |r| {
808 let labels = namespaced_request!(r);
809 r.extensions_mut().insert(labels);
810 if r.get_ref().wait_new_event {
811 r.extensions_mut().insert(IsUserLongPoll);
812 }
813 }
814 );
815 (
816 get_workflow_execution_history_reverse,
817 GetWorkflowExecutionHistoryReverseRequest,
818 GetWorkflowExecutionHistoryReverseResponse,
819 |r| {
820 let labels = namespaced_request!(r);
821 r.extensions_mut().insert(labels);
822 }
823 );
824 (
825 poll_workflow_task_queue,
826 PollWorkflowTaskQueueRequest,
827 PollWorkflowTaskQueueResponse,
828 |r| {
829 let mut labels = namespaced_request!(r);
830 labels.task_q(r.get_ref().task_queue.clone());
831 r.extensions_mut().insert(labels);
832 }
833 );
834 (
835 respond_workflow_task_completed,
836 RespondWorkflowTaskCompletedRequest,
837 RespondWorkflowTaskCompletedResponse,
838 |r| {
839 let labels = namespaced_request!(r);
840 r.extensions_mut().insert(labels);
841 }
842 );
843 (
844 respond_workflow_task_failed,
845 RespondWorkflowTaskFailedRequest,
846 RespondWorkflowTaskFailedResponse,
847 |r| {
848 let labels = namespaced_request!(r);
849 r.extensions_mut().insert(labels);
850 }
851 );
852 (
853 poll_activity_task_queue,
854 PollActivityTaskQueueRequest,
855 PollActivityTaskQueueResponse,
856 |r| {
857 let mut labels = namespaced_request!(r);
858 labels.task_q(r.get_ref().task_queue.clone());
859 r.extensions_mut().insert(labels);
860 }
861 );
862 (
863 record_activity_task_heartbeat,
864 RecordActivityTaskHeartbeatRequest,
865 RecordActivityTaskHeartbeatResponse,
866 |r| {
867 let labels = namespaced_request!(r);
868 r.extensions_mut().insert(labels);
869 }
870 );
871 (
872 record_activity_task_heartbeat_by_id,
873 RecordActivityTaskHeartbeatByIdRequest,
874 RecordActivityTaskHeartbeatByIdResponse,
875 |r| {
876 let labels = namespaced_request!(r);
877 r.extensions_mut().insert(labels);
878 }
879 );
880 (
881 respond_activity_task_completed,
882 RespondActivityTaskCompletedRequest,
883 RespondActivityTaskCompletedResponse,
884 |r| {
885 let labels = namespaced_request!(r);
886 r.extensions_mut().insert(labels);
887 }
888 );
889 (
890 respond_activity_task_completed_by_id,
891 RespondActivityTaskCompletedByIdRequest,
892 RespondActivityTaskCompletedByIdResponse,
893 |r| {
894 let labels = namespaced_request!(r);
895 r.extensions_mut().insert(labels);
896 }
897 );
898
899 (
900 respond_activity_task_failed,
901 RespondActivityTaskFailedRequest,
902 RespondActivityTaskFailedResponse,
903 |r| {
904 let labels = namespaced_request!(r);
905 r.extensions_mut().insert(labels);
906 }
907 );
908 (
909 respond_activity_task_failed_by_id,
910 RespondActivityTaskFailedByIdRequest,
911 RespondActivityTaskFailedByIdResponse,
912 |r| {
913 let labels = namespaced_request!(r);
914 r.extensions_mut().insert(labels);
915 }
916 );
917 (
918 respond_activity_task_canceled,
919 RespondActivityTaskCanceledRequest,
920 RespondActivityTaskCanceledResponse,
921 |r| {
922 let labels = namespaced_request!(r);
923 r.extensions_mut().insert(labels);
924 }
925 );
926 (
927 respond_activity_task_canceled_by_id,
928 RespondActivityTaskCanceledByIdRequest,
929 RespondActivityTaskCanceledByIdResponse,
930 |r| {
931 let labels = namespaced_request!(r);
932 r.extensions_mut().insert(labels);
933 }
934 );
935 (
936 request_cancel_workflow_execution,
937 RequestCancelWorkflowExecutionRequest,
938 RequestCancelWorkflowExecutionResponse,
939 |r| {
940 let labels = namespaced_request!(r);
941 r.extensions_mut().insert(labels);
942 }
943 );
944 (
945 signal_workflow_execution,
946 SignalWorkflowExecutionRequest,
947 SignalWorkflowExecutionResponse,
948 |r| {
949 let labels = namespaced_request!(r);
950 r.extensions_mut().insert(labels);
951 }
952 );
953 (
954 signal_with_start_workflow_execution,
955 SignalWithStartWorkflowExecutionRequest,
956 SignalWithStartWorkflowExecutionResponse,
957 |r| {
958 let mut labels = namespaced_request!(r);
959 labels.task_q(r.get_ref().task_queue.clone());
960 r.extensions_mut().insert(labels);
961 }
962 );
963 (
964 reset_workflow_execution,
965 ResetWorkflowExecutionRequest,
966 ResetWorkflowExecutionResponse,
967 |r| {
968 let labels = namespaced_request!(r);
969 r.extensions_mut().insert(labels);
970 }
971 );
972 (
973 terminate_workflow_execution,
974 TerminateWorkflowExecutionRequest,
975 TerminateWorkflowExecutionResponse,
976 |r| {
977 let labels = namespaced_request!(r);
978 r.extensions_mut().insert(labels);
979 }
980 );
981 (
982 delete_workflow_execution,
983 DeleteWorkflowExecutionRequest,
984 DeleteWorkflowExecutionResponse,
985 |r| {
986 let labels = namespaced_request!(r);
987 r.extensions_mut().insert(labels);
988 }
989 );
990 (
991 list_open_workflow_executions,
992 ListOpenWorkflowExecutionsRequest,
993 ListOpenWorkflowExecutionsResponse,
994 |r| {
995 let labels = namespaced_request!(r);
996 r.extensions_mut().insert(labels);
997 }
998 );
999 (
1000 list_closed_workflow_executions,
1001 ListClosedWorkflowExecutionsRequest,
1002 ListClosedWorkflowExecutionsResponse,
1003 |r| {
1004 let labels = namespaced_request!(r);
1005 r.extensions_mut().insert(labels);
1006 }
1007 );
1008 (
1009 list_workflow_executions,
1010 ListWorkflowExecutionsRequest,
1011 ListWorkflowExecutionsResponse,
1012 |r| {
1013 let labels = namespaced_request!(r);
1014 r.extensions_mut().insert(labels);
1015 }
1016 );
1017 (
1018 list_archived_workflow_executions,
1019 ListArchivedWorkflowExecutionsRequest,
1020 ListArchivedWorkflowExecutionsResponse,
1021 |r| {
1022 let labels = namespaced_request!(r);
1023 r.extensions_mut().insert(labels);
1024 }
1025 );
1026 (
1027 scan_workflow_executions,
1028 ScanWorkflowExecutionsRequest,
1029 ScanWorkflowExecutionsResponse,
1030 |r| {
1031 let labels = namespaced_request!(r);
1032 r.extensions_mut().insert(labels);
1033 }
1034 );
1035 (
1036 count_workflow_executions,
1037 CountWorkflowExecutionsRequest,
1038 CountWorkflowExecutionsResponse,
1039 |r| {
1040 let labels = namespaced_request!(r);
1041 r.extensions_mut().insert(labels);
1042 }
1043 );
1044 (
1045 create_workflow_rule,
1046 CreateWorkflowRuleRequest,
1047 CreateWorkflowRuleResponse,
1048 |r| {
1049 let labels = namespaced_request!(r);
1050 r.extensions_mut().insert(labels);
1051 }
1052 );
1053 (
1054 describe_workflow_rule,
1055 DescribeWorkflowRuleRequest,
1056 DescribeWorkflowRuleResponse,
1057 |r| {
1058 let labels = namespaced_request!(r);
1059 r.extensions_mut().insert(labels);
1060 }
1061 );
1062 (
1063 delete_workflow_rule,
1064 DeleteWorkflowRuleRequest,
1065 DeleteWorkflowRuleResponse,
1066 |r| {
1067 let labels = namespaced_request!(r);
1068 r.extensions_mut().insert(labels);
1069 }
1070 );
1071 (
1072 list_workflow_rules,
1073 ListWorkflowRulesRequest,
1074 ListWorkflowRulesResponse,
1075 |r| {
1076 let labels = namespaced_request!(r);
1077 r.extensions_mut().insert(labels);
1078 }
1079 );
1080 (
1081 trigger_workflow_rule,
1082 TriggerWorkflowRuleRequest,
1083 TriggerWorkflowRuleResponse,
1084 |r| {
1085 let labels = namespaced_request!(r);
1086 r.extensions_mut().insert(labels);
1087 }
1088 );
1089 (
1090 get_search_attributes,
1091 GetSearchAttributesRequest,
1092 GetSearchAttributesResponse
1093 );
1094 (
1095 respond_query_task_completed,
1096 RespondQueryTaskCompletedRequest,
1097 RespondQueryTaskCompletedResponse,
1098 |r| {
1099 let labels = namespaced_request!(r);
1100 r.extensions_mut().insert(labels);
1101 }
1102 );
1103 (
1104 reset_sticky_task_queue,
1105 ResetStickyTaskQueueRequest,
1106 ResetStickyTaskQueueResponse,
1107 |r| {
1108 let labels = namespaced_request!(r);
1109 r.extensions_mut().insert(labels);
1110 }
1111 );
1112 (
1113 query_workflow,
1114 QueryWorkflowRequest,
1115 QueryWorkflowResponse,
1116 |r| {
1117 let labels = namespaced_request!(r);
1118 r.extensions_mut().insert(labels);
1119 }
1120 );
1121 (
1122 describe_workflow_execution,
1123 DescribeWorkflowExecutionRequest,
1124 DescribeWorkflowExecutionResponse,
1125 |r| {
1126 let labels = namespaced_request!(r);
1127 r.extensions_mut().insert(labels);
1128 }
1129 );
1130 (
1131 describe_task_queue,
1132 DescribeTaskQueueRequest,
1133 DescribeTaskQueueResponse,
1134 |r| {
1135 let mut labels = namespaced_request!(r);
1136 labels.task_q(r.get_ref().task_queue.clone());
1137 r.extensions_mut().insert(labels);
1138 }
1139 );
1140 (
1141 get_cluster_info,
1142 GetClusterInfoRequest,
1143 GetClusterInfoResponse
1144 );
1145 (
1146 get_system_info,
1147 GetSystemInfoRequest,
1148 GetSystemInfoResponse
1149 );
1150 (
1151 list_task_queue_partitions,
1152 ListTaskQueuePartitionsRequest,
1153 ListTaskQueuePartitionsResponse,
1154 |r| {
1155 let mut labels = namespaced_request!(r);
1156 labels.task_q(r.get_ref().task_queue.clone());
1157 r.extensions_mut().insert(labels);
1158 }
1159 );
1160 (
1161 create_schedule,
1162 CreateScheduleRequest,
1163 CreateScheduleResponse,
1164 |r| {
1165 let labels = namespaced_request!(r);
1166 r.extensions_mut().insert(labels);
1167 }
1168 );
1169 (
1170 describe_schedule,
1171 DescribeScheduleRequest,
1172 DescribeScheduleResponse,
1173 |r| {
1174 let labels = namespaced_request!(r);
1175 r.extensions_mut().insert(labels);
1176 }
1177 );
1178 (
1179 update_schedule,
1180 UpdateScheduleRequest,
1181 UpdateScheduleResponse,
1182 |r| {
1183 let labels = namespaced_request!(r);
1184 r.extensions_mut().insert(labels);
1185 }
1186 );
1187 (
1188 patch_schedule,
1189 PatchScheduleRequest,
1190 PatchScheduleResponse,
1191 |r| {
1192 let labels = namespaced_request!(r);
1193 r.extensions_mut().insert(labels);
1194 }
1195 );
1196 (
1197 list_schedule_matching_times,
1198 ListScheduleMatchingTimesRequest,
1199 ListScheduleMatchingTimesResponse,
1200 |r| {
1201 let labels = namespaced_request!(r);
1202 r.extensions_mut().insert(labels);
1203 }
1204 );
1205 (
1206 delete_schedule,
1207 DeleteScheduleRequest,
1208 DeleteScheduleResponse,
1209 |r| {
1210 let labels = namespaced_request!(r);
1211 r.extensions_mut().insert(labels);
1212 }
1213 );
1214 (
1215 list_schedules,
1216 ListSchedulesRequest,
1217 ListSchedulesResponse,
1218 |r| {
1219 let labels = namespaced_request!(r);
1220 r.extensions_mut().insert(labels);
1221 }
1222 );
1223 (
1224 count_schedules,
1225 CountSchedulesRequest,
1226 CountSchedulesResponse,
1227 |r| {
1228 let labels = namespaced_request!(r);
1229 r.extensions_mut().insert(labels);
1230 }
1231 );
1232 (
1233 update_worker_build_id_compatibility,
1234 UpdateWorkerBuildIdCompatibilityRequest,
1235 UpdateWorkerBuildIdCompatibilityResponse,
1236 |r| {
1237 let mut labels = namespaced_request!(r);
1238 labels.task_q_str(r.get_ref().task_queue.clone());
1239 r.extensions_mut().insert(labels);
1240 }
1241 );
1242 (
1243 get_worker_build_id_compatibility,
1244 GetWorkerBuildIdCompatibilityRequest,
1245 GetWorkerBuildIdCompatibilityResponse,
1246 |r| {
1247 let mut labels = namespaced_request!(r);
1248 labels.task_q_str(r.get_ref().task_queue.clone());
1249 r.extensions_mut().insert(labels);
1250 }
1251 );
1252 (
1253 get_worker_task_reachability,
1254 GetWorkerTaskReachabilityRequest,
1255 GetWorkerTaskReachabilityResponse,
1256 |r| {
1257 let labels = namespaced_request!(r);
1258 r.extensions_mut().insert(labels);
1259 }
1260 );
1261 (
1262 update_workflow_execution,
1263 UpdateWorkflowExecutionRequest,
1264 UpdateWorkflowExecutionResponse,
1265 |r| {
1266 let labels = namespaced_request!(r);
1267 let exts = r.extensions_mut();
1268 exts.insert(labels);
1269 exts.insert(IsUserLongPoll);
1270 }
1271 );
1272 (
1273 poll_workflow_execution_update,
1274 PollWorkflowExecutionUpdateRequest,
1275 PollWorkflowExecutionUpdateResponse,
1276 |r| {
1277 let labels = namespaced_request!(r);
1278 r.extensions_mut().insert(labels);
1279 }
1280 );
1281 (
1282 start_batch_operation,
1283 StartBatchOperationRequest,
1284 StartBatchOperationResponse,
1285 |r| {
1286 let labels = namespaced_request!(r);
1287 r.extensions_mut().insert(labels);
1288 }
1289 );
1290 (
1291 stop_batch_operation,
1292 StopBatchOperationRequest,
1293 StopBatchOperationResponse,
1294 |r| {
1295 let labels = namespaced_request!(r);
1296 r.extensions_mut().insert(labels);
1297 }
1298 );
1299 (
1300 describe_batch_operation,
1301 DescribeBatchOperationRequest,
1302 DescribeBatchOperationResponse,
1303 |r| {
1304 let labels = namespaced_request!(r);
1305 r.extensions_mut().insert(labels);
1306 }
1307 );
1308 (
1309 describe_deployment,
1310 DescribeDeploymentRequest,
1311 DescribeDeploymentResponse,
1312 |r| {
1313 let labels = namespaced_request!(r);
1314 r.extensions_mut().insert(labels);
1315 }
1316 );
1317 (
1318 list_batch_operations,
1319 ListBatchOperationsRequest,
1320 ListBatchOperationsResponse,
1321 |r| {
1322 let labels = namespaced_request!(r);
1323 r.extensions_mut().insert(labels);
1324 }
1325 );
1326 (
1327 list_deployments,
1328 ListDeploymentsRequest,
1329 ListDeploymentsResponse,
1330 |r| {
1331 let labels = namespaced_request!(r);
1332 r.extensions_mut().insert(labels);
1333 }
1334 );
1335 (
1336 execute_multi_operation,
1337 ExecuteMultiOperationRequest,
1338 ExecuteMultiOperationResponse,
1339 |r| {
1340 let labels = namespaced_request!(r);
1341 r.extensions_mut().insert(labels);
1342 }
1343 );
1344 (
1345 get_current_deployment,
1346 GetCurrentDeploymentRequest,
1347 GetCurrentDeploymentResponse,
1348 |r| {
1349 let labels = namespaced_request!(r);
1350 r.extensions_mut().insert(labels);
1351 }
1352 );
1353 (
1354 get_deployment_reachability,
1355 GetDeploymentReachabilityRequest,
1356 GetDeploymentReachabilityResponse,
1357 |r| {
1358 let labels = namespaced_request!(r);
1359 r.extensions_mut().insert(labels);
1360 }
1361 );
1362 (
1363 get_worker_versioning_rules,
1364 GetWorkerVersioningRulesRequest,
1365 GetWorkerVersioningRulesResponse,
1366 |r| {
1367 let mut labels = namespaced_request!(r);
1368 labels.task_q_str(&r.get_ref().task_queue);
1369 r.extensions_mut().insert(labels);
1370 }
1371 );
1372 (
1373 update_worker_versioning_rules,
1374 UpdateWorkerVersioningRulesRequest,
1375 UpdateWorkerVersioningRulesResponse,
1376 |r| {
1377 let mut labels = namespaced_request!(r);
1378 labels.task_q_str(&r.get_ref().task_queue);
1379 r.extensions_mut().insert(labels);
1380 }
1381 );
1382 (
1383 poll_nexus_task_queue,
1384 PollNexusTaskQueueRequest,
1385 PollNexusTaskQueueResponse,
1386 |r| {
1387 let mut labels = namespaced_request!(r);
1388 labels.task_q(r.get_ref().task_queue.clone());
1389 r.extensions_mut().insert(labels);
1390 }
1391 );
1392 (
1393 respond_nexus_task_completed,
1394 RespondNexusTaskCompletedRequest,
1395 RespondNexusTaskCompletedResponse,
1396 |r| {
1397 let labels = namespaced_request!(r);
1398 r.extensions_mut().insert(labels);
1399 }
1400 );
1401 (
1402 respond_nexus_task_failed,
1403 RespondNexusTaskFailedRequest,
1404 RespondNexusTaskFailedResponse,
1405 |r| {
1406 let labels = namespaced_request!(r);
1407 r.extensions_mut().insert(labels);
1408 }
1409 );
1410 (
1411 set_current_deployment,
1412 SetCurrentDeploymentRequest,
1413 SetCurrentDeploymentResponse,
1414 |r| {
1415 let labels = namespaced_request!(r);
1416 r.extensions_mut().insert(labels);
1417 }
1418 );
1419 (
1420 shutdown_worker,
1421 ShutdownWorkerRequest,
1422 ShutdownWorkerResponse,
1423 |r| {
1424 let labels = namespaced_request!(r);
1425 r.extensions_mut().insert(labels);
1426 }
1427 );
1428 (
1429 update_activity_options,
1430 UpdateActivityOptionsRequest,
1431 UpdateActivityOptionsResponse,
1432 |r| {
1433 let labels = namespaced_request!(r);
1434 r.extensions_mut().insert(labels);
1435 }
1436 );
1437 (
1438 pause_activity,
1439 PauseActivityRequest,
1440 PauseActivityResponse,
1441 |r| {
1442 let labels = namespaced_request!(r);
1443 r.extensions_mut().insert(labels);
1444 }
1445 );
1446 (
1447 unpause_activity,
1448 UnpauseActivityRequest,
1449 UnpauseActivityResponse,
1450 |r| {
1451 let labels = namespaced_request!(r);
1452 r.extensions_mut().insert(labels);
1453 }
1454 );
1455 (
1456 update_workflow_execution_options,
1457 UpdateWorkflowExecutionOptionsRequest,
1458 UpdateWorkflowExecutionOptionsResponse,
1459 |r| {
1460 let labels = namespaced_request!(r);
1461 r.extensions_mut().insert(labels);
1462 }
1463 );
1464 (
1465 reset_activity,
1466 ResetActivityRequest,
1467 ResetActivityResponse,
1468 |r| {
1469 let labels = namespaced_request!(r);
1470 r.extensions_mut().insert(labels);
1471 }
1472 );
1473 (
1474 delete_worker_deployment,
1475 DeleteWorkerDeploymentRequest,
1476 DeleteWorkerDeploymentResponse,
1477 |r| {
1478 let labels = namespaced_request!(r);
1479 r.extensions_mut().insert(labels);
1480 }
1481 );
1482 (
1483 delete_worker_deployment_version,
1484 DeleteWorkerDeploymentVersionRequest,
1485 DeleteWorkerDeploymentVersionResponse,
1486 |r| {
1487 let labels = namespaced_request!(r);
1488 r.extensions_mut().insert(labels);
1489 }
1490 );
1491 (
1492 describe_worker_deployment,
1493 DescribeWorkerDeploymentRequest,
1494 DescribeWorkerDeploymentResponse,
1495 |r| {
1496 let labels = namespaced_request!(r);
1497 r.extensions_mut().insert(labels);
1498 }
1499 );
1500 (
1501 describe_worker_deployment_version,
1502 DescribeWorkerDeploymentVersionRequest,
1503 DescribeWorkerDeploymentVersionResponse,
1504 |r| {
1505 let labels = namespaced_request!(r);
1506 r.extensions_mut().insert(labels);
1507 }
1508 );
1509 (
1510 list_worker_deployments,
1511 ListWorkerDeploymentsRequest,
1512 ListWorkerDeploymentsResponse,
1513 |r| {
1514 let labels = namespaced_request!(r);
1515 r.extensions_mut().insert(labels);
1516 }
1517 );
1518 (
1519 set_worker_deployment_current_version,
1520 SetWorkerDeploymentCurrentVersionRequest,
1521 SetWorkerDeploymentCurrentVersionResponse,
1522 |r| {
1523 let labels = namespaced_request!(r);
1524 r.extensions_mut().insert(labels);
1525 }
1526 );
1527 (
1528 set_worker_deployment_ramping_version,
1529 SetWorkerDeploymentRampingVersionRequest,
1530 SetWorkerDeploymentRampingVersionResponse,
1531 |r| {
1532 let labels = namespaced_request!(r);
1533 r.extensions_mut().insert(labels);
1534 }
1535 );
1536 (
1537 update_worker_deployment_version_metadata,
1538 UpdateWorkerDeploymentVersionMetadataRequest,
1539 UpdateWorkerDeploymentVersionMetadataResponse,
1540 |r| {
1541 let labels = namespaced_request!(r);
1542 r.extensions_mut().insert(labels);
1543 }
1544 );
1545 (
1546 list_workers,
1547 ListWorkersRequest,
1548 ListWorkersResponse,
1549 |r| {
1550 let labels = namespaced_request!(r);
1551 r.extensions_mut().insert(labels);
1552 }
1553 );
1554 (
1555 count_workers,
1556 CountWorkersRequest,
1557 CountWorkersResponse,
1558 |r| {
1559 let labels = namespaced_request!(r);
1560 r.extensions_mut().insert(labels);
1561 }
1562 );
1563 (
1564 record_worker_heartbeat,
1565 RecordWorkerHeartbeatRequest,
1566 RecordWorkerHeartbeatResponse,
1567 |r| {
1568 let labels = namespaced_request!(r);
1569 r.extensions_mut().insert(labels);
1570 }
1571 );
1572 (
1573 update_task_queue_config,
1574 UpdateTaskQueueConfigRequest,
1575 UpdateTaskQueueConfigResponse,
1576 |r| {
1577 let mut labels = namespaced_request!(r);
1578 labels.task_q_str(r.get_ref().task_queue.clone());
1579 r.extensions_mut().insert(labels);
1580 }
1581 );
1582 (
1583 fetch_worker_config,
1584 FetchWorkerConfigRequest,
1585 FetchWorkerConfigResponse,
1586 |r| {
1587 let labels = namespaced_request!(r);
1588 r.extensions_mut().insert(labels);
1589 }
1590 );
1591 (
1592 update_worker_config,
1593 UpdateWorkerConfigRequest,
1594 UpdateWorkerConfigResponse,
1595 |r| {
1596 let labels = namespaced_request!(r);
1597 r.extensions_mut().insert(labels);
1598 }
1599 );
1600 (
1601 describe_worker,
1602 DescribeWorkerRequest,
1603 DescribeWorkerResponse,
1604 |r| {
1605 let labels = namespaced_request!(r);
1606 r.extensions_mut().insert(labels);
1607 }
1608 );
1609 (
1610 set_worker_deployment_manager,
1611 SetWorkerDeploymentManagerRequest,
1612 SetWorkerDeploymentManagerResponse,
1613 |r| {
1614 let labels = namespaced_request!(r);
1615 r.extensions_mut().insert(labels);
1616 }
1617 );
1618 (
1619 pause_workflow_execution,
1620 PauseWorkflowExecutionRequest,
1621 PauseWorkflowExecutionResponse,
1622 |r| {
1623 let labels = namespaced_request!(r);
1624 r.extensions_mut().insert(labels);
1625 }
1626 );
1627 (
1628 unpause_workflow_execution,
1629 UnpauseWorkflowExecutionRequest,
1630 UnpauseWorkflowExecutionResponse,
1631 |r| {
1632 let labels = namespaced_request!(r);
1633 r.extensions_mut().insert(labels);
1634 }
1635 );
1636 (
1637 start_activity_execution,
1638 StartActivityExecutionRequest,
1639 StartActivityExecutionResponse,
1640 |r| {
1641 let labels = namespaced_request!(r);
1642 r.extensions_mut().insert(labels);
1643 }
1644 );
1645 (
1646 describe_activity_execution,
1647 DescribeActivityExecutionRequest,
1648 DescribeActivityExecutionResponse,
1649 |r| {
1650 let labels = namespaced_request!(r);
1651 r.extensions_mut().insert(labels);
1652 }
1653 );
1654 (
1655 poll_activity_execution,
1656 PollActivityExecutionRequest,
1657 PollActivityExecutionResponse,
1658 |r| {
1659 let labels = namespaced_request!(r);
1660 r.extensions_mut().insert(labels);
1661 }
1662 );
1663 (
1664 list_activity_executions,
1665 ListActivityExecutionsRequest,
1666 ListActivityExecutionsResponse,
1667 |r| {
1668 let labels = namespaced_request!(r);
1669 r.extensions_mut().insert(labels);
1670 }
1671 );
1672 (
1673 count_activity_executions,
1674 CountActivityExecutionsRequest,
1675 CountActivityExecutionsResponse,
1676 |r| {
1677 let labels = namespaced_request!(r);
1678 r.extensions_mut().insert(labels);
1679 }
1680 );
1681 (
1682 request_cancel_activity_execution,
1683 RequestCancelActivityExecutionRequest,
1684 RequestCancelActivityExecutionResponse,
1685 |r| {
1686 let labels = namespaced_request!(r);
1687 r.extensions_mut().insert(labels);
1688 }
1689 );
1690 (
1691 terminate_activity_execution,
1692 TerminateActivityExecutionRequest,
1693 TerminateActivityExecutionResponse,
1694 |r| {
1695 let labels = namespaced_request!(r);
1696 r.extensions_mut().insert(labels);
1697 }
1698 );
1699 (
1700 delete_activity_execution,
1701 DeleteActivityExecutionRequest,
1702 DeleteActivityExecutionResponse,
1703 |r| {
1704 let labels = namespaced_request!(r);
1705 r.extensions_mut().insert(labels);
1706 }
1707 );
1708 (
1709 pause_activity_execution,
1710 PauseActivityExecutionRequest,
1711 PauseActivityExecutionResponse,
1712 |r| {
1713 let labels = namespaced_request!(r);
1714 r.extensions_mut().insert(labels);
1715 }
1716 );
1717 (
1718 unpause_activity_execution,
1719 UnpauseActivityExecutionRequest,
1720 UnpauseActivityExecutionResponse,
1721 |r| {
1722 let labels = namespaced_request!(r);
1723 r.extensions_mut().insert(labels);
1724 }
1725 );
1726 (
1727 reset_activity_execution,
1728 ResetActivityExecutionRequest,
1729 ResetActivityExecutionResponse,
1730 |r| {
1731 let labels = namespaced_request!(r);
1732 r.extensions_mut().insert(labels);
1733 }
1734 );
1735 (
1736 update_activity_execution_options,
1737 UpdateActivityExecutionOptionsRequest,
1738 UpdateActivityExecutionOptionsResponse,
1739 |r| {
1740 let labels = namespaced_request!(r);
1741 r.extensions_mut().insert(labels);
1742 }
1743 );
1744 (
1745 count_nexus_operation_executions,
1746 CountNexusOperationExecutionsRequest,
1747 CountNexusOperationExecutionsResponse,
1748 |r| {
1749 let labels = namespaced_request!(r);
1750 r.extensions_mut().insert(labels);
1751 }
1752 );
1753 (
1754 create_worker_deployment,
1755 CreateWorkerDeploymentRequest,
1756 CreateWorkerDeploymentResponse,
1757 |r| {
1758 let labels = namespaced_request!(r);
1759 r.extensions_mut().insert(labels);
1760 }
1761 );
1762 (
1763 create_worker_deployment_version,
1764 CreateWorkerDeploymentVersionRequest,
1765 CreateWorkerDeploymentVersionResponse,
1766 |r| {
1767 let labels = namespaced_request!(r);
1768 r.extensions_mut().insert(labels);
1769 }
1770 );
1771 (
1772 delete_nexus_operation_execution,
1773 DeleteNexusOperationExecutionRequest,
1774 DeleteNexusOperationExecutionResponse,
1775 |r| {
1776 let labels = namespaced_request!(r);
1777 r.extensions_mut().insert(labels);
1778 }
1779 );
1780 (
1781 describe_nexus_operation_execution,
1782 DescribeNexusOperationExecutionRequest,
1783 DescribeNexusOperationExecutionResponse,
1784 |r| {
1785 let labels = namespaced_request!(r);
1786 r.extensions_mut().insert(labels);
1787 }
1788 );
1789 (
1790 list_nexus_operation_executions,
1791 ListNexusOperationExecutionsRequest,
1792 ListNexusOperationExecutionsResponse,
1793 |r| {
1794 let labels = namespaced_request!(r);
1795 r.extensions_mut().insert(labels);
1796 }
1797 );
1798 (
1799 poll_nexus_operation_execution,
1800 PollNexusOperationExecutionRequest,
1801 PollNexusOperationExecutionResponse,
1802 |r| {
1803 let labels = namespaced_request!(r);
1804 r.extensions_mut().insert(labels);
1805 }
1806 );
1807 (
1808 poll_workflow_execution_time_skipping,
1809 PollWorkflowExecutionTimeSkippingRequest,
1810 PollWorkflowExecutionTimeSkippingResponse,
1811 |r| {
1812 let labels = namespaced_request!(r);
1813 r.extensions_mut().insert(labels);
1814 r.extensions_mut().insert(IsUserLongPoll);
1815 }
1816 );
1817 (
1818 request_cancel_nexus_operation_execution,
1819 RequestCancelNexusOperationExecutionRequest,
1820 RequestCancelNexusOperationExecutionResponse,
1821 |r| {
1822 let labels = namespaced_request!(r);
1823 r.extensions_mut().insert(labels);
1824 }
1825 );
1826 (
1827 start_nexus_operation_execution,
1828 StartNexusOperationExecutionRequest,
1829 StartNexusOperationExecutionResponse,
1830 |r| {
1831 let labels = namespaced_request!(r);
1832 r.extensions_mut().insert(labels);
1833 }
1834 );
1835 (
1836 terminate_nexus_operation_execution,
1837 TerminateNexusOperationExecutionRequest,
1838 TerminateNexusOperationExecutionResponse,
1839 |r| {
1840 let labels = namespaced_request!(r);
1841 r.extensions_mut().insert(labels);
1842 }
1843 );
1844 (
1845 update_worker_deployment_version_compute_config,
1846 UpdateWorkerDeploymentVersionComputeConfigRequest,
1847 UpdateWorkerDeploymentVersionComputeConfigResponse,
1848 |r| {
1849 let labels = namespaced_request!(r);
1850 r.extensions_mut().insert(labels);
1851 }
1852 );
1853 (
1854 validate_worker_deployment_version_compute_config,
1855 ValidateWorkerDeploymentVersionComputeConfigRequest,
1856 ValidateWorkerDeploymentVersionComputeConfigResponse,
1857 |r| {
1858 let labels = namespaced_request!(r);
1859 r.extensions_mut().insert(labels);
1860 }
1861 );
1862}
1863
1864proxier! {
1865 OperatorService; ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS; OperatorServiceClient; operator_client; defaults;
1866 (add_search_attributes, AddSearchAttributesRequest, AddSearchAttributesResponse);
1867 (remove_search_attributes, RemoveSearchAttributesRequest, RemoveSearchAttributesResponse);
1868 (list_search_attributes, ListSearchAttributesRequest, ListSearchAttributesResponse);
1869 (delete_namespace, DeleteNamespaceRequest, DeleteNamespaceResponse,
1870 |r| {
1871 let labels = namespaced_request!(r);
1872 r.extensions_mut().insert(labels);
1873 }
1874 );
1875 (add_or_update_remote_cluster, AddOrUpdateRemoteClusterRequest, AddOrUpdateRemoteClusterResponse);
1876 (remove_remote_cluster, RemoveRemoteClusterRequest, RemoveRemoteClusterResponse);
1877 (list_clusters, ListClustersRequest, ListClustersResponse);
1878 (get_nexus_endpoint, GetNexusEndpointRequest, GetNexusEndpointResponse);
1879 (create_nexus_endpoint, CreateNexusEndpointRequest, CreateNexusEndpointResponse);
1880 (update_nexus_endpoint, UpdateNexusEndpointRequest, UpdateNexusEndpointResponse);
1881 (delete_nexus_endpoint, DeleteNexusEndpointRequest, DeleteNexusEndpointResponse);
1882 (list_nexus_endpoints, ListNexusEndpointsRequest, ListNexusEndpointsResponse);
1883}
1884
1885proxier! {
1886 CloudService; ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS; CloudServiceClient; cloud_client; defaults;
1887 (get_users, cloudreq::GetUsersRequest, cloudreq::GetUsersResponse);
1888 (get_user, cloudreq::GetUserRequest, cloudreq::GetUserResponse);
1889 (create_user, cloudreq::CreateUserRequest, cloudreq::CreateUserResponse);
1890 (update_user, cloudreq::UpdateUserRequest, cloudreq::UpdateUserResponse);
1891 (delete_user, cloudreq::DeleteUserRequest, cloudreq::DeleteUserResponse);
1892 (set_user_namespace_access, cloudreq::SetUserNamespaceAccessRequest, cloudreq::SetUserNamespaceAccessResponse);
1893 (get_async_operation, cloudreq::GetAsyncOperationRequest, cloudreq::GetAsyncOperationResponse);
1894 (create_namespace, cloudreq::CreateNamespaceRequest, cloudreq::CreateNamespaceResponse);
1895 (get_namespaces, cloudreq::GetNamespacesRequest, cloudreq::GetNamespacesResponse);
1896 (get_namespace, cloudreq::GetNamespaceRequest, cloudreq::GetNamespaceResponse,
1897 |r| {
1898 let labels = namespaced_request!(r);
1899 r.extensions_mut().insert(labels);
1900 }
1901 );
1902 (update_namespace, cloudreq::UpdateNamespaceRequest, cloudreq::UpdateNamespaceResponse,
1903 |r| {
1904 let labels = namespaced_request!(r);
1905 r.extensions_mut().insert(labels);
1906 }
1907 );
1908 (rename_custom_search_attribute, cloudreq::RenameCustomSearchAttributeRequest, cloudreq::RenameCustomSearchAttributeResponse);
1909 (delete_namespace, cloudreq::DeleteNamespaceRequest, cloudreq::DeleteNamespaceResponse,
1910 |r| {
1911 let labels = namespaced_request!(r);
1912 r.extensions_mut().insert(labels);
1913 }
1914 );
1915 (failover_namespace_region, cloudreq::FailoverNamespaceRegionRequest, cloudreq::FailoverNamespaceRegionResponse);
1916 (add_namespace_region, cloudreq::AddNamespaceRegionRequest, cloudreq::AddNamespaceRegionResponse);
1917 (delete_namespace_region, cloudreq::DeleteNamespaceRegionRequest, cloudreq::DeleteNamespaceRegionResponse);
1918 (get_regions, cloudreq::GetRegionsRequest, cloudreq::GetRegionsResponse);
1919 (get_region, cloudreq::GetRegionRequest, cloudreq::GetRegionResponse);
1920 (get_api_keys, cloudreq::GetApiKeysRequest, cloudreq::GetApiKeysResponse);
1921 (get_api_key, cloudreq::GetApiKeyRequest, cloudreq::GetApiKeyResponse);
1922 (create_api_key, cloudreq::CreateApiKeyRequest, cloudreq::CreateApiKeyResponse);
1923 (update_api_key, cloudreq::UpdateApiKeyRequest, cloudreq::UpdateApiKeyResponse);
1924 (delete_api_key, cloudreq::DeleteApiKeyRequest, cloudreq::DeleteApiKeyResponse);
1925 (get_nexus_endpoints, cloudreq::GetNexusEndpointsRequest, cloudreq::GetNexusEndpointsResponse);
1926 (get_nexus_endpoint, cloudreq::GetNexusEndpointRequest, cloudreq::GetNexusEndpointResponse);
1927 (create_nexus_endpoint, cloudreq::CreateNexusEndpointRequest, cloudreq::CreateNexusEndpointResponse);
1928 (update_nexus_endpoint, cloudreq::UpdateNexusEndpointRequest, cloudreq::UpdateNexusEndpointResponse);
1929 (delete_nexus_endpoint, cloudreq::DeleteNexusEndpointRequest, cloudreq::DeleteNexusEndpointResponse);
1930 (get_user_groups, cloudreq::GetUserGroupsRequest, cloudreq::GetUserGroupsResponse);
1931 (get_user_group, cloudreq::GetUserGroupRequest, cloudreq::GetUserGroupResponse);
1932 (create_user_group, cloudreq::CreateUserGroupRequest, cloudreq::CreateUserGroupResponse);
1933 (update_user_group, cloudreq::UpdateUserGroupRequest, cloudreq::UpdateUserGroupResponse);
1934 (delete_user_group, cloudreq::DeleteUserGroupRequest, cloudreq::DeleteUserGroupResponse);
1935 (add_user_group_member, cloudreq::AddUserGroupMemberRequest, cloudreq::AddUserGroupMemberResponse);
1936 (remove_user_group_member, cloudreq::RemoveUserGroupMemberRequest, cloudreq::RemoveUserGroupMemberResponse);
1937 (get_user_group_members, cloudreq::GetUserGroupMembersRequest, cloudreq::GetUserGroupMembersResponse);
1938 (set_user_group_namespace_access, cloudreq::SetUserGroupNamespaceAccessRequest, cloudreq::SetUserGroupNamespaceAccessResponse);
1939 (create_service_account, cloudreq::CreateServiceAccountRequest, cloudreq::CreateServiceAccountResponse);
1940 (get_service_account, cloudreq::GetServiceAccountRequest, cloudreq::GetServiceAccountResponse);
1941 (get_service_accounts, cloudreq::GetServiceAccountsRequest, cloudreq::GetServiceAccountsResponse);
1942 (update_service_account, cloudreq::UpdateServiceAccountRequest, cloudreq::UpdateServiceAccountResponse);
1943 (delete_service_account, cloudreq::DeleteServiceAccountRequest, cloudreq::DeleteServiceAccountResponse);
1944 (get_usage, cloudreq::GetUsageRequest, cloudreq::GetUsageResponse);
1945 (get_account, cloudreq::GetAccountRequest, cloudreq::GetAccountResponse);
1946 (update_account, cloudreq::UpdateAccountRequest, cloudreq::UpdateAccountResponse);
1947 (create_namespace_export_sink, cloudreq::CreateNamespaceExportSinkRequest, cloudreq::CreateNamespaceExportSinkResponse);
1948 (get_namespace_export_sink, cloudreq::GetNamespaceExportSinkRequest, cloudreq::GetNamespaceExportSinkResponse);
1949 (get_namespace_export_sinks, cloudreq::GetNamespaceExportSinksRequest, cloudreq::GetNamespaceExportSinksResponse);
1950 (update_namespace_export_sink, cloudreq::UpdateNamespaceExportSinkRequest, cloudreq::UpdateNamespaceExportSinkResponse);
1951 (delete_namespace_export_sink, cloudreq::DeleteNamespaceExportSinkRequest, cloudreq::DeleteNamespaceExportSinkResponse);
1952 (validate_namespace_export_sink, cloudreq::ValidateNamespaceExportSinkRequest, cloudreq::ValidateNamespaceExportSinkResponse);
1953 (update_namespace_tags, cloudreq::UpdateNamespaceTagsRequest, cloudreq::UpdateNamespaceTagsResponse);
1954 (create_connectivity_rule, cloudreq::CreateConnectivityRuleRequest, cloudreq::CreateConnectivityRuleResponse);
1955 (get_connectivity_rule, cloudreq::GetConnectivityRuleRequest, cloudreq::GetConnectivityRuleResponse);
1956 (get_connectivity_rules, cloudreq::GetConnectivityRulesRequest, cloudreq::GetConnectivityRulesResponse);
1957 (delete_connectivity_rule, cloudreq::DeleteConnectivityRuleRequest, cloudreq::DeleteConnectivityRuleResponse);
1958 (set_service_account_namespace_access, cloudreq::SetServiceAccountNamespaceAccessRequest, cloudreq::SetServiceAccountNamespaceAccessResponse);
1959 (validate_account_audit_log_sink, cloudreq::ValidateAccountAuditLogSinkRequest, cloudreq::ValidateAccountAuditLogSinkResponse);
1960 (get_current_identity, cloudreq::GetCurrentIdentityRequest, cloudreq::GetCurrentIdentityResponse);
1961 (get_audit_logs, cloudreq::GetAuditLogsRequest, cloudreq::GetAuditLogsResponse);
1962 (create_account_audit_log_sink, cloudreq::CreateAccountAuditLogSinkRequest, cloudreq::CreateAccountAuditLogSinkResponse);
1963 (get_account_audit_log_sink, cloudreq::GetAccountAuditLogSinkRequest, cloudreq::GetAccountAuditLogSinkResponse);
1964 (get_account_audit_log_sinks, cloudreq::GetAccountAuditLogSinksRequest, cloudreq::GetAccountAuditLogSinksResponse);
1965 (update_account_audit_log_sink, cloudreq::UpdateAccountAuditLogSinkRequest, cloudreq::UpdateAccountAuditLogSinkResponse);
1966 (delete_account_audit_log_sink, cloudreq::DeleteAccountAuditLogSinkRequest, cloudreq::DeleteAccountAuditLogSinkResponse);
1967 (get_namespace_capacity_info, cloudreq::GetNamespaceCapacityInfoRequest, cloudreq::GetNamespaceCapacityInfoResponse);
1968 (create_billing_report, cloudreq::CreateBillingReportRequest, cloudreq::CreateBillingReportResponse);
1969 (get_billing_report, cloudreq::GetBillingReportRequest, cloudreq::GetBillingReportResponse);
1970 (get_custom_roles, cloudreq::GetCustomRolesRequest, cloudreq::GetCustomRolesResponse);
1971 (get_custom_role, cloudreq::GetCustomRoleRequest, cloudreq::GetCustomRoleResponse);
1972 (create_custom_role, cloudreq::CreateCustomRoleRequest, cloudreq::CreateCustomRoleResponse);
1973 (update_custom_role, cloudreq::UpdateCustomRoleRequest, cloudreq::UpdateCustomRoleResponse);
1974 (delete_custom_role, cloudreq::DeleteCustomRoleRequest, cloudreq::DeleteCustomRoleResponse);
1975 (get_user_namespace_assignments, cloudreq::GetUserNamespaceAssignmentsRequest, cloudreq::GetUserNamespaceAssignmentsResponse);
1976 (get_service_account_namespace_assignments, cloudreq::GetServiceAccountNamespaceAssignmentsRequest, cloudreq::GetServiceAccountNamespaceAssignmentsResponse);
1977 (get_user_group_namespace_assignments, cloudreq::GetUserGroupNamespaceAssignmentsRequest, cloudreq::GetUserGroupNamespaceAssignmentsResponse);
1978}
1979
1980proxier! {
1981 TestService; ALL_IMPLEMENTED_TEST_SERVICE_RPCS; TestServiceClient; test_client; defaults;
1982 (lock_time_skipping, LockTimeSkippingRequest, LockTimeSkippingResponse);
1983 (unlock_time_skipping, UnlockTimeSkippingRequest, UnlockTimeSkippingResponse);
1984 (sleep, SleepRequest, SleepResponse);
1985 (sleep_until, SleepUntilRequest, SleepResponse);
1986 (unlock_time_skipping_with_sleep, SleepRequest, SleepResponse);
1987 (get_current_time, (), GetCurrentTimeResponse);
1988}
1989
1990proxier! {
1991 HealthService; ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS; HealthClient; health_client;
1992 (check, HealthCheckRequest, HealthCheckResponse);
1993 (watch, HealthCheckRequest, tonic::codec::Streaming<HealthCheckResponse>);
1994}
1995
1996#[cfg(test)]
1997mod tests {
1998 use super::*;
1999 use crate::{ClientOptions, ConnectionOptions};
2000 use std::collections::HashSet;
2001 use temporalio_common::{
2002 protos::temporal::api::{
2003 operatorservice::v1::DeleteNamespaceRequest, workflowservice::v1::ListNamespacesRequest,
2004 },
2005 worker::WorkerTaskTypes,
2006 };
2007 use tonic::IntoRequest;
2008 use url::Url;
2009 use uuid::Uuid;
2010
2011 #[test]
2012 fn payload_limits_warn_only_vs_error() {
2013 use temporalio_common::protos::temporal::api::{
2014 common::v1::{Payload, Payloads},
2015 workflowservice::v1::StartWorkflowExecutionRequest,
2016 };
2017 let big = Payloads {
2018 payloads: vec![Payload {
2019 data: vec![0u8; 1000],
2020 ..Default::default()
2021 }],
2022 };
2023 let new_req = || {
2024 StartWorkflowExecutionRequest {
2025 input: Some(big.clone()),
2026 ..Default::default()
2027 }
2028 .into_request()
2029 };
2030
2031 assert!(validate_request_payload_limits(&new_req(), 1, 1).is_ok());
2033
2034 let mut req = new_req();
2036 req.extensions_mut()
2037 .insert(PayloadErrorLimits { blob: 10, memo: 10 });
2038 let err = validate_request_payload_limits(&req, 1, 1).unwrap_err();
2039 let violation =
2040 crate::payload_limit_violation_from(&err).expect("violation carried on status");
2041 assert_eq!(violation.path, "input");
2042 assert_eq!(
2043 violation.class,
2044 temporalio_common::payload_limits::LimitClass::Blob
2045 );
2046 assert!(violation.size > violation.limit);
2047
2048 let mut req = new_req();
2050 req.extensions_mut()
2051 .insert(PayloadErrorLimits { blob: 0, memo: 0 });
2052 assert!(validate_request_payload_limits(&req, 1, 1).is_ok());
2053
2054 assert!(validate_request_payload_limits(&new_req(), 0, 0).is_ok());
2056 }
2057
2058 #[allow(dead_code)]
2060 async fn raw_client_retry_compiles() {
2061 let opts = ConnectionOptions::new(Url::parse("http://localhost:7233").unwrap())
2062 .client_name("test")
2063 .client_version("0.0.0")
2064 .build();
2065 let connection = Connection::connect(opts).await.unwrap();
2066 let mut client = Client::new(connection, ClientOptions::new("default").build()).unwrap();
2067
2068 let list_ns_req = ListNamespacesRequest::default();
2069 let wf_client = client.workflow_client();
2070 let fact = move |req| {
2071 let mut c = wf_client.clone();
2072 async move { c.list_namespaces(req).await }.boxed()
2073 };
2074 client
2075 .call("whatever", fact, Request::new(list_ns_req.clone()))
2076 .await
2077 .unwrap();
2078
2079 let op_del_ns_req = DeleteNamespaceRequest::default();
2081 let op_client = client.operator_client();
2082 let fact = move |req| {
2083 let mut c = op_client.clone();
2084 async move { c.delete_namespace(req).await }.boxed()
2085 };
2086 client
2087 .call("whatever", fact, Request::new(op_del_ns_req.clone()))
2088 .await
2089 .unwrap();
2090
2091 let cloud_del_ns_req = cloudreq::DeleteNamespaceRequest::default();
2093 let cloud_client = client.cloud_client();
2094 let fact = move |req| {
2095 let mut c = cloud_client.clone();
2096 async move { c.delete_namespace(req).await }.boxed()
2097 };
2098 client
2099 .call("whatever", fact, Request::new(cloud_del_ns_req.clone()))
2100 .await
2101 .unwrap();
2102
2103 client
2105 .list_namespaces(list_ns_req.into_request())
2106 .await
2107 .unwrap();
2108 OperatorService::delete_namespace(&mut client, op_del_ns_req.into_request())
2110 .await
2111 .unwrap();
2112 CloudService::delete_namespace(&mut client, cloud_del_ns_req.into_request())
2113 .await
2114 .unwrap();
2115 client.get_current_time(().into_request()).await.unwrap();
2116 client
2117 .check(HealthCheckRequest::default().into_request())
2118 .await
2119 .unwrap();
2120 }
2121
2122 fn verify_methods(proto_def_str: &str, impl_list: &[&str]) {
2123 let methods: Vec<_> = proto_def_str
2124 .lines()
2125 .map(|l| l.trim())
2126 .filter(|l| l.starts_with("rpc"))
2127 .map(|l| {
2128 let stripped = l.strip_prefix("rpc ").unwrap();
2129 stripped[..stripped.find('(').unwrap()].trim()
2130 })
2131 .collect();
2132 let no_underscores: HashSet<_> = impl_list.iter().map(|x| x.replace('_', "")).collect();
2133 let mut not_implemented = vec![];
2134 for method in methods {
2135 if !no_underscores.contains(&method.to_lowercase()) {
2136 not_implemented.push(method);
2137 }
2138 }
2139 if !not_implemented.is_empty() {
2140 panic!(
2141 "The following RPC methods are not implemented by raw client: {not_implemented:?}"
2142 );
2143 }
2144 }
2145 #[test]
2146 fn verify_all_workflow_service_methods_implemented() {
2147 let proto_def = include_str!(
2149 "../../protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto"
2150 );
2151 verify_methods(proto_def, ALL_IMPLEMENTED_WORKFLOW_SERVICE_RPCS);
2152 }
2153
2154 #[test]
2155 fn verify_all_operator_service_methods_implemented() {
2156 let proto_def = include_str!(
2157 "../../protos/protos/api_upstream/temporal/api/operatorservice/v1/service.proto"
2158 );
2159 verify_methods(proto_def, ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS);
2160 }
2161
2162 #[test]
2163 fn verify_all_cloud_service_methods_implemented() {
2164 let proto_def = include_str!(
2165 "../../protos/protos/api_cloud_upstream/temporal/api/cloud/cloudservice/v1/service.proto"
2166 );
2167 verify_methods(proto_def, ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS);
2168 }
2169
2170 #[test]
2171 fn verify_all_test_service_methods_implemented() {
2172 let proto_def = include_str!(
2173 "../../protos/protos/testsrv_upstream/temporal/api/testservice/v1/service.proto"
2174 );
2175 verify_methods(proto_def, ALL_IMPLEMENTED_TEST_SERVICE_RPCS);
2176 }
2177
2178 #[test]
2179 fn verify_all_health_service_methods_implemented() {
2180 let proto_def = include_str!("../../protos/protos/grpc/health/v1/health.proto");
2181 verify_methods(proto_def, ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS);
2182 }
2183
2184 #[tokio::test]
2185 async fn can_mock_services() {
2186 #[derive(Clone)]
2187 struct MyFakeServices {}
2188 impl RawGrpcCaller for MyFakeServices {}
2189 impl WorkflowService for MyFakeServices {
2190 fn list_namespaces(
2191 &mut self,
2192 _request: Request<ListNamespacesRequest>,
2193 ) -> BoxFuture<'_, Result<Response<ListNamespacesResponse>, Status>> {
2194 async {
2195 Ok(Response::new(ListNamespacesResponse {
2196 namespaces: vec![DescribeNamespaceResponse {
2197 failover_version: 12345,
2198 ..Default::default()
2199 }],
2200 ..Default::default()
2201 }))
2202 }
2203 .boxed()
2204 }
2205 }
2206 impl OperatorService for MyFakeServices {}
2207 impl CloudService for MyFakeServices {}
2208 impl TestService for MyFakeServices {}
2209 impl HealthService for MyFakeServices {
2211 fn check(
2212 &mut self,
2213 _request: tonic::Request<HealthCheckRequest>,
2214 ) -> BoxFuture<'_, Result<tonic::Response<HealthCheckResponse>, tonic::Status>>
2215 {
2216 todo!()
2217 }
2218 fn watch(
2219 &mut self,
2220 _request: tonic::Request<HealthCheckRequest>,
2221 ) -> BoxFuture<
2222 '_,
2223 Result<
2224 tonic::Response<tonic::codec::Streaming<HealthCheckResponse>>,
2225 tonic::Status,
2226 >,
2227 > {
2228 todo!()
2229 }
2230 }
2231 let mut mocked_client = TemporalServiceClient::from_services(
2232 Box::new(MyFakeServices {}),
2233 Box::new(MyFakeServices {}),
2234 Box::new(MyFakeServices {}),
2235 Box::new(MyFakeServices {}),
2236 Box::new(MyFakeServices {}),
2237 );
2238 let r = mocked_client
2239 .list_namespaces(ListNamespacesRequest::default().into_request())
2240 .await
2241 .unwrap();
2242 assert_eq!(r.into_inner().namespaces[0].failover_version, 12345);
2243 }
2244
2245 #[rstest::rstest]
2246 #[case::with_versioning(true)]
2247 #[case::without_versioning(false)]
2248 #[tokio::test]
2249 async fn eager_reservations_attach_deployment_options(#[case] use_worker_versioning: bool) {
2250 use crate::worker::{MockClientWorker, MockSlot};
2251 use temporalio_common::{
2252 protos::temporal::api::enums::v1::WorkerVersioningMode,
2253 worker::{WorkerDeploymentOptions, WorkerDeploymentVersion},
2254 };
2255
2256 let expected_mode = if use_worker_versioning {
2257 WorkerVersioningMode::Versioned
2258 } else {
2259 WorkerVersioningMode::Unversioned
2260 };
2261
2262 #[derive(Clone)]
2263 struct MyFakeServices {
2264 client_worker_set: Arc<ClientWorkerSet>,
2265 expected_mode: WorkerVersioningMode,
2266 }
2267 impl RawGrpcCaller for MyFakeServices {}
2268 impl RawClientProducer for MyFakeServices {
2269 fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
2270 Some(self.client_worker_set.clone())
2271 }
2272 fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
2273 Box::new(MyFakeWfClient {
2274 expected_mode: self.expected_mode,
2275 })
2276 }
2277 fn operator_client(&mut self) -> Box<dyn OperatorService> {
2278 unimplemented!()
2279 }
2280 fn cloud_client(&mut self) -> Box<dyn CloudService> {
2281 unimplemented!()
2282 }
2283 fn test_client(&mut self) -> Box<dyn TestService> {
2284 unimplemented!()
2285 }
2286 fn health_client(&mut self) -> Box<dyn HealthService> {
2287 unimplemented!()
2288 }
2289 }
2290
2291 let deployment_opts = WorkerDeploymentOptions {
2292 version: WorkerDeploymentVersion {
2293 deployment_name: "test-deployment".to_string(),
2294 build_id: "test-build-123".to_string(),
2295 },
2296 use_worker_versioning,
2297 default_versioning_behavior: None,
2298 };
2299
2300 let mut mock_provider = MockClientWorker::new();
2301 mock_provider
2302 .expect_namespace()
2303 .return_const("test-namespace".to_string());
2304 mock_provider
2305 .expect_task_queue()
2306 .return_const("test-task-queue".to_string());
2307 let mut mock_slot = MockSlot::new();
2308 mock_slot.expect_schedule_wft().returning(|_| Ok(()));
2309 mock_provider
2310 .expect_try_reserve_wft_slot()
2311 .return_once(|| Some(Box::new(mock_slot)));
2312 mock_provider
2313 .expect_deployment_options()
2314 .return_const(Some(deployment_opts.clone()));
2315 mock_provider.expect_heartbeat_enabled().return_const(false);
2316 let uuid = Uuid::new_v4();
2317 mock_provider
2318 .expect_worker_instance_key()
2319 .return_const(uuid);
2320 mock_provider
2321 .expect_worker_task_types()
2322 .return_const(WorkerTaskTypes {
2323 enable_workflows: true,
2324 enable_local_activities: true,
2325 enable_remote_activities: true,
2326 enable_nexus: true,
2327 });
2328
2329 let client_worker_set = Arc::new(ClientWorkerSet::new());
2330 client_worker_set
2331 .register_worker(Arc::new(mock_provider), true)
2332 .unwrap();
2333
2334 #[derive(Clone)]
2335 struct MyFakeWfClient {
2336 expected_mode: WorkerVersioningMode,
2337 }
2338 impl WorkflowService for MyFakeWfClient {
2339 fn start_workflow_execution(
2340 &mut self,
2341 request: tonic::Request<StartWorkflowExecutionRequest>,
2342 ) -> BoxFuture<'_, Result<tonic::Response<StartWorkflowExecutionResponse>, tonic::Status>>
2343 {
2344 let req = request.into_inner();
2345 let expected_mode = self.expected_mode;
2346
2347 assert!(
2348 req.eager_worker_deployment_options.is_some(),
2349 "eager_worker_deployment_options should be populated"
2350 );
2351
2352 let opts = req.eager_worker_deployment_options.as_ref().unwrap();
2353 assert_eq!(opts.deployment_name, "test-deployment");
2354 assert_eq!(opts.build_id, "test-build-123");
2355 assert_eq!(opts.worker_versioning_mode, expected_mode as i32);
2356
2357 async { Ok(Response::new(StartWorkflowExecutionResponse::default())) }.boxed()
2358 }
2359 }
2360
2361 let mut mfs = MyFakeServices {
2362 client_worker_set,
2363 expected_mode,
2364 };
2365
2366 let req = StartWorkflowExecutionRequest {
2368 namespace: "test-namespace".to_string(),
2369 workflow_id: "test-wf-id".to_string(),
2370 workflow_type: Some(
2371 temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2372 name: "test-workflow".to_string(),
2373 },
2374 ),
2375 task_queue: Some(TaskQueue {
2376 name: "test-task-queue".to_string(),
2377 kind: 0,
2378 normal_name: String::new(),
2379 }),
2380 request_eager_execution: true,
2381 ..Default::default()
2382 };
2383
2384 mfs.start_workflow_execution(req.into_request())
2385 .await
2386 .unwrap();
2387 }
2388
2389 #[tokio::test]
2392 async fn connection_eager_start_dispatches_wft() {
2393 use crate::{
2394 ConnectionOptions,
2395 callback_based::{CallbackBasedGrpcService, GrpcSuccessResponse},
2396 worker::{MockClientWorker, MockSlot},
2397 };
2398 use prost::Message;
2399 use std::sync::atomic::{AtomicBool, Ordering};
2400 use temporalio_common::protos::temporal::api::workflowservice::v1::PollWorkflowTaskQueueResponse;
2401
2402 let dispatched = Arc::new(AtomicBool::new(false));
2403 let dispatched_clone = dispatched.clone();
2404
2405 let service_override = CallbackBasedGrpcService {
2407 callback: Arc::new(|_req| {
2408 Box::pin(async {
2409 let resp = StartWorkflowExecutionResponse {
2410 run_id: "test-run-id".to_string(),
2411 eager_workflow_task: Some(PollWorkflowTaskQueueResponse {
2412 task_token: vec![1, 2, 3],
2413 ..Default::default()
2414 }),
2415 ..Default::default()
2416 };
2417 let proto = resp.encode_to_vec();
2418 Ok(GrpcSuccessResponse {
2419 headers: Default::default(),
2420 proto,
2421 })
2422 })
2423 }),
2424 };
2425
2426 let opts = ConnectionOptions::new(url::Url::parse("http://localhost:7233").unwrap())
2427 .skip_get_system_info(true)
2428 .service_override(service_override)
2429 .dns_load_balancing(None)
2430 .build();
2431 let mut connection = crate::Connection::connect(opts).await.unwrap();
2432
2433 let mut mock_worker = MockClientWorker::new();
2435 mock_worker
2436 .expect_namespace()
2437 .return_const("default".to_string());
2438 mock_worker
2439 .expect_task_queue()
2440 .return_const("test-tq".to_string());
2441 mock_worker
2442 .expect_deployment_options()
2443 .return_const(None::<temporalio_common::worker::WorkerDeploymentOptions>);
2444 mock_worker.expect_heartbeat_enabled().return_const(false);
2445 let uuid = Uuid::new_v4();
2446 mock_worker.expect_worker_instance_key().return_const(uuid);
2447 mock_worker
2448 .expect_worker_task_types()
2449 .return_const(WorkerTaskTypes {
2450 enable_workflows: true,
2451 enable_local_activities: false,
2452 enable_remote_activities: false,
2453 enable_nexus: false,
2454 });
2455
2456 let mut mock_slot = MockSlot::new();
2457 mock_slot.expect_schedule_wft().returning(move |_| {
2458 dispatched_clone.store(true, Ordering::SeqCst);
2459 Ok(())
2460 });
2461 mock_worker
2462 .expect_try_reserve_wft_slot()
2463 .return_once(|| Some(Box::new(mock_slot)));
2464
2465 connection
2466 .workers()
2467 .register_worker(Arc::new(mock_worker), true)
2468 .unwrap();
2469
2470 let req = StartWorkflowExecutionRequest {
2472 namespace: "default".to_string(),
2473 workflow_id: "test-wf".to_string(),
2474 workflow_type: Some(
2475 temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2476 name: "test-workflow".to_string(),
2477 },
2478 ),
2479 task_queue: Some(TaskQueue {
2480 name: "test-tq".to_string(),
2481 kind: 0,
2482 normal_name: String::new(),
2483 }),
2484 request_eager_execution: true,
2485 ..Default::default()
2486 };
2487
2488 connection
2489 .start_workflow_execution(req.into_request())
2490 .await
2491 .unwrap();
2492
2493 assert!(
2494 dispatched.load(Ordering::SeqCst),
2495 "Eager workflow task should have been dispatched to the worker"
2496 );
2497 }
2498}