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

1"""Framework-neutral task/worker context for DataExcept observability. 

2 

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""" 

7 

8from __future__ import annotations 

9 

10from collections.abc import Mapping, Sequence 

11 

12from .observability import OperationContext 

13from .trace_context import W3CTraceContext, trace_context_from_mapping 

14 

15__all__ = ["WorkerContext", "worker_context_from_task"] 

16 

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) 

23 

24 

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 

32 

33 

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 

41 

42 

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 

49 

50 

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") 

54 

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) 

64 

65 

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 

75 

76 

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 

85 

86 

87class WorkerContext: 

88 """Task operation context plus optional trace and retry metadata.""" 

89 

90 __slots__ = ("attempt", "operation_context", "trace_context") 

91 

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 

101 

102 

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. 

116 

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") 

129 

130 if metadata is None: 

131 metadata = {} 

132 if not isinstance(metadata, Mapping): 

133 raise TypeError("metadata must be a mapping or None") 

134 

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 ) 

141 

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 ) 

147 

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 ) 

158 

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) 

162 

163 return WorkerContext( 

164 operation_context, 

165 trace_context, 

166 _validate_attempt(attempt), 

167 )