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