#!/usr/bin/env python3 """test_offline_queue.py — H2: writes survive an outage; cold start degrades to cache. Phase 1 — vault DOWN: `put` / `append` are queued (exit 0, "queued"), not lost. Phase 2 — vault UP: `flush` replays them idempotently; the writes land. Phase 3 — vault UP: `load` caches the orientation reads. Phase 4 — vault DOWN: `load` serves the last-known-good cache and says OFFLINE. Drives the real echo.py against eval/mock_olrapi.py. No creds, no live vault. Run: python test_offline_queue.py [--port 8840] """ from __future__ import annotations import argparse import os import subprocess import sys import tempfile import time import urllib.request from pathlib import Path HERE = Path(__file__).resolve().parent ECHO = HERE.parent / "echo-memory.plugin.src" / "skills" / "echo-memory" / "scripts" / "echo.py" KEY = "test-key-not-a-real-secret" DEAD = "http://127.0.0.1:59998" # nothing listens here -> connection refused failures = [] def check(name, cond, detail=""): print(f"{'ok ' if cond else 'FAIL'} {name}" + (f" -- {detail}" if not cond else "")) if not cond: failures.append(name) def http(method, url, body=None): data = body.encode() if isinstance(body, str) else body req = urllib.request.Request(url, data=data, method=method, headers={"Authorization": f"Bearer {KEY}"}) try: with urllib.request.urlopen(req, timeout=10) as r: return r.status, r.read().decode("utf-8", "replace") except Exception as e: # noqa: BLE001 return getattr(e, "code", 0), "" def tmp(content): f = tempfile.NamedTemporaryFile("w", suffix=".md", delete=False, encoding="utf-8") f.write(content); f.close() return f.name def main(): ap = argparse.ArgumentParser() ap.add_argument("--port", type=int, default=8840) a = ap.parse_args() mock_base = f"http://127.0.0.1:{a.port}" state = tempfile.mkdtemp() def echo(base, *args): env = dict(os.environ, ECHO_BASE=base, ECHO_KEY=KEY, ECHO_STATE_DIR=state, ECHO_VERIFY="0", ECHO_TODAY="2026-06-22") return subprocess.run([sys.executable, str(ECHO), *args], capture_output=True, text=True, env=env) def start_mock(): p = subprocess.Popen([sys.executable, str(HERE / "mock_olrapi.py"), "--port", str(a.port)], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True) for _ in range(50): try: urllib.request.urlopen(f"{mock_base}/__debug__reset", data=b"", timeout=1); break except Exception: time.sleep(0.1) return p def ground(path): _, body = http("GET", f"{mock_base}/__debug__?path={path}") return None if body == "<>" else body srv = None try: # --- Phase 1: vault DOWN -> writes are queued, not lost ------------------ r = echo(DEAD, "put", "_agent/sessions/2026-06-22-1200-x.md", tmp("EVALMARK-session\n")) check("offline put exits 0 (queued)", r.returncode == 0, r.stderr) check("offline put reports queued", "queued (offline)" in r.stdout, r.stdout) r = echo(DEAD, "append", "inbox/captures/inbox.md", "- 2026-06-22: EVALMARK-inbox") check("offline append reports queued", "queued (offline)" in r.stdout, r.stdout) outbox = Path(state) / "outbox.ndjson" check("two writes are persisted to the outbox", outbox.exists() and len(outbox.read_text(encoding="utf-8").splitlines()) == 2) # --- Phase 2: vault UP -> flush replays them ----------------------------- srv = start_mock() r = echo(mock_base, "flush") check("flush replays the queue", "flushed 2" in r.stdout, r.stdout + r.stderr) check("queued PUT landed in the vault", (ground("_agent/sessions/2026-06-22-1200-x.md") or "").find("EVALMARK-session") >= 0) check("queued APPEND landed in the vault", "EVALMARK-inbox" in (ground("inbox/captures/inbox.md") or "")) check("outbox is empty after a clean flush", not outbox.exists() or not outbox.read_text(encoding="utf-8").strip()) # idempotent replay: a second flush of the same content must not duplicate. echo(DEAD, "append", "inbox/captures/inbox.md", "- 2026-06-22: EVALMARK-inbox") # re-queue same line echo(mock_base, "flush") inbox = ground("inbox/captures/inbox.md") or "" check("replay is idempotent (no duplicate line)", inbox.count("EVALMARK-inbox") == 1, inbox) # --- Phase 3: vault UP -> load caches the orientation reads --------------- http("PUT", f"{mock_base}/vault/_agent/echo-vault.md", "schema_version: 4\nEVALMARK-marker\n") r = echo(mock_base, "load") check("online load shows marker", "EVALMARK-marker" in r.stdout, r.stdout) check("load wrote a cache dir", (Path(state) / "cache").exists()) # --- Phase 4: vault DOWN -> load serves the cache ------------------------ srv.terminate(); srv.wait(timeout=5); srv = None r = echo(DEAD, "load") check("offline load flags OFFLINE", "OFFLINE" in r.stdout, r.stdout) check("offline load serves cached marker", "EVALMARK-marker" in r.stdout, r.stdout) print(f"\n{len(failures)} failure(s)" if failures else "\nall offline-queue tests passed") return 1 if failures else 0 finally: if srv: srv.terminate() if __name__ == "__main__": raise SystemExit(main())