From 9185e7aa1230288ef9648c72bdea81f709eef8a8 Mon Sep 17 00:00:00 2001 From: Baris Ozbas Date: Sat, 19 Sep 2026 13:00:58 +0200 Subject: [PATCH] fix(sensor): acknowledge OTLP delivery independently of local capture Summary: Drain bounded log batches and reconcile submitted/exported counts so the SDK queue cannot silently drop large captures. Validate protobuf acknowledgements, including HTTP-success partial rejection, before marking delivery successful. Track successful session snapshots in atomic hash-only destination checkpoints, independently of local JSON files. Retry unacknowledged sessions on a later run; continue sending sensor health even when sessions are already acknowledged. Document at-least-once semantics, credential configuration, and retry limitations. This builds on #131 (which builds on #130). Review and merge in that order; the main-only CI matrix will run after retargeting this PR to main. Test Plan: Synthetic tests exercise queue overflow protection, dropped-record detection, partial collector acknowledgements, retries, checkpoint failures, authentication scope, and health-only runs. No live sessions or collectors are used. Revert Plan: Revert this commit to restore the previous exporter. Hash-only checkpoint files remain harmless and are ignored by the previous version. --- Sensor/README.md | 40 ++- Sensor/adr_sensor/cli.py | 26 +- Sensor/adr_sensor/diagnostics.py | 4 +- .../exporters/delivery_checkpoint.py | 120 +++++++ Sensor/adr_sensor/exporters/opentelemetry.py | 159 +++++++-- Sensor/tests/test_delivery_checkpoint.py | 332 ++++++++++++++++++ Sensor/tests/test_opentelemetry_exporter.py | 174 +++++++++ 7 files changed, 821 insertions(+), 34 deletions(-) create mode 100644 Sensor/adr_sensor/exporters/delivery_checkpoint.py create mode 100644 Sensor/tests/test_delivery_checkpoint.py diff --git a/Sensor/README.md b/Sensor/README.md index d156ca1..44d71bf 100644 --- a/Sensor/README.md +++ b/Sensor/README.md @@ -354,11 +354,43 @@ no redaction or field projection, so prompts, responses, tool arguments, tool results, usernames, hostnames, and local paths can be transmitted. Any normalization already performed by a source parser still applies. -System-configuration records are sent as `adr.system.configuration` logs. Runs -are not checkpointed specifically for OTLP: repeated runs can resend the same -records, and consumers can use `adr.event.uuid` to deduplicate them. +System-configuration records are sent as `adr.system.configuration` logs on each +run. Sensor health logs are also sent on every run, even when all session snapshots +are already acknowledged. With `--save-sessions`, successful session delivery is tracked independently +of local session files. A failed export is retried on the next run, even when the +local JSON already exists. A session is skipped only when its complete normalized +payload was successfully exported to the same destination configuration. Changes +to tool results, destination settings, or configured authentication headers cause +a resend. The checkpoint also accounts for effective OTLP environment headers and +mTLS client certificate/key paths. It does not read credential files: after +changing certificate or key contents in place, remove the destination's checkpoint +to resend sessions. Dynamic HTTP credential-provider plugins +(`OTEL_PYTHON_EXPORTER_OTLP_HTTP_CREDENTIAL_PROVIDER` and its `LOGS` variant) +are unsupported and cause an explicit error; use configured headers or mTLS. + +Delivery checkpoints are hidden `.adr-otel-delivery..json` files in the +session output directory. They contain only hashes, including a destination hash +that accounts for authentication headers; they do not store raw URLs, credentials, +session identifiers, or payloads. Missing, unreadable, or corrupt checkpoints cause +sessions to be retried. The checkpoint is replaced atomically only after flush and +shutdown succeed; a checkpoint write failure exits with an error. `--no-save` +disables checkpoint reads and writes. Without `--save-sessions`, every run exports +all collected sessions. + +The one-shot Sensor process drains bounded batches and reconciles submitted and +successfully exported counts before reporting success, so a full SDK queue cannot +silently drop records. HTTP success is also checked for an OTLP acknowledgement: +partial rejection or a malformed response fails delivery and leaves the affected +run unacknowledged. Resolve persistent collector rejection before rerunning: OTLP +does not identify individual rejected records, so retrying can resend accepted +records too. Export or checkpoint failures exit with a nonzero status. +Delivery is at least once: a collector may receive data before a timeout, process +interruption, or checkpoint write failure, so retries can duplicate records. +Checkpointing only covers sessions that are collected again on a later run; it is +not a persistent payload queue. Consumers can use `adr.event.uuid` and a full +payload digest to identify repeated snapshots, since a session UUID alone does not +necessarily change when tool results change. -The one-shot Sensor process flushes and shuts down the exporter before exiting. Use an OpenTelemetry Collector when vendor-specific routing, transformation, retry, or persistent queuing is needed. diff --git a/Sensor/adr_sensor/cli.py b/Sensor/adr_sensor/cli.py index ddf6dda..0ed99b5 100644 --- a/Sensor/adr_sensor/cli.py +++ b/Sensor/adr_sensor/cli.py @@ -26,6 +26,7 @@ from . import __version__ from .diagnostics import health_record, write_health_records from .exporters import OpenTelemetryConfigError, load_opentelemetry_config +from .exporters.delivery_checkpoint import DeliveryCheckpoint, DeliveryCheckpointError from .exporters.opentelemetry import OpenTelemetryExportError, OpenTelemetryLogExporter from .observer import AgentObserver @@ -150,7 +151,21 @@ def main(): stage = "parse" entries, system_config_data = observer.ingest_all(args.source) + # A local file is not an OTLP acknowledgement. Keep remote candidates + # independent of local incremental filtering, including after a failed run. + otel_entries = entries + delivery_checkpoint = None + if otel_config is not None and args.save_sessions and not args.no_save: + stage = "export" + checkpoint_dir = args.output_dir if args.output_dir is not None else observer._get_default_session_dir() + delivery_checkpoint = DeliveryCheckpoint(checkpoint_dir, otel_config) + otel_entries = delivery_checkpoint.pending_entries(entries) + if delivery_checkpoint.load_failed: + observer.record_failure("export", "checkpoint_read_error") + print("OpenTelemetry delivery checkpoint unreadable or invalid; retrying sessions.", file=sys.stderr) + # Apply incremental filtering + stage = "save" if args.save_sessions and entries: print("\nSession-based incremental mode: Checking existing session files...") original_count = len(entries) @@ -186,10 +201,12 @@ def main(): stage = "export" otel_exporter = OpenTelemetryLogExporter(otel_config, service_version=get_version()) try: - exported_count = otel_exporter.export(entries, system_config_data) + exported_count = otel_exporter.export(otel_entries, system_config_data) otel_exporter.export_diagnostics(observer.get_diagnostic_records()) finally: otel_exporter.shutdown() + if delivery_checkpoint is not None: + delivery_checkpoint.commit() print(f"\nOpenTelemetry session/configuration logs sent: {exported_count}") success = observer.has_errors is not True @@ -207,6 +224,13 @@ def main(): print(f"OpenTelemetry export failed: {exc}", file=sys.stderr) raise SystemExit(1) + except DeliveryCheckpointError as exc: + success = False + if observer is not None: + observer.record_failure("export", "checkpoint_write_error") + print(f"OpenTelemetry checkpoint failed: {exc}", file=sys.stderr) + raise SystemExit(1) + except Exception: success = False if observer is not None: diff --git a/Sensor/adr_sensor/diagnostics.py b/Sensor/adr_sensor/diagnostics.py index 3b756dc..08e8ddb 100644 --- a/Sensor/adr_sensor/diagnostics.py +++ b/Sensor/adr_sensor/diagnostics.py @@ -15,7 +15,9 @@ {"sensor", "claude", "claude_desktop", "cursor", "cline", "codex", "copilot", "dsh", "gemini", "opencode", "warp"} ) DIAGNOSTIC_STAGES = frozenset({"parse", "save", "save_session", "export", "startup"}) -OPERATIONAL_REASONS = frozenset({"parser_error", "write_error", "export_error", "startup_error"}) +OPERATIONAL_REASONS = frozenset( + {"parser_error", "write_error", "export_error", "startup_error", "checkpoint_read_error", "checkpoint_write_error"} +) COUNT_FIELDS = frozenset({"events_returned", "events_emitted", "events_filtered", "attempted", "succeeded", "failed"}) MAX_LOG_BYTES = 1024 * 1024 LOG_BACKUP_COUNT = 2 diff --git a/Sensor/adr_sensor/exporters/delivery_checkpoint.py b/Sensor/adr_sensor/exporters/delivery_checkpoint.py new file mode 100644 index 0000000..cadeef8 --- /dev/null +++ b/Sensor/adr_sensor/exporters/delivery_checkpoint.py @@ -0,0 +1,120 @@ +"""Hash-only acknowledgements for incremental OTLP session delivery.""" + +import hashlib +import json +import os +import re +import tempfile +from dataclasses import asdict +from pathlib import Path +from typing import Dict, List + +from ..schemas.agent_event_schema import AgentEvent +from .config import OpenTelemetryConfig +from .opentelemetry import SCHEMA_VERSION, _validate_credential_provider + +_DIGEST = re.compile(r"[0-9a-f]{64}\Z") + + +class DeliveryCheckpointError(RuntimeError): + """Raised when a successful export cannot be durably checkpointed.""" + + +def _fingerprint(value: object) -> str: + serialized = json.dumps(value, sort_keys=True, ensure_ascii=False, separators=(",", ":")) + return hashlib.sha256(serialized.encode("utf-8")).hexdigest() + + +class DeliveryCheckpoint: + """Remember delivered session snapshots separately from their local JSON files. + + Select pending entries before export, then call commit only after flush and + shutdown succeed. Losing or corrupting this cache causes retries, never skips. + """ + + def __init__(self, output_dir: Path, config: OpenTelemetryConfig): + _validate_credential_provider() + destination = { + "config": asdict(config), + "schema_version": SCHEMA_VERSION, + "checkpoint_version": 1, + } + if not config.headers: + # The HTTP exporter falls back to these variables for empty headers. + destination["environment_headers"] = os.environ.get( + "OTEL_EXPORTER_OTLP_LOGS_HEADERS", os.environ.get("OTEL_EXPORTER_OTLP_HEADERS", "") + ) + for setting in ("CLIENT_CERTIFICATE", "CLIENT_KEY"): + destination[setting] = os.environ.get( + f"OTEL_EXPORTER_OTLP_LOGS_{setting}", os.environ.get(f"OTEL_EXPORTER_OTLP_{setting}", "") + ) + self.path = Path(output_dir) / f".adr-otel-delivery.{_fingerprint(destination)}.json" + self.load_failed = False + self._delivered = self._load() + self._pending: Dict[str, str] = {} + + def _load(self) -> Dict[str, str]: + try: + with self.path.open(encoding="utf-8") as checkpoint_file: + data = json.load(checkpoint_file) + if not isinstance(data, dict) or not all( + isinstance(key, str) and _DIGEST.fullmatch(key) and isinstance(value, str) and _DIGEST.fullmatch(value) + for key, value in data.items() + ): + raise ValueError("invalid checkpoint fingerprints") + return data + except FileNotFoundError: + return {} + except (OSError, ValueError): + self.load_failed = True + return {} + + def pending_entries(self, entries: List[AgentEvent]) -> List[AgentEvent]: + """Select snapshots not acknowledged for this destination, including results.""" + self._pending = {} + pending_entries = [] + for entry in entries: + identity = _fingerprint([entry.source, entry.session_id, entry.hostname, entry.username]) + payload = _fingerprint(entry.get_non_null_fields()) + if self._delivered.get(identity) != payload: + pending_entries.append(entry) + self._pending[identity] = payload + return pending_entries + + def commit(self) -> None: + """Atomically persist only the snapshots selected for a successful export.""" + if not self._pending: + return + temporary_path = None + try: + self.path.parent.mkdir(parents=True, exist_ok=True) + # Merge recent acknowledgements; concurrent writers may cause extra + # retries, but can never acknowledge a payload they have not exported. + delivered = self._load() + delivered.update(self._pending) + fd, temporary_name = tempfile.mkstemp(prefix=f"{self.path.name}.", suffix=".tmp", dir=self.path.parent) + temporary_path = Path(temporary_name) + with os.fdopen(fd, "w", encoding="utf-8") as checkpoint_file: + json.dump(delivered, checkpoint_file, sort_keys=True, separators=(",", ":")) + checkpoint_file.flush() + os.fsync(checkpoint_file.fileno()) + os.replace(temporary_path, self.path) + temporary_path = None + if os.name != "nt": + directory_fd = os.open(self.path.parent, os.O_RDONLY) + try: + os.fsync(directory_fd) + finally: + os.close(directory_fd) + self._delivered = delivered + self._pending = {} + except OSError as exc: + raise DeliveryCheckpointError( + "could not save the OpenTelemetry delivery checkpoint; a later run may resend delivered sessions" + ) from exc + finally: + if temporary_path is not None: + try: + temporary_path.unlink() + except OSError: + pass diff --git a/Sensor/adr_sensor/exporters/opentelemetry.py b/Sensor/adr_sensor/exporters/opentelemetry.py index 7037df9..a3618fe 100644 --- a/Sensor/adr_sensor/exporters/opentelemetry.py +++ b/Sensor/adr_sensor/exporters/opentelemetry.py @@ -1,5 +1,6 @@ """OTLP/HTTP logs exporter for normalized ADR Sensor records.""" +import os import socket import threading import time @@ -12,6 +13,7 @@ from .config import OpenTelemetryConfig SCHEMA_VERSION = "1" +_MAX_BATCH_SIZE = 512 class OpenTelemetryExportError(RuntimeError): @@ -29,6 +31,7 @@ def __init__( _log_record_exporter: Optional[Any] = None, _processor_factory: Optional[Any] = None, ): + _validate_credential_provider() components = _load_opentelemetry_components() ( logger_provider_cls, @@ -55,25 +58,37 @@ def __init__( certificate_file=config.certificate_file, headers=config.headers, timeout=config.timeout_seconds, + session=_create_otlp_session(), ) self._record_exporter = _ExportStatusTracker(record_exporter, export_result_cls.SUCCESS) - processor_factory = _processor_factory or processor_cls - processor = processor_factory(self._record_exporter) + if _processor_factory is None: + # Bound the queue explicitly instead of inheriting environment defaults. + # _emit() drains it before submitting another batch, providing backpressure. + processor = processor_cls( + self._record_exporter, + max_queue_size=_MAX_BATCH_SIZE, + max_export_batch_size=_MAX_BATCH_SIZE, + ) + else: + processor = _processor_factory(self._record_exporter) self._provider.add_log_record_processor(processor) self._logger = self._provider.get_logger("adr_sensor", service_version) self._info_severity = severity_number_cls.INFO self._warning_severity = severity_number_cls.WARN self._flush_timeout_millis = int(config.flush_timeout_seconds * 1000) self._closed = False + self._submitted_count = 0 + self._flushed_count = 0 + self._emit_lock = threading.Lock() def export( self, entries: List[AgentEvent], system_config_data: List[SystemConfiguration], ) -> int: - """Queue complete, unredacted Sensor records for OTLP export.""" + """Submit complete Sensor records; shutdown must succeed to confirm delivery.""" for entry in entries: self._emit( body=entry.get_non_null_fields(), @@ -108,9 +123,8 @@ def export_diagnostics(self, records: List[dict]) -> int: for record in records: body = sanitize_health_record(record) degraded = body["status"] in {"partial", "failed"} - self._logger.emit( - timestamp=_datetime_to_unix_nanos(datetime.fromisoformat(body["timestamp"])), - observed_timestamp=time.time_ns(), + self._emit( + timestamp=datetime.fromisoformat(body["timestamp"]), severity_number=self._warning_severity if degraded else self._info_severity, severity_text="WARN" if degraded else "INFO", body=body, @@ -127,31 +141,65 @@ def export_diagnostics(self, records: List[dict]) -> int: def shutdown(self) -> None: """Flush pending records and stop the provider's worker thread.""" - if self._closed: - return + with self._emit_lock: + if self._closed: + return + + self._closed = True + try: + self._flush() + finally: + try: + self._provider.shutdown() + except Exception as exc: + raise OpenTelemetryExportError("OpenTelemetry exporter shutdown failed") from exc + self._check_delivery() + + def _check_delivery(self) -> None: + if self._record_exporter.failed: + raise OpenTelemetryExportError("the OTLP endpoint did not accept one or more log batches") + exported_count = self._record_exporter.exported_count + if exported_count != self._submitted_count: + raise OpenTelemetryExportError( + f"OpenTelemetry delivery count mismatch: submitted {self._submitted_count}, " + f"successfully exported {exported_count}" + ) - self._closed = True - flushed = False + def _flush(self) -> None: try: flushed = self._provider.force_flush(timeout_millis=self._flush_timeout_millis) - finally: - self._provider.shutdown() - + except Exception as exc: + raise OpenTelemetryExportError("OpenTelemetry logs could not be flushed") from exc if not flushed: raise OpenTelemetryExportError(f"OpenTelemetry logs did not flush within {self._flush_timeout_millis} ms") - if self._record_exporter.failed: - raise OpenTelemetryExportError("the OTLP endpoint did not accept one or more log batches") + self._check_delivery() + self._flushed_count = self._submitted_count - def _emit(self, body: dict, timestamp: datetime, event_name: str, attributes: dict) -> None: - self._logger.emit( - timestamp=_datetime_to_unix_nanos(timestamp), - observed_timestamp=time.time_ns(), - severity_number=self._info_severity, - severity_text="INFO", - body=body, - attributes=attributes, - event_name=event_name, - ) + def _emit( + self, + body: dict, + timestamp: datetime, + event_name: str, + attributes: dict, + *, + severity_number: Optional[Any] = None, + severity_text: str = "INFO", + ) -> None: + with self._emit_lock: + if self._closed: + raise OpenTelemetryExportError("OpenTelemetry exporter is already shut down") + self._logger.emit( + timestamp=_datetime_to_unix_nanos(timestamp), + observed_timestamp=time.time_ns(), + severity_number=self._info_severity if severity_number is None else severity_number, + severity_text=severity_text, + body=body, + attributes=attributes, + event_name=event_name, + ) + self._submitted_count += 1 + if self._submitted_count - self._flushed_count >= _MAX_BATCH_SIZE: + self._flush() def _datetime_to_unix_nanos(value: datetime) -> int: @@ -163,13 +211,61 @@ def _datetime_to_unix_nanos(value: datetime) -> int: return ((delta.days * 86400 + delta.seconds) * 1_000_000_000) + (delta.microseconds * 1000) +def _validate_credential_provider() -> None: + """Reject opaque authentication that cannot be identified in checkpoints.""" + if os.environ.get("OTEL_PYTHON_EXPORTER_OTLP_HTTP_CREDENTIAL_PROVIDER") or os.environ.get( + "OTEL_PYTHON_EXPORTER_OTLP_HTTP_LOGS_CREDENTIAL_PROVIDER" + ): + raise OpenTelemetryExportError( + "OpenTelemetry HTTP credential-provider plugins are unsupported; use explicit headers or mTLS settings" + ) + + +def _create_otlp_session() -> Any: + """Validate the collector acknowledgement before the SDK counts HTTP success. + + The pinned SDK treats any successful HTTP response as full batch success, + including OTLP partial rejection. Its public session hook lets us check the + protobuf response without overriding the SDK's private transport methods. + """ + import requests + from google.protobuf.message import DecodeError + from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ExportLogsServiceResponse + + def validate_response(response: Any, *args: Any, **kwargs: Any) -> Any: + if response.is_redirect: + return response + if not 200 <= response.status_code < 300: + # Requests follows real redirects; other 3xx responses are not OTLP + # acknowledgements even though the SDK's Response.ok accepts them. + if response.ok: + raise OpenTelemetryExportError("the OTLP endpoint returned an invalid acknowledgement") + return response + if not response.content: + return response + acknowledgement = ExportLogsServiceResponse() + try: + acknowledgement.ParseFromString(response.content) + except DecodeError: + raise OpenTelemetryExportError("the OTLP endpoint returned an invalid acknowledgement") from None + if acknowledgement.partial_success.rejected_log_records != 0: + # Do not log the collector's error_message; it may contain payloads. + raise OpenTelemetryExportError("the OTLP endpoint rejected one or more log records") + return response + + session = requests.Session() + session.hooks["response"].append(validate_response) + return session + + class _ExportStatusTracker: - """Track exporter failures that OpenTelemetry's batch processor otherwise ignores.""" + """Count successful records and track failures the batch processor ignores.""" def __init__(self, exporter: Any, success_result: Any): self._exporter = exporter self._success_result = success_result self._failed = False + self._exported_count = 0 self._lock = threading.Lock() @property @@ -177,6 +273,11 @@ def failed(self) -> bool: with self._lock: return self._failed + @property + def exported_count(self) -> int: + with self._lock: + return self._exported_count + def export(self, batch: Any) -> Any: try: result = self._exporter.export(batch) @@ -185,8 +286,10 @@ def export(self, batch: Any) -> Any: self._failed = True raise - if result != self._success_result: - with self._lock: + with self._lock: + if result == self._success_result: + self._exported_count += len(batch) + else: self._failed = True return result diff --git a/Sensor/tests/test_delivery_checkpoint.py b/Sensor/tests/test_delivery_checkpoint.py new file mode 100644 index 0000000..6f2ccad --- /dev/null +++ b/Sensor/tests/test_delivery_checkpoint.py @@ -0,0 +1,332 @@ +"""Synthetic incremental-delivery tests; no collector or real sessions needed.""" + +import json +import re +from dataclasses import replace +from datetime import datetime, timezone +from unittest.mock import patch + +import pytest +from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ExportLogsServiceRequest, ExportLogsServiceResponse +from requests import Response + +from adr_sensor.cli import main +from adr_sensor.diagnostics import health_record +from adr_sensor.exporters.config import OpenTelemetryConfig +from adr_sensor.exporters.delivery_checkpoint import DeliveryCheckpoint, DeliveryCheckpointError +from adr_sensor.exporters.opentelemetry import OpenTelemetryExportError +from adr_sensor.schemas.agent_event_schema import AgentEvent, ChatMessage, ToolUsage + + +@pytest.fixture +def event(): + return AgentEvent( + timestamp=datetime(2026, 9, 9, tzinfo=timezone.utc), + source="codex", + session_id="synthetic-private-session", + hostname="synthetic-host", + username="synthetic-user", + chat_history=[ + ChatMessage( + role="assistant", + content="synthetic private text", + tools=[ToolUsage(tool_name="shell", tool_type="custom", result="first synthetic result")], + ) + ], + ) + + +@pytest.fixture +def config(): + return OpenTelemetryConfig( + endpoint="https://collector.example.invalid/v1/logs", + headers={"Authorization": "Bearer synthetic-secret"}, + ) + + +def test_checkpoint_skips_only_successfully_committed_payloads(tmp_path, config, event): + checkpoint = DeliveryCheckpoint(tmp_path, config) + assert checkpoint.pending_entries([event]) == [event] + assert not checkpoint.path.exists() + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [event] + + checkpoint.commit() + + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [] + + +@pytest.mark.parametrize("changed_field", ["tool_result", "tool_status", "model", "session_context"]) +def test_checkpoint_hashes_complete_normalized_payload(tmp_path, config, event, changed_field): + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + original_uuid = event.uuid + if changed_field == "tool_result": + event.chat_history[0].tools[0] = replace(event.chat_history[0].tools[0], result="updated synthetic result") + elif changed_field == "tool_status": + event.chat_history[0].tools[0] = replace(event.chat_history[0].tools[0], status="error") + elif changed_field == "model": + event = replace(event, model="changed-model") + else: + event = replace(event, session_context={"changed": True}) + + assert event.uuid == original_uuid + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [event] + + +@pytest.mark.parametrize( + "settings", + [ + {"endpoint": "https://other.example.invalid/v1/logs"}, + {"headers": {"Authorization": "Bearer other-synthetic-secret"}}, + {"service_name": "another-service"}, + ], +) +def test_destination_or_authentication_change_resends_sessions(tmp_path, config, event, settings): + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + changed = DeliveryCheckpoint(tmp_path, replace(config, **settings)) + + assert changed.path != checkpoint.path + assert changed.pending_entries([event]) == [event] + + +def test_environment_authentication_change_resends_sessions(tmp_path, event, monkeypatch): + config = OpenTelemetryConfig(endpoint="https://collector.example.invalid/v1/logs") + monkeypatch.setenv("OTEL_EXPORTER_OTLP_LOGS_HEADERS", "Authorization=first-synthetic-token") + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + monkeypatch.setenv("OTEL_EXPORTER_OTLP_LOGS_HEADERS", "Authorization=second-synthetic-token") + + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [event] + + +@pytest.mark.parametrize("prefix", ["OTEL_EXPORTER_OTLP", "OTEL_EXPORTER_OTLP_LOGS"]) +@pytest.mark.parametrize("setting", ["CLIENT_CERTIFICATE", "CLIENT_KEY"]) +def test_environment_mtls_path_change_resends_sessions(tmp_path, config, event, monkeypatch, prefix, setting): + variable = f"{prefix}_{setting}" + monkeypatch.setenv(variable, "/synthetic/first-credential.pem") + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + monkeypatch.setenv(variable, "/synthetic/second-credential.pem") + changed = DeliveryCheckpoint(tmp_path, config) + + assert changed.path != checkpoint.path + assert changed.pending_entries([event]) == [event] + + +def test_checkpoint_uses_effective_mtls_settings(tmp_path, config, monkeypatch): + monkeypatch.setenv("OTEL_EXPORTER_OTLP_CLIENT_CERTIFICATE", "/synthetic/generic-cert.pem") + monkeypatch.setenv("OTEL_EXPORTER_OTLP_LOGS_CLIENT_CERTIFICATE", "/synthetic/logs-cert.pem") + checkpoint = DeliveryCheckpoint(tmp_path, config) + monkeypatch.setenv("OTEL_EXPORTER_OTLP_CLIENT_CERTIFICATE", "/synthetic/unused-cert.pem") + + assert DeliveryCheckpoint(tmp_path, config).path == checkpoint.path + + +def test_opaque_credential_provider_cannot_skip_previously_delivered_sessions(tmp_path, config, event, monkeypatch): + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + monkeypatch.setenv("OTEL_PYTHON_EXPORTER_OTLP_HTTP_LOGS_CREDENTIAL_PROVIDER", "synthetic-private-provider") + + with pytest.raises(OpenTelemetryExportError, match="credential-provider plugins are unsupported"): + DeliveryCheckpoint(tmp_path, config) + + +def test_checkpoint_persists_only_hashes(tmp_path, config, event): + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + + assert re.fullmatch(r"\.adr-otel-delivery\.[0-9a-f]{64}\.json", checkpoint.path.name) + data = json.loads(checkpoint.path.read_text()) + assert len(data) == 1 + assert all(re.fullmatch(r"[0-9a-f]{64}", digest) for pair in data.items() for digest in pair) + + +@pytest.mark.parametrize("contents", ["broken JSON", "[]", '{"raw-session": "raw-content"}', "\udcff"]) +def test_corrupt_checkpoint_retries_and_can_be_replaced(tmp_path, config, event, contents): + path = DeliveryCheckpoint(tmp_path, config).path + path.write_bytes(contents.encode("utf-8", errors="surrogateescape")) + + checkpoint = DeliveryCheckpoint(tmp_path, config) + assert checkpoint.load_failed + assert checkpoint.pending_entries([event]) == [event] + checkpoint.commit() + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [] + + +def test_unreadable_checkpoint_retries(tmp_path, config, event): + with patch("pathlib.Path.open", side_effect=PermissionError("synthetic private path")): + checkpoint = DeliveryCheckpoint(tmp_path, config) + assert checkpoint.load_failed + assert checkpoint.pending_entries([event]) == [event] + + +def test_atomic_replace_failure_keeps_previous_checkpoint_and_retries(tmp_path, config, event): + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + previous = checkpoint.path.read_bytes() + event.chat_history[0].tools[0] = replace(event.chat_history[0].tools[0], result="changed result") + checkpoint.pending_entries([event]) + + with ( + patch("adr_sensor.exporters.delivery_checkpoint.os.replace", side_effect=OSError("synthetic private path")), + pytest.raises(DeliveryCheckpointError, match="could not save") as error, + ): + checkpoint.commit() + + assert "synthetic private path" not in str(error.value) + assert checkpoint.path.read_bytes() == previous + assert not list(tmp_path.glob("*.tmp")) + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [event] + + +def _prepare_cli(observer, event, tmp_path, monkeypatch, extra_args=()): + observer.ingest_all.return_value = ([event], []) + observer.get_diagnostic_records.return_value = [health_record("codex", "parse", counts={"events_emitted": 1})] + local_file = tmp_path / "synthetic-session.json" + observer.filter_entries_by_existing_files.side_effect = lambda entries, _: [] if local_file.exists() else entries + + def save_sessions(entries, output_dir): + local_file.write_text(json.dumps(entries[0].get_non_null_fields())) + return [local_file] + + observer.save_sessions_to_individual_files.side_effect = save_sessions + monkeypatch.setattr( + "sys.argv", + [ + "adr-sensor", + "--save-sessions", + "--output-dir", + str(tmp_path), + "--otel-config", + "synthetic.json", + *extra_args, + ], + ) + return local_file + + +@patch("adr_sensor.cli.OpenTelemetryLogExporter") +@patch("adr_sensor.cli.load_opentelemetry_config") +@patch("adr_sensor.cli.AgentObserver") +def test_cli_failed_delivery_retries_then_skips_unchanged( + observer_cls, load_config, exporter_cls, tmp_path, config, event, monkeypatch +): + load_config.return_value = config + observer = observer_cls.return_value + local_file = _prepare_cli(observer, event, tmp_path, monkeypatch) + exporter = exporter_cls.return_value + exporter.export.return_value = 1 + exporter.shutdown.side_effect = [OpenTelemetryExportError("synthetic delivery failure"), None, None] + + with pytest.raises(SystemExit) as error: + main() + assert error.value.code == 1 + assert local_file.exists() + assert not list(tmp_path.glob(".adr-otel-delivery.*.json")) + + main() + main() + + assert exporter_cls.call_count == 3 + assert exporter.export.call_count == 3 + assert [call.args for call in exporter.export.call_args_list] == [([event], []), ([event], []), ([], [])] + assert exporter.export_diagnostics.call_count == 3 + assert observer.save_sessions_to_individual_files.call_count == 1 + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [] + + +@patch("adr_sensor.cli.OpenTelemetryLogExporter") +@patch("adr_sensor.cli.load_opentelemetry_config") +@patch("adr_sensor.cli.AgentObserver") +def test_cli_checkpoint_write_failure_is_reported_and_retried( + observer_cls, load_config, exporter_cls, tmp_path, config, event, monkeypatch, capsys +): + load_config.return_value = config + _prepare_cli(observer_cls.return_value, event, tmp_path, monkeypatch) + exporter_cls.return_value.export.return_value = 1 + with patch("adr_sensor.exporters.delivery_checkpoint.os.replace", side_effect=OSError("synthetic private path")): + with pytest.raises(SystemExit) as error: + main() + + assert error.value.code == 1 + stderr = capsys.readouterr().err + assert "OpenTelemetry checkpoint failed" in stderr + assert "synthetic private path" not in stderr + assert not list(tmp_path.glob(".adr-otel-delivery.*.json")) + main() + assert exporter_cls.return_value.export.call_count == 2 + + +@patch("adr_sensor.cli.OpenTelemetryLogExporter") +@patch("adr_sensor.cli.load_opentelemetry_config") +@patch("adr_sensor.cli.AgentObserver") +def test_cli_no_save_does_not_read_or_write_delivery_state( + observer_cls, load_config, exporter_cls, tmp_path, config, event, monkeypatch +): + load_config.return_value = config + _prepare_cli(observer_cls.return_value, event, tmp_path, monkeypatch, extra_args=["--no-save"]) + checkpoint = DeliveryCheckpoint(tmp_path, config) + checkpoint.pending_entries([event]) + checkpoint.commit() + original = checkpoint.path.read_bytes() + exporter_cls.return_value.export.return_value = 1 + + main() + + exporter_cls.return_value.export.assert_called_once_with([event], []) + observer_cls.return_value.save_sessions_to_individual_files.assert_not_called() + assert checkpoint.path.read_bytes() == original + + +@patch("adr_sensor.cli.load_opentelemetry_config") +@patch("adr_sensor.cli.AgentObserver") +def test_cli_partial_rejection_does_not_checkpoint_and_next_run_retries( + observer_cls, load_config, tmp_path, config, event, monkeypatch, caplog +): + load_config.return_value = config + local_file = _prepare_cli(observer_cls.return_value, event, tmp_path, monkeypatch) + rejected = Response() + rejected.status_code = 200 + rejected._content = ExportLogsServiceResponse( + partial_success={"rejected_log_records": 1, "error_message": "synthetic private collector response"} + ).SerializeToString() + accepted = Response() + accepted.status_code = 200 + accepted._content = b"" + + with patch("requests.adapters.HTTPAdapter.send", side_effect=[rejected, accepted, accepted]) as transport: + with pytest.raises(SystemExit) as error: + main() + assert error.value.code == 1 + assert local_file.exists() + assert not list(tmp_path.glob(".adr-otel-delivery.*.json")) + + main() + main() + + assert transport.call_count == 3 + requests = [ExportLogsServiceRequest.FromString(call.args[0].body) for call in transport.call_args_list] + event_names = [ + [ + record.event_name + for resource in request.resource_logs + for scope in resource.scope_logs + for record in scope.log_records + ] + for request in requests + ] + assert event_names == [ + ["adr.agent.session", "adr.sensor.health"], + ["adr.agent.session", "adr.sensor.health"], + ["adr.sensor.health"], + ] + assert DeliveryCheckpoint(tmp_path, config).pending_entries([event]) == [] + assert "synthetic private collector response" not in caplog.text diff --git a/Sensor/tests/test_opentelemetry_exporter.py b/Sensor/tests/test_opentelemetry_exporter.py index 4c48759..c06ee41 100644 --- a/Sensor/tests/test_opentelemetry_exporter.py +++ b/Sensor/tests/test_opentelemetry_exporter.py @@ -1,14 +1,19 @@ """Tests for OTLP conversion of ADR Sensor records.""" +import time from datetime import datetime, timezone +from unittest.mock import patch import pytest +from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ExportLogsServiceResponse from opentelemetry.sdk._logs.export import ( InMemoryLogRecordExporter, LogRecordExportResult, SimpleLogRecordProcessor, ) +from requests import Response +from adr_sensor.diagnostics import health_record from adr_sensor.exporters.config import OpenTelemetryConfig from adr_sensor.exporters.opentelemetry import ( OpenTelemetryExportError, @@ -115,3 +120,172 @@ def test_datetime_to_unix_nanos_treats_naive_datetime_as_utc(): naive = aware.replace(tzinfo=None) assert _datetime_to_unix_nanos(naive) == _datetime_to_unix_nanos(aware) + + +def test_large_export_drains_bounded_batches_without_losing_records(monkeypatch): + # The old default queue lost records above 2048 when its worker was busy. + # Environment settings must not silently shrink our explicit queue bound. + monkeypatch.setenv("OTEL_BLRP_MAX_QUEUE_SIZE", "1") + monkeypatch.setenv("OTEL_BLRP_MAX_EXPORT_BATCH_SIZE", "1") + memory_exporter = InMemoryLogRecordExporter() + batch_sizes = [] + original_export = memory_exporter.export + + def record_batch(batch): + time.sleep(0.005) + batch_sizes.append(len(batch)) + return original_export(batch) + + memory_exporter.export = record_batch + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="http://localhost:4318/v1/logs"), + service_version="1.2.3", + _log_record_exporter=memory_exporter, + ) + entries = [_event()] * 4097 + health = [health_record("codex", "parse", counts={"events_emitted": 1})] * 1025 + + assert exporter.export(entries, []) == len(entries) + assert exporter.export_diagnostics(health) == len(health) + exporter.shutdown() + + assert len(memory_exporter.get_finished_logs()) == len(entries) + len(health) + assert sum(batch_sizes) == len(entries) + len(health) + assert max(batch_sizes) <= 512 + + +def test_shutdown_detects_silently_dropped_records(): + class DroppingProcessor(SimpleLogRecordProcessor): + def on_emit(self, log_record): + pass + + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="http://localhost:4318/v1/logs"), + service_version="1.2.3", + _log_record_exporter=InMemoryLogRecordExporter(), + _processor_factory=DroppingProcessor, + ) + exporter.export([_event()], []) + + with pytest.raises(OpenTelemetryExportError, match="submitted 1, successfully exported 0"): + exporter.shutdown() + + +def test_failed_intermediate_batch_stops_large_exports(): + class FailingExporter(InMemoryLogRecordExporter): + def export(self, batch): + return LogRecordExportResult.FAILURE + + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="http://localhost:4318/v1/logs"), + service_version="1.2.3", + _log_record_exporter=FailingExporter(), + ) + with pytest.raises(OpenTelemetryExportError, match="did not accept"): + exporter.export([_event()] * 4097, []) + with pytest.raises(OpenTelemetryExportError, match="did not accept"): + exporter.shutdown() + + +def test_shutdown_reports_flush_timeout_and_still_stops_provider(): + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="http://localhost:4318/v1/logs"), + service_version="1.2.3", + _log_record_exporter=InMemoryLogRecordExporter(), + ) + with ( + patch.object(exporter._provider, "force_flush", return_value=False), + patch.object(exporter._provider, "shutdown", wraps=exporter._provider.shutdown) as shutdown, + pytest.raises(OpenTelemetryExportError, match="did not flush"), + ): + exporter.shutdown() + shutdown.assert_called_once_with() + + +def test_export_after_shutdown_is_rejected(): + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="http://localhost:4318/v1/logs"), + service_version="1.2.3", + _log_record_exporter=InMemoryLogRecordExporter(), + ) + exporter.shutdown() + with pytest.raises(OpenTelemetryExportError, match="already shut down"): + exporter.export([_event()], []) + + +def _http_response(body, status_code=200): + response = Response() + response.status_code = status_code + response._content = body + response.headers["Content-Type"] = "application/x-protobuf" + return response + + +@pytest.mark.parametrize( + "body", + [ + b"", + ExportLogsServiceResponse(partial_success={"rejected_log_records": 0}).SerializeToString(), + ExportLogsServiceResponse( + partial_success={"rejected_log_records": 0, "error_message": "synthetic private warning"} + ).SerializeToString(), + b"\x10\x01", # Unknown protobuf fields remain forward compatible. + ], +) +def test_http_acknowledgement_accepts_full_success_and_zero_rejection_warnings(body, caplog): + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="https://collector.example.invalid/v1/logs"), service_version="1.2.3" + ) + with patch("requests.adapters.HTTPAdapter.send", return_value=_http_response(body)) as transport: + assert exporter.export([_event()], []) == 1 + exporter.shutdown() + + transport.assert_called_once() + assert exporter._record_exporter.exported_count == 1 + assert "synthetic private warning" not in caplog.text + + +@pytest.mark.parametrize( + "body,status_code", + [ + ( + ExportLogsServiceResponse( + partial_success={"rejected_log_records": 1, "error_message": "synthetic private rejection"} + ).SerializeToString(), + 200, + ), + (ExportLogsServiceResponse(partial_success={"rejected_log_records": -1}).SerializeToString(), 200), + (b"synthetic private invalid response", 200), + (b"", 304), + ], +) +def test_http_acknowledgement_rejects_partial_or_invalid_success_without_private_text(body, status_code, caplog): + exporter = OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="https://collector.example.invalid/v1/logs"), service_version="1.2.3" + ) + with patch("requests.adapters.HTTPAdapter.send", return_value=_http_response(body, status_code)) as transport: + exporter.export([_event()], []) + with pytest.raises(OpenTelemetryExportError, match="did not accept") as error: + exporter.shutdown() + + transport.assert_called_once() + assert exporter._record_exporter.exported_count == 0 + assert "synthetic private" not in caplog.text + assert "synthetic private" not in str(error.value) + + +@pytest.mark.parametrize( + "variable", + ["OTEL_PYTHON_EXPORTER_OTLP_HTTP_CREDENTIAL_PROVIDER", "OTEL_PYTHON_EXPORTER_OTLP_HTTP_LOGS_CREDENTIAL_PROVIDER"], +) +def test_opaque_credential_providers_fail_explicitly_without_loading_plugin(variable, monkeypatch): + monkeypatch.setenv(variable, "synthetic-private-provider") + with ( + patch("adr_sensor.exporters.opentelemetry._load_opentelemetry_components") as load_components, + pytest.raises(OpenTelemetryExportError, match="credential-provider plugins are unsupported") as error, + ): + OpenTelemetryLogExporter( + OpenTelemetryConfig(endpoint="https://collector.example.invalid/v1/logs"), service_version="1.2.3" + ) + load_components.assert_not_called() + assert "synthetic-private-provider" not in str(error.value)