Skip to content

Graph and state#

The agent topology: nodes wired by edges, entered at a hub. See Customize agents and tools for the narrative version.

adda.Graph dataclass #

Agent graph passed to :class:~agent_runtime.AgenticRun.

Declares which agents exist, which directed delegation edges connect them, and which agent starts the run. Loops are permitted.

Parameters:

Name Type Description Default
nodes dict[str, Agent]

Maps unique agent names to :class:Agent instances.

required
edges sequence of Edge

Directed delegation edges. An agent with no outgoing edges receives no Delegate tool.

()
entry str

Name of the agent that receives the initial briefing.

None

Raises:

Type Description
ValueError

If any edge endpoint names an undeclared node, or if entry is not declared.

Source code in src/adda/_src/backends/base.py
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
@dataclass
class Graph:
    """Agent graph passed to :class:`~agent_runtime.AgenticRun`.

    Declares which agents exist, which directed delegation edges connect
    them, and which agent starts the run.  Loops are permitted.

    Parameters
    ----------
    nodes : dict[str, Agent]
        Maps unique agent names to :class:`Agent` instances.
    edges : sequence of Edge
        Directed delegation edges.  An agent with no outgoing edges
        receives no ``Delegate`` tool.
    entry : str
        Name of the agent that receives the initial briefing.

    Raises
    ------
    ValueError
        If any edge endpoint names an undeclared node, or if *entry* is
        not declared.
    """

    nodes: dict  # dict[str, Agent]
    edges: tuple = ()
    entry: str = None  # type: ignore[assignment]

    def __post_init__(self) -> None:
        self.edges = tuple(self.edges)
        if self.entry is None:
            raise ValueError("Graph.entry is required and must not be None.")
        bad = [k for k, v in self.nodes.items() if not isinstance(v, Agent)]
        if bad:
            raise TypeError(
                f"Graph nodes values must be Agent instances; "
                f"got invalid values for keys: {sorted(bad)}"
            )
        names = set(self.nodes)
        for e in self.edges:
            if e.source not in names or e.target not in names:
                raise ValueError(
                    f"Edge {e!r} references undeclared node. "
                    f"Declared names: {sorted(names)}"
                )
        if self.entry not in names:
            raise ValueError(
                f"entry={self.entry!r} not in nodes (entry node undeclared). "
                f"Declared names: {sorted(names)}"
            )
        missing_desc = [n for n, a in self.nodes.items() if not a.description]
        if missing_desc:
            raise ValueError(
                f"All agents must define a non-empty description. "
                f"Missing in: {sorted(missing_desc)}"
            )

    def outgoing(self, name: str) -> list[str]:
        """Return target names for all edges out of *name*."""
        return [e.target for e in self.edges if e.source == name]

    def incoming(self, name: str) -> list[str]:
        """Return source names for all edges into *name*."""
        return [e.source for e in self.edges if e.target == name]

    def edge(self, source: str, target: str) -> Edge | None:
        """Return the Edge from *source* to *target*, or None if absent."""
        for e in self.edges:
            if e.source == source and e.target == target:
                return e
        return None

    # Curated semantic colours (fill, stroke) by agent class name;
    # unknown classes cycle a fallback palette so any graph stays legible.
    _MERMAID_COLOURS = {
        "StrategizerAgent": ("#1d4ed8", "#1e3a8a"),
        "LiteratureReviewAgent": ("#6d28d9", "#4c1d95"),
        "DataGeneratorAgent": ("#b45309", "#7c2d12"),
        "F3dasmImplementerAgent": ("#15803d", "#14532d"),
        "AdversarialCritiqueAgent": ("#b91c1c", "#7f1d1d"),
        "DebuggerAgent": ("#475569", "#1e293b"),
    }
    _MERMAID_FALLBACK = [
        ("#0f766e", "#134e4a"), ("#a16207", "#713f12"),
        ("#be185d", "#831843"), ("#4338ca", "#312e81"),
    ]

    def to_mermaid(self) -> str:
        """Return a styled Mermaid flowchart for this graph.

        Generated from the live spec: node labels carry each agent's
        class, role, and a short description; nodes are coloured by agent
        class; edges from the entry node render as solid delegation
        arrows and edges from worker nodes as dotted consultation arrows.
        Paste at https://mermaid.live or any Mermaid-aware renderer
        (GitHub markdown, Jupyter, VS Code).
        """
        def _clean(text: str, n: int = 46) -> str:
            t = " ".join(str(text).split()).replace('"', "'")
            return (t[: n - 1] + "…") if len(t) > n else t

        lines = ["flowchart TD"]
        colours: dict = {}
        members: dict = {}
        fb = 0
        for name in self.nodes:
            agent = self.nodes[name]
            cls = type(agent).__name__
            if cls not in colours:
                if cls in self._MERMAID_COLOURS:
                    colours[cls] = self._MERMAID_COLOURS[cls]
                else:
                    colours[cls] = self._MERMAID_FALLBACK[
                        fb % len(self._MERMAID_FALLBACK)]
                    fb += 1
            members.setdefault(cls, []).append(name)
            desc = _clean(
                getattr(agent, "description", "") or agent.role)
            label = f"<b>{name}</b><br/>{cls}<br/><i>{desc}</i>"
            if name == self.entry:
                lines.append(f'    {name}(["{label}"])')
            else:
                lines.append(f'    {name}["{label}"]')
        for e in self.edges:
            arrow = "-->" if e.source == self.entry else "-.->"
            if e.preamble:
                snippet = _clean(e.preamble, 35)
                lines.append(
                    f'    {e.source} {arrow}|"{snippet}"| {e.target}')
            else:
                lines.append(f'    {e.source} {arrow} {e.target}')
        for cls, (fill, stroke) in colours.items():
            lines.append(
                f"    classDef {cls} fill:{fill},stroke:{stroke},"
                "color:#fff,stroke-width:1px;")
        for cls, names in members.items():
            lines.append(f"    class {','.join(names)} {cls}")
        return "\n".join(lines)

    def __repr__(self) -> str:
        out_map: dict[str, list] = {}
        for e in self.edges:
            out_map.setdefault(e.source, []).append(e)
        lines = [f"Graph(entry={self.entry!r})"]
        for name in self.nodes:
            tag = " [entry]" if name == self.entry else ""
            edges = out_map.get(name, [])
            if edges:
                for e in edges:
                    preamble = f"  # {e.preamble!r}" if e.preamble else ""
                    lines.append(f"  {name}{tag}  ──▶  {e.target}{preamble}")
                    tag = ""  # only label first edge row
            else:
                lines.append(f"  {name}{tag}  (leaf)")
        return "\n".join(lines)
nodes: dict instance-attribute #
edges: tuple = () class-attribute instance-attribute #
entry: str = None class-attribute instance-attribute #
_MERMAID_COLOURS = {'StrategizerAgent': ('#1d4ed8', '#1e3a8a'), 'LiteratureReviewAgent': ('#6d28d9', '#4c1d95'), 'DataGeneratorAgent': ('#b45309', '#7c2d12'), 'F3dasmImplementerAgent': ('#15803d', '#14532d'), 'AdversarialCritiqueAgent': ('#b91c1c', '#7f1d1d'), 'DebuggerAgent': ('#475569', '#1e293b')} class-attribute instance-attribute #
_MERMAID_FALLBACK = [('#0f766e', '#134e4a'), ('#a16207', '#713f12'), ('#be185d', '#831843'), ('#4338ca', '#312e81')] class-attribute instance-attribute #
outgoing(name: str) -> list[str] #

Return target names for all edges out of name.

Source code in src/adda/_src/backends/base.py
450
451
452
def outgoing(self, name: str) -> list[str]:
    """Return target names for all edges out of *name*."""
    return [e.target for e in self.edges if e.source == name]
incoming(name: str) -> list[str] #

Return source names for all edges into name.

Source code in src/adda/_src/backends/base.py
454
455
456
def incoming(self, name: str) -> list[str]:
    """Return source names for all edges into *name*."""
    return [e.source for e in self.edges if e.target == name]
edge(source: str, target: str) -> Edge | None #

Return the Edge from source to target, or None if absent.

Source code in src/adda/_src/backends/base.py
458
459
460
461
462
463
def edge(self, source: str, target: str) -> Edge | None:
    """Return the Edge from *source* to *target*, or None if absent."""
    for e in self.edges:
        if e.source == source and e.target == target:
            return e
    return None
to_mermaid() -> str #

Return a styled Mermaid flowchart for this graph.

Generated from the live spec: node labels carry each agent's class, role, and a short description; nodes are coloured by agent class; edges from the entry node render as solid delegation arrows and edges from worker nodes as dotted consultation arrows. Paste at https://mermaid.live or any Mermaid-aware renderer (GitHub markdown, Jupyter, VS Code).

Source code in src/adda/_src/backends/base.py
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
def to_mermaid(self) -> str:
    """Return a styled Mermaid flowchart for this graph.

    Generated from the live spec: node labels carry each agent's
    class, role, and a short description; nodes are coloured by agent
    class; edges from the entry node render as solid delegation
    arrows and edges from worker nodes as dotted consultation arrows.
    Paste at https://mermaid.live or any Mermaid-aware renderer
    (GitHub markdown, Jupyter, VS Code).
    """
    def _clean(text: str, n: int = 46) -> str:
        t = " ".join(str(text).split()).replace('"', "'")
        return (t[: n - 1] + "…") if len(t) > n else t

    lines = ["flowchart TD"]
    colours: dict = {}
    members: dict = {}
    fb = 0
    for name in self.nodes:
        agent = self.nodes[name]
        cls = type(agent).__name__
        if cls not in colours:
            if cls in self._MERMAID_COLOURS:
                colours[cls] = self._MERMAID_COLOURS[cls]
            else:
                colours[cls] = self._MERMAID_FALLBACK[
                    fb % len(self._MERMAID_FALLBACK)]
                fb += 1
        members.setdefault(cls, []).append(name)
        desc = _clean(
            getattr(agent, "description", "") or agent.role)
        label = f"<b>{name}</b><br/>{cls}<br/><i>{desc}</i>"
        if name == self.entry:
            lines.append(f'    {name}(["{label}"])')
        else:
            lines.append(f'    {name}["{label}"]')
    for e in self.edges:
        arrow = "-->" if e.source == self.entry else "-.->"
        if e.preamble:
            snippet = _clean(e.preamble, 35)
            lines.append(
                f'    {e.source} {arrow}|"{snippet}"| {e.target}')
        else:
            lines.append(f'    {e.source} {arrow} {e.target}')
    for cls, (fill, stroke) in colours.items():
        lines.append(
            f"    classDef {cls} fill:{fill},stroke:{stroke},"
            "color:#fff,stroke-width:1px;")
    for cls, names in members.items():
        lines.append(f"    class {','.join(names)} {cls}")
    return "\n".join(lines)

adda.Edge dataclass #

A directed delegation edge between two named agents.

Parameters:

Name Type Description Default
source str

Name of the agent that is allowed to call Delegate.

required
target str

Name of the agent that receives the delegated task.

required
preamble str

Text prepended to the task message when this edge is traversed.

''
Source code in src/adda/_src/backends/base.py
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
@dataclass(frozen=True)
class Edge:
    """A directed delegation edge between two named agents.

    Parameters
    ----------
    source : str
        Name of the agent that is allowed to call ``Delegate``.
    target : str
        Name of the agent that receives the delegated task.
    preamble : str
        Text prepended to the task message when this edge is traversed.
    """

    source: str
    target: str
    preamble: str = ""
source: str instance-attribute #
target: str instance-attribute #
preamble: str = '' class-attribute instance-attribute #

adda.Agent #

Base class for all agentic nodes in a Graph.

Subclasses override class-level attributes (system_prompt, tools, etc.) to configure behaviour. Behavioural differences belong in class attributes, not constructor arguments.

Tool system — three categories:

tools: frozenset[str] declares tools from two categories:

  1. Native backend tools — names from :data:NATIVE_TOOL_NAMES ("Bash", "Read", "Write", etc.). The runtime passes these to the backend session's native tool executor.

  2. Protocol closure tools — names from :data:PROTOCOL_CLOSURE_NAMES ("Done", "WriteNote", "ReadNote"). The runtime builds Python callables for these and passes them as closure_tools to the session factory.

  3. Topology-injected tools — "Delegate", "Parallel", "Debate", "Retry" (outgoing edges), "Ask" (entry node), and "SendMessage" (every node, regardless of edges). Never declare these in Agent.tools. The runtime injects them automatically from the graph topology; any declaration here is ignored.

  4. External MCP server tools — names declared in extra_allowed_tools (e.g. "arxiv_search_papers"). The runtime passes these to the backend together with mcp_servers, a dict of {server_name: McpStdioServerConfig} that declares which external MCP servers to start.

Default is frozenset() — no tools (opt-in, conservative).

Parameters:

Name Type Description Default
model str or None

Model identifier. None delegates to the backend default.

None
Source code in src/adda/_src/backends/base.py
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
class Agent:
    """Base class for all agentic nodes in a Graph.

    Subclasses override class-level attributes (``system_prompt``,
    ``tools``, etc.) to configure behaviour.  Behavioural differences
    belong in class attributes, not constructor arguments.

    **Tool system — three categories:**

    ``tools: frozenset[str]`` declares tools from two categories:

    1. **Native backend tools** — names from :data:`NATIVE_TOOL_NAMES`
       (``"Bash"``, ``"Read"``, ``"Write"``, etc.).  The runtime passes
       these to the backend session's native tool executor.

    2. **Protocol closure tools** — names from
       :data:`PROTOCOL_CLOSURE_NAMES` (``"Done"``, ``"WriteNote"``,
       ``"ReadNote"``).  The runtime builds Python callables for these and
       passes them as ``closure_tools`` to the session factory.

    3. **Topology-injected tools** — ``"Delegate"``, ``"Parallel"``,
       ``"Debate"``, ``"Retry"`` (outgoing edges), ``"Ask"`` (entry node),
       and ``"SendMessage"`` (every node, regardless of edges).  **Never
       declare these in** ``Agent.tools``.  The runtime injects them
       automatically from the graph topology; any declaration here is
       ignored.

    4. **External MCP server tools** — names declared in
       ``extra_allowed_tools`` (e.g.
       ``"arxiv_search_papers"``).  The runtime passes these to the
       backend together with ``mcp_servers``, a dict of
       ``{server_name: McpStdioServerConfig}`` that declares which external
       MCP servers to start.

    Default is ``frozenset()`` — no tools (opt-in, conservative).

    Parameters
    ----------
    model : str or None
        Model identifier.  ``None`` delegates to the backend default.
    """

    system_prompt: str = ""
    tools: frozenset[str] = frozenset()
    reset_on_checkpoint: bool = True
    description: str = ""
    # Neutral default: a subclass that forgets to declare its role must NOT
    # silently inherit "implementer" and pick up implementer-only behavior
    # (the milestone gate, the eval-parallelism nudge). Every shipped agent
    # declares its own role; this default only guards future ones.
    role: str = "worker"
    # Declared as a CLASS attribute, like its sibling `backend`, so the
    # documented idiom above ("behavioural differences belong in class
    # attributes") actually works for it. It previously existed only as an
    # instance attribute assigned in __init__, so a subclass setting
    # `model = "..."` had it silently overwritten with None while a subclass
    # setting `backend = "..."` was honoured — the asymmetry produced the
    # worst possible outcome, a node pinned to one backend and left on
    # another backend's default model id.
    model: str | None = None
    backend: str | None = None
    #: Endpoint override (OpenAI-compatible backends); config.yaml sets it per node.
    base_url: str | None = None
    #: ``DEFAULT_PROMPT`` keeps the backend's own system prompt and appends
    #: this node's text to it; None replaces it. config.yaml sets it per node.
    base_prompt: str | None = None
    mcp_servers: dict = {}
    extra_allowed_tools: frozenset[str] = frozenset()
    max_history_pairs: int = 5
    report_sections: tuple[str, ...] = (
        "### Actions taken",
        "### Conclusions",
        "### Numbers",
    )

    def __init__(self, model: str | None = None) -> None:
        # Only an EXPLICIT argument overrides the class attribute; assigning
        # unconditionally would re-introduce the silent clobber described
        # above, since the default is indistinguishable from "not passed".
        if model is not None:
            self.model = model

    def forward(self) -> None:
        """ADAS hook — override for inspectable Python orchestration."""

    def build_closure_tools(
        self,
        study_dir: Path,
        delegation_id: str | None = None,
        lit_reviewer_notes_dir: Path | None = None,
    ) -> dict:
        """Return runtime closure tools for this agent.

        Called by the runtime when constructing the worker adapter so agents can
        inject Python callables (e.g. corpus management tools) without declaring
        them in Agent.tools.

        The default gives EVERY agent read-only lookup against the study's
        persistent literature corpus (ConsultLiterature) —
        the corpus is shared, queryable infrastructure (see LiteratureCorpus),
        the same way QueryStore lets every node read the canonical evaluation
        ledger without delegating to the data generator. ACQUIRING a new paper
        (CorpusAdd, external search) stays literature_reviewer-only: finding
        and vetting a new paper needs judgment a raw tool call can't supply,
        so it stays gated behind an actual delegation — see
        LiteratureReviewAgent.build_closure_tools, which overrides this
        method entirely (the same ConsultLiterature plus CorpusAdd) and
        does not call super().

        Override in a subclass to replace this default entirely.
        """
        from pathlib import Path as _Path

        try:
            from ..literature.literature_corpus import LiteratureCorpus
        except ImportError:
            return {}

        corpus_dir = (
            _Path(lit_reviewer_notes_dir) if lit_reviewer_notes_dir is not None
            else _Path(study_dir) / "runs" / "lit_reviewer_notes"
        )
        corpus = LiteratureCorpus(corpus_dir)
        from ..agents.literature_tools.corpus import build_corpus_read_closures
        return build_corpus_read_closures(corpus)
system_prompt: str = '' class-attribute instance-attribute #
tools: frozenset[str] = frozenset() class-attribute instance-attribute #
reset_on_checkpoint: bool = True class-attribute instance-attribute #
description: str = '' class-attribute instance-attribute #
role: str = 'worker' class-attribute instance-attribute #
model: str | None = None class-attribute instance-attribute #
backend: str | None = None class-attribute instance-attribute #
base_url: str | None = None class-attribute instance-attribute #
base_prompt: str | None = None class-attribute instance-attribute #
mcp_servers: dict = {} class-attribute instance-attribute #
extra_allowed_tools: frozenset[str] = frozenset() class-attribute instance-attribute #
max_history_pairs: int = 5 class-attribute instance-attribute #
report_sections: tuple[str, ...] = ('### Actions taken', '### Conclusions', '### Numbers') class-attribute instance-attribute #
forward() -> None #

ADAS hook — override for inspectable Python orchestration.

Source code in src/adda/_src/backends/base.py
324
325
def forward(self) -> None:
    """ADAS hook — override for inspectable Python orchestration."""
build_closure_tools(study_dir: Path, delegation_id: str | None = None, lit_reviewer_notes_dir: Path | None = None) -> dict #

Return runtime closure tools for this agent.

Called by the runtime when constructing the worker adapter so agents can inject Python callables (e.g. corpus management tools) without declaring them in Agent.tools.

The default gives EVERY agent read-only lookup against the study's persistent literature corpus (ConsultLiterature) — the corpus is shared, queryable infrastructure (see LiteratureCorpus), the same way QueryStore lets every node read the canonical evaluation ledger without delegating to the data generator. ACQUIRING a new paper (CorpusAdd, external search) stays literature_reviewer-only: finding and vetting a new paper needs judgment a raw tool call can't supply, so it stays gated behind an actual delegation — see LiteratureReviewAgent.build_closure_tools, which overrides this method entirely (the same ConsultLiterature plus CorpusAdd) and does not call super().

Override in a subclass to replace this default entirely.

Source code in src/adda/_src/backends/base.py
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
def build_closure_tools(
    self,
    study_dir: Path,
    delegation_id: str | None = None,
    lit_reviewer_notes_dir: Path | None = None,
) -> dict:
    """Return runtime closure tools for this agent.

    Called by the runtime when constructing the worker adapter so agents can
    inject Python callables (e.g. corpus management tools) without declaring
    them in Agent.tools.

    The default gives EVERY agent read-only lookup against the study's
    persistent literature corpus (ConsultLiterature) —
    the corpus is shared, queryable infrastructure (see LiteratureCorpus),
    the same way QueryStore lets every node read the canonical evaluation
    ledger without delegating to the data generator. ACQUIRING a new paper
    (CorpusAdd, external search) stays literature_reviewer-only: finding
    and vetting a new paper needs judgment a raw tool call can't supply,
    so it stays gated behind an actual delegation — see
    LiteratureReviewAgent.build_closure_tools, which overrides this
    method entirely (the same ConsultLiterature plus CorpusAdd) and
    does not call super().

    Override in a subclass to replace this default entirely.
    """
    from pathlib import Path as _Path

    try:
        from ..literature.literature_corpus import LiteratureCorpus
    except ImportError:
        return {}

    corpus_dir = (
        _Path(lit_reviewer_notes_dir) if lit_reviewer_notes_dir is not None
        else _Path(study_dir) / "runs" / "lit_reviewer_notes"
    )
    corpus = LiteratureCorpus(corpus_dir)
    from ..agents.literature_tools.corpus import build_corpus_read_closures
    return build_corpus_read_closures(corpus)

adda.Node #

One node in the agent graph.

inspect.getsource(Node._orchestrate) reads the routing logic every node runs.

Parameters:

Name Type Description Default
adapter Any

The LLM adapter this node speaks through.

required
name str

The node's name in the graph. This is what distinguishes one agent from another — there is no per-agent node class.

'node'
outgoing list[str] | tuple[str, ...]

Names of the nodes this one may delegate to. Empty means Delegate() has nowhere to send work — the only thing that differs structurally.

()
spec Any

The whole :class:~..backends.base.Graph, so a node can read its peers' roles and descriptions.

None
report_sections tuple[str, ...] | None

This agent's own declared report contract and toolset. agent_tools gates which capability closures :meth:_init_capabilities grants; report_sections is retained for callers that validate a worker's report against its own contract (parsing._classify_response).

None
agent_tools tuple[str, ...] | None

This agent's own declared report contract and toolset. agent_tools gates which capability closures :meth:_init_capabilities grants; report_sections is retained for callers that validate a worker's report against its own contract (parsing._classify_response).

None
Source code in src/adda/_src/nodes/node.py
 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
class Node(
    RecordingMixin,
    CriticGateMixin,
    LifecycleMixin,
    ReproductionGateMixin,
    OrchestrationMixin,
    StopMixin,
):
    """One node in the agent graph.

    ``inspect.getsource(Node._orchestrate)`` reads the routing logic every
    node runs.

    Parameters
    ----------
    adapter
        The LLM adapter this node speaks through.
    name
        The node's name in the graph. This is what distinguishes one agent
        from another — there is no per-agent node class.
    outgoing
        Names of the nodes this one may delegate to. Empty means ``Delegate()``
        has nowhere to send work — the only thing that differs structurally.
    spec
        The whole :class:`~..backends.base.Graph`, so a node can read its
        peers' roles and descriptions.
    report_sections, agent_tools
        This agent's own declared report contract and toolset. ``agent_tools``
        gates which capability closures :meth:`_init_capabilities` grants;
        ``report_sections`` is retained for callers that validate a worker's
        report against its own contract (``parsing._classify_response``).
    """

    _awake_slots = None  # run-wide AwakeSlots, set by graph_builder

    def __init__(
        self,
        adapter: Any,
        *,
        name: str = "node",
        outgoing: list[str] | tuple[str, ...] = (),
        spec: Any = None,
        study_dir: Any = None,
        interactive: bool = False,
        max_ask: int = 1,
        worker_adapters: dict | None = None,
        notes_dir: Any = None,
        workspace_dir: Any = None,
        delegation_log: DelegationLog | None = None,
        report_sections: tuple[str, ...] | None = None,
        agent_tools: frozenset[str] | None = None,
    ) -> None:
        self.adapter = adapter
        self.adapter.on_retry = self._record_llm_retry
        self._name = name
        self._outgoing = list(outgoing)
        self._spec = spec
        self._init_recording()
        # Universal capability setup — every node's adapter, regardless of
        # whether it has outgoing edges (see the module docstring for why
        # this matters even for a node that never takes a turn of its own).
        self._init_capabilities(
            study_dir=study_dir, workspace_dir=workspace_dir,
            delegation_log=delegation_log,
            report_sections=report_sections, agent_tools=agent_tools,
        )
        # Delegation registry, ledgers and routing closures. Harmless for a
        # node with no outgoing edges (Delegate() simply has no valid
        # target); epistemic OWNERSHIP still requires outgoing edges (see
        # _init_orchestration / _install_epistemics) — a leaf never acquires
        # the ledgers just because notes_dir was passed to every node.
        self._init_orchestration(
            name=name, outgoing=self._outgoing, spec=spec,
            study_dir=study_dir, interactive=interactive, max_ask=max_ask,
            worker_adapters=worker_adapters, notes_dir=notes_dir,
            workspace_dir=workspace_dir, delegation_log=delegation_log,
        )

    def __call__(self, state: AgenticState) -> Any:
        """Take one turn. Every node runs the same delegate-or-close loop."""
        return self._orchestrate(state)

    def _raise_if_abandoned(self) -> None:
        """End the calling thread once the run has stopped waiting for it."""
        if self._abandon.is_set():
            raise RunAbandoned(self._name)
        raise_if_stopped(self._name)

    def _init_recording(self) -> None:
        """Establish the state ``RecordingMixin`` writes to, on EVERY node.

        Recording is a property of a node — any node's tools can return an
        ERROR, and any node's adapter reports token usage — so its substrate
        belongs here rather than in one behaviour's setup. It used to be
        created only by the orchestration setup, which was invisible while
        leaves were a separate class that did not carry RecordingMixin at all:
        the moment a leaf reached any recording call it raised AttributeError
        on ``_registry_lock``. The orchestration setup still assigns these
        itself, to the same values, so an orchestrating node is unchanged.
        """
        self._registry_lock = threading.Lock()
        self._abandon = threading.Event()
        self._notifications: list[str] = []
        self._notifications_lock = threading.Lock()
        self._error_counts: dict[str, int] = {}
        self._token_totals: dict = {
            "input_tokens": 0,
            "output_tokens": 0,
            "cache_read_input_tokens": 0,
            "cache_creation_input_tokens": 0,
            # The backend-independent schema (infra/telemetry.py): `legacy_calls`
            # counts calls whose backend reported none of it.
            "fresh_input": 0,
            "cache_read": 0,
            "cache_write": 0,
            "output": 0,
            "normalized_calls": 0,
            "legacy_calls": 0,
            "total_cost_usd": 0.0,
        }
        self._cost_observed: bool = False
        self._current_notes_dir: Path | None = None
        self._telemetry: Any = None
        # Time-budget wrap-up ladder (nodes/_constants.py:budget_band_due):
        # every node's OWN 10%-of-budget bands already reported, so an
        # escalating message fires once per band whether this node
        # orchestrates or answers — a property of any node, like recording.
        self._budget_bands_fired: set[int] = set()

    def _init_capabilities(
        self,
        *,
        study_dir: Any = None,
        workspace_dir: Any = None,
        delegation_log: DelegationLog | None = None,
        report_sections: tuple[str, ...] | None = None,
        agent_tools: frozenset[str] | None = None,
    ) -> None:
        """Wire this node's sandboxed Write and its declared read-only tools.

        Runs for EVERY node, not only one reachable as a standalone entry
        point — see the module docstring: whichever node ends up as someone
        else's ``Delegate()`` target is dispatched through a copy of THIS adapter
        (``adapter.copy()`` carries the ``closure_tools`` set here), so this is where a
        dispatched specialist's baseline capabilities actually come from,
        before ``WorkerSession.install_worker_tools``/``_sandbox_worker_writes``
        layer the per-delegation overrides (ReportEvals's own record callback,
        a delegation-scoped Write) on top.
        """
        self._study_dir = study_dir
        self._workspace_dir = Path(workspace_dir) if workspace_dir else None
        # The agent's declared tools — the single source of truth for which
        # capability closures this node is granted (read-only ledger/store
        # tools). Kept as a frozenset for membership checks.
        self._agent_tools: frozenset[str] = frozenset(agent_tools or ())
        # A node built without a declared toolset has nothing to narrow by.
        self._tools_declared = agent_tools is not None
        # This agent's declared report sections (e.g. the critic's
        # Findings/Verdict, not the implementer's Conclusions/Files touched).
        # Used to validate a worker's report against ITS OWN contract instead
        # of the implementer-shaped default — otherwise a correct critic or
        # literature report is wrongly flagged malformed (audit BF-10/O40).
        self._report_sections = report_sections
        self._delegation_log = delegation_log
        self._evals_reported: dict = {}
        self._setup_sandboxed_write()
        self.adapter.closure_tools.update(self._build_eval_closures())
        if delegation_log is not None:
            from .tools.routing import build_recall_history
            self.adapter.closure_tools["RecallHistory"] = build_recall_history(self)
        # Declaration-gated read-only ledger/store tools — the SAME builder
        # every node uses, so a specialist (e.g. the critic) gets an
        # identical, working QueryStore/HypothesisList
        # surface whenever it declares them. Resolves the run via the shared
        # Node._resolve_run_dir (delegation-log path).
        from ..runtime import features as _features
        from .tools.routing import build_declared_shared_closures
        self.adapter.closure_tools.update(
            build_declared_shared_closures(
                self, self._agent_tools - _features.disabled_tool_names()))
        # Closures from Agent.build_closure_tools() (ConsultLiterature,
        # CorpusAdd, ...) are installed on the adapter BEFORE this node
        # exists (agent_runtime.py), so they never pass through
        # orchestration.py's _wrap_closure the way routing tools do. A
        # closure tagged _adda_diagnostic_source (a fact about the run
        # environment a tool discovers, not about one call's return value —
        # see LiteratureCorpus.pop_diagnostic_event) needs that wrapper
        # regardless, so re-wrap it here, once, by the tag rather than by
        # name — any future tool needing the same reporting just carries the
        # same tag.
        for _cname, _cfn in list(self.adapter.closure_tools.items()):
            if getattr(_cfn, "_adda_diagnostic_source", None) is not None:
                self.adapter.closure_tools[_cname] = self._wrap_closure(
                    _cfn, self._name)
        # A config.yaml `nodes.<name>.tools` list is the whole tool set: the
        # always-on closures it does not name are withheld.
        from ..runtime.node_tools import withheld_closures
        _agent = self._spec.nodes.get(self._name) if self._spec is not None else None
        for _cname in withheld_closures(_agent):
            self.adapter.closure_tools.pop(_cname, None)

    def _setup_sandboxed_write(self) -> None:
        """Replace native Write with a workspace-sandboxed closure.

        Removes 'Write' from native_tools so the SDK doesn't expose it,
        then installs the same builder a dispatched worker gets
        (:func:`build_sandboxed_write`), scoped to this node's whole
        workspace instead of one delegation's subfolder — a node reached via
        real graph routing (the module docstring) has no delegation_id of
        its own to scope tighter than that.

        Bash is kept native but cwd is already set to study_dir by the adapter;
        the prompt further constrains it to the workspace.
        """
        if self._workspace_dir is None:
            return  # no sandboxing if study_dir unknown (e.g. tests)

        # Remove native Write so the SDK doesn't expose an unrestricted version
        if hasattr(self.adapter, "native_tools") and "Write" in self.adapter.native_tools:
            self.adapter.native_tools = [
                t for t in self.adapter.native_tools if t != "Write"
            ]

        from .tools.routing import build_sandboxed_write

        # Bound to a plain name (not inlined into the assignment) so
        # internal/tools/promptmap.py's injected_tool_docs() scanner — which
        # resolves a closure_tools["Write"] = <name> rebind to the function
        # <name> refers to — can still find this tool's docstring after the
        # move into the shared builder.
        _study_ws = (
            Path(self._study_dir) / "workspace"
            if self._study_dir is not None else None
        )
        Write = build_sandboxed_write(
            self._workspace_dir, study_workspace=_study_ws)
        self.adapter.closure_tools["Write"] = Write

    def _build_eval_closures(self) -> dict:
        from .tools.routing import build_report_evals

        evals = self._evals_reported
        return {
            "ReportEvals": build_report_evals(
                record=lambda n: evals.__setitem__("count", n)
            )
        }

    # ── Run-context resolution (shared by every node) ────────────────────────
    # The read tools (QueryStore/HypothesisList) may be granted
    # to any node — the entry node, a mid-tier delegating node, or a leaf such
    # as the critic. They need the run's store/ledger paths, which are resolved
    # here so the tools work identically wherever they are granted.
    def _commit_workspace(self, message: str) -> str | None:
        """Commit the run workspace and return the sha (spec 11).

        On Node rather than on WorkerSession because a delegation is recorded
        from four places, not one: a worker finishing (ok or error), the
        critic gate, an AskForFeedback audit, and the close-time
        reconciliation of a delegation still running when the run ended. Every
        one of them appends a row to the delegation log, so every one of them
        owes that row a sha — otherwise ``workspace_sha: null`` means both
        "this delegation changed nothing" and "nobody looked", which is the
        ambiguity --allow-empty exists to prevent.

        The interrupted case is the one with teeth: a delegation killed
        mid-flight never reaches _finish_ok/_finish_error, so its partial file
        writes stay uncommitted and are swept into whichever delegation
        commits NEXT — attributing one delegation's work to another. Silence
        would be better than that; a commit is better still.

        Never raises — see infra/workspace_vcs.
        """
        run_dir = self._resolve_run_dir()
        if run_dir is None:
            return None
        from ..infra.workspace_vcs import commit_workspace
        return commit_workspace(run_dir / "debug" / "delegations", message)

    def _resolve_run_dir(self) -> Path | None:
        """Best-effort run_dir, valid on any node.

        Every node's ``_current_notes_dir`` is set at construction (from
        ``notes_dir``, passed to every node — see ``graph_builder.py``) and
        re-pointed each turn by ``_absorb_state`` if the run's real notes dir
        differs. It can still be ``None`` (e.g. a bare ``Node()`` built
        without one, as in unit tests) — every node DOES hold the shared
        delegation log at run_dir/debug/delegation_log.jsonl, so derive
        run_dir from that when the notes dir is unavailable.
        """
        notes = getattr(self, "_current_notes_dir", None)
        if notes is not None:
            return Path(notes).parent.parent
        dlog = getattr(self, "_delegation_log", None)
        p = getattr(dlog, "_path", None)
        return Path(p).parent.parent if p is not None else None

    def _read_ledger(self) -> HypothesisLedger | None:
        """The hypothesis ledger for READ access, resolved on any node.

        Returns the node's own bound ledger when it has one (the entry node);
        otherwise resolves a read-only view from the run's strategizer_notes
        when hypotheses.json exists. HypothesisLedger.__init__ performs no I/O,
        and only READ callers use this, so there is no write race with the
        entry node that owns the file.
        """
        own = getattr(self, "_ledger", None)
        if own is not None:
            return own
        rd = self._resolve_run_dir()
        if rd is None:
            return None
        notes = rd / "debug" / "strategizer_notes"
        if (notes / "hypotheses.json").exists():
            from ..epistemics.hypothesis_ledger import HypothesisLedger
            return HypothesisLedger(notes)
        return None
_current_run_dir property #

This run's root, derived from the notes dir the graph state sets.

run_dir arrives on the state, not on the node, and only _current_notes_dir (<run>/debug/strategizer_notes) is kept from it — so the run root is that path's grandparent. None before the first invoke, which callers must tolerate.

_awake_slots = None class-attribute instance-attribute #
adapter = adapter instance-attribute #
_name = name instance-attribute #
_outgoing = list(outgoing) instance-attribute #
_spec = spec instance-attribute #
_stop_tick() -> str #

Raw notice text for the entry node, "" when there is nothing new (the caller marks it as an adda notice).

Idempotent and cheap: one small file read until a request is seen. Only the node that runs a turn (the one holding the run's start time) acts; a worker node's checkpoints leave the request alone.

Source code in src/adda/_src/nodes/stop.py
 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
def _stop_tick(self) -> str:
    """Raw notice text for the entry node, "" when there is nothing new
    (the caller marks it as an adda notice).

    Idempotent and cheap: one small file read until a request is seen.
    Only the node that runs a turn (the one holding the run's start
    time) acts; a worker node's checkpoints leave the request alone.
    """
    stop = self._stop
    if stop is None:
        run_dir = self._current_run_dir
        if run_dir is None or self._run_start is None:
            return ""
        req = read_stop_request(run_dir, since=self._run_start)
        if req is None:
            return ""
        stop = self._stop = {
            **req,
            "deadline": time.time() + req["grace_s"],
            "cancelled": False,
        }
        wound = self._stop_wind_down()
        if req.get("termination") == terminal.CRASHED:
            self._log_lost_workers()
        self._record_intervention(
            "RUN_STOP", "(run)",
            f"stop requested by {req['by']}"
            + (f": {req['reason']}" if req["reason"] else "")
            + f"; winding down {wound or 'no delegations'}, "
            f"grace {req['grace_s']:g}s",
        )
        return (
            f"[RUN STOP — requested by {req['by']}"
            + (f" ({req['reason']})" if req["reason"] else "")
            + ". New delegations are refused. "
            + (f"Winding down {', '.join(wound)}: Wait() for their "
               "reports, then " if wound else "")
            + "call Done() with what you have — it will skip the "
            "usual gates and ask only for your retrospective.]")
    if not stop["cancelled"] and time.time() >= stop["deadline"]:
        cancelled = self._stop_cancel_stragglers()
        if cancelled:
            return (
                "[RUN STOP — grace expired; cancelled "
                f"{', '.join(cancelled)} (no report). Call Done() now.]")
    return ""
_stop_nodes() -> list[Any] #
Source code in src/adda/_src/nodes/stop.py
144
145
def _stop_nodes(self) -> list[Any]:
    return [self, *getattr(self, "_peers", {}).values()]
_stop_wind_down() -> list[str] #

Tell every live delegation to report now. Returns their ids.

Source code in src/adda/_src/nodes/stop.py
147
148
149
150
151
152
153
154
155
156
157
158
159
160
def _stop_wind_down(self) -> list[str]:
    """Tell every live delegation to report now. Returns their ids."""
    ids: list[str] = []
    for n in self._stop_nodes():
        with n._registry_lock:
            live = [d for d, e in n._registry.items()
                    if e.get("status") in _LIVE]
        ids.extend(live)
    for did in ids:
        for n in self._stop_nodes():
            with n._pending_worker_msgs_lock:
                n._pending_worker_msgs.setdefault(did, []).append(
                    wind_down_notice(self._stop))
    return sorted(ids)
_stop_cancel_stragglers() -> list[str] #

Past the grace: detach what has not reported, and say so.

No placeholder retrospective is written for a straggler — it never gave one, and the record must say that rather than invent it.

Source code in src/adda/_src/nodes/stop.py
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
def _stop_cancel_stragglers(self) -> list[str]:
    """Past the grace: detach what has not reported, and say so.

    No placeholder retrospective is written for a straggler — it never
    gave one, and the record must say that rather than invent it.
    """
    self._stop["cancelled"] = True
    cancelled: list[str] = []
    now = datetime.now(tz=timezone.utc).isoformat(timespec="seconds")
    for n in self._stop_nodes():
        with n._registry_lock:
            for did, e in n._registry.items():
                if e.get("status") in _LIVE:
                    e["status"] = "Cancelled"
                    cancelled.append((did, e))
    for did, e in cancelled:
        if self._delegation_log is not None:
            self._delegation_log.record(
                id=did,
                from_node=(
                    self._name if e.get("parent") in (None, "entry")
                    else e["parent"]),
                to_node=e.get("target") or "",
                task=str(e.get("task", "")),
                hypothesis_ids=list(e.get("hypothesis_ids") or []),
                deliverable=(
                    "run stop: cancelled after the "
                    f"{self._stop['grace_s']:g}s grace without a "
                    "report; no retrospective was given"),
                started_at=e.get("started_at") or now, completed_at=now,
                status="CANCELLED",
                tokens_in=0, tokens_out=0, cost_usd=None,
            )
    ids = sorted(d for d, _ in cancelled)
    if ids:
        self._record_intervention(
            "RUN_STOP_CANCELLED", "(run)",
            f"cancelled after grace, no report: {', '.join(ids)}")
    return ids
_stop_refusal() -> str | None #

The Delegate refusal while a stop is active, else None.

Source code in src/adda/_src/nodes/stop.py
202
203
204
205
206
207
208
209
def _stop_refusal(self) -> str | None:
    """The ``Delegate`` refusal while a stop is active, else None."""
    if not any(n._stop is not None for n in self._stop_nodes()):
        return None
    return (
        "run stop: new delegations are refused. This delegation "
        "was NOT started. Wait() for any still reporting, then call "
        "Done().")
_stop_consume() -> None #
Source code in src/adda/_src/nodes/stop.py
211
212
213
214
def _stop_consume(self) -> None:
    run_dir = self._current_run_dir
    if run_dir is not None:
        consume_stop_request(run_dir)
_request_wind_down(*, reason: str, termination: str, run_dir: Any = None, grace_s: float = DEFAULT_GRACE_S) -> bool #

A backstop asks for the graceful path instead of halting. False when it cannot (no run directory, or the file cannot be written), in which case the caller halts as it always did.

Source code in src/adda/_src/nodes/stop.py
216
217
218
219
220
221
222
223
224
225
226
227
228
def _request_wind_down(
    self, *, reason: str, termination: str, run_dir: Any = None,
    grace_s: float = DEFAULT_GRACE_S,
) -> bool:
    """A backstop asks for the graceful path instead of halting. False when
    it cannot (no run directory, or the file cannot be written), in which
    case the caller halts as it always did."""
    run_dir = Path(run_dir) if run_dir else self._current_run_dir
    if run_dir is None or self._run_start is None:
        return False
    return write_stop_request(
        run_dir, by="backstop", reason=reason, grace_s=grace_s,
        termination=termination)
_wind_down_overdue() -> bool #

The grace for the workers plus an equal one for the entry node's own retrospective round has passed: the hard halt may fire.

Source code in src/adda/_src/nodes/stop.py
230
231
232
233
234
def _wind_down_overdue(self) -> bool:
    """The grace for the workers plus an equal one for the entry node's own
    retrospective round has passed: the hard halt may fire."""
    stop = self._stop
    return bool(stop) and time.time() > stop["deadline"] + stop["grace_s"]
_stop_termination() -> str #
Source code in src/adda/_src/nodes/stop.py
236
237
def _stop_termination(self) -> str:
    return (self._stop or {}).get("termination") or terminal.STOPPED
_recorded_retrospective_ids() -> set[str] #
Source code in src/adda/_src/nodes/stop.py
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
def _recorded_retrospective_ids(self) -> set[str]:
    import json

    run_dir = self._current_run_dir
    have: set[str] = set()
    if run_dir is None:
        return have
    try:
        with open(run_dir / "debug" / "retrospectives.jsonl",
                  encoding="utf-8") as f:
            for line in f:
                try:
                    have.add(str(json.loads(line).get("source_id")))
                except (json.JSONDecodeError, AttributeError):
                    continue
    except OSError:
        pass
    return have
_missing_retrospectives() -> list[str] #

Delegations that ended without a recorded retrospective.

Source code in src/adda/_src/nodes/stop.py
258
259
260
261
262
263
264
265
266
267
268
269
270
271
def _missing_retrospectives(self) -> list[str]:
    """Delegations that ended without a recorded retrospective."""
    run_dir = self._current_run_dir
    if run_dir is None:
        return []
    have = self._recorded_retrospective_ids()
    ids: set[str] = set()
    for n in self._stop_nodes():
        with n._registry_lock:
            ids.update(n._registry)
    missing = sorted(ids - have)
    if "DONE" not in have:
        missing.append("DONE")
    return missing
_log_lost_workers() -> None #

A resume after a crash: delegations the delegation log last saw RUNNING died with the process. They are named, never given invented retrospective text.

Source code in src/adda/_src/nodes/stop.py
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
def _log_lost_workers(self) -> None:
    """A resume after a crash: delegations the delegation log last saw
    RUNNING died with the process. They are named, never given invented
    retrospective text."""
    log = getattr(self, "_delegation_log", None)
    if log is None:
        return
    have = self._recorded_retrospective_ids()
    lost = sorted(
        r["id"] for r in log.query_all()
        if r.get("status") == "RUNNING" and r["id"] not in have)
    if lost:
        self._record_intervention(
            "RETROSPECTIVES_MISSING", "(run)",
            "process lost: no retrospective from "
            f"{', '.join(lost)} (they were running when the process "
            "died; none is synthesized)")
_log_missing_retrospectives(why: str) -> None #
Source code in src/adda/_src/nodes/stop.py
291
292
293
294
295
296
def _log_missing_retrospectives(self, why: str) -> None:
    missing = self._missing_retrospectives()
    if missing:
        self._record_intervention(
            "RETROSPECTIVES_MISSING", "(run)",
            f"{why}; no retrospective from: {', '.join(missing)}")
_init_orchestration(*, name: str, outgoing: list[str], spec: Any, study_dir: Any = None, interactive: bool = False, max_ask: int = 1, worker_adapters: dict | None = None, notes_dir: Any = None, workspace_dir: Any = None, delegation_log: DelegationLog | None = None) -> None #

Set up the delegation registry, the ledgers and the routing tools.

Runs for every node (Node.__init__ no longer branches on outgoing). The registry/routing state below is inert when outgoing is empty — Delegate() just has no valid target — so there is no harm running it unconditionally. Epistemic OWNERSHIP is the one thing that must stay gated on having outgoing edges (see _owns_epistemics below): notes_dir is passed identically to every node by graph_builder.py, and a node with nobody to delegate to must never acquire the hypothesis/milestone ledgers just because it was handed a path.

study_dir/workspace_dir/delegation_log are already set by :meth:Node._init_capabilities, which runs first for every node — not re-set here.

Source code in src/adda/_src/nodes/orchestration.py
 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
def _init_orchestration(
    self,
    *,
    name: str,
    outgoing: list[str],
    spec: Any,
    study_dir: Any = None,
    interactive: bool = False,
    max_ask: int = 1,
    worker_adapters: dict | None = None,
    notes_dir: Any = None,
    workspace_dir: Any = None,
    delegation_log: DelegationLog | None = None,
) -> None:
    """Set up the delegation registry, the ledgers and the routing tools.

    Runs for every node (``Node.__init__`` no longer branches on
    ``outgoing``). The registry/routing state below is inert when
    ``outgoing`` is empty — ``Delegate()`` just has no valid target — so
    there is no harm running it unconditionally. Epistemic OWNERSHIP is
    the one thing that must stay gated on having outgoing edges (see
    ``_owns_epistemics`` below): ``notes_dir`` is passed identically to
    every node by ``graph_builder.py``, and a node with nobody to
    delegate to must never acquire the hypothesis/milestone ledgers just
    because it was handed a path.

    ``study_dir``/``workspace_dir``/``delegation_log`` are already set by
    :meth:`Node._init_capabilities`, which runs first for every node —
    not re-set here.
    """
    self._route: dict = {}
    self._interactive = interactive
    self._max_ask = max_ask
    self._ask_count = 0
    self._current_notes_dir: Path | None = (
        Path(notes_dir) if notes_dir is not None else None
    )
    # Parallel delegation registry: id → {"status", "result", "evals", "hypothesis_ids", "started_at"}
    self._worker_adapters: dict[str, Any] = worker_adapters or {}
    self._registry: dict[str, dict] = {}
    self._registry_lock = threading.Lock()
    self._threads: dict[str, threading.Thread] = {}
    # SendMessage/review (spec 12, peer_interaction): the live
    # WorkerSession for each OPEN-FOR-REVIEW delegation, so an
    # approving SendMessage can finalize it (call its own _finish_ok)
    # without re-deriving everything _finish_ok needs from the
    # registry alone. Popped once finalized -- not a leak over a long
    # run with many delegations.
    self._worker_sessions: dict[str, Any] = {}
    # SendMessage (spec 12, peer_interaction feature, not yet default-on):
    # one threading.Condition PER DELEGATOR IDENTITY -- keyed by that
    # delegator's OWN delegation_id, or "entry" for the orchestrating
    # node itself -- not per Node. Two concurrent delegations of the
    # SAME role sharing this Node object (e.g. two "implementer"
    # delegations D001/D002, each possibly delegating further) are
    # DIFFERENT delegator identities and must never wake each other's
    # waits or see each other's children's messages. Lazily created
    # (one dict entry per identity that ever waits), guarded by its own
    # lock since Condition creation itself must be race-free.
    self._delegator_conds: dict[str, threading.Condition] = {}
    self._delegator_conds_lock = threading.Lock()
    # Push notifications: background threads append here; tool calls drain it.
    self._notifications: list[str] = []
    self._notifications_lock = threading.Lock()
    # Budget state — set at the start of each __call__ from AgenticState
    self._budget_seconds: float | None = None
    self._run_start: float | None = None
    self._stop: dict | None = None
    # Hard USD cost ceiling (None = inactive). Set each __call__ from state.
    self._budget_usd: float | None = None
    # True once any LLM call reports a real cost (claude). Stays False under
    # ollama (cost is None) → the USD ceiling is treated as inactive.
    self._cost_observed: bool = False
    self._usd_inactive_warned: bool = False
    # Consecutive Errored delegations per target (reset on that target's
    # next success). Drives the repeated-errors resumable halt.
    self._consecutive_errors: dict[str, int] = {}
    # Per-delegation pending messages (budget warnings) to prepend to
    # worker tool results.  Keyed by delegation_id; drained on next call.
    self._pending_worker_msgs: dict[str, list[str]] = {}
    self._pending_worker_msgs_lock = threading.Lock()
    # Tracks which budget % thresholds (80, 90, 100, 110 …) have already
    # been broadcast to workers so each is sent exactly once.
    self._budget_notified_pcts: set[int] = set()
    # Whether THIS node owns the run's epistemic ledgers. Decided once,
    # here, from what graph_builder passed: it hands the real notes_dir
    # to every orchestrating node (any node with outgoing edges), not
    # just the entry node — a delegating node is a node that needs
    # help from another node, nothing more, and telemetry / the science
    # monitor / hypothesis-ledger READ access matter for every role.
    # WRITE access (HypothesisPropose/Update, Milestone*) is gated
    # separately, by each Agent's own declared `tools`. Ownership is
    # never acquired later — see _install_epistemics, which then narrows it
    # to the records whose tools the node holds. Gated on `outgoing`
    # (not just `notes_dir is not None`) because graph_builder passes
    # the same notes_dir to every node's constructor, including a node
    # with no outgoing edges — that node must not acquire ledger
    # ownership merely because a path was handed to it.
    self._owns_epistemics: bool = bool(outgoing) and notes_dir is not None
    self._install_epistemics()
    # Running total of delegations at the START of the current __call__
    # Used as a seed for the delegation sequence counter.
    self._state_total_delegations: int = 0
    # Snapshot of _delegation_seq at the start of the current turn, so
    # the terminal Command counts only THIS turn's new delegations. Set
    # again by _absorb_state every turn; initialised here so the node's
    # attribute set is complete before any turn has run.
    self._seq_at_turn_start: int = 0
    # Monotonic per-node delegation counter — never reset within a
    # run.  Seeded from _state_total_delegations on first __call__
    # so checkpoint-resumed runs continue from the correct offset.
    # Because it never resets, it avoids the ID collision that
    # occurs when completed delegations are pruned from the registry
    # but _state_total_delegations has not yet accumulated them.
    self._delegation_seq: int = 0
    # Two-shot Done() gate: first call warns, second call closes.
    # Resets to False whenever a new Delegate() fires.
    self._done_warned: bool = False
    # Science monitor fires once per turn; reset at __call__ start.
    self._science_injected_this_turn: bool = False
    # Post-Done exit interview: set after the critic accepts; the next
    # Done() carries only the retrospective. _final_summary holds the real
    # conclusion so the recorded summary is the science, not the interview.
    self._awaiting_retro: bool = False
    self._final_summary: str | None = None
    # Consecutive non-PASS critic verdicts; after 3, the gate closes
    # gracefully UNGATED (bounded escape) instead of looping forever.
    self._revise_count: int = 0
    # Eval budget for this run (stashed each turn from state).
    self._eval_budget: int | None = None
    # Cumulative cap on the "no canonical source registered" nudge (soft).
    self._no_source_nudges: int = 0
    # Bounded re-prompt counter: incremented each time the node loops back
    # due to an unaccepted termination (no Done or refused Done).  NOT reset
    # in the A1/A2 per-turn block — it persists across loopbacks within one
    # run.  After 3 loopbacks the run terminates UNGATED.
    self._finish_attempts: int = 0
    # Accumulated token usage across strategizer + all workers this run.
    self._token_totals: dict = {
        "input_tokens": 0,
        "output_tokens": 0,
        "cache_read_input_tokens": 0,
        "cache_creation_input_tokens": 0,
        # The backend-independent schema (infra/telemetry.py): `legacy_calls`
        # counts calls whose backend reported none of it.
        "fresh_input": 0,
        "cache_read": 0,
        "cache_write": 0,
        "output": 0,
        "normalized_calls": 0,
        "legacy_calls": 0,
        "total_cost_usd": 0.0,
    }
    # Per-node raw tool-call error count: any ERROR: return or raised
    # exception from any injected closure counts as one error for that node.
    self._error_counts: dict[str, int] = {}
    self.adapter.closure_tools.update(self._build_routing_closures())
    self.adapter.route_watcher = lambda: self._route.get("kind") == "done"
_install_epistemics() -> None #

Build (or rebuild at a new path) everything notes_dir owns.

THE one construction site for the hypothesis ledger, the milestone ledger, the science monitor and telemetry. There used to be two: the constructor built the ledger, and _absorb_state — which runs at the start of every turn — rebuilt it with if self._ledger is None, knowing nothing about why the constructor had left it None. Three consequences, all of them silent:

  • graph_builder used to hand notes_dir to the entry node alone, so it alone owned the ledgers; any other orchestrating node re-acquired one on its first turn, undoing that decision. (It now hands the real notes_dir to every orchestrating node, so this particular inconsistency no longer applies — kept here as the historical reason _install_epistemics exists as one site.)
  • Only the ledger was re-pointed when the run's notes dir differed from the constructor's; the milestone ledger and telemetry kept writing to the stale path.
  • The science monitor, built once in the constructor, was never rebuilt at all — so a node whose ledger was resurrected ran with the ledger ON and the monitor OFF.

Ownership is fixed at construction (_owns_epistemics); this only ever rebuilds at a corrected path, never grants ownership.

Source code in src/adda/_src/nodes/orchestration.py
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
def _install_epistemics(self) -> None:
    """Build (or rebuild at a new path) everything ``notes_dir`` owns.

    THE one construction site for the hypothesis ledger, the milestone
    ledger, the science monitor and telemetry. There used to be two: the
    constructor built the ledger, and ``_absorb_state`` — which runs at the
    start of every turn — rebuilt it with ``if self._ledger is None``,
    knowing nothing about why the constructor had left it None. Three
    consequences, all of them silent:

    * ``graph_builder`` used to hand ``notes_dir`` to the entry node
      alone, so it alone owned the ledgers; any other orchestrating node
      re-acquired one on its first turn, undoing that decision. (It now
      hands the real ``notes_dir`` to every orchestrating node, so this
      particular inconsistency no longer applies — kept here as the
      historical reason ``_install_epistemics`` exists as one site.)
    * Only the ledger was re-pointed when the run's notes dir differed from
      the constructor's; the milestone ledger and telemetry kept writing to
      the stale path.
    * The science monitor, built once in the constructor, was never
      rebuilt at all — so a node whose ledger was resurrected ran with the
      ledger ON and the monitor OFF.

    Ownership is fixed at construction (``_owns_epistemics``); this only
    ever rebuilds at a corrected path, never grants ownership.
    """
    from ..epistemics.milestones import MilestoneLedger
    from ..infra.telemetry import Telemetry
    from ..runtime import features

    notes = self._current_notes_dir
    if not self._owns_epistemics or notes is None:
        self._ledger = None
        self._milestones = None
        self._science_monitor = None
        self._telemetry = None
        return

    # A node owns an epistemic record only if it holds the tools that
    # write to it. A ledger nobody can write is not a ledger, and a
    # milestone gate on a node that cannot set a milestone blocks a close
    # it has no means to unblock. The monitor watches what those tools
    # produce, so it follows either family.
    tools = self._agent_tools
    holds_ledger = not self._tools_declared or bool(
        tools & features.by_key("hypothesis_ledger").tools)
    holds_milestones = not self._tools_declared or bool(
        tools & features.by_key("milestones_enabled").tools)

    # Hypothesis ledger — persists hypotheses.json. Switchable: its tools
    # and its prompt section are withheld by the same knob (runtime.
    # features), so turning it off does not leave the agent commanded to
    # use tools that error.
    self._ledger: HypothesisLedger | None = (
        HypothesisLedger(notes)
        if holds_ledger and features.enabled("hypothesis_ledger") else None
    )

    # Milestone ledger (process policy) — persists milestones.json.
    # Seeded with the config default gates unless disabled. DISTINCT from
    # the hypothesis ledger (epistemics): process vs what's-true.
    self._milestones: MilestoneLedger | None = None
    if holds_milestones and features.enabled("milestones_enabled"):
        self._milestones = MilestoneLedger(notes)
        # C3 switchable: the draft-pipeline gate seeds only when the
        # pipeline-deliverable knob is on (off = byte-identical to today).
        # The oracle-gold-state milestone likewise seeds only when the
        # reproduction gate itself is on — its whole reason to exist is
        # that gate's store-row precondition.
        self._milestones.seed_defaults(
            include_pipeline=features.enabled("pipeline_deliverable"),
            include_reproduction_gate=features.enabled(
                "reproduction_gate"))

    # Science drift monitor — needs the delegation log and nothing else.
    # It used to be gated on the hypothesis ledger too, via a constructor
    # argument it stored and never read, so disabling the ledger disabled
    # the monitor as well.
    self._science_monitor: ScienceMonitor | None = None
    if (self._delegation_log is not None
            and (holds_ledger or holds_milestones)
            and features.enabled("science_monitor")):
        self._science_monitor = ScienceMonitor(
            self._delegation_log,
            diagnostics_writer=self._record_science_drift,
            role_of=self._role_of,
        )

    # Separable per-call telemetry — additive, off the decision path.
    # Lives under debug/telemetry/ (notes is debug/strategizer_notes).
    self._telemetry: Telemetry | None = Telemetry(notes.parent)
_log_status(delegation_id: str) -> tuple[str | None, str] #

Return (status, deliverable) for delegation_id from the persistent log, or (None, "") if the log has no such delegation.

Source code in src/adda/_src/nodes/orchestration.py
290
291
292
293
294
295
296
297
298
def _log_status(self, delegation_id: str) -> tuple[str | None, str]:
    """Return (status, deliverable) for *delegation_id* from the persistent
    log, or (None, "") if the log has no such delegation."""
    if self._delegation_log is None:
        return None, ""
    for r in self._delegation_log.query_all():
        if r.get("id") == delegation_id:
            return r.get("status"), (r.get("deliverable") or "")
    return None, ""
_pending_delegations() -> list[str] #

In-flight delegations, reconciled against the authoritative log.

A delegation the log shows terminal (DONE/FAILED) is never reported pending, even if the in-memory cache still says "Working". That stale state is what made Done()'s liveness gate refuse forever and kill run4 by watchdog after it had already found the optimum (audit BF-0/BF-2).

Source code in src/adda/_src/nodes/orchestration.py
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
def _pending_delegations(self) -> list[str]:
    """In-flight delegations, reconciled against the authoritative log.

    A delegation the log shows terminal (DONE/FAILED) is never reported
    pending, even if the in-memory cache still says "Working". That stale
    state is what made Done()'s liveness gate refuse forever and kill run4
    by watchdog after it had already found the optimum (audit BF-0/BF-2).
    """
    with self._registry_lock:
        pending = [
            d for d, e in self._registry.items()
            if e.get("status") == "Working"
        ]
    if self._delegation_log is not None:
        terminal = {
            r["id"] for r in self._delegation_log.query_all()
            if r.get("status") in ("DONE", "FAILED")
        }
        pending = [d for d in pending if d not in terminal]
    return pending
_find_datagenerator_name() -> str | None #

Name of the first connected datagenerator worker, or None.

Its presence means a canonical ground-truth source CAN be authored and registered for this study.

Source code in src/adda/_src/nodes/orchestration.py
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
def _find_datagenerator_name(self) -> str | None:
    """Name of the first connected datagenerator worker, or None.

    Its presence means a canonical ground-truth source CAN be authored
    and registered for this study.
    """
    spec = self._spec
    if spec is None or not hasattr(spec, "nodes"):
        return None
    for target in self._outgoing:
        agent = spec.nodes.get(target)
        if (
            agent is not None
            and getattr(agent, "role", None) == "datagenerator"
            and target in self._worker_adapters
        ):
            return target
    return None
_canonical_source_registered() -> bool #

True if a canonical ground-truth source is resolvable.

Reads run_config.json: an evaluator entrypoint OR a lookup pool counts as a registered source. Best-effort — on any read failure, assume not registered (the nudge is soft, so a false 'no' just costs one notice).

Source code in src/adda/_src/nodes/orchestration.py
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
def _canonical_source_registered(self) -> bool:
    """True if a canonical ground-truth source is resolvable.

    Reads run_config.json: an evaluator entrypoint OR a lookup pool counts
    as a registered source. Best-effort — on any read failure, assume not
    registered (the nudge is soft, so a false 'no' just costs one notice).
    """
    notes = self._current_notes_dir
    if notes is None:
        return False
    try:
        import json as _json
        cfg = _json.loads(
            (Path(notes).parent / "run_config.json").read_text())
        return bool(
            cfg.get("evaluator_entrypoint")
            or cfg.get("evaluator_lookup")
        )
    except Exception:  # noqa: BLE001
        return False
_drain_operator_notes() -> str #

Operator notes queued by a human in the viewer, marked as coming from the operator; "" when none. Notes addressed to a RUNNING delegation are routed to that worker instead. Claims the queue destructively, so whatever calls this owns delivery.

Source code in src/adda/_src/nodes/orchestration.py
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
def _drain_operator_notes(self) -> str:
    """Operator notes queued by a human in the viewer, marked as coming
    from the operator; "" when none. Notes addressed to a RUNNING
    delegation are routed to that worker instead. Claims the queue
    destructively, so whatever calls this owns delivery."""
    text = ""
    # Operator notes: messages a human queued in the viewer while the run
    # was working. Delivered here, on the same path as the budget warnings, so a note
    # reaches the agent at its next tool call rather than interrupting a
    # turn in progress. Marked as coming from the operator because the
    # agent should weigh it differently from another agent's message —
    # it is the one voice in the run that is not itself an agent.
    run_dir = self._current_run_dir
    if run_dir is not None:
        from ..infra.operator_channel import drain_note_rows
        rows = drain_note_rows(run_dir)
        mine: list[str] = []
        for _row in rows:
            _to = _row.get("to_node") or ""
            _note = (
                "[OPERATOR NOTE — from the human running this study. "
                "Weigh it as a briefing correction, not as another "
                f"agent's opinion: {_row['text']}]"
            )
            # A note addressed to a RUNNING delegation goes to that
            # worker, on the same per-delegation queue the
            # budget warnings use — so a human can correct work already
            # in flight instead of waiting for a wrong result. The queue
            # is claimed destructively, so this is the only place that
            # may drain it: routing here is what keeps an addressed note
            # from being swallowed by the orchestrator's own delivery.
            if _to:
                with self._registry_lock:
                    _entry = self._registry.get(_to)
                    _live = bool(_entry) and _entry.get("status") == "Working"
                if _live:
                    with self._pending_worker_msgs_lock:
                        self._pending_worker_msgs.setdefault(
                            _to, []).append(_note)
                    continue
                # Addressed to something not running: the human still
                # said it, so it must not vanish — hand it to the
                # orchestrator with the intended recipient named.
                _note = (
                    f"[OPERATOR NOTE addressed to {_to}, which is not "
                    f"running — delivered to you instead: "
                    f"{_row['text']}]"
                )
            mine.append(_note)
        if mine:
            text = "\n\n".join(mine) + "\n\n" + text
    return text
_drain_entry_notices() -> str #

Notices a campaign process queued for the entry node (the oracle was edited, say); "" for any other node or when none. The entry node holds no delegation id, so the backends' post-tool hook never reaches its queue; its next tool call does.

Source code in src/adda/_src/nodes/orchestration.py
426
427
428
429
430
431
432
433
434
435
436
437
def _drain_entry_notices(self) -> str:
    """Notices a campaign process queued for the entry node (the oracle
    was edited, say); "" for any other node or when none. The entry
    node holds no delegation id, so the backends' post-tool hook never
    reaches its queue; its next tool call does."""
    run_dir = self._current_run_dir
    entry = getattr(self._spec, "entry", None)
    if run_dir is None or entry is None or entry != self._name:
        return ""
    from ..infra.pending_notices import drain
    queued = drain(run_dir / "debug", "entry")
    return "\n".join(queued) + "\n\n" if queued else ""
_drain_notifications() -> str #

Return and clear any pending push notifications, or empty string.

Source code in src/adda/_src/nodes/orchestration.py
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
def _drain_notifications(self) -> str:
    """Return and clear any pending push notifications, or empty
    string."""
    with self._notifications_lock:
        if not self._notifications:
            text = ""
        else:
            msgs = list(self._notifications)
            self._notifications.clear()
            text = "\n".join(msgs) + "\n\n"
    text = self._drain_entry_notices() + text
    text = self._drain_operator_notes() + text
    text = self._stop_tick() + text
    if self._science_monitor is not None:
        offenders = self._science_monitor.escalation_due()
        _critic_name = self._find_critic_name()
        if offenders and _critic_name is not None:
            # Escalation fires: perform bookkeeping-only drain (discard
            # text) so the critic findings are the sole corrective
            # payload — regular drift messages would pollute context.
            self._science_monitor.drain()
            task_msg = self._build_feedback_task_msg(offenders)
            findings = self._invoke_critic(task_msg)
            self._science_monitor.note_escalated()
            if self._delegation_log is not None:
                _fb_id = (
                    "FB"
                    + datetime.now(
                        tz=timezone.utc
                    ).strftime("%H%M%S")
                )
                self._delegation_log.record(
                    id=_fb_id,
                    from_node=self._name,
                    to_node=_critic_name,
                    task="ScienceMonitor escalation audit",
                    deliverable=findings,
                    hypothesis_ids=offenders,
                    workspace_sha=self._commit_workspace(
                        f"{_fb_id} {self._name} -> {_critic_name} "
                        "[ESCALATION]"),
                    started_at=datetime.now(
                        tz=timezone.utc
                    ).isoformat(timespec="seconds"),
                    completed_at=datetime.now(
                        tz=timezone.utc
                    ).isoformat(timespec="seconds"),
                    status="FEEDBACK",
                    tokens_in=0,
                    tokens_out=0,
                    cost_usd=None,
                    critic_review=self._last_critic_review,
                )
            text += (
                "[SCIENCE MONITOR — ESCALATION] Repeated drift "
                f"on {', '.join(offenders)}. Critic audit "
                f"findings:\n{findings}\n"
            )
        else:
            # No escalation: inject at most once per strategizer turn
            # to avoid the same warning appearing on every tool call.
            if not self._science_injected_this_turn:
                drift = self._science_monitor.drain()
                if drift:
                    text += drift
                    self._science_injected_this_turn = True
    # Everything accumulated above is adda speaking to the agent, not
    # a tool's output — mark it so both the agent and the viewer can
    # tell the difference (see nodes/notices.py).
    return wrap_notice(text)
_delegation_entry(delegation_id: str) -> tuple[Any, dict] | None #

(owning node, registry entry) of a delegation, or None.

The entry lives in the registry of the node that DELEGATED (its delegator), not of the node running the work; the tools a worker calls are bound to the worker's own node, so a worker looking up its own delegation must search across the graph's nodes. Delegation ids come from one graph-wide log, so at most one node owns an id.

Source code in src/adda/_src/nodes/orchestration.py
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
def _delegation_entry(self, delegation_id: str) -> tuple[Any, dict] | None:
    """``(owning node, registry entry)`` of a delegation, or None.

    The entry lives in the registry of the node that DELEGATED (its
    delegator), not of the node running the work; the tools a worker
    calls are bound to the worker's own node, so a worker looking up its
    own delegation must search across the graph's nodes. Delegation ids
    come from one graph-wide log, so at most one node owns an id.
    """
    for n in (self, *getattr(self, "_peers", {}).values()):
        with n._registry_lock:
            entry = n._registry.get(delegation_id)
        if entry is not None:
            return n, entry
    return None
_get_delegator_cond(identity: str) -> threading.Condition #

The one Condition a given delegator identity waits on and is woken through (SendMessage, spec 12). identity is a delegation_id or "entry" — see __init__'s comment on _delegator_conds. Lazily created, race-free.

Source code in src/adda/_src/nodes/orchestration.py
526
527
528
529
530
531
532
533
534
535
536
def _get_delegator_cond(self, identity: str) -> threading.Condition:
    """The one Condition a given delegator identity waits on and is
    woken through (SendMessage, spec 12). ``identity`` is a
    delegation_id or ``"entry"`` — see ``__init__``'s comment on
    ``_delegator_conds``. Lazily created, race-free."""
    with self._delegator_conds_lock:
        cond = self._delegator_conds.get(identity)
        if cond is None:
            cond = threading.Condition()
            self._delegator_conds[identity] = cond
        return cond
_truncate_for_notice(text: str, n: int = 100) -> str staticmethod #
Source code in src/adda/_src/nodes/orchestration.py
538
539
540
541
@staticmethod
def _truncate_for_notice(text: str, n: int = 100) -> str:
    text = text.strip()
    return text if len(text) <= n else text[:n] + "..."
_pending_for_you(identity: str) -> str #

What identity (a delegator: "entry" or a delegation id of a node that is ITSELF delegating further) currently owes, computed FRESH on every call -- never a drained queue, since "nothing pending" must read as silence indefinitely, not just until the first drain (spec 12 design item 10(a): "at minimum, every agent must always KNOW what it currently owes ... or nothing (in which case: silence, not a notice for its own sake)").

Four kinds. As a DELEGATOR (entry.get("parent") == identity -- never a sibling's or a nested child's): a report open for review (OpenForReview), a finished delegation (Done/Errored) not yet collected via Wait, and an unread SendMessage question sitting in a child's own to_delegator queue. As a WORKER (identity is itself a registry entry, i.e. this node is mid-delegation): an unread SendMessage from ITS OWN delegator sitting in that entry's to_worker queue. Every queue peek (never a pop -- that stays Wait's/SendMessage's job) shares ONE lock acquisition with its own check, under the SAME Condition/lock the queue's real consumer uses (the check-then-pop race rule applies just as much to a check-then-PEEK). Returns "" when none apply -- the common case, deliberately silent.

Source code in src/adda/_src/nodes/orchestration.py
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
def _pending_for_you(self, identity: str) -> str:
    """What ``identity`` (a delegator: ``"entry"`` or a delegation id
    of a node that is ITSELF delegating further) currently owes,
    computed FRESH on every call -- never a drained queue, since
    "nothing pending" must read as silence indefinitely, not just
    until the first drain (spec 12 design item 10(a): "at minimum,
    every agent must always KNOW what it currently owes ... or
    nothing (in which case: silence, not a notice for its own
    sake)").

    Four kinds. As a DELEGATOR (``entry.get("parent") == identity``
    -- never a sibling's or a nested child's): a report open for
    review (`OpenForReview`), a finished delegation (`Done`/`Errored`)
    not yet collected via `Wait`, and an unread `SendMessage` question
    sitting in a child's own `to_delegator` queue. As a WORKER (``identity`` is itself a
    registry entry, i.e. this node is mid-delegation): an unread
    `SendMessage` from ITS OWN delegator sitting in that entry's
    `to_worker` queue. Every queue peek (never a pop -- that stays
    Wait's/SendMessage's job) shares ONE lock acquisition with its
    own check, under the SAME Condition/lock the queue's real
    consumer uses (the check-then-pop race rule applies just as much
    to a check-then-PEEK). Returns "" when none apply -- the common
    case, deliberately silent.
    """
    with self._registry_lock:
        mine = [
            (did, dict(e)) for did, e in self._registry.items()
            if e.get("parent") == identity
        ]
    _owned = self._delegation_entry(identity)
    my_own_entry = _owned[1] if _owned else None
    reviews = sorted(
        did for did, e in mine if e.get("status") == "OpenForReview")
    uncollected = sorted(
        did for did, e in mine
        if e.get("status") in ("Done", "Errored") and not e.get("waited")
    )

    cond = self._get_delegator_cond(identity)
    with cond:
        questions = [
            (did, e.get("target", did), e["to_delegator"][0][1])
            for did, e in mine if e.get("to_delegator")
        ]

    worker_message = None
    if my_own_entry is not None:
        with my_own_entry["worker_cond"]:
            if my_own_entry.get("to_worker"):
                worker_message = my_own_entry["to_worker"][0][1]

    if not (
        reviews or uncollected or questions
        or worker_message
    ):
        return ""
    bits = []
    if reviews:
        bits.append(
            f"open for review: {', '.join(reviews)} -- SendMessage(id, "
            "..., approve=True) to finalize, or ask a question first"
        )
    if uncollected:
        bits.append(
            f"finished, not yet collected: {', '.join(uncollected)} "
            "-- Wait(id)"
        )
    if questions:
        bits.append("; ".join(
            f"{did} ({target}) asked you: "
            f"{self._truncate_for_notice(msg)}"
            for did, target, msg in questions
        ))
    if worker_message:
        bits.append(
            "your delegator sent you a message: "
            f"{self._truncate_for_notice(worker_message)}"
        )
    return "You have pending items -- " + "; ".join(bits) + "."
_build_routing_closures() -> dict #
Source code in src/adda/_src/nodes/orchestration.py
623
624
625
def _build_routing_closures(self) -> dict:
    from .tools.routing import build_routing_tools
    return build_routing_tools(self)
_wrap_closure(fn: Any, node_name: str) -> Any #

Return a version of fn that records ERROR returns and exceptions.

Uses functools.wraps so inspect.signature() follows wrapped to the original function — _infer_schema_from_callable must see the real parameter names, not (args, *kwargs).

Also coerces string-typed arguments to int/float/bool when the function annotation requests it (handles Ollama passing "5" for an int parameter).

And it is where pending notices reach the agent: whatever _drain_notifications holds is APPENDED to the tool's result — after it, so a result's leading word stays what callers dispatch on ("Done", "Working", "ERROR:", an id like "H3"). That used to be hand-copied into 18 tool bodies and missing from the rest — the store tools never delivered a notice — so whether the agent heard about a finished delegation depended on which tool it happened to call. A tool that places notices itself (the status poll keeps its status word first; Wait drains while it blocks; Done composes its own reply) still does, and finds the queue already empty here. A worker's closures pass through here too, and are never drained into.

Source code in src/adda/_src/nodes/orchestration.py
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
def _wrap_closure(self, fn: Any, node_name: str) -> Any:
    """Return a version of *fn* that records ERROR returns and exceptions.

    Uses functools.wraps so inspect.signature() follows __wrapped__ to the
    original function — _infer_schema_from_callable must see the real
    parameter names, not (*args, **kwargs).

    Also coerces string-typed arguments to int/float/bool when the
    function annotation requests it (handles Ollama passing "5" for
    an int parameter).

    And it is where pending notices reach the agent: whatever
    ``_drain_notifications`` holds is APPENDED to the tool's result — after
    it, so a result's leading word stays what callers dispatch on ("Done",
    "Working", "ERROR:", an id like "H3"). That
    used to be hand-copied into 18 tool bodies and missing from the rest —
    the store tools never delivered a notice — so whether the agent heard
    about a finished delegation depended on which tool it happened to
    call. A tool that places notices itself (the status poll keeps its
    status word first; Wait drains while it blocks; Done composes its own
    reply) still does, and finds the queue already empty here. A
    worker's closures pass through here too, and are never drained into.
    """
    import functools as _functools
    import inspect as _inspect
    import typing as _typing

    node = self
    tool_name = getattr(fn, "__name__", repr(fn))

    # Resolve type hints once; fall back to {} if any forward ref
    # cannot be resolved (e.g. "DelegationLog | None").
    try:
        _hints = _typing.get_type_hints(fn)
    except Exception:  # noqa: BLE001
        _hints = {}
    _COERCIBLE = {int, float, bool}

    def _coerce(name: str, value: Any) -> Any:
        target = _hints.get(name)
        if target not in _COERCIBLE or not isinstance(value, str):
            return value
        if target is bool:
            low = value.strip().lower()
            if low in ("true", "1", "yes"):
                return True
            if low in ("false", "0", "no"):
                return False
            return value
        try:
            return target(value)
        except ValueError:
            return value

    @_functools.wraps(fn)
    def _wrapped(*args, **kwargs):
        node._raise_if_abandoned()
        # Coerce string args before calling the real function.
        call_args = dict(kwargs)
        try:
            bound = _inspect.signature(fn).bind_partial(
                *args, **kwargs
            )
            for pname in list(bound.arguments):
                bound.arguments[pname] = _coerce(
                    pname, bound.arguments[pname]
                )
            args, kwargs = bound.args, bound.kwargs
            call_args = dict(bound.arguments)
        except TypeError:
            pass  # signature mismatch: let fn raise its own error

        try:
            result = fn(*args, **kwargs)
            if (
                isinstance(result, str)
                and result.lstrip().startswith("ERROR:")
            ):
                node._record_tool_error(
                    node_name,
                    tool_name,
                    "ERROR_RETURN",
                    result[:2000],
                    args=call_args,
                )
            # A tool can carry a diagnostic source (e.g. ConsultLiterature
            # tags itself with its LiteratureCorpus) for facts that are
            # true of the environment, not of one call's return value —
            # e.g. the dense embedder being unavailable — and so would
            # otherwise never reach diagnostics.jsonl or the agent. This
            # is the one place with both the per-run diagnostics path
            # and, via the tag, a handle back to the source; it fires
            # once (the source pops its own event) rather than on every
            # call.
            _diag_src = getattr(fn, "_adda_diagnostic_source", None)
            if _diag_src is not None:
                _event = _diag_src.pop_diagnostic_event()
                if _event is not None:
                    _kind, _msg, *_extra = _event
                    node._record_intervention(
                        _kind, node_name, _msg, **(_extra[0] if _extra else {}))
                    if isinstance(result, str):
                        result = wrap_notice(_msg) + result
            # Only for this node's OWN tools. A dispatched worker's
            # closures are wrapped by the delegating node too (it records
            # their errors), and draining there would hand the
            # orchestrator's notices to the worker.
            if isinstance(result, str) and node_name == node._name:
                notices = node._drain_notifications()
                if notices.strip():
                    result = insert_notice(result, notices)
                # Spec 12 design item 10(a): a pending-for-you notice
                # on every tool result, scoped to THIS CALL's own
                # delegator identity (never a sibling's or a nested
                # child's). Gated behind peer_interaction (on by
                # default since the migration-sweep commit) -- with the
                # feature off (the no-peer-messaging arm),
                # OpenForReview cannot exist and this stays silent.
                from ..runtime import features as _features
                if _features.enabled("peer_interaction"):
                    from ..backends.base import get_delegation_id
                    identity = get_delegation_id() or "entry"
                    pending = node._pending_for_you(identity)
                    if pending:
                        result = insert_notice(
                            result, wrap_notice(pending))
            return result
        except Exception as exc:
            node._record_tool_error(
                node_name,
                tool_name,
                type(exc).__name__,
                str(exc)[:2000],
                tb=traceback.format_exc(),
                args=call_args,
            )
            raise

    return _wrapped
_orchestrate(state: AgenticState) -> Any #

One orchestration turn, in the order it happens.

Absorb the run state onto the node, work out what the agent must be TOLD this turn, invoke it once, then route on what it did. Every step is a method below, named for its step; a turn ends in exactly one of three ways, which :meth:_route_turn states.

Source code in src/adda/_src/nodes/orchestration.py
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
def _orchestrate(self, state: AgenticState) -> Any:
    """One orchestration turn, in the order it happens.

    Absorb the run state onto the node, work out what the agent must be
    TOLD this turn, invoke it once, then route on what it did. Every step
    is a method below, named for its step; a turn ends in exactly one of
    three ways, which :meth:`_route_turn` states.
    """
    self._absorb_state(state)
    budget_warnings = self._budget_warnings(state)
    halt = self._check_unrecoverable(
        state, self._budget_seconds, self._run_start)
    if halt is not None:
        return halt
    pending_notifs = self._reset_for_turn()
    stop_notice = self._stop_tick()
    if stop_notice:
        pending_notifs.append(stop_notice)
    messages = self._compose_messages(state, budget_warnings, pending_notifs)
    try:
        ai_msg = self._invoke_turn(messages)
    except Exception as exc:  # noqa: BLE001
        if self._stop is None:
            raise
        return self._close_after_wind_down_error(state, exc)
    return self._route_turn(state, ai_msg)
_absorb_state(state: AgenticState) -> None #

Copy the run state the node's TOOLS read onto the node itself.

The closures reach their context through node.…, not through AgenticState, so anything a tool needs has to land here first.

Source code in src/adda/_src/nodes/orchestration.py
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
def _absorb_state(self, state: AgenticState) -> None:
    """Copy the run state the node's TOOLS read onto the node itself.

    The closures reach their context through ``node.…``, not through
    AgenticState, so anything a tool needs has to land here first.
    """
    # Update notes_dir from current state run_dir
    run_dir = state.get("run_dir")
    if run_dir:
        _notes = Path(run_dir) / "debug" / "strategizer_notes"
        if _notes != self._current_notes_dir:
            # The run's real notes dir differs from the one the
            # constructor saw. Re-point EVERYTHING that lives there, not
            # just the hypothesis ledger — the milestone ledger and
            # telemetry used to keep writing to the stale path. Nodes that
            # do not own the ledgers still track the path (the critic gate
            # reads run files through it) but acquire nothing.
            self._current_notes_dir = _notes
            self._install_epistemics()
        # Wire canonical store dir into ScienceMonitor lazily.
        # store_dir is the ExperimentData *project_dir* (run_dir/
        # experiment_data), NOT the folder holding the CSVs. ExperimentData
        # appends its own EXPERIMENTDATA_SUBFOLDER ("experiment_data"), so
        # the rows live one level deeper at
        # run_dir/experiment_data/experiment_data/output.csv — hence the
        # apparent double directory is correct, not a typo.
        if self._science_monitor is not None:
            self._science_monitor.store_dir = (
                self._current_notes_dir.parent.parent
                / "experiment_data"
            )

    # Eval budget → available to the Done() critic gate for budget-aware
    # framing (judge the best honest conclusion within evals spent).
    self._eval_budget = state.get("eval_budget")

    # Required aux deliverables (config.yaml) → on the node so WriteDeliverable
    # may write them: the gate REQUIRES them, so the writing tool must accept
    # them (else gate-vs-tool deadlock — audit run 20260624T021359).
    self._required_deliverables = state.get("required_deliverables") or []

    # Store on node so the status poll can compute delegation timeout
    self._budget_seconds = state.get("budget_seconds")
    self._run_start = state.get("start_time")
    self._budget_usd = state.get("budget_usd")

    # Capture total_delegations so Delegate() can seed the counter.
    self._state_total_delegations = state.get("total_delegations", 0)
    # Seed the monotonic counter from state on first turn (or after a
    # checkpoint rebuild).  Never decremented — ensures IDs are unique
    # even when the registry is pruned between turns.
    if self._delegation_seq < self._state_total_delegations:
        self._delegation_seq = self._state_total_delegations
    # Snapshot seq at turn start so total_new counts only THIS turn.
    self._seq_at_turn_start: int = self._delegation_seq
_ledgered_eval_total(floor: int) -> int #

The run's eval count, preferring the ledger over an accumulator.

The canonical ledger is the source of truth: a killed/cancelled delegation flushes rows the state accumulator never sees, so the accumulator undercounts. Summed across the canonical store AND every design namespace — namespace evals were invisible to the run total and to the soft budget. floor keeps the accumulator for lookup-direct studies with no instrumented store.

Source code in src/adda/_src/nodes/orchestration.py
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
def _ledgered_eval_total(self, floor: int) -> int:
    """The run's eval count, preferring the ledger over an accumulator.

    The canonical ledger is the source of truth: a killed/cancelled
    delegation flushes rows the state accumulator never sees, so the
    accumulator undercounts. Summed across the canonical store AND every
    design namespace — namespace evals were invisible to the run total and
    to the soft budget. ``floor`` keeps the accumulator for lookup-direct
    studies with no instrumented store.
    """
    try:
        from ..evaluation.ledger_summary import total_ledgered_evals
        _nd = getattr(self, "_current_notes_dir", None)
        if _nd is not None:
            return max(
                floor,
                int(total_ledgered_evals(
                    _nd.parent.parent / "experiment_data")),
            )
    except Exception:  # noqa: BLE001
        pass
    return floor
_budget_warnings(state: AgenticState) -> list[dict] #

Advisory budget messages for this turn.

The time budget is a SOFT constraint — warnings only; the run is never force-terminated for exceeding it. A separate run-level backstop (RUN_BACKSTOP_MULTIPLE x budget) bounds runaway cost.

Source code in src/adda/_src/nodes/orchestration.py
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
def _budget_warnings(self, state: AgenticState) -> list[dict]:
    """Advisory budget messages for this turn.

    The time budget is a SOFT constraint — warnings only; the run is never
    force-terminated for exceeding it. A separate run-level backstop
    (RUN_BACKSTOP_MULTIPLE x budget) bounds runaway cost.
    """
    import time

    from ..runtime import features
    warnings: list[dict] = []
    if not features.enabled("budget_notes"):
        return warnings
    budget, start = self._budget_seconds, self._run_start
    if budget is not None and start is not None:
        elapsed = time.time() - start
        pct = elapsed / budget
        if pct >= 1.0:
            # Escalating ladder, once per newly-crossed 10% band (100,
            # 110, 120, …) — not every turn, which a model learns to
            # skip. Shared with leaf.py/delegation.py so a worker gets
            # the same signal (nodes/_constants.py).
            if budget_band_due(elapsed, budget, self._budget_bands_fired):
                warnings.append({
                    "role": "user",
                    "content": budget_wrapup_message(
                        elapsed, budget, can_call_done=True),
                })
        elif pct >= 0.95:
            warnings.append({
                "role": "user",
                "content": (
                    f"Warning: time budget at {pct*100:.0f}% "
                    f"({elapsed:.0f}s / {budget:.0f}s). "
                    "Begin wrapping up — call Done() soon. Don't cancel a "
                    "progressing delegation under time pressure; its evals "
                    "are already ledgered and cancelling only loses its "
                    "report (Wait(id, block=False) shows whether it's "
                    "progressing)."
                ),
            })

    eval_budget = state.get("eval_budget")
    evals_used = self._ledgered_eval_total(state.get("evals_used", 0))
    if eval_budget is not None and evals_used >= eval_budget:
        warnings.append({
            "role": "user",
            "content": (
                f"Warning: eval budget exceeded"
                f" ({evals_used} used / {eval_budget} budget)."
                f" Do not run further evaluations."
            ),
        })
    return warnings
_reset_for_turn() -> list[str] #

A1/A2: clear per-turn state; return the notifications to deliver.

Working entries are preserved so loopbacks don't orphan live delegations whose background threads are still running.

Source code in src/adda/_src/nodes/orchestration.py
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
def _reset_for_turn(self) -> list[str]:
    """A1/A2: clear per-turn state; return the notifications to deliver.

    Working entries are preserved so loopbacks don't orphan live
    delegations whose background threads are still running.
    """
    self._route.clear()
    self._ask_count = 0
    self._done_warned = False
    self._science_injected_this_turn = False
    with self._registry_lock:
        self._registry = {
            d: e for d, e in self._registry.items()
            if e["status"] in ("Working", "Done")
        }
        self._threads = {
            d: t for d, t in self._threads.items()
            if d in self._registry
        }
    with self._notifications_lock:
        pending = list(self._notifications)
        self._notifications.clear()
    return pending
_compose_messages(state: AgenticState, budget_warnings: list[dict], pending_notifs: list[str]) -> list[dict] #

The conversation plus everything injected in-band this turn.

Everything after the conversation itself is adda speaking to the agent — a budget warning, the no-source nudge, the milestone backlog, a pushed notification — arriving in the SAME role ("user") the human's own task arrives in. Marked, for the same reason tool-result notices are (nodes/notices.py): unmarked, neither the agent nor a reader can tell the runtime's nudge from the human's brief, and the viewer cannot style it as anything else.

Source code in src/adda/_src/nodes/orchestration.py
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
def _compose_messages(
    self,
    state: AgenticState,
    budget_warnings: list[dict],
    pending_notifs: list[str],
) -> list[dict]:
    """The conversation plus everything injected in-band this turn.

    Everything after the conversation itself is adda speaking to the
    agent — a budget warning, the no-source nudge, the milestone backlog,
    a pushed notification — arriving in the SAME role ("user") the human's
    own task arrives in. Marked, for the same reason tool-result notices
    are (nodes/notices.py): unmarked, neither the agent nor a reader can
    tell the runtime's nudge from the human's brief, and the viewer cannot
    style it as anything else.
    """
    injected = (
        self._constraint_refresh()
        + budget_warnings
        + self._no_source_nudge()
        + self._backlog_announcement()
        + [{"role": "user", "content": n} for n in pending_notifs]
    )
    return _to_adapter_messages(state["messages"]) + [
        {**m, "content": wrap_notice(str(m.get("content", "")),
                                     trailing="")}
        for m in injected
        if str(m.get("content", "")).strip()
    ]
_constraint_refresh() -> list[dict] #

This turn's constraint snapshot, recomputed NOW.

The entry node used to be the one call site that did not re-snapshot. agent_runtime rendered a snapshot once at run start and concatenated it onto the problem statement, and that string is the standing first user turn — re-sent verbatim on every later turn, so the numbers inside it could never advance. Observed on run 20260917T141603: four consecutive strategizer turns spanning 23.5 minutes all read "2.7min/60.0min used (5%)", while the true figure at turn 4 was 26.2min (44%).

The orchestrator was not blind -- every delegation report carries a fresh snapshot (_append_budget_report) -- which made this a CONTRADICTION rather than an absence: a frozen block and live blocks in one context. Worse, the frozen copy lived in the first user turn, which the context trim pins and never evicts, so it was the one guaranteed to survive while the fresh ones aged out.

Recomputing per turn is what constraint_snapshot.py already asks of every caller: "call this at every delegation boundary rather than caching a value from earlier ... its entire purpose depends on being current". The orchestrator is the node that decides how much more to attempt, so it is the node that most needs the live clock.

Source code in src/adda/_src/nodes/orchestration.py
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
def _constraint_refresh(self) -> list[dict]:
    """This turn's constraint snapshot, recomputed NOW.

    The entry node used to be the one call site that did not re-snapshot.
    ``agent_runtime`` rendered a snapshot once at run start and
    concatenated it onto the problem statement, and that string is the
    standing first user turn — re-sent verbatim on every later turn, so
    the numbers inside it could never advance. Observed on run
    20260917T141603: four consecutive strategizer turns spanning 23.5
    minutes all read "2.7min/60.0min used (5%)", while the true figure at
    turn 4 was 26.2min (44%).

    The orchestrator was not blind -- every delegation report carries a
    fresh snapshot (``_append_budget_report``) -- which made this a
    CONTRADICTION rather than an absence: a frozen block and live blocks
    in one context. Worse, the frozen copy lived in the first user turn,
    which the context trim pins and never evicts, so it was the one
    guaranteed to survive while the fresh ones aged out.

    Recomputing per turn is what ``constraint_snapshot.py`` already asks
    of every caller: "call this at every delegation boundary rather than
    caching a value from earlier ... its entire purpose depends on being
    current". The orchestrator is the node that decides how much more to
    attempt, so it is the node that most needs the live clock.
    """
    from ..runtime import features
    from ..runtime.constraint_snapshot import snapshot_for_node
    if not features.enabled("budget_notes"):
        return []
    try:
        text = snapshot_for_node(self).as_text()
    except Exception:  # noqa: BLE001 — a missing snapshot never fails a turn
        return []
    return [{"role": "user", "content": text}] if text.strip() else []
_no_source_nudge() -> list[dict] #

Recommend registering a canonical source (soft, ≤3×).

If this graph has a datagenerator (so a canonical ground-truth source CAN be authored) but none is registered, recommend delegating to it. Without a registered source every evaluation lands off-ledger and nothing is reproducible from the canonical store. Soft and capped — never blocks; the strategizer may ignore it for a genuinely source-free study.

Source code in src/adda/_src/nodes/orchestration.py
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
def _no_source_nudge(self) -> list[dict]:
    """Recommend registering a canonical source (soft, ≤3×).

    If this graph has a datagenerator (so a canonical ground-truth source
    CAN be authored) but none is registered, recommend delegating to it.
    Without a registered source every evaluation lands off-ledger and
    nothing is reproducible from the canonical store. Soft and capped —
    never blocks; the strategizer may ignore it for a genuinely
    source-free study.
    """
    if (
        self._no_source_nudges >= 3
        or self._find_datagenerator_name() is None
        or self._canonical_source_registered()
    ):
        return []
    self._no_source_nudges += 1
    _dg = self._find_datagenerator_name()
    self._record_intervention(
        "NO_SOURCE_NUDGE", self._name,
        "No canonical source registered; recommended delegating to "
        f"'{_dg}'.",
        notice=self._no_source_nudges, cap=3,
    )
    return [{
        "role": "user",
        "content": (
            "[SETUP] No canonical ground-truth source is registered "
            "for this study (no evaluator entrypoint or lookup pool). "
            f"A '{_dg}' agent is available — delegate to it to author "
            "and register the source, so evaluations flow through "
            "get_evaluator(), land in the canonical store, and the "
            "result is reproducible. If this is intentionally a "
            "source-free (surrogate-only) study, disregard this. "
            f"(notice {self._no_source_nudges}/3)"
        ),
    }]
_backlog_announcement() -> list[dict] #

Announce the process backlog ONCE, at the start of the run.

As a conversation message, so the agent cannot claim it didn't know these gate the implementer. Injected the first time this node runs.

Source code in src/adda/_src/nodes/orchestration.py
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
def _backlog_announcement(self) -> list[dict]:
    """Announce the process backlog ONCE, at the start of the run.

    As a conversation message, so the agent cannot claim it didn't know
    these gate the implementer. Injected the first time this node runs.
    """
    if self._milestones is None or getattr(self, "_backlog_announced", False):
        return []
    from ..epistemics.milestones import render_backlog
    _bl = render_backlog(self._milestones)
    self._backlog_announced = True
    return [{"role": "user", "content": _bl}] if _bl else []
_invoke_turn(messages: list[dict]) -> Any #

Run one model turn and account for what it spent.

Source code in src/adda/_src/nodes/orchestration.py
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
def _invoke_turn(self, messages: list[dict]) -> Any:
    """Run one model turn and account for what it spent."""
    from langchain_core.messages import AIMessage

    # DEBUG: stream this turn's full reasoning + tool-calls to
    # debug/transcripts/strategizer/turn_NNN.jsonl.
    from ..backends.base import (
        bind_run_context as _bind_rc,
    )
    from ..backends.base import (
        debug_enabled as _dbg,
    )
    from ..backends.base import (
        set_transcript_sink as _set_sink,
    )
    self._turn_count = getattr(self, "_turn_count", 0) + 1
    if _dbg() and self._current_notes_dir is not None:
        _set_sink(str(
            self._current_notes_dir.parent / "transcripts"
            / "strategizer" / f"turn_{self._turn_count:03d}.jsonl"))
    _notes = self._current_notes_dir
    _rc_path = (
        str(_notes.parent / "run_config.json") if _notes is not None else None
    )
    with _bind_rc(f"{self._name}-turn-{self._turn_count:03d}", _rc_path):
        text = self.adapter.invoke(messages)
    # Accumulate this node's own token usage.
    self._record_usage(
        getattr(self.adapter, "last_usage", {}) or {},
        role=self._role_of(self._name),
        model=getattr(self.adapter, "model", None),
        phase="strategizer_turn",
        delegation_id=None,
    )
    return AIMessage(content=text)
_route_turn(state: AgenticState, ai_msg: Any) -> Any #

A turn ends in exactly one of three ways.

  1. Work is still in flight → re-prompt, free (no attempt spent)
  2. It stopped without closing → re-prompt, bounded (3 attempts)
  3. Otherwise → the run ends

Reproduction is owned entirely by the Done() gate (it runs the controlled gate before any close and declares a FAILED run after a bounded number of sighted attempts — see RunNotebook(gate=True)), so there is no separate post-accept repro check here; this handles only deliverable presence and un-accepted termination.

Source code in src/adda/_src/nodes/orchestration.py
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
def _route_turn(self, state: AgenticState, ai_msg: Any) -> Any:
    """A turn ends in exactly one of three ways.

    1. Work is still in flight     → re-prompt, free (no attempt spent)
    2. It stopped without closing  → re-prompt, bounded (3 attempts)
    3. Otherwise                   → the run ends

    Reproduction is owned entirely by the Done() gate (it runs the
    controlled gate before any close and declares a FAILED run after a
    bounded number of sighted attempts — see RunNotebook(gate=True)), so there is
    no separate post-accept repro check here; this handles only deliverable
    presence and un-accepted termination.
    """
    accepted = self._route.get("kind") == "done"
    missing = self._missing_deliverables(state)
    for router in (self._reprompt_while_working,
                   self._reprompt_unfinished):
        held = router(ai_msg, accepted, missing)
        if held is not None:
            return held
    return self._terminate_run(state, ai_msg, accepted, missing)
_reprompt_while_working(ai_msg: Any, accepted: bool, missing: list) -> Any | None #

Delegations still running: re-prompt WITHOUT spending an attempt.

A healthy delegation still in flight is WORK IN PROGRESS, not a failed finish: the deliverables usually depend on its result, and it WILL report. Spending a bounded finish-attempt on it means a slow-but-healthy delegation (run-4: D004 at ~2.5 evals/s, ~100s from done, with wall budget to spare) burns 3 "finish attempts" across turns and force- terminates the run UNGATED. The run's time backstop (run_backstop_multiple x budget, checked each turn) bounds a delegation that truly hangs.

Source code in src/adda/_src/nodes/orchestration.py
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
def _reprompt_while_working(
    self, ai_msg: Any, accepted: bool, missing: list
) -> Any | None:
    """Delegations still running: re-prompt WITHOUT spending an attempt.

    A healthy delegation still in flight is WORK IN PROGRESS, not a failed
    finish: the deliverables usually depend on its result, and it WILL
    report. Spending a bounded finish-attempt on it means a slow-but-healthy
    delegation (run-4: D004 at ~2.5 evals/s, ~100s from done, with wall
    budget to spare) burns 3 "finish attempts" across turns and force-
    terminates the run UNGATED. The run's time backstop
    (run_backstop_multiple x budget, checked each turn) bounds a delegation
    that truly hangs.
    """
    from langchain_core.messages import HumanMessage
    from langgraph.types import Command

    if accepted:
        return None
    working = self._working_delegations()
    if not working:
        return None
    msg = (
        f"Delegations still running: {working}. They are"
        " progressing — collect them with Wait and call Done() only"
        " once they report (then write any remaining deliverables"
        " from their results). Do NOT close early. This wait does"
        " NOT count against your finish attempts; the run's time"
        " budget is the backstop."
    )
    if missing:
        msg += "\n\nStill to write AFTER they finish: " + ", ".join(missing)
    return Command(
        goto=self._name,
        update={"messages": [ai_msg, HumanMessage(content=msg)]},
    )
_reprompt_unfinished(ai_msg: Any, accepted: bool, missing: list) -> Any | None #

Bounded re-prompt when the turn ended without an accepted close.

Source code in src/adda/_src/nodes/orchestration.py
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
def _reprompt_unfinished(
    self, ai_msg: Any, accepted: bool, missing: list
) -> Any | None:
    """Bounded re-prompt when the turn ended without an accepted close."""
    from langchain_core.messages import HumanMessage
    from langgraph.types import Command

    from ..runtime import features
    if (accepted and not missing) or self._finish_attempts >= 3:
        return None
    if not features.enabled("reprompt_unfinished"):
        return None
    self._finish_attempts += 1
    problems: list[str] = []
    if missing:
        missing_list = "\n".join(f"- {p}" for p in missing)
        problems.append(
            "Required deliverables are missing from the"
            f" study directory:\n{missing_list}\n"
            "Write them via WriteDeliverable() before"
            " calling Done()."
        )
    if not accepted:
        working = self._working_delegations()
        if working:
            problems.append(
                f"Delegations still running: {working}."
                " Collect them with Wait and call Done()"
                " once they finish."
            )
        else:
            problems.append(
                "You ended your turn without an accepted"
                " Done(). If Done() was refused (critic"
                " verdict, two-shot confirmation, or another"
                " gate), address the refusal and call Done()"
                " again. A run only closes through an"
                " accepted Done()."
            )
    return Command(
        goto=self._name,
        update={
            "messages": [
                ai_msg,
                HumanMessage(content=(
                    "Run cannot complete"
                    f" (attempt {self._finish_attempts}/3):\n"
                    + "\n\n".join(problems)
                )),
            ],
        },
    )
_working_delegations() -> list[str] #

Delegation ids still in flight right now.

Source code in src/adda/_src/nodes/orchestration.py
1222
1223
1224
1225
1226
1227
1228
def _working_delegations(self) -> list[str]:
    """Delegation ids still in flight right now."""
    with self._registry_lock:
        return [
            d for d, e in self._registry.items()
            if e["status"] == "Working"
        ]
_terminate_run(state: AgenticState, ai_msg: Any, accepted: bool, missing: list) -> Any #

Close the run: final counts, banner, ghost flush, terminal Command.

Source code in src/adda/_src/nodes/orchestration.py
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
def _terminate_run(
    self, state: AgenticState, ai_msg: Any, accepted: bool, missing: list
) -> Any:
    """Close the run: final counts, banner, ghost flush, terminal Command."""
    from langgraph.graph import END
    from langgraph.types import Command

    from ..runtime import terminal

    # total_new: only delegations created THIS turn (seq delta vs the
    # snapshot taken at turn start), not Done entries from prior turns.
    with self._registry_lock:
        total_new = self._delegation_seq - self._seq_at_turn_start
        evals_new = sum(e["evals"] for e in self._registry.values())

    summary = self._banner(
        self._route.get("summary") or ai_msg.content, accepted, missing)
    self._flush_ghost_delegations()

    # The persisted (reported) eval total prefers the ledger aggregate over
    # the accumulator: the accumulator can drop evals a namespace-blind
    # guard mis-flagged as off-ledger, and never saw namespace stores at
    # all. The ledger across all namespaces is authoritative → run_status.
    _evals_persist = self._ledgered_eval_total(
        state.get("evals_used", 0) + evals_new)
    # The terminal triple, decided here rather than inferred from the
    # banner later. An un-accepted close never reached a gate, and a close
    # missing deliverables was never validated against them — both are
    # UNGATED regardless of what the route recorded. terminal.resolve
    # fails safe (unrecorded → UNGATED) and refuses GATED without a review.
    _outcome, _termination, _reviewed = terminal.resolve(
        self._route.get("outcome")
        if accepted and not missing else terminal.UNGATED,
        self._route.get("termination") if accepted else terminal.NO_CLOSE,
        self._route.get("reviewed"),
    )
    return Command(
        goto=END,
        update={
            "messages": [ai_msg],
            "done": True,
            "last_report": summary,
            "total_delegations": state["total_delegations"] + total_new,
            "evals_used": _evals_persist,
            "token_totals": dict(self._token_totals),
            "error_counts": dict(self._error_counts),
            "outcome": _outcome,
            "termination": _termination,
            "reviewed": _reviewed,
        },
    )
_banner(summary: str, accepted: bool, missing: list) -> str #

Prepend the UNGATED banner when the run ends without an accepted Done().

A FAILED-reproduction close carries its own ⛔ banner in the route summary and IS accepted=done, so it is not re-banner'd here.

Source code in src/adda/_src/nodes/orchestration.py
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
def _banner(self, summary: str, accepted: bool, missing: list) -> str:
    """Prepend the UNGATED banner when the run ends without an accepted Done().

    A FAILED-reproduction close carries its own ⛔ banner in the route
    summary and IS accepted=done, so it is not re-banner'd here.
    """
    from ..runtime import features, terminal

    if (accepted and not missing) or not features.enabled(
            "reprompt_unfinished"):
        return summary
    flags = []
    if not accepted:
        flags.append(
            "the run terminated WITHOUT an accepted Done() —"
            " the final conclusions did NOT pass the"
            " adversarial critic gate"
        )
    if missing:
        flags.append(f"required deliverables missing: {missing}")
    return terminal.ungated_banner(flags) + summary
_flush_ghost_delegations() -> None #

Close out delegations whose threads die with the interpreter.

Daemon threads still alive when the run closes are killed at process exit — their run() never reaches the DONE/FAILED record write, leaving orphan RUNNING entries in the log. Write an INTERRUPTED terminal record for each so query_all() (last-wins) collapses to a closed state instead of RUNNING.

Source code in src/adda/_src/nodes/orchestration.py
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
def _flush_ghost_delegations(self) -> None:
    """Close out delegations whose threads die with the interpreter.

    Daemon threads still alive when the run closes are killed at process
    exit — their run() never reaches the DONE/FAILED record write, leaving
    orphan RUNNING entries in the log. Write an INTERRUPTED terminal record
    for each so query_all() (last-wins) collapses to a closed state instead
    of RUNNING.
    """
    with self._registry_lock:
        live = [
            (did, dict(entry))
            for did, entry in self._registry.items()
            if entry.get("status") == "Working"
        ]
    if not live or self._delegation_log is None:
        return
    _now = datetime.now(tz=timezone.utc).isoformat(timespec="seconds")
    for _did, _entry in live:
        self._delegation_log.record(
            id=_did,
            from_node=self._name,
            to_node=_entry.get("target", "unknown"),
            # An interrupted delegation never reached _finish_ok/_error,
            # so its partial writes are still uncommitted. Commit them
            # HERE, against the delegation that made them, or the next
            # delegation to commit absorbs them and the history says the
            # wrong worker wrote those files.
            workspace_sha=self._commit_workspace(
                f"{_did} {self._name} -> "
                f"{_entry.get('target', 'unknown')} [INTERRUPTED]"),
            task="",
            deliverable=(
                "INTERRUPTED: run closed while this delegation was "
                "still running (background thread killed at process exit)"
            ),
            hypothesis_ids=_entry.get("hypothesis_ids") or [],
            started_at=_entry.get("started_at") or "",
            completed_at=_now,
            status="INTERRUPTED",
            tokens_in=0,
            tokens_out=0,
            cost_usd=None,
            is_falsification_attempt=bool(
                _entry.get("is_falsification_attempt")
            ),
            evals=_entry.get("evals", 0),
            phase=_entry.get("phase"),
        )
_missing_deliverables(state: AgenticState) -> list[str] #

Return required deliverable paths not present at study_dir yet.

The single deliverable (pipeline.ipynb) is required, authored before Done() is accepted, UNLESS the study turns off pipeline_deliverable — it is the human-readable recipe AND the reproduction in one notebook: the runtime executes it lazily (see _reproduction_gate) to verify the headline re-derives from the ledger with zero new evals, which makes no sense for a study with no ledger at all. Requiring it unconditionally (BACKLOG #30) meant a pure-derivation study's strategizer could never satisfy Done() regardless of what its own tools/prompt said — this was the third of three places that assumption was baked in (the other two, the injected notebook_deliverable_spec() preamble and the notebook- authoring tools themselves, are already gated the same way). Additional paths can still be declared in state['required_deliverables'] regardless of this flag.

Source code in src/adda/_src/nodes/reproduction_gate.py
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
def _missing_deliverables(self, state: AgenticState) -> list[str]:
    """Return required deliverable paths not present at study_dir yet.

    The single deliverable (pipeline.ipynb) is required, authored before
    Done() is accepted, UNLESS the study turns off pipeline_deliverable —
    it is the human-readable recipe AND the reproduction in one notebook:
    the runtime executes it lazily (see _reproduction_gate) to verify the
    headline re-derives from the ledger with zero new evals, which makes
    no sense for a study with no ledger at all. Requiring it unconditionally
    (BACKLOG #30) meant a pure-derivation study's strategizer could never
    satisfy Done() regardless of what its own tools/prompt said — this was
    the third of three places that assumption was baked in (the other two,
    the injected notebook_deliverable_spec() preamble and the notebook-
    authoring tools themselves, are already gated the same way). Additional
    paths can still be declared in state['required_deliverables']
    regardless of this flag.
    """
    from ..evaluation.notebook_exec import required_deliverable_name
    from ..runtime import features
    from ..runtime.features import NOTEBOOK_TOOLS
    study_dir = Path(state.get("study_dir", "."))
    # WriteDeliverable writes BARE names to study_dir/ (it rejects path
    # separators). Normalise any configured path to its basename so a stray
    # 'workspace/…' prefix in a study config can't spuriously flag a present
    # deliverable as missing.
    required = list(state.get("required_deliverables") or [])
    # Only a node that holds the notebook tools can author the notebook,
    # so only it can be required to have.
    if features.enabled("pipeline_deliverable") and (
            not self._tools_declared or self._agent_tools & NOTEBOOK_TOOLS):
        required = [required_deliverable_name()] + required
    seen: set[str] = set()
    missing: list[str] = []
    for p in required:
        name = Path(p).name
        if name in seen:
            continue
        seen.add(name)
        if not (study_dir / name).exists():
            missing.append(name)
    return missing
_reproduction_gate(state: AgenticState | None = None) -> str | None #

Before Done() can close a run — and on every RunNotebook(gate=True) dry run — pipeline.ipynb must satisfy every one of these, checked in order. This is a MECHANICAL check this exact function executes every time, not a judgement call, and nothing narrative or intent-based can satisfy it in a check's place.

  1. The canonical store must already hold at least one oracle row. Zero rows means no campaign has been evaluated yet, so there is nothing for the notebook to reproduce FROM — the gate refuses outright, before even running the notebook.
  2. The notebook must finish cleanly within a time ceiling (no heavy from-scratch computation — a reproduction is lazy, not a re-run).
  3. It must add ZERO new oracle rows: it may only LOAD the ledger (ExperimentData.from_file) and reach the oracle through get_evaluator(), which skips every already-FINISHED row.
  4. It must NOT modify or delete any existing ledger row — no faking a zero-delta by delete-then-re-add or by rewriting a value.
  5. The notebook prints a freshly-computed REPRODUCED: <value> and the CLAIMED_HEADLINE: <value> its write-up states. The gate fails the notebook when the two differ. It does not fail a notebook for a missing print.

Passing (0)-(4) means the deliverable is a faithful, lightweight, lazy reproduction of a real, already-evaluated campaign — not a script doing something unrelated to validating the pipeline.

Source code in src/adda/_src/nodes/reproduction_gate.py
 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
def _reproduction_gate(self, state: AgenticState | None = None) -> str | None:
    """Before Done() can close a run — and on every RunNotebook(gate=True)
    dry run — pipeline.ipynb must satisfy every one of these, checked in
    order. This is a MECHANICAL check this exact function executes every
    time, not a judgement call, and nothing narrative or intent-based can
    satisfy it in a check's place.

      0. The canonical store must already hold at least one oracle row.
         Zero rows means no campaign has been evaluated yet, so there is
         nothing for the notebook to reproduce FROM — the gate refuses
         outright, before even running the notebook.
      1. The notebook must finish cleanly within a time ceiling (no heavy
         from-scratch computation — a reproduction is lazy, not a re-run).
      2. It must add ZERO new oracle rows: it may only LOAD the ledger
         (``ExperimentData.from_file``) and reach the oracle through
         ``get_evaluator()``, which skips every already-FINISHED row.
      3. It must NOT modify or delete any existing ledger row — no faking
         a zero-delta by delete-then-re-add or by rewriting a value.
      4. The notebook prints a freshly-computed ``REPRODUCED: <value>``
         and the ``CLAIMED_HEADLINE: <value>`` its write-up states. The
         gate fails the notebook when the two differ. It does not fail
         a notebook for a missing print.

    Passing (0)-(4) means the deliverable is a faithful, lightweight,
    lazy reproduction of a real, already-evaluated campaign — not a
    script doing something unrelated to validating the pipeline.
    """
    # Developer notes (not part of the agent-facing contract above, see
    # gate_contract() below): returns None on PASS (and stashes
    # self._repro_ok_detail), else a problem string describing which
    # check failed. Skips silently when there is no run context; callable
    # without `state` (study dir comes from self._study_dir). (d)'s
    # REPRODUCED marker is informational for the critic/human — the
    # runtime does not gate on the value itself, only on (4)'s
    # self-consistency; an independent runtime extremum match used to
    # wrongly reject legitimate constrained optima.
    from ..runtime import features
    if not features.enabled("reproduction_gate"):
        return None

    import re
    import subprocess

    study_dir = (
        Path(self._study_dir) if getattr(self, "_study_dir", None) is not None
        else Path((state or {}).get("study_dir", "."))
    )
    # The deliverable is pipeline.ipynb. (A .py is still executable by the
    # executor-agnostic gate, kept only as a fallback for gate-logic tests;
    # the notebook is preferred when both are present.) Absence is left to
    # _missing_deliverables.
    deliverable = next(
        (study_dir / n for n in ("pipeline.ipynb", "pipeline.py")
         if (study_dir / n).exists()),
        None,
    )
    if deliverable is None:
        return None  # absence is handled by _missing_deliverables
    if deliverable.suffix == ".ipynb":
        # NOTEBOOK-LEDGER SYNC, by construction: refresh the hypotheses
        # cell's ledger-status block before anything else runs. This one
        # hook covers BOTH call sites that reach here — RunNotebook
        # (gate=True) and Done()'s pre-critic check — so a stale status
        # is impossible the moment either reads the notebook, rather
        # than merely detected once it's already been read.
        from .tools.routing.notebook import refresh_hypotheses_ledger_block
        # Reassigned every call (None on a no-op refresh) so a stale
        # rev from an EARLIER call is never re-surfaced by _gate_check.
        self._hypotheses_rev_after_refresh = refresh_hypotheses_ledger_block(
            study_dir, self._read_ledger())
    # One resolver for "where is this run", shared with the store tools:
    # they used to compute it separately and could disagree.
    if self._current_notes_dir is None:
        return None  # no run dir context (e.g. non-debug) — skip the gate
    run_dir = self._resolve_run_dir()          # …/runs/<id>
    store_dir = run_dir / "experiment_data"
    run_config = run_dir / "debug" / "run_config.json"

    from ..evaluation.notebook_exec import (
        ledger_snapshot as _ledger_snapshot,
    )

    # ── HERMETIC SANDBOX ──────────────────────────────────────────────────
    # CRITICAL: run the deliverable against a COPY of the canonical store, never
    # the live one. A faithful lazy pipeline adds nothing; a NON-lazy one
    # (re-evaluating) writes its evals into the THROWAWAY copy — we detect
    # that as "not lazy" while the real ledger stays pristine. Without this,
    # checking a non-lazy pipeline pollutes + inflates the canonical store
    # (and the gate check could be looped to balloon it without bound).
    before_n, before_hash = _ledger_snapshot(store_dir)
    if before_n == 0:
        return (
            "Canonical store has no rows — the campaign has not been "
            "evaluated yet. Run the delegation pipeline first so the "
            "ledger is populated, then the notebook can be reproduced "
            "lazily against those rows.")
    from ..evaluation.notebook_exec import replay_sandbox
    with replay_sandbox(
            store_dir, run_config, self._study_dir,
    ) as (sandbox, sb_store, env):
        _timeout = (
            max(0.1 * self._budget_seconds, 180.0)
            if self._budget_seconds else 300.0
        )
        try:
            # Executor-agnostic: a .ipynb runs via nbclient (in-env kernel),
            # a .py via subprocess — both return a CompletedProcess and raise
            # TimeoutExpired on timeout, so the asserts below are unchanged.
            from ..evaluation.notebook_exec import run_deliverable
            proc = run_deliverable(
                deliverable, cwd=sandbox, env=env, timeout=_timeout)
        except subprocess.TimeoutExpired:
            return (
                f"{deliverable.name} did not finish within {_timeout:.0f}s. A "
                "reproduction must be lightweight — load the ledger and skip "
                "finished evals and heavy refits (cache-or-load surrogates). "
                "Make it lazy.")
        after_n, after_hash = _ledger_snapshot(sb_store)

    # (a) clean exit — surface a generous stderr tail for sighted debugging.
    if proc.returncode != 0:
        return (
            f"{deliverable.name} FAILED to run (exit {proc.returncode}). It must "
            "load the ledger and derive the headline cleanly. Stderr:\n"
            + (proc.stderr or "")[-3000:]
            + ("\n\nStdout tail:\n" + proc.stdout[-800:]
               if proc.stdout else ""))
    # (b) zero new evals (lazy).
    if after_n != before_n:
        return (
            f"{deliverable.name} is NOT lazy: re-running it changed the ledger row "
            f"count ({before_n} → {after_n}). It must LOAD the ledger "
            "(ExperimentData.from_file) and reach the oracle only via "
            "get_evaluator() so FINISHED rows are skipped — zero new evals.")
    # (c) integrity — existing rows unchanged.
    if before_hash and after_hash and before_hash != after_hash:
        return (
            f"{deliverable.name} MODIFIED existing ledger rows. A reproduction must "
            "read the ledger READ-ONLY (it may re-store identical rows, but "
            "must not rewrite values or delete+re-add). Do not tamper with "
            "the canonical store.")
    # (d) The printed ``REPRODUCED:`` line is an informational headline
    # marker for the critic / human reader — the runtime no longer gates on
    # it. Headline GROUNDING (the value traces to a real ledger row) is
    # owned by the critic's HEADLINE PROVENANCE check; an independent
    # runtime extremum match wrongly rejected legitimate CONSTRAINED optima
    # (a constrained best is, by definition, not an objective extremum), so
    # it forced studies to headline their infeasible unconstrained extremum
    # — see audit run 20260624T021359.
    # (e) internal consistency — if the notebook declares CLAIMED_HEADLINE
    # (the value its write-up states), it must equal the freshly-computed
    # REPRODUCED. Lenient: skips when the marker is absent.
    _hc = _headline_consistency(proc.stdout or "")
    if _hc is not None:
        return _hc
    m = re.search(r"REPRODUCED:\s*([-+]?[0-9]*\.?[0-9]+(?:[eE][-+]?[0-9]+)?)",
                  proc.stdout or "")
    headline = f", REPRODUCED={m.group(1)}" if m else ""
    self._repro_ok_detail = (
        f"reproduced cleanly ({before_n} rows, unchanged, 0 new evals, "
        f"ran in <{_timeout:.0f}s{headline})")
    return None
_halt_resumable(state: Any, *, reason: str, termination: str, status: str = 'halted', extra_update: dict | None = None) #

Checkpoint-and-halt cleanly on an unrecoverable condition.

Instead of crashing, write debug/run_status.json (so tooling and the human can see the run is resumable) and return a Command to END whose last_report is prefixed with a HALTED banner. The durable SqliteSaver checkpoint + the persisted thread_id are what make the run resumable via AgenticRun(resume_from=...) — no new serialized state is introduced here.

termination names WHICH unrecoverable condition fired. It is carried on the state so the close path records it verbatim instead of inferring an outcome from the HALTED banner — which matched none of the old banner patterns and so read as GATED, logging every backstop kill as a validated success.

Source code in src/adda/_src/nodes/lifecycle.py
 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
def _halt_resumable(
    self,
    state: Any,
    *,
    reason: str,
    termination: str,
    status: str = "halted",
    extra_update: dict | None = None,
):
    """Checkpoint-and-halt cleanly on an unrecoverable condition.

    Instead of crashing, write ``debug/run_status.json`` (so tooling and
    the human can see the run is resumable) and return a ``Command`` to
    ``END`` whose ``last_report`` is prefixed with a HALTED banner.  The
    durable SqliteSaver checkpoint + the persisted ``thread_id`` are what
    make the run resumable via ``AgenticRun(resume_from=...)`` — no new
    serialized state is introduced here.

    ``termination`` names WHICH unrecoverable condition fired. It is
    carried on the state so the close path records it verbatim instead of
    inferring an outcome from the HALTED banner — which matched none of the
    old banner patterns and so read as GATED, logging every backstop kill
    as a validated success.
    """
    import json as _json

    from langchain_core.messages import AIMessage
    from langgraph.graph import END
    from langgraph.types import Command

    from ..runtime import features as _features

    run_dir = state.get("run_dir")
    thread_id = None
    if run_dir:
        debug_dir = Path(run_dir) / "debug"
        tid_path = debug_dir / "thread_id"
        try:
            if tid_path.exists():
                thread_id = tid_path.read_text().strip()
        except OSError:
            thread_id = None
        try:
            debug_dir.mkdir(parents=True, exist_ok=True)
            (debug_dir / "run_status.json").write_text(
                _json.dumps(
                    {
                        "status": status,
                        "reason": reason,
                        "resumable": True,
                        "thread_id": thread_id,
                        "outcome": terminal.UNGATED,
                        "termination": termination,
                        "reviewed": False,
                        "arms": _features.arm_config(),
                    },
                    indent=2,
                ),
                encoding="utf-8",
            )
        except OSError:
            pass

    # Preserve the latest conclusion below the banner.
    prior_text = ""
    for _m in reversed(state["messages"]):
        if isinstance(_m, AIMessage):
            prior_text = str(_m.content)
            break

    banner = f"## ⚠ HALTED (resumable) — {reason}\n\n"
    update = {
        "messages": [],
        "done": True,
        "last_report": banner + (prior_text or "(no prior report)"),
        "token_totals": dict(self._token_totals),
        "error_counts": dict(self._error_counts),
        "outcome": terminal.UNGATED,
        "termination": termination,
        "reviewed": False,
    }
    if extra_update:
        update.update(extra_update)
    return Command(goto=END, update=update)
_halt_tallies(state: Any) -> dict #
Source code in src/adda/_src/nodes/lifecycle.py
102
103
104
105
106
107
108
109
def _halt_tallies(self, state: Any) -> dict:
    with self._registry_lock:
        _total_new = len(self._registry)
        _evals_new = sum(e["evals"] for e in self._registry.values())
    return {
        "total_delegations": state["total_delegations"] + _total_new,
        "evals_used": state.get("evals_used", 0) + _evals_new,
    }
_close_after_wind_down_error(state: Any, exc: Exception) #

The turn that was to close a wind-down raised. The run closes anyway, with the halt's own termination, and says which retrospectives it lacks.

Source code in src/adda/_src/nodes/lifecycle.py
111
112
113
114
115
116
117
118
119
def _close_after_wind_down_error(self, state: Any, exc: Exception):
    """The turn that was to close a wind-down raised. The run closes anyway,
    with the halt's own termination, and says which retrospectives it lacks."""
    reason = (f"the wind-down turn failed ({type(exc).__name__}: "
              f"{str(exc)[:200]})")
    self._log_missing_retrospectives(reason)
    return self._halt_resumable(
        state, reason=reason, termination=self._stop_termination(),
        extra_update=self._halt_tallies(state))
_trip(state: Any, *, drift: dict, reason: str, termination: str, tallies: dict) -> Command | None #

A backstop condition holds. Ask the run to wind down — every agent gives its retrospective, the run closes with termination — instead of jumping to END with none. While the wind-down is in progress this returns None (the condition keeps holding; that is expected). Only a wind-down that overruns its grace, or one that cannot be asked for, falls through to the hard halt, so a halt is always bounded.

Source code in src/adda/_src/nodes/lifecycle.py
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
def _trip(
    self, state: Any, *, drift: dict, reason: str, termination: str,
    tallies: dict,
) -> Command | None:
    """A backstop condition holds. Ask the run to wind down — every agent
    gives its retrospective, the run closes with ``termination`` — instead of
    jumping to END with none. While the wind-down is in progress this
    returns None (the condition keeps holding; that is expected). Only a
    wind-down that overruns its grace, or one that cannot be asked for,
    falls through to the hard halt, so a halt is always bounded."""
    if self._stop is None:
        self._record_science_drift(drift)
        if self._request_wind_down(
            reason=reason, termination=termination,
            run_dir=state.get("run_dir")):
            self._record_intervention(
                "BACKSTOP_WIND_DOWN", "(run)",
                f"{reason}; winding down for retrospectives, closing "
                f"{termination}")
            return None
    elif not self._wind_down_overdue():
        return None
    else:
        self._record_science_drift({**drift, "wind_down": "overran"})
        self._log_missing_retrospectives(
            f"{reason}; the wind-down overran its grace")
    return self._halt_resumable(
        state, reason=reason, termination=termination,
        extra_update=tallies)
_check_unrecoverable(state: Any, budget: float | None, start: float | None) -> Command | None #

Return a halt Command if an unrecoverable condition is met, else None. Extracted verbatim from call (USD ceiling → repeated errors → time backstop).

Source code in src/adda/_src/nodes/lifecycle.py
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
def _check_unrecoverable(self, state: Any, budget: float | None, start: float | None) -> Command | None:
    """Return a halt Command if an unrecoverable condition is met, else None.
    Extracted verbatim from __call__ (USD ceiling → repeated errors → time
    backstop)."""

    # (1) USD cost ceiling. Hard, resumable (raise budget_usd and resume).
    # Inactive under ollama (no cost data): warn once, never halt.
    _budget_usd = self._budget_usd
    if _budget_usd is not None and _budget_usd > 0:
        _spent = self._token_totals.get("total_cost_usd") or 0.0
        if not self._cost_observed:
            if (
                getattr(self, "_turn_count", 0) >= 1
                and not self._usd_inactive_warned
            ):
                self._usd_inactive_warned = True
                self._record_science_drift({
                    "error_type": "USD_BUDGET_INACTIVE",
                    "budget_usd": _budget_usd,
                    "note": "no per-call cost reported (e.g. ollama); "
                            "USD ceiling treated as inactive",
                })
        elif _spent >= _budget_usd:
            return self._trip(
                state,
                drift={
                    "error_type": "USD_BACKSTOP",
                    "spent_usd": _spent,
                    "budget_usd": _budget_usd,
                },
                reason=(
                    f"USD budget exhausted "
                    f"(${_spent:.4f} / ${_budget_usd:.4f})"
                ),
                termination=terminal.BACKSTOP_USD,
                tallies=self._halt_tallies(state),
            )

    # (2) Repeated errors: a target failing N times in a row (genuine
    # worker EXCEPTIONS, not REVISE loops or poor results — those reset the
    # streak on any success) is not going to self-heal by looping the
    # strategizer at it again. The default is deliberately CONSERVATIVE: a
    # legitimate run "stuck" in the scientific process loops on critic
    # verdicts and slow delegations, none of which count here — only hard
    # consecutive crashes do. Knob: max_consecutive_errors (config.yaml
    # runtime block; F3DASM_MAX_CONSECUTIVE_ERRORS overrides); 0 disables.
    from ..runtime.settings import get_int
    _max_err = get_int("max_consecutive_errors", 12)
    if _max_err > 0:
        with self._registry_lock:
            _stuck = [
                (t, n) for t, n in self._consecutive_errors.items()
                if n >= _max_err
            ]
        if _stuck:
            _t, _n = _stuck[0]
            return self._trip(
                state,
                drift={
                    "error_type": "REPEATED_ERRORS",
                    "target": _t,
                    "consecutive": _n,
                },
                reason=(
                    f"repeated errors: {_t} failed {_n}x consecutively"
                ),
                termination=terminal.REPEATED_ERRORS,
                tallies=self._halt_tallies(state),
            )

    # (3) Time backstop: past run_backstop_multiple x the (soft) time
    # budget, bound runaway cost. Now resumable (raise budget + resume).
    _backstop_mult = run_backstop_multiple()
    if backstop_enabled() and budget is not None and start is not None:
        _elapsed_now = time.time() - start
        if _elapsed_now > budget * _backstop_mult:
            with self._registry_lock:
                _abandoned = [
                    d for d, e in self._registry.items()
                    if e["status"] == "Working"
                ]
            return self._trip(
                state,
                drift={
                    "error_type": "RUN_BACKSTOP",
                    "elapsed": _elapsed_now,
                    "budget": budget,
                    "multiple": _backstop_mult,
                    "abandoned": _abandoned,
                },
                reason=(
                    f"time backstop: {int(_backstop_mult)}x budget "
                    f"exceeded ({_elapsed_now:.0f}s / {budget:.0f}s)"
                ),
                termination=terminal.BACKSTOP_TIME,
                tallies=self._halt_tallies(state),
            )

    return None
_find_critic_name() -> str | None #

Name of the first connected critic worker, or None.

Source code in src/adda/_src/nodes/critic_gate.py
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
def _find_critic_name(self) -> str | None:
    """Name of the first connected critic worker, or None."""
    spec = self._spec
    if spec is None or not hasattr(spec, "nodes"):
        return None
    for target in self._outgoing:
        agent = spec.nodes.get(target)
        if (
            agent is not None
            and getattr(agent, "role", None) == "critic"
            and target in self._worker_adapters
        ):
            return target
    return None
_prior_reviews_digest() -> str #

A bounded <prior_reviews_this_run> block of THIS run's earlier critic reviews (verdict + findings only), or "" if none.

The critic is invoked one-shot per gate with no live session and no RecallHistory, so without this it cannot see what it already ruled and can silently contradict an earlier verdict (the H1 SUPPORTED->FALSIFIED ->back whipsaw that drove the REVISE-spin). Echoing its standing objections back forces consistency: it may still reverse, but only by saying so and citing new evidence — never silently.

Source code in src/adda/_src/nodes/critic_gate.py
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
def _prior_reviews_digest(self) -> str:
    """A bounded `<prior_reviews_this_run>` block of THIS run's earlier
    critic reviews (verdict + findings only), or "" if none.

    The critic is invoked one-shot per gate with no live session and no
    RecallHistory, so without this it cannot see what it already ruled and
    can silently contradict an earlier verdict (the H1 SUPPORTED->FALSIFIED
    ->back whipsaw that drove the REVISE-spin). Echoing its standing
    objections back forces consistency: it may still reverse, but only by
    saying so and citing new evidence — never silently.
    """
    notes = self._current_notes_dir
    if notes is None:
        return ""
    review_dir = Path(notes).parent / "critic_reviews"
    if not review_dir.is_dir():
        return ""
    files = sorted(review_dir.glob("call_*.md"))
    if not files:
        return ""
    elided = max(0, len(files) - _MAX_PRIOR_REVIEWS)
    kept = files[-_MAX_PRIOR_REVIEWS:]
    blocks: list[str] = []
    for f in kept:
        try:
            text = f.read_text(encoding="utf-8")
        except Exception:  # noqa: BLE001
            continue
        n = f.stem.replace("call_", "")
        verdict = (
            _extract_report_section(text, "Verdict")
            or _parse_verdict(text)
        )
        findings = (
            _extract_report_section(text, "Findings")
            or text.strip()[:_PRIOR_REVIEW_BODY_CAP]
            or "(none)"
        )
        blocks.append(
            f"--- call_{n} ---\n"
            f"Verdict: {verdict}\n"
            f"Findings:\n{findings}"
        )
    # Char budget: drop oldest kept blocks until under budget.
    while blocks and sum(len(b) for b in blocks) > _PRIOR_REVIEWS_CHAR_BUDGET:
        blocks.pop(0)
        elided += 1
    if not blocks:
        return ""
    elided_note = (
        f"\n({elided} earlier review(s) elided for brevity.)"
        if elided else ""
    )
    return (
        "<prior_reviews_this_run>\n"
        "You have already reviewed this run. Your standing verdicts and "
        "objections are below. Be CONSISTENT with them: do not silently "
        "contradict a verdict you reached earlier, and do not re-raise an "
        "objection the strategizer has since resolved. You MAY reverse a "
        "prior position, but only by stating which call you are reversing "
        "and citing the Charter clause and the NEW evidence that justifies "
        "it.\n"
        + "\n\n".join(blocks)
        + elided_note
        + "\n</prior_reviews_this_run>\n\n"
    )
_invoke_critic(task_msg: str) -> str #

Synchronously invoke the connected critic; returns its text or an ERROR string.

The critic is a worker too: under runtime: debug its full transcript is streamed to disk, its verdict/review is ALWAYS persisted (the PASS branch doesn't echo it to the strategizer, so this is the only place the deciding verdict is auditable), and its ### Retrospective is recorded like every other node's (#7).

Source code in src/adda/_src/nodes/critic_gate.py
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
def _invoke_critic(self, task_msg: str) -> str:
    """Synchronously invoke the connected critic; returns its
    text or an ERROR string.

    The critic is a worker too: under `runtime: debug` its full transcript is
    streamed to disk, its verdict/review is ALWAYS persisted (the PASS
    branch doesn't echo it to the strategizer, so this is the only place
    the deciding verdict is auditable), and its ### Retrospective is
    recorded like every other node's (#7).
    """
    critic_name = self._find_critic_name()
    if critic_name is None:
        return "ERROR: no critic connected."
    # Cross-round memory: prepend this run's earlier reviews so the critic
    # stays consistent instead of whipsawing its own verdicts (Fix A).
    task_msg = self._prior_reviews_digest() + task_msg
    adapter = self._worker_adapters[critic_name]
    worker = (
        adapter.copy() if hasattr(adapter, "copy") else adapter
    )
    self._critic_calls = getattr(self, "_critic_calls", 0) + 1
    _n = self._critic_calls
    from ..backends.base import (
        bind_run_context as _bind_rc,
    )
    from ..backends.base import (
        debug_enabled as _dbg,
    )
    from ..backends.base import (
        get_transcript_sink as _get_sink,
    )
    from ..backends.base import (
        set_transcript_sink as _set_sink,
    )
    _notes = self._current_notes_dir
    _prev_sink = _get_sink()
    _rc_path = (
        str(_notes.parent / "run_config.json") if _notes is not None else None
    )
    # Cleared up front so a failed invoke cannot leave the PREVIOUS
    # call's usage to be logged against this delegation.
    self._last_critic_usage = {}
    self._last_critic_review = None
    _ok = False
    if _dbg() and _notes is not None:
        _set_sink(str(
            _notes.parent / "transcripts" / "critic"
            / f"call_{_n:03d}.jsonl"))
    try:
        with _bind_rc(f"critic-{_n}", _rc_path):
            critique = worker.invoke(
                [{"role": "user", "content": task_msg}]
            )
        _ok = True
    except Exception as _exc:  # noqa: BLE001
        # An infrastructure failure invoking the critic — NOT a problem with
        # your deliverable. Give a one-line cause, not a raw traceback, and a
        # constructive next step.
        critique = (
            "ERROR: the critic could not be invoked due to an "
            f"infrastructure error ({type(_exc).__name__}: "
            f"{str(_exc)[:200]}). This is not a defect in your deliverable. "
            "Re-call Done() to retry the gate; if it persists, the run will "
            "close without a critic PASS (the failure is on record)."
        )
    finally:
        # Restore the strategizer's own sink (same thread-local).
        if _dbg() and _notes is not None:
            _set_sink(_prev_sink)
    # Account the critic's tokens/cost — critic consults are real LLM calls
    # and must land in token_totals AND telemetry (they were previously
    # uncounted, undercounting run cost and omitting the 'critic' role).
    #
    # Only when the call actually succeeded. ``last_usage`` lives on the
    # adapter and survives a failed invoke, so reading it unconditionally
    # bills THIS call for the previous one's tokens — a failed gate would
    # be logged at the cost of the gate before it.
    _usage = (getattr(worker, "last_usage", {}) or {}) if _ok else {}
    # Also published for the CALLER to put on its delegation row. Both
    # call sites (the Done() GATE check and AskForFeedback) log a
    # strategizer -> critic delegation, and both used to hardcode
    # tokens 0 / cost None on it — so the row representing the single
    # interaction that decides whether a run closes carried no
    # accounting at all, and anything summing cost_usd over
    # delegation_log.jsonl silently omitted the critic's whole budget.
    self._last_critic_usage = _usage
    self._record_usage(
        _usage,
        role=self._role_of(critic_name),
        model=getattr(worker, "model", None),
        phase="critic_review",
        delegation_id=f"critic-{_n}",
    )
    # Always-on: persist the verdict/review to disk + record retrospective.
    self._last_critic_review = self._persist_critic_review(_n, critique)
    self._record_retrospective("critic", f"critic-{_n}", critique)
    return critique
_invoke_verdict_validator(prompt: str) -> str #

One-shot LLM judgement reusing the CRITIC's adapter (same judge as the gate), but lighter than _invoke_critic — it does NOT persist a critic review or a retrospective. Records token usage so the call lands in the run cost. Best-effort: returns "" if no critic is connected or the call fails (the validator is advisory; a missing judge must not break the run).

Source code in src/adda/_src/nodes/critic_gate.py
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
def _invoke_verdict_validator(self, prompt: str) -> str:
    """One-shot LLM judgement reusing the CRITIC's adapter (same judge as the
    gate), but lighter than ``_invoke_critic`` — it does NOT persist a critic
    review or a retrospective. Records token usage so the call lands in the
    run cost. Best-effort: returns "" if no critic is connected or the call
    fails (the validator is advisory; a missing judge must not break the run).
    """
    critic_name = self._find_critic_name()
    if critic_name is None:
        return ""
    adapter = self._worker_adapters[critic_name]
    worker = adapter.copy() if hasattr(adapter, "copy") else adapter
    from ..backends.base import (
        bind_run_context as _bind_rc,
    )
    _notes = self._current_notes_dir
    _rc_path = (
        str(_notes.parent / "run_config.json") if _notes is not None else None
    )
    # Tight budget: this advisory judge must NOT inherit a real agent turn's
    # 5×600s stream/retry budget. A hung CLI stream once froze a whole run
    # for ~89 min here (run 20260627T211310). idle=120 + retry_max=1 abort
    # ~2 min after the stream goes silent; the call is advisory, so on any
    # failure the verdict stands (see _run_verdict_validator).
    try:
        with _bind_rc("verdict-validator", _rc_path):
            try:
                reply = worker.invoke(
                    [{"role": "user", "content": prompt}],
                    idle_timeout=120.0, retry_max=1,
                )
            except TypeError:
                # A backend/stub without the budget kwargs — fall back gracefully.
                reply = worker.invoke([{"role": "user", "content": prompt}])
    except Exception:  # noqa: BLE001
        return ""
    self._record_usage(
        getattr(worker, "last_usage", {}) or {},
        role=self._role_of(critic_name),
        model=getattr(worker, "model", None),
        phase="verdict_validation",
        delegation_id="verdict-validator",
    )
    return reply or ""
_cited_delegation_brief(evidence: dict | None) -> str | None #

A short brief on the cited delegation — what it was asked to do and whether it was flagged a falsification attempt — so the judge can weigh §2 attempt-adequacy. The numeric result the verdict rests on is carried separately in evidence['numbers']. None if no delegation is cited or the record is absent.

Source code in src/adda/_src/nodes/critic_gate.py
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
def _cited_delegation_brief(self, evidence: dict | None) -> str | None:
    """A short brief on the cited delegation — what it was asked to do and
    whether it was flagged a falsification attempt — so the judge can weigh
    §2 attempt-adequacy. The numeric result the verdict rests on is carried
    separately in ``evidence['numbers']``. None if no delegation is cited or
    the record is absent.
    """
    d = (evidence or {}).get("delegation")
    if not d or self._delegation_log is None:
        return None
    rec = next(
        (r for r in self._delegation_log.query_all() if r.get("id") == d),
        None,
    )
    if rec is None:
        return None
    return (
        f"delegation {d} (to {rec.get('to_node')}, "
        f"is_falsification_attempt={rec.get('is_falsification_attempt')}, "
        f"evals={rec.get('evals')}):\n"
        f"  task: {rec.get('task', '')}\n"
        f"  deliverable: {rec.get('deliverable', '')}"
    )
_run_verdict_validator(h_id: str, status: str, comment: str, evidence: dict | None) -> str #

Advise (never block) on a closing verdict's substance against the charter. Persists the critique on the verdict it judged, emits a VERDICT_SUBSTANCE_FLAG diagnostics event on a flag, and on the Nth repeat flag of the same hypothesis appends a louder gate-critic warning. Returns the text to append to the HypothesisUpdate result ("" = no concern / validator unavailable). NEVER raises — the update it annotates must always stand (Q1=(B), advise-with-teeth).

Source code in src/adda/_src/nodes/critic_gate.py
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
def _run_verdict_validator(
    self, h_id: str, status: str, comment: str, evidence: dict | None,
) -> str:
    """Advise (never block) on a closing verdict's substance against the
    charter. Persists the critique on the verdict it judged, emits a
    ``VERDICT_SUBSTANCE_FLAG`` diagnostics event on a flag, and on the Nth
    repeat flag of the same hypothesis appends a louder gate-critic warning.
    Returns the text to append to the HypothesisUpdate result ("" = no concern
    / validator unavailable). NEVER raises — the update it annotates must
    always stand (Q1=(B), advise-with-teeth).
    """
    if not verdict_validator_enabled():
        return ""  # kill switch (verdict_validator: false) — fully bypassed
    try:
        from ..epistemics.verdict_validator import (
            build_judge_prompt,
            parse_judge_reply,
        )
        h = (self._ledger.get(h_id) or {}) if self._ledger else {}
        if not h:
            return ""
        prompt = build_judge_prompt(
            statement=h.get("statement", ""),
            prediction=h.get("prediction", ""),
            criterion=h.get("falsification_criterion", ""),
            status=status,
            comment=comment,
            evidence=evidence,
            delegation_report=self._cited_delegation_brief(evidence),
            prior_rulings=_prior_rulings_digest(h),
        )
        reply = self._invoke_verdict_validator(prompt)
        if not reply:
            return ""  # no judge / call failed — silent; the update stands
        flagged, critique = parse_judge_reply(reply)
        if not flagged:
            if self._ledger:
                self._ledger.annotate_last(h_id, "validated: no charter concern")
            return ""
        # Flagged: persist on the verdict, surface to the agent, count, escalate.
        if self._ledger:
            self._ledger.annotate_last(h_id, critique)
        self._record_science_drift({
            "error_type": "VERDICT_SUBSTANCE_FLAG",
            "hypothesis": h_id,
            "status": status,
            "message": critique[:300],
        })
        counts = self.__dict__.setdefault("_verdict_flag_counts", {})
        counts[h_id] = counts.get(h_id, 0) + 1
        msg = f"[VERDICT VALIDATOR] {critique}"
        if counts[h_id] >= VERDICT_FLAG_ESCALATE_AFTER:
            msg += (
                f"\n[VERDICT VALIDATOR] {h_id}'s verdict has now been flagged "
                f"{counts[h_id]}× — the gate critic will scrutinise this. "
                "Re-examine the cited charter clause before relying on it."
            )
        return msg
    except Exception:  # noqa: BLE001
        return ""  # advisory must never break the update
_persist_critic_review(n: int, critique_text: str) -> str | None #

Write the critic's full review to debug/critic_reviews/ so the deciding verdict is auditable regardless of PASS/REVISE. Best-effort. Returns the file name written (the caller's delegation row records it), or None when nothing was written.

Source code in src/adda/_src/nodes/critic_gate.py
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
def _persist_critic_review(self, n: int, critique_text: str) -> str | None:
    """Write the critic's full review to debug/critic_reviews/ so the
    deciding verdict is auditable regardless of PASS/REVISE. Best-effort.
    Returns the file name written (the caller's delegation row records
    it), or None when nothing was written."""
    try:
        notes = self._current_notes_dir
        if notes is None:
            return None
        d = Path(notes).parent / "critic_reviews"
        d.mkdir(parents=True, exist_ok=True)
        name = f"call_{n:03d}.md"
        (d / name).write_text(critique_text or "", encoding="utf-8")
        return name
    except Exception:  # noqa: BLE001
        return None
_build_feedback_task_msg(h_ids: list, *, constraints_text: str = '') -> str #

FEEDBACK task message with paths block.

constraints_text is the run's ConstraintSnapshot.as_text() (see constraint_snapshot.py) — this is a delegation like any other (strategizer -> critic), so it carries the same live budget facts automatically, not something the critic has to go query for.

Source code in src/adda/_src/nodes/critic_gate.py
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
def _build_feedback_task_msg(
    self, h_ids: list, *, constraints_text: str = ""
) -> str:
    """<mode>FEEDBACK</mode> task message with paths block.

    ``constraints_text`` is the run's ConstraintSnapshot.as_text()
    (see constraint_snapshot.py) — this is a delegation like any other
    (strategizer -> critic), so it carries the same live budget facts
    automatically, not something the critic has to go query for.
    """
    notes_path = str(self._current_notes_dir or "")
    _notes_dir = self._current_notes_dir
    _study_dir = self._study_dir
    _debug_dir = (
        _notes_dir.parent if _notes_dir is not None else None
    )
    return (
        "<mode>FEEDBACK</mode>\n\n"
        "Perform a synchronous find-only adversarial audit.  "
        "PASS is not an available verdict here — return REVISE or "
        "REJECT with your findings, or NOTED if you find no CRITICAL "
        "or MAJOR issue (NOTED is not an acceptance; only Done()'s "
        "GATE-mode review can close the run).\n\n"
        + problem_statement_block(_study_dir)
        + (constraints_text + "\n\n" if constraints_text else "")
        + "<paths>\n"
        f"study_dir             = {_study_dir}\n"
        f"debug_dir             = {_debug_dir}\n"
        f"delegation_log        = {_debug_dir}/delegation_log.jsonl\n"
        f"diagnostics           = {_debug_dir}/diagnostics.jsonl\n"
        f"strategizer_notes     = {notes_path}\n"
        f"delegations_workspace = {_debug_dir}/delegations/\n"
        f"deliverable            = {_study_dir}/pipeline.ipynb "
        "(the runtime EXECUTES the notebook lazily after this gate to verify "
        "the headline re-derives from the ledger with zero new evals; the "
        "notebook's own markdown cells ARE the writeup — there is no "
        "solution.md, do NOT flag it as missing)\n"
        "</paths>\n\n"
        + evidence_index_block(_debug_dir)
        + f"Focus hypotheses: {h_ids if h_ids else 'all'}"
    )
_role_of(target: str) -> str #

Configured role of a connected node, falling back to its name.

Source code in src/adda/_src/nodes/recording.py
50
51
52
def _role_of(self, target: str) -> str:
    """Configured role of a connected node, falling back to its name."""
    return getattr(self._spec.nodes.get(target), "role", None) or target
_record_usage(usage: dict, *, role: str | None, model: str | None, phase: str | None, delegation_id: str | None) -> None #

Accumulate token totals (decision path) AND emit one telemetry row (additive, off the decision path). Telemetry failures are swallowed inside record_call so they can never break a run.

Source code in src/adda/_src/nodes/recording.py
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
def _record_usage(
    self,
    usage: dict,
    *,
    role: str | None,
    model: str | None,
    phase: str | None,
    delegation_id: str | None,
) -> None:
    """Accumulate token totals (decision path) AND emit one telemetry row
    (additive, off the decision path).  Telemetry failures are swallowed
    inside ``record_call`` so they can never break a run."""
    self._accumulate_usage(usage)
    if self._telemetry is not None:
        self._telemetry.record_call(
            role=role, model=model, phase=phase,
            delegation_id=delegation_id, usage=usage,
        )
_accumulate_usage(usage: dict) -> None #

Thread-safe accumulation of token counts from adapter.last_usage.

Source code in src/adda/_src/nodes/recording.py
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
def _accumulate_usage(self, usage: dict) -> None:
    """Thread-safe accumulation of token counts from adapter.last_usage."""
    with self._registry_lock:
        self._token_totals["input_tokens"] += usage.get("input_tokens", 0) or 0
        self._token_totals["output_tokens"] += usage.get("output_tokens", 0) or 0
        self._token_totals["cache_read_input_tokens"] += usage.get("cache_read_input_tokens", 0) or 0
        self._token_totals["cache_creation_input_tokens"] += usage.get("cache_creation_input_tokens", 0) or 0
        t = self._token_totals
        if has_normalized_usage(usage):
            t["normalized_calls"] = t.get("normalized_calls", 0) + 1
            for f in NORMALIZED_FIELDS:
                t[f] = t.get(f, 0) + int(usage[f])
        else:
            t["legacy_calls"] = t.get("legacy_calls", 0) + 1
        cost = usage.get("total_cost_usd")
        if cost is not None:
            self._token_totals["total_cost_usd"] += cost
            self._cost_observed = True
_record_tool_error(node_name: str, tool_name: str, error_type: str, message: str, tb: str | None = None, args: dict | None = None) -> None #

Increment error counter and append to diagnostics.jsonl (thread-safe).

args are the call's arguments; each value is capped at 2,000 chars.

Source code in src/adda/_src/nodes/recording.py
 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
def _record_tool_error(
    self,
    node_name: str,
    tool_name: str,
    error_type: str,
    message: str,
    tb: str | None = None,
    args: dict | None = None,
) -> None:
    """Increment error counter and append to diagnostics.jsonl (thread-safe).

    ``args`` are the call's arguments; each value is capped at 2,000 chars.
    """
    import json as _json

    # Classify fault: system (transient API/network) vs agent (bad usage).
    fault = _classify_fault(error_type, message)

    with self._registry_lock:
        self._error_counts[node_name] = self._error_counts.get(node_name, 0) + 1
    notes = self._current_notes_dir
    if notes is None:
        return
    # _current_notes_dir is debug/strategizer_notes/; parent is debug/
    debug_dir = Path(notes).parent
    record: dict = {
        "ts": datetime.now(tz=timezone.utc).isoformat(timespec="seconds"),
        "node": node_name,
        "tool": tool_name,
        "error_type": error_type,
        "fault": fault,
        "message": message,
    }
    if tb:
        record["traceback"] = tb
    if args:
        record["args"] = {k: str(v)[:2000] for k, v in args.items()}
    try:
        with (debug_dir / "diagnostics.jsonl").open("a", encoding="utf-8") as f:
            f.write(_json.dumps(record) + "\n")
    except Exception:  # noqa: BLE001
        pass
_record_llm_retry(attempt: int, max_attempts: int, exc: BaseException, delay: float) -> None #

One diagnostics.jsonl row per transient-error retry of an LLM call.

Not an error (no _error_counts bump): the call is retried and may succeed. The row exists so a retry is visible -- how often, which exception, how long the backoff -- instead of silent.

Source code in src/adda/_src/nodes/recording.py
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
def _record_llm_retry(
    self, attempt: int, max_attempts: int, exc: BaseException,
    delay: float,
) -> None:
    """One diagnostics.jsonl row per transient-error retry of an LLM call.

    Not an error (no ``_error_counts`` bump): the call is retried and may
    succeed. The row exists so a retry is visible -- how often, which
    exception, how long the backoff -- instead of silent.
    """
    import json as _json

    from ..backends.base import get_delegation_id

    notes = self._current_notes_dir
    if notes is None:
        return
    record = {
        "ts": datetime.now(tz=timezone.utc).isoformat(timespec="seconds"),
        "node": self._name,
        "delegation_id": get_delegation_id(),
        "error_type": "LLM_RETRY",
        "event": "LLM_RETRY",
        "attempt": attempt,
        "max_attempts": max_attempts,
        "exception": type(exc).__name__,
        "message": str(exc)[:300],
        "delay_s": round(delay, 2),
    }
    try:
        with (Path(notes).parent / "diagnostics.jsonl").open(
                "a", encoding="utf-8") as f:
            f.write(_json.dumps(record) + "\n")
    except Exception:  # noqa: BLE001
        pass
_record_retrospective(role: str, source_id: str, report_text: str, deliverable_version: int | None = None) -> None #

Capture a node's end-of-life ### Retrospective.

Every node has a 'my job is done' moment — workers at delegation completion, the strategizer (and any future orchestrator) at Done(). Parse the section, persist to retrospectives.jsonl, and surface a diagnostic + strategizer notification the instant it flags contradictory system instructions (the cheapest, highest-value failure mode to catch). Best-effort; never raises.

A Done() call that ACTUALLY ARRIVED with non-empty text is always recorded here, even when _extract_report_section fails to find a ### Retrospective heading in it — never silently dropped. Only genuinely empty text (nothing submitted at all) returns without writing anything. This matters because write_fallback_retrospective (infra/watchdog_cleanup.py) trusts "no entry for this role in retrospectives.jsonl" to mean "the retrospective never arrived" and synthesizes a claim saying so; if a real, substantive reply arrived but merely failed to parse (run 20260926T124841's regex bug, now fixed — see _extract_report_section), silently dropping it here made that downstream claim FALSE — the run's own record said a retrospective "never arrived" when one, in fact, had. Recording a parse-failure entry (with the raw text preserved) keeps that claim honest without needing write_fallback_retrospective itself to know the difference.

Source code in src/adda/_src/nodes/recording.py
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
def _record_retrospective(
    self, role: str, source_id: str, report_text: str,
    deliverable_version: int | None = None,
) -> None:
    """Capture a node's end-of-life ### Retrospective.

    Every node has a 'my job is done' moment — workers at delegation
    completion, the strategizer (and any future orchestrator) at Done().
    Parse the section, persist to retrospectives.jsonl, and surface a
    diagnostic + strategizer notification the instant it flags
    contradictory system instructions (the cheapest, highest-value
    failure mode to catch). Best-effort; never raises.

    A Done() call that ACTUALLY ARRIVED with non-empty text is always
    recorded here, even when ``_extract_report_section`` fails to find
    a ``### Retrospective`` heading in it — never silently dropped. Only
    genuinely empty text (nothing submitted at all) returns without
    writing anything. This matters because ``write_fallback_retrospective``
    (infra/watchdog_cleanup.py) trusts "no entry for this role in
    retrospectives.jsonl" to mean "the retrospective never arrived" and
    synthesizes a claim saying so; if a real, substantive reply arrived
    but merely failed to parse (run 20260926T124841's regex bug, now
    fixed — see _extract_report_section), silently dropping it here made
    that downstream claim FALSE — the run's own record said a retrospective
    "never arrived" when one, in fact, had. Recording a parse-failure
    entry (with the raw text preserved) keeps that claim honest without
    needing write_fallback_retrospective itself to know the difference.
    """
    try:
        report_text = report_text or ""
        retro = _extract_retrospective(report_text)
        parse_failed = bool(report_text.strip()) and not retro
        if not retro and not parse_failed:
            return
        notes = self._current_notes_dir
        if notes is None:
            return
        import json as _json
        import re as _re
        debug_dir = Path(notes).parent
        scan_text = retro or report_text
        flagged = bool(_re.search(r"CONSISTENCY:\s*flagged", scan_text, _re.I))
        now = datetime.now(tz=timezone.utc).isoformat(timespec="seconds")
        if parse_failed:
            body = (
                "PARSE FAILURE: the response contained no ### Retrospective "
                f"section. Raw text follows.\n\n{report_text}"
            )
        else:
            body = retro
        rec = {
            "ts": now, "source_id": source_id, "role": role,
            "flagged": flagged, "text": body[:_RETRO_TEXT_CAP],
            "parse_failed": parse_failed,
        }
        if deliverable_version is not None:
            rec["deliverable_version"] = deliverable_version
        with (debug_dir / "retrospectives.jsonl").open(
                "a", encoding="utf-8") as f:
            f.write(_json.dumps(rec) + "\n")
        if flagged:
            drec = {
                "ts": now, "node": role, "tool": "Retrospective",
                "error_type": "CONSISTENCY_FLAG", "fault": "system",
                "message": body[:300],
            }
            with (debug_dir / "diagnostics.jsonl").open(
                    "a", encoding="utf-8") as f:
                f.write(_json.dumps(drec) + "\n")
            with self._notifications_lock:
                self._notifications.append(
                    f"[CONSISTENCY FLAG — {role} ({source_id}) reported "
                    "contradictory system instructions; see "
                    "debug/retrospectives.jsonl]"
                )
    except Exception:  # noqa: BLE001
        pass
_record_intervention(kind: str, target: str, message: str, *, fault: str = 'nudge', **extra) -> None #

Log a scientific-correction event (a nudge/bounce firing) to diagnostics.jsonl — direct evidence the self-healing layer acted, and, when extra carries before/after state, whether it worked.

Neutral classification: fault='nudge' — a nudge is a correction, not an agent error, so this must NOT bump the error/escalation counters (unlike _record_tool_error). fault='observation' is for a row that only records what was seen (adda said nothing to the agent); it also bumps no counter. Best-effort; never raises.

Source code in src/adda/_src/nodes/recording.py
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
def _record_intervention(
    self, kind: str, target: str, message: str, *,
    fault: str = "nudge", **extra
) -> None:
    """Log a scientific-correction event (a nudge/bounce firing) to
    diagnostics.jsonl — direct evidence the self-healing layer acted,
    and, when ``extra`` carries before/after state, whether it worked.

    Neutral classification: fault='nudge' — a nudge is a correction, not
    an agent error, so this must NOT bump the error/escalation counters
    (unlike _record_tool_error). ``fault='observation'`` is for a row that
    only records what was seen (adda said nothing to the agent); it also
    bumps no counter. Best-effort; never raises.
    """
    notes = self._current_notes_dir
    if notes is None:
        return
    import json as _json
    debug_dir = Path(notes).parent
    record: dict = {
        "ts": datetime.now(tz=timezone.utc).isoformat(
            timespec="seconds"),
        "node": target,
        "tool": kind,
        "error_type": kind,
        "fault": fault,
        "message": message,
    }
    record.update(extra)
    try:
        with (debug_dir / "diagnostics.jsonl").open(
                "a", encoding="utf-8") as f:
            f.write(_json.dumps(record) + "\n")
    except Exception:  # noqa: BLE001
        pass
_record_science_drift(payload: dict) -> None #

Append a SCIENCE_DRIFT record to diagnostics.jsonl.

Source code in src/adda/_src/nodes/recording.py
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
def _record_science_drift(self, payload: dict) -> None:
    """Append a SCIENCE_DRIFT record to diagnostics.jsonl."""
    import json as _json
    notes = self._current_notes_dir
    if notes is None:
        return
    record = {
        "ts": datetime.now(tz=timezone.utc).isoformat(
            timespec="seconds"),
        "node": self._name,
        "error_type": "SCIENCE_DRIFT",
        **payload,
    }
    try:
        path = Path(notes).parent / "diagnostics.jsonl"
        with path.open("a", encoding="utf-8") as f:
            f.write(_json.dumps(record) + "\n")
    except Exception:  # noqa: BLE001
        pass
_raise_if_abandoned() -> None #

End the calling thread once the run has stopped waiting for it.

Source code in src/adda/_src/nodes/node.py
147
148
149
150
151
def _raise_if_abandoned(self) -> None:
    """End the calling thread once the run has stopped waiting for it."""
    if self._abandon.is_set():
        raise RunAbandoned(self._name)
    raise_if_stopped(self._name)
_init_recording() -> None #

Establish the state RecordingMixin writes to, on EVERY node.

Recording is a property of a node — any node's tools can return an ERROR, and any node's adapter reports token usage — so its substrate belongs here rather than in one behaviour's setup. It used to be created only by the orchestration setup, which was invisible while leaves were a separate class that did not carry RecordingMixin at all: the moment a leaf reached any recording call it raised AttributeError on _registry_lock. The orchestration setup still assigns these itself, to the same values, so an orchestrating node is unchanged.

Source code in src/adda/_src/nodes/node.py
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
def _init_recording(self) -> None:
    """Establish the state ``RecordingMixin`` writes to, on EVERY node.

    Recording is a property of a node — any node's tools can return an
    ERROR, and any node's adapter reports token usage — so its substrate
    belongs here rather than in one behaviour's setup. It used to be
    created only by the orchestration setup, which was invisible while
    leaves were a separate class that did not carry RecordingMixin at all:
    the moment a leaf reached any recording call it raised AttributeError
    on ``_registry_lock``. The orchestration setup still assigns these
    itself, to the same values, so an orchestrating node is unchanged.
    """
    self._registry_lock = threading.Lock()
    self._abandon = threading.Event()
    self._notifications: list[str] = []
    self._notifications_lock = threading.Lock()
    self._error_counts: dict[str, int] = {}
    self._token_totals: dict = {
        "input_tokens": 0,
        "output_tokens": 0,
        "cache_read_input_tokens": 0,
        "cache_creation_input_tokens": 0,
        # The backend-independent schema (infra/telemetry.py): `legacy_calls`
        # counts calls whose backend reported none of it.
        "fresh_input": 0,
        "cache_read": 0,
        "cache_write": 0,
        "output": 0,
        "normalized_calls": 0,
        "legacy_calls": 0,
        "total_cost_usd": 0.0,
    }
    self._cost_observed: bool = False
    self._current_notes_dir: Path | None = None
    self._telemetry: Any = None
    # Time-budget wrap-up ladder (nodes/_constants.py:budget_band_due):
    # every node's OWN 10%-of-budget bands already reported, so an
    # escalating message fires once per band whether this node
    # orchestrates or answers — a property of any node, like recording.
    self._budget_bands_fired: set[int] = set()
_init_capabilities(*, study_dir: Any = None, workspace_dir: Any = None, delegation_log: DelegationLog | None = None, report_sections: tuple[str, ...] | None = None, agent_tools: frozenset[str] | None = None) -> None #

Wire this node's sandboxed Write and its declared read-only tools.

Runs for EVERY node, not only one reachable as a standalone entry point — see the module docstring: whichever node ends up as someone else's Delegate() target is dispatched through a copy of THIS adapter (adapter.copy() carries the closure_tools set here), so this is where a dispatched specialist's baseline capabilities actually come from, before WorkerSession.install_worker_tools/_sandbox_worker_writes layer the per-delegation overrides (ReportEvals's own record callback, a delegation-scoped Write) on top.

Source code in src/adda/_src/nodes/node.py
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
def _init_capabilities(
    self,
    *,
    study_dir: Any = None,
    workspace_dir: Any = None,
    delegation_log: DelegationLog | None = None,
    report_sections: tuple[str, ...] | None = None,
    agent_tools: frozenset[str] | None = None,
) -> None:
    """Wire this node's sandboxed Write and its declared read-only tools.

    Runs for EVERY node, not only one reachable as a standalone entry
    point — see the module docstring: whichever node ends up as someone
    else's ``Delegate()`` target is dispatched through a copy of THIS adapter
    (``adapter.copy()`` carries the ``closure_tools`` set here), so this is where a
    dispatched specialist's baseline capabilities actually come from,
    before ``WorkerSession.install_worker_tools``/``_sandbox_worker_writes``
    layer the per-delegation overrides (ReportEvals's own record callback,
    a delegation-scoped Write) on top.
    """
    self._study_dir = study_dir
    self._workspace_dir = Path(workspace_dir) if workspace_dir else None
    # The agent's declared tools — the single source of truth for which
    # capability closures this node is granted (read-only ledger/store
    # tools). Kept as a frozenset for membership checks.
    self._agent_tools: frozenset[str] = frozenset(agent_tools or ())
    # A node built without a declared toolset has nothing to narrow by.
    self._tools_declared = agent_tools is not None
    # This agent's declared report sections (e.g. the critic's
    # Findings/Verdict, not the implementer's Conclusions/Files touched).
    # Used to validate a worker's report against ITS OWN contract instead
    # of the implementer-shaped default — otherwise a correct critic or
    # literature report is wrongly flagged malformed (audit BF-10/O40).
    self._report_sections = report_sections
    self._delegation_log = delegation_log
    self._evals_reported: dict = {}
    self._setup_sandboxed_write()
    self.adapter.closure_tools.update(self._build_eval_closures())
    if delegation_log is not None:
        from .tools.routing import build_recall_history
        self.adapter.closure_tools["RecallHistory"] = build_recall_history(self)
    # Declaration-gated read-only ledger/store tools — the SAME builder
    # every node uses, so a specialist (e.g. the critic) gets an
    # identical, working QueryStore/HypothesisList
    # surface whenever it declares them. Resolves the run via the shared
    # Node._resolve_run_dir (delegation-log path).
    from ..runtime import features as _features
    from .tools.routing import build_declared_shared_closures
    self.adapter.closure_tools.update(
        build_declared_shared_closures(
            self, self._agent_tools - _features.disabled_tool_names()))
    # Closures from Agent.build_closure_tools() (ConsultLiterature,
    # CorpusAdd, ...) are installed on the adapter BEFORE this node
    # exists (agent_runtime.py), so they never pass through
    # orchestration.py's _wrap_closure the way routing tools do. A
    # closure tagged _adda_diagnostic_source (a fact about the run
    # environment a tool discovers, not about one call's return value —
    # see LiteratureCorpus.pop_diagnostic_event) needs that wrapper
    # regardless, so re-wrap it here, once, by the tag rather than by
    # name — any future tool needing the same reporting just carries the
    # same tag.
    for _cname, _cfn in list(self.adapter.closure_tools.items()):
        if getattr(_cfn, "_adda_diagnostic_source", None) is not None:
            self.adapter.closure_tools[_cname] = self._wrap_closure(
                _cfn, self._name)
    # A config.yaml `nodes.<name>.tools` list is the whole tool set: the
    # always-on closures it does not name are withheld.
    from ..runtime.node_tools import withheld_closures
    _agent = self._spec.nodes.get(self._name) if self._spec is not None else None
    for _cname in withheld_closures(_agent):
        self.adapter.closure_tools.pop(_cname, None)
_setup_sandboxed_write() -> None #

Replace native Write with a workspace-sandboxed closure.

Removes 'Write' from native_tools so the SDK doesn't expose it, then installs the same builder a dispatched worker gets (:func:build_sandboxed_write), scoped to this node's whole workspace instead of one delegation's subfolder — a node reached via real graph routing (the module docstring) has no delegation_id of its own to scope tighter than that.

Bash is kept native but cwd is already set to study_dir by the adapter; the prompt further constrains it to the workspace.

Source code in src/adda/_src/nodes/node.py
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
def _setup_sandboxed_write(self) -> None:
    """Replace native Write with a workspace-sandboxed closure.

    Removes 'Write' from native_tools so the SDK doesn't expose it,
    then installs the same builder a dispatched worker gets
    (:func:`build_sandboxed_write`), scoped to this node's whole
    workspace instead of one delegation's subfolder — a node reached via
    real graph routing (the module docstring) has no delegation_id of
    its own to scope tighter than that.

    Bash is kept native but cwd is already set to study_dir by the adapter;
    the prompt further constrains it to the workspace.
    """
    if self._workspace_dir is None:
        return  # no sandboxing if study_dir unknown (e.g. tests)

    # Remove native Write so the SDK doesn't expose an unrestricted version
    if hasattr(self.adapter, "native_tools") and "Write" in self.adapter.native_tools:
        self.adapter.native_tools = [
            t for t in self.adapter.native_tools if t != "Write"
        ]

    from .tools.routing import build_sandboxed_write

    # Bound to a plain name (not inlined into the assignment) so
    # internal/tools/promptmap.py's injected_tool_docs() scanner — which
    # resolves a closure_tools["Write"] = <name> rebind to the function
    # <name> refers to — can still find this tool's docstring after the
    # move into the shared builder.
    _study_ws = (
        Path(self._study_dir) / "workspace"
        if self._study_dir is not None else None
    )
    Write = build_sandboxed_write(
        self._workspace_dir, study_workspace=_study_ws)
    self.adapter.closure_tools["Write"] = Write
_build_eval_closures() -> dict #
Source code in src/adda/_src/nodes/node.py
303
304
305
306
307
308
309
310
311
def _build_eval_closures(self) -> dict:
    from .tools.routing import build_report_evals

    evals = self._evals_reported
    return {
        "ReportEvals": build_report_evals(
            record=lambda n: evals.__setitem__("count", n)
        )
    }
_commit_workspace(message: str) -> str | None #

Commit the run workspace and return the sha (spec 11).

On Node rather than on WorkerSession because a delegation is recorded from four places, not one: a worker finishing (ok or error), the critic gate, an AskForFeedback audit, and the close-time reconciliation of a delegation still running when the run ended. Every one of them appends a row to the delegation log, so every one of them owes that row a sha — otherwise workspace_sha: null means both "this delegation changed nothing" and "nobody looked", which is the ambiguity --allow-empty exists to prevent.

The interrupted case is the one with teeth: a delegation killed mid-flight never reaches _finish_ok/_finish_error, so its partial file writes stay uncommitted and are swept into whichever delegation commits NEXT — attributing one delegation's work to another. Silence would be better than that; a commit is better still.

Never raises — see infra/workspace_vcs.

Source code in src/adda/_src/nodes/node.py
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
def _commit_workspace(self, message: str) -> str | None:
    """Commit the run workspace and return the sha (spec 11).

    On Node rather than on WorkerSession because a delegation is recorded
    from four places, not one: a worker finishing (ok or error), the
    critic gate, an AskForFeedback audit, and the close-time
    reconciliation of a delegation still running when the run ended. Every
    one of them appends a row to the delegation log, so every one of them
    owes that row a sha — otherwise ``workspace_sha: null`` means both
    "this delegation changed nothing" and "nobody looked", which is the
    ambiguity --allow-empty exists to prevent.

    The interrupted case is the one with teeth: a delegation killed
    mid-flight never reaches _finish_ok/_finish_error, so its partial file
    writes stay uncommitted and are swept into whichever delegation
    commits NEXT — attributing one delegation's work to another. Silence
    would be better than that; a commit is better still.

    Never raises — see infra/workspace_vcs.
    """
    run_dir = self._resolve_run_dir()
    if run_dir is None:
        return None
    from ..infra.workspace_vcs import commit_workspace
    return commit_workspace(run_dir / "debug" / "delegations", message)
_resolve_run_dir() -> Path | None #

Best-effort run_dir, valid on any node.

Every node's _current_notes_dir is set at construction (from notes_dir, passed to every node — see graph_builder.py) and re-pointed each turn by _absorb_state if the run's real notes dir differs. It can still be None (e.g. a bare Node() built without one, as in unit tests) — every node DOES hold the shared delegation log at run_dir/debug/delegation_log.jsonl, so derive run_dir from that when the notes dir is unavailable.

Source code in src/adda/_src/nodes/node.py
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
def _resolve_run_dir(self) -> Path | None:
    """Best-effort run_dir, valid on any node.

    Every node's ``_current_notes_dir`` is set at construction (from
    ``notes_dir``, passed to every node — see ``graph_builder.py``) and
    re-pointed each turn by ``_absorb_state`` if the run's real notes dir
    differs. It can still be ``None`` (e.g. a bare ``Node()`` built
    without one, as in unit tests) — every node DOES hold the shared
    delegation log at run_dir/debug/delegation_log.jsonl, so derive
    run_dir from that when the notes dir is unavailable.
    """
    notes = getattr(self, "_current_notes_dir", None)
    if notes is not None:
        return Path(notes).parent.parent
    dlog = getattr(self, "_delegation_log", None)
    p = getattr(dlog, "_path", None)
    return Path(p).parent.parent if p is not None else None
_read_ledger() -> HypothesisLedger | None #

The hypothesis ledger for READ access, resolved on any node.

Returns the node's own bound ledger when it has one (the entry node); otherwise resolves a read-only view from the run's strategizer_notes when hypotheses.json exists. HypothesisLedger.init performs no I/O, and only READ callers use this, so there is no write race with the entry node that owns the file.

Source code in src/adda/_src/nodes/node.py
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
def _read_ledger(self) -> HypothesisLedger | None:
    """The hypothesis ledger for READ access, resolved on any node.

    Returns the node's own bound ledger when it has one (the entry node);
    otherwise resolves a read-only view from the run's strategizer_notes
    when hypotheses.json exists. HypothesisLedger.__init__ performs no I/O,
    and only READ callers use this, so there is no write race with the
    entry node that owns the file.
    """
    own = getattr(self, "_ledger", None)
    if own is not None:
        return own
    rd = self._resolve_run_dir()
    if rd is None:
        return None
    notes = rd / "debug" / "strategizer_notes"
    if (notes / "hypotheses.json").exists():
        from ..epistemics.hypothesis_ledger import HypothesisLedger
        return HypothesisLedger(notes)
    return None

adda.build_graph(spec: Graph, make_adapter: Callable[[str, Agent], Any], checkpointer: Any = None, study_dir: Any = None, interactive: bool = False, max_ask: int = 1, notes_dir: Any = None, workspace_dir: Any = None, delegation_log: DelegationLog | None = None, node_registry: dict[str, Any] | None = None) -> Any #

Build and compile a LangGraph StateGraph from a Graph spec.

Parameters:

Name Type Description Default
spec Graph

Agent graph specification (nodes, edges, entry).

required
make_adapter callable

(name: str, agent: Agent) -> adapter — factory that produces a ClaudeAdapter or OllamaAdapter for the given node.

required
checkpointer any

LangGraph checkpointer. Defaults to an in-memory :class:MemorySaver.

None
delegation_log DelegationLog

Graph-wide delegation log for episodic memory (RecallHistory tool).

None
node_registry dict

If given, populated in place with {name: Node instance} as each node is constructed — an explicit, caller-owned way to reach the live Node objects afterward (e.g. spec 12's close-time open-review sweep) instead of reaching through the COMPILED graph's own internals (compiled.nodes[name].bound.func), which is a LangGraph implementation detail that could silently break on an upgrade.

None

Returns:

Type Description
CompiledGraph

A compiled LangGraph graph ready to invoke.

Source code in src/adda/_src/runtime/graph_builder.py
 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
def build_graph(
    spec: Graph,
    make_adapter: Callable[[str, Agent], Any],
    checkpointer: Any = None,
    study_dir: Any = None,
    interactive: bool = False,
    max_ask: int = 1,
    notes_dir: Any = None,
    workspace_dir: Any = None,
    delegation_log: DelegationLog | None = None,
    node_registry: dict[str, Any] | None = None,
) -> Any:
    """Build and compile a LangGraph StateGraph from a Graph spec.

    Parameters
    ----------
    spec : Graph
        Agent graph specification (nodes, edges, entry).
    make_adapter : callable
        ``(name: str, agent: Agent) -> adapter`` — factory that produces a
        ``ClaudeAdapter`` or ``OllamaAdapter`` for the given node.
    checkpointer : any, optional
        LangGraph checkpointer.  Defaults to an in-memory :class:`MemorySaver`.
    delegation_log : DelegationLog, optional
        Graph-wide delegation log for episodic memory (RecallHistory tool).
    node_registry : dict, optional
        If given, populated in place with ``{name: Node instance}`` as each
        node is constructed — an explicit, caller-owned way to reach the
        live Node objects afterward (e.g. spec 12's close-time open-review
        sweep) instead of reaching through the COMPILED graph's own
        internals (``compiled.nodes[name].bound.func``), which is a
        LangGraph implementation detail that could silently break on an
        upgrade.

    Returns
    -------
    CompiledGraph
        A compiled LangGraph graph ready to invoke.
    """
    settings.set_graph_nodes(spec.nodes)
    builder = StateGraph(AgenticState)

    # ONE adapter per named node — shared across all orchestrating nodes.
    node_adapters = {n: make_adapter(n, spec.nodes[n]) for n in spec.nodes}

    live_nodes: dict[str, Any] = (
        node_registry if node_registry is not None else {})

    for name, agent in spec.nodes.items():
        adapter = node_adapters[name]  # shared instance, NOT make_adapter() again
        outgoing = spec.outgoing(name)

        # ONE node class; every node runs the same turn loop
        # (Node._orchestrate). What a node may DO follows from what its
        # Agent declares in `tools` — never from its topology. notes_dir is
        # passed to every node: a delegating node is simply a node that needs
        # help from another node (CLAUDE.md "all nodes are equal"), and
        # telemetry / the science monitor / hypothesis-ledger READ access
        # matter for every role, not only the entry node. WRITE access
        # (HypothesisPropose/Update, Milestone*) stays gated separately, by
        # each Agent's own declared `tools` (see
        # nodes/tools/routing/__init__.py) — passing notes_dir here grants no
        # capability a node has not already declared. A node with no
        # outgoing edges still receives this argument (Node.__init__ has one
        # init path for every node), but never ACQUIRES ledger ownership from
        # it — Node._init_orchestration gates `_owns_epistemics` on having
        # outgoing edges, not merely on notes_dir being set, so this only
        # takes effect for a node that itself has outgoing edges.
        node = Node(
            adapter,
            name=name,
            outgoing=outgoing,
            spec=spec,
            study_dir=study_dir,
            interactive=interactive,
            max_ask=max_ask,
            worker_adapters={n: node_adapters[n] for n in outgoing},
            notes_dir=notes_dir,
            workspace_dir=workspace_dir,
            delegation_log=delegation_log,
            report_sections=getattr(agent, "report_sections", None),
            agent_tools=getattr(agent, "tools", None),
        )
        live_nodes[name] = node

        builder.add_node(name, node)

    # A delegation's registry entry lives in its DELEGATOR's Node, while the
    # worker calls the tool closures bound to its OWN Node -- so every node
    # must be able to find the others (Node._delegation_entry).
    from ..nodes.slots import AwakeSlots
    from . import features
    awake_slots = AwakeSlots(features.max_awake_nodes())
    for node in live_nodes.values():
        node._peers = live_nodes
        node._awake_slots = awake_slots

    builder.set_entry_point(spec.entry)

    return builder.compile(checkpointer=checkpointer or MemorySaver())

Run state#

The types that move through the graph while it runs. You rarely construct these yourself; they are what you see in a transcript or a checkpoint.

adda.AgenticState #

LangGraph state for one agentic run.

Inherits messages: Annotated[list[AnyMessage], add_messages] from :class:~langgraph.graph.MessagesState.

Source code in src/adda/_src/runtime/graph_state.py
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
class AgenticState(MessagesState):
    """LangGraph state for one agentic run.

    Inherits ``messages: Annotated[list[AnyMessage], add_messages]``
    from :class:`~langgraph.graph.MessagesState`.
    """

    study_dir: str
    done: bool
    last_report: str | None
    total_delegations: int
    budget_seconds: float | None
    budget_usd: float | None     # hard USD cost ceiling; None = no USD ceiling
    run_dir: str | None          # absolute path to runs/<timestamp>/
    eval_budget: int | None      # max function evaluations (from config)
    evals_used: int              # running count across delegations
    # time.time() at run start, for budget enforcement
    start_time: float | None
    # paths relative to study_dir; checked before Done accepted
    required_deliverables: list | None
    # Token usage accumulated across all agents; set by the entry node on Done
    token_totals: dict | None
    # Per-node tool-call error counts (ERROR: returns + exceptions)
    error_counts: dict | None
    # Canonical run-level ExperimentData store project_dir.
    # Set by AgenticRun.execute() so nodes and workers can locate the
    # shared store without re-deriving it from run_dir.
    experiment_data_dir: str | None
    # Structured terminal state, recorded where the decision is MADE and
    # carried out as data — never re-derived by grepping the final report's
    # banner, which is how a backstop halt used to read as a validated
    # success. Vocabulary + fail-safe resolution live in runtime.terminal.
    outcome: str | None       # GATED | UNGATED | FAILED
    termination: str | None   # done | backstop_time | crashed | …
    reviewed: bool | None     # a critic gate actually ran
study_dir: str instance-attribute #
done: bool instance-attribute #
last_report: str | None instance-attribute #
total_delegations: int instance-attribute #
budget_seconds: float | None instance-attribute #
budget_usd: float | None instance-attribute #
run_dir: str | None instance-attribute #
eval_budget: int | None instance-attribute #
evals_used: int instance-attribute #
start_time: float | None instance-attribute #
required_deliverables: list | None instance-attribute #
token_totals: dict | None instance-attribute #
error_counts: dict | None instance-attribute #
experiment_data_dir: str | None instance-attribute #
outcome: str | None instance-attribute #
termination: str | None instance-attribute #
reviewed: bool | None instance-attribute #

adda.StudyConfig dataclass #

Lightweight config loaded from the study directory.

Source code in src/adda/_src/runtime/graph_state.py
82
83
84
85
86
87
88
@dataclass
class StudyConfig:
    """Lightweight config loaded from the study directory."""

    study_dir: str
    model: str = "claude-opus-4-5"
    budget_seconds: float | None = None
study_dir: str instance-attribute #
model: str = 'claude-opus-4-5' class-attribute instance-attribute #
budget_seconds: float | None = None class-attribute instance-attribute #

adda.Task dataclass #

A unit of work delegated from one agent to another.

Source code in src/adda/_src/runtime/graph_state.py
52
53
54
55
56
57
58
@dataclass
class Task:
    """A unit of work delegated from one agent to another."""

    intent: str
    expected_report: str
    target: str
intent: str instance-attribute #
expected_report: str instance-attribute #
target: str instance-attribute #

adda.Report dataclass #

The result returned by an implementer agent.

Source code in src/adda/_src/runtime/graph_state.py
61
62
63
64
65
66
@dataclass
class Report:
    """The result returned by an implementer agent."""

    content: str
    target: str | None = None
content: str instance-attribute #
target: str | None = None class-attribute instance-attribute #

adda.Delegation dataclass #

A delegation request parsed from an agent's response.

Source code in src/adda/_src/runtime/graph_state.py
69
70
71
72
73
74
75
76
77
78
79
@dataclass
class Delegation:
    """A delegation request parsed from an agent's response."""

    target: str
    task: str
    expected_report: str = ""
    # Optional design namespace this delegation operates in (Axis 3a). None →
    # the canonical single-study oracle + ledger (today's behavior). A non-None
    # value scopes the worker to that namespace's oracle via F3DASM_NAMESPACE.
    namespace: str | None = None
target: str instance-attribute #
task: str instance-attribute #
expected_report: str = '' class-attribute instance-attribute #
namespace: str | None = None class-attribute instance-attribute #