|
11 | 11 | import shutil |
12 | 12 | import subprocess |
13 | 13 | import sys |
| 14 | +import threading |
14 | 15 | import time |
15 | 16 | from pathlib import Path |
16 | 17 |
|
17 | 18 | from core import count_tokens |
| 19 | +from services.circuit_breaker import CircuitBreaker |
18 | 20 |
|
19 | 21 | log = logging.getLogger(__name__) |
20 | 22 |
|
@@ -252,11 +254,20 @@ def _handle_claude_delegate(task: str, task_type: str, context: str, |
252 | 254 | file_path: str, svc, dcfg: dict, finalize) -> str: |
253 | 255 | """Handle delegation via Claude Code CLI.""" |
254 | 256 | timeout = int(dcfg.get("claude_timeout", 90)) |
| 257 | + breaker = _backend_breaker("claude", dcfg) |
| 258 | + if not breaker.allow(): |
| 259 | + return finalize("c3_delegate", {"task_type": task_type, "backend": "claude"}, |
| 260 | + "[delegate:degraded] Claude skipped after repeated failures; retrying in " |
| 261 | + f"~{breaker.cooldown_remaining()}s. Run 'claude --version' to diagnose.", |
| 262 | + "degraded") |
255 | 263 | _log_progress(svc, f"[delegate] Routing {task_type} → Claude CLI...") |
256 | 264 | output, ok = _run_claude(task, context, cwd=str(svc.project_path), timeout=timeout) |
257 | 265 | if not ok: |
| 266 | + if breaker.record_failure(): |
| 267 | + _notify_backend_degraded(svc, "claude", breaker) |
258 | 268 | return finalize("c3_delegate", {"task_type": task_type, "backend": "claude"}, |
259 | 269 | output, "error") |
| 270 | + breaker.record_success() |
260 | 271 | return finalize("c3_delegate", {"task_type": task_type, "backend": "claude"}, |
261 | 272 | output, "ok") |
262 | 273 |
|
@@ -681,6 +692,51 @@ def _run_codex_resume(follow_up: str, timeout: int = 120, |
681 | 692 | _delegate_cache: dict[str, tuple[str, int]] = {} |
682 | 693 | _delegate_metrics = {"total_calls": 0, "tokens_saved": 0} |
683 | 694 |
|
| 695 | +# Per-backend runtime circuit breakers. Distinct from the install-status flags |
| 696 | +# (_gemini_available etc., which only answer "is the CLI on PATH"): these track |
| 697 | +# *runtime* health so a broken-but-installed backend (expired auth, repeated |
| 698 | +# timeouts) stops re-spawning a 90-120s subprocess on every call. Keyed by |
| 699 | +# backend name and intentionally process-global — backend health (auth, CLI |
| 700 | +# version) is a property of the host, not of any single project. |
| 701 | +_backend_breakers: dict[str, CircuitBreaker] = {} |
| 702 | +_backend_breakers_lock = threading.Lock() |
| 703 | + |
| 704 | + |
| 705 | +def _backend_breaker(name: str, dcfg: dict | None = None) -> CircuitBreaker: |
| 706 | + """Return (creating on first use) the runtime circuit breaker for a backend.""" |
| 707 | + with _backend_breakers_lock: |
| 708 | + breaker = _backend_breakers.get(name) |
| 709 | + if breaker is None: |
| 710 | + cfg = dcfg or {} |
| 711 | + breaker = CircuitBreaker( |
| 712 | + name, |
| 713 | + failure_threshold=int(cfg.get("breaker_failure_threshold", 3) or 3), |
| 714 | + cooldown_seconds=float(cfg.get("breaker_cooldown_seconds", 60) or 60), |
| 715 | + ) |
| 716 | + _backend_breakers[name] = breaker |
| 717 | + return breaker |
| 718 | + |
| 719 | + |
| 720 | +def _notify_backend_degraded(svc, name: str, breaker: CircuitBreaker) -> None: |
| 721 | + """Surface a backend trip via the NotificationStore (best-effort, never raises).""" |
| 722 | + notifications = getattr(svc, "notifications", None) |
| 723 | + if notifications is None: |
| 724 | + return |
| 725 | + try: |
| 726 | + notifications.add( |
| 727 | + agent="c3", |
| 728 | + severity="warning", |
| 729 | + title=f"Delegate backend degraded: {name}", |
| 730 | + message=( |
| 731 | + f"{name} failed {breaker.failure_threshold}x consecutively; c3_delegate " |
| 732 | + f"will skip it for ~{int(breaker.cooldown_seconds)}s instead of re-spawning " |
| 733 | + f"the CLI. Run '{name} --version' to diagnose." |
| 734 | + ), |
| 735 | + replace_if_unacked=True, |
| 736 | + ) |
| 737 | + except Exception: |
| 738 | + pass |
| 739 | + |
684 | 740 |
|
685 | 741 | def get_delegate_metrics() -> dict: |
686 | 742 | return dict(_delegate_metrics) |
@@ -765,6 +821,13 @@ def _handle_codex_delegate(task: str, task_type: str, context: str, |
765 | 821 | "[delegate:error] Codex CLI not available. Run 'codex --version' to diagnose.", |
766 | 822 | "unavailable") |
767 | 823 |
|
| 824 | + breaker = _backend_breaker("codex", dcfg) |
| 825 | + if not breaker.allow(): |
| 826 | + return finalize("c3_delegate", {"task_type": task_type, "backend": "codex"}, |
| 827 | + "[delegate:degraded] Codex skipped after repeated failures; retrying in " |
| 828 | + f"~{breaker.cooldown_remaining()}s. Run 'codex --version' to diagnose.", |
| 829 | + "degraded") |
| 830 | + |
768 | 831 | # Resolve model/sandbox/reasoning from config or defaults |
769 | 832 | cdef = CODEX_MODELS.get(task_type, CODEX_MODELS.get("ask", {})) |
770 | 833 | model = dcfg.get("codex_default_model") or cdef.get("model", "gpt-5.3-codex-spark") |
@@ -807,10 +870,13 @@ def _handle_codex_delegate(task: str, task_type: str, context: str, |
807 | 870 | elapsed = round(time.monotonic() - t0, 1) |
808 | 871 |
|
809 | 872 | if not ok: |
| 873 | + if breaker.record_failure(): |
| 874 | + _notify_backend_degraded(svc, "codex", breaker) |
810 | 875 | return finalize("c3_delegate", |
811 | 876 | {"task_type": task_type, "backend": "codex", "model": model, "elapsed": f"{elapsed}s"}, |
812 | 877 | output, "error") |
813 | 878 |
|
| 879 | + breaker.record_success() |
814 | 880 | _delegate_metrics["total_calls"] += 1 |
815 | 881 | _delegate_cache[ckey] = (output, count_tokens(output)) |
816 | 882 |
|
@@ -880,6 +946,13 @@ def _handle_gemini_delegate(task: str, task_type: str, context: str, |
880 | 946 | "[delegate:error] Gemini CLI not available. Run 'gemini --version' to diagnose.", |
881 | 947 | "unavailable") |
882 | 948 |
|
| 949 | + breaker = _backend_breaker("gemini", dcfg) |
| 950 | + if not breaker.allow(): |
| 951 | + return finalize("c3_delegate", {"task_type": task_type, "backend": "gemini"}, |
| 952 | + "[delegate:degraded] Gemini skipped after repeated failures; retrying in " |
| 953 | + f"~{breaker.cooldown_remaining()}s. Run 'gemini --version' to diagnose.", |
| 954 | + "degraded") |
| 955 | + |
883 | 956 | # Resolve model from config or defaults |
884 | 957 | gdef = GEMINI_MODELS.get(task_type, GEMINI_MODELS.get("ask", {})) |
885 | 958 | model = dcfg.get("gemini_default_model") or gdef.get("model", "gemini-2.5-flash") |
@@ -919,10 +992,13 @@ def _handle_gemini_delegate(task: str, task_type: str, context: str, |
919 | 992 | elapsed = round(time.monotonic() - t0, 1) |
920 | 993 |
|
921 | 994 | if not ok: |
| 995 | + if breaker.record_failure(): |
| 996 | + _notify_backend_degraded(svc, "gemini", breaker) |
922 | 997 | return finalize("c3_delegate", |
923 | 998 | {"task_type": task_type, "backend": "gemini", "model": model, "elapsed": f"{elapsed}s"}, |
924 | 999 | output, "error") |
925 | 1000 |
|
| 1001 | + breaker.record_success() |
926 | 1002 | _delegate_metrics["total_calls"] += 1 |
927 | 1003 | _delegate_cache[ckey] = (output, count_tokens(output)) |
928 | 1004 |
|
@@ -1068,9 +1144,11 @@ def _check_claude(): |
1068 | 1144 | _gemini_avail = (_gemini_available is True) or ( |
1069 | 1145 | _gemini_available is None and task_type not in _light_tasks and _is_gemini_on_path() |
1070 | 1146 | ) |
1071 | | - if task_type in heavy_codex and _codex_avail and _codex_available is not False: |
| 1147 | + if (task_type in heavy_codex and _codex_avail and _codex_available is not False |
| 1148 | + and _backend_breaker("codex", dcfg).allow()): |
1072 | 1149 | backend = "codex" |
1073 | | - elif task_type in heavy_gemini and _gemini_avail and _gemini_available is not False: |
| 1150 | + elif (task_type in heavy_gemini and _gemini_avail and _gemini_available is not False |
| 1151 | + and _backend_breaker("gemini", dcfg).allow()): |
1074 | 1152 | backend = "gemini" |
1075 | 1153 | else: |
1076 | 1154 | backend = "ollama" |
|
0 commit comments