railtracks.observability
Observability submodule: streaming Event pipeline with per-writer queues, plus a process-wide default Observer.
1"""Observability submodule: streaming Event pipeline with per-writer queues, 2plus a process-wide default Observer. 3""" 4 5from .configure import ( 6 configure_writers, 7 ensure_started, 8 shutdown, 9) 10from .models import ( 11 SCOPE_EVALUATION, 12 SCOPE_RETRIEVAL, 13 SCOPE_SESSION, 14 Event, 15 Timestamp, 16) 17from .observer import Observer, QueuePolicy 18from .publish import publish_event 19from .storage import EVENTS_DIR_ENV, EVENTS_SUBDIR, resolve_events_dir 20from .writers import JsonlWriter, Writer 21 22__all__ = [ 23 "Event", 24 "Timestamp", 25 "Observer", 26 "QueuePolicy", 27 "Writer", 28 "JsonlWriter", 29 "SCOPE_SESSION", 30 "SCOPE_RETRIEVAL", 31 "SCOPE_EVALUATION", 32 "configure_writers", 33 "publish_event", 34 "ensure_started", 35 "shutdown", 36 "EVENTS_DIR_ENV", 37 "EVENTS_SUBDIR", 38 "resolve_events_dir", 39]
23@dataclass 24class Event: 25 event_type: str 26 scope_type: str 27 scope_id: str 28 event_id: str = field(default_factory=lambda: str(uuid.uuid4())) 29 stamp: datetime = field(default_factory=Timestamp.now) 30 parent_scope_id: str | None = None 31 payload: dict[str, Any] = field(default_factory=dict) 32 33 def encode(self) -> dict[str, Any]: 34 """Encode the event to a JSON-serializable dictionary.""" 35 return { 36 "event_type": self.event_type, 37 "scope_type": self.scope_type, 38 "scope_id": self.scope_id, 39 "event_id": self.event_id, 40 "stamp": self.stamp.isoformat(), 41 "parent_scope_id": self.parent_scope_id, 42 "payload": self.payload, 43 }
33 def encode(self) -> dict[str, Any]: 34 """Encode the event to a JSON-serializable dictionary.""" 35 return { 36 "event_type": self.event_type, 37 "scope_type": self.scope_type, 38 "scope_id": self.scope_id, 39 "event_id": self.event_id, 40 "stamp": self.stamp.isoformat(), 41 "parent_scope_id": self.parent_scope_id, 42 "payload": self.payload, 43 }
Encode the event to a JSON-serializable dictionary.
14class Timestamp: 15 """Namespace helper for constructing `Event.stamp`. The field itself is a plain tz-aware UTC datetime.""" 16 17 # Doing it this way to make potential changes easier 18 @staticmethod 19 def now() -> datetime: 20 return datetime.now(timezone.utc)
Namespace helper for constructing Event.stamp. The field itself is a plain tz-aware UTC datetime.
38class Observer: 39 def __init__(self) -> None: 40 self._writers: dict[str, _Entry] = {} 41 self._drops: dict[str, int] = {} # dropped events per writer 42 self._running = False 43 self._pending_writers: list[Writer] = [] 44 self._start_lock: asyncio.Lock = asyncio.Lock() 45 self._loop: asyncio.AbstractEventLoop | None = None 46 47 # async context manager support added for now, this will become more clear 48 # once we move to integrating with the other modules 49 async def __aenter__(self) -> Observer: 50 await self.start() 51 return self 52 53 async def __aexit__(self, exc_type, exc, tb) -> None: 54 await self.shutdown() 55 56 def configure_writers(self, writers: list[Writer]) -> None: 57 """Register writers in bulk, to be brought up on the next `start()`. 58 59 Must be called before `start()` — raises `RuntimeError` if the observer 60 is already running. Replaces any previously configured pending writers. 61 """ 62 if self._running: 63 raise RuntimeError( 64 "configure_writers must be called before start(); use register() to add writers after." 65 ) 66 self._pending_writers = list(writers) 67 68 async def start(self) -> None: 69 """Bring up the observer. 70 71 This is idempotent safe to call multiple times. 72 """ 73 if self._running: 74 return 75 async with self._start_lock: 76 if self._running: 77 return 78 self._loop = asyncio.get_running_loop() 79 for i, writer in enumerate(self._pending_writers): 80 name = f"writer-{i}" 81 try: 82 await self._register_impl(writer, name) 83 except Exception as exc: 84 logger.warning( 85 "observability writer %r failed to start (%s: %s); " 86 "continuing without it.", 87 name, 88 type(exc).__name__, 89 exc, 90 ) 91 self._running = True 92 93 async def shutdown(self) -> None: 94 if not self._running: 95 return 96 self._running = False 97 self._loop = None 98 for name in list(self._writers.keys()): 99 await self._teardown(name) 100 101 async def register( 102 self, 103 writer: Writer, 104 name: str, 105 maxsize: int = 10_000, 106 policy: QueuePolicy = QueuePolicy.DROP_OLDEST, 107 ) -> None: 108 """Register a writer on a running observer. Post-start only. 109 110 For a pre-start batch, use `configure_writers()` and let `start()` 111 register them. 112 """ 113 if not self._running: 114 raise RuntimeError( 115 "Observer is not running; call start() first, or use configure_writers() " 116 "for pre-start batch registration." 117 ) 118 await self._register_impl(writer, name, maxsize=maxsize, policy=policy) 119 120 async def _register_impl( 121 self, 122 writer: Writer, 123 name: str, 124 maxsize: int = 10_000, 125 policy: QueuePolicy = QueuePolicy.DROP_OLDEST, 126 ) -> None: 127 """Internal registration path used by both public `register` (post-start) 128 and `start()` (registering pending writers during startup). Skips the 129 `_running` check so `start()` can register pending writers before 130 flipping the flag.""" 131 if name in self._writers: 132 raise ValueError(f"Writer {name!r} is already registered.") 133 await writer.start() 134 queue: asyncio.Queue[_QueueItem] = asyncio.Queue( 135 maxsize=maxsize 136 ) # Each writer has its own queue 137 task = asyncio.create_task( 138 self._consumer_loop(name, writer, queue), 139 name=f"observer-consumer:{name}", 140 ) 141 self._writers[name] = _Entry( 142 writer=writer, queue=queue, task=task, policy=policy 143 ) 144 self._drops[name] = 0 145 146 async def unregister(self, name: str) -> None: 147 if name not in self._writers: 148 raise KeyError(f"No writer registered as {name!r}.") 149 await self._teardown(name) 150 151 @property 152 def is_observing(self) -> bool: 153 """Whether a published event would reach at least one writer.""" 154 return self._running and bool(self._writers) 155 156 @property 157 def loop(self) -> asyncio.AbstractEventLoop | None: 158 """The loop the writer tasks run on, or None while stopped.""" 159 return self._loop 160 161 async def publish(self, event: Event) -> None: 162 """Fan the event out to every registered writer's queue. 163 164 `async` on this method is contract-enforcement, the body doesn't `await` anything. 165 requiring callers to be inside a coroutine means they're on the same running loop 166 167 Args: 168 event: The event to publish. 169 """ 170 self.publish_nowait(event) 171 172 def publish_nowait(self, event: Event) -> None: 173 """Sync body of `publish`. Call it on `loop`; other threads go through 174 `publish_event_nowait`. 175 176 Args: 177 event: The event to publish. 178 """ 179 if not self._running: 180 raise RuntimeError("Observer is not running.") 181 for name, entry in self._writers.items(): 182 try: 183 entry.queue.put_nowait(event) 184 except asyncio.QueueFull: 185 self._handle_full_queue(name, entry, event) 186 187 def _handle_full_queue(self, name: str, entry: _Entry, event: Event) -> None: 188 match entry.policy: 189 case QueuePolicy.DROP_OLDEST: 190 self._drop_oldest(name, entry, event) 191 192 def _drop_oldest(self, name: str, entry: _Entry, event: Event) -> None: 193 try: 194 entry.queue.get_nowait() 195 except asyncio.QueueEmpty: 196 pass 197 entry.queue.put_nowait(event) 198 self._drops[name] += 1 199 logger.warning( 200 "observability writer %r queue full; dropped oldest event " 201 "(policy=%s, total drops for this writer: %d)", 202 name, 203 entry.policy.value, 204 self._drops[name], 205 ) 206 207 async def _teardown(self, name: str) -> None: 208 entry = self._writers.pop(name) 209 self._drops.pop(name, None) 210 _enqueue_end(entry.queue) 211 await entry.task 212 await entry.writer.shutdown() 213 214 async def _consumer_loop( 215 self, name: str, writer: Writer, queue: asyncio.Queue[_QueueItem] 216 ) -> None: 217 while True: 218 item = await queue.get() 219 if isinstance(item, _EndOfStream): 220 return 221 try: 222 await writer.write(item) 223 except Exception as exc: 224 logger.warning( 225 "observability writer %r failed on event %s: %s", 226 name, 227 item.event_id, 228 exc, 229 )
56 def configure_writers(self, writers: list[Writer]) -> None: 57 """Register writers in bulk, to be brought up on the next `start()`. 58 59 Must be called before `start()` — raises `RuntimeError` if the observer 60 is already running. Replaces any previously configured pending writers. 61 """ 62 if self._running: 63 raise RuntimeError( 64 "configure_writers must be called before start(); use register() to add writers after." 65 ) 66 self._pending_writers = list(writers)
68 async def start(self) -> None: 69 """Bring up the observer. 70 71 This is idempotent safe to call multiple times. 72 """ 73 if self._running: 74 return 75 async with self._start_lock: 76 if self._running: 77 return 78 self._loop = asyncio.get_running_loop() 79 for i, writer in enumerate(self._pending_writers): 80 name = f"writer-{i}" 81 try: 82 await self._register_impl(writer, name) 83 except Exception as exc: 84 logger.warning( 85 "observability writer %r failed to start (%s: %s); " 86 "continuing without it.", 87 name, 88 type(exc).__name__, 89 exc, 90 ) 91 self._running = True
Bring up the observer.
This is idempotent safe to call multiple times.
101 async def register( 102 self, 103 writer: Writer, 104 name: str, 105 maxsize: int = 10_000, 106 policy: QueuePolicy = QueuePolicy.DROP_OLDEST, 107 ) -> None: 108 """Register a writer on a running observer. Post-start only. 109 110 For a pre-start batch, use `configure_writers()` and let `start()` 111 register them. 112 """ 113 if not self._running: 114 raise RuntimeError( 115 "Observer is not running; call start() first, or use configure_writers() " 116 "for pre-start batch registration." 117 ) 118 await self._register_impl(writer, name, maxsize=maxsize, policy=policy)
Register a writer on a running observer. Post-start only.
For a pre-start batch, use configure_writers() and let start()
register them.
151 @property 152 def is_observing(self) -> bool: 153 """Whether a published event would reach at least one writer.""" 154 return self._running and bool(self._writers)
Whether a published event would reach at least one writer.
156 @property 157 def loop(self) -> asyncio.AbstractEventLoop | None: 158 """The loop the writer tasks run on, or None while stopped.""" 159 return self._loop
The loop the writer tasks run on, or None while stopped.
161 async def publish(self, event: Event) -> None: 162 """Fan the event out to every registered writer's queue. 163 164 `async` on this method is contract-enforcement, the body doesn't `await` anything. 165 requiring callers to be inside a coroutine means they're on the same running loop 166 167 Args: 168 event: The event to publish. 169 """ 170 self.publish_nowait(event)
Fan the event out to every registered writer's queue.
async on this method is contract-enforcement, the body doesn't await anything.
requiring callers to be inside a coroutine means they're on the same running loop
Arguments:
- event: The event to publish.
172 def publish_nowait(self, event: Event) -> None: 173 """Sync body of `publish`. Call it on `loop`; other threads go through 174 `publish_event_nowait`. 175 176 Args: 177 event: The event to publish. 178 """ 179 if not self._running: 180 raise RuntimeError("Observer is not running.") 181 for name, entry in self._writers.items(): 182 try: 183 entry.queue.put_nowait(event) 184 except asyncio.QueueFull: 185 self._handle_full_queue(name, entry, event)
15class QueuePolicy(Enum): 16 """How a writer's queue behaves when it's full at publish time.""" 17 18 DROP_OLDEST = "drop_oldest"
How a writer's queue behaves when it's full at publish time.
9class Writer(Protocol): 10 async def start(self) -> None: ... 11 async def write(self, event: Event) -> None: ... 12 async def shutdown(self) -> None: ...
Base class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol):
def meth(self) -> int:
...
Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing), for example::
class C:
def meth(self) -> int:
return 0
def func(x: Proto) -> int:
return x.meth()
func(C()) # Passes static type check
See PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]):
def meth(self) -> T:
...
1431def _no_init_or_replace_init(self, *args, **kwargs): 1432 cls = type(self) 1433 1434 if cls._is_protocol: 1435 raise TypeError('Protocols cannot be instantiated') 1436 1437 # Already using a custom `__init__`. No need to calculate correct 1438 # `__init__` to call. This can lead to RecursionError. See bpo-45121. 1439 if cls.__init__ is not _no_init_or_replace_init: 1440 return 1441 1442 # Initially, `__init__` of a protocol subclass is set to `_no_init_or_replace_init`. 1443 # The first instantiation of the subclass will call `_no_init_or_replace_init` which 1444 # searches for a proper new `__init__` in the MRO. The new `__init__` 1445 # replaces the subclass' old `__init__` (ie `_no_init_or_replace_init`). Subsequent 1446 # instantiation of the protocol subclass will thus use the new 1447 # `__init__` and no longer call `_no_init_or_replace_init`. 1448 for base in cls.__mro__: 1449 init = base.__dict__.get('__init__', _no_init_or_replace_init) 1450 if init is not _no_init_or_replace_init: 1451 cls.__init__ = init 1452 break 1453 else: 1454 # should not happen 1455 cls.__init__ = object.__init__ 1456 1457 cls.__init__(self, *args, **kwargs)
13class JsonlWriter: 14 def __init__(self, directory: Path | None = None): 15 """Write events to ``directory`` or the shared visualizer event store.""" 16 self._directory = directory if directory is not None else resolve_events_dir() 17 self._files: dict[str, TextIO] = {} 18 19 async def start(self) -> None: 20 self._directory.mkdir(parents=True, exist_ok=True) 21 22 async def write(self, event: Event) -> None: 23 handle = self._files.get(event.scope_id) 24 if handle is None: 25 _check_safe_scope_id(event.scope_id) 26 handle = (self._directory / f"{event.scope_id}.jsonl").open( 27 "a", encoding="utf-8" 28 ) 29 self._files[event.scope_id] = handle 30 handle.write(_serialize(event) + "\n") 31 handle.flush() 32 33 async def shutdown(self) -> None: 34 for handle in self._files.values(): 35 handle.flush() 36 handle.close() 37 self._files.clear()
14 def __init__(self, directory: Path | None = None): 15 """Write events to ``directory`` or the shared visualizer event store.""" 16 self._directory = directory if directory is not None else resolve_events_dir() 17 self._files: dict[str, TextIO] = {}
Write events to directory or the shared visualizer event store.
22 async def write(self, event: Event) -> None: 23 handle = self._files.get(event.scope_id) 24 if handle is None: 25 _check_safe_scope_id(event.scope_id) 26 handle = (self._directory / f"{event.scope_id}.jsonl").open( 27 "a", encoding="utf-8" 28 ) 29 self._files[event.scope_id] = handle 30 handle.write(_serialize(event) + "\n") 31 handle.flush()
23def configure_writers(writers: list[Writer]) -> None: 24 """Set the writers to register on the singleton Observer on first start(). 25 26 Delegates to `observer.configure_writers`. Must be called before the 27 observer has started; raises `RuntimeError` otherwise. 28 """ 29 observer.configure_writers(writers)
Set the writers to register on the singleton Observer on first start().
Delegates to observer.configure_writers. Must be called before the
observer has started; raises RuntimeError otherwise.
16async def publish_event(event: Event) -> None: 17 """Convenience wrapper to publish an Event via the process-wide singleton Observer. 18 19 Inline listeners run first and are isolated from each other: one raising must 20 neither lose the event for the others nor stop it reaching the Observer. 21 """ 22 _run_inline_listeners(event) 23 24 await configure.observer.publish(event)
Convenience wrapper to publish an Event via the process-wide singleton Observer.
Inline listeners run first and are isolated from each other: one raising must neither lose the event for the others nor stop it reaching the Observer.
44async def ensure_started() -> Observer: 45 """Start the singleton observer if not already started, return it. 46 47 Auto-registers a default `JsonlWriter` when no writers have been configured 48 and `RAILTRACKS_DISABLE_EVENTS` is unset. 49 """ 50 if ( 51 not observer._running 52 and not observer._pending_writers 53 and not _disable_events() 54 ): 55 observer.configure_writers([JsonlWriter()]) 56 await observer.start() 57 return observer
Start the singleton observer if not already started, return it.
Auto-registers a default JsonlWriter when no writers have been configured
and RAILTRACKS_DISABLE_EVENTS is unset.
60async def shutdown() -> None: 61 """Drain per-writer queues and stop the singleton Observer's consumer tasks. 62 63 Safe to call when the observer isn't running. 64 """ 65 await observer.shutdown()
Drain per-writer queues and stop the singleton Observer's consumer tasks.
Safe to call when the observer isn't running.
15def resolve_events_dir() -> Path: 16 """Return the JSONL event store used by the observer and visualizer. 17 18 ``RAILTRACKS_EVENTS_DIR`` may point at an alternate store for development, 19 imports, or tests. Relative overrides are resolved from the current working 20 directory; without one, events live under the resolved Railtracks home. 21 """ 22 override = os.environ.get(EVENTS_DIR_ENV) 23 if override: 24 return Path(override).expanduser().resolve() 25 return resolve_railtracks_home() / EVENTS_SUBDIR
Return the JSONL event store used by the observer and visualizer.
RAILTRACKS_EVENTS_DIR may point at an alternate store for development,
imports, or tests. Relative overrides are resolved from the current working
directory; without one, events live under the resolved Railtracks home.