Skip to content
This repository was archived by the owner on Jul 12, 2026. It is now read-only.
Draft
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
3 changes: 3 additions & 0 deletions .jules/bolt.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
## 2024-07-12 - Concurrent PR open check filtering
**Learning:** `_filter_to_still_open_prs` checks PR states sequentially, leading to an N+1 performance bottleneck when handling many PRs in fan-out mode because `_pr_is_still_open` triggers a slow subprocess call (`gh pr view`).
**Action:** Use `concurrent.futures.ThreadPoolExecutor` combined with `executor.map` to fetch PR states concurrently while preserving order and catching exceptions without interleaving main thread logs.
90 changes: 44 additions & 46 deletions ralph_loop/cli.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
"""Command-line interface and top-level orchestration."""

from __future__ import annotations

import argparse
Expand All @@ -7,6 +8,7 @@
import shlex
import signal
import subprocess
import concurrent.futures
import sys
import threading
import time
Expand Down Expand Up @@ -335,9 +337,7 @@ def _validate_pr_metadata(
) -> str:
if pr_data.get("state") != "OPEN":
raise CommandError(
"PR {} is not open (state={}).".format(
pr_number, pr_data.get("state")
)
"PR {} is not open (state={}).".format(pr_number, pr_data.get("state"))
)
if pr_data.get("isDraft"):
raise CommandError(
Expand Down Expand Up @@ -544,23 +544,36 @@ def _filter_to_still_open_prs(pr_numbers: List[int]) -> List[int]:
swallow stale PRs because of a flaky network.
"""
kept: List[int] = []
for pr in pr_numbers:
if not pr_numbers:
return kept

def check_pr(pr: int) -> Tuple[int, bool, Optional[Exception]]:
try:
still_open = _pr_is_still_open(pr)
return pr, _pr_is_still_open(pr), None
except CommandError as exc:
_print_step(
"Could not confirm PR #{} open state ({}); keeping it in the "
"fan-out set.".format(pr, exc)
)
kept.append(pr)
continue
if still_open:
kept.append(pr)
else:
_print_step(
"PR #{} is no longer open (per gh pr view); skipping "
"fan-out spawn.".format(pr)
)
return pr, True, exc

# ⚑ Bolt Optimization: Use ThreadPoolExecutor for concurrent PR state fetching
# This prevents the N+1 performance bottleneck of waiting on sequential `gh pr view` calls,
# significantly speeding up the supervisor startup and re-evaluation cycle.
with concurrent.futures.ThreadPoolExecutor(
max_workers=min(10, len(pr_numbers))
) as executor:
for pr, still_open, exc in executor.map(check_pr, pr_numbers):
if exc:
_print_step(
"Could not confirm PR #{} open state ({}); keeping it in the "
"fan-out set.".format(pr, exc)
)
kept.append(pr)
continue
if still_open:
kept.append(pr)
else:
_print_step(
"PR #{} is no longer open (per gh pr view); skipping "
"fan-out spawn.".format(pr)
)
return kept


Expand Down Expand Up @@ -605,9 +618,7 @@ def _fan_out_all_prs(
if args.fan_out_log_dir:
log_root = os.path.abspath(os.path.expanduser(args.fan_out_log_dir))
else:
log_root = os.path.join(
os.path.dirname(script_path), ".ralph-logs", "fan-out"
)
log_root = os.path.join(os.path.dirname(script_path), ".ralph-logs", "fan-out")
os.makedirs(log_root, exist_ok=True)
stuck_timeout = max(60, args.fan_out_stuck_timeout_seconds)
respawn_backoff = max(1, args.fan_out_respawn_backoff_seconds)
Expand Down Expand Up @@ -700,9 +711,7 @@ def _request_reload(_signum, _frame):
reason = "ordinary exit"
_print_step(
"PR #{} loop exited with code {} ({}); respawning after "
"{}s (log: {})".format(
pr, rc, reason, backoff_for_pr, log_path
)
"{}s (log: {})".format(pr, rc, reason, backoff_for_pr, log_path)
)
try:
with open(log_path, "ab", buffering=0) as marker:
Expand Down Expand Up @@ -796,15 +805,11 @@ def _request_reload(_signum, _frame):
# ordinary-exit and stuck-timeout branches above, which both
# rewrite pending_backoff to ``respawn_backoff``.
_print_step(
"Respawned PR #{} pid={} (log: {})".format(
pr, proc.pid, log_path
)
"Respawned PR #{} pid={} (log: {})".format(pr, proc.pid, log_path)
)
finally:
if reload_requested["flag"]:
_print_step(
"Reload requested via SIGHUP; re-exec'ing supervisor"
)
_print_step("Reload requested via SIGHUP; re-exec'ing supervisor")
for pr, (proc, _log_path, log_handle, _spawned_at) in list(
children.items()
):
Expand All @@ -822,9 +827,7 @@ def _request_reload(_signum, _frame):
signal.signal(signal.SIGHUP, previous_hup)
sys.stdout.flush()
sys.stderr.flush()
os.execv(
sys.executable, [sys.executable, script_path] + sys.argv[1:]
)
os.execv(sys.executable, [sys.executable, script_path] + sys.argv[1:])
for pr, (proc, _log_path, log_handle, _spawned_at) in list(children.items()):
try:
proc.terminate()
Expand Down Expand Up @@ -854,9 +857,7 @@ def _request_reload(_signum, _frame):
return 0


def _resolve_target_directories(
raw_dirs: List[str], recursive: bool
) -> List[str]:
def _resolve_target_directories(raw_dirs: List[str], recursive: bool) -> List[str]:
"""Expand the user-supplied directory args into concrete repo paths.

- If ``recursive`` is set, each input path is treated as a parent and we
Expand All @@ -878,9 +879,7 @@ def _push(path: str) -> None:
path = os.path.abspath(os.path.expanduser(raw))
if not os.path.isdir(path):
raise CommandError(
"Target directory does not exist or is not a directory: {}".format(
raw
)
"Target directory does not exist or is not a directory: {}".format(raw)
)
if not recursive:
_push(path)
Expand Down Expand Up @@ -933,7 +932,10 @@ def _fan_out_across_directories(
base_args.append(token)
procs: List[Tuple[str, subprocess.Popen, str, Any]] = []
for target_dir in target_dirs:
slug = re.sub(r"[^A-Za-z0-9._-]+", "-", os.path.basename(target_dir.rstrip("/"))) or "repo"
slug = (
re.sub(r"[^A-Za-z0-9._-]+", "-", os.path.basename(target_dir.rstrip("/")))
or "repo"
)
log_path = os.path.join(log_root, "{}.log".format(slug))
log_handle = open(log_path, "ab", buffering=0)
cmd = [sys.executable, script_path] + base_args + [target_dir]
Expand Down Expand Up @@ -1127,9 +1129,7 @@ def _handle_shutdown(signum, _frame):
)

pr_target = str(pr_number)
_print_step(
"Using PR #{} {}".format(pr_number, pr_data.get("url", ""))
)
_print_step("Using PR #{} {}".format(pr_number, pr_data.get("url", "")))
_mark_pr_needs_review(pr_target)
if not args.skip_rebase:
_print_step("Initial rebase before review/fix loop")
Expand Down Expand Up @@ -1345,9 +1345,7 @@ def _handle_shutdown(signum, _frame):

if not args.skip_merge:
fresh_pr_data = _pr_view(str(pr_number))
fresh_branch = _validate_pr_metadata(
fresh_pr_data, pr_number, args.base
)
fresh_branch = _validate_pr_metadata(fresh_pr_data, pr_number, args.base)
if fresh_branch != branch:
raise CommandError(
"PR #{} head branch changed from '{}' to '{}' during the run.".format(
Expand Down