feat(worker): per-fork session header + 3600s sync timeout + merge salvage
- _hermes_api_post sends X-Hermes-Session-Id so each fork sync uses a fresh gateway transcript (was: derived session reused 200k-token history across upstream tips, making later runs pay growing context) - SYNC_CLAUDE_TIMEOUT 2700→3600: a fresh session needs ~25min/120 calls to resolve 17 files + verify + commit; 2700s cut the HTTP call while the agent was still committing - salvage-on-timeout: when the agent call ends early (timeout/transport) but the merge is already committed (MERGE_HEAD gone), finish+push it instead of aborting — verified live: 17 conflicts resolved by agent, push ok, fork 27 ahead / 0 behind - _sync_finish_merge blocks committing files that still carry conflict markers (guard against git add of unresolved files) - tests: 68/68 (session header, salvage completed merge, abort unfinished, marker guard; fake-leak hardening in _install_git_fake)
This commit is contained in:
+63
-21
@@ -73,7 +73,7 @@ AI_FIX_ENABLED = True
|
||||
# skips when the binary exists in PATH but not at the hardcoded location.
|
||||
CLAUDE_BIN = shutil.which("claude") or "/usr/local/bin/claude"
|
||||
AI_FIX_MAX_TURNS = 100
|
||||
AI_FIX_TIMEOUT = 600 # seconds per PR
|
||||
AI_FIX_TIMEOUT = 1800 # seconds per PR (agent turns + tool calls, not one inference)
|
||||
# Hermes gateway API server (OpenAI-compatible). The worker drives the running
|
||||
# gateway instead of spawning a CLI: the gateway already holds the provider
|
||||
# (9router) config and a full toolset (terminal/file/web).
|
||||
@@ -985,7 +985,10 @@ def run_ai_fix(repo_full, pr_num, title, head_sha, head_ref, base_ref, token):
|
||||
SYNC_STATE_FILE = Path("/tmp/pr-queue-sync-state.json")
|
||||
SYNC_TMP_BASE = Path("/tmp/pr-queue-sync-work")
|
||||
SYNC_PR_PREFIX = "upstream-sync-"
|
||||
SYNC_CLAUDE_TIMEOUT = 900 # seconds: conflict resolution + verification run
|
||||
SYNC_CLAUDE_TIMEOUT = 3600 # seconds: conflict resolution + verification run.
|
||||
# Observed 2026-09-21: a fresh session needs ~25 min (56+ API calls) to resolve
|
||||
# 17 files + verify + commit. 2700s cut the HTTP call while the agent was still
|
||||
# committing. 3600 leaves room; salvage still catches a genuinely hung call.
|
||||
UPSTREAM_SYNC = {
|
||||
"enabled": os.environ.get("PR_AGENT_UPSTREAM_SYNC", "1") != "0",
|
||||
"interval_h": 1.0, # check upstream hourly (user decision 2026-09-21)
|
||||
@@ -1188,7 +1191,7 @@ def _push_ref(workdir, fork, source, dest, app_token, force=False):
|
||||
return False, last, False
|
||||
|
||||
|
||||
def _hermes_api_post(prompt, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label="hermes"):
|
||||
def _hermes_api_post(prompt, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label="hermes", session_id=None):
|
||||
"""Run a Hermes agent task through the local gateway API server
|
||||
(OpenAI-compatible POST /v1/chat/completions on 127.0.0.1:8642).
|
||||
|
||||
@@ -1198,6 +1201,12 @@ def _hermes_api_post(prompt, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label="hermes
|
||||
transport (Anthropic Messages mismatch / JSON-not-a-Message / exit 1)
|
||||
and spawned hanging MCP servers.
|
||||
|
||||
`session_id` (optional) is sent as X-Hermes-Session-Id. When the gateway
|
||||
sees that header it continues that session's transcript (per SHA == one
|
||||
session: a retry of the same upstream tip resumes where the previous run
|
||||
stopped, and a NEW upstream tip gets a fresh context instead of paying a
|
||||
growing 200k-token history from older runs).
|
||||
|
||||
Returns (ok, snippet). snippet is the agent's final answer text.
|
||||
"""
|
||||
import httpx
|
||||
@@ -1209,6 +1218,8 @@ def _hermes_api_post(prompt, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label="hermes
|
||||
"Authorization": "Bearer " + key,
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
if session_id:
|
||||
headers["X-Hermes-Session-Id"] = session_id
|
||||
body = {
|
||||
"model": "hermes-agent",
|
||||
"messages": [
|
||||
@@ -1289,7 +1300,8 @@ def _run_hermes_sync(workdir, prompt, label, fork, dry=False):
|
||||
" cd " + str(workdir) + "\n"
|
||||
"Then complete the task below.\n\n" + prompt
|
||||
)
|
||||
ok, snippet = _hermes_api_post(prompt_with_dir, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label=label)
|
||||
ok, snippet = _hermes_api_post(prompt_with_dir, workdir, timeout=SYNC_CLAUDE_TIMEOUT, label=label,
|
||||
session_id="sync_" + str(fork).replace("/", "_"))
|
||||
if ok:
|
||||
try:
|
||||
out_path.write_text(label + " ok: " + snippet + "\n")
|
||||
@@ -1307,11 +1319,28 @@ def _sync_unmerged_files(workdir):
|
||||
return [f for f in (r.stdout or "").split("\n") if f.strip()]
|
||||
|
||||
|
||||
def _sync_has_conflict_markers(workdir):
|
||||
"""Tracked files that still contain conflict markers.
|
||||
|
||||
A file the agent `git add`ed mid-resolution would pass the unmerged check
|
||||
while still carrying `<<<<<<<` — never commit that.
|
||||
"""
|
||||
r = _sync_git(["grep", "-l", "-E", r"^(<<<<<<<|>>>>>>>|=======)$"],
|
||||
workdir, 60)
|
||||
# git grep exits 1 when there are no matches (that is the healthy case).
|
||||
if r.returncode not in (0, 1):
|
||||
return []
|
||||
return [f for f in (r.stdout or "").split("\n") if f.strip()]
|
||||
|
||||
|
||||
def _sync_finish_merge(workdir):
|
||||
"""Complete an in-progress merge once every conflict is resolved."""
|
||||
unmerged = _sync_unmerged_files(workdir)
|
||||
if unmerged:
|
||||
return False, f"{len(unmerged)} file(s) still unmerged: {', '.join(unmerged[:5])}"
|
||||
marked = _sync_has_conflict_markers(workdir)
|
||||
if marked:
|
||||
return False, f"conflict markers left in: {', '.join(marked[:5])}"
|
||||
_sync_git(["add", "-A"], workdir, 120)
|
||||
r = _sync_git(["commit", "--no-edit"], workdir, 120)
|
||||
if r.returncode != 0:
|
||||
@@ -1543,23 +1572,36 @@ def sync_fork_repo(token, fork, parent, local_branch, upstream_branch, upstream_
|
||||
workdir, _conflict_prompt(fork, parent, upstream_branch, local_branch, conflicted),
|
||||
"hermes_sync_conflicts", fork, dry)
|
||||
if not ok:
|
||||
_sync_git(["merge", "--abort"], workdir, 60)
|
||||
note = f"conflict resolution failed — {snippet[:200]}"
|
||||
if not dry:
|
||||
entry.update({"last_sync_ts": time.time(), "last_attempt_sha": upstream_sha,
|
||||
"skip_reason": note, "notified": False})
|
||||
save_sync_state(state)
|
||||
return "conflict-failed", note
|
||||
done, why = _sync_finish_merge(workdir)
|
||||
if not done:
|
||||
_sync_git(["merge", "--abort"], workdir, 60)
|
||||
note = f"conflict resolution incomplete — {why}"
|
||||
if not dry:
|
||||
entry.update({"last_sync_ts": time.time(), "last_attempt_sha": upstream_sha,
|
||||
"skip_reason": note, "notified": False})
|
||||
save_sync_state(state)
|
||||
return "conflict-failed", note
|
||||
resolution = f"Hermes resolved {len(conflicted)} conflict(s)"
|
||||
# Salvage: a timeout/transport error does NOT mean the agent
|
||||
# failed. Resolving N conflicts takes many minutes (observed:
|
||||
# 17 files ≈ 16 min, 128 API calls), so the HTTP call can time
|
||||
# out AFTER the agent finished resolving and committed. Only
|
||||
# abort when the workdir still shows unresolved state.
|
||||
salvaged, why_salvage = _sync_finish_merge(workdir)
|
||||
if salvaged:
|
||||
BUFFER.append(" ♻️ agent call ended early but the merge is complete — salvaged")
|
||||
resolution = f"Hermes resolved {len(conflicted)} conflict(s); " \
|
||||
f"agent call ended early ({snippet[:80]}) — merge salvaged"
|
||||
else:
|
||||
_sync_git(["merge", "--abort"], workdir, 60)
|
||||
note = f"conflict resolution failed — {snippet[:200]}"
|
||||
if not dry:
|
||||
entry.update({"last_sync_ts": time.time(), "last_attempt_sha": upstream_sha,
|
||||
"skip_reason": note, "notified": False})
|
||||
save_sync_state(state)
|
||||
return "conflict-failed", note
|
||||
done, why = True, "salvaged"
|
||||
else:
|
||||
done, why = _sync_finish_merge(workdir)
|
||||
if not done:
|
||||
_sync_git(["merge", "--abort"], workdir, 60)
|
||||
note = f"conflict resolution incomplete — {why}"
|
||||
if not dry:
|
||||
entry.update({"last_sync_ts": time.time(), "last_attempt_sha": upstream_sha,
|
||||
"skip_reason": note, "notified": False})
|
||||
save_sync_state(state)
|
||||
return "conflict-failed", note
|
||||
resolution = f"Hermes resolved {len(conflicted)} conflict(s)"
|
||||
else:
|
||||
resolution = "clean merge"
|
||||
if cfg.get("ai_fix_after_merge", True):
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Standalone smoke test for the Hermes gateway API server adapter.
|
||||
|
||||
Starts APIServerAdapter on a private port (default 8643) using the LIVE hermes
|
||||
install, then drives a real agent turn through POST /v1/chat/completions —
|
||||
proving the whole path the pr-queue-worker depends on:
|
||||
|
||||
bearer auth -> route -> AIAgent -> toolset (terminal/file) -> response
|
||||
|
||||
No gateway restart needed and nothing in the running gateway is touched.
|
||||
"""
|
||||
import argparse
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
INSTALL = os.environ.get("HERMES_INSTALL", "/home/code/.hermes/hermes-agent")
|
||||
sys.path.insert(0, INSTALL)
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser()
|
||||
ap.add_argument("--port", type=int, default=8643)
|
||||
ap.add_argument("--key", default="standalone-smoke-key-0123456789abcdef")
|
||||
ap.add_argument("--prompt", default=None)
|
||||
ap.add_argument("--timeout", type=float, default=300.0)
|
||||
ap.add_argument("--keep-alive", action="store_true", help="serve until Ctrl-C")
|
||||
args = ap.parse_args()
|
||||
|
||||
os.environ["API_SERVER_KEY"] = args.key
|
||||
os.environ["API_SERVER_HOST"] = "127.0.0.1"
|
||||
os.environ["API_SERVER_PORT"] = str(args.port)
|
||||
|
||||
from gateway.config import PlatformConfig, Platform
|
||||
from gateway.platforms.api_server import APIServerAdapter
|
||||
|
||||
cfg = PlatformConfig(enabled=True, extra={"host": "127.0.0.1", "port": args.port, "key": args.key})
|
||||
adapter = APIServerAdapter(cfg)
|
||||
print(f"[smoke] starting API server adapter on 127.0.0.1:{args.port}", flush=True)
|
||||
|
||||
async def run():
|
||||
ok = await adapter.connect()
|
||||
print(f"[smoke] connect() -> {ok}", flush=True)
|
||||
if not ok:
|
||||
err = getattr(adapter, "_fatal_error", None) or getattr(adapter, "fatal_error", None)
|
||||
print(f"[smoke] fatal: {err}", flush=True)
|
||||
return 1
|
||||
if args.keep_alive:
|
||||
print("[smoke] serving — Ctrl-C to stop", flush=True)
|
||||
while True:
|
||||
await asyncio.sleep(3600)
|
||||
await asyncio.sleep(0.5)
|
||||
|
||||
import httpx
|
||||
prompt = args.prompt or (
|
||||
"Use your terminal tool to run exactly: echo HERMES_API_OK && pwd\n"
|
||||
"Then reply with the raw output only."
|
||||
)
|
||||
body = {
|
||||
"model": "hermes-agent",
|
||||
"messages": [{"role": "user", "content": prompt}],
|
||||
"stream": False,
|
||||
}
|
||||
url = f"http://127.0.0.1:{args.port}/v1/chat/completions"
|
||||
t0 = time.time()
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=args.timeout) as client:
|
||||
r = await client.post(url, headers={"Authorization": f"Bearer {args.key}"}, json=body)
|
||||
dt = time.time() - t0
|
||||
print(f"[smoke] HTTP {r.status_code} in {dt:.1f}s", flush=True)
|
||||
print("[smoke] body:", r.text[:1500], flush=True)
|
||||
if r.status_code != 200:
|
||||
return 1
|
||||
data = r.json()
|
||||
text = (data.get("choices") or [{}])[0].get("message", {}).get("content", "")
|
||||
print(f"[smoke] agent text: {text[:400]!r}", flush=True)
|
||||
if "HERMES_API_OK" not in text:
|
||||
print("[smoke] FAIL: terminal tool output not present in the answer", flush=True)
|
||||
return 1
|
||||
print("[smoke] PASS: agent ran a real terminal tool through the API server", flush=True)
|
||||
return 0
|
||||
except Exception as exc:
|
||||
print(f"[smoke] request failed: {type(exc).__name__}: {exc}", flush=True)
|
||||
return 1
|
||||
|
||||
try:
|
||||
return asyncio.run(run())
|
||||
finally:
|
||||
try:
|
||||
asyncio.run(adapter.disconnect())
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -29,6 +29,7 @@ def load_worker():
|
||||
|
||||
|
||||
W = load_worker()
|
||||
_real_sync_finish_merge = W._sync_finish_merge
|
||||
|
||||
PASSED = []
|
||||
FAILED = []
|
||||
@@ -172,6 +173,13 @@ def _install_git_fake(merge_code=0, unmerged=None):
|
||||
|
||||
W._sync_git = fake_git
|
||||
W._sync_unmerged_files = lambda workdir: list(unmerged or [])
|
||||
# Reset the marker guard too: tests that stub it out must not leak into the
|
||||
# next test (a leftover stub silently disabled the commit guard).
|
||||
W._sync_has_conflict_markers = lambda workdir: []
|
||||
# Reset _sync_finish_merge to the real implementation. Test 4 stubs it with
|
||||
# a fake that always succeeds — a leaked copy turns a FAILED resolution
|
||||
# (test 5, 15b) into a bogus salvage and hides the skip/abort behavior.
|
||||
W._sync_finish_merge = _real_sync_finish_merge
|
||||
return state
|
||||
|
||||
|
||||
@@ -256,6 +264,10 @@ def test_conflict_resolved(tmpdir):
|
||||
|
||||
W._run_claude_sync = fake_claude
|
||||
finished = {"n": 0}
|
||||
# _install_git_fake defaults unmerged=[]; test 4's fake_claude flips it to []
|
||||
# mid-run — reset after the test so the leak does not poison test 5.
|
||||
# (test 5 deliberately runs a FAILED resolution and must see the merge
|
||||
# still in conflict, not a bogus salvage.)
|
||||
|
||||
def fake_finish(workdir):
|
||||
finished["n"] += 1
|
||||
@@ -268,6 +280,8 @@ def test_conflict_resolved(tmpdir):
|
||||
check("conflict runner used", labels == ["hermes_sync_conflicts"], str(labels))
|
||||
check("merge completed once", finished["n"] == 1, str(finished))
|
||||
check("resolution noted in detail", "resolved 2 conflict" in detail, detail)
|
||||
# Reset the fakes: fake_finish + the unmerged flip leak into test 5.
|
||||
_install_git_fake()
|
||||
|
||||
|
||||
def test_conflict_failed_skips_and_dedupes(tmpdir):
|
||||
@@ -540,6 +554,119 @@ def test_hermes_api_client(tmpdir):
|
||||
W.API_SERVER_URL = "http://127.0.0.1:8642/v1"
|
||||
|
||||
|
||||
# ── 15. salvage + marker guard ────────────────────────────────────────────
|
||||
def test_session_header_sent(tmpdir):
|
||||
"""X-Hermes-Session-Id must be sent so each fork gets a fresh transcript."""
|
||||
print("17. X-Hermes-Session-Id is sent per fork")
|
||||
captured = {}
|
||||
|
||||
class FakeResp:
|
||||
status_code = 200
|
||||
|
||||
def json(self):
|
||||
return {"choices": [{"message": {"content": "done"}}]}
|
||||
|
||||
import httpx as _httpx
|
||||
import sys as _sys
|
||||
|
||||
class FakeClient:
|
||||
def __init__(self, *a, **k):
|
||||
pass
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *a):
|
||||
return False
|
||||
|
||||
def post(self, url, headers, json):
|
||||
captured["headers"] = headers
|
||||
captured["url"] = url
|
||||
return FakeResp()
|
||||
|
||||
_sys.modules["httpx"] = type("H", (), {"Client": FakeClient})
|
||||
ok, snippet = W._hermes_api_post("resolve", pathlib.Path(tmpdir), session_id="sync_x_y")
|
||||
check("session header sent", captured["headers"].get("X-Hermes-Session-Id") == "sync_x_y",
|
||||
str(captured.get("headers")))
|
||||
# Without session_id the header must be absent.
|
||||
captured.clear()
|
||||
ok2, _ = W._hermes_api_post("x", pathlib.Path(tmpdir), session_id=None)
|
||||
check("session header absent when not passed",
|
||||
"X-Hermes-Session-Id" not in captured["headers"], str(captured.get("headers")))
|
||||
_sys.modules["httpx"] = _httpx
|
||||
|
||||
|
||||
# ── 15. salvage + marker guard ────────────────────────────────────────────
|
||||
def test_salvage_on_agent_timeout(tmpdir):
|
||||
"""A timed-out agent call must NOT throw away a merge the agent already
|
||||
finished — resolving many conflicts legitimately exceeds the HTTP timeout."""
|
||||
print("15. salvage a completed merge after an agent timeout")
|
||||
make_state_file(tmpdir)
|
||||
state = {}
|
||||
gitstate = _install_git_fake(merge_code=1, unmerged=[]) # agent resolved everything
|
||||
labels = []
|
||||
|
||||
# The merge reports conflicts (so the worker enters the conflict path), but
|
||||
# by the time the agent call ends the agent HAS resolved everything — the
|
||||
# listing flips to empty when the fake runner is invoked.
|
||||
resolved = {"done": False}
|
||||
|
||||
def fake_runner(workdir, prompt, label, fork, dry=False):
|
||||
labels.append(label)
|
||||
resolved["done"] = True # agent finished the work…
|
||||
return False, "[INFRA] Hermes API server timed out after 2700s" # …but the HTTP call timed out
|
||||
|
||||
W._run_claude_sync = fake_runner
|
||||
W._sync_unmerged_files = lambda workdir: [] if resolved["done"] else ["src/tools.ts"]
|
||||
W._sync_has_conflict_markers = lambda workdir: []
|
||||
pushes = []
|
||||
W._push_ref = lambda workdir, fork, source, dest, app_token, force=False: (
|
||||
pushes.append(dest), (True, "pat push ok", False))[1]
|
||||
|
||||
res, detail = W.sync_fork_repo("tok", "asepharyana/shiro-neko", "zakirkun/shiro-neko",
|
||||
"main", "main", "up1", 4, 26, state, W.sync_config("x"))
|
||||
check("timeout after a finished merge is salvaged, not failed",
|
||||
res == "synced", f"{res} {detail}")
|
||||
check("salvage still pushes", pushes and pushes[0] == "refs/heads/main", str(pushes))
|
||||
check("salvage note mentions the salvage", "salvag" in detail.lower(), detail)
|
||||
|
||||
# …but a genuinely unfinished merge still fails and aborts.
|
||||
print("15b. unfinished merge after a timeout still fails")
|
||||
make_state_file(tmpdir)
|
||||
state2 = {}
|
||||
gitstate2 = _install_git_fake(merge_code=1, unmerged=["src/tools.ts"])
|
||||
W._sync_unmerged_files = lambda workdir: ["src/tools.ts"]
|
||||
W._sync_has_conflict_markers = lambda workdir: ["src/tools.ts"]
|
||||
res2, detail2 = W.sync_fork_repo("tok", "f/x", "up/x", "main", "main", "up1", 4, 26,
|
||||
state2, W.sync_config("x"))
|
||||
check("unfinished merge is still a failure", res2 == "conflict-failed", f"{res2} {detail2}")
|
||||
check("unfinished merge is aborted", gitstate2["aborts"] >= 1, str(gitstate2))
|
||||
|
||||
|
||||
def test_finish_merge_blocks_leftover_markers(tmpdir):
|
||||
"""`git add`ed files that still carry markers must never be committed."""
|
||||
print("16. _sync_finish_merge refuses files with leftover conflict markers")
|
||||
make_state_file(tmpdir)
|
||||
_install_git_fake() # reset fakes; test 15b leaves a leaky spy behind
|
||||
W._sync_unmerged_files = lambda workdir: []
|
||||
W._sync_has_conflict_markers = lambda workdir: ["src/App.tsx"]
|
||||
commits = []
|
||||
real_git = W._sync_git
|
||||
|
||||
def git_spy(args, workdir, timeout):
|
||||
if str(args[0]) == "commit":
|
||||
commits.append(list(args))
|
||||
return real_git(args, workdir, timeout)
|
||||
|
||||
W._sync_git = git_spy
|
||||
ok, why = W._sync_finish_merge(pathlib.Path(tmpdir))
|
||||
check("marker guard blocks the commit", (not ok) and "marker" in why, f"{ok} {why}")
|
||||
check("no commit attempted", not commits, str(commits))
|
||||
W._sync_git = real_git
|
||||
W._sync_has_conflict_markers = lambda workdir: []
|
||||
W._sync_unmerged_files = lambda workdir: []
|
||||
|
||||
|
||||
def main():
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
test_upstream_status_parsing()
|
||||
@@ -557,6 +684,9 @@ def main():
|
||||
test_push_error_classification()
|
||||
test_sync_config_override()
|
||||
test_hermes_api_client(pathlib.Path(tmpdir))
|
||||
test_session_header_sent(tmpdir)
|
||||
test_salvage_on_agent_timeout(tmpdir)
|
||||
test_finish_merge_blocks_leftover_markers(tmpdir)
|
||||
print(f"\n{len(PASSED)} passed, {len(FAILED)} failed")
|
||||
if FAILED:
|
||||
print("failed: " + ", ".join(FAILED))
|
||||
|
||||
Reference in New Issue
Block a user