feat(worker): upstream fork auto-sync — pull+merge fork repos hourly

Merge new upstream (parent) commits into every fork in the App installation,
gated by a per-repo interval (default 1h), inside the existing 5-minute tick
(STEP 0, max 2 forks/tick, oldest-first).

- Conflicted merges are resolved by Claude Code (merge-reconciler rules:
  never wholesale --ours/--theirs, verify with the repo's own
  typecheck+tests, commit --no-edit; Claude never pushes — harness does).
- Clean merges get a single Claude Code quality pass commit.
- Push path: owner PAT (gh CLI) first — the App lacks workflows:write and a
  workflows-touching merge is rejected for the App token; App token fallback.
- Protected default branch: detected from the push result (GH006 /
  required-status-check) → upstream-sync-<ts> branch + PR through the normal
  pipeline; duplicate open sync PRs are skipped.
- CI safety: after a direct push, ticks verify the fork CI at our merge sha;
  red CI at OUR merge (still the tip) → sha-guarded force-revert to
  pre-merge sha + Discord notify; never reverts foreign commits.
- Discord: synced / PR opened / reverted / skipped-once on pr-agent-ops.
- merge_pr gains the same PAT fallback (a PR merge touching workflows is a
  workflow-file push).
- CLI: --sync-status, --sync-only <repo> [--dry].
- Tests: scripts/test_pr_queue_sync.py (46 assertions, monkeypatched, no
  network); py_compile clean.
- Plan: .hermes/plans/2026-09-21-upstream-auto-sync.md
This commit is contained in:
asepharyana
2026-09-21 17:04:39 +07:00
parent 1dc4d098fa
commit c48adea6c3
4 changed files with 1364 additions and 9 deletions
@@ -0,0 +1,101 @@
# Upstream Auto-Sync — Auto Pull + Merge for Fork Repos
Status: implemented in `scripts/pr-queue-worker.py` (2026-09-21).
User decisions: interval **1 hour**, sync PRs **do** get the Claude Code auto-fix.
## Goal
For every repo the GitHub App is installed on whose GitHub metadata says
`fork: true`, pull the commits the upstream (parent) repo has that the fork lacks
and MERGE them into the fork's default branch — with conflict resolution, CI
verification and deduped notifications. No new cron/service: the existing
`pr-queue-worker.py` tick (every 5 min) runs it, gated by a per-repo interval.
## Fork state at build time (live)
| Fork | Parent | To merge | Fork divergence |
|---|---|---|---|
| `asepharyana/shiro-neko` | `zakirkun/shiro-neko` | 4 | 26 |
| `asepharyana/hermes-agent-mission-control` | `sharbelxyz/hermes-agent-mission-control` | 0 | 8 |
Both default branches are unprotected; the real shiro-neko merge conflicts in 14
files (`src/tools.ts`, `src/ui/App.tsx`, `src/tests…`).
## Design decisions (with reasons)
1. **Merge, never rebase** — preserves fork history/divergence (skill
`fork-upstream-sync`).
2. **Inside `pr-queue-worker.py`** — the worker already owns the atomic cron
lock, installation tokens and Discord plumbing; the sync is STEP 0 of the same
tick, bounded to `max_per_tick = 2` forks, oldest-attempt-first so the budget
rotates fairly.
3. **1h interval per repo** — the compare API is checked every tick (cheap), a
merge attempt only when the interval elapsed AND the upstream tip changed
(`last_attempt_sha`), so a failed/skipped sync is not retried every 5 minutes.
4. **Conflicted merges go to Claude Code** — `resolve_conflicts: true`. The
prompt encodes the `fork-upstream-sync` rules (merge hunks by hand, never
wholesale `--ours/--theirs`, strip UTF-8 BOM, run the repo's typecheck+tests
before committing) plus the `merge-reconciler` impartiality contract (classify
every hunk: disjoint-intent / same-question-different-answer / superseded;
change nothing outside conflict markers; report every decision).
Claude never pushes — the harness pushes.
5. **Clean merges get a quality pass** — `ai_fix_after_merge: true` runs one
Claude Code pass over the merged files and commits
`fix: auto-fix code quality [skip ci]`. On the PR path the normal worker AI
fix applies (user decision), so sync PRs are treated like any other PR.
6. **Push path: owner PAT first, App token second** — the App has
`contents:write` but NOT `workflows:write`, and the shiro-neko merge touches
`.github/workflows/*`, so an App-token push is rejected. Clones use the App
token; pushes use the gh CLI PAT (`_fetch_gh_token()`).
7. **Protected branch → PR path** — App tokens get 403 (not 404) on the
branch-protection endpoint, so protection is detected from the PUSH result
(GH006 / "protected branch" / required-status-check markers) and falls back to
pushing a `upstream-sync-<ts>` branch + opening a PR that the normal pipeline
finishes. Duplicate PRs are avoided by scanning open PRs for that head prefix.
8. **CI verify with sha-guarded auto-revert** — `verify_ci: true`. After a direct
push, later ticks check the fork's CI at our merge sha; if it is still the tip
and CI failed, the branch is force-pushed back to `pre_merge_sha` and Discord
is notified. Never reverts when a human pushed on top; never applies to the
protected/PR path; watch stops after `verify_ci_max_age_h = 6h`.
9. **Discord**: sync events go to the same `pr-agent-ops` webhook as the run
reports (synced / PR opened / reverted / skipped-once).
10. **`merge_pr` PAT fallback** — a PR merge that changes workflow files is also
a workflow-file push, so a 403 on the App merge retries with the PAT.
## Scope
- `scripts/pr-queue-worker.py` — new `# Upstream Fork Auto-Sync` section (config,
discovery, compare, merge, conflict-resolve, push, PR path, verify/revert,
orchestrator) + STEP 0 in `main()` + CLI flags.
- `scripts/test_pr_queue_sync.py` — new: 46 assertions over the orchestration.
- State: `/tmp/pr-queue-sync-state.json`; workdirs `/tmp/pr-queue-sync-work/`.
- No changes to `server/` (Bun), systemd units or the cron list.
## State shape
```json
{"owner/fork": {
"last_sync_ts": 0, "last_attempt_sha": "", "last_merged_upstream_sha": "",
"skip_reason": "", "notified": false,
"pending_verify": {"sha": "", "pre_merge_sha": "", "branch": "main", "pushed_at": 0}
}}
```
## CLI (manual/dev)
```bash
python3 scripts/pr-queue-worker.py --sync-status
python3 scripts/pr-queue-worker.py --sync-only asepharyana/shiro-neko --dry
python3 scripts/pr-queue-worker.py --sync-only asepharyana/shiro-neko
```
`--dry` runs the whole flow (clone, merge, Claude conflict resolution) and stops
before pushing, leaving the prepared workdir in `/tmp/pr-queue-sync-work/<repo>`.
## Verification
- `python3 -m py_compile` clean.
- `python3 scripts/test_pr_queue_sync.py` → 46/46 (gating, clean merge, conflict
resolution, failed resolution skip+dedupe, protected→PR, CI green/red/foreign-
commit/running, fork discovery, push-error classification, config override).
- Live dry run against the real fork: clone + upstream fetch + conflicted merge +
Claude Code resolution, verified with the repo's own `bun run typecheck`/`bun test`.
- Live production: the 5-minute cron picks it up on the next tick after the push
(the cron wrapper ff-only pulls this repo and execs the repo copy).
- Regression guard for the bug the dry run caught: compare head ref must be
`{owner}:{branch}` — a full `owner/repo:branch` head 404s.
## Out of scope
- Syncing non-default branches; pushing fork→upstream; scheduled rebase option.
+37 -3
View File
@@ -10,11 +10,45 @@ repos where the PR-Agent GitHub App is installed and drives the full lifecycle:
of waiting on them forever of waiting on them forever
3. **AI auto-fix** → runs Claude Code (`-p`) on the PR head for up to 3. **AI auto-fix** → runs Claude Code (`-p`) on the PR head for up to
`AI_FIX_MAX_TURNS` turns, then pushes the fix `AI_FIX_MAX_TURNS` turns, then pushes the fix
4. **Safety analysis** → parses the review body for security/major-issue 5. **Safety analysis** → parses the review body for security/major-issue
blockers; score must be ≥ 6/10 blockers; score must be ≥ 6/10
5. **CI gate** → waits for the required check to pass (closes stale dependabot 6. **CI gate** → waits for the required check to pass (closes stale dependabot
PRs stuck failing CI for > `STALE_CI_CLOSE_DAYS`) PRs stuck failing CI for > `STALE_CI_CLOSE_DAYS`)
6. **Approve + merge** 7. **Approve + merge**
## Upstream Fork Auto-Sync
Since 2026-09-21 the same 5-minute tick also syncs every repo in the App
installation whose GitHub metadata says `fork: true`: new upstream (parent)
commits are **merged** (never rebased) into the fork's default branch, gated by
a per-repo interval (default **1 hour**; `UPSTREAM_SYNC` config block).
- **Conflicted merge** → Claude Code resolves it (merge-reconciler rules: merge
hunks by hand, never wholesale `--ours/--theirs`, run the repo's own
typecheck/tests before committing). Claude never pushes — the harness does.
- **Clean merge** → one Claude Code quality pass over the merged files,
committed as `fix: auto-fix code quality [skip ci]`.
- **Protected default branch** → detect from the push result (GH006 /
required-status-check) and fall back to opening an `upstream-sync-*` PR that
the normal pipeline (review → AI fix → CI → approve → merge) finishes.
- **CI safety** → after a direct push the worker verifies the fork's CI at our
merge commit; a red CI at OUR merge sha (still the tip, no human commits on
top) force-reverts to the pre-merge sha. Watch stops after 6 h.
- **Push credentials** → owner PAT (gh CLI) first because the App lacks
`workflows:write`; App token is the fallback. Clones use the App token.
- **State** → `/tmp/pr-queue-sync-state.json` (skip/interval/pending-verify),
workdirs `/tmp/pr-queue-sync-work/`.
- **Notifications** → same `pr-agent-ops` Discord webhook: synced, PR opened,
reverted, skipped-once.
Manual/dev:
```bash
python3 scripts/pr-queue-worker.py --sync-status
python3 scripts/pr-queue-worker.py --sync-only asepharyana/shiro-neko --dry # stops before push
python3 scripts/pr-queue-worker.py --sync-only asepharyana/shiro-neko
```
Tests: `python3 scripts/test_pr_queue_sync.py` (46 assertions; no network —
gh_api/git/Claude/push are monkeypatched).
## Deployment ## Deployment
+774 -6
View File
@@ -454,8 +454,15 @@ def approve_pr(token, repo_full, pr_num):
return status return status
def merge_pr(token, repo_full, pr_num, sha): def merge_pr(token, repo_full, pr_num, sha):
status, data = gh_api("PUT", f"/repos/{repo_full}/pulls/{pr_num}/merge", token=token, payload = {"commit_title": f"Auto-merge PR #{pr_num}", "merge_method": "merge", "sha": sha}
json_data={"commit_title": f"Auto-merge PR #{pr_num}", "merge_method": "merge", "sha": sha}) status, data = gh_api("PUT", f"/repos/{repo_full}/pulls/{pr_num}/merge", token=token, json_data=payload)
if status == 403:
# GitHub Apps without the `workflows` permission cannot merge a PR that
# changes .github/workflows/* — the merge IS a push of those files. Retry
# as the repo owner (gh CLI PAT), which can, and report that result.
pat = _fetch_gh_token()
if pat:
return gh_api("PUT", f"/repos/{repo_full}/pulls/{pr_num}/merge", token=pat, json_data=payload)
return status, data return status, data
def post_discord_notification(repo_full, pr_num, status, summary="", score="", url=""): def post_discord_notification(repo_full, pr_num, status, summary="", score="", url=""):
@@ -978,6 +985,733 @@ def run_ai_fix(repo_full, pr_num, title, head_sha, head_ref, base_ref, token):
snippet = claude_output[:300].replace("\n", " ") snippet = claude_output[:300].replace("\n", " ")
return False, "Claude ran but no push. Output: " + snippet return False, "Claude ran but no push. Output: " + snippet
# ══════════════════════════════════════════════════════════════════════════
# Upstream Fork Auto-Sync (2026-09-21)
# ══════════════════════════════════════════════════════════════════════════
# For every repo the GitHub App is installed on whose metadata says `fork: true`,
# pull the commits the upstream (parent) repo has that the fork lacks and MERGE
# them into the fork's default branch — MERGE, never rebase, so the fork keeps
# its local divergence intact (skill: fork-upstream-sync).
#
# Why it lives in this worker: the worker already holds the atomic cron lock, has
# installation-token + Discord plumbing and runs every 5 minutes. The sync itself
# is gated by a per-repo interval (default 1h — user decision 2026-09-21).
#
# GitHub facts this code depends on (probed live 2026-09-21, see findings):
# • /repos/{fork}/compare/{local}...{parent}:{upstream}
# ahead_by = upstream commits MISSING from the fork → what we merge
# behind_by = fork-only commits (divergence) → never discarded
# Verified against `git rev-list --count` on a real clone.
# • The App has contents:write but NOT workflows:write. A merge that touches
# .github/workflows/* can therefore only be pushed with the owner PAT
# (gh CLI token, _fetch_gh_token()). Clones still use the App token.
# • App tokens get 403 (not 404) on the branch-protection endpoint, so a
# protected branch is detected from the PUSH result and falls back to a PR.
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
UPSTREAM_SYNC = {
"enabled": os.environ.get("PR_AGENT_UPSTREAM_SYNC", "1") != "0",
"interval_h": 1.0, # check upstream hourly (user decision 2026-09-21)
"max_per_tick": 2, # bound work per 5-minute tick
"resolve_conflicts": True, # hand a conflicted merge to Claude Code
"ai_fix_after_merge": True, # Claude Code quality pass on a clean merge
"verify_ci": True, # revert our own merge commit if the fork CI fails
"verify_ci_max_age_h": 6.0, # stop watching (keep the merge) after this
"no_ci_grace_s": 600, # wait this long before concluding "no CI here"
"repos": {}, # per-repo overrides, e.g.
# {"owner/fork": {"interval_h": 6, "verify_ci": False,
# "branches": {"local": "main", "upstream": "main"}}}
}
def load_sync_state():
"""Sync bookkeeping per fork: {repo_full: {last_sync_ts, last_attempt_sha,
last_merged_upstream_sha, skip_reason, notified, pending_verify}}."""
if SYNC_STATE_FILE.exists():
try:
return json.loads(SYNC_STATE_FILE.read_text())
except (json.JSONDecodeError, OSError):
return {}
return {}
def save_sync_state(state):
try:
SYNC_STATE_FILE.write_text(json.dumps(state, indent=1))
except OSError:
pass
def _sync_entry(state, repo_full):
return state.setdefault(repo_full, {})
def sync_config(repo_full):
"""Worker-wide sync defaults merged with the per-repo override."""
cfg = {k: v for k, v in UPSTREAM_SYNC.items() if k != "repos"}
cfg.update(UPSTREAM_SYNC["repos"].get(repo_full) or {})
return cfg
def post_sync_discord(title, lines, color=0x5865F2):
"""Post a fork-sync embed to the pr-agent-ops webhook (same channel as the
run reports). Best-effort: never raises, never blocks a tick."""
try:
cfg = json.loads((pathlib.Path.home() / ".hermes/.ops-webhooks.json").read_text())
url = cfg.get("pr-agent-ops")
if not url:
return False
import httpx
with httpx.Client(timeout=15) as client:
r = client.post(url, json={
"username": "PR-Agent Ops",
"embeds": [{"title": title, "description": "\n".join(lines)[:4000], "color": color}],
})
return r.status_code in (200, 204)
except Exception:
return False
def _sync_git(args, cwd=None, timeout=180):
"""Run git, never raising: timeouts become exit 124, missing git exit 127."""
cmd = ["git"] + [str(a) for a in args]
try:
return subprocess.run(
cmd, cwd=str(cwd) if cwd else None,
capture_output=True, text=True, timeout=timeout,
)
except subprocess.TimeoutExpired as exc:
return subprocess.CompletedProcess(cmd, 124, "", f"timeout after {timeout}s: {exc}")
except FileNotFoundError:
return subprocess.CompletedProcess(cmd, 127, "", "git not found")
def list_fork_repos():
"""[(app_token, fork_full_name, parent_full_name, default_branch)] for every
fork in the App installation. The repo metadata lookup is per repo because
the parent block is authoritative on the repo object."""
out = []
_, installs = gh_api("GET", "/app/installations")
for inst in installs if isinstance(installs, list) else []:
token = get_installation_token(inst["id"])
if not token:
continue
_, repos = gh_api("GET", "/installation/repositories?per_page=100", token=token)
for r in (repos.get("repositories", []) if isinstance(repos, dict) else []):
full = r.get("full_name", "")
if not full or not r.get("fork"):
continue
status, meta = gh_api("GET", f"/repos/{full}", token=token)
parent = (meta.get("parent") or {}) if isinstance(meta, dict) else {}
if status != 200 or not parent.get("full_name"):
continue
out.append((token, full, parent["full_name"], meta.get("default_branch") or "main"))
return out
def upstream_status(token, fork, parent, local_branch, upstream_branch):
"""(merge_count, fork_divergence, upstream_tip_sha) for the fork branch vs the
upstream branch, or None when the comparison is unavailable.
merge_count = commits upstream has that the fork lacks (ahead_by on
compare/{local}...{parent}:{upstream}) — i.e. what a sync would merge.
NOTE the cross-repo compare syntax is `{owner}:{branch}` — passing the full
`owner/repo` yields a 404 (verified live 2026-09-21), so only the owner part
of the parent's full name goes into the head ref."""
owner = parent.split("/")[0] if "/" in parent else parent
path = f"/repos/{fork}/compare/{local_branch}...{owner}:{upstream_branch}"
status, data = gh_api("GET", path, token=token)
if status != 200:
pat = _fetch_gh_token()
if pat:
status, data = gh_api("GET", path, token=pat)
if status != 200 or not isinstance(data, dict):
return None
commits = data.get("commits") or []
tip = commits[-1].get("sha", "") if commits else ""
if not tip:
s2, d2 = gh_api("GET", f"/repos/{parent}/commits/{upstream_branch}", token=token)
if s2 == 200 and isinstance(d2, dict):
tip = d2.get("sha", "")
return int(data.get("ahead_by", 0)), int(data.get("behind_by", 0)), tip
def _sync_open_pr(token, fork, base_branch):
"""Number of an already-open upstream-sync PR targeting base_branch, else 0."""
_, prs = gh_api("GET", f"/repos/{fork}/pulls?state=open&per_page=50", token=token)
for pr in (prs if isinstance(prs, list) else []):
head_ref = str((pr.get("head") or {}).get("ref", ""))
if head_ref.startswith(SYNC_PR_PREFIX) and (pr.get("base") or {}).get("ref") == base_branch:
return pr.get("number", 0)
return 0
def _sync_push_urls(fork, app_token):
"""Push credentials, best first. The owner PAT comes first on purpose: the
App lacks `workflows` permission, so a merge touching .github/workflows/*
(very common when syncing) is rejected for the App token."""
urls = []
pat = _fetch_gh_token()
if pat:
urls.append((f"https://x-access-token:{pat}@github.com/{fork}.git", "pat"))
if app_token:
urls.append((f"https://x-access-token:{app_token}@github.com/{fork}.git", "app"))
return urls
def _sync_fetch_url(parent):
"""Read credentials for the upstream fetch — the PAT when available (private
upstreams + rate limits), plain HTTPS otherwise (public read needs no auth)."""
pat = _fetch_gh_token()
if pat:
return f"https://x-access-token:{pat}@github.com/{parent}.git"
return f"https://github.com/{parent}.git"
_PROTECTED_PUSH_MARKERS = (
"protected branch",
"gh006",
"required status check",
"branch protection",
"protected_branch",
)
_WORKFLOW_PUSH_MARKERS = (
"workflows permission",
"workflows` permission",
"create or update workflow",
)
def _is_protected_push_error(err):
low = (err or "").lower()
return any(m in low for m in _PROTECTED_PUSH_MARKERS)
def _is_workflow_push_error(err):
low = (err or "").lower()
return any(m in low for m in _WORKFLOW_PUSH_MARKERS)
def _push_ref(workdir, fork, source, dest, app_token, force=False):
"""Push `source` to `dest` on the fork. Returns (ok, detail, protected)."""
refspec = ("+" if force else "") + f"{source}:{dest}"
last = "no push credentials available"
for url, kind in _sync_push_urls(fork, app_token):
r = _sync_git(["push", url, refspec], workdir, 300)
if r.returncode == 0:
return True, f"{kind} push ok", False
err = ((r.stderr or "") + (r.stdout or "")).strip()
last = f"{kind}: {err[-260:]}"
if _is_protected_push_error(err):
return False, last, True
if _is_workflow_push_error(err):
# PAT missing/insufficient: the App cannot push workflow changes.
return False, f"needs the `workflows` App permission or the gh PAT ({last})", False
return False, last, False
def _run_claude_sync(workdir, prompt, label, fork):
"""Run Claude Code inside the sync workdir. Returns (ok, snippet).
Mirrors run_ai_fix's invocation (same binary resolution + provider env) but
with its own log names, a longer timeout (conflict resolution runs a full
test suite) and no push expectations — the harness pushes."""
claude_bin = shutil.which("claude") or CLAUDE_BIN
if not os.path.isfile(claude_bin):
return False, "[INFRA] Claude Code CLI not found at " + str(claude_bin)
try:
(workdir / (label + ".prompt.txt")).write_text(prompt)
except OSError:
pass
env = {k: v for k, v in _claude_env().items() if k in (
"PATH", "HOME", "HERMES_HOME",
"ANTHROPIC_API_KEY", "ANTHROPIC_BASE_URL", "ANTHROPIC_URL",
)}
out_path = workdir / (label + ".out.log")
try:
with open(out_path, "w") as lf:
res = subprocess.run(
[claude_bin, "-p", prompt,
"--allowedTools", "Read,Edit,Bash,Write",
"--max-turns", str(AI_FIX_MAX_TURNS)],
cwd=str(workdir), stdout=lf, stderr=subprocess.STDOUT,
text=True, env=env, timeout=SYNC_CLAUDE_TIMEOUT,
)
saved = out_path.read_text() if out_path.exists() else ""
except subprocess.TimeoutExpired:
return False, f"[INFRA] Claude Code timed out after {SYNC_CLAUDE_TIMEOUT}s"
except FileNotFoundError:
return False, "[INFRA] Claude Code CLI not found"
except OSError as exc:
return False, f"[INFRA] Claude Code could not run: {exc}"
snippet = saved[-400:].replace("\n", " ") if saved else ""
if res.returncode != 0:
return False, f"claude exited {res.returncode}: {snippet}"
return True, snippet
def _sync_unmerged_files(workdir):
r = _sync_git(["diff", "--name-only", "--diff-filter=U"], workdir, 60)
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])}"
_sync_git(["add", "-A"], workdir, 120)
r = _sync_git(["commit", "--no-edit"], workdir, 120)
if r.returncode != 0:
combined = ((r.stdout or "") + (r.stderr or "")).lower()
if "nothing to commit" in combined:
return True, ""
return False, ((r.stderr or r.stdout) or "").strip()[:200]
return True, ""
def _sync_commit_if_dirty(workdir, message):
"""Commit pending edits made by a quality pass. Returns True if a commit was
created (Claude normally commits itself; this is the safety net)."""
st = _sync_git(["status", "--porcelain"], workdir, 60)
if not (st.stdout or "").strip():
return False
_sync_git(["add", "-A"], workdir, 120)
_sync_git(["commit", "--message", message], workdir, 120)
return True
def _conflict_prompt(fork, parent, upstream_branch, local_branch, conflicted):
"""Prompt for resolving the halted upstream→fork merge.
Encodes the fork-upstream-sync rules (merge both sides by hand, never
wholesale --ours/--theirs, strip BOM, verify before committing) plus the
merge-reconciler impartiality contract (classify each hunk, no drive-by
edits, report every decision)."""
files = "\n".join(" - " + f for f in conflicted[:40])
more = "" if len(conflicted) <= 40 else f"\n ... and {len(conflicted) - 40} more"
return (
"A `git merge` is IN PROGRESS inside this repository and stopped with conflicts.\n"
f"Direction: {parent} (branch {upstream_branch}) → {fork} (branch {local_branch}).\n"
"You are the neutral reconciler: neither side may be dropped.\n\n"
f"Conflicted files:\n{files}{more}\n\n"
"TASK: resolve every conflict, then finish the merge.\n\n"
"CLASSIFY EVERY HUNK before editing, then resolve by class:\n"
" • disjoint-intent — the two changes serve different goals → keep BOTH.\n"
" • same-question-different-answer — both sides answered one question\n"
" differently → pick the one matching the fork's stated intent and note\n"
" the decision; never invent a hybrid nobody asked for.\n"
" • superseded — one side's premise no longer holds after the other change\n"
" → keep the surviving side and record why.\n\n"
"RULES (non-negotiable):\n"
"1. Merge conflicting hunks by hand. NEVER `git checkout --ours/--theirs`\n"
" wholesale, never `git rebase`, never delete a side without saying why.\n"
"2. The fork carries local features upstream does not know about: every\n"
" fork-only feature must still work after the merge.\n"
"3. Change NOTHING outside conflict markers — no reformatting, no renames,\n"
" no opportunistic fixes.\n"
"4. If a file starts with a UTF-8 BOM (bytes EF BB BF), strip it.\n"
"5. VERIFY before committing: run the repository's own checks when they\n"
" exist (package.json scripts: `bun run typecheck`, `bun test`; else\n"
" `npm test`, `cargo test`, `pytest -q`). Fix what YOUR resolution broke\n"
" until they pass.\n"
"6. Complete the merge: `git add -A && git commit --no-edit`\n"
"7. Do NOT push — the harness pushes after you finish.\n\n"
"Finish with a report listing, per file, the hunks you resolved and which\n"
"side(s) you kept, the verification commands you ran, and their results."
)
def _quality_prompt(fork, parent, upstream_branch, files):
file_list = "\n".join(" - " + f for f in files[:30])
if len(files) > 30:
file_list += f"\n ... and {len(files) - 30} more"
return (
f"The fork {fork} just merged {parent} (branch {upstream_branch}) into its\n"
"default branch. The merge itself is already committed and correct.\n\n"
f"Files the merge changed:\n{file_list}\n\n"
"TASK: improve code quality of ONLY these files — naming, DRY, error\n"
"handling, missing types, docstrings, clear anti-patterns.\n\n"
"RULES:\n"
"- Behavior must stay identical. Do NOT add features or change logic.\n"
"- Do NOT rewrite the upstream architecture; this is a fresh merge.\n"
"- Run the repository's checks when they exist (`bun run typecheck`,\n"
" `bun test`, else `npm test`/`cargo test`/`pytest -q`) and keep them green.\n"
"- Commit exactly one commit: `git add -A && git commit --message=\"fix: auto-fix code quality [skip ci]\"`\n"
"- Do NOT push — the harness pushes after you finish."
)
def _sync_revert_merge(app_token, fork, branch, pre_merge_sha, reason):
"""Force the fork branch back to the pre-merge commit. Only ever called when
the branch tip IS our own merge commit (sha-guarded by the caller)."""
workdir = SYNC_TMP_BASE / (fork.replace("/", "_") + "_revert")
if workdir.exists():
subprocess.run(["rm", "-rf", str(workdir)], timeout=60)
workdir.parent.mkdir(parents=True, exist_ok=True)
clone_url = f"https://x-access-token:{app_token}@github.com/{fork}.git" if app_token else f"https://github.com/{fork}.git"
r = _sync_git(["clone", clone_url, str(workdir), "--branch", branch], None, 300)
if r.returncode != 0:
return False, f"revert clone failed: {(r.stderr or '').strip()[:160]}"
try:
have = _sync_git(["cat-file", "-e", pre_merge_sha + "^{commit}"], workdir, 60)
if have.returncode != 0:
return False, f"pre-merge commit {pre_merge_sha[:8]} not found in clone"
ok, detail, _protected = _push_ref(
workdir, fork, pre_merge_sha, f"refs/heads/{branch}", app_token, force=True)
if ok:
return True, f"reverted {branch} to {pre_merge_sha[:8]} ({reason})"
return False, f"revert push failed: {detail[:200]}"
finally:
subprocess.run(["rm", "-rf", str(workdir)], timeout=60)
def verify_pending_syncs(state, token_by_repo, dry=False):
"""CI-verify merges we pushed directly. A failed fork CI at OUR merge commit
(while it is still the tip) reverts the merge instead of leaving a red main.
Returns report lines."""
lines = []
for fork in list(state.keys()):
entry = state.get(fork) or {}
pending = entry.get("pending_verify") or {}
if not pending.get("sha"):
continue
token = token_by_repo.get(fork, "")
if not token:
continue
cfg = sync_config(fork)
branch = pending.get("branch") or "main"
if not cfg.get("verify_ci", True) or dry:
entry["pending_verify"] = None
continue
age = time.time() - float(pending.get("pushed_at") or 0)
if age > cfg["verify_ci_max_age_h"] * 3600:
lines.append(f" ⌛ {fork}: merge {pending['sha'][:8]} unverified for "
f"{age / 3600:.1f}h — keeping it, stopping the watch")
entry["pending_verify"] = None
save_sync_state(state)
continue
s, data = gh_api("GET", f"/repos/{fork}/commits/{branch}?per_page=1", token=token)
tip = data[0]["sha"] if isinstance(data, list) and data else ""
if tip and tip != pending["sha"]:
lines.append(f" ✓ {fork}: {branch} moved past our merge "
f"({pending['sha'][:8]} → {tip[:8]}) — nothing to verify")
entry["pending_verify"] = None
save_sync_state(state)
continue
s, data = gh_api("GET", f"/repos/{fork}/commits/{pending['sha']}/check-runs", token=token)
checks = data.get("check_runs", []) if isinstance(data, dict) else []
failed = [c for c in checks if c.get("conclusion") == "failure"]
running = [c for c in checks if c.get("status") != "completed"]
if failed:
names = ", ".join(c.get("name", "?") for c in failed[:3])
ok, detail = _sync_revert_merge(token, fork, branch, pending.get("pre_merge_sha", ""), f"CI failed: {names}")
lines.append(f" ↩️ {fork}: {detail}")
post_sync_discord(
f"↩️ Fork sync reverted: {fork}",
[f"CI failed at our merge commit `{pending['sha'][:8]}` ({names}).",
detail,
f"https://github.com/{fork}/commits/{branch}"],
0xE74C3C,
)
entry["pending_verify"] = None
entry["last_merged_upstream_sha"] = ""
save_sync_state(state)
continue
if not checks:
if age < cfg["no_ci_grace_s"]:
continue # check-runs may not have registered yet
entry["pending_verify"] = None
save_sync_state(state)
continue
if running:
continue # still running — verify on a later tick
lines.append(f" ✅ {fork}: merge {pending['sha'][:8]} verified green "
f"({len(checks)} check(s))")
entry["pending_verify"] = None
save_sync_state(state)
return lines
def sync_fork_repo(token, fork, parent, local_branch, upstream_branch, upstream_sha,
merge_count, divergence, state, cfg, dry=False):
"""One sync attempt for one fork: clone → merge upstream → Claude Code when
needed → push (or open a PR when the branch is protected). Returns
(status, detail) with status ∈ synced|pr-opened|conflict-failed|push-failed|error|dry."""
entry = _sync_entry(state, fork)
workdir = SYNC_TMP_BASE / fork.replace("/", "_")
try:
if workdir.exists():
subprocess.run(["rm", "-rf", str(workdir)], timeout=60)
workdir.parent.mkdir(parents=True, exist_ok=True)
clone_url = f"https://x-access-token:{token}@github.com/{fork}.git"
r = _sync_git(["clone", clone_url, str(workdir), "--branch", local_branch], None, 300)
if r.returncode != 0:
entry["last_sync_ts"] = time.time()
save_sync_state(state)
return "error", f"clone failed: {(r.stderr or '').strip()[:200]}"
_sync_git(["config", "user.name", "mytheclipsebotreview"], workdir, 30)
_sync_git(["config", "user.email", "bot@users.noreply.github.com"], workdir, 30)
pre_merge_sha = (_sync_git(["rev-parse", "HEAD"], workdir, 30).stdout or "").strip()
upstream_ref = f"refs/remotes/upstream/{upstream_branch}"
r = _sync_git(["fetch", "--no-tags", _sync_fetch_url(parent),
f"{upstream_branch}:{upstream_ref}"], workdir, 300)
if r.returncode != 0:
entry["last_sync_ts"] = time.time()
save_sync_state(state)
return "error", f"upstream fetch failed: {(r.stderr or '').strip()[:200]}"
r = _sync_git(["merge", upstream_ref, "--no-edit"], workdir, 300)
conflicted = _sync_unmerged_files(workdir)
if r.returncode != 0 and not conflicted:
entry["last_sync_ts"] = time.time()
save_sync_state(state)
return "error", f"merge error: {((r.stderr or '') + (r.stdout or '')).strip()[:200]}"
if conflicted:
if not cfg.get("resolve_conflicts", True):
_sync_git(["merge", "--abort"], workdir, 60)
note = (f"{len(conflicted)} conflicting file(s) and conflict resolution is "
f"disabled: {', '.join(conflicted[:5])}")
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
BUFFER.append(f" 🤖 resolving {len(conflicted)} conflict(s) with Claude Code "
f"(up to {SYNC_CLAUDE_TIMEOUT}s)...")
ok, snippet = _run_claude_sync(
workdir, _conflict_prompt(fork, parent, upstream_branch, local_branch, conflicted),
"claude_sync_conflicts", fork)
if not ok:
_sync_git(["merge", "--abort"], workdir, 60)
note = f"conflict resolution failed — {snippet[:200]}"
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}"
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"Claude Code resolved {len(conflicted)} conflict(s)"
else:
resolution = "clean merge"
if cfg.get("ai_fix_after_merge", True):
diff = _sync_git(["diff", "--name-only", f"{pre_merge_sha}..HEAD"], workdir, 120)
merged_files = [f for f in (diff.stdout or "").split("\n") if f.strip()]
if merged_files:
BUFFER.append(f" 🤖 Claude Code quality pass on {len(merged_files)} merged file(s)...")
ok, snippet = _run_claude_sync(
workdir, _quality_prompt(fork, parent, upstream_branch, merged_files),
"claude_sync_quality", fork)
if ok:
if _sync_commit_if_dirty(workdir, "fix: auto-fix code quality [skip ci]"):
resolution += " + Claude Code quality pass committed"
else:
resolution += " + quality pass made no changes"
else:
resolution += f" (quality pass skipped: {snippet[:120]})"
still_unmerged = _sync_unmerged_files(workdir)
if still_unmerged:
note = f"unmerged paths remain: {', '.join(still_unmerged[:5])}"
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
merged_sha = (_sync_git(["rev-parse", "HEAD"], workdir, 30).stdout or "").strip()
if dry:
return "dry", (f"prepared in {workdir} — pre_merge {pre_merge_sha[:8]}, "
f"upstream {upstream_sha[:8]}, head {merged_sha[:8]}, {resolution}")
ok, detail, protected = _push_ref(
workdir, fork, "HEAD", f"refs/heads/{local_branch}", token)
if ok:
entry.update({
"last_sync_ts": time.time(),
"last_attempt_sha": upstream_sha,
"last_merged_upstream_sha": upstream_sha,
"skip_reason": "",
"notified": False,
"pending_verify": {
"sha": merged_sha, "pre_merge_sha": pre_merge_sha,
"branch": local_branch, "pushed_at": time.time(),
},
})
save_sync_state(state)
return "synced", (f"{merge_count} upstream commit(s) merged into {local_branch} "
f"({resolution}); head {merged_sha[:8]}; {detail}")
if protected:
pr_num = open_sync_pr(token, workdir, fork, parent, local_branch, upstream_branch,
upstream_sha, merge_count, divergence, resolution)
if pr_num:
entry.update({
"last_sync_ts": time.time(),
"last_attempt_sha": upstream_sha,
"last_merged_upstream_sha": "",
"skip_reason": "",
"notified": False,
})
save_sync_state(state)
return "pr-opened", (f"{local_branch} is protected — opened PR #{pr_num} "
f"({merge_count} upstream commit(s), {resolution})")
entry["last_sync_ts"] = time.time()
save_sync_state(state)
return "push-failed", f"protected branch and PR creation failed: {detail[:200]}"
entry["last_sync_ts"] = time.time()
entry["last_attempt_sha"] = upstream_sha
entry["skip_reason"] = detail[:200]
entry["notified"] = False
save_sync_state(state)
return "push-failed", detail[:300]
except Exception as exc: # a sync must never take the whole tick down
return "error", f"{type(exc).__name__}: {exc}"
finally:
if not dry:
subprocess.run(["rm", "-rf", str(workdir)], timeout=60)
def open_sync_pr(token, workdir, fork, parent, base_branch, upstream_branch,
upstream_sha, merge_count, divergence, resolution):
"""Push the merged workdir as a `upstream-sync-<ts>` branch and open a PR
against base_branch. The normal worker pipeline (review → CI → approve →
merge) finishes the job. Returns the PR number, or 0 on failure."""
sync_branch = SYNC_PR_PREFIX + time.strftime("%Y%m%d-%H%M%S")
r = _sync_git(["checkout", "-b", sync_branch], workdir, 60)
if r.returncode != 0:
return 0
ok, _detail, _protected = _push_ref(workdir, fork, "HEAD", f"refs/heads/{sync_branch}", token)
if not ok:
return 0
body = (
f"⬆️ Automated upstream sync from `{parent}` (branch `{upstream_branch}`).\n\n"
f"- upstream commits merged: **{merge_count}**\n"
f"- fork-only commits preserved: **{divergence}**\n"
f"- upstream tip: `{upstream_sha}`\n"
f"- merge: {resolution}\n\n"
f"Opened by pr-queue-worker because `{base_branch}` is a protected branch, so the\n"
"merge goes through the normal pipeline (PR-Agent review → AI fix → CI → approve → merge).\n\n"
f"Compare: https://github.com/{fork}/compare/{base_branch}...{parent}:{upstream_branch}"
)
status, data = gh_api("POST", f"/repos/{fork}/pulls", token=token, json_data={
"title": f"⬆️ upstream-sync: merge {parent}@{upstream_sha[:8]} into {base_branch}",
"head": sync_branch,
"base": base_branch,
"body": body,
})
if status in (200, 201) and isinstance(data, dict):
return data.get("number", 0)
return 0
def run_upstream_sync(only=None, dry=False):
"""Sync every fork in the installation (bounded per tick), verify merges we
pushed earlier, and return the report lines. Never raises."""
lines = []
if not UPSTREAM_SYNC["enabled"] and not only:
return lines
try:
state = load_sync_state()
forks = list_fork_repos()
if only:
forks = [f for f in forks if f[1] == only]
token_by_repo = {f[1]: f[0] for f in forks}
lines += verify_pending_syncs(state, token_by_repo, dry)
now = time.time()
# oldest attempt first so the per-tick budget rotates fairly across forks
forks.sort(key=lambda f: _sync_entry(state, f[1]).get("last_sync_ts") or 0)
budget = int(UPSTREAM_SYNC["max_per_tick"])
for token, fork, parent, branch in forks:
cfg = sync_config(fork)
if not cfg.get("enabled", True):
continue
if budget <= 0 and not only:
break
branches = cfg.get("branches") or {}
s, meta = gh_api("GET", f"/repos/{fork}", token=token)
parent_meta = (meta.get("parent") or {}) if isinstance(meta, dict) else {}
local_branch = branches.get("local") or branch
upstream_branch = (branches.get("upstream")
or parent_meta.get("default_branch") or local_branch)
info = upstream_status(token, fork, parent, local_branch, upstream_branch)
if info is None:
lines.append(f"🔁 {fork}: upstream comparison unavailable — will retry")
continue
merge_count, divergence, upstream_sha = info
entry = _sync_entry(state, fork)
if merge_count <= 0:
if entry.get("skip_reason"):
entry["skip_reason"] = ""
save_sync_state(state)
continue # in sync — silent
if entry.get("last_attempt_sha") == upstream_sha:
continue # this upstream tip was already handled (synced/skipped)
if now - float(entry.get("last_sync_ts") or 0) < cfg["interval_h"] * 3600:
continue # interval not due yet — silent
open_pr = _sync_open_pr(token, fork, local_branch)
if open_pr:
lines.append(f"🔁 {fork}: upstream-sync PR #{open_pr} already open — waiting")
continue
lines.append(f"🔁 {fork}: `{parent}` has {merge_count} commit(s) the fork lacks "
f"(fork divergence {divergence}) — syncing into {local_branch}...")
res, detail = sync_fork_repo(
token, fork, parent, local_branch, upstream_branch, upstream_sha,
merge_count, divergence, state, cfg, dry)
entry = _sync_entry(state, fork)
if res == "synced":
merged_sha = (entry.get("pending_verify") or {}).get("sha", "")[:8]
lines.append(f" ✅ merged + pushed: {detail}")
post_sync_discord(
f"🔁 Fork synced: {fork}",
[f"⬆️ {merge_count} upstream commit(s) from `{parent}` merged into `{local_branch}`.",
f"merge head `{merged_sha}` (CI-verified on the next ticks; reverted automatically if red).",
f"https://github.com/{fork}"],
)
budget -= 1
elif res == "pr-opened":
lines.append(f" 📬 {detail}")
post_sync_discord(
f"📬 Fork sync PR opened: {fork}",
[detail, f"https://github.com/{fork}/pulls"],
)
budget -= 1
elif res == "dry":
lines.append(f" 🧪 dry run: {detail}")
budget -= 1
elif res == "conflict-failed":
lines.append(f" ⏭️ {detail}")
if not entry.get("notified"):
entry["notified"] = True
save_sync_state(state)
post_sync_discord(
f"⏭️ Fork sync skipped: {fork}",
[f"❗ {detail}",
f"upstream `{parent}@{upstream_sha[:8]}` — retried when upstream moves or the state file is cleared.",
f"https://github.com/{fork}"],
0xE67E22,
)
budget -= 1
else:
lines.append(f" ⚠️ {res}: {detail}")
budget -= 1
save_sync_state(state)
except Exception as exc:
lines.append(f"⚠️ Upstream sync error: {type(exc).__name__}: {exc}")
return lines
# ── Main ── # ── Main ──
def main(): def main():
start = time.time() start = time.time()
@@ -987,15 +1721,23 @@ def main():
return return
try: try:
# ── STEP 0: Upstream fork auto-sync (hourly per fork, bounded per tick) ──
# Runs before the PR loop so a fork's default branch is refreshed while
# the same tick still processes PRs. Never raises (returns report lines).
sync_lines = run_upstream_sync()
all_prs = gather_open_prs() all_prs = gather_open_prs()
if not all_prs: if not all_prs and not sync_lines:
return # truly silent return # truly silent
BUFFER.append(f"🔍 PR Queue Worker — {time.ctime()}") BUFFER.append(f"🔍 PR Queue Worker — {time.ctime()}")
BUFFER.append(f"{'='*50}") BUFFER.append(f"{'='*50}")
BUFFER.append(f"📋 Found {len(all_prs)} open PR(s) to process") if sync_lines:
if AI_FIX_ENABLED: BUFFER.extend(sync_lines)
BUFFER.append(f" ✨ AI auto-fix: ENABLED (Claude Code)") if all_prs:
BUFFER.append(f"📋 Found {len(all_prs)} open PR(s) to process")
if AI_FIX_ENABLED:
BUFFER.append(f" ✨ AI auto-fix: ENABLED (Claude Code)")
merged_count = 0 merged_count = 0
triggered_count = 0 triggered_count = 0
@@ -1337,4 +2079,30 @@ if __name__ == "__main__":
raise SystemExit(128 + signum) raise SystemExit(128 + signum)
_signal.signal(_signal.SIGTERM, _term_handler) _signal.signal(_signal.SIGTERM, _term_handler)
# ── Manual/dev entrypoints for the upstream fork sync ──
# python3 pr-queue-worker.py --sync-status
# python3 pr-queue-worker.py --sync-only owner/fork [--dry]
# --dry prepares the merge (including the Claude Code conflict resolution)
# in /tmp/pr-queue-sync-work/<repo> and stops before pushing.
_argv = sys.argv[1:]
if "--sync-status" in _argv:
print(json.dumps(load_sync_state(), indent=2))
sys.exit(0)
if "--sync-only" in _argv:
_idx = _argv.index("--sync-only")
_target = (_argv[_idx + 1]
if len(_argv) > _idx + 1 and not _argv[_idx + 1].startswith("-")
else None)
_dry = "--dry" in _argv
if not get_lock():
print("another pr-queue-worker run holds the lock — try again shortly")
sys.exit(1)
try:
_report = run_upstream_sync(only=_target, dry=_dry)
finally:
release_lock()
print("\n".join(_report) if _report else "(nothing to sync)")
sys.exit(0)
main() main()
+452
View File
@@ -0,0 +1,452 @@
#!/usr/bin/env python3
"""
Unit tests for the pr-queue-worker upstream fork auto-sync section.
Run: python3 scripts/test_pr_queue_sync.py
No network, no GitHub token, no real git: `gh_api`, `_sync_git`, the Claude
runner and the push helper are monkeypatched so the worker's *orchestration*
(merge path, conflict path, protected-branch path, CI verify/revert, gating and
state bookkeeping) is exercised deterministically. The git mechanics themselves
are covered by the live E2E (`--sync-only <repo> --dry`), not here.
"""
import importlib.util
import json
import pathlib
import subprocess
import sys
import tempfile
import time
HERE = pathlib.Path(__file__).resolve().parent
def load_worker():
spec = importlib.util.spec_from_file_location("pr_queue_worker", HERE / "pr-queue-worker.py")
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
return mod
W = load_worker()
PASSED = []
FAILED = []
def check(name, cond, detail=""):
(PASSED if cond else FAILED).append(name)
print((" PASS " if cond else " FAIL ") + name + ("" if cond else f" ← {detail}"))
def cp(args, code=0, out="", err=""):
return subprocess.CompletedProcess(args, code, out, err)
def make_state_file(tmpdir):
"""Point the worker's sync state at a temp file so tests never touch /tmp."""
W.SYNC_STATE_FILE = pathlib.Path(tmpdir) / "sync-state.json"
W.SYNC_TMP_BASE = pathlib.Path(tmpdir) / "sync-work"
W.sync_config # noqa: B018 — keep the reference obvious for readers
# ── 1. compare-API parsing ────────────────────────────────────────────────
def test_upstream_status_parsing():
print("1. upstream_status parses ahead_by/behind_by/tip")
calls = []
def fake_gh(method, path, token=None, json_data=None, retries=3):
calls.append(path)
return 200, {
"ahead_by": 4, "behind_by": 26,
"commits": [{"sha": "aaa"}, {"sha": "bbb"}, {"sha": "ccc"}, {"sha": "4eafc064"}],
}
W.gh_api = fake_gh
info = W.upstream_status("tok", "asepharyana/shiro-neko", "zakirkun/shiro-neko", "main", "main")
check("merge_count is ahead_by", info == (4, 26, "4eafc064"), str(info))
check("single compare call", len(calls) == 1, str(calls))
check("compare path uses owner:branch head (not owner/repo:branch)",
calls[0] == "/repos/asepharyana/shiro-neko/compare/main...zakirkun:main", str(calls))
def failing_gh(method, path, token=None, json_data=None, retries=3):
return 404, {"message": "Not Found"}
W.gh_api = failing_gh
check("unavailable compare → None", W.upstream_status("tok", "f", "p", "main", "main") is None)
# ── 2. gating: same upstream tip / interval not due ───────────────────────
def test_gating(tmpdir):
print("2. run_upstream_sync gating (interval + same upstream tip)")
make_state_file(tmpdir)
attempts = []
orig_list_forks = W.list_fork_repos
W.list_fork_repos = lambda: [("tok", "asepharyana/shiro-neko", "zakirkun/shiro-neko", "main")]
def fake_gh(method, path, token=None, json_data=None, retries=3):
if "/compare/" in path:
return 200, {"ahead_by": 4, "behind_by": 26, "commits": [{"sha": "up1"}]}
if path.startswith("/repos/asepharyana/shiro-neko/pulls"):
return 200, []
if path.endswith("/repos/asepharyana/shiro-neko"):
return 200, {"fork": True, "parent": {"full_name": "zakirkun/shiro-neko",
"default_branch": "main"},
"default_branch": "main"}
return 200, {}
W.gh_api = fake_gh
W.post_sync_discord = lambda *a, **k: True
orig_sync = W.sync_fork_repo
def fake_sync(token, fork, parent, local_branch, upstream_branch, upstream_sha,
merge_count, divergence, state, cfg, dry=False):
# mirror what the real function records so the gating is exercised
attempts.append((upstream_sha, dry))
entry = W._sync_entry(state, fork)
entry["last_sync_ts"] = time.time()
entry["last_attempt_sha"] = upstream_sha
W.save_sync_state(state)
return "synced", "fake"
W.sync_fork_repo = fake_sync
lines = W.run_upstream_sync()
check("first tick attempts the sync", len(attempts) == 1, str(lines))
lines = W.run_upstream_sync()
check("second tick not attempted (same upstream tip)", len(attempts) == 1, str(lines))
# upstream moves, but the interval has not elapsed → still no attempt
def fake_gh2(method, path, token=None, json_data=None, retries=3):
if "/compare/" in path:
return 200, {"ahead_by": 5, "behind_by": 26, "commits": [{"sha": "up2"}]}
if path == "/repos/asepharyana/shiro-neko":
return 200, {"fork": True, "parent": {"full_name": "zakirkun/shiro-neko",
"default_branch": "main"},
"default_branch": "main"}
return 200, []
W.gh_api = fake_gh2
lines = W.run_upstream_sync()
check("interval gate blocks a fresh retry", len(attempts) == 1, str(lines) + str(attempts))
# backdate the last attempt → the new upstream tip is picked up
st = json.loads(W.SYNC_STATE_FILE.read_text())
st["asepharyana/shiro-neko"]["last_sync_ts"] = time.time() - 2 * 3600
W.SYNC_STATE_FILE.write_text(json.dumps(st))
W.run_upstream_sync()
check("new upstream tip attempts again once the interval elapsed",
[a[0] for a in attempts] == ["up1", "up2"], str(attempts))
W.sync_fork_repo = orig_sync # later tests exercise the real implementation
W.list_fork_repos = orig_list_forks
# ── 3. clean merge → push → pending_verify ────────────────────────────────
def _install_git_fake(merge_code=0, unmerged=None):
"""Fake `git` for sync_fork_repo: clone/config/rev-parse/fetch/merge."""
state = {"head": "pre" + "0" * 37, "merges": [], "aborts": 0}
def fake_git(args, cwd=None, timeout=180):
a = [str(x) for x in args]
if a[0] == "clone":
return cp(a)
if a[0] == "rev-parse":
return cp(a, 0, state["head"] + "\n")
if a[0] == "merge":
if "--abort" in a:
state["aborts"] += 1
return cp(a)
state["merges"].append(a)
if merge_code == 0:
state["head"] = "merge" + "1" * 36
return cp(a, merge_code, "", "CONFLICT" if merge_code else "")
if a[0] == "diff":
# the merge diff (used for the quality pass) lists files; the
# unmerged listing is served by the monkeypatched _sync_unmerged_files
return cp(a, 0, "src/a.ts\nsrc/b.ts\n")
if a[0] == "status":
return cp(a, 0, "")
return cp(a)
W._sync_git = fake_git
W._sync_unmerged_files = lambda workdir: list(unmerged or [])
return state
def test_clean_merge_synced(tmpdir):
print("3. clean merge pushes and records pending_verify")
make_state_file(tmpdir)
state = {}
gitstate = _install_git_fake(merge_code=0, unmerged=[])
claude_calls = []
W._run_claude_sync = lambda workdir, prompt, label, fork: (claude_calls.append(label), (True, "ok"))[1]
W._sync_commit_if_dirty = lambda workdir, msg: False
pushes = []
def fake_push(workdir, fork, source, dest, app_token, force=False):
pushes.append((source, dest, force))
return True, "pat push ok", False
W._push_ref = fake_push
res, detail = W.sync_fork_repo("tok", "asepharyana/shiro-neko", "zakirkun/shiro-neko",
"main", "main", "up1", 4, 26, state, W.sync_config("x"))
entry = state["asepharyana/shiro-neko"]
check("status synced", res == "synced", f"{res} {detail}")
check("pushed to the local branch", pushes and pushes[0][1] == "refs/heads/main", str(pushes))
check("pending_verify records merge sha", entry["pending_verify"]["sha"] == gitstate["head"], str(entry))
check("pending_verify records pre_merge sha", entry["pending_verify"]["pre_merge_sha"].startswith("pre"), str(entry))
check("last_merged_upstream_sha recorded", entry["last_merged_upstream_sha"] == "up1", str(entry))
check("clean merge gets a quality pass", claude_calls == ["claude_sync_quality"], str(claude_calls))
# ── 4. conflicted merge resolved by Claude ────────────────────────────────
def test_conflict_resolved(tmpdir):
print("4. conflicted merge handed to Claude Code, then pushed")
make_state_file(tmpdir)
state = {}
_install_git_fake(merge_code=1, unmerged=["src/tools.ts", "src/ui/App.tsx"])
labels = []
def fake_claude(workdir, prompt, label, fork):
labels.append(label)
# after resolution the tree is clean → _sync_finish_merge must be called
W._sync_unmerged_files = lambda w: []
check("conflict prompt forbids --ours/--theirs",
"--ours/--theirs" in prompt and "rebase" in prompt)
return True, "resolved"
W._run_claude_sync = fake_claude
finished = {"n": 0}
def fake_finish(workdir):
finished["n"] += 1
return True, ""
W._sync_finish_merge = fake_finish
W._push_ref = lambda *a, **k: (True, "pat push ok", False)
res, detail = W.sync_fork_repo("tok", "f/x", "up/x", "main", "main", "up1", 4, 26, state, W.sync_config("x"))
check("status synced", res == "synced", f"{res} {detail}")
check("conflict runner used", labels == ["claude_sync_conflicts"], str(labels))
check("merge completed once", finished["n"] == 1, str(finished))
check("resolution noted in detail", "resolved 2 conflict" in detail, detail)
def test_conflict_failed_skips_and_dedupes(tmpdir):
print("5. failed conflict resolution → skip once, no retry at same upstream tip")
make_state_file(tmpdir)
state = {}
_install_git_fake(merge_code=1, unmerged=["src/tools.ts"])
W._run_claude_sync = lambda *a, **k: (False, "[INFRA] Claude Code timed out")
W._push_ref = lambda *a, **k: (_ for _ in ()).throw(AssertionError("must not push after failed resolution"))
res, detail = W.sync_fork_repo("tok", "f/x", "up/x", "main", "main", "up1", 4, 26, state, W.sync_config("x"))
entry = state["f/x"]
check("status conflict-failed", res == "conflict-failed", f"{res} {detail}")
check("upstream tip recorded as attempted", entry["last_attempt_sha"] == "up1", str(entry))
check("skip reason recorded", "timed out" in entry.get("skip_reason", ""), str(entry))
check("pending_verify not set", not entry.get("pending_verify"), str(entry))
# ── 6. protected branch → PR path ─────────────────────────────────────────
def test_protected_branch_opens_pr(tmpdir):
print("6. protected branch falls back to an upstream-sync PR")
make_state_file(tmpdir)
state = {}
_install_git_fake(merge_code=0, unmerged=[])
W._run_claude_sync = lambda *a, **k: (True, "ok")
W._sync_commit_if_dirty = lambda *a, **k: False
pushes = []
def fake_push(workdir, fork, source, dest, app_token, force=False):
pushes.append(dest)
if dest.endswith("refs/heads/main"):
return False, "pat: remote: error: GH006 protected branch", True
return True, "pat push ok", False
W._push_ref = fake_push
W.open_sync_pr = lambda *a, **k: 42
res, detail = W.sync_fork_repo("tok", "f/x", "up/x", "main", "main", "up1", 4, 26, state, W.sync_config("x"))
check("status pr-opened", res == "pr-opened", f"{res} {detail}")
check("direct push attempted first", pushes and pushes[0] == "refs/heads/main", str(pushes))
check("no pending_verify on the PR path", not state["f/x"].get("pending_verify"), str(state["f/x"]))
# ── 7. CI verify / auto-revert ────────────────────────────────────────────
def test_verify_green_clears(tmpdir):
print("7. verify_pending_syncs: green CI clears the watch")
make_state_file(tmpdir)
state = {"f/x": {"pending_verify": {"sha": "abc", "pre_merge_sha": "pre", "branch": "main",
"pushed_at": time.time()}}}
def fake_gh(method, path, token=None, json_data=None, retries=3):
if "/commits/main" in path:
return 200, [{"sha": "abc"}]
if path.endswith("/check-runs"):
return 200, {"check_runs": [{"name": "build", "status": "completed", "conclusion": "success"}]}
return 200, {}
W.gh_api = fake_gh
reverted = []
W._sync_revert_merge = lambda *a, **k: (reverted.append(a), (True, "reverted"))[1]
lines = W.verify_pending_syncs(state, {"f/x": "tok"})
check("green verifies", any("verified green" in l for l in lines), str(lines))
check("watch cleared", state["f/x"]["pending_verify"] is None, str(state))
check("no revert on green", not reverted, str(reverted))
def test_verify_red_reverts(tmpdir):
print("8. verify_pending_syncs: red CI at our merge sha reverts it")
make_state_file(tmpdir)
state = {"f/x": {"pending_verify": {"sha": "abc", "pre_merge_sha": "pre", "branch": "main",
"pushed_at": time.time()}}}
def fake_gh(method, path, token=None, json_data=None, retries=3):
if "/commits/main" in path:
return 200, [{"sha": "abc"}]
if path.endswith("/check-runs"):
return 200, {"check_runs": [{"name": "build", "status": "completed", "conclusion": "failure"}]}
return 200, {}
W.gh_api = fake_gh
reverted = []
def fake_revert(token, fork, branch, pre_merge_sha, reason):
reverted.append((fork, branch, pre_merge_sha, reason))
return True, "reverted main to pre"
W._sync_revert_merge = fake_revert
W.post_sync_discord = lambda *a, **k: True
lines = W.verify_pending_syncs(state, {"f/x": "tok"})
check("revert invoked with the pre-merge sha",
reverted and reverted[0][2] == "pre" and reverted[0][1] == "main", str(reverted))
check("revert reported", any("reverted" in l for l in lines), str(lines))
check("watch cleared after revert", state["f/x"]["pending_verify"] is None, str(state))
def test_verify_never_reverts_foreign_commits(tmpdir):
print("9. verify_pending_syncs: never reverts when someone pushed on top")
make_state_file(tmpdir)
state = {"f/x": {"pending_verify": {"sha": "abc", "pre_merge_sha": "pre", "branch": "main",
"pushed_at": time.time()}}}
def fake_gh(method, path, token=None, json_data=None, retries=3):
if "/commits/main" in path:
return 200, [{"sha": "humancommit"}]
if path.endswith("/check-runs"):
return 200, {"check_runs": [{"name": "build", "status": "completed", "conclusion": "failure"}]}
return 200, {}
W.gh_api = fake_gh
reverted = []
W._sync_revert_merge = lambda *a, **k: (reverted.append(a), (True, "x"))[1]
lines = W.verify_pending_syncs(state, {"f/x": "tok"})
check("no revert when the tip moved", not reverted, str(reverted))
check("tip change reported", any("moved past our merge" in l for l in lines), str(lines))
check("watch cleared", state["f/x"]["pending_verify"] is None, str(state))
def test_verify_waits_for_running_ci(tmpdir):
print("10. verify_pending_syncs: waits while checks run, respects no-CI grace")
make_state_file(tmpdir)
state = {"f/x": {"pending_verify": {"sha": "abc", "pre_merge_sha": "pre", "branch": "main",
"pushed_at": time.time()}}}
def fake_gh(method, path, token=None, json_data=None, retries=3):
if "/commits/main" in path:
return 200, [{"sha": "abc"}]
if path.endswith("/check-runs"):
return 200, {"check_runs": [{"name": "build", "status": "in_progress"}]}
return 200, {}
W.gh_api = fake_gh
W.verify_pending_syncs(state, {"f/x": "tok"})
check("running CI keeps the watch", state["f/x"]["pending_verify"] is not None, str(state))
# no checks at all and young → still pending (check-runs lag)
def fake_gh_none(method, path, token=None, json_data=None, retries=3):
if "/commits/main" in path:
return 200, [{"sha": "abc"}]
return 200, {"check_runs": []}
W.gh_api = fake_gh_none
W.verify_pending_syncs(state, {"f/x": "tok"})
check("fresh no-CI merge stays pending", state["f/x"]["pending_verify"] is not None, str(state))
# old + no CI → stop watching
state["f/x"]["pending_verify"]["pushed_at"] = time.time() - 3600
W.verify_pending_syncs(state, {"f/x": "tok"})
check("stale no-CI merge stops the watch", state["f/x"]["pending_verify"] is None, str(state))
# ── 11. fork discovery ────────────────────────────────────────────────────
def test_list_fork_repos():
print("11. list_fork_repos returns only forks with a resolvable parent")
def fake_gh(method, path, token=None, json_data=None, retries=3):
if path == "/app/installations":
return 200, [{"id": 1}]
if path.startswith("/installation/repositories"):
return 200, {"repositories": [
{"full_name": "o/plain-repo", "fork": False},
{"full_name": "o/fork-a", "fork": True},
{"full_name": "o/fork-b", "fork": True},
]}
if path == "/repos/o/fork-a":
return 200, {"fork": True, "default_branch": "main", "parent": {"full_name": "up/a"}}
if path == "/repos/o/fork-b":
return 200, {"fork": True, "default_branch": "main"} # parent stripped
return 404, {}
W.gh_api = fake_gh
W.get_installation_token = lambda inst: "tok"
forks = W.list_fork_repos()
check("only the resolvable fork returned", forks == [("tok", "o/fork-a", "up/a", "main")], str(forks))
# ── 12. push-error classification ─────────────────────────────────────────
def test_push_error_classification():
print("12. push error classification")
check("GH006 → protected", W._is_protected_push_error(
"remote: error: GH006: Protected branch update failed for refs/heads/main."))
check("required status checks → protected", W._is_protected_push_error(
"remote: error: Required status check \"ci\" is expected."))
check("workflows permission detected", W._is_workflow_push_error(
"! [remote rejected] main -> main (refusing to allow a GitHub App to create or update workflow `.github/workflows/ci.yml` without `workflows` permission)"))
check("plain rejected push is not 'protected'", not W._is_protected_push_error(
"fatal: could not read Username for 'https://github.com'"))
# ── 13. per-repo config override ──────────────────────────────────────────
def test_sync_config_override():
print("13. per-repo config override merges over defaults")
W.UPSTREAM_SYNC["repos"] = {"o/fork": {"interval_h": 6, "verify_ci": False}}
cfg = W.sync_config("o/fork")
check("override applied", cfg["interval_h"] == 6 and cfg["verify_ci"] is False, str(cfg))
check("defaults preserved", cfg["resolve_conflicts"] is True, str(cfg))
check("repos key not leaked", "repos" not in cfg, str(cfg))
W.UPSTREAM_SYNC["repos"] = {}
def main():
with tempfile.TemporaryDirectory() as tmpdir:
test_upstream_status_parsing()
test_gating(tmpdir)
test_clean_merge_synced(tmpdir)
test_conflict_resolved(tmpdir)
test_conflict_failed_skips_and_dedupes(tmpdir)
test_protected_branch_opens_pr(tmpdir)
test_verify_green_clears(tmpdir)
test_verify_red_reverts(tmpdir)
test_verify_never_reverts_foreign_commits(tmpdir)
test_verify_waits_for_running_ci(tmpdir)
test_list_fork_repos()
test_push_error_classification()
test_sync_config_override()
print(f"\n{len(PASSED)} passed, {len(FAILED)} failed")
if FAILED:
print("failed: " + ", ".join(FAILED))
return 1 if FAILED else 0
if __name__ == "__main__":
sys.exit(main())