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 labels = namespaced_request!(r);
1346 r.extensions_mut().insert(labels);
1347 }
1348 );
1349 (
1350 get_current_deployment,
1351 GetCurrentDeploymentRequest,
1352 GetCurrentDeploymentResponse,
1353 |r| {
1354 let labels = namespaced_request!(r);
1355 r.extensions_mut().insert(labels);
1356 }
1357 );
1358 (
1359 get_deployment_reachability,
1360 GetDeploymentReachabilityRequest,
1361 GetDeploymentReachabilityResponse,
1362 |r| {
1363 let labels = namespaced_request!(r);
1364 r.extensions_mut().insert(labels);
1365 }
1366 );
1367 (
1368 get_worker_versioning_rules,
1369 GetWorkerVersioningRulesRequest,
1370 GetWorkerVersioningRulesResponse,
1371 |r| {
1372 let mut labels = namespaced_request!(r);
1373 labels.task_q_str(&r.get_ref().task_queue);
1374 r.extensions_mut().insert(labels);
1375 }
1376 );
1377 (
1378 update_worker_versioning_rules,
1379 UpdateWorkerVersioningRulesRequest,
1380 UpdateWorkerVersioningRulesResponse,
1381 |r| {
1382 let mut labels = namespaced_request!(r);
1383 labels.task_q_str(&r.get_ref().task_queue);
1384 r.extensions_mut().insert(labels);
1385 }
1386 );
1387 (
1388 poll_nexus_task_queue,
1389 PollNexusTaskQueueRequest,
1390 PollNexusTaskQueueResponse,
1391 |r| {
1392 let mut labels = namespaced_request!(r);
1393 labels.task_q(r.get_ref().task_queue.clone());
1394 r.extensions_mut().insert(labels);
1395 }
1396 );
1397 (
1398 respond_nexus_task_completed,
1399 RespondNexusTaskCompletedRequest,
1400 RespondNexusTaskCompletedResponse,
1401 |r| {
1402 let labels = namespaced_request!(r);
1403 r.extensions_mut().insert(labels);
1404 }
1405 );
1406 (
1407 respond_nexus_task_failed,
1408 RespondNexusTaskFailedRequest,
1409 RespondNexusTaskFailedResponse,
1410 |r| {
1411 let labels = namespaced_request!(r);
1412 r.extensions_mut().insert(labels);
1413 }
1414 );
1415 (
1416 set_current_deployment,
1417 SetCurrentDeploymentRequest,
1418 SetCurrentDeploymentResponse,
1419 |r| {
1420 let labels = namespaced_request!(r);
1421 r.extensions_mut().insert(labels);
1422 }
1423 );
1424 (
1425 shutdown_worker,
1426 ShutdownWorkerRequest,
1427 ShutdownWorkerResponse,
1428 |r| {
1429 let labels = namespaced_request!(r);
1430 r.extensions_mut().insert(labels);
1431 }
1432 );
1433 (
1434 update_activity_options,
1435 UpdateActivityOptionsRequest,
1436 UpdateActivityOptionsResponse,
1437 |r| {
1438 let labels = namespaced_request!(r);
1439 r.extensions_mut().insert(labels);
1440 }
1441 );
1442 (
1443 pause_activity,
1444 PauseActivityRequest,
1445 PauseActivityResponse,
1446 |r| {
1447 let labels = namespaced_request!(r);
1448 r.extensions_mut().insert(labels);
1449 }
1450 );
1451 (
1452 unpause_activity,
1453 UnpauseActivityRequest,
1454 UnpauseActivityResponse,
1455 |r| {
1456 let labels = namespaced_request!(r);
1457 r.extensions_mut().insert(labels);
1458 }
1459 );
1460 (
1461 update_workflow_execution_options,
1462 UpdateWorkflowExecutionOptionsRequest,
1463 UpdateWorkflowExecutionOptionsResponse,
1464 |r| {
1465 let labels = namespaced_request!(r);
1466 r.extensions_mut().insert(labels);
1467 }
1468 );
1469 (
1470 reset_activity,
1471 ResetActivityRequest,
1472 ResetActivityResponse,
1473 |r| {
1474 let labels = namespaced_request!(r);
1475 r.extensions_mut().insert(labels);
1476 }
1477 );
1478 (
1479 delete_worker_deployment,
1480 DeleteWorkerDeploymentRequest,
1481 DeleteWorkerDeploymentResponse,
1482 |r| {
1483 let labels = namespaced_request!(r);
1484 r.extensions_mut().insert(labels);
1485 }
1486 );
1487 (
1488 delete_worker_deployment_version,
1489 DeleteWorkerDeploymentVersionRequest,
1490 DeleteWorkerDeploymentVersionResponse,
1491 |r| {
1492 let labels = namespaced_request!(r);
1493 r.extensions_mut().insert(labels);
1494 }
1495 );
1496 (
1497 describe_worker_deployment,
1498 DescribeWorkerDeploymentRequest,
1499 DescribeWorkerDeploymentResponse,
1500 |r| {
1501 let labels = namespaced_request!(r);
1502 r.extensions_mut().insert(labels);
1503 }
1504 );
1505 (
1506 describe_worker_deployment_version,
1507 DescribeWorkerDeploymentVersionRequest,
1508 DescribeWorkerDeploymentVersionResponse,
1509 |r| {
1510 let labels = namespaced_request!(r);
1511 r.extensions_mut().insert(labels);
1512 }
1513 );
1514 (
1515 list_worker_deployments,
1516 ListWorkerDeploymentsRequest,
1517 ListWorkerDeploymentsResponse,
1518 |r| {
1519 let labels = namespaced_request!(r);
1520 r.extensions_mut().insert(labels);
1521 }
1522 );
1523 (
1524 set_worker_deployment_current_version,
1525 SetWorkerDeploymentCurrentVersionRequest,
1526 SetWorkerDeploymentCurrentVersionResponse,
1527 |r| {
1528 let labels = namespaced_request!(r);
1529 r.extensions_mut().insert(labels);
1530 }
1531 );
1532 (
1533 set_worker_deployment_ramping_version,
1534 SetWorkerDeploymentRampingVersionRequest,
1535 SetWorkerDeploymentRampingVersionResponse,
1536 |r| {
1537 let labels = namespaced_request!(r);
1538 r.extensions_mut().insert(labels);
1539 }
1540 );
1541 (
1542 update_worker_deployment_version_metadata,
1543 UpdateWorkerDeploymentVersionMetadataRequest,
1544 UpdateWorkerDeploymentVersionMetadataResponse,
1545 |r| {
1546 let labels = namespaced_request!(r);
1547 r.extensions_mut().insert(labels);
1548 }
1549 );
1550 (
1551 list_workers,
1552 ListWorkersRequest,
1553 ListWorkersResponse,
1554 |r| {
1555 let labels = namespaced_request!(r);
1556 r.extensions_mut().insert(labels);
1557 }
1558 );
1559 (
1560 count_workers,
1561 CountWorkersRequest,
1562 CountWorkersResponse,
1563 |r| {
1564 let labels = namespaced_request!(r);
1565 r.extensions_mut().insert(labels);
1566 }
1567 );
1568 (
1569 record_worker_heartbeat,
1570 RecordWorkerHeartbeatRequest,
1571 RecordWorkerHeartbeatResponse,
1572 |r| {
1573 let labels = namespaced_request!(r);
1574 r.extensions_mut().insert(labels);
1575 }
1576 );
1577 (
1578 update_task_queue_config,
1579 UpdateTaskQueueConfigRequest,
1580 UpdateTaskQueueConfigResponse,
1581 |r| {
1582 let mut labels = namespaced_request!(r);
1583 labels.task_q_str(r.get_ref().task_queue.clone());
1584 r.extensions_mut().insert(labels);
1585 }
1586 );
1587 (
1588 fetch_worker_config,
1589 FetchWorkerConfigRequest,
1590 FetchWorkerConfigResponse,
1591 |r| {
1592 let labels = namespaced_request!(r);
1593 r.extensions_mut().insert(labels);
1594 }
1595 );
1596 (
1597 update_worker_config,
1598 UpdateWorkerConfigRequest,
1599 UpdateWorkerConfigResponse,
1600 |r| {
1601 let labels = namespaced_request!(r);
1602 r.extensions_mut().insert(labels);
1603 }
1604 );
1605 (
1606 describe_worker,
1607 DescribeWorkerRequest,
1608 DescribeWorkerResponse,
1609 |r| {
1610 let labels = namespaced_request!(r);
1611 r.extensions_mut().insert(labels);
1612 }
1613 );
1614 (
1615 set_worker_deployment_manager,
1616 SetWorkerDeploymentManagerRequest,
1617 SetWorkerDeploymentManagerResponse,
1618 |r| {
1619 let labels = namespaced_request!(r);
1620 r.extensions_mut().insert(labels);
1621 }
1622 );
1623 (
1624 pause_workflow_execution,
1625 PauseWorkflowExecutionRequest,
1626 PauseWorkflowExecutionResponse,
1627 |r| {
1628 let labels = namespaced_request!(r);
1629 r.extensions_mut().insert(labels);
1630 }
1631 );
1632 (
1633 unpause_workflow_execution,
1634 UnpauseWorkflowExecutionRequest,
1635 UnpauseWorkflowExecutionResponse,
1636 |r| {
1637 let labels = namespaced_request!(r);
1638 r.extensions_mut().insert(labels);
1639 }
1640 );
1641 (
1642 start_activity_execution,
1643 StartActivityExecutionRequest,
1644 StartActivityExecutionResponse,
1645 |r| {
1646 let labels = namespaced_request!(r);
1647 r.extensions_mut().insert(labels);
1648 }
1649 );
1650 (
1651 describe_activity_execution,
1652 DescribeActivityExecutionRequest,
1653 DescribeActivityExecutionResponse,
1654 |r| {
1655 let labels = namespaced_request!(r);
1656 r.extensions_mut().insert(labels);
1657 }
1658 );
1659 (
1660 poll_activity_execution,
1661 PollActivityExecutionRequest,
1662 PollActivityExecutionResponse,
1663 |r| {
1664 let labels = namespaced_request!(r);
1665 r.extensions_mut().insert(labels);
1666 }
1667 );
1668 (
1669 list_activity_executions,
1670 ListActivityExecutionsRequest,
1671 ListActivityExecutionsResponse,
1672 |r| {
1673 let labels = namespaced_request!(r);
1674 r.extensions_mut().insert(labels);
1675 }
1676 );
1677 (
1678 count_activity_executions,
1679 CountActivityExecutionsRequest,
1680 CountActivityExecutionsResponse,
1681 |r| {
1682 let labels = namespaced_request!(r);
1683 r.extensions_mut().insert(labels);
1684 }
1685 );
1686 (
1687 request_cancel_activity_execution,
1688 RequestCancelActivityExecutionRequest,
1689 RequestCancelActivityExecutionResponse,
1690 |r| {
1691 let labels = namespaced_request!(r);
1692 r.extensions_mut().insert(labels);
1693 }
1694 );
1695 (
1696 terminate_activity_execution,
1697 TerminateActivityExecutionRequest,
1698 TerminateActivityExecutionResponse,
1699 |r| {
1700 let labels = namespaced_request!(r);
1701 r.extensions_mut().insert(labels);
1702 }
1703 );
1704 (
1705 delete_activity_execution,
1706 DeleteActivityExecutionRequest,
1707 DeleteActivityExecutionResponse,
1708 |r| {
1709 let labels = namespaced_request!(r);
1710 r.extensions_mut().insert(labels);
1711 }
1712 );
1713 (
1714 pause_activity_execution,
1715 PauseActivityExecutionRequest,
1716 PauseActivityExecutionResponse,
1717 |r| {
1718 let labels = namespaced_request!(r);
1719 r.extensions_mut().insert(labels);
1720 }
1721 );
1722 (
1723 unpause_activity_execution,
1724 UnpauseActivityExecutionRequest,
1725 UnpauseActivityExecutionResponse,
1726 |r| {
1727 let labels = namespaced_request!(r);
1728 r.extensions_mut().insert(labels);
1729 }
1730 );
1731 (
1732 reset_activity_execution,
1733 ResetActivityExecutionRequest,
1734 ResetActivityExecutionResponse,
1735 |r| {
1736 let labels = namespaced_request!(r);
1737 r.extensions_mut().insert(labels);
1738 }
1739 );
1740 (
1741 update_activity_execution_options,
1742 UpdateActivityExecutionOptionsRequest,
1743 UpdateActivityExecutionOptionsResponse,
1744 |r| {
1745 let labels = namespaced_request!(r);
1746 r.extensions_mut().insert(labels);
1747 }
1748 );
1749 (
1750 count_nexus_operation_executions,
1751 CountNexusOperationExecutionsRequest,
1752 CountNexusOperationExecutionsResponse,
1753 |r| {
1754 let labels = namespaced_request!(r);
1755 r.extensions_mut().insert(labels);
1756 }
1757 );
1758 (
1759 create_worker_deployment,
1760 CreateWorkerDeploymentRequest,
1761 CreateWorkerDeploymentResponse,
1762 |r| {
1763 let labels = namespaced_request!(r);
1764 r.extensions_mut().insert(labels);
1765 }
1766 );
1767 (
1768 create_worker_deployment_version,
1769 CreateWorkerDeploymentVersionRequest,
1770 CreateWorkerDeploymentVersionResponse,
1771 |r| {
1772 let labels = namespaced_request!(r);
1773 r.extensions_mut().insert(labels);
1774 }
1775 );
1776 (
1777 delete_nexus_operation_execution,
1778 DeleteNexusOperationExecutionRequest,
1779 DeleteNexusOperationExecutionResponse,
1780 |r| {
1781 let labels = namespaced_request!(r);
1782 r.extensions_mut().insert(labels);
1783 }
1784 );
1785 (
1786 describe_nexus_operation_execution,
1787 DescribeNexusOperationExecutionRequest,
1788 DescribeNexusOperationExecutionResponse,
1789 |r| {
1790 let labels = namespaced_request!(r);
1791 r.extensions_mut().insert(labels);
1792 }
1793 );
1794 (
1795 list_nexus_operation_executions,
1796 ListNexusOperationExecutionsRequest,
1797 ListNexusOperationExecutionsResponse,
1798 |r| {
1799 let labels = namespaced_request!(r);
1800 r.extensions_mut().insert(labels);
1801 }
1802 );
1803 (
1804 poll_nexus_operation_execution,
1805 PollNexusOperationExecutionRequest,
1806 PollNexusOperationExecutionResponse,
1807 |r| {
1808 let labels = namespaced_request!(r);
1809 r.extensions_mut().insert(labels);
1810 }
1811 );
1812 (
1813 poll_workflow_execution_time_skipping,
1814 PollWorkflowExecutionTimeSkippingRequest,
1815 PollWorkflowExecutionTimeSkippingResponse,
1816 |r| {
1817 let labels = namespaced_request!(r);
1818 r.extensions_mut().insert(labels);
1819 r.extensions_mut().insert(IsUserLongPoll);
1820 }
1821 );
1822 (
1823 request_cancel_nexus_operation_execution,
1824 RequestCancelNexusOperationExecutionRequest,
1825 RequestCancelNexusOperationExecutionResponse,
1826 |r| {
1827 let labels = namespaced_request!(r);
1828 r.extensions_mut().insert(labels);
1829 }
1830 );
1831 (
1832 start_nexus_operation_execution,
1833 StartNexusOperationExecutionRequest,
1834 StartNexusOperationExecutionResponse,
1835 |r| {
1836 let labels = namespaced_request!(r);
1837 r.extensions_mut().insert(labels);
1838 }
1839 );
1840 (
1841 terminate_nexus_operation_execution,
1842 TerminateNexusOperationExecutionRequest,
1843 TerminateNexusOperationExecutionResponse,
1844 |r| {
1845 let labels = namespaced_request!(r);
1846 r.extensions_mut().insert(labels);
1847 }
1848 );
1849 (
1850 update_worker_deployment_version_compute_config,
1851 UpdateWorkerDeploymentVersionComputeConfigRequest,
1852 UpdateWorkerDeploymentVersionComputeConfigResponse,
1853 |r| {
1854 let labels = namespaced_request!(r);
1855 r.extensions_mut().insert(labels);
1856 }
1857 );
1858 (
1859 validate_worker_deployment_version_compute_config,
1860 ValidateWorkerDeploymentVersionComputeConfigRequest,
1861 ValidateWorkerDeploymentVersionComputeConfigResponse,
1862 |r| {
1863 let labels = namespaced_request!(r);
1864 r.extensions_mut().insert(labels);
1865 }
1866 );
1867}
1868
1869proxier! {
1870 OperatorService; ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS; OperatorServiceClient; operator_client; defaults;
1871 (add_search_attributes, AddSearchAttributesRequest, AddSearchAttributesResponse);
1872 (remove_search_attributes, RemoveSearchAttributesRequest, RemoveSearchAttributesResponse);
1873 (list_search_attributes, ListSearchAttributesRequest, ListSearchAttributesResponse);
1874 (delete_namespace, DeleteNamespaceRequest, DeleteNamespaceResponse,
1875 |r| {
1876 let labels = namespaced_request!(r);
1877 r.extensions_mut().insert(labels);
1878 }
1879 );
1880 (add_or_update_remote_cluster, AddOrUpdateRemoteClusterRequest, AddOrUpdateRemoteClusterResponse);
1881 (remove_remote_cluster, RemoveRemoteClusterRequest, RemoveRemoteClusterResponse);
1882 (list_clusters, ListClustersRequest, ListClustersResponse);
1883 (get_nexus_endpoint, GetNexusEndpointRequest, GetNexusEndpointResponse);
1884 (create_nexus_endpoint, CreateNexusEndpointRequest, CreateNexusEndpointResponse);
1885 (update_nexus_endpoint, UpdateNexusEndpointRequest, UpdateNexusEndpointResponse);
1886 (delete_nexus_endpoint, DeleteNexusEndpointRequest, DeleteNexusEndpointResponse);
1887 (list_nexus_endpoints, ListNexusEndpointsRequest, ListNexusEndpointsResponse);
1888}
1889
1890proxier! {
1891 CloudService; ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS; CloudServiceClient; cloud_client; defaults;
1892 (get_users, cloudreq::GetUsersRequest, cloudreq::GetUsersResponse);
1893 (get_user, cloudreq::GetUserRequest, cloudreq::GetUserResponse);
1894 (create_user, cloudreq::CreateUserRequest, cloudreq::CreateUserResponse);
1895 (update_user, cloudreq::UpdateUserRequest, cloudreq::UpdateUserResponse);
1896 (delete_user, cloudreq::DeleteUserRequest, cloudreq::DeleteUserResponse);
1897 (set_user_namespace_access, cloudreq::SetUserNamespaceAccessRequest, cloudreq::SetUserNamespaceAccessResponse);
1898 (get_async_operation, cloudreq::GetAsyncOperationRequest, cloudreq::GetAsyncOperationResponse);
1899 (create_namespace, cloudreq::CreateNamespaceRequest, cloudreq::CreateNamespaceResponse);
1900 (get_namespaces, cloudreq::GetNamespacesRequest, cloudreq::GetNamespacesResponse);
1901 (get_namespace, cloudreq::GetNamespaceRequest, cloudreq::GetNamespaceResponse,
1902 |r| {
1903 let labels = namespaced_request!(r);
1904 r.extensions_mut().insert(labels);
1905 }
1906 );
1907 (update_namespace, cloudreq::UpdateNamespaceRequest, cloudreq::UpdateNamespaceResponse,
1908 |r| {
1909 let labels = namespaced_request!(r);
1910 r.extensions_mut().insert(labels);
1911 }
1912 );
1913 (rename_custom_search_attribute, cloudreq::RenameCustomSearchAttributeRequest, cloudreq::RenameCustomSearchAttributeResponse);
1914 (delete_namespace, cloudreq::DeleteNamespaceRequest, cloudreq::DeleteNamespaceResponse,
1915 |r| {
1916 let labels = namespaced_request!(r);
1917 r.extensions_mut().insert(labels);
1918 }
1919 );
1920 (failover_namespace_region, cloudreq::FailoverNamespaceRegionRequest, cloudreq::FailoverNamespaceRegionResponse);
1921 (add_namespace_region, cloudreq::AddNamespaceRegionRequest, cloudreq::AddNamespaceRegionResponse);
1922 (delete_namespace_region, cloudreq::DeleteNamespaceRegionRequest, cloudreq::DeleteNamespaceRegionResponse);
1923 (get_regions, cloudreq::GetRegionsRequest, cloudreq::GetRegionsResponse);
1924 (get_region, cloudreq::GetRegionRequest, cloudreq::GetRegionResponse);
1925 (get_api_keys, cloudreq::GetApiKeysRequest, cloudreq::GetApiKeysResponse);
1926 (get_api_key, cloudreq::GetApiKeyRequest, cloudreq::GetApiKeyResponse);
1927 (create_api_key, cloudreq::CreateApiKeyRequest, cloudreq::CreateApiKeyResponse);
1928 (update_api_key, cloudreq::UpdateApiKeyRequest, cloudreq::UpdateApiKeyResponse);
1929 (delete_api_key, cloudreq::DeleteApiKeyRequest, cloudreq::DeleteApiKeyResponse);
1930 (get_nexus_endpoints, cloudreq::GetNexusEndpointsRequest, cloudreq::GetNexusEndpointsResponse);
1931 (get_nexus_endpoint, cloudreq::GetNexusEndpointRequest, cloudreq::GetNexusEndpointResponse);
1932 (create_nexus_endpoint, cloudreq::CreateNexusEndpointRequest, cloudreq::CreateNexusEndpointResponse);
1933 (update_nexus_endpoint, cloudreq::UpdateNexusEndpointRequest, cloudreq::UpdateNexusEndpointResponse);
1934 (delete_nexus_endpoint, cloudreq::DeleteNexusEndpointRequest, cloudreq::DeleteNexusEndpointResponse);
1935 (get_user_groups, cloudreq::GetUserGroupsRequest, cloudreq::GetUserGroupsResponse);
1936 (get_user_group, cloudreq::GetUserGroupRequest, cloudreq::GetUserGroupResponse);
1937 (create_user_group, cloudreq::CreateUserGroupRequest, cloudreq::CreateUserGroupResponse);
1938 (update_user_group, cloudreq::UpdateUserGroupRequest, cloudreq::UpdateUserGroupResponse);
1939 (delete_user_group, cloudreq::DeleteUserGroupRequest, cloudreq::DeleteUserGroupResponse);
1940 (add_user_group_member, cloudreq::AddUserGroupMemberRequest, cloudreq::AddUserGroupMemberResponse);
1941 (remove_user_group_member, cloudreq::RemoveUserGroupMemberRequest, cloudreq::RemoveUserGroupMemberResponse);
1942 (get_user_group_members, cloudreq::GetUserGroupMembersRequest, cloudreq::GetUserGroupMembersResponse);
1943 (set_user_group_namespace_access, cloudreq::SetUserGroupNamespaceAccessRequest, cloudreq::SetUserGroupNamespaceAccessResponse);
1944 (create_service_account, cloudreq::CreateServiceAccountRequest, cloudreq::CreateServiceAccountResponse);
1945 (get_service_account, cloudreq::GetServiceAccountRequest, cloudreq::GetServiceAccountResponse);
1946 (get_service_accounts, cloudreq::GetServiceAccountsRequest, cloudreq::GetServiceAccountsResponse);
1947 (update_service_account, cloudreq::UpdateServiceAccountRequest, cloudreq::UpdateServiceAccountResponse);
1948 (delete_service_account, cloudreq::DeleteServiceAccountRequest, cloudreq::DeleteServiceAccountResponse);
1949 (get_usage, cloudreq::GetUsageRequest, cloudreq::GetUsageResponse);
1950 (get_account, cloudreq::GetAccountRequest, cloudreq::GetAccountResponse);
1951 (update_account, cloudreq::UpdateAccountRequest, cloudreq::UpdateAccountResponse);
1952 (create_namespace_export_sink, cloudreq::CreateNamespaceExportSinkRequest, cloudreq::CreateNamespaceExportSinkResponse);
1953 (get_namespace_export_sink, cloudreq::GetNamespaceExportSinkRequest, cloudreq::GetNamespaceExportSinkResponse);
1954 (get_namespace_export_sinks, cloudreq::GetNamespaceExportSinksRequest, cloudreq::GetNamespaceExportSinksResponse);
1955 (update_namespace_export_sink, cloudreq::UpdateNamespaceExportSinkRequest, cloudreq::UpdateNamespaceExportSinkResponse);
1956 (delete_namespace_export_sink, cloudreq::DeleteNamespaceExportSinkRequest, cloudreq::DeleteNamespaceExportSinkResponse);
1957 (validate_namespace_export_sink, cloudreq::ValidateNamespaceExportSinkRequest, cloudreq::ValidateNamespaceExportSinkResponse);
1958 (update_namespace_tags, cloudreq::UpdateNamespaceTagsRequest, cloudreq::UpdateNamespaceTagsResponse);
1959 (create_connectivity_rule, cloudreq::CreateConnectivityRuleRequest, cloudreq::CreateConnectivityRuleResponse);
1960 (get_connectivity_rule, cloudreq::GetConnectivityRuleRequest, cloudreq::GetConnectivityRuleResponse);
1961 (get_connectivity_rules, cloudreq::GetConnectivityRulesRequest, cloudreq::GetConnectivityRulesResponse);
1962 (delete_connectivity_rule, cloudreq::DeleteConnectivityRuleRequest, cloudreq::DeleteConnectivityRuleResponse);
1963 (set_service_account_namespace_access, cloudreq::SetServiceAccountNamespaceAccessRequest, cloudreq::SetServiceAccountNamespaceAccessResponse);
1964 (validate_account_audit_log_sink, cloudreq::ValidateAccountAuditLogSinkRequest, cloudreq::ValidateAccountAuditLogSinkResponse);
1965 (get_current_identity, cloudreq::GetCurrentIdentityRequest, cloudreq::GetCurrentIdentityResponse);
1966 (get_audit_logs, cloudreq::GetAuditLogsRequest, cloudreq::GetAuditLogsResponse);
1967 (create_account_audit_log_sink, cloudreq::CreateAccountAuditLogSinkRequest, cloudreq::CreateAccountAuditLogSinkResponse);
1968 (get_account_audit_log_sink, cloudreq::GetAccountAuditLogSinkRequest, cloudreq::GetAccountAuditLogSinkResponse);
1969 (get_account_audit_log_sinks, cloudreq::GetAccountAuditLogSinksRequest, cloudreq::GetAccountAuditLogSinksResponse);
1970 (update_account_audit_log_sink, cloudreq::UpdateAccountAuditLogSinkRequest, cloudreq::UpdateAccountAuditLogSinkResponse);
1971 (delete_account_audit_log_sink, cloudreq::DeleteAccountAuditLogSinkRequest, cloudreq::DeleteAccountAuditLogSinkResponse);
1972 (get_namespace_capacity_info, cloudreq::GetNamespaceCapacityInfoRequest, cloudreq::GetNamespaceCapacityInfoResponse);
1973 (create_billing_report, cloudreq::CreateBillingReportRequest, cloudreq::CreateBillingReportResponse);
1974 (get_billing_report, cloudreq::GetBillingReportRequest, cloudreq::GetBillingReportResponse);
1975 (get_custom_roles, cloudreq::GetCustomRolesRequest, cloudreq::GetCustomRolesResponse);
1976 (get_custom_role, cloudreq::GetCustomRoleRequest, cloudreq::GetCustomRoleResponse);
1977 (create_custom_role, cloudreq::CreateCustomRoleRequest, cloudreq::CreateCustomRoleResponse);
1978 (update_custom_role, cloudreq::UpdateCustomRoleRequest, cloudreq::UpdateCustomRoleResponse);
1979 (delete_custom_role, cloudreq::DeleteCustomRoleRequest, cloudreq::DeleteCustomRoleResponse);
1980 (get_user_namespace_assignments, cloudreq::GetUserNamespaceAssignmentsRequest, cloudreq::GetUserNamespaceAssignmentsResponse);
1981 (get_service_account_namespace_assignments, cloudreq::GetServiceAccountNamespaceAssignmentsRequest, cloudreq::GetServiceAccountNamespaceAssignmentsResponse);
1982 (get_user_group_namespace_assignments, cloudreq::GetUserGroupNamespaceAssignmentsRequest, cloudreq::GetUserGroupNamespaceAssignmentsResponse);
1983}
1984
1985proxier! {
1986 TestService; ALL_IMPLEMENTED_TEST_SERVICE_RPCS; TestServiceClient; test_client; defaults;
1987 (lock_time_skipping, LockTimeSkippingRequest, LockTimeSkippingResponse);
1988 (unlock_time_skipping, UnlockTimeSkippingRequest, UnlockTimeSkippingResponse);
1989 (sleep, SleepRequest, SleepResponse);
1990 (sleep_until, SleepUntilRequest, SleepResponse);
1991 (unlock_time_skipping_with_sleep, SleepRequest, SleepResponse);
1992 (get_current_time, (), GetCurrentTimeResponse);
1993}
1994
1995proxier! {
1996 HealthService; ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS; HealthClient; health_client;
1997 (check, HealthCheckRequest, HealthCheckResponse);
1998 (watch, HealthCheckRequest, tonic::codec::Streaming<HealthCheckResponse>);
1999}
2000
2001#[cfg(test)]
2002mod tests {
2003 use super::*;
2004 use crate::{ClientOptions, ConnectionOptions};
2005 use std::collections::HashSet;
2006 use temporalio_common::{
2007 protos::temporal::api::{
2008 operatorservice::v1::DeleteNamespaceRequest, workflowservice::v1::ListNamespacesRequest,
2009 },
2010 worker::WorkerTaskTypes,
2011 };
2012 use tonic::IntoRequest;
2013 use url::Url;
2014 use uuid::Uuid;
2015
2016 #[test]
2017 fn payload_limits_warn_only_vs_error() {
2018 use temporalio_common::protos::temporal::api::{
2019 common::v1::{Payload, Payloads},
2020 workflowservice::v1::StartWorkflowExecutionRequest,
2021 };
2022 let big = Payloads {
2023 payloads: vec![Payload {
2024 data: vec![0u8; 1000],
2025 ..Default::default()
2026 }],
2027 };
2028 let new_req = || {
2029 StartWorkflowExecutionRequest {
2030 input: Some(big.clone()),
2031 ..Default::default()
2032 }
2033 .into_request()
2034 };
2035
2036 assert!(validate_request_payload_limits(&new_req(), 1, 1).is_ok());
2038
2039 let mut req = new_req();
2041 req.extensions_mut()
2042 .insert(PayloadErrorLimits { blob: 10, memo: 10 });
2043 let err = validate_request_payload_limits(&req, 1, 1).unwrap_err();
2044 let violation =
2045 crate::payload_limit_violation_from(&err).expect("violation carried on status");
2046 assert_eq!(violation.path, "input");
2047 assert_eq!(
2048 violation.class,
2049 temporalio_common::payload_limits::LimitClass::Blob
2050 );
2051 assert!(violation.size > violation.limit);
2052
2053 let mut req = new_req();
2055 req.extensions_mut()
2056 .insert(PayloadErrorLimits { blob: 0, memo: 0 });
2057 assert!(validate_request_payload_limits(&req, 1, 1).is_ok());
2058
2059 assert!(validate_request_payload_limits(&new_req(), 0, 0).is_ok());
2061 }
2062
2063 #[allow(dead_code)]
2065 async fn raw_client_retry_compiles() {
2066 let opts = ConnectionOptions::new(Url::parse("http://localhost:7233").unwrap())
2067 .client_name("test")
2068 .client_version("0.0.0")
2069 .build();
2070 let connection = Connection::connect(opts).await.unwrap();
2071 let mut client = Client::new(connection, ClientOptions::new("default").build()).unwrap();
2072
2073 let list_ns_req = ListNamespacesRequest::default();
2074 let wf_client = client.workflow_client();
2075 let fact = move |req| {
2076 let mut c = wf_client.clone();
2077 async move { c.list_namespaces(req).await }.boxed()
2078 };
2079 client
2080 .call("whatever", fact, Request::new(list_ns_req.clone()))
2081 .await
2082 .unwrap();
2083
2084 let op_del_ns_req = DeleteNamespaceRequest::default();
2086 let op_client = client.operator_client();
2087 let fact = move |req| {
2088 let mut c = op_client.clone();
2089 async move { c.delete_namespace(req).await }.boxed()
2090 };
2091 client
2092 .call("whatever", fact, Request::new(op_del_ns_req.clone()))
2093 .await
2094 .unwrap();
2095
2096 let cloud_del_ns_req = cloudreq::DeleteNamespaceRequest::default();
2098 let cloud_client = client.cloud_client();
2099 let fact = move |req| {
2100 let mut c = cloud_client.clone();
2101 async move { c.delete_namespace(req).await }.boxed()
2102 };
2103 client
2104 .call("whatever", fact, Request::new(cloud_del_ns_req.clone()))
2105 .await
2106 .unwrap();
2107
2108 client
2110 .list_namespaces(list_ns_req.into_request())
2111 .await
2112 .unwrap();
2113 OperatorService::delete_namespace(&mut client, op_del_ns_req.into_request())
2115 .await
2116 .unwrap();
2117 CloudService::delete_namespace(&mut client, cloud_del_ns_req.into_request())
2118 .await
2119 .unwrap();
2120 client.get_current_time(().into_request()).await.unwrap();
2121 client
2122 .check(HealthCheckRequest::default().into_request())
2123 .await
2124 .unwrap();
2125 }
2126
2127 fn verify_methods(proto_def_str: &str, impl_list: &[&str]) {
2128 let methods: Vec<_> = proto_def_str
2129 .lines()
2130 .map(|l| l.trim())
2131 .filter(|l| l.starts_with("rpc"))
2132 .map(|l| {
2133 let stripped = l.strip_prefix("rpc ").unwrap();
2134 stripped[..stripped.find('(').unwrap()].trim()
2135 })
2136 .collect();
2137 let no_underscores: HashSet<_> = impl_list.iter().map(|x| x.replace('_', "")).collect();
2138 let mut not_implemented = vec![];
2139 for method in methods {
2140 if !no_underscores.contains(&method.to_lowercase()) {
2141 not_implemented.push(method);
2142 }
2143 }
2144 if !not_implemented.is_empty() {
2145 panic!(
2146 "The following RPC methods are not implemented by raw client: {not_implemented:?}"
2147 );
2148 }
2149 }
2150 #[test]
2151 fn verify_all_workflow_service_methods_implemented() {
2152 let proto_def = include_str!(
2154 "../../protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto"
2155 );
2156 verify_methods(proto_def, ALL_IMPLEMENTED_WORKFLOW_SERVICE_RPCS);
2157 }
2158
2159 #[test]
2160 fn verify_all_operator_service_methods_implemented() {
2161 let proto_def = include_str!(
2162 "../../protos/protos/api_upstream/temporal/api/operatorservice/v1/service.proto"
2163 );
2164 verify_methods(proto_def, ALL_IMPLEMENTED_OPERATOR_SERVICE_RPCS);
2165 }
2166
2167 #[test]
2168 fn verify_all_cloud_service_methods_implemented() {
2169 let proto_def = include_str!(
2170 "../../protos/protos/api_cloud_upstream/temporal/api/cloud/cloudservice/v1/service.proto"
2171 );
2172 verify_methods(proto_def, ALL_IMPLEMENTED_CLOUD_SERVICE_RPCS);
2173 }
2174
2175 #[test]
2176 fn verify_all_test_service_methods_implemented() {
2177 let proto_def = include_str!(
2178 "../../protos/protos/testsrv_upstream/temporal/api/testservice/v1/service.proto"
2179 );
2180 verify_methods(proto_def, ALL_IMPLEMENTED_TEST_SERVICE_RPCS);
2181 }
2182
2183 #[test]
2184 fn verify_all_health_service_methods_implemented() {
2185 let proto_def = include_str!("../../protos/protos/grpc/health/v1/health.proto");
2186 verify_methods(proto_def, ALL_IMPLEMENTED_HEALTH_SERVICE_RPCS);
2187 }
2188
2189 #[tokio::test]
2190 async fn can_mock_services() {
2191 #[derive(Clone)]
2192 struct MyFakeServices {}
2193 impl RawGrpcCaller for MyFakeServices {}
2194 impl WorkflowService for MyFakeServices {
2195 fn list_namespaces(
2196 &mut self,
2197 _request: Request<ListNamespacesRequest>,
2198 ) -> BoxFuture<'_, Result<Response<ListNamespacesResponse>, Status>> {
2199 async {
2200 Ok(Response::new(ListNamespacesResponse {
2201 namespaces: vec![DescribeNamespaceResponse {
2202 failover_version: 12345,
2203 ..Default::default()
2204 }],
2205 ..Default::default()
2206 }))
2207 }
2208 .boxed()
2209 }
2210 }
2211 impl OperatorService for MyFakeServices {}
2212 impl CloudService for MyFakeServices {}
2213 impl TestService for MyFakeServices {}
2214 impl HealthService for MyFakeServices {
2216 fn check(
2217 &mut self,
2218 _request: tonic::Request<HealthCheckRequest>,
2219 ) -> BoxFuture<'_, Result<tonic::Response<HealthCheckResponse>, tonic::Status>>
2220 {
2221 todo!()
2222 }
2223 fn watch(
2224 &mut self,
2225 _request: tonic::Request<HealthCheckRequest>,
2226 ) -> BoxFuture<
2227 '_,
2228 Result<
2229 tonic::Response<tonic::codec::Streaming<HealthCheckResponse>>,
2230 tonic::Status,
2231 >,
2232 > {
2233 todo!()
2234 }
2235 }
2236 let mut mocked_client = TemporalServiceClient::from_services(
2237 Box::new(MyFakeServices {}),
2238 Box::new(MyFakeServices {}),
2239 Box::new(MyFakeServices {}),
2240 Box::new(MyFakeServices {}),
2241 Box::new(MyFakeServices {}),
2242 );
2243 let r = mocked_client
2244 .list_namespaces(ListNamespacesRequest::default().into_request())
2245 .await
2246 .unwrap();
2247 assert_eq!(r.into_inner().namespaces[0].failover_version, 12345);
2248 }
2249
2250 #[rstest::rstest]
2251 #[case::with_versioning(true)]
2252 #[case::without_versioning(false)]
2253 #[tokio::test]
2254 async fn eager_reservations_attach_deployment_options(#[case] use_worker_versioning: bool) {
2255 use crate::worker::{MockClientWorker, MockSlot};
2256 use temporalio_common::{
2257 protos::temporal::api::enums::v1::WorkerVersioningMode,
2258 worker::{WorkerDeploymentOptions, WorkerDeploymentVersion},
2259 };
2260
2261 let expected_mode = if use_worker_versioning {
2262 WorkerVersioningMode::Versioned
2263 } else {
2264 WorkerVersioningMode::Unversioned
2265 };
2266
2267 #[derive(Clone)]
2268 struct MyFakeServices {
2269 client_worker_set: Arc<ClientWorkerSet>,
2270 expected_mode: WorkerVersioningMode,
2271 }
2272 impl RawGrpcCaller for MyFakeServices {}
2273 impl RawClientProducer for MyFakeServices {
2274 fn get_workers_info(&self) -> Option<Arc<ClientWorkerSet>> {
2275 Some(self.client_worker_set.clone())
2276 }
2277 fn workflow_client(&mut self) -> Box<dyn WorkflowService> {
2278 Box::new(MyFakeWfClient {
2279 expected_mode: self.expected_mode,
2280 })
2281 }
2282 fn operator_client(&mut self) -> Box<dyn OperatorService> {
2283 unimplemented!()
2284 }
2285 fn cloud_client(&mut self) -> Box<dyn CloudService> {
2286 unimplemented!()
2287 }
2288 fn test_client(&mut self) -> Box<dyn TestService> {
2289 unimplemented!()
2290 }
2291 fn health_client(&mut self) -> Box<dyn HealthService> {
2292 unimplemented!()
2293 }
2294 }
2295
2296 let deployment_opts = WorkerDeploymentOptions::new(WorkerDeploymentVersion {
2297 deployment_name: "test-deployment".to_string(),
2298 build_id: "test-build-123".to_string(),
2299 })
2300 .use_worker_versioning(use_worker_versioning)
2301 .build();
2302
2303 let mut mock_provider = MockClientWorker::new();
2304 mock_provider
2305 .expect_namespace()
2306 .return_const("test-namespace".to_string());
2307 mock_provider
2308 .expect_task_queue()
2309 .return_const("test-task-queue".to_string());
2310 let mut mock_slot = MockSlot::new();
2311 mock_slot.expect_schedule_wft().returning(|_| Ok(()));
2312 mock_provider
2313 .expect_try_reserve_wft_slot()
2314 .return_once(|| Some(Box::new(mock_slot)));
2315 mock_provider
2316 .expect_deployment_options()
2317 .return_const(Some(deployment_opts.clone()));
2318 mock_provider.expect_heartbeat_enabled().return_const(false);
2319 let uuid = Uuid::new_v4();
2320 mock_provider
2321 .expect_worker_instance_key()
2322 .return_const(uuid);
2323 mock_provider
2324 .expect_worker_task_types()
2325 .return_const(WorkerTaskTypes {
2326 enable_workflows: true,
2327 enable_local_activities: true,
2328 enable_remote_activities: true,
2329 enable_nexus: true,
2330 });
2331
2332 let client_worker_set = Arc::new(ClientWorkerSet::new());
2333 client_worker_set
2334 .register_worker(Arc::new(mock_provider), true)
2335 .unwrap();
2336
2337 #[derive(Clone)]
2338 struct MyFakeWfClient {
2339 expected_mode: WorkerVersioningMode,
2340 }
2341 impl WorkflowService for MyFakeWfClient {
2342 fn start_workflow_execution(
2343 &mut self,
2344 request: tonic::Request<StartWorkflowExecutionRequest>,
2345 ) -> BoxFuture<'_, Result<tonic::Response<StartWorkflowExecutionResponse>, tonic::Status>>
2346 {
2347 let req = request.into_inner();
2348 let expected_mode = self.expected_mode;
2349
2350 assert!(
2351 req.eager_worker_deployment_options.is_some(),
2352 "eager_worker_deployment_options should be populated"
2353 );
2354
2355 let opts = req.eager_worker_deployment_options.as_ref().unwrap();
2356 assert_eq!(opts.deployment_name, "test-deployment");
2357 assert_eq!(opts.build_id, "test-build-123");
2358 assert_eq!(opts.worker_versioning_mode, expected_mode as i32);
2359
2360 async { Ok(Response::new(StartWorkflowExecutionResponse::default())) }.boxed()
2361 }
2362 }
2363
2364 let mut mfs = MyFakeServices {
2365 client_worker_set,
2366 expected_mode,
2367 };
2368
2369 let req = StartWorkflowExecutionRequest {
2371 namespace: "test-namespace".to_string(),
2372 workflow_id: "test-wf-id".to_string(),
2373 workflow_type: Some(
2374 temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2375 name: "test-workflow".to_string(),
2376 },
2377 ),
2378 task_queue: Some(TaskQueue {
2379 name: "test-task-queue".to_string(),
2380 kind: 0,
2381 normal_name: String::new(),
2382 }),
2383 request_eager_execution: true,
2384 ..Default::default()
2385 };
2386
2387 mfs.start_workflow_execution(req.into_request())
2388 .await
2389 .unwrap();
2390 }
2391
2392 #[tokio::test]
2395 async fn connection_eager_start_dispatches_wft() {
2396 use crate::{
2397 ConnectionOptions,
2398 callback_based::{CallbackBasedGrpcService, GrpcSuccessResponse},
2399 worker::{MockClientWorker, MockSlot},
2400 };
2401 use prost::Message;
2402 use std::sync::atomic::{AtomicBool, Ordering};
2403 use temporalio_common::protos::temporal::api::workflowservice::v1::PollWorkflowTaskQueueResponse;
2404
2405 let dispatched = Arc::new(AtomicBool::new(false));
2406 let dispatched_clone = dispatched.clone();
2407
2408 let service_override = CallbackBasedGrpcService {
2410 callback: Arc::new(|_req| {
2411 Box::pin(async {
2412 let resp = StartWorkflowExecutionResponse {
2413 run_id: "test-run-id".to_string(),
2414 eager_workflow_task: Some(PollWorkflowTaskQueueResponse {
2415 task_token: vec![1, 2, 3],
2416 ..Default::default()
2417 }),
2418 ..Default::default()
2419 };
2420 let proto = resp.encode_to_vec();
2421 Ok(GrpcSuccessResponse {
2422 headers: Default::default(),
2423 proto,
2424 })
2425 })
2426 }),
2427 };
2428
2429 let opts = ConnectionOptions::new(url::Url::parse("http://localhost:7233").unwrap())
2430 .skip_get_system_info(true)
2431 .service_override(service_override)
2432 .dns_load_balancing(None)
2433 .build();
2434 let mut connection = crate::Connection::connect(opts).await.unwrap();
2435
2436 let mut mock_worker = MockClientWorker::new();
2438 mock_worker
2439 .expect_namespace()
2440 .return_const("default".to_string());
2441 mock_worker
2442 .expect_task_queue()
2443 .return_const("test-tq".to_string());
2444 mock_worker
2445 .expect_deployment_options()
2446 .return_const(None::<temporalio_common::worker::WorkerDeploymentOptions>);
2447 mock_worker.expect_heartbeat_enabled().return_const(false);
2448 let uuid = Uuid::new_v4();
2449 mock_worker.expect_worker_instance_key().return_const(uuid);
2450 mock_worker
2451 .expect_worker_task_types()
2452 .return_const(WorkerTaskTypes {
2453 enable_workflows: true,
2454 enable_local_activities: false,
2455 enable_remote_activities: false,
2456 enable_nexus: false,
2457 });
2458
2459 let mut mock_slot = MockSlot::new();
2460 mock_slot.expect_schedule_wft().returning(move |_| {
2461 dispatched_clone.store(true, Ordering::SeqCst);
2462 Ok(())
2463 });
2464 mock_worker
2465 .expect_try_reserve_wft_slot()
2466 .return_once(|| Some(Box::new(mock_slot)));
2467
2468 connection
2469 .workers()
2470 .register_worker(Arc::new(mock_worker), true)
2471 .unwrap();
2472
2473 let req = StartWorkflowExecutionRequest {
2475 namespace: "default".to_string(),
2476 workflow_id: "test-wf".to_string(),
2477 workflow_type: Some(
2478 temporalio_common::protos::temporal::api::common::v1::WorkflowType {
2479 name: "test-workflow".to_string(),
2480 },
2481 ),
2482 task_queue: Some(TaskQueue {
2483 name: "test-tq".to_string(),
2484 kind: 0,
2485 normal_name: String::new(),
2486 }),
2487 request_eager_execution: true,
2488 ..Default::default()
2489 };
2490
2491 connection
2492 .start_workflow_execution(req.into_request())
2493 .await
2494 .unwrap();
2495
2496 assert!(
2497 dispatched.load(Ordering::SeqCst),
2498 "Eager workflow task should have been dispatched to the worker"
2499 );
2500 }
2501}