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]
@dataclass
class Event:
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        }
Event( event_type: str, scope_type: str, scope_id: str, event_id: str = <factory>, stamp: datetime.datetime = <factory>, parent_scope_id: str | None = None, payload: dict[str, typing.Any] = <factory>)
event_type: str
scope_type: str
scope_id: str
event_id: str
stamp: datetime.datetime
parent_scope_id: str | None = None
payload: dict[str, typing.Any]
def encode(self) -> dict[str, typing.Any]:
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.

class Timestamp:
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.

@staticmethod
def now() -> datetime.datetime:
18    @staticmethod
19    def now() -> datetime:
20        return datetime.now(timezone.utc)
class Observer:
 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                )
def configure_writers( self, writers: list[Writer]) -> None:
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)

Register writers in bulk, to be brought up on the next start().

Must be called before start() — raises RuntimeError if the observer is already running. Replaces any previously configured pending writers.

async def start(self) -> None:
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.

async def shutdown(self) -> None:
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)
async def register( self, writer: Writer, name: str, maxsize: int = 10000, policy: QueuePolicy = <QueuePolicy.DROP_OLDEST: 'drop_oldest'>) -> None:
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.

async def unregister(self, name: str) -> None:
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)
is_observing: bool
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.

loop: asyncio.events.AbstractEventLoop | None
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.

async def publish(self, event: Event) -> None:
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.
def publish_nowait(self, event: Event) -> None:
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)

Sync body of publish. Call it on loop; other threads go through publish_event_nowait.

Arguments:
  • event: The event to publish.
class QueuePolicy(enum.Enum):
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.

DROP_OLDEST = <QueuePolicy.DROP_OLDEST: 'drop_oldest'>
class Writer(typing.Protocol):
 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:
        ...
Writer(*args, **kwargs)
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)
async def start(self) -> None:
10    async def start(self) -> None: ...
async def write(self, event: Event) -> None:
11    async def write(self, event: Event) -> None: ...
async def shutdown(self) -> None:
12    async def shutdown(self) -> None: ...
class JsonlWriter:
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()
JsonlWriter(directory: pathlib.Path | None = None)
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.

async def start(self) -> None:
19    async def start(self) -> None:
20        self._directory.mkdir(parents=True, exist_ok=True)
async def write(self, event: Event) -> None:
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()
async def shutdown(self) -> None:
33    async def shutdown(self) -> None:
34        for handle in self._files.values():
35            handle.flush()
36            handle.close()
37        self._files.clear()
SCOPE_SESSION = 'session'
SCOPE_RETRIEVAL = 'retrieval'
SCOPE_EVALUATION = 'evaluation'
def configure_writers(writers: list[Writer]) -> None:
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.

async def publish_event(event: Event) -> None:
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.

async def ensure_started() -> 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.

async def shutdown() -> None:
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.

EVENTS_DIR_ENV = 'RAILTRACKS_EVENTS_DIR'
EVENTS_SUBDIR = PosixPath('data/events')
def resolve_events_dir() -> pathlib.Path:
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.