Source code for stormlog.telemetry_sink

"""Append-only telemetry sink with rollover and retention bounds."""

from __future__ import annotations

import json
import logging
import os
import threading
import time
from dataclasses import asdict, dataclass, field
from pathlib import Path
from typing import TYPE_CHECKING, Any, Mapping, cast

from .collector_health import collector_retry_delay_seconds
from .session import (
    SESSION_STATUS_COMPLETED,
    SESSION_STATUS_INTERRUPTED,
    SESSION_STATUS_RUNNING,
    SessionSummary,
    create_session_summary,
    now_ns,
    session_summary_from_dict,
    session_summary_to_dict,
    update_session_summary,
)

if TYPE_CHECKING:
    from .telemetry_rollups import RollupCoverage

MANIFEST_FILENAME = "manifest.json"
SEGMENT_PREFIX = "segment-"
SEGMENT_SUFFIX = ".jsonl"
SINK_SCHEMA_VERSION = 2
_LOGGER = logging.getLogger(__name__)
# How much of a segment the repair of its last line reads at a time.
_TAIL_CHUNK_BYTES = 64 * 1024


[docs] @dataclass class TelemetrySinkConfig: """Runtime policy for append-only telemetry persistence.""" root_dir: Path flush_every_events: int = 50 flush_every_seconds: float = 2.0 rollover_max_bytes: int = 64 * 1024 * 1024 rollover_max_events: int = 10000 retention_max_files: int = 8 retention_max_total_bytes: int = 512 * 1024 * 1024 write_rollups: bool = True rollup_window_seconds: int = 60 # Bounded mode, off by default: at most this many bytes wait in memory # for a flush. A record that would go over is dropped and counted, and a # failed flush is counted and retried after a growing backoff instead of # raising, so a full or failing disk cannot grow the caller's memory. max_buffer_bytes: int | None = None failure_backoff_seconds: float = 1.0 failure_backoff_max_seconds: float = 60.0 def __post_init__(self) -> None: self.root_dir = Path(self.root_dir) if self.flush_every_events <= 0: raise ValueError("flush_every_events must be >= 1") if self.flush_every_seconds <= 0: raise ValueError("flush_every_seconds must be > 0") if self.rollover_max_bytes <= 0: raise ValueError("rollover_max_bytes must be >= 1") if self.rollover_max_events <= 0: raise ValueError("rollover_max_events must be >= 1") if self.retention_max_files <= 0: raise ValueError("retention_max_files must be >= 1") if self.retention_max_total_bytes <= 0: raise ValueError("retention_max_total_bytes must be >= 1") if self.retention_max_total_bytes < self.rollover_max_bytes: raise ValueError("retention_max_total_bytes must be >= rollover_max_bytes") if self.rollup_window_seconds <= 0: raise ValueError("rollup_window_seconds must be >= 1") self._validate_bounded_mode() def _validate_bounded_mode(self) -> None: if self.max_buffer_bytes is not None and self.max_buffer_bytes <= 0: raise ValueError("max_buffer_bytes must be >= 1") if self.failure_backoff_seconds <= 0: raise ValueError("failure_backoff_seconds must be > 0") if self.failure_backoff_max_seconds < self.failure_backoff_seconds: raise ValueError( "failure_backoff_max_seconds must be >= failure_backoff_seconds" )
[docs] @dataclass class TelemetrySinkSegment: """One JSONL segment tracked by the append-only sink manifest.""" filename: str event_count: int size_bytes: int closed: bool session_id: str | None = None
[docs] @dataclass class TelemetrySinkManifest: """Parsed append-only sink manifest with session ledger and segments.""" schema_version: int format: str sessions: list[SessionSummary] = field(default_factory=list) segments: list[TelemetrySinkSegment] = field(default_factory=list)
[docs] class AppendOnlyTelemetrySink: """Write telemetry records to newline-delimited JSON segments.""" def __init__(self, config: TelemetrySinkConfig) -> None: self.config = config self.root_dir = config.root_dir self.root_dir.mkdir(parents=True, exist_ok=True) self._manifest_path = self.root_dir / MANIFEST_FILENAME self._segments: list[TelemetrySinkSegment] = [] self._sessions: dict[str, SessionSummary] = {} self._active_session_id: str | None = None self._next_segment_index = 1 # One buffer: what it holds is what buffered_bytes counts, and a # flush writes it as it is, with no joined copy. self._buffer = bytearray() self._buffered_event_count = 0 self._buffered_bytes = 0 self._fd: int | None = None # A segment a failed write left longer than it was, and the size to # cut it back to, when the cut-back itself failed. self._pending_cut_back: tuple[str, int] | None = None self._lock = threading.Lock() self._flush_stop_event = threading.Event() self._flush_thread: threading.Thread | None = None self._last_flush_monotonic = time.monotonic() self._closed = False self._rollover_count = 0 self._pruned_segment_count = 0 self._pruned_bytes = 0 self._dropped_records = 0 self._dropped_bytes = 0 self._flush_failures = 0 self._consecutive_flush_failures = 0 self._flush_retry_at = 0.0 self._last_flush_error: str | None = None self._load_existing_state()
[docs] def start_session(self, summary: SessionSummary | None = None) -> SessionSummary: """Register the active session for subsequent records.""" with self._lock: if self._active_session_id is not None: active = self._sessions.get(self._active_session_id) if active is not None: return active resolved = summary or create_session_summary( source="stormlog.telemetry_sink", status=SESSION_STATUS_RUNNING, ) resolved = update_session_summary( resolved, status=SESSION_STATUS_RUNNING, ended_at_ns=None, ) self._sessions[resolved.session_id] = resolved self._active_session_id = resolved.session_id self._closed = False self._write_manifest_checked_locked() return resolved
[docs] def current_session(self) -> SessionSummary | None: """Return the active session summary, if any.""" with self._lock: if self._active_session_id is None: return None return self._sessions.get(self._active_session_id)
[docs] def append(self, record: Mapping[str, Any]) -> None: line = (json.dumps(dict(record), sort_keys=True) + "\n").encode("utf-8") with self._lock: self._ensure_flush_thread_locked() self._closed = False self._ensure_active_session_locked(record) limit = self.config.max_buffer_bytes if limit is not None and self._buffered_bytes + len(line) > limit: # Make room by flushing now, unless a failed flush is backing # off; drop only what still does not fit. if time.monotonic() >= self._flush_retry_at: self._flush_locked(force=True) if self._buffered_bytes + len(line) > limit: self._dropped_records += 1 self._dropped_bytes += len(line) return self._buffer += line self._buffered_event_count += 1 self._buffered_bytes += len(line) self._flush_locked(force=False)
[docs] def flush(self, force: bool = False) -> None: with self._lock: self._flush_locked(force=force)
[docs] def close(self, session_status: str = SESSION_STATUS_COMPLETED) -> None: rollup_inputs: tuple[TelemetrySinkManifest, RollupCoverage] | None = None try: with self._lock: self._flush_locked(force=True) self._drop_unflushed_locked() self._close_fd_locked() current = self._current_segment() if current is not None and not current.closed: current.closed = True if self._active_session_id is not None: active = self._sessions.get(self._active_session_id) if active is not None: self._sessions[self._active_session_id] = ( update_session_summary( active, status=session_status, ended_at_ns=now_ns(), ) ) self._active_session_id = None self._write_manifest_checked_locked() rollup_inputs = self._rollup_inputs_locked() self._closed = True if rollup_inputs is not None: self._write_rollups(*rollup_inputs) finally: self._stop_flush_thread()
[docs] def get_diagnostics(self) -> dict[str, int]: """Return runtime retention and rollover diagnostics.""" with self._lock: return self._diagnostics_locked()
[docs] def failure_diagnostics(self) -> dict[str, int | str | None]: """Buffer and flush-failure counters; the bounded mode's evidence.""" with self._lock: return { "buffered_records": self._buffered_event_count, "buffered_bytes": self._buffered_bytes, "dropped_records": self._dropped_records, "dropped_bytes": self._dropped_bytes, "flush_failures": self._flush_failures, "consecutive_flush_failures": self._consecutive_flush_failures, "last_flush_error": self._last_flush_error, }
@property def _bounded(self) -> bool: return self.config.max_buffer_bytes is not None def _ensure_active_session_locked(self, record: Mapping[str, Any]) -> None: record_session_id = record.get("session_id") if record_session_id is not None and not isinstance(record_session_id, str): raise ValueError("telemetry sink record session_id must be a string") if self._active_session_id is None: resolved = create_session_summary( source="stormlog.telemetry_sink", status=SESSION_STATUS_RUNNING, session_id=( record_session_id if isinstance(record_session_id, str) else None ), ) self._sessions[resolved.session_id] = resolved self._active_session_id = resolved.session_id self._write_manifest_checked_locked() return if ( isinstance(record_session_id, str) and record_session_id != self._active_session_id ): raise ValueError( "telemetry sink record session_id does not match the active session" ) def _flush_locked(self, force: bool) -> None: if not self._buffer: return now = time.monotonic() if not force: if self._buffered_event_count < self.config.flush_every_events and ( now - self._last_flush_monotonic < self.config.flush_every_seconds ): return if self._bounded and now < self._flush_retry_at: return current = self._ensure_current_segment_locked() try: with memoryview(self._buffer) as payload: # released however it ends self._write_payload_locked(current, payload) except OSError as exc: self._record_flush_failure_locked(exc, now) if not self._bounded: raise return current.event_count += self._buffered_event_count current.size_bytes += len(self._buffer) self._buffer = bytearray() self._buffered_event_count = 0 self._buffered_bytes = 0 self._last_flush_monotonic = now self._consecutive_flush_failures = 0 self._rollover_locked(current) self._prune_retention_locked() # The records are durable; a failed manifest is rewritten next flush. self._write_manifest_checked_locked() def _write_manifest_checked_locked(self) -> None: """Write the manifest; in bounded mode a failure is counted, not raised. Every manifest write after construction goes through here, so a bounded sink never raises from append, flush, start_session or close. """ try: self._write_manifest_locked() except OSError as exc: self._record_flush_failure_locked(exc, time.monotonic()) if not self._bounded: raise def _write_payload_locked( self, current: TelemetrySinkSegment, payload: memoryview ) -> None: """Append every byte and fsync, or cut the segment back as it was. A failed write would otherwise leave a partial line that the next successful write extends into a corrupt record, or whole lines it writes again. The cut-back goes to the file's own size before the write, never to a remembered one. When it works, the descriptor is kept, so a failing disk costs no reopen and no read of the segment. When the cut-back itself fails, the descriptor is closed, and the segment is cut back to that size before it is written again. """ fd = self._ensure_fd_locked(current) before = os.fstat(fd).st_size try: _write_all(fd, payload) os.fsync(fd) except BaseException: # a KeyboardInterrupt mid-write too try: os.ftruncate(fd, before) except OSError: self._pending_cut_back = (current.filename, before) self._close_fd_locked() raise def _record_flush_failure_locked(self, exc: OSError, now: float) -> None: self._flush_failures += 1 self._consecutive_flush_failures += 1 self._last_flush_error = f"{type(exc).__name__}: {exc}" delay = collector_retry_delay_seconds( self._consecutive_flush_failures, initial_delay_s=self.config.failure_backoff_seconds, factor=2.0, max_delay_s=self.config.failure_backoff_max_seconds, ) self._flush_retry_at = now + delay if self._consecutive_flush_failures == 1: _LOGGER.warning("telemetry sink flush failed: %s", self._last_flush_error) def _drop_unflushed_locked(self) -> None: """Count what a bounded sink could not write before closing.""" if not self._bounded or not self._buffer: return self._dropped_records += self._buffered_event_count self._dropped_bytes += self._buffered_bytes self._buffer = bytearray() self._buffered_event_count = 0 self._buffered_bytes = 0 def _close_fd_locked(self) -> None: if self._fd is not None: try: os.close(self._fd) except OSError: pass self._fd = None def _rollover_locked(self, current: TelemetrySinkSegment) -> None: if ( current.event_count < self.config.rollover_max_events and current.size_bytes < self.config.rollover_max_bytes ): return current.closed = True self._rollover_count += 1 self._close_fd_locked() def _prune_retention_locked(self) -> None: while True: total_bytes = sum(segment.size_bytes for segment in self._segments) over_file_limit = len(self._segments) > self.config.retention_max_files over_size_limit = total_bytes > self.config.retention_max_total_bytes if not over_file_limit and not over_size_limit: return removable = next( (segment for segment in self._segments if segment.closed), None, ) if removable is None: return path = self.root_dir / removable.filename try: path.unlink(missing_ok=True) except OSError as exc: # Kept, and retried at the next flush. self._record_flush_failure_locked(exc, time.monotonic()) if not self._bounded: raise return self._segments.remove(removable) self._pruned_segment_count += 1 self._pruned_bytes += removable.size_bytes def _current_segment(self) -> TelemetrySinkSegment | None: if not self._segments: return None current = self._segments[-1] if current.closed: return None if ( self._active_session_id is not None and current.session_id != self._active_session_id ): return None return current def _ensure_flush_thread_locked(self) -> None: if self._flush_thread is not None and self._flush_thread.is_alive(): return self._flush_stop_event = threading.Event() self._flush_thread = threading.Thread( target=self._run_flush_loop, name="stormlog-telemetry-sink-flush", daemon=True, ) self._flush_thread.start() def _stop_flush_thread(self) -> None: thread = self._flush_thread if thread is None: return self._flush_stop_event.set() thread.join(timeout=self.config.flush_every_seconds + 1.0) self._flush_thread = None def _run_flush_loop(self) -> None: while not self._flush_stop_event.wait(timeout=self.config.flush_every_seconds): with self._lock: self._flush_locked(force=False) def _ensure_current_segment_locked(self) -> TelemetrySinkSegment: current = self._current_segment() if current is not None: return current segment = TelemetrySinkSegment( filename=f"{SEGMENT_PREFIX}{self._next_segment_index:06d}{SEGMENT_SUFFIX}", event_count=0, size_bytes=0, closed=False, session_id=self._active_session_id, ) self._next_segment_index += 1 self._segments.append(segment) return segment def _ensure_fd_locked(self, current: TelemetrySinkSegment) -> int: if self._fd is None: if not self._finish_cut_back_locked(current): segment_path = self.root_dir / current.filename self._recover_segment_tail_locked(segment_path, current) self._fd = os.open( self.root_dir / current.filename, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o666, ) return self._fd def _finish_cut_back_locked(self, current: TelemetrySinkSegment) -> bool: """Cut a segment back where a failed write's own cut-back could not; True when that was the current segment, whose counts then still hold, so it need not be read.""" if self._pending_cut_back is None: return False filename, size = self._pending_cut_back path = self.root_dir / filename try: cut = path.stat().st_size >= size # never extended with zeros if cut: os.truncate(path, size) except FileNotFoundError: cut = False # Any other error is a failed flush, and the cut-back stays pending. self._pending_cut_back = None return cut and filename == current.filename def _load_existing_state(self) -> None: discovered = _discover_segment_paths(self.root_dir) manifest = read_telemetry_sink_manifest(self.root_dir) if manifest is not None: self._segments = self._merge_segment_state(manifest.segments, discovered) self._sessions = { summary.session_id: summary for summary in manifest.sessions } manifest_needs_rewrite = self._interrupt_running_sessions() self._next_segment_index = self._compute_next_segment_index() if manifest.schema_version != SINK_SCHEMA_VERSION: manifest_needs_rewrite = True if manifest_needs_rewrite: self._write_manifest_locked() return if not discovered: return for path in discovered: self._segments.append( TelemetrySinkSegment( filename=path.name, event_count=self._count_records(path), size_bytes=path.stat().st_size, closed=True, session_id=None, ) ) self._next_segment_index = self._compute_next_segment_index() self._write_manifest_locked() def _interrupt_running_sessions(self) -> bool: manifest_needs_rewrite = False for session_id, summary in list(self._sessions.items()): if summary.status == SESSION_STATUS_RUNNING: self._sessions[session_id] = update_session_summary( summary, status=SESSION_STATUS_INTERRUPTED, ended_at_ns=now_ns(), ) manifest_needs_rewrite = True for segment in self._segments: if segment.session_id in self._sessions and ( self._sessions[segment.session_id].status == SESSION_STATUS_INTERRUPTED ): segment.closed = True return manifest_needs_rewrite def _compute_next_segment_index(self) -> int: max_index = 0 for segment in self._segments: stem = Path(segment.filename).stem if not stem.startswith(SEGMENT_PREFIX): continue suffix = stem[len(SEGMENT_PREFIX) :] if suffix.isdigit(): max_index = max(max_index, int(suffix)) return max_index + 1 def _write_manifest_locked(self) -> None: payload = { "schema_version": SINK_SCHEMA_VERSION, "format": "stormlog.append_only_telemetry_sink", "sessions": [ session_summary_to_dict(summary) for summary in sorted( self._sessions.values(), key=lambda session: (session.started_at_ns, session.session_id), ) ], "segments": [asdict(segment) for segment in self._segments], } temp_path = self._manifest_path.with_suffix(".tmp") temp_path.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8") temp_path.replace(self._manifest_path) def _diagnostics_locked(self) -> dict[str, int]: retained_bytes = sum(segment.size_bytes for segment in self._segments) return { "rollover_count": self._rollover_count, "pruned_segment_count": self._pruned_segment_count, "pruned_bytes": self._pruned_bytes, "final_retained_files": len(self._segments), "final_retained_bytes": retained_bytes, } def _rollup_inputs_locked( self, ) -> tuple[TelemetrySinkManifest, RollupCoverage] | None: if not self.config.write_rollups: return None from .telemetry_rollups import rollup_coverage_from_manifest manifest = TelemetrySinkManifest( schema_version=SINK_SCHEMA_VERSION, format="stormlog.append_only_telemetry_sink", sessions=list(self._sessions.values()), segments=list(self._segments), ) coverage = rollup_coverage_from_manifest( manifest, pruned_segment_count=self._pruned_segment_count, pruned_bytes=self._pruned_bytes, ) return manifest, coverage def _write_rollups( self, manifest: TelemetrySinkManifest, coverage: RollupCoverage, ) -> None: try: from .telemetry import load_telemetry_sessions from .telemetry_rollups import ( build_telemetry_rollups, write_telemetry_rollups, ) sessions = load_telemetry_sessions(self.root_dir) rollups = build_telemetry_rollups( sessions, manifest, window_duration_ns=(self.config.rollup_window_seconds * 1_000_000_000), coverage=coverage, ) write_telemetry_rollups(self.root_dir, rollups) except Exception as exc: _LOGGER.warning("telemetry rollup write failed: %s", exc, exc_info=True) @staticmethod def _count_records(path: Path) -> int: with path.open("r", encoding="utf-8") as handle: return sum(1 for line in handle if line.strip()) def _recover_segment_tail_locked( self, segment_path: Path, current: TelemetrySinkSegment, ) -> None: if not segment_path.exists(): current.event_count = 0 current.size_bytes = 0 return size = _whole_lines_size(segment_path) if size < segment_path.stat().st_size: # In place: rewriting the file would risk the whole lines too. os.truncate(segment_path, size) current.size_bytes = size current.event_count = self._count_records(segment_path) def _merge_segment_state( self, manifest_segments: list[TelemetrySinkSegment], discovered_segments: list[Path], ) -> list[TelemetrySinkSegment]: manifest_by_name = { segment.filename: TelemetrySinkSegment( filename=segment.filename, event_count=segment.event_count, size_bytes=segment.size_bytes, closed=segment.closed, session_id=segment.session_id, ) for segment in manifest_segments } discovered_by_name = {path.name: path for path in discovered_segments} merged_names = sorted(set(manifest_by_name) | set(discovered_by_name)) merged: list[TelemetrySinkSegment] = [] for name in merged_names: manifest_segment = manifest_by_name.get(name) if manifest_segment is not None: merged.append(manifest_segment) continue path = discovered_by_name[name] merged.append( TelemetrySinkSegment( filename=name, event_count=self._count_records(path), size_bytes=path.stat().st_size, closed=True, session_id=None, ) ) return merged
def _whole_lines_size(path: Path) -> int: """How long a file is up to its last newline, read backwards a chunk at a time, so a long segment is never held whole.""" with path.open("rb") as handle: position = handle.seek(0, os.SEEK_END) while position > 0: start = max(0, position - _TAIL_CHUNK_BYTES) handle.seek(start) newline = handle.read(position - start).rfind(b"\n") if newline >= 0: return start + newline + 1 position = start return 0 def _write_all(fd: int, payload: bytes | memoryview) -> None: """Write every byte; a short write continues, an error raises. One view, released however this ends, with each slice only an argument to the write: a view left in this frame would live on in the traceback of the error raised here, and while a caller kept that error, the sink's bytearray buffer could not be resized and the next append would raise BufferError. """ with memoryview(payload) as view: offset = 0 while offset < len(view): written = os.write(fd, view[offset:]) if written <= 0: raise OSError("write made no progress") offset += written
[docs] def resolve_telemetry_sink_segment_paths(path: str | Path) -> list[Path]: """Resolve append-only sink inputs to ordered JSONL segment paths.""" resolved_path = Path(path) if resolved_path.is_file(): if resolved_path.suffix == SEGMENT_SUFFIX: return [resolved_path] if resolved_path.name == MANIFEST_FILENAME: manifest = read_telemetry_sink_manifest(resolved_path) manifest_segments = _manifest_segment_paths(resolved_path.parent, manifest) return _merge_segment_paths( manifest_segments, _discover_segment_paths(resolved_path.parent), ) return [] if not resolved_path.is_dir(): return [] manifest = read_telemetry_sink_manifest(resolved_path) if manifest is not None: manifest_segments = _manifest_segment_paths(resolved_path, manifest) return _merge_segment_paths( manifest_segments, _discover_segment_paths(resolved_path), ) return _discover_segment_paths(resolved_path)
[docs] def resolve_telemetry_sink_manifest_path(path: str | Path) -> Path | None: """Resolve a file or directory path to a sink manifest path, if present.""" resolved = Path(path) if resolved.is_file(): if resolved.name == MANIFEST_FILENAME: return resolved if resolved.suffix == SEGMENT_SUFFIX: candidate = resolved.parent / MANIFEST_FILENAME return candidate if candidate.exists() else None return None if resolved.is_dir(): candidate = resolved / MANIFEST_FILENAME return candidate if candidate.exists() else None return None
def _coerce_manifest_int( value: object, *, default: int, minimum: int | None = None, ) -> int: try: coerced = int(cast(Any, value)) except (TypeError, ValueError): return default if minimum is not None and coerced < minimum: return default return int(coerced) def _coerce_manifest_entries(value: object) -> list[object]: return value if isinstance(value, list) else []
[docs] def read_telemetry_sink_manifest(path: str | Path) -> TelemetrySinkManifest | None: """Read a sink manifest from a sink directory, manifest file, or segment file.""" manifest_path = resolve_telemetry_sink_manifest_path(path) if manifest_path is None or not manifest_path.exists(): return None try: payload = json.loads(manifest_path.read_text(encoding="utf-8")) except Exception: return None if not isinstance(payload, Mapping): return None schema_version = _coerce_manifest_int( payload.get("schema_version", 1), default=1, minimum=1, ) fmt_value = payload.get("format", "stormlog.append_only_telemetry_sink") fmt = ( fmt_value if isinstance(fmt_value, str) and fmt_value else "stormlog.append_only_telemetry_sink" ) sessions = _manifest_sessions(payload, schema_version) segments = _manifest_segments(payload) return TelemetrySinkManifest( schema_version=schema_version, format=fmt, sessions=sessions, segments=segments, )
def _discover_segment_paths(root_dir: Path) -> list[Path]: return sorted(root_dir.glob(f"{SEGMENT_PREFIX}*{SEGMENT_SUFFIX}")) def _merge_segment_paths( manifest_segments: list[Path], discovered_segments: list[Path], ) -> list[Path]: merged_by_name = {path.name: path for path in discovered_segments} for path in manifest_segments: merged_by_name[path.name] = path return [merged_by_name[name] for name in sorted(merged_by_name)] def _manifest_sessions( payload: Mapping[str, Any], schema_version: int ) -> list[SessionSummary]: sessions: list[SessionSummary] = [] if schema_version >= 2: for raw_session in _coerce_manifest_entries(payload.get("sessions", [])): if not isinstance(raw_session, Mapping): continue try: sessions.append(session_summary_from_dict(raw_session)) except Exception: continue return sessions def _manifest_segments(payload: Mapping[str, Any]) -> list[TelemetrySinkSegment]: segments: list[TelemetrySinkSegment] = [] for raw_segment in _coerce_manifest_entries(payload.get("segments", [])): if not isinstance(raw_segment, Mapping): continue filename = raw_segment.get("filename") if not isinstance(filename, str): continue try: segments.append( TelemetrySinkSegment( filename=filename, event_count=_coerce_manifest_int( raw_segment.get("event_count", 0), default=0, minimum=0, ), size_bytes=_coerce_manifest_int( raw_segment.get("size_bytes", 0), default=0, minimum=0, ), closed=bool(raw_segment.get("closed", False)), session_id=( raw_segment.get("session_id") if isinstance(raw_segment.get("session_id"), str) else None ), ) ) except Exception: continue return segments def _manifest_segment_paths( root: Path, manifest: TelemetrySinkManifest | None ) -> list[Path]: return [ root / segment.filename for segment in (manifest.segments if manifest is not None else []) if (root / segment.filename).exists() ] __all__ = [ "AppendOnlyTelemetrySink", "MANIFEST_FILENAME", "SEGMENT_PREFIX", "SEGMENT_SUFFIX", "SINK_SCHEMA_VERSION", "TelemetrySinkConfig", "TelemetrySinkManifest", "TelemetrySinkSegment", "read_telemetry_sink_manifest", "resolve_telemetry_sink_manifest_path", "resolve_telemetry_sink_segment_paths", ]