Coverage for dataexcept/broker_context.py: 93%
58 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 broker and stream-processing observability context.
3The adapter accepts plain message metadata so Kafka, RabbitMQ, Pulsar, NATS,
4stream processors and custom brokers can share one dependency-free boundary
5model. Message bodies and application payloads 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__ = ["BrokerContext", "broker_context_from_message"]
17_VALID_OPERATIONS = {"publish", "consume", "acknowledge"}
20def _single_line(value: str, field: str) -> str:
21 if not isinstance(value, str): 21 ↛ 22line 21 didn't jump to line 22 because the condition on line 21 was never true
22 raise TypeError(f"{field} must be a string")
23 normalized = value.strip()
24 if not normalized:
25 raise ValueError(f"{field} must not be empty")
26 if "\r" in value or "\n" in value:
27 raise ValueError(f"{field} must be a single line")
28 return normalized
31def _optional_single_line(value: str | None, field: str) -> str | None:
32 if value is None:
33 return None
34 return _single_line(value, field)
37def _optional_non_negative_int(value: int | None, field: str) -> int | None:
38 if value is None:
39 return None
40 if isinstance(value, bool) or not isinstance(value, int): 40 ↛ 41line 40 didn't jump to line 41 because the condition on line 40 was never true
41 raise TypeError(f"{field} must be an integer or None")
42 if value < 0:
43 raise ValueError(f"{field} must be non-negative")
44 return value
47class BrokerContext:
48 """Broker operation context plus optional message coordinates."""
50 __slots__ = (
51 "consumer_group",
52 "message_id",
53 "offset",
54 "operation_context",
55 "partition",
56 "topic",
57 "trace_context",
58 )
60 def __init__(
61 self,
62 operation_context: OperationContext,
63 trace_context: W3CTraceContext | None,
64 *,
65 topic: str,
66 partition: int | None,
67 offset: int | None,
68 consumer_group: str | None,
69 message_id: str | None,
70 ) -> None:
71 self.operation_context = operation_context
72 self.trace_context = trace_context
73 self.topic = topic
74 self.partition = partition
75 self.offset = offset
76 self.consumer_group = consumer_group
77 self.message_id = message_id
80def broker_context_from_message(
81 operation: str,
82 topic: str,
83 *,
84 metadata: Mapping[str, object] | None = None,
85 system: str | None = "broker",
86 component: str | None = None,
87 correlation_id: str | None = None,
88 partition: int | None = None,
89 offset: int | None = None,
90 consumer_group: str | None = None,
91 message_id: str | None = None,
92) -> BrokerContext:
93 """Build observability context for a broker or stream-processing boundary.
95 The stable operation identity is "<operation> <topic>", where operation is
96 one of publish, consume or acknowledge. Partition, offset, consumer-group
97 and message identifiers remain wrapper metadata instead of becoming
98 low-cardinality operation labels. Incoming W3C Trace Context is preserved
99 when present.
100 """
101 if not isinstance(operation, str): 101 ↛ 102line 101 didn't jump to line 102 because the condition on line 101 was never true
102 raise TypeError("operation must be a string")
103 normalized_operation = operation.strip().lower()
104 if normalized_operation not in _VALID_OPERATIONS:
105 raise ValueError("operation must be one of publish, consume or acknowledge")
107 topic_name = _single_line(topic, "topic")
108 correlation_identifier = _optional_single_line(
109 correlation_id,
110 "correlation_id",
111 )
112 group = _optional_single_line(consumer_group, "consumer_group")
113 message = _optional_single_line(message_id, "message_id")
114 partition_number = _optional_non_negative_int(partition, "partition")
115 offset_number = _optional_non_negative_int(offset, "offset")
117 if metadata is None:
118 metadata = {}
119 if not isinstance(metadata, Mapping):
120 raise TypeError("metadata must be a mapping or None")
122 operation_context = OperationContext(
123 system=system,
124 component=component,
125 operation=f"{normalized_operation} {topic_name}",
126 correlation_id=correlation_identifier,
127 )
129 trace_context = trace_context_from_mapping(metadata)
130 if trace_context is not None:
131 operation_context = trace_context.to_operation_context(operation_context)
133 return BrokerContext(
134 operation_context,
135 trace_context,
136 topic=topic_name,
137 partition=partition_number,
138 offset=offset_number,
139 consumer_group=group,
140 message_id=message,
141 )