Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
52 changes: 52 additions & 0 deletions scripts/livesync-a1-watch.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
"""A1 live proof: cross-process commit→fetch latency through session.changes.

Polls the dashboard's session.changes (the desktop's poll path) at the shipped
2.5s interval against a session the MESSAGING GATEWAY (other process) is
writing. For each newly returned row, latency = fetch_time - row.timestamp.
This is the gateway-written/dashboard-read A1 proof with measured numbers.
"""
import asyncio, json, sys, time
sys.path.insert(0, 'scripts')
probe = __import__('livesync-rc1-probe')

SID = sys.argv[4] if len(sys.argv) > 4 else ""
DURATION = 90

if not SID:
raise SystemExit(
"usage: livesync-a1-watch.py <user> <pw> <base_url> <stored_session_id>\n"
"stored_session_id is REQUIRED (state.db key of the session being "
"written by the messaging gateway)."
)

async def main():
c = probe.Client('a1')
await c.connect()
r = await c.call("session.changes", {"session_id": SID, "since_message_id": 0})
Comment thread
Kyzcreig marked this conversation as resolved.
res = r.get("result") or {}
cursor = res.get("last_id", 0)
print(f"A1 watch start: session={SID} cursor={cursor} preloaded={len(res.get('messages', []))} rows")
t_end = time.time() + DURATION
lat = []
while time.time() < t_end:
await asyncio.sleep(2.5)
r = await c.call("session.changes", {"session_id": SID, "since_message_id": cursor})
res = r.get("result") or {}
rows = res.get("messages", [])
now = time.time()
for m in rows:
ts = m.get("timestamp")
if isinstance(ts, (int, float)) and ts > 1e9:
d = now - ts
lat.append(d)
print(f" row id={m.get('id')} role={m.get('role')} commit→fetch={d:.2f}s")
cursor = res.get("last_id", cursor)
await c.close()
if lat:
print(f"\nA1 RESULT: {len(lat)} rows; max latency {max(lat):.2f}s; "
f"median {sorted(lat)[len(lat)//2]:.2f}s; SLO ≤3s: "
f"{'PASS' if max(lat) <= 3.0 else 'CHECK (relax to ≤4s per spec if 3<x≤4)'}")
else:
print("\nA1: no new rows observed in window — re-run while the session is active")

asyncio.run(main())
4 changes: 2 additions & 2 deletions scripts/livesync-derive-tsilence.py
Original file line number Diff line number Diff line change
Expand Up @@ -121,9 +121,9 @@ async def call(method, params, timeout=120):
last_frame[0] = time.monotonic()
await call("prompt.submit", {"session_id": sid, "text": prompt}, timeout=30)
try:
await asyncio.wait_for(turn_done.wait(), timeout=300)
await asyncio.wait_for(turn_done.wait(), timeout=120)
except asyncio.TimeoutError:
print(" [warn] turn did not settle in 300s; gaps recorded anyway")
print(" [warn] no settle event in 120s; gaps recorded anyway")
await asyncio.sleep(2)

try:
Expand Down
66 changes: 46 additions & 20 deletions scripts/livesync-rc1-probe.py
Original file line number Diff line number Diff line change
Expand Up @@ -109,20 +109,33 @@ async def main() -> int:
ok_all = True
sid = ""
try:
# Disposable seeded session (per dashboard-ws-rpc-harness).
seed = [
{"role": r_, "content": f"{r_} {i}: probe filler " * 5}
for i in range(3)
for r_ in ("user", "assistant")
]
r = await c1.call(
"session.create", {"messages": seed, "title": "rc1-probe", "source": "tui"}
)
sid = (r.get("result") or {}).get("session_id", "")
# Probe over an EXISTING stored session (read-only: resume/changes/
# active_list only — never delete). `session.create` seeds only the
# live registry; its state.db row persists lazily on the first real
# turn, so a synthetic session has no committed rows for
# `session.changes` (found live 2026-07-10). session.changes /
# session.resume take the STORED id (state.db key). Pass argv[4] to
# pin one; otherwise pick the most recent small idle cli/tui session.
import os
import sqlite3

if len(sys.argv) > 4:
sid = sys.argv[4]
else:
sdb = sqlite3.connect(os.path.expanduser("~/.hermes/state.db"))
pick = sdb.execute(
"SELECT m.session_id, COUNT(*) c FROM messages m "
"JOIN sessions s ON s.id = m.session_id "
"WHERE s.source IN ('cli','tui') "
"GROUP BY m.session_id HAVING c BETWEEN 2 AND 50 "
"ORDER BY MAX(m.id) DESC LIMIT 1"
).fetchone()
sdb.close()
sid = pick[0] if pick else ""
if not sid:
print(f"FATAL: session.create failed: {json.dumps(r)[:300]}")
print("FATAL: no suitable idle stored session; pass one as argv[4]")
return 1
print(f"probe session: {sid}")
print(f"probe session (existing, read-only): {sid}")

# Client-1 resumes (becomes transport holder), takes a cursor.
r = await c1.call("session.resume", {"session_id": sid})
Expand All @@ -136,13 +149,17 @@ async def main() -> int:
if m.get("id") is not None
]
ok_all &= row(
"c1 pre-steal session.changes", "result" in r, f"{len(pre_ids)} rows"
"c1 pre-steal session.changes returns committed rows",
"result" in r and len(pre_ids) > 0,
f"{len(pre_ids)} rows",
)
cursor = (r.get("result") or {}).get("last_id", 0)

# THE STEAL: client-2 resumes the same session. WATCH-ONLY — no send.
r = await c2.call("session.resume", {"session_id": sid})
ok_all &= row("c2 steal (session.resume)", "result" in r)
# c2 is now the holder — its LIVE registry id keys the active_list row.
c2_live_id = (r.get("result") or {}).get("session_id", "")
await asyncio.sleep(1.0)

# (a) ROUTING: c1's polls/probes still round-trip on ITS socket.
Expand All @@ -165,12 +182,25 @@ async def main() -> int:

# (b) SEMANTIC: status for the stolen session settles idle/waiting
# (no turn is running; a lost frame must read as settled, not working).
# active_list rows key by the LIVE registry id (captured from c2's
# steal above). Match rows on stored-id fields OR that live id — NO
# loose fallback: an unrelated idle session must not be able to green
# the semantic (Greptile #272).
status = None
deadline = time.monotonic() + 20
while time.monotonic() < deadline:
status = None # reset per iteration: only the CURRENT read counts
r = await c1.call("session.active_list", {}, timeout=15)
for s in (r.get("result") or {}).get("sessions", []):
if s.get("session_id") == sid or s.get("id") == sid:
rows_ = (r.get("result") or {}).get("sessions", [])
for s in rows_:
candidates = (
s.get("session_id"),
s.get("id"),
s.get("stored_session_id"),
s.get("resumed"),
s.get("session_key"),
)
if sid in candidates or (c2_live_id and c2_live_id in candidates):
status = s.get("status")
if status in ("idle", "waiting"):
break
Expand All @@ -181,11 +211,7 @@ async def main() -> int:
f"status={status!r}",
)
finally:
if sid:
try:
await c1.call("session.delete", {"session_id": sid}, timeout=30)
except Exception: # noqa: BLE001
print(" [warn] probe session cleanup failed — delete manually:", sid)
# Read-only probe over a pre-existing session — nothing to delete.
await c1.close()
await c2.close()

Expand Down
Loading