Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .amplifier/digital-twin-universe/profiles/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ seams without those; that is the point of an end-to-end test.
| `context-intelligence-write-fanout-validation.yaml` | **WRITE fan-out** — one session, TWO `destinations`; proves BOTH servers received the events (observes existing hook fan-out; never modifies it) | Incus + Gitea mirror + LLM + **2 CI servers** | `... launch .../context-intelligence-write-fanout-validation.yaml --var gitea_host=... --var ci_server_a=... --var ci_server_b=...` |
| `context-intelligence-query-validation.yaml` | **EXECUTE queries** (read side) — after logging, drives `graph_query` (Cypher) + `blob_read` (`ci-blob://`) via the `graph-analyst` agent; proves real rows/content come back with the `source` provenance naming the server | Incus + Gitea mirror + LLM + **CI server** | `... launch .../context-intelligence-query-validation.yaml --var gitea_host=... --var ci_server_url=...` |
| `context-intelligence-upload-format-validation.yaml` | **Legacy hooks-logging IMPORT** — `--format logging-hook` ingests a shipped neutral synthetic legacy fixture; proves discrimination, runtime slug parity/no fork (graph workspaces exactly equal the runtime-derived set), coexistence/dedupe/idempotency (node count captured at runtime does not grow); self-contained, no host data, no pinned counts | Incus + Gitea mirror + **CI backend** (`context-intelligence-backend.yaml` launched fresh) — no LLM | `... launch .../context-intelligence-upload-format-validation.yaml --var gitea_host=... --var server_url=http://<gateway-ip>:38000 --var server_token=...` |
| `context-intelligence-self-healing-replay-validation.yaml` | Self-healing event-replay (issue #403) — S1 reproduces the shortfall through a latency-injecting proxy in front of a real server, S2 a later session heals it, S3 replay is duplicate-free (fixed point on Neo4j node counts), S4 a real outage stays loud, S5 the age bound strands old events audibly | Incus + Gitea mirror + LLM + CI backend | `... launch .../context-intelligence-self-healing-replay-validation.yaml --var gitea_host=... --var server_url=... --var server_token=... --var run_id=... --var latency_ms=... --var queue_capacity=...` |
| `example-dtu-external-server.yaml` | *Not a test* — reference profile: point the client hook at an **external CI server** with a tagged workspace | Incus + running CI server (below) | see below |

**Self-contained smoke to prove the harness works on your host:**
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
#!/usr/bin/env python3
"""Latency-injecting HTTP proxy — the piece that makes the bug reproducible.

The delivery shortfall this profile validates only exists when the destination
is SLOWER than the producer. A localhost POST is sub-millisecond, so the
single-in-flight dispatcher keeps up trivially and NOTHING reproduces. Sitting
this in front of a real Context-Intelligence server reproduces the measured
Azure/APIM round-trip (~250-300 ms) without needing Azure.

PER-REQUEST latency, not per-connection. httpx reuses keep-alive connections, so
a connection-level sleep would delay only the first request on each connection
and the queue would never build. That is why this speaks HTTP/1.1 rather than
piping bytes between sockets.

Stdlib only (http.server + urllib): the DTU container has python3 but is not
guaranteed to have httpx before the bundle install runs, and this proxy must be
able to start first.

python3 latency_proxy.py --listen 8100 --upstream http://10.0.0.1:38000 \
--latency-ms 250 --stats /root/proxy-stats.json
"""

from __future__ import annotations

import argparse
import json
import threading
import time
import urllib.error
import urllib.request
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer

STATS = {"requests": 0, "bytes_in": 0, "errors": 0, "started_at": time.time()}
STATS_LOCK = threading.Lock()

UPSTREAM = ""
LATENCY_S = 0.0
STATS_PATH = ""


def _record(key: str, amount: int = 1) -> None:
with STATS_LOCK:
STATS[key] += amount
if STATS_PATH:
try:
with open(STATS_PATH, "w") as fh:
json.dump(STATS, fh)
except OSError:
pass


class Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"

def log_message(self, *args: object) -> None: # silence access logs
return

def _proxy(self, method: str) -> None:
length = int(self.headers.get("Content-Length") or 0)
body = self.rfile.read(length) if length else b""
_record("requests")
_record("bytes_in", len(body))

# THE POINT OF THIS FILE.
if LATENCY_S > 0:
time.sleep(LATENCY_S)

request = urllib.request.Request(
f"{UPSTREAM.rstrip('/')}{self.path}", data=body or None, method=method
)
for key, value in self.headers.items():
if key.lower() in ("host", "content-length", "connection", "transfer-encoding"):
continue
request.add_header(key, value)

try:
with urllib.request.urlopen(request, timeout=60) as response:
payload = response.read()
status = response.status
ctype = response.headers.get("Content-Type", "application/json")
except urllib.error.HTTPError as exc:
payload, status = exc.read(), exc.code
ctype = exc.headers.get("Content-Type", "application/json")
except Exception as exc: # upstream down -> an honest 502, never a hang
_record("errors")
payload = f"proxy upstream error: {exc}".encode()
status, ctype = 502, "text/plain"

self.send_response(status)
self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(len(payload)))
self.end_headers()
self.wfile.write(payload)

def do_POST(self) -> None: # BaseHTTPRequestHandler API
self._proxy("POST")

def do_GET(self) -> None: # BaseHTTPRequestHandler API
self._proxy("GET")

def do_DELETE(self) -> None: # BaseHTTPRequestHandler API
self._proxy("DELETE")


def main() -> None:
global UPSTREAM, LATENCY_S, STATS_PATH
parser = argparse.ArgumentParser()
parser.add_argument("--listen", type=int, default=8100)
parser.add_argument("--upstream", required=True)
parser.add_argument("--latency-ms", type=float, default=250.0)
parser.add_argument("--stats", default="")
args = parser.parse_args()

UPSTREAM = args.upstream
LATENCY_S = args.latency_ms / 1000.0
STATS_PATH = args.stats

server = ThreadingHTTPServer(("127.0.0.1", args.listen), Handler)
print(
f"latency-proxy :{args.listen} -> {UPSTREAM} (+{args.latency_ms}ms/request)",
flush=True,
)
server.serve_forever()


if __name__ == "__main__":
main()
Original file line number Diff line number Diff line change
@@ -0,0 +1,253 @@
#!/usr/bin/env python3
"""Assertions for the self-healing replay profile. Stdlib only.

Three lessons from the spike are encoded here rather than left to the caller,
because each one produced a WRONG result first:

1. **Wait for the drain.** ``POST /events`` returns 202 after a durable append,
BEFORE the async flush to Neo4j. Counting immediately reads a number that is
still climbing, so a duplicate check passes or fails for the wrong reason.

2. **Never assert ``:Event count == record count``.** Not every record becomes an
``:Event`` node -- some become ``ContentBlock`` / ``SST_EVENT`` nodes. The
obvious assertion fails at 399 vs 400 on a perfectly healthy system.
Convergence is proven as a FIXED POINT: one more full replay changes nothing,
which demonstrates complete delivery AND idempotence in one move.

3. **Namespace the workspace per run.** The graph persists between runs. Reusing
a workspace name makes run 2 read run 1's nodes and report "bug did not
reproduce" -- a passing test that asserts nothing.

Usage:
verify.py nodes --server URL --token T --workspace WS [--wait]
verify.py watermark --session-dir DIR --destination NAME
verify.py lines --session-dir DIR
verify.py reproduced --session-dir DIR --server URL --token T --workspace WS
verify.py healed --session-dir DIR --server URL --token T --workspace WS
verify.py stalled --session-dir DIR --destination NAME
"""

from __future__ import annotations

import argparse
import json
import sys
import time
import urllib.request
from pathlib import Path


def fail(message: str) -> None:
print(f"FAIL: {message}")
sys.exit(1)


def ok(message: str) -> None:
print(f"PASS: {message}")


def cypher(server: str, token: str, query: str, params: dict) -> list[dict]:
body = json.dumps({"query": query, "params": params}).encode()
request = urllib.request.Request(
f"{server.rstrip('/')}/cypher",
data=body,
method="POST",
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"},
)
with urllib.request.urlopen(request, timeout=60) as response:
return json.loads(response.read()).get("results", [])


def count_nodes(server: str, token: str, workspace: str) -> int:
rows = cypher(
server, token, "MATCH (n) WHERE n.workspace = $ws RETURN count(n) AS c", {"ws": workspace}
)
return int(rows[0]["c"]) if rows else 0


def count_events(server: str, token: str, workspace: str) -> int:
rows = cypher(
server,
token,
"MATCH (n:Event) WHERE n.workspace = $ws RETURN count(n) AS c",
{"ws": workspace},
)
return int(rows[0]["c"]) if rows else 0


def wait_drained(server: str, token: str, workspace: str, timeout: float = 240.0) -> int:
"""Poll until the node count stops moving (lesson 1)."""
deadline = time.monotonic() + timeout
last, stable_at = -1, 0.0
while time.monotonic() < deadline:
current = count_nodes(server, token, workspace)
if current == last and current > 0:
if stable_at and (time.monotonic() - stable_at) >= 3.0:
return current
stable_at = stable_at or time.monotonic()
else:
last, stable_at = current, 0.0
time.sleep(1.0)
return count_nodes(server, token, workspace)


def count_session_nodes(server: str, token: str, workspace: str, session_id: str) -> int:
"""Nodes belonging to ONE session.

The fixed-point check must be scoped this way. A workspace-wide count is
wrong the moment the scenario runs further sessions to drive the sweep --
their brand-new events inflate the total and look exactly like duplicates.
node_id is deterministic and prefixed with the session id (see the server's
make_node_id), so the prefix is an exact per-session filter.
"""
rows = cypher(
server,
token,
"MATCH (n) WHERE n.workspace = $ws AND n.node_id STARTS WITH $prefix RETURN count(n) AS c",
{"ws": workspace, "prefix": f"{session_id}__"},
)
return int(rows[0]["c"]) if rows else 0


def count_lines(session_dir: Path) -> int:
path = session_dir / "events.jsonl"
if not path.exists():
return 0
with path.open("rb") as fh:
return sum(1 for line in fh if line.strip())


def load_watermark(session_dir: Path, destination: str) -> dict:
path = session_dir / "delivery" / f"{destination}.json"
if not path.exists():
fail(f"no watermark at {path} -- Increment 0 did not write one")
return json.loads(path.read_text())


def newest_session(root: Path) -> Path:
candidates = sorted(
root.glob("*/sessions/*/context-intelligence"),
key=lambda p: (p / "events.jsonl").stat().st_mtime if (p / "events.jsonl").exists() else 0,
)
if not candidates:
fail(f"no session directory under {root}")
return candidates[-1]


def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("command")
parser.add_argument("--server", default="")
parser.add_argument("--token", default="")
parser.add_argument("--workspace", default="")
parser.add_argument("--session-dir", default="")
parser.add_argument("--destination", default="main")
parser.add_argument("--expect-nodes", type=int, default=-1)
parser.add_argument("--wait", action="store_true")
args = parser.parse_args()

session_dir = Path(args.session_dir) if args.session_dir else None

if args.command == "nodes":
total = (
wait_drained(args.server, args.token, args.workspace)
if args.wait
else count_nodes(args.server, args.token, args.workspace)
)
events = count_events(args.server, args.token, args.workspace)
print(json.dumps({"nodes": total, "events": events, "workspace": args.workspace}))
return

if args.command == "lines":
assert session_dir
print(json.dumps({"lines": count_lines(session_dir)}))
return

if args.command == "watermark":
assert session_dir
print(json.dumps(load_watermark(session_dir, args.destination)))
return

if args.command == "reproduced":
# S1: the live path must have FAILED to deliver everything.
assert session_dir
lines = count_lines(session_dir)
delivered = wait_drained(args.server, args.token, args.workspace, timeout=60)
events = count_events(args.server, args.token, args.workspace)
watermark = load_watermark(session_dir, args.destination)
size = (session_dir / "events.jsonl").stat().st_size
print(
json.dumps(
{"lines": lines, "nodes": delivered, "events": events, "watermark": watermark}
)
)
if events >= lines:
fail(
f"bug did NOT reproduce: server holds {events} :Event for {lines} records. "
"Raise the proxy latency or lower dispatch_queue_capacity."
)
if watermark["offset"] >= size:
fail("watermark claims EOF but the server is behind -- the cursor is lying")
ok(
f"reproduced: {events} :Event on the server for {lines} local records; "
f"watermark at {watermark['offset']}/{size}"
)
return

if args.command == "session-nodes":
assert session_dir
session_id = session_dir.parent.name
print(
json.dumps(
{
"session_id": session_id,
"nodes": count_session_nodes(
args.server, args.token, args.workspace, session_id
),
}
)
)
return

if args.command == "healed":
# S2/S3: converge, watermark at EOF, and a further replay is a no-op.
assert session_dir
size = (session_dir / "events.jsonl").stat().st_size
total = wait_drained(args.server, args.token, args.workspace)
watermark = load_watermark(session_dir, args.destination)
if watermark["offset"] != size:
fail(f"watermark {watermark['offset']} != EOF {size} -- backlog not fully swept")
if args.expect_nodes >= 0:
session_id = session_dir.parent.name
scoped = count_session_nodes(args.server, args.token, args.workspace, session_id)
if scoped != args.expect_nodes:
fail(
f"NOT a fixed point for session {session_id}: "
f"{args.expect_nodes} -> {scoped} nodes after a full replay"
)
print(json.dumps({"session_nodes": scoped, "workspace_nodes": total}))
print(json.dumps({"nodes": total, "watermark": watermark}))
ok(f"healed: watermark at EOF ({size}), {total} nodes for {args.workspace}")
return

if args.command == "stalled":
# S4: a real outage must leave the cursor where it was, loudly.
assert session_dir
watermark = load_watermark(session_dir, args.destination)
print(json.dumps(watermark))
if watermark["offset"] != 0:
fail(f"watermark advanced to {watermark['offset']} against a broken destination")
if watermark.get("last_outcome") != "no_progress":
fail(f"expected last_outcome=no_progress, got {watermark.get('last_outcome')!r}")
if not watermark.get("last_error"):
fail("no last_error recorded -- the failure would be invisible")
if not watermark.get("consecutive_sweeps_without_progress"):
fail("no-progress counter did not increment")
ok(f"outage stayed loud: offset 0, last_error={watermark['last_error'][:60]!r}")
return

fail(f"unknown command {args.command!r}")


if __name__ == "__main__":
main()
Loading
Loading