@@ -419,7 +419,8 @@ def __init__(self, *,
419419 resiliency_options : GrpcClientResiliencyOptions | None = None ,
420420 default_version : str | None = None ,
421421 payload_store : PayloadStore | None = None ,
422- data_converter : DataConverter | None = None ):
422+ data_converter : DataConverter | None = None ,
423+ emit_trace_spans : bool = True ):
423424
424425 self ._owns_channel = channel is None
425426 self ._data_converter = data_converter if data_converter is not None else JsonDataConverter ()
@@ -484,6 +485,12 @@ def __init__(self, *,
484485 self ._logger = shared .get_logger ("client" , log_handler , log_formatter )
485486 self .default_version = default_version
486487 self ._payload_store = payload_store
488+ self ._emit_trace_spans = emit_trace_spans
489+
490+ @property
491+ def emit_trace_spans (self ) -> bool :
492+ """Return whether this client emits Durable Task lifecycle spans."""
493+ return self ._emit_trace_spans
487494
488495 def _compose_interceptors (
489496 self ,
@@ -619,30 +626,30 @@ def schedule_new_orchestration(self, orchestrator: task.Orchestrator[TInput, TOu
619626 resolved_instance_id = instance_id if instance_id else uuid .uuid4 ().hex
620627 resolved_version = version if version else self .default_version
621628
622- with tracing .start_create_orchestration_span (
623- name , resolved_instance_id , version = resolved_version ,
624- ):
625- req = build_schedule_new_orchestration_req (
626- orchestrator , input = input , instance_id = instance_id , start_at = start_at ,
627- reuse_id_policy = reuse_id_policy , tags = tags ,
628- version = version if version else self . default_version ,
629- data_converter = self ._data_converter )
630-
631- # Inject the active PRODUCER span context into the request so the sidecar
632- # stores it in the executionStarted event and the worker can parent all
633- # orchestration/activity/timer spans under this trace .
634- parent_trace_ctx = tracing .get_current_trace_context ()
635- if parent_trace_ctx is not None :
636- req .parentTraceContext .CopyFrom (parent_trace_ctx )
637-
638- self ._logger .info (f"Starting new '{ req .name } ' instance with ID = '{ req .instanceId } '." )
639- # Externalize any large payloads in the request
640- if self ._payload_store is not None :
641- payload_helpers .externalize_payloads (
642- req , self ._payload_store , instance_id = req .instanceId ,
643- )
644- res : pb .CreateInstanceResponse = self ._stub .StartInstance (req )
645- return res .instanceId
629+ with tracing .suppress_span_emission ( not self . _emit_trace_spans ):
630+ with tracing . start_create_orchestration_span (
631+ name , resolved_instance_id , version = resolved_version ,
632+ ):
633+ req = build_schedule_new_orchestration_req (
634+ orchestrator , input = input , instance_id = instance_id , start_at = start_at ,
635+ reuse_id_policy = reuse_id_policy , tags = tags ,
636+ version = version if version else self .default_version ,
637+ data_converter = self . _data_converter )
638+
639+ # Inject the active PRODUCER span context, or the ambient context
640+ # when lifecycle span emission is disabled .
641+ parent_trace_ctx = tracing .get_current_trace_context ()
642+ if parent_trace_ctx is not None :
643+ req .parentTraceContext .CopyFrom (parent_trace_ctx )
644+
645+ self ._logger .info (f"Starting new '{ req .name } ' instance with ID = '{ req .instanceId } '." )
646+ # Externalize any large payloads in the request
647+ if self ._payload_store is not None :
648+ payload_helpers .externalize_payloads (
649+ req , self ._payload_store , instance_id = req .instanceId ,
650+ )
651+ res : pb .CreateInstanceResponse = self ._stub .StartInstance (req )
652+ return res .instanceId
646653
647654 def get_orchestration_state (self , instance_id : str , * , fetch_payloads : bool = True ) -> OrchestrationState | None :
648655 req = pb .GetInstanceRequest (instanceId = instance_id , getInputsAndOutputs = fetch_payloads )
@@ -757,14 +764,15 @@ def wait_for_orchestration_completion(self, instance_id: str, *,
757764
758765 def raise_orchestration_event (self , instance_id : str , event_name : str , * ,
759766 data : Any | None = None ) -> None :
760- with tracing .start_raise_event_span (event_name , instance_id ):
761- req = build_raise_event_req (instance_id , event_name , data , self ._data_converter )
762- self ._logger .info (f"Raising event '{ event_name } ' for instance '{ instance_id } '." )
763- if self ._payload_store is not None :
764- payload_helpers .externalize_payloads (
765- req , self ._payload_store , instance_id = instance_id ,
766- )
767- self ._stub .RaiseEvent (req )
767+ with tracing .suppress_span_emission (not self ._emit_trace_spans ):
768+ with tracing .start_raise_event_span (event_name , instance_id ):
769+ req = build_raise_event_req (instance_id , event_name , data , self ._data_converter )
770+ self ._logger .info (f"Raising event '{ event_name } ' for instance '{ instance_id } '." )
771+ if self ._payload_store is not None :
772+ payload_helpers .externalize_payloads (
773+ req , self ._payload_store , instance_id = instance_id ,
774+ )
775+ self ._stub .RaiseEvent (req )
768776
769777 def terminate_orchestration (self , instance_id : str , * ,
770778 output : Any | None = None ,
@@ -940,7 +948,8 @@ def __init__(self, *,
940948 resiliency_options : GrpcClientResiliencyOptions | None = None ,
941949 default_version : str | None = None ,
942950 payload_store : PayloadStore | None = None ,
943- data_converter : DataConverter | None = None ):
951+ data_converter : DataConverter | None = None ,
952+ emit_trace_spans : bool = True ):
944953
945954 self ._owns_channel = channel is None
946955 self ._data_converter = data_converter if data_converter is not None else JsonDataConverter ()
@@ -1005,6 +1014,12 @@ def __init__(self, *,
10051014 self ._logger = shared .get_logger ("async_client" , log_handler , log_formatter )
10061015 self .default_version = default_version
10071016 self ._payload_store = payload_store
1017+ self ._emit_trace_spans = emit_trace_spans
1018+
1019+ @property
1020+ def emit_trace_spans (self ) -> bool :
1021+ """Return whether this client emits Durable Task lifecycle spans."""
1022+ return self ._emit_trace_spans
10081023
10091024 def _compose_interceptors (
10101025 self ,
@@ -1128,27 +1143,28 @@ async def schedule_new_orchestration(self, orchestrator: task.Orchestrator[TInpu
11281143 resolved_instance_id = instance_id if instance_id else uuid .uuid4 ().hex
11291144 resolved_version = version if version else self .default_version
11301145
1131- with tracing .start_create_orchestration_span (
1132- name , resolved_instance_id , version = resolved_version ,
1133- ):
1134- req = build_schedule_new_orchestration_req (
1135- orchestrator , input = input , instance_id = instance_id , start_at = start_at ,
1136- reuse_id_policy = reuse_id_policy , tags = tags ,
1137- version = version if version else self .default_version ,
1138- data_converter = self ._data_converter )
1139-
1140- parent_trace_ctx = tracing .get_current_trace_context ()
1141- if parent_trace_ctx is not None :
1142- req .parentTraceContext .CopyFrom (parent_trace_ctx )
1143-
1144- self ._logger .info (f"Starting new '{ req .name } ' instance with ID = '{ req .instanceId } '." )
1145- # Externalize any large payloads in the request
1146- if self ._payload_store is not None :
1147- await payload_helpers .externalize_payloads_async (
1148- req , self ._payload_store , instance_id = req .instanceId ,
1149- )
1150- res : pb .CreateInstanceResponse = await self ._stub .StartInstance (req )
1151- return res .instanceId
1146+ with tracing .suppress_span_emission (not self ._emit_trace_spans ):
1147+ with tracing .start_create_orchestration_span (
1148+ name , resolved_instance_id , version = resolved_version ,
1149+ ):
1150+ req = build_schedule_new_orchestration_req (
1151+ orchestrator , input = input , instance_id = instance_id , start_at = start_at ,
1152+ reuse_id_policy = reuse_id_policy , tags = tags ,
1153+ version = version if version else self .default_version ,
1154+ data_converter = self ._data_converter )
1155+
1156+ parent_trace_ctx = tracing .get_current_trace_context ()
1157+ if parent_trace_ctx is not None :
1158+ req .parentTraceContext .CopyFrom (parent_trace_ctx )
1159+
1160+ self ._logger .info (f"Starting new '{ req .name } ' instance with ID = '{ req .instanceId } '." )
1161+ # Externalize any large payloads in the request
1162+ if self ._payload_store is not None :
1163+ await payload_helpers .externalize_payloads_async (
1164+ req , self ._payload_store , instance_id = req .instanceId ,
1165+ )
1166+ res : pb .CreateInstanceResponse = await self ._stub .StartInstance (req )
1167+ return res .instanceId
11521168
11531169 async def get_orchestration_state (self , instance_id : str , * ,
11541170 fetch_payloads : bool = True ) -> OrchestrationState | None :
@@ -1262,14 +1278,15 @@ async def wait_for_orchestration_completion(self, instance_id: str, *,
12621278
12631279 async def raise_orchestration_event (self , instance_id : str , event_name : str , * ,
12641280 data : Any | None = None ) -> None :
1265- with tracing .start_raise_event_span (event_name , instance_id ):
1266- req = build_raise_event_req (instance_id , event_name , data , self ._data_converter )
1267- self ._logger .info (f"Raising event '{ event_name } ' for instance '{ instance_id } '." )
1268- if self ._payload_store is not None :
1269- await payload_helpers .externalize_payloads_async (
1270- req , self ._payload_store , instance_id = instance_id ,
1271- )
1272- await self ._stub .RaiseEvent (req )
1281+ with tracing .suppress_span_emission (not self ._emit_trace_spans ):
1282+ with tracing .start_raise_event_span (event_name , instance_id ):
1283+ req = build_raise_event_req (instance_id , event_name , data , self ._data_converter )
1284+ self ._logger .info (f"Raising event '{ event_name } ' for instance '{ instance_id } '." )
1285+ if self ._payload_store is not None :
1286+ await payload_helpers .externalize_payloads_async (
1287+ req , self ._payload_store , instance_id = instance_id ,
1288+ )
1289+ await self ._stub .RaiseEvent (req )
12731290
12741291 async def terminate_orchestration (self , instance_id : str , * ,
12751292 output : Any | None = None ,
0 commit comments