"""Snapshot lifecycle: collect → sign → send → buffer-on-failure, plus replay. These coroutines operate on an ``AissurancePlugin`` instance (duck-typed to avoid an import cycle) and contain the per-cycle decision logic. The periodic timing loops live in plugin.py. """ import asyncio import time from typing import TYPE_CHECKING from compliance import payload as payload_mod from compliance.collector import collect_snapshot from compliance.log import log from compliance.transport import Outcome, SendResult if TYPE_CHECKING: # pragma: no cover from compliance.plugin import AissurancePlugin # Snapshot collection budget (spec §9.2). Exceeding it skips the cycle. _SNAPSHOT_TIMEOUT = 10.0 async def _send_once(plugin: "AissurancePlugin", body: dict) -> SendResult: result = await plugin.transport.send_payload(body) if result.outcome is Outcome.TOO_LARGE: # 413 — split by batch_size and send each chunk; worst chunk outcome wins. batches = payload_mod.split_into_batches(body, plugin.cfg.batch_size, plugin.cfg) if len(batches) > 1: worst = SendResult(Outcome.OK, 202) for chunk in batches: r = await plugin.transport.send_payload(chunk) if r.outcome is not Outcome.OK: worst = r return worst return result async def send_with_retry(plugin: "AissurancePlugin", body: dict) -> Outcome: """Send one payload, retrying transient failures with exponential backoff. Returns the final outcome. The caller decides whether to buffer. """ attempts = plugin.cfg.max_retry_attempts base = plugin.cfg.retry_backoff_base for attempt in range(max(attempts, 1)): result = await _send_once(plugin, body) if result.outcome is Outcome.OK: plugin.sent_total += 1 log("payload_sent", tenant_id=plugin.cfg.tenant_id, router_id=plugin.cfg.router_id, models_count=len(body.get("models", [])), response_code=result.status_code, latency_ms=result.latency_ms) return Outcome.OK if result.outcome in (Outcome.AUTH_FAILED, Outcome.RATE_LIMITED): # Non-retryable within this cycle. return result.outcome # RETRY (network/5xx): back off unless this was the last attempt. if attempt < attempts - 1 and not plugin.stopping: log("payload_retry", level="warning", attempt=attempt + 1, response_code=result.status_code) await plugin.sleep(min(base ** attempt, 60)) return Outcome.RETRY async def replay_buffer(plugin: "AissurancePlugin") -> None: """Replay buffered payloads oldest-first. Stops at the first non-OK send.""" for path in plugin.buffer.pending(): if plugin.stopping: return body = plugin.buffer.load(path) if body is None: continue # corrupt file already dropped by load() result = await plugin.transport.send_payload(body) if result.outcome is Outcome.OK: plugin.buffer.delete(path) plugin.sent_total += 1 log("payload_sent", router_id=plugin.cfg.router_id, response_code=result.status_code, latency_ms=result.latency_ms, replayed=True) elif result.outcome is Outcome.AUTH_FAILED: await plugin.handle_auth_failure() return else: # Still offline / rate limited — keep the rest buffered for next time. return async def run_discovery_cycle(plugin: "AissurancePlugin") -> None: """One scheduled discovery snapshot.""" try: snapshot = await asyncio.wait_for( collect_snapshot(plugin.cfg_obj, plugin.started_at), timeout=_SNAPSHOT_TIMEOUT ) except asyncio.TimeoutError: log("snapshot_too_slow", level="warning", timeout_s=_SNAPSHOT_TIMEOUT) return except Exception as exc: # serialization / state read error — skip this cycle log("snapshot_failed", level="error", error=str(exc)) return body = payload_mod.build_payload(plugin.cfg, snapshot) outcome = await send_with_retry(plugin, body) if outcome is Outcome.OK: await replay_buffer(plugin) elif outcome is Outcome.AUTH_FAILED: await plugin.handle_auth_failure() elif outcome is Outcome.RATE_LIMITED: # Server already rate-limited us; drop this snapshot (next is as good). log("payload_rate_limited", level="warning") else: # RETRY exhausted — persist for later. path = plugin.buffer.save(body, int(time.time() * 1000)) log("payload_buffered", level="warning", buffered=plugin.buffer.count(), path=str(path) if path else None) async def run_config_report(plugin: "AissurancePlugin") -> None: """POST current config to /config/report (handshake / on change).""" result = await plugin.transport.send_config_report( payload_mod.build_config_report(plugin.cfg) ) if result.outcome is Outcome.AUTH_FAILED: await plugin.handle_auth_failure() return suggested = (result.body or {}).get("server_suggested_polling_interval") log("config_reported", response_code=result.status_code, server_suggested_polling_interval=suggested) async def run_health_cycle(plugin: "AissurancePlugin") -> None: """Periodic health/keepalive + config-sync probe.""" result = await plugin.transport.send_health() if result.outcome is Outcome.AUTH_FAILED: await plugin.handle_auth_failure()