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 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#[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
544macro_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)]
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 $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 AttachMetricLabels::namespace(ns_str)
701 }};
702}
703
704fn 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 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 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 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 assert!(validate_request_payload_limits(&new_req(), 1, 1).is_ok());
2052
2053 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 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 assert!(validate_request_payload_limits(&new_req(), 0, 0).is_ok());
2075 }
2076
2077 #[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 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 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 client
2124 .list_namespaces(list_ns_req.into_request())
2125 .await
2126 .unwrap();
2127 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 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 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 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 #[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 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 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 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}