Repository navigation
Expand file tree
/
Copy pathdispatch.py
More file actions
2837 lines (2547 loc) · 138 KB
/
Copy pathdispatch.py
File metadata and controls
2837 lines (2547 loc) · 138 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""Dispatch: local workers for an outside coordinator (Fable / Claude / any
API caller), so the expensive model plans and reviews and never reads a tool
transcript.
job = start(owner, {"tasks": [...], "workspace": "D:/proj", ...})
...
compact(job) → a few hundred tokens: per task status, what changed,
the verification, the worker's last words; never the
transcript.
Each job runs the SAME machinery as `/agents` in a chat (delegate_agents:
src/agent_tools/subagent_tools.py — file locks, watchdog, supervisor, the
machine-wide GPU semaphore, lean toolset, the control board), inside a
"Workers" chat of its own so the human can open the board, steer or stop a
worker, and read the transcripts. The model is `dispatch_endpoint_id` /
`dispatch_model` from settings (falls back to the utility, then the default
chat model), or the request's `model` on that endpoint; the request never
names a URL.
What makes the answer trustworthy (none of it comes from the worker):
* evidence — the workspace is checkpointed before the job (the harness's
shadow repo) and diffed after it: `changes` is what really
changed on disk, `claimed_only` what a worker said it changed
but did not;
* verification — Faustus runs the project's tests (or the `verify` command
the coordinator gave) in the workspace after the workers,
compares failures with the checkpoint (pre-existing vs new),
and, when they fail, sends ONE fixer worker with the failure
output before verifying again (`fix_rounds` — a MAXIMUM: the
loop also stops by itself when the rounds stop producing
change, see src/convergence.py);
* proof — the step after observing (src/prove.py): "a mutation is not
the completion of the objective". `result.proof` reconciles
what was observed with what was claimed and answers `proved`,
`partial`, `unproved` (no runner and nothing observable
changed — honest, not a failure) or `contradicted`, with a
confidence and a NAMED reason for every point it is missing;
* honest status — `done` only when every worker finished and the
verification passed (or could not run); otherwise `partial`
with the reason in `verdict`; a cancelled job still reports
what changed.
Jobs in the same workspace run one at a time (a second one waits, `queued`);
a retried POST with the same `Idempotency-Key` returns the first job.
A SEQUENTIAL job whose tasks name objectives (`OBJ-3`) runs them in the order
the project's objectives graph ranks those objectives, not the order they were
typed (`order_tasks_by_impact`, `agent_objective_ordering`). That is the whole
of it: ONE job's task list, recorded in `task_order` and said in the verdict.
There is no queue across jobs and no scheduler here.
Jobs live in memory with a JSON mirror under DATA_DIR/dispatch/ (rotated at
MAX_JOBS_KEPT) so a finished job can still be read after a restart (a
running one is reported as `interrupted`).
Waiting is a CONDITION, not a sleep (`wait_for`): a caller blocks until the
job is done, reaches a phase, a worker enters a state read off its own output
(src/output_rules.py), an event says something, or the workspace changes on
disk — and it resolves the moment the condition holds, because the job's own
progress updates set the `asyncio.Event` the waiter sleeps on. There is no
poll tick inside. A timeout is not an error: it answers `met: False` with how
long it waited, the way the four-value outcomes in this tree already do.
Those same updates feed the live event stream
(`GET /api/dispatch/{id}/events?stream=1`, routes/dispatch_routes.py).
A worker detected `rate_limited` or `waiting_for_input` is REPORTED, never
killed: the state and the literal that proves it land in that worker's
`progress` entry while it runs, and the existing supervisor/ceiling logic
still owns every decision to stop anything.
A task may name a `runner` (src/agent_runners.py): an agent Faustus did not
write — Claude Code, OpenCode, Qwen Code, whatever `ollama launch` knows —
run as a worker through src/external_worker.py. Everything above still
happens around it: the checkpoint before, the diff after, `claimed_only`,
Faustus's own verification, the fix round, the proof. What does NOT happen is
the one thing this app is otherwise careful about: **Faustus's command guard
cannot see inside another agent's own shell**, so a job that used one carries
an explicit `external_agent_unguarded` entry in its proof's uncertainty list
and says so in its verdict. A task with no `runner` runs exactly as it always
has, byte for byte.
"""
from __future__ import annotations
import asyncio
import json
import logging
import os
import re
import time
import uuid
from collections import deque
from typing import Any, Callable, Deque, Dict, List, Optional, Tuple
logger = logging.getLogger(__name__)
MAX_TASKS = 4 # delegate_agents' own cap (MAX_SUBAGENTS)
MAX_JOBS_KEPT = 200
SUMMARY_CHARS = 1200 # the worker's last words, per task
EVENTS_KEPT = 400
CHANGES_LISTED = 60 # paths listed per change kind in the compact result
_DEFAULT_TIMEOUT_S = 900
_DEFAULT_MAX_ROUNDS = 20
_DEFAULT_FIX_ROUNDS = 1
_MAX_FIX_ROUNDS = 2
# With the convergence detector on, the fix loop ends itself as soon as the
# rounds stop producing change, so a caller may ask for more of them without
# buying a fixed number of pointless workers.
_MAX_FIX_ROUNDS_CONVERGENCE = 4
_VERIFY_TIMEOUT_S = 300
_VERIFY_CMD_CHARS = 500
_IDEMPOTENCY_TTL_S = 3600
_LIVE = ("queued", "running", "verifying", "cancelling")
# How much of each worker's own output is kept for the state rules. The rules
# only read their own tail (src/output_rules.py); this bounds what a job holds.
_OUTPUT_TAIL_CHARS = 8192
# The event keys that carry a worker's OWN words (its live command tail, a
# tool's output, its last words, its error) — the only text the rules read.
_OUTPUT_KEYS = ("tail", "output", "final_text", "message")
# A task run by an agent Faustus did not write (src/external_worker.py): the
# named reason that goes into the proof's uncertainty list. The proof exists to
# name every reason its confidence is not 1; an unguarded external shell is the
# biggest one this app can have. The sentence and the penalty live in
# src/prove.py with `note_external_gate`, which now writes the entry — two
# copies of a sentence a reader compares between a payload and a proof is how
# they drift apart.
EXTERNAL_UNGUARDED = "external_agent_unguarded"
# How much of an external agent's output is kept as its "last words".
EXTERNAL_SUMMARY_CHARS = 2000
# `wait_for(condition=...)`: the prefixes a condition may take.
WAIT_CONDITIONS = ("done", "changed", "phase:<name>", "worker:<label>:<state>", "event:<text>")
# `changed` re-reads the workspace on every job update; two scans closer
# together than this reuse the previous answer (a flooding worker must not
# turn one wait into a tree walk per line).
_CHANGED_SCAN_MIN_INTERVAL_S = 0.25
# Live events (SSE): a comment every STREAM_HEARTBEAT_S so no proxy closes an
# idle stream, and the stream itself never outlives the job's own ceiling by
# more than STREAM_MARGIN_S.
STREAM_HEARTBEAT_S = 15.0
STREAM_MARGIN_S = 120.0
_SNAPSHOT_MAX_FILES = 60_000
_SNAPSHOT_SKIP = frozenset({".git", "node_modules", "__pycache__", ".venv", "venv", "env", ".mypy_cache", ".pytest_cache",
".ruff_cache", ".tox", "dist", "build", ".idea", ".vscode", "target", ".next", ".cache",
".faustus_checkpoints", ".odysseus_checkpoints", ".faustus", ".odysseus"})
_jobs: Dict[str, "DispatchJob"] = {}
_lock = asyncio.Lock()
_idempotent: Dict[Tuple[str, str], Tuple[str, float]] = {} # (owner, key) → (job id, ts)
_loaded_all_at = 0.0
def _data_dir() -> str:
try:
from src.constants import DATA_DIR
return os.path.join(DATA_DIR, "dispatch")
except Exception: # pragma: no cover
return os.path.join(os.getcwd(), "data", "dispatch")
def _setting(key: str, default: Any) -> Any:
try:
from src.settings import get_setting
return get_setting(key, default)
except Exception: # noqa: BLE001 - a job never fails over a settings read
return default
def _convergence_on() -> bool:
"""`agent_fix_round_convergence`. Off = the fixed fix-round counter."""
return bool(_setting("agent_fix_round_convergence", True))
def state_detection_on() -> bool:
"""`agent_worker_state_detection`. Off = a worker's output is never read
for a state and `progress` carries exactly what it carried before."""
return bool(_setting("agent_worker_state_detection", True))
def sse_on() -> bool:
"""`agent_dispatch_sse`. Off = `/{id}/events?stream=1` answers with the
same JSON body the endpoint has always returned."""
return bool(_setting("agent_dispatch_sse", True))
def prove_on() -> bool:
"""`agent_dispatch_prove`. Off = the result payload and the verdict line
are exactly what they were before src/prove.py existed."""
try:
from src import prove
return prove.enabled()
except Exception: # noqa: BLE001
return False
def objective_ordering_on() -> bool:
"""`agent_objective_ordering`. Off = the tasks of a job run in the order
they were written, which is what they have always done."""
return bool(_setting("agent_objective_ordering", True))
def _outcomes_on() -> bool:
"""`agent_tool_outcomes`. Off = a stopped worker counts as an error."""
try:
from src import tool_outcome
return tool_outcome.enabled()
except Exception: # noqa: BLE001
return False
def external_runners_on() -> bool:
"""`agent_external_runners`. **Off by default** — it runs third-party
binaries on this machine. Off = a task naming a `runner` is refused with
the reason, and a job without one is untouched."""
try:
from src import agent_runners
return agent_runners.enabled()
except Exception: # noqa: BLE001 - never raise into a hot path
return False
def _max_fix_rounds() -> int:
return _MAX_FIX_ROUNDS_CONVERGENCE if _convergence_on() else _MAX_FIX_ROUNDS
def _timing_key(workspace: Optional[str]) -> str:
"""The adaptive-timeout bucket a job's duration is remembered under: one
per workspace, because how long a job takes is a property of the project."""
return "dispatch:" + (_ws_key(workspace) or "-")
def _adaptive_ceiling(workspace: Optional[str], fixed: int) -> int:
"""The fixed ceiling, raised to what jobs in this workspace really take."""
try:
from src import adaptive_timeout as at
if not at.enabled():
return fixed
key = _timing_key(workspace)
value = at.idle_timeout(key, fixed, lo=fixed, hi=max(fixed, fixed * 3))
if value > fixed:
at.note_difference(key, value, fixed, what="job ceiling")
return int(value)
except Exception as e: # noqa: BLE001 - a poller hint, never load-bearing
logger.debug("dispatch: adaptive ceiling unavailable: %s", e)
return fixed
def _record_job_duration(job: "DispatchJob") -> None:
"""Feed this job's wall-clock into the adaptive recorder (best effort)."""
try:
from src import adaptive_timeout as at
if not at.enabled() or not job.started or not job.finished:
return
at.record(_timing_key(job.workspace), float(job.finished) - float(job.started))
except Exception as e: # noqa: BLE001
logger.debug("dispatch: could not record the job duration: %s", e)
def _round_artifact(job: "DispatchJob", v: Dict[str, Any]) -> str:
"""What one fix round LEFT BEHIND, as the convergence detector reads it:
the verification's verdict and failures, the tail of its output and the
files that changed on disk. Two rounds that leave the same artifact
changed nothing between them."""
changes = job.changes or {}
touched = sorted(list(changes.get("added") or []) + list(changes.get("modified") or [])
+ list(changes.get("deleted") or []))[:CHANGES_LISTED]
parts = [
str(v.get("summary") or ""),
"\n".join(str(f) for f in (v.get("failures") or [])[:20]),
str(v.get("output_tail") or "")[-2000:],
"changed: " + ", ".join(touched),
]
return "\n".join(p for p in parts if p.strip())
def _assess_convergence(rounds: List[str]) -> Optional[Dict[str, Any]]:
try:
from src import convergence
return convergence.assess(rounds)
except Exception as e: # noqa: BLE001 - the fix loop runs without it
logger.debug("dispatch: convergence unavailable: %s", e)
return None
class DispatchJob:
def __init__(self, owner: Optional[str], args: Dict[str, Any], workspace: Optional[str],
endpoint_url: str, model: str, headers: Optional[Dict[str, str]],
title: str, gen_overrides: Optional[Dict[str, Any]] = None,
verify: str = "auto", verify_scope: str = "related", fix_rounds: int = _DEFAULT_FIX_ROUNDS,
verify_timeout_s: float = _VERIFY_TIMEOUT_S, expected_output: Optional[str] = None):
self.id = uuid.uuid4().hex[:12]
self.owner = owner
self.args = args
self.workspace = workspace
self.endpoint_url = endpoint_url
self.model = model
self.headers = headers
self.title = title
self.gen_overrides = gen_overrides
self.verify = verify # "auto" | "none" | a shell command
self.verify_scope = verify_scope # "related" | "all"
self.fix_rounds = fix_rounds
self.verify_timeout_s = verify_timeout_s
# A landmark from the verification's output, declared with the plan.
# It only proves anything because it was written down BEFORE the run:
# a string chosen after seeing the output is a description, not a test.
self.expected_output = expected_output
self.created = time.time()
self.started: Optional[float] = None
self.finished: Optional[float] = None
# queued | running | verifying | done | partial | error | cancelling | cancelled | interrupted
self.status = "queued"
self.error: Optional[str] = None
self.verdict: Optional[str] = None
self.session_id: Optional[str] = None
self.result: Optional[Dict[str, Any]] = None
self.changes: Optional[Dict[str, Any]] = None # observed by Faustus, not claimed by a worker
self.verification: Optional[Dict[str, Any]] = None
self.checkpoint: Optional[str] = None
# PLAN-02: the JOB's own acceptance criteria (docs/spec/v2/backlog.json
# PLAN-02 — "objetivo acotado ... con evidencias requeridas"), distinct
# from a task's per-child `criteria` (src/agent_tools/subagent_tools.py,
# verify_criteria — already checked per worker by SubagentRun.report()).
# A job can pass every worker and still not have delivered what the
# COORDINATOR actually asked for across the whole job (e.g. two workers
# each touch their own file but neither the one the criterion names);
# `_settle` verifies this against the job's own OBSERVED diff, never a
# worker's claim. `None` while the caller declared no job-level
# criteria — that is the common case and leaves `to_dict()` unchanged.
self.criteria_check: Optional[Dict[str, Any]] = None
# Convergence of the fix loop (src/convergence.py) and, when the loop
# ended for a reason other than "the rounds ran out", which one.
self.convergence: Optional[Dict[str, Any]] = None
self.stopped_by: Optional[str] = None
# Spec lint (src/dispatch_spec_lint.py, `dispatch_spec_lint`): what the
# coordinator left out of the spec, attached so the job still runs in
# `warn` mode. None when the spec was complete or the lint is off.
self.spec_lint: Optional[Dict[str, Any]] = None
# Started by unattended work (a night shift, a scheduled run): its
# failure counts for the breaker (src/unattended_breaker.py) and it
# waits while an interactive turn is live (src/period_budget.py).
self.unattended: bool = False
# What the objectives graph did to the task list before the job ran
# ({"by": "impact", "from": [...], "to": [...]}), so the reordering is
# visible and auditable instead of silent. None when nothing moved and
# while `agent_objective_ordering` is off.
self.task_order: Optional[Dict[str, Any]] = None
# The proof packet (src/prove.py): what the evidence and the
# verification really SHOW, with every reason the confidence is not 1
# named. None while `agent_dispatch_prove` is off.
self.proof: Optional[Dict[str, Any]] = None
# The external agents this job ran, if any (src/agent_runners.py). An
# empty list is the normal case and adds NOTHING to the payload: a job
# with no runner is byte-identical to one from before this existed.
self.runners_used: List[str] = []
# What Faustus's own guard did in front of each of them
# (src/agent_gate.py): runner key → the run's ledger. Empty for every
# runner whose row says `gate: "none"`, and an empty one is what keeps
# `prove.note_external_gate` on the blanket answer this path has always
# given — not persisted, because a ledger is about a run, not a job.
self.runner_gates: Dict[str, Dict[str, Any]] = {}
self.events: Deque[Dict[str, Any]] = deque(maxlen=EVENTS_KEPT)
self.task: Optional[asyncio.Task] = None
self._waiters: List[asyncio.Event] = []
self._entered = False # _run has begun (its finally will settle the job)
# What each worker's own output says about it (src/output_rules.py):
# name → {state, why, matched, confidence, seen: [...]}. Never a
# reason to kill anything — it is what `progress` reports.
self.worker_states: Dict[str, Dict[str, Any]] = {}
self._output: Dict[str, str] = {} # name → the tail of its output
# Every event ever appended (the deque rotates at EVENTS_KEPT): a
# stream that has sent N of them knows what is new without stamping
# a sequence number onto the events themselves.
self.events_produced = 0
self._updates: List[asyncio.Event] = [] # woken on every change (wait_for, SSE)
# ── views ────────────────────────────────────────────────────────────
def to_dict(self, *, include_result: bool = True, brief: bool = False) -> Dict[str, Any]:
d = {
"id": self.id, "owner": self.owner, "title": self.title, "status": self.status,
"error": self.error, "verdict": self.verdict, "workspace": self.workspace, "model": self.model,
"session_id": self.session_id, "chat_url": f"/#{self.session_id}" if self.session_id else None,
"created": self.created, "started": self.started, "finished": self.finished,
"duration_s": round((self.finished or time.time()) - (self.started or self.created), 1),
"tasks": [{"name": t.get("name"), "instruction": _squash(t.get("instruction"), 200) if brief else t.get("instruction"),
"files": t.get("files") or [], "model": t.get("model") or None}
for t in self.args.get("tasks") or []],
"parallel": bool(self.args.get("parallel")), "reviewer": bool(self.args.get("reviewer")),
"max_rounds": self.args.get("max_rounds"), "timeout_s": self.args.get("timeout_s"),
"verify": self.verify, "verify_scope": self.verify_scope, "fix_rounds": self.fix_rounds,
}
if self.expected_output:
# On the mirror, so the declaration survives a crash: what a
# recovery plan re-pins has to include the thing the run was going
# to be judged against, or the resumed job proves less than the
# original one would have. Absent when nothing was declared, which
# keeps every existing payload byte-identical.
d["expected_output_contains"] = self.expected_output
if self.task_order is not None:
d["task_order"] = self.task_order
if self.spec_lint is not None:
d["spec_lint"] = self.spec_lint
if self.unattended:
d["unattended"] = True
if self.args.get("criteria"):
# Declared alongside `expected_output_contains` above: what the
# job was pinned to prove, visible even on a brief/queued view.
d["criteria"] = self.args["criteria"]
if include_result:
d["result"] = self.result
d["changes"] = self.changes
d["verification"] = self.verification
d["checkpoint"] = self.checkpoint
if self.criteria_check is not None:
d["criteria_check"] = self.criteria_check
if self.convergence is not None:
d["convergence"] = self.convergence
if self.stopped_by:
d["stopped_by"] = self.stopped_by
if self.proof is not None:
d["proof"] = self.proof
if self.runners_used:
# Only ever present when an external agent really ran: what the
# command guard could not see has to be readable in the payload,
# not only in the proof.
d["runners"] = list(self.runners_used)
d["unguarded"] = True
return d
def ceiling_s(self) -> int:
"""The most wall-clock the job can take: every worker's timeout in
turn at the configured parallelism, a reviewer, the verification and
the fix loop — so a coordinator knows how long to keep waiting.
With `agent_adaptive_idle_timeout` on, jobs that really took longer
than this estimate in this workspace raise it (never lower it: a
coordinator that stops waiting early re-dispatches work that is still
running)."""
tasks = len(self.args.get("tasks") or [])
try:
from src.agent_tools.subagent_tools import _setting as _sa_setting
par = max(1, int(_sa_setting("agent_subagent_max_parallel", 2) or 1))
except Exception:
par = 2
per = int(self.args.get("timeout_s") or _DEFAULT_TIMEOUT_S)
waves = -(-tasks // par) if self.args.get("parallel") else tasks
n = per * max(1, waves) + (per if self.args.get("reviewer") else 0)
if self.verify != "none":
n += int(self.verify_timeout_s) * (1 + self.fix_rounds) + per * self.fix_rounds
return int(_adaptive_ceiling(self.workspace, n))
def _notify(self) -> None:
for ev in self._waiters:
ev.set()
self._waiters.clear()
self._wake()
# ── live updates: what makes a wait resolve at once ──────────────────
def _wake(self) -> None:
"""Wake every condition waiter and every open stream. Called from the
job's own progress path, so a `wait_for` resolves on the update that
made its condition true instead of on a poll tick."""
try:
for ev in list(self._updates):
ev.set()
except Exception as e: # noqa: BLE001 - a notification never breaks a job
logger.debug("dispatch %s: wake failed: %s", self.id, e)
def subscribe(self) -> asyncio.Event:
"""An Event set on every change to this job. The caller MUST call
:meth:`unsubscribe` in a finally (a client that disconnects mid-stream
must not leave a waiter behind)."""
ev = asyncio.Event()
self._updates.append(ev)
return ev
def unsubscribe(self, ev: asyncio.Event) -> None:
try:
self._updates.remove(ev)
except ValueError:
pass
def note_worker_event(self, ev: Dict[str, Any]) -> None:
"""One board event from a worker: kept as it is (the events endpoint
answers exactly what it always did), read for a state, and broadcast."""
self.events.append(ev)
self.events_produced += 1
try:
self._detect(ev)
except Exception as e: # noqa: BLE001 - the rules never break a job
logger.debug("dispatch %s: state detection failed: %s", self.id, e)
self._wake()
def _detect(self, ev: Dict[str, Any]) -> None:
"""Classify the newest words of one worker (`agent_worker_state_detection`).
A `rate_limited` or `waiting_for_input` worker is RECORDED here and
surfaced in `progress`; nothing in this path stops a worker — that
stays with the supervisor and the job ceiling.
"""
if not state_detection_on():
return
name = str(ev.get("name") or "").strip()
if not name or name == "job":
return
chunk = "\n".join(str(ev.get(k)) for k in _OUTPUT_KEYS if ev.get(k)).strip()
if not chunk:
return
from src import output_rules
buf = (self._output.get(name, "") + "\n" + chunk)[-_OUTPUT_TAIL_CHARS:]
self._output[name] = buf
verdict = output_rules.classify_output(buf)
entry = self.worker_states.setdefault(name, {"seen": []})
states = list(verdict.get("states") or [])
if not states:
for key in ("state", "why", "matched", "confidence"):
entry.pop(key, None)
return
matches = verdict.get("matches") or []
entry["state"] = states[0]
entry["states"] = states
entry["why"] = output_rules.why(verdict, states[0])
entry["matched"] = str((matches[0] or {}).get("literal") or "") if matches else ""
entry["confidence"] = verdict.get("confidence")
entry["ts"] = time.time()
for s in states:
if s not in entry["seen"]:
entry["seen"].append(s)
def _persist(self) -> None:
try:
d = _data_dir()
os.makedirs(d, exist_ok=True)
tmp = os.path.join(d, f".{self.id}.tmp")
# The pid goes in the MIRROR only, never in the API payload: it is
# what src/crash_recovery.py probes before calling a job that was
# left `running` interrupted. A pid that still answers means the
# job belongs to a process that is alive, not to a power cut.
doc = dict(self.to_dict(), pid=os.getpid())
with open(tmp, "w", encoding="utf-8") as fh:
json.dump(doc, fh, ensure_ascii=False, indent=1)
os.replace(tmp, os.path.join(d, f"{self.id}.json"))
except Exception as e: # noqa: BLE001 — a mirror, never load-bearing
logger.debug("dispatch: persist failed: %s", e)
def _event(self, **ev: Any) -> None:
ev.setdefault("ts", time.time())
ev.setdefault("name", "job")
self.events.append(ev)
self.events_produced += 1
self._wake()
# ── the compact answer ──────────────────────────────────────────────────────
_WS_RE = re.compile(r"\s+")
def _squash(text: Any, limit: int) -> str:
s = _WS_RE.sub(" ", str(text or "")).strip()
if len(s) <= limit:
return s
return s[: limit - 1].rstrip() + "…"
def _compact_static_checks(sc: Any) -> Any:
"""Only the failures, bounded — a list of 40 `ok: True` rows per worker is
what turned the "few hundred tokens" answer into 14k."""
if isinstance(sc, list):
bad = [x for x in sc if isinstance(x, dict) and x.get("ok") is False]
return {"checked": len(sc), "failed": [{"path": str(x.get("path"))[:120], "error": _squash(x.get("error"), 200)}
for x in bad[:10]]}
if isinstance(sc, dict):
return {k: (_squash(v, 200) if isinstance(v, str) else v) for k, v in list(sc.items())[:8]}
return _squash(sc, 200)
def _compact_refusals(items: Any) -> List[Dict[str, Any]]:
"""A refused tool call as three fields: what it tried, why it was stopped,
how often. Not the sentence the model was handed and not its prose about
it — this projection is why a coordinator costs 1.5k tokens instead of
118k (FAUSTUS.md §21), and a refusal earns its bytes only by being the
thing that stops the next pointless re-dispatch."""
out: List[Dict[str, Any]] = []
for x in list(items or [])[:10]:
if not isinstance(x, dict) or not x.get("tool"):
continue
out.append({"tool": str(x.get("tool"))[:60], "why": _squash(x.get("why"), 120),
"count": int(x.get("count") or 1)})
return out
def compact_from_result(result: Optional[Dict[str, Any]], *, summary_chars: int = SUMMARY_CHARS) -> Dict[str, Any]:
"""What an outside coordinator needs and nothing more: per worker the
status, the files it says it changed, its tool/round/token counts, the
static-check failures and its last words — never the transcript. The
per-worker `git` snapshot is dropped: it was the WHOLE tree's status
repeated once per worker; the job's `changes` block is the honest one."""
out: Dict[str, Any] = {"workers": [], "files_changed": [], "totals": {
"tool_calls": 0, "failed_calls": 0, "rounds": 0, "input_tokens": 0, "output_tokens": 0, "errors": 0}}
if not isinstance(result, dict):
return out
outcomes_on = _outcomes_on()
cancelled = 0
changed: List[str] = []
for r in result.get("subagents") or []:
if not isinstance(r, dict):
continue
w = {
"name": r.get("name"), "role": r.get("role") or "worker", "status": r.get("status"),
"stop_reason": r.get("stop_reason"), "error": _squash(r.get("error"), 300) or None,
"rounds": int(r.get("rounds") or 0), "tool_calls": int(r.get("tool_calls") or 0),
"failed_calls": int(r.get("failed_calls") or 0),
"files_changed": [str(p) for p in list(r.get("mutations") or [])[:40]],
"input_tokens": int(r.get("input_tokens") or 0), "output_tokens": int(r.get("output_tokens") or 0),
"duration_s": r.get("duration_s"), "model": r.get("model"),
"summary": _squash(r.get("final_text"), summary_chars),
"session_id": r.get("session_id"),
}
if outcomes_on:
w["outcome"] = _worker_outcome(r)
if r.get("refusals"):
# Only when there are any: a worker nothing refused keeps the row
# byte-for-byte the one a coordinator has been reading.
refusals = _compact_refusals(r.get("refusals"))
if refusals:
w["refusals"] = refusals
sc = r.get("static_checks")
if sc:
w["static_checks"] = _compact_static_checks(sc)
if r.get("supervisor"):
w["supervisor"] = [_squash(x, 160) for x in list(r.get("supervisor") or [])[:4]]
if r.get("runner"):
# An external agent (src/external_worker.py). These keys exist ONLY
# on such a row: a built-in worker's row is what it always was.
w["runner"] = str(r.get("runner"))
# Read through from the run's own report, not asserted: a runner
# Faustus can gate (src/agent_gate.py) says False, and the compact
# answer must not contradict the worker row it was built from.
w["unguarded"] = bool(r.get("unguarded", True))
if r.get("argv_shown"):
w["argv_shown"] = _squash(r.get("argv_shown"), 400)
if r.get("state"):
w["state"] = str(r.get("state"))
if r.get("why"):
w["why"] = _squash(r.get("why"), 200)
out["workers"].append(w)
for p in w["files_changed"]:
if p not in changed:
changed.append(p)
t = out["totals"]
t["tool_calls"] += w["tool_calls"]
t["failed_calls"] += w["failed_calls"]
t["rounds"] += w["rounds"]
t["input_tokens"] += w["input_tokens"]
t["output_tokens"] += w["output_tokens"]
# any worker that did not finish its task counts — stalled, stopped
# and timed-out workers have no `error` and used to be 0 errors.
# A worker the USER stopped is the exception (four-value outcomes): it
# is `cancelled`, counted apart, and never blamed on the model.
if outcomes_on and w.get("outcome") == "cancelled":
cancelled += 1
elif w["error"] or (w["status"] not in ("done", None)):
t["errors"] += 1
if cancelled:
out["totals"]["cancelled"] = cancelled
out["files_changed"] = changed
if result.get("lock_conflicts"):
out["lock_conflicts"] = [f"{c.get('worker')} → {c.get('path')}" for c in list(result["lock_conflicts"])[:10]
if isinstance(c, dict)]
if result.get("dropped_tasks"):
out["dropped_tasks"] = int(result["dropped_tasks"])
out["exit_code"] = result.get("exit_code")
return out
def _worker_outcome(r: Dict[str, Any]) -> Optional[str]:
"""The four-value outcome of one worker report: the one the worker itself
recorded when it is there, else read from its status and error."""
try:
from src import tool_outcome
known = tool_outcome.value_of(r.get("outcome"))
if known:
return known
return tool_outcome.classify_status(r.get("status") or r.get("stop_reason"),
error=r.get("error")).value
except Exception: # noqa: BLE001 - the compact answer never fails over this
return None
def _seed_progress(job: DispatchJob) -> Dict[str, Dict[str, Any]]:
"""Every task appears in `progress` from the start (`queued`), so a
4-task job at max_parallel 2 never looks like a 2-worker job."""
latest: Dict[str, Dict[str, Any]] = {}
for i, t in enumerate(job.args.get("tasks") or []):
name = str(t.get("name") or f"worker-{i + 1}")
latest[name] = {"last_event": "queued"}
return latest
def _attach_outcome_verification(job: DispatchJob, res: Dict[str, Any]) -> None:
"""The one verified / unverified / uncertain answer for a finished job
(src/outcome_verification.py), fed by what Faustus observed -- the diff,
`claimed_only`, its own verification run, the proof -- and by each worker
row's status and whether an external agent ran outside the command guard.
Additive; `agent_outcome_verification` off leaves the payload as it was."""
try:
from src import outcome_verification
if not outcome_verification.enabled():
return
rows: List[Dict[str, Any]] = [{
"runner": "job", "status": job.status, "exit_code": res.get("exit_code"),
"claimed_only": res.get("claimed_only"), "changes": res.get("changes"),
"verification": res.get("verification"), "proof": res.get("proof"),
}]
for w in res.get("workers") or []:
if isinstance(w, dict):
rows.append({"runner": str(w.get("name") or w.get("runner") or "worker"),
"status": w.get("status"), "error": w.get("error"),
"unguarded": w.get("unguarded") if w.get("runner") else None,
"cancelled": w.get("outcome") == "cancelled"})
res["outcome_verification"] = outcome_verification.verify_outcome(workers=rows)
except Exception: # noqa: BLE001 - a report about the job never breaks it
pass
def compact(job: DispatchJob) -> Dict[str, Any]:
d = job.to_dict(include_result=False)
res = compact_from_result(job.result)
if job.changes is not None:
# what Faustus SAW on disk beats what a worker SAID it wrote
observed = list(job.changes.get("added") or []) + list(job.changes.get("modified") or []) + list(job.changes.get("deleted") or [])
claimed = list(res["files_changed"])
res["files_changed"] = observed[: CHANGES_LISTED * 3]
res["claimed_only"] = [p for p in claimed if not _observed(p, job.changes)][:20]
res["changes"] = job.changes
if job.verification is not None:
res["verification"] = job.verification
if job.convergence is not None:
res["convergence"] = job.convergence
if job.stopped_by:
res["stopped_by"] = job.stopped_by
# The proof of what the job really did (src/prove.py). Absent while
# `agent_dispatch_prove` is off, so the payload stays byte-for-byte the one
# a coordinator has been reading.
if job.proof is not None:
res["proof"] = job.proof
if job.status not in _LIVE:
res["exit_code"] = 0 if job.status == "done" else 1
_attach_outcome_verification(job, res)
d["result"] = res
# a running job: the board's latest tick per worker, so a poller sees
# progress without the event stream
if job.status in _LIVE:
latest = _seed_progress(job)
for ev in job.events:
name = str(ev.get("name") or ev.get("id") or "")
if not name or name == "job":
continue
kind = ev.get("event")
if kind in ("queued", "started", "round", "tool", "tick", "steer", "supervisor", "done", "error", "guard", "harness"):
cur = latest.setdefault(name, {})
cur["last_event"] = kind
for k in ("round", "elapsed_s", "idle_s", "last_tool", "tool", "stalled", "stall_reason", "status"):
if k in ev:
cur[k] = ev[k]
# What the worker's own output says about it, while it runs rather
# than after (`agent_worker_state_detection`; off = nothing is added).
for name, st in (job.worker_states or {}).items():
if not st.get("state") or name not in latest:
continue
latest[name]["state"] = st["state"]
if st.get("why"):
latest[name]["why"] = st["why"]
d["progress"] = latest
d["wait_again"] = True
d["ceiling_s"] = job.ceiling_s()
for ev in reversed(job.events):
if ev.get("name") == "job" and ev.get("message"):
d["phase"] = ev["message"]
break
return d
def _observed(path: str, changes: Dict[str, Any]) -> bool:
p = str(path or "").replace("\\", "/").strip("/").lower()
for kind in ("added", "modified", "deleted"):
for q in changes.get(kind) or []:
qq = str(q).replace("\\", "/").strip("/").lower()
if qq == p or qq.endswith("/" + p) or p.endswith("/" + qq):
return True
return bool(changes.get("truncated"))
# ── evidence: what changed on disk ──────────────────────────────────────────
def _snapshot(workspace: str) -> Tuple[Dict[str, Tuple[int, int]], bool]:
"""rel path → (mtime_ns, size) for the workspace tree (bounded, skips the
usual generated folders). The fallback when the harness's checkpoints are
unavailable (no git on the box)."""
files: Dict[str, Tuple[int, int]] = {}
truncated = False
root = os.path.realpath(workspace)
stack = [root]
while stack:
cur = stack.pop()
try:
with os.scandir(cur) as it:
for entry in it:
try:
if entry.is_dir(follow_symlinks=False):
if entry.name not in _SNAPSHOT_SKIP:
stack.append(entry.path)
elif entry.is_file(follow_symlinks=False):
st = entry.stat(follow_symlinks=False)
rel = os.path.relpath(entry.path, root).replace("\\", "/")
files[rel] = (int(st.st_mtime_ns), int(st.st_size))
if len(files) >= _SNAPSHOT_MAX_FILES:
return files, True
except OSError:
continue
except OSError:
continue
return files, truncated
def _diff_snapshots(before: Dict[str, Tuple[int, int]], after: Dict[str, Tuple[int, int]], truncated: bool) -> Dict[str, Any]:
added = sorted(p for p in after if p not in before)
deleted = sorted(p for p in before if p not in after)
modified = sorted(p for p in after if p in before and after[p] != before[p])
return _changes_block(added, modified, deleted, "mtime", truncated)
def _changes_block(added: List[str], modified: List[str], deleted: List[str], source: str, truncated: bool = False) -> Dict[str, Any]:
n = len(added) + len(modified) + len(deleted)
return {"source": source, "count": n, "added": added[:CHANGES_LISTED], "modified": modified[:CHANGES_LISTED],
"deleted": deleted[:CHANGES_LISTED],
"truncated": bool(truncated or n > CHANGES_LISTED * 3 or max(len(added), len(modified), len(deleted)) > CHANGES_LISTED)}
def _checkpoint(workspace: str, label: str) -> Optional[str]:
try:
from src import workspace_checkpoints as wc
if not wc.enabled():
return None
cp = wc.checkpoint(workspace, label=label)
return str(cp.get("sha")) if cp and cp.get("sha") else None
except Exception as e: # noqa: BLE001
logger.debug("dispatch: checkpoint failed: %s", e)
return None
def _changes_since(workspace: str, sha: str) -> Optional[Dict[str, Any]]:
try:
from src import workspace_checkpoints as wc
rows = wc.changed_since(workspace, sha)
except Exception as e: # noqa: BLE001
logger.debug("dispatch: changed_since failed: %s", e)
return None
if rows is None:
return None
added = sorted(r["path"] for r in rows if r.get("status") == "A")
deleted = sorted(r["path"] for r in rows if r.get("status") == "D")
modified = sorted(r["path"] for r in rows if r.get("status") not in ("A", "D"))
block = _changes_block(added, modified, deleted, "checkpoint")
block["checkpoint"] = sha[:12]
return block
def _git_facts(workspace: str) -> Optional[Dict[str, Any]]:
"""The user's own repo, once per job: is it one, how dirty is it now."""
try:
from src.agent_harness import git_change_summary
g = git_change_summary(workspace)
except Exception:
return None
if not g:
return None
return {"repo": True, "dirty_count": int(g.get("changed_count") or 0), "shortstat": _squash(g.get("shortstat"), 200)}
class _Evidence:
"""Before/after the job: a checkpoint (content-exact, via the harness's
shadow repo) or an mtime snapshot of the tree."""
def __init__(self, workspace: Optional[str], job_id: str):
self.workspace = workspace
self.sha: Optional[str] = None
self.snap: Optional[Dict[str, Tuple[int, int]]] = None
self.snap_truncated = False
self.label = f"dispatch {job_id}"
def before(self) -> None:
if not self.workspace:
return
self.sha = _checkpoint(self.workspace, self.label)
if not self.sha:
self.snap, self.snap_truncated = _snapshot(self.workspace)
def after(self) -> Optional[Dict[str, Any]]:
if not self.workspace:
return None
block = _changes_since(self.workspace, self.sha) if self.sha else None
if block is None and self.snap is not None:
now, trunc = _snapshot(self.workspace)
block = _diff_snapshots(self.snap, now, self.snap_truncated or trunc)
if block is not None:
git = _git_facts(self.workspace)
if git:
block["git"] = git
return block
# ── verification: Faustus runs the proof, not the worker ────────────────────
def _verification_spec(workspace: str, verify: str) -> Tuple[Optional[Dict[str, Any]], str]:
from src import project_tests as pt
if verify == "none":
return None, "off"
if verify and verify != "auto":
return pt.detect_test_command(workspace, override=verify), "command"
override = ""
try:
from src.settings import get_setting
override = str(get_setting("agent_project_test_command", "") or "").strip()
except Exception:
override = ""
return pt.detect_test_command(workspace, override=override or None), "auto"
def run_verification(workspace: Optional[str], verify: str, changed: List[str], *, scope: str = "related",
timeout_s: float = _VERIFY_TIMEOUT_S, checkpoint_sha: Optional[str] = None,
expected_output: Optional[str] = None) -> Dict[str, Any]:
"""Run the project's tests (or `verify`) in the workspace, bounded, and
return a compact verdict. `ok` is None when nothing could be run — that is
"not verified", never "passed".
`expected_output` is a landmark the request DECLARED when it asked for the
job, and it travels down to `project_tests.run_tests` unchanged. Exit 0 is
not evidence a suite ran — a collection that found nothing reports it too —
so a declaration made after the run would prove nothing. This one was made
before it (src/output_oracle.py)."""
if not workspace:
return {"mode": "off", "ran": False, "ok": None, "summary": "no workspace"}
try:
spec, mode = _verification_spec(workspace, verify)
except Exception as e: # noqa: BLE001
return {"mode": "auto", "ran": False, "ok": None, "summary": f"verification unavailable: {e}"[:200]}
if mode == "off":
return {"mode": "off", "ran": False, "ok": None, "summary": "verification disabled by the request (verify: none)"}
if not spec:
return {"mode": mode, "ran": False, "ok": None,
"summary": "no test runner detected in the workspace (give `verify` a command that proves the task)"}
spec["expected_output_contains"] = expected_output
from src import project_tests as pt
res = pt.run_tests(workspace, spec, changed=list(changed or []), scope=scope, timeout_s=timeout_s)
if res.get("ran") and res.get("ok") is False and not res.get("inconclusive") and checkpoint_sha:
try:
res = pt.compare_with_baseline(workspace, checkpoint_sha, spec, res, changed=list(changed or []))
except Exception as e: # noqa: BLE001
logger.debug("dispatch: baseline comparison failed: %s", e)
out = {
"mode": mode, "ran": bool(res.get("ran")), "ok": res.get("ok"), "inconclusive": bool(res.get("inconclusive")),
"kind": res.get("kind"), "command": _squash(res.get("command"), 300), "scope": res.get("scope"),
"exit_code": res.get("exit_code"), "timed_out": bool(res.get("timed_out")), "duration_s": res.get("duration_s"),
"summary": _squash(res.get("summary"), 300), "failures": [_squash(f, 200) for f in (res.get("failures") or [])[:10]],
"output_tail": str(res.get("output_tail") or "")[-1500:],
}
if res.get("output_matched") is not None:
# Only ever present when something was declared. Without it a run the
# oracle overturned reads as `ok: False` beside a summary saying "47
# passed", and the reason is nowhere in the payload.
out["output_matched"] = bool(res["output_matched"])
out["expected_output_contains"] = str(expected_output or "")
if res.get("related_files"):
out["related_files"] = [str(p) for p in res["related_files"][:12]]
for k in ("new_failures", "pre_existing"):
if res.get(k):
out[k] = [_squash(f, 200) for f in res[k][:10]]
if res.get("pre_existing_only"):
out["pre_existing_only"] = True
return out
def verification_failed(v: Optional[Dict[str, Any]]) -> bool:
"""A verdict that should block `done`: it ran, it failed, it is
conclusive, and the failures are not all pre-existing."""
return bool(v and v.get("ran") and v.get("ok") is False and not v.get("inconclusive") and not v.get("pre_existing_only"))
def _fixer_instruction(job: DispatchJob, v: Dict[str, Any], attempt: int) -> str:
"""What the fix round asks for — the same words whether or not the worker
is resumed.
It was tempting to drop the task recap for a resumed worker, since its own
session already carries it. It is not safe: this side decides to resume
OPTIMISTICALLY and the worker side is the one that finds out whether the
session still has any history (it may have been pruned, or the manager may
keep none). Trimming here would produce the one outcome worse than today's
— a thin instruction AND no recovered context. Resume stays purely