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

1"""Framework-neutral broker and stream-processing observability context. 

2 

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

7 

8from __future__ import annotations 

9 

10from collections.abc import Mapping 

11 

12from .observability import OperationContext 

13from .trace_context import W3CTraceContext, trace_context_from_mapping 

14 

15__all__ = ["BrokerContext", "broker_context_from_message"] 

16 

17_VALID_OPERATIONS = {"publish", "consume", "acknowledge"} 

18 

19 

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 

29 

30 

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) 

35 

36 

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 

45 

46 

47class BrokerContext: 

48 """Broker operation context plus optional message coordinates.""" 

49 

50 __slots__ = ( 

51 "consumer_group", 

52 "message_id", 

53 "offset", 

54 "operation_context", 

55 "partition", 

56 "topic", 

57 "trace_context", 

58 ) 

59 

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 

78 

79 

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. 

94 

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

106 

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

116 

117 if metadata is None: 

118 metadata = {} 

119 if not isinstance(metadata, Mapping): 

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

121 

122 operation_context = OperationContext( 

123 system=system, 

124 component=component, 

125 operation=f"{normalized_operation} {topic_name}", 

126 correlation_id=correlation_identifier, 

127 ) 

128 

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) 

132 

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 )