Coverage for dataexcept/orchestrator_context.py: 94%
48 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-27 14:44 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-27 14:44 +0000
1"""Framework-neutral workflow/orchestrator observability context.
3The adapter accepts plain workflow metadata so Airflow, Dagster, Prefect,
4Argo and custom schedulers can share the same dependency-free boundary model.
5Payloads and task arguments are intentionally excluded.
6"""
8from __future__ import annotations
10from collections.abc import Mapping
12from .observability import OperationContext
13from .trace_context import W3CTraceContext, trace_context_from_mapping
15__all__ = ["OrchestratorContext", "orchestrator_context_from_step"]
18def _single_line(value: str, field: str) -> str:
19 if not isinstance(value, str): 19 ↛ 20line 19 didn't jump to line 20 because the condition on line 19 was never true
20 raise TypeError(f"{field} must be a string")
21 normalized = value.strip()
22 if not normalized:
23 raise ValueError(f"{field} must not be empty")
24 if "\r" in value or "\n" in value:
25 raise ValueError(f"{field} must be a single line")
26 return normalized
29def _optional_single_line(value: str | None, field: str) -> str | None:
30 if value is None:
31 return None
32 return _single_line(value, field)
35def _validate_attempt(attempt: int | None) -> int | None:
36 if attempt is None:
37 return None
38 if isinstance(attempt, bool) or not isinstance(attempt, int): 38 ↛ 39line 38 didn't jump to line 39 because the condition on line 38 was never true
39 raise TypeError("attempt must be an integer or None")
40 if attempt < 0:
41 raise ValueError("attempt must be non-negative")
42 return attempt
45class OrchestratorContext:
46 """Workflow-step context plus optional trace and retry metadata."""
48 __slots__ = ("attempt", "operation_context", "step_run_id", "trace_context")
50 def __init__(
51 self,
52 operation_context: OperationContext,
53 trace_context: W3CTraceContext | None,
54 *,
55 step_run_id: str | None,
56 attempt: int | None,
57 ) -> None:
58 self.operation_context = operation_context
59 self.trace_context = trace_context
60 self.step_run_id = step_run_id
61 self.attempt = attempt
64def orchestrator_context_from_step(
65 workflow: str,
66 step: str,
67 *,
68 run_id: str | None = None,
69 step_run_id: str | None = None,
70 metadata: Mapping[str, object] | None = None,
71 system: str | None = "orchestrator",
72 component: str | None = None,
73 correlation_id: str | None = None,
74 attempt: int | None = None,
75) -> OrchestratorContext:
76 """Build observability context for one workflow/orchestrator step.
78 The operation name is the stable ``workflow:step`` pair. Run identifiers
79 remain correlation metadata and are never folded into the operation name.
80 Incoming W3C trace context is preserved when present.
81 """
82 workflow_name = _single_line(workflow, "workflow")
83 step_name = _single_line(step, "step")
84 run_identifier = _optional_single_line(run_id, "run_id")
85 step_run_identifier = _optional_single_line(step_run_id, "step_run_id")
86 correlation_identifier = _optional_single_line(correlation_id, "correlation_id")
88 if metadata is None:
89 metadata = {}
90 if not isinstance(metadata, Mapping):
91 raise TypeError("metadata must be a mapping or None")
93 operation_context = OperationContext(
94 system=system,
95 component=component,
96 operation=f"{workflow_name}:{step_name}",
97 job_id=run_identifier,
98 correlation_id=correlation_identifier,
99 )
101 trace_context = trace_context_from_mapping(metadata)
102 if trace_context is not None:
103 operation_context = trace_context.to_operation_context(operation_context)
105 return OrchestratorContext(
106 operation_context,
107 trace_context,
108 step_run_id=step_run_identifier,
109 attempt=_validate_attempt(attempt),
110 )