File size: 20,743 Bytes
9d0fd45
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
"""task_board β€” a per-run task board for the coordinator (main agent)."""

from __future__ import annotations

import logging
from typing import Any

from frontier_agent.components.agent_bus import AgentBus
from frontier_agent.components.observers.task_board import TaskBoardObserver
from frontier_agent.components.task_board_types import (
    BOARD_TOOLS,
    RESOLUTION_MARKS,
    VALID_RESOLUTION,
    BoardCounts,
    count_resolutions,
)
from frontier_agent.core.execution_context import get_current_execution_scope
from frontier_agent.core.runtime import registry
from frontier_agent.core.tool import tool
from plugins.tools._bus_scope import resolve_bus_task_id
from plugins.tools._coerce import coerce_json_list, coerce_json_object

logger = logging.getLogger(__name__)

# task_id -> {"seq": int, "tasks": {id: {description, resolution, owners, group}}}
_BOARDS: dict[str, dict[str, Any]] = {}

# task_id -> list of pending board write-ops the stream observer has not yet
# drained. Each op is ``{"op": "add"|"update"|"finish_planning", "ids": [...],
# "phase": "planning"|"execution"}``. The board tools append here on every write;
# A task-board stream observer drains them and emits one
# ``response.swarm.task_board`` frame per op. Recording is unconditional + cheap (the
# eval / HTTP paths simply never drain), and ``clear_board`` empties it at run end
# so it can never leak across trials.
_PENDING_OPS: dict[str, list[dict[str, Any]]] = {}


def _record_op(task_id: str, op: str, ids: list[str] | None = None) -> None:
    """Append a board write-op for the stream observer to drain.

    The op stamps the phase AT WRITE TIME. In the agent-team two-loop profile the
    stream observer is not attached to the planning loop, so ops written there are
    drained late (once the execution loop runs) β€” by then the phase has flipped to
    ``execution``. Freezing the write-time phase keeps a planning-loop ``add_task``
    rendered as ``phase=planning`` rather than the phase current at drain time.
    """
    _PENDING_OPS.setdefault(task_id, []).append(
        {"op": op, "ids": list(ids or []), "phase": current_phase(task_id)}
    )


# ── State helpers (also used by the observer + finalize_answer gate) ────────

def _board(task_id: str) -> dict[str, Any]:
    return _BOARDS.setdefault(task_id, {"seq": 0, "tasks": {}})


def board_size(task_id: str) -> int:
    """Number of tasks on this task's board (0 if no board)."""
    b = _BOARDS.get(task_id)
    return len(b["tasks"]) if b else 0


def build_task_board_observer(*, cooldown_turns: int = 5) -> TaskBoardObserver:
    """Bind the shared observer to this plugin's task-board state."""
    return TaskBoardObserver(
        board_size=board_size,
        render_board=lambda task_id, bus_task_id: render_board(
            task_id, bus_task_id=bus_task_id,
        ),
        resolve_bus_task_id=resolve_bus_task_id,
        cooldown_turns=cooldown_turns,
    )


def clear_board(task_id: str) -> None:
    """Drop a task's board + phase β€” called at the end of the main-agent run."""
    _BOARDS.pop(task_id, None)
    _PHASE.pop(task_id, None)
    _PENDING_OPS.pop(task_id, None)


# ── Task-board stream projection β€” pure read helpers ─────────────────────────
# Consumed by the protocol layer's task-board stream observer
# to build the ``task_board.*`` wire frames. Kept here (next to the state they
# read) and side-effect-free so the observer stays a thin emitter.

def current_phase(task_id: str) -> str:
    """The board's coordinator phase β€” ``planning`` or ``execution``.

    Defaults to ``execution`` when the run never entered Planning Mode
    (``planning_mode:false`` profiles β€” apodex production β€” never call
    ``start_planning``, so ``_PHASE`` has no entry and the board is live in
    execution from the first ``add_task``)."""
    return _PHASE.get(task_id, "execution")


def serialize_tasks(task_id: str, ids: list[str]) -> list[dict[str, Any]]:
    """Render the given task ids to the wire shape, reading CURRENT board state.

    Reading the live ``resolution`` (rather than assuming ``open``) matters: an
    ``add_task`` that de-dups onto an already-``resolved`` id must carry that
    real resolution, or the frontend's whole-row replace would roll the task
    back to ``open`` and regress progress. Ids no longer on the
    board (raced clear) are skipped rather than emitted as ghosts."""
    b = _BOARDS.get(task_id)
    if not b:
        return []
    out: list[dict[str, Any]] = []
    for tid in ids:
        t = b["tasks"].get(tid)
        if t is None:
            continue
        out.append({
            "id": tid,
            "description": t.get("description", ""),
            "resolution": t.get("resolution", "open"),
            "owners": list(t.get("owners") or []),
            "group": t.get("group", ""),
        })
    return out


def snapshot_tasks(task_id: str) -> list[dict[str, Any]]:
    """Return the complete board in display order for UI projections.

    Unlike :func:`render_board`, this keeps task data structured so a client
    can render a real board instead of scraping the human-readable tool result.
    """
    b = _BOARDS.get(task_id)
    if not b:
        return []
    return serialize_tasks(task_id, list(b["tasks"]))


def drain_board_ops(task_id: str) -> list[dict[str, Any]]:
    """Pop and return all pending board write-ops for this task (FIFO)."""
    return _PENDING_OPS.pop(task_id, [])


# ── Planning Mode (two-phase state machine) ─────────────────────────────────
# task_id -> "planning" | "execution". Default (no entry) = execution, so ONLY
# the agent-team main agent (which calls start_planning) is ever gated; every
# other caller / workflow is unaffected.
_PHASE: dict[str, str] = {}
# In Planning Mode the agent is a PLANNER, not a solver: it may use only
# READ-ONLY tools (to understand the problem / look up a basic term) plus the
# board tools. EVERYTHING ELSE β€” team building, dispatch, code/file writes,
# finalize β€” is blocked until it calls finish_planning. This is an ALLOWLIST
# (the complement is blocked) so a newly-added write tool is denied by default.
_PLANNING_ALLOWED = (
    # read-only inspection / understanding
    "grep_search", "glob_search", "read_file", "read_text", "view_image",
    "web_search", "web_fetch",
    # board tools
    "add_task", "update_task", "finish_planning",
)


def start_planning(task_id: str) -> None:
    _PHASE[task_id] = "planning"


def force_finish_planning(task_id: str) -> None:
    """Flip a task out of Planning Mode WITHOUT the empty-board guard the
    ``finish_planning`` tool enforces. Used by the planning-turn-cap path
    (auto-finish at ``planning_max_turns``) and by the two-loop driver after a
    planning loop that ended on max_turns rather than an explicit finish."""
    if task_id in _PHASE:
        _PHASE[task_id] = "execution"


def is_planning_allowed(tool_name: str) -> bool:
    """True if ``tool_name`` may run during Planning Mode (read-only / board)."""
    return tool_name in _PLANNING_ALLOWED


def in_planning(task_id: str) -> bool:
    return _PHASE.get(task_id) == "planning"


def planning_enabled(task_id: str) -> bool:
    """True if this run enabled Planning Mode (start_planning was called) β€”
    stays True after finish_planning, so the no-solo / verify finalize gates
    keep applying through execution. ``False`` for non-planning runs."""
    return task_id in _PHASE


def planning_block_message(task_id: str, tool_name: str) -> str | None:
    """Error string to return when ``tool_name`` is called during Planning Mode;
    ``None`` when the call is allowed. No-op unless start_planning() was called
    for this task (i.e. only the agent-team main agent)."""
    if in_planning(task_id) and not is_planning_allowed(tool_name):
        return (
            f"Blocked: `{tool_name}` is unavailable in PLANNING MODE. You are a "
            "PLANNER here β€” use only READ-ONLY tools (grep_search / glob_search "
            "/ web_search to understand the problem) and the board tools "
            "(add_task / update_task). List EVERY sub-question via add_task, "
            "then call finish_planning() to unlock the team."
        )
    return None


def unresolved_task_ids(task_id: str) -> list[str]:
    """Ids not yet finished β€” open or in_progress. ``[]`` if no board.

    Cancelled tasks are retracted work (dropped on purpose / created in error),
    so they do NOT count as unresolved and never hold back the finalize gate."""
    b = _BOARDS.get(task_id)
    if not b:
        return []
    return [
        tid for tid, t in b["tasks"].items()
        if t.get("resolution") not in ("resolved", "cancelled")
    ]


def _exec_status_by_owner(task_id: str) -> dict[str, str]:
    """Map owner sub-agent name -> coarse execution status, derived from the bus.

    This is the system-owned half of the board: the model never writes it, so
    the model's ``resolution`` verdict can never clobber the live run state.
    """
    bus = registry.get_optional(AgentBus)
    if bus is None:
        return {}
    out: dict[str, str] = {}
    try:
        sessions = bus.list_sessions_for_task(task_id)
    except Exception:
        return {}
    for s in sessions:
        if getattr(s, "current_job_id", None) is not None:
            st = "running"
        elif getattr(s, "pending_tasks", None):
            st = "queued"
        elif getattr(s, "total_task_count", 0) == 0:
            st = "created"
        else:
            st = "reported"
        out[s.name] = st
    return out


def render_board(task_id: str, *, bus_task_id: str | None = None) -> str:
    """Render the board, joining model fields with live bus exec-status."""
    b = _BOARDS.get(task_id)
    if not b or not b["tasks"]:
        return "[task board] empty β€” call add_task to register sub-questions."
    tasks = b["tasks"]
    c = count_resolutions(t["resolution"] for t in tasks.values())
    lines = [
        f"[task board] resolved {c.resolved}/{c.active} Β· "
        f"in-progress {c.in_progress} Β· "
        f"open {c.open} Β· cancelled {c.cancelled}"
    ]
    # The owners/exec column only makes sense when there IS a team: in a solo
    # run (e.g. react) no task is ever assigned, so showing "agents=[unassigned]"
    # on every row is pure noise. Drop the whole column unless at least one task
    # has an owner β€” then skip the bus query too (nothing to join against).
    show_owners = any(t.get("owners") for t in tasks.values())
    exec_by_owner = _exec_status_by_owner(bus_task_id or task_id) if show_owners else {}
    for tid, t in tasks.items():
        row = (
            f"  {RESOLUTION_MARKS.get(t['resolution'], 'β—‹')} {tid} "
            f"{t['resolution']:<11} "
        )
        if show_owners:
            owners = t.get("owners") or []
            # one "name:exec_status" per owning agent, so the coordinator sees
            # exactly who is on this task and how far each has got.
            agents = (
                " Β· ".join(f"{o}:{exec_by_owner.get(o, '?')}" for o in owners)
                if owners else "unassigned"
            )
            row += f"agents=[{agents}]  "
        lines.append(f"{row}{t['description'][:80]}")
    return "\n".join(lines)


# ── Tools ───────────────────────────────────────────────────────────────────

def _as_owner_list(val: Any) -> list[str]:
    """Normalise an ``owner`` field (a name, a comma-string, or a list of names)
    into a clean list of agent names. One task can have MANY owners β€” several
    agents attacking the SAME sub-question from different angles for
    corroboration are all owners of that one task (not separate tasks)."""
    if val is None:
        return []
    items = val if isinstance(val, list) else str(val).split(",")
    out: list[str] = []
    for x in items:
        name = str(x).strip()
        if name and name not in out:
            out.append(name)
    return out


@tool
async def add_task(tasks: list[Any]) -> str:
    """Register the work items / sub-questions for this run on the task board.

    This is your plan and external memory. Call it up front β€” before doing any
    real work (fetching, running code, gathering evidence) β€” once you've broken
    the question into the concrete steps you'll work through, and again whenever
    a new sub-question emerges. (Reasoning, and a few clarifying searches to
    understand the problem, may come first.) Duplicate descriptions are
    de-duplicated (you get the existing id).

    Args:
        tasks: a list, each item ``{"description": str}``:
            - description (required): ONE concrete, checkable work item, e.g.
              "Verify the Markov condition Δω ≫ system rate holds" β€” not a vague
              area like "look into the math".
            - owner (OPTIONAL, multi-agent runs only): if a teammate sub-agent is
              already assigned to this item, name it here (a name, comma-string,
              or list β€” a task may have many). In a solo run, just OMIT it.

    Returns:
        The assigned ids plus the rendered board.
    """
    items = coerce_json_list(tasks) or []
    scope = get_current_execution_scope()
    if scope is None:
        return "Error: add_task can only be called inside an active run."
    if not items:
        return "Error: add_task requires at least one {description} item."
    b = _board(scope.task_id)
    existing = {t["description"].strip(): tid for tid, t in b["tasks"].items()}
    new_ids: list[str] = []
    skipped = 0  # items that weren't usable {description} objects
    for raw in items:
        it = coerce_json_object(raw)
        if it is None:
            skipped += 1
            continue
        desc = str(it.get("description", "")).strip()
        if not desc:
            skipped += 1
            continue
        if desc in existing:  # dedup re-decomposition
            new_ids.append(existing[desc])
            continue
        b["seq"] += 1
        tid = f"t{b['seq']}"
        b["tasks"][tid] = {
            "description": desc,
            "resolution": "open",
            "owners": _as_owner_list(it.get("owner") or it.get("owners")),
            "group": str(it.get("group", "")).strip(),
        }
        existing[desc] = tid
        new_ids.append(tid)
    logger.info(
        "add_task(task=%s): +%d skipped=%d (ids=%s)",
        scope.task_id, len(new_ids), skipped, new_ids,
    )
    # Nothing landed but items were passed β†’ the shape was wrong. Return a
    # corrective error (not a silent "Added []") so the model re-sends the right
    # shape instead of burning turns repeating the mistake.
    if not new_ids:
        return (
            'Error: add_task expects a list of objects like '
            '[{"description": "..."}]; none of the items were usable, so nothing '
            'was added. Re-call with each task as its own '
            '{"description": "<one concrete work item>"}.'
        )
    _record_op(scope.task_id, "add", new_ids)
    msg = f"Added {new_ids}.\n{render_board(scope.task_id, bus_task_id=resolve_bus_task_id(scope))}"
    if skipped:
        msg += (
            f'\nNote: {skipped} item(s) were skipped (not a {{"description": ...}} '
            "object). Re-add them with that shape if still needed."
        )
    return msg


@tool
async def update_task(updates: list[Any]) -> str:
    """Mark progress on the board β€” call this the MOMENT a task finishes.

    Update each task as soon as it is done, before starting the next one; don't
    let finished tasks pile up. The arg is a LIST only to cover the case where
    two tasks finished in the SAME turn (resolve both at once) β€” it is NOT a
    reason to batch resolutions across turns. Returns the full updated board.

    Args:
        updates: a list of ``{"id", "resolution"?}`` items:
            - id (required): task id from add_task, e.g. "t3".
            - resolution: "in_progress" (set this the MOMENT you start working a
              task β€” it marks the one you are on now) | "resolved" (the work item
              is answered AND corroborated β€” your judgment, not merely "I glanced
              at it") | "cancelled" (retract a task you no longer need or created
              in error β€” it stops counting toward unresolved work; the id stays
              on the board for the trail) | "open" (the default; not started).
            - owner (OPTIONAL, multi-agent runs only): teammate sub-agent
              name(s) now working it β€” ADDED to the task's owner list (a name,
              comma-string, or list). Set "replace_owners": true to overwrite
              instead of add. Omit entirely in a solo run.

    Returns:
        The full updated task board (+ any per-item errors).
    """
    items = coerce_json_list(updates) or []
    scope = get_current_execution_scope()
    if scope is None:
        return "Error: update_task can only be called inside an active run."
    b = _BOARDS.get(scope.task_id)
    if not b:
        return "Error: no task board yet β€” call add_task first."
    changed: list[str] = []
    errors: list[str] = []
    for raw in items:
        u = coerce_json_object(raw)
        if u is None:
            continue
        tid = str(u.get("id", ""))
        if tid not in b["tasks"]:
            errors.append(f"{tid or '?'}: no such task")
            continue
        res = str(u.get("resolution", "")).strip()
        if res and res not in VALID_RESOLUTION:
            errors.append(f"{tid}: bad resolution {res!r} (use {VALID_RESOLUTION})")
            continue
        t = b["tasks"][tid]
        t.setdefault("owners", [])
        if res:
            t["resolution"] = res
        # owners ACCUMULATE β€” assigning another agent to the same task adds it,
        # it does not replace the existing owner(s). ``replace_owners: true``
        # overwrites (e.g. to drop an agent that was reassigned elsewhere).
        new_owners = _as_owner_list(u.get("owner") or u.get("owners"))
        if new_owners:
            if u.get("replace_owners"):
                t["owners"] = new_owners
            else:
                for o in new_owners:
                    if o not in t["owners"]:
                        t["owners"].append(o)
        changed.append(tid)
    if not changed and not errors:
        return ("Error: update_task needs a list of "
                "{id, resolution?, owner?} items.")
    # Emit ONLY the actually-changed ids β€” rejected ids (bad id /
    # bad resolution) stay off the wire so the frontend never builds ghost rows.
    if changed:
        _record_op(scope.task_id, "update", changed)
    msg = f"Updated {changed}.\n{render_board(scope.task_id, bus_task_id=resolve_bus_task_id(scope))}"
    if errors:
        msg += "\nerrors: " + "; ".join(errors)
    return msg


@tool
async def finish_planning() -> str:
    """Leave Planning Mode and start building the team.

    While planning, only add_task / update_task are available. Call this once
    your task board lists EVERY sub-question; afterwards create_subagent /
    assign_task / collect_reports become available (you can still add_task /
    update_task to refine the plan as the investigation unfolds).
    """
    scope = get_current_execution_scope()
    if scope is None:
        return "Error: finish_planning can only be called inside an active run."
    if board_size(scope.task_id) == 0:
        return (
            "Error: the task board is empty β€” call add_task to list the "
            "sub-questions before finishing planning."
        )
    _PHASE[scope.task_id] = "execution"
    _record_op(scope.task_id, "finish_planning")
    return (
        "Planning complete β€” now in EXECUTION mode. You may create_subagent / "
        "assign_task to build and dispatch the team.\n"
        + render_board(scope.task_id, bus_task_id=resolve_bus_task_id(scope))
    )


__all__ = [
    "BOARD_TOOLS",
    "RESOLUTION_MARKS",
    "VALID_RESOLUTION",
    "BoardCounts",
    "add_task",
    "board_size",
    "build_task_board_observer",
    "clear_board",
    "count_resolutions",
    "current_phase",
    "drain_board_ops",
    "finish_planning",
    "force_finish_planning",
    "in_planning",
    "is_planning_allowed",
    "planning_block_message",
    "planning_enabled",
    "render_board",
    "serialize_tasks",
    "snapshot_tasks",
    "start_planning",
    "unresolved_task_ids",
    "update_task",
]