Fix critical coordination bugs in implementation pool supervisor
CI / status-check (push) Blocked by required conditions
CI / docker (push) Blocked by required conditions
CI / lint (push) Waiting to run
CI / typecheck (push) Waiting to run
CI / security (push) Waiting to run
CI / quality (push) Waiting to run
CI / unit_tests (push) Waiting to run
CI / integration_tests (push) Waiting to run
CI / helm (push) Waiting to run
CI / e2e_tests (push) Waiting to run
CI / coverage (push) Blocked by required conditions
CI / benchmark-regression (push) Blocked by required conditions
CI / benchmark-publish (push) Waiting to run
CI / build (push) Waiting to run

- Fix session adoption logic with correct title patterns for both
  worker-issue-impl and worker-pr-fix sessions
- Add PR worker adoption to coordinate orphaned PR fix workers
- Enhance worker verification with comprehensive status checking,
  retry logic, and proper error handling
- Add defensive programming with worker count enforcement and 
  state validation to prevent coordination drift
- Improve JSON parsing with safe error handling throughout
- Add periodic maintenance cycle (every 5 iterations) for
  worker state validation and limit enforcement

These fixes resolve the core issue where the implementation pool
supervisor was not properly coordinating 40+ existing workers,
causing worker count to exceed the designed limit of 32.
This commit is contained in:
2026-04-05 19:17:35 +00:00
parent 1bd0c7999d
commit 88cfc33ab2
+348 -91
View File
@@ -349,23 +349,41 @@ pattern with configurable parallelism and **speculative pre-cloning** for
maximum throughput. maximum throughput.
``` ```
max_workers = CA_MAX_PARALLEL_WORKERS (from env, or ask user if unset) # Initialize core state variables with proper error checking
queue = [] # prioritized list of unblocked issues to work on try:
active_pr_workers = {} # pr_number -> {session_id, work_type, assigned_at, issue_number} max_workers = int(bash("echo $CA_MAX_PARALLEL_WORKERS", timeout=5000) or "4")
active_issue_workers = {} # issue_number -> session_id if max_workers <= 0 or max_workers > 64:
pr_work_queue = [] # PRs needing work but no worker assigned yet print(f"[WARNING] Invalid max_workers {max_workers}, using default 32")
completed = [] # list of {issue_number, branch, pr_number, pr_url, ...} max_workers = 32
completed_issues = set() # Set of completed issue numbers except:
failed = {} # issue_number -> consecutive_failure_count print(f"[WARNING] Could not parse CA_MAX_PARALLEL_WORKERS, using default 32")
ref_summary = result from ca-ref-reader max_workers = 32
cycle = 0
SERVER = "http://localhost:4096" print(f"[CONFIG] Max parallel workers: {max_workers}")
FORGEJO_USERNAME = <from env or user>
owner = "cleveragents" # Core tracking dictionaries
repo = "cleveragents-core" queue = [] # prioritized list of unblocked issues to work on
targeted_issue_numbers = [] # If user specifies specific issues active_pr_workers = {} # pr_number -> {session_id, work_type, assigned_at, issue_number}
active_issue_workers = {} # issue_number -> session_id
pr_work_queue = [] # PRs needing work but no worker assigned yet
completed = [] # list of {issue_number, branch, pr_number, pr_url, ...}
completed_issues = set() # Set of completed issue numbers
failed = {} # issue_number -> consecutive_failure_count
ref_summary = result_from_ca_ref_reader # Will be set after ref-reader completes
cycle = 0
# Constants
SERVER = "http://localhost:4096"
FORGEJO_USERNAME = forgejo_username # From startup sequence
owner = "cleveragents"
repo = "cleveragents-core"
targeted_issue_numbers = [] # If user specifies specific issues
selective_issue_targeting = False selective_issue_targeting = False
print(f"[CONFIG] Repository: {owner}/{repo}")
print(f"[CONFIG] Forgejo username: {FORGEJO_USERNAME}")
print(f"[CONFIG] OpenCode server: {SERVER}")
# Helper function to extract branch from issue body metadata # Helper function to extract branch from issue body metadata
function extract_branch_from_issue_body(body): function extract_branch_from_issue_body(body):
# Look for branch in metadata section # Look for branch in metadata section
@@ -391,26 +409,168 @@ function now():
# active tracking instead of launching duplicates. # active tracking instead of launching duplicates.
# To start fresh, run ca-session-cleanup BEFORE the product-builder. # To start fresh, run ca-session-cleanup BEFORE the product-builder.
EXISTING_WORKERS = bash("curl -s ${SERVER}/session | python3 -c \" def adopt_existing_workers():
import sys, json """Adopt orphaned worker sessions into coordination."""
for s in json.loads(sys.stdin.read()): print("[ADOPTION] Checking for existing worker sessions to adopt...")
title = s.get('title','')
if title.startswith('[CA-AUTO] worker-impl:'): try:
# Extract issue number from title # Get all active sessions
issue_num = title.replace('[CA-AUTO] worker-impl: issue-','') sessions_response = bash("curl -s ${SERVER}/session", timeout=30000)
print(issue_num + '=' + s['id']) sessions = json.loads(sessions_response) if sessions_response else []
\"", timeout=30000)
# Get session status for verification
status_response = bash("curl -s ${SERVER}/session/status", timeout=30000)
status_data = json.loads(status_response) if status_response else []
active_session_ids = {s["id"] for s in status_data if s.get("status") == "active"}
adopted_issue_workers = 0
adopted_pr_workers = 0
for session in sessions:
title = session.get('title', '')
session_id = session.get('id', '')
# Only adopt active sessions
if session_id not in active_session_ids:
continue
# Adopt issue implementation workers
if title.startswith('[CA-AUTO] worker-issue-impl: issue-'):
try:
# Extract issue number: "[CA-AUTO] worker-issue-impl: issue-123"
issue_num_str = title.replace('[CA-AUTO] worker-issue-impl: issue-', '').strip()
issue_number = int(issue_num_str)
active_issue_workers[issue_number] = session_id
adopted_issue_workers += 1
print(f"[ADOPTION] Adopted issue worker: Issue #{issue_number} -> {session_id[:8]}...")
except (ValueError, IndexError) as e:
print(f"[WARNING] Could not parse issue number from title: {title}")
# Adopt PR fix workers
elif title.startswith('[CA-AUTO] worker-pr-fix: PR-'):
try:
# Extract PR number: "[CA-AUTO] worker-pr-fix: PR-456"
pr_num_str = title.replace('[CA-AUTO] worker-pr-fix: PR-', '').strip()
pr_number = int(pr_num_str)
# We need to reconstruct the PR worker info - get PR data from Forgejo
pr_data = forgejo_get_pull_request_by_index(owner, repo, pr_number)
if pr_data:
# Extract issue number from PR body
issue_number = None
if "Closes #" in pr_data.body:
issue_number = int(pr_data.body.split("Closes #")[1].split()[0])
elif "Fixes #" in pr_data.body:
issue_number = int(pr_data.body.split("Fixes #")[1].split()[0])
active_pr_workers[pr_number] = {
"session_id": session_id,
"work_type": "unknown", # Will be re-determined in next cycle
"assigned_at": now(),
"issue_number": issue_number
}
adopted_pr_workers += 1
print(f"[ADOPTION] Adopted PR worker: PR #{pr_number} -> {session_id[:8]}...")
except (ValueError, IndexError) as e:
print(f"[WARNING] Could not parse PR number from title: {title}")
print(f"[ADOPTION] Adopted {adopted_issue_workers} issue workers, {adopted_pr_workers} PR workers")
return adopted_issue_workers + adopted_pr_workers
except Exception as e:
print(f"[ERROR] Session adoption failed: {e}")
return 0
STATUS = bash("curl -s ${SERVER}/session/status", timeout=30000) # Call adoption function
for line in EXISTING_WORKERS: total_adopted = adopt_existing_workers()
issue_number, session_id = line.split("=")
if session_id is active in STATUS:
active[int(issue_number)] = session_id # Adopt into active tracking
# Remove from queue if present (already being worked on)
queue = [i for i in queue if i.number != int(issue_number)]
# ── Helper: launch one worker via prompt_async ─────────────────── # ── Defensive Programming: Worker Count Enforcement ──────────────
def enforce_worker_limits():
"""Enforce hard worker count limits and clean up excess workers."""
total_workers = len(active_pr_workers) + len(active_issue_workers)
if total_workers > max_workers:
print(f"[WARNING] Worker count ({total_workers}) exceeds limit ({max_workers})")
# Clean up oldest workers first (LIFO cleanup to preserve recent work)
excess = total_workers - max_workers
# First clean up issue workers (PRs have higher priority)
issue_items = list(active_issue_workers.items())
for i, (issue_num, session_id) in enumerate(issue_items[-excess:]):
print(f"[CLEANUP] Terminating excess issue worker: Issue #{issue_num}")
bash(f"curl -s -X DELETE ${SERVER}/session/{session_id}", timeout=15000)
del active_issue_workers[issue_num]
excess -= 1
if excess <= 0:
break
# If still over limit, clean up PR workers
if excess > 0:
pr_items = list(active_pr_workers.items())
for i, (pr_num, pr_info) in enumerate(pr_items[-excess:]):
print(f"[CLEANUP] Terminating excess PR worker: PR #{pr_num}")
bash(f"curl -s -X DELETE ${SERVER}/session/{pr_info['session_id']}", timeout=15000)
del active_pr_workers[pr_num]
excess -= 1
if excess <= 0:
break
def validate_worker_state():
"""Validate that tracked workers actually exist and are active."""
print("[VALIDATION] Validating worker state consistency...")
try:
status_response = bash("curl -s ${SERVER}/session/status", timeout=30000)
if not status_response:
print("[WARNING] Could not get session status for validation")
return
status_data = json.loads(status_response)
active_session_ids = {s["id"] for s in status_data if s.get("status") == "active"}
# Validate issue workers
dead_issue_workers = []
for issue_num, session_id in active_issue_workers.items():
if session_id not in active_session_ids:
print(f"[VALIDATION] Issue worker #{issue_num} session {session_id[:8]}... is dead")
dead_issue_workers.append(issue_num)
# Clean up dead issue workers
for issue_num in dead_issue_workers:
del active_issue_workers[issue_num]
print(f"[CLEANUP] Removed dead issue worker tracking for Issue #{issue_num}")
# Validate PR workers
dead_pr_workers = []
for pr_num, pr_info in active_pr_workers.items():
session_id = pr_info["session_id"]
if session_id not in active_session_ids:
print(f"[VALIDATION] PR worker #{pr_num} session {session_id[:8]}... is dead")
dead_pr_workers.append(pr_num)
# Clean up dead PR workers
for pr_num in dead_pr_workers:
del active_pr_workers[pr_num]
print(f"[CLEANUP] Removed dead PR worker tracking for PR #{pr_num}")
cleaned_count = len(dead_issue_workers) + len(dead_pr_workers)
if cleaned_count > 0:
print(f"[VALIDATION] Cleaned up {cleaned_count} dead worker references")
else:
print(f"[VALIDATION] ✓ All {len(active_issue_workers)} issue + {len(active_pr_workers)} PR workers are active")
except Exception as e:
print(f"[ERROR] Worker state validation failed: {e}")
# Run initial worker limit enforcement and validation
enforce_worker_limits()
validate_worker_state()
# ── Helper: launch one worker via prompt_async with robust verification ──
function dispatch_worker(mode, work_item, ref_summary): function dispatch_worker(mode, work_item, ref_summary):
"""Dispatch a worker with comprehensive verification and retry logic."""
if mode == "pr-fix": if mode == "pr-fix":
# PR fix mode # PR fix mode
pr = work_item["pr"] pr = work_item["pr"]
@@ -442,6 +602,7 @@ CRITICAL: You own this PR until it is merged. Monitor and handle all feedback.
Do not exit until the PR is merged or blocked by human feedback (needs feedback label).""" Do not exit until the PR is merged or blocked by human feedback (needs feedback label)."""
title = f"[CA-AUTO] worker-pr-fix: PR-{pr.number}" title = f"[CA-AUTO] worker-pr-fix: PR-{pr.number}"
work_id = f"PR-{pr.number}"
else: # issue-impl mode else: # issue-impl mode
issue = work_item issue = work_item
@@ -468,25 +629,141 @@ CRITICAL: You own this issue from implementation through PR merge.
Do not exit until the PR is merged. Monitor and handle all review feedback.""" Do not exit until the PR is merged. Monitor and handle all review feedback."""
title = f"[CA-AUTO] worker-issue-impl: issue-{issue['number']}" title = f"[CA-AUTO] worker-issue-impl: issue-{issue['number']}"
work_id = f"issue-{issue['number']}"
# Create session
SESSION_ID = bash("curl -s -X POST ${SERVER}/session \
-H 'Content-Type: application/json' \
-d '{\"title\": \"${title}\"}' \
| python3 -c \"import sys,json; print(json.loads(sys.stdin.read())['id'])\"",
timeout=30000)
# Launch worker def verify_worker_started(session_id, retries=3):
bash("curl -s -X POST ${SERVER}/session/${SESSION_ID}/prompt_async \ """Comprehensively verify worker is running."""
-H 'Content-Type: application/json' \ for attempt in range(retries):
-d '{\"agent\": \"ca-issue-worker\", \ try:
\"parts\": [{\"type\": \"text\", \"text\": \"${prompt}\"}]}'", # Check session exists and is active
timeout=30000) session_response = bash(f"curl -s ${SERVER}/session/{session_id}", timeout=15000)
if not session_response or "not found" in session_response.lower():
print(f"[VERIFY] Attempt {attempt+1}: Session {session_id[:8]}... not found")
bash("sleep 2", timeout=5000)
continue
# Parse JSON safely
try:
session_data = json.loads(session_response)
except:
print(f"[VERIFY] Attempt {attempt+1}: Invalid JSON response for session {session_id[:8]}...")
bash("sleep 2", timeout=5000)
continue
# Check if session is active
status_response = bash("curl -s ${SERVER}/session/status", timeout=15000)
if status_response:
try:
status_data = json.loads(status_response)
session_active = any(s["id"] == session_id and s.get("status") == "active"
for s in status_data)
if session_active:
print(f"[VERIFY] ✓ Worker {work_id} verified active (session {session_id[:8]}...)")
return True
else:
print(f"[VERIFY] Attempt {attempt+1}: Session {session_id[:8]}... not active")
except:
print(f"[VERIFY] Attempt {attempt+1}: Could not parse status response")
bash("sleep 3", timeout=5000)
except Exception as e:
print(f"[VERIFY] Attempt {attempt+1}: Verification error: {e}")
bash("sleep 3", timeout=5000)
print(f"[ERROR] Worker verification failed after {retries} attempts for {work_id}")
return False
return SESSION_ID # Main dispatch logic with retry
max_attempts = 2
for attempt in range(max_attempts):
try:
print(f"[DISPATCH] Attempt {attempt+1}: Launching worker for {work_id}")
# Create session with error checking
session_create_cmd = f"""curl -s -X POST ${SERVER}/session \\
-H 'Content-Type: application/json' \\
-d '{{"title": "{title}"}}' """
session_response = bash(session_create_cmd, timeout=30000)
if not session_response:
print(f"[ERROR] Empty response from session creation for {work_id}")
if attempt < max_attempts - 1:
bash("sleep 5", timeout=10000)
continue
else:
return None
# Parse session ID safely
try:
session_data = json.loads(session_response)
session_id = session_data.get('id')
if not session_id:
print(f"[ERROR] No session ID in response for {work_id}")
if attempt < max_attempts - 1:
bash("sleep 5", timeout=10000)
continue
else:
return None
except json.JSONDecodeError as e:
print(f"[ERROR] Invalid JSON from session creation for {work_id}: {e}")
if attempt < max_attempts - 1:
bash("sleep 5", timeout=10000)
continue
else:
return None
# Launch worker with escaped prompt
escaped_prompt = prompt.replace('"', '\\"').replace('\n', '\\n')
launch_cmd = f"""curl -s -X POST ${SERVER}/session/{session_id}/prompt_async \\
-H 'Content-Type: application/json' \\
-d '{{"agent": "ca-issue-worker", "parts": [{{"type": "text", "text": "{escaped_prompt}"}}]}}'"""
launch_response = bash(launch_cmd, timeout=30000)
# Brief wait for worker to initialize
bash("sleep 3", timeout=10000)
# Verify worker is actually running
if verify_worker_started(session_id):
return session_id
else:
print(f"[ERROR] Worker failed to start properly for {work_id}, cleaning up session {session_id[:8]}...")
bash(f"curl -s -X DELETE ${SERVER}/session/{session_id}", timeout=15000)
if attempt < max_attempts - 1:
bash("sleep 10", timeout=15000) # Longer wait before retry
continue
else:
return None
except Exception as e:
print(f"[ERROR] Dispatch attempt {attempt+1} failed for {work_id}: {e}")
if attempt < max_attempts - 1:
bash("sleep 10", timeout=15000)
print(f"[CRITICAL] All dispatch attempts failed for {work_id}")
return None
# ── Main dispatch + monitoring loop ────────────────────────────── # ── Main dispatch + monitoring loop ──────────────────────────────
# Import required modules for JSON handling
import json
import datetime
# Helper function to safely parse JSON responses
def safe_json_parse(response, default=None):
"""Safely parse JSON response with error handling."""
if not response:
return default or []
try:
return json.loads(response)
except (json.JSONDecodeError, TypeError) as e:
print(f"[JSON ERROR] Failed to parse response: {e}")
return default or []
LOOP FOREVER: LOOP FOREVER:
cycle += 1 cycle += 1
@@ -520,19 +797,11 @@ LOOP FOREVER:
# Actually dispatch PR fix worker using the helper function # Actually dispatch PR fix worker using the helper function
session_id = dispatch_worker("pr-fix", pr_work, ref_summary) session_id = dispatch_worker("pr-fix", pr_work, ref_summary)
# Verify the worker actually started if session_id is None:
bash("sleep 2", timeout=5000) # Brief wait print(f"[CRITICAL] Failed to dispatch worker for PR #{pr_work['pr'].number} - skipping")
verify_status = bash(f"curl -s ${SERVER}/session/{session_id}", timeout=30000) # Put work back in queue for next cycle
pr_work_queue.insert(0, pr_work)
if "error" in verify_status.lower() or "not found" in verify_status.lower(): break # Stop trying to dispatch more workers this cycle
print(f"[ERROR] Worker dispatch failed for PR #{pr_work['pr'].number}")
# Retry dispatch
session_id = dispatch_worker("pr-fix", pr_work, ref_summary)
bash("sleep 2", timeout=5000)
verify_status = bash(f"curl -s ${SERVER}/session/{session_id}", timeout=30000)
if "error" in verify_status.lower() or "not found" in verify_status.lower():
print(f"[ERROR] Retry failed - skipping PR #{pr_work['pr'].number}")
continue
active_pr_workers[pr_work["pr"].number] = { active_pr_workers[pr_work["pr"].number] = {
"session_id": session_id, "session_id": session_id,
@@ -543,7 +812,7 @@ LOOP FOREVER:
slots_available -= 1 slots_available -= 1
# Log dispatch # Log dispatch
print(f"[{now()}] Dispatched PR-fix worker for PR #{pr_work['pr'].number} ({pr_work['work_type']}) - verified active") print(f"[{now()}] Dispatched PR-fix worker for PR #{pr_work['pr'].number} ({pr_work['work_type']}) - session {session_id[:8]}...")
# ── STEP 3: Only dispatch to issues if ALL PRs have workers ────── # ── STEP 3: Only dispatch to issues if ALL PRs have workers ──────
# CRITICAL: Block ALL issue work if ANY PR needs attention # CRITICAL: Block ALL issue work if ANY PR needs attention
@@ -628,33 +897,22 @@ LOOP FOREVER:
# Actually dispatch issue implementation worker # Actually dispatch issue implementation worker
session_id = dispatch_worker("issue-impl", issue_work, ref_summary) session_id = dispatch_worker("issue-impl", issue_work, ref_summary)
# Verify the worker actually started if session_id is None:
bash("sleep 2", timeout=5000) # Brief wait print(f"[CRITICAL] Failed to dispatch worker for Issue #{issue.number} - skipping")
verify_status = bash(f"curl -s ${SERVER}/session/{session_id}", timeout=30000) # Put issue back in queue for next cycle
queue.insert(0, issue)
if "error" in verify_status.lower() or "not found" in verify_status.lower(): break # Stop trying to dispatch more workers this cycle
print(f"[ERROR] Worker dispatch failed for Issue #{issue.number}")
# Retry dispatch
session_id = dispatch_worker("issue-impl", issue_work, ref_summary)
bash("sleep 2", timeout=5000)
verify_status = bash(f"curl -s ${SERVER}/session/{session_id}", timeout=30000)
if "error" in verify_status.lower() or "not found" in verify_status.lower():
print(f"[ERROR] Retry failed - skipping Issue #{issue.number}")
continue
active_issue_workers[issue.number] = session_id active_issue_workers[issue.number] = session_id
slots_available -= 1 slots_available -= 1
print(f"[{now()}] Dispatched issue worker for Issue #{issue.number} - verified active") print(f"[{now()}] Dispatched issue worker for Issue #{issue.number} - session {session_id[:8]}...")
# ── Monitor active workers: poll every 10 seconds ──────────── # ── Monitor active workers: poll every 10 seconds ────────────
bash("sleep 10", timeout=30000) bash("sleep 10", timeout=30000)
# Get session status from server # Get session status from server
try: status_response = bash("curl -s ${SERVER}/session/status", timeout=30000)
status_response = bash("curl -s ${SERVER}/session/status", timeout=30000) all_sessions = safe_json_parse(status_response, [])
all_sessions = json.loads(status_response) if status_response else []
except:
all_sessions = []
# Monitor PR workers # Monitor PR workers
for pr_number, pr_info in list(active_pr_workers.items()): for pr_number, pr_info in list(active_pr_workers.items()):
@@ -666,12 +924,9 @@ LOOP FOREVER:
if not session_active: if not session_active:
# Session completed or died - check final result # Session completed or died - check final result
try: final_response = bash(f"curl -s ${SERVER}/session/{session_id}/messages", timeout=30000)
final_response = bash(f"curl -s ${SERVER}/session/{session_id}/messages", timeout=30000) messages = safe_json_parse(final_response, [])
messages = json.loads(final_response) if final_response else [] final_msg = messages[-1].get("content", "ERROR") if messages else "ERROR: No messages found"
final_msg = messages[-1]["content"] if messages else "ERROR"
except:
final_msg = "ERROR: Could not retrieve final message"
# Parse result # Parse result
if "PR merged successfully" in final_msg: if "PR merged successfully" in final_msg:
@@ -705,12 +960,9 @@ LOOP FOREVER:
if not session_active: if not session_active:
# Session completed or died - check final result # Session completed or died - check final result
try: final_response = bash(f"curl -s ${SERVER}/session/{session_id}/messages", timeout=30000)
final_response = bash(f"curl -s ${SERVER}/session/{session_id}/messages", timeout=30000) messages = safe_json_parse(final_response, [])
messages = json.loads(final_response) if final_response else [] final_msg = messages[-1].get("content", "ERROR") if messages else "ERROR: No messages found"
final_msg = messages[-1]["content"] if messages else "ERROR"
except:
final_msg = "ERROR: Could not retrieve final message"
if "PR merged successfully" in final_msg or "PR created successfully" in final_msg: if "PR merged successfully" in final_msg or "PR created successfully" in final_msg:
# Success - issue implemented and PR created/merged # Success - issue implemented and PR created/merged
@@ -751,6 +1003,11 @@ LOOP FOREVER:
bash(f"curl -s -X DELETE ${SERVER}/session/{session_id}", timeout=15000) bash(f"curl -s -X DELETE ${SERVER}/session/{session_id}", timeout=15000)
del active_issue_workers[issue_number] del active_issue_workers[issue_number]
# ── Periodic maintenance every 5 cycles ──────────────────────
if cycle % 5 == 0:
validate_worker_state()
enforce_worker_limits()
# ── Health signal every 10 cycles ───────────────────────────── # ── Health signal every 10 cycles ─────────────────────────────
if cycle % 10 == 0: if cycle % 10 == 0:
# Helper functions to format worker details # Helper functions to format worker details