Coverage for dataexcept/worker_context.py: 93%
81 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 task/worker context for DataExcept observability.
3The adapter accepts plain task metadata so Celery, RQ, Arq, Dramatiq, custom
4workers and orchestration executors can share one dependency-free boundary
5model. Task arguments and payloads are intentionally out of scope.
6"""
8from __future__ import annotations
10from collections.abc import Mapping, Sequence
12from .observability import OperationContext
13from .trace_context import W3CTraceContext, trace_context_from_mapping
15__all__ = ["WorkerContext", "worker_context_from_task"]
17_DEFAULT_JOB_ID_KEYS = ("job_id", "task_id", "id")
18_DEFAULT_CORRELATION_ID_KEYS = (
19 "correlation_id",
20 "x-correlation-id",
21 "correlation-id",
22)
25def _safe_text(value: object) -> str | None:
26 if not isinstance(value, str):
27 return None
28 stripped = value.strip()
29 if not stripped or "\r" in value or "\n" in value:
30 return None
31 return stripped
34def _validated_optional_text(value: str | None, field: str) -> str | None:
35 if value is None:
36 return None
37 normalized = _safe_text(value)
38 if normalized is None:
39 raise ValueError(f"{field} must be non-empty single-line text or None")
40 return normalized
43def _normalized_metadata(metadata: Mapping[str, object]) -> dict[str, object]:
44 normalized: dict[str, object] = {}
45 for key, value in metadata.items():
46 if isinstance(key, str): 46 ↛ 45line 46 didn't jump to line 45 because the condition on line 46 was always true
47 normalized[key.lower()] = value
48 return normalized
51def _validated_keys(keys: Sequence[str], field: str) -> tuple[str, ...]:
52 if isinstance(keys, (str, bytes)): 52 ↛ 53line 52 didn't jump to line 53 because the condition on line 52 was never true
53 raise TypeError(f"{field} must be a sequence of metadata keys")
55 result: list[str] = []
56 for key in keys:
57 if not isinstance(key, str): 57 ↛ 58line 57 didn't jump to line 58 because the condition on line 57 was never true
58 raise TypeError(f"{field} entries must be strings")
59 normalized = key.strip().lower()
60 if not normalized: 60 ↛ 61line 60 didn't jump to line 61 because the condition on line 60 was never true
61 raise ValueError(f"{field} entries must not be empty")
62 result.append(normalized)
63 return tuple(result)
66def _first_metadata_value(
67 metadata: Mapping[str, object],
68 keys: Sequence[str],
69) -> str | None:
70 for key in keys:
71 value = _safe_text(metadata.get(key.lower()))
72 if value is not None:
73 return value
74 return None
77def _validate_attempt(attempt: int | None) -> int | None:
78 if attempt is None:
79 return None
80 if isinstance(attempt, bool) or not isinstance(attempt, int):
81 raise TypeError("attempt must be an integer or None")
82 if attempt < 0:
83 raise ValueError("attempt must be non-negative")
84 return attempt
87class WorkerContext:
88 """Task operation context plus optional trace and retry metadata."""
90 __slots__ = ("attempt", "operation_context", "trace_context")
92 def __init__(
93 self,
94 operation_context: OperationContext,
95 trace_context: W3CTraceContext | None,
96 attempt: int | None,
97 ) -> None:
98 self.operation_context = operation_context
99 self.trace_context = trace_context
100 self.attempt = attempt
103def worker_context_from_task(
104 task_name: str,
105 *,
106 job_id: str | None = None,
107 metadata: Mapping[str, object] | None = None,
108 system: str | None = "worker",
109 component: str | None = None,
110 correlation_id: str | None = None,
111 attempt: int | None = None,
112 job_id_keys: Sequence[str] = _DEFAULT_JOB_ID_KEYS,
113 correlation_id_keys: Sequence[str] = _DEFAULT_CORRELATION_ID_KEYS,
114) -> WorkerContext:
115 """Build worker observability context from plain task metadata.
117 ``task_name`` should be the stable registered task name, never arguments or
118 a rendered payload. Explicit ``job_id`` and ``correlation_id`` values take
119 precedence over metadata fallbacks. W3C trace context is propagated when
120 present and no identifiers are generated when it is absent.
121 """
122 if not isinstance(task_name, str): 122 ↛ 123line 122 didn't jump to line 123 because the condition on line 122 was never true
123 raise TypeError("task_name must be a string")
124 normalized_task = task_name.strip()
125 if not normalized_task:
126 raise ValueError("task_name must not be empty")
127 if "\r" in task_name or "\n" in task_name:
128 raise ValueError("task_name must be a single line")
130 if metadata is None:
131 metadata = {}
132 if not isinstance(metadata, Mapping):
133 raise TypeError("metadata must be a mapping or None")
135 normalized = _normalized_metadata(metadata)
136 job_keys = _validated_keys(job_id_keys, "job_id_keys")
137 correlation_keys = _validated_keys(
138 correlation_id_keys,
139 "correlation_id_keys",
140 )
142 explicit_job_id = _validated_optional_text(job_id, "job_id")
143 explicit_correlation_id = _validated_optional_text(
144 correlation_id,
145 "correlation_id",
146 )
148 operation_context = OperationContext(
149 system=system,
150 component=component,
151 operation=normalized_task,
152 job_id=explicit_job_id or _first_metadata_value(normalized, job_keys),
153 correlation_id=(
154 explicit_correlation_id
155 or _first_metadata_value(normalized, correlation_keys)
156 ),
157 )
159 trace_context = trace_context_from_mapping(normalized)
160 if trace_context is not None:
161 operation_context = trace_context.to_operation_context(operation_context)
163 return WorkerContext(
164 operation_context,
165 trace_context,
166 _validate_attempt(attempt),
167 )