diff --git a/docs/adr/010-delegation/ADR-0010-delegation.md b/docs/adr/010-delegation/ADR-0010-delegation.md index f95a8ad..fdf6fd2 100644 --- a/docs/adr/010-delegation/ADR-0010-delegation.md +++ b/docs/adr/010-delegation/ADR-0010-delegation.md @@ -31,7 +31,7 @@ turn's rule ([[ADR-0002-the-run#^c1|ADR-0002/C1]]) with nothing awaited inline ([[ADR-0005-command-handling#^c2|ADR-0005/C2]]); the notification names what arrived ([[ADR-0006-the-notification#^c2|ADR-0006/C2]]). Source: `Actor.queue`, `Actor._take_inbox`, the handle branch of dispatch, `Actor._deliver` and the `_ACTORS` registry in `lionagi/actor.py`; -`examples/two_actors.py`. +`notebooks/03_two_actors.ipynb`. ## Definitions @@ -136,7 +136,7 @@ Serves C1. `queue` carries content and provenance only; the peer's profile, requ the asking handler's to choose, and it starts the peer's run as its own task. For mail, the long-running agent is the something outside the loop that runs the actor. -- **Landing evidence**: `examples/two_actors.py` (the escalation handler queues, starts the +- **Landing evidence**: `notebooks/03_two_actors.ipynb` (the escalation handler queues, starts the researcher's run and returns the handle); the tests' `_pair` helper does the same, and the round-trip test asserts the peer's INPUT names the asker. @@ -156,8 +156,8 @@ long-running agent is the something outside the loop that runs the actor. it lands. - **S2**: The registry is per process, keyed by name: a second actor made with the same name takes the name, and answers to the first go to it. -- **S3**: No bench exercises this path; the tests drive it with scripted backends, and - `examples/two_actors.py` runs it against a live model. +- **S3**: No bench exercises this path; the tests and `notebooks/03_two_actors.ipynb` drive it with + scripted backends; nothing runs it against a live model. - **S4**: No budget is shared across peers, and a peer run started as its own task outlives an asker that ends: cancellation at the asking run's end reaches the ask's future, never the peer's run. - **S5**: The boundary: who may ask whom, the lineage of an ask, ending a peer with the asker, diff --git a/examples/backends.py b/examples/backends.py index 4a114aa..34c5c65 100644 --- a/examples/backends.py +++ b/examples/backends.py @@ -90,32 +90,44 @@ async def longest(req: Longest, ctx: Context) -> str: return max(req.text.split(), key=len) backend = build_backend(kind, model) - run = await actor.run( - Profile( - "counter", system="Use the commands; point at values; finish with OUT{} when both numbers are in." - ), - RunRequest( - inputs=(f"How many words, and which is the longest? Text: {TEXT!r}",), - emits=(Report,), - max_rounds=6, - time_budget=180, - ), - backend, - ) + try: + run = await actor.run( + Profile( + "counter", + system="Use the commands; point at values; finish with OUT{} when both numbers are in.", + ), + RunRequest( + inputs=(f"How many words, and which is the longest? Text: {TEXT!r}",), + emits=(Report,), + max_rounds=6, + time_budget=180, + ), + backend, + ) + except BaseException: # an interrupt too: the calls that completed before it were paid for + if backend.calls: + spent(backend) + raise print(f"{kind}: {type(run.outcome).__name__} after {run.round} rounds") for e in run.record: if e.kind is Kind.TEXT: print(f" [{e.name}] {e.content.strip()[:160]!r}") if isinstance(run.outcome, Success): print(" output:", run.output["report"].model_dump()) + spent(backend) + +def spent(backend) -> None: + """One line per call the backend made, whatever the run came to.""" print( f"\n{'call':>4} {'ms':>6} {'input':>6} {'cached':>6} {'output':>6} {'think':>6} " f"{'cost $':>9} finish / provider" ) - def cost(v) -> str: - return "unknown" if v is None else f"{float(v):.5f}" + def cost(v, complete) -> str: + # a subtotal with an attempt uncounted (the call was cut off, or an attempt answered without + # a price) is at least what was counted, and says so + return "unknown" if v is None else f"{'' if complete is not False else '>='}{float(v):.5f}" for i, c in enumerate(backend.calls, 1): u = c.get("usage") or {} @@ -123,10 +135,15 @@ def cost(v) -> str: print( f"{i:>4} {ms:>6} {u.get('input_tokens', 0):>6} " f"{u.get('cache_read_input_tokens', 0):>6} {u.get('output_tokens', 0):>6} " - f"{u.get('reasoning_output_tokens') or 0:>6} {cost(c.get('total_cost_usd')):>9} " - f"{c.get('finish_reason')} / {c.get('provider')}" + f"{u.get('reasoning_output_tokens') or 0:>6} " + f"{cost(c.get('total_cost_usd'), c.get('cost_complete')):>9} " + + ( # one line per call: a CLI's error can run to several + c["error"].splitlines()[0] + if c.get("error") + else f"{c.get('finish_reason')} / {c.get('provider')}" + ) ) - print("cost is None on a subscription CLI: the spend is unknown, not zero") + print("cost is None on a subscription CLI: the spend is unknown, not zero; >= marks a subtotal") if __name__ == "__main__": diff --git a/examples/notes_and_context.py b/examples/notes_and_context.py index fb1d11f..9b4ff4b 100644 --- a/examples/notes_and_context.py +++ b/examples/notes_and_context.py @@ -13,6 +13,8 @@ from __future__ import annotations import asyncio +import atexit +import shutil import tempfile from pathlib import Path @@ -21,6 +23,7 @@ from lionagi import Actor, Context, Kind, Profile, RunRequest, fold root = Path(tempfile.mkdtemp()) +atexit.register(shutil.rmtree, root, ignore_errors=True) # the scratch goes with the process (root / "spec.md").write_text("# spec\n\nthe port is 8080 and the workers are 4\n" * 20) diff --git a/examples/tools_and_gates.py b/examples/tools_and_gates.py index b205404..532885c 100644 --- a/examples/tools_and_gates.py +++ b/examples/tools_and_gates.py @@ -13,6 +13,8 @@ from __future__ import annotations import asyncio +import atexit +import shutil import tempfile from pathlib import Path @@ -21,6 +23,7 @@ from lionagi import Actor, Context, Kind, Profile, ProfileError, RunRequest root = Path(tempfile.mkdtemp()) +atexit.register(shutil.rmtree, root, ignore_errors=True) # the scratch goes with the process (root / "a.txt").write_text("alpha\n") diff --git a/notebooks/03_two_actors.ipynb b/notebooks/03_two_actors.ipynb index 56dc8ae..b545f26 100644 --- a/notebooks/03_two_actors.ipynb +++ b/notebooks/03_two_actors.ipynb @@ -101,6 +101,7 @@ " \"OUT{finding: [a1, e1]}\",\n", ")\n", "peer_runs: dict[str, Run] = {}\n", + "peer_tasks: set[asyncio.Task] = set()\n", "\n", "\n", "@reviewer.handler(Escalation, requires={\"escalate\"})\n", @@ -114,7 +115,7 @@ " RESEARCHER, RunRequest(emits=(Finding,), max_rounds=3), researcher_turns\n", " )\n", "\n", - " asyncio.create_task(serve())\n", + " peer_tasks.add(asyncio.create_task(serve())) # held here, awaited below: its error is ours\n", " return handle" ] }, @@ -158,6 +159,7 @@ "metadata": {}, "outputs": [], "source": [ + "await asyncio.gather(*peer_tasks) # a peer that raised raises here, not as a warning at exit\n", "for handle, peer in peer_runs.items():\n", " print(f\"=== researcher run for handle {handle}\")\n", " show(peer)" diff --git a/tests/test_examples.py b/tests/test_examples.py index 4a1b9c3..12eb8db 100644 --- a/tests/test_examples.py +++ b/tests/test_examples.py @@ -1,7 +1,11 @@ """The scripts under examples/ keep working: each one imports, each scripted one runs to its end, and the live one, with no key and no CLI to reach a model, stops with one line naming what is missing.""" +import ast +import asyncio import importlib.util +import inspect +import json import os import subprocess import sys @@ -12,6 +16,7 @@ ROOT = Path(__file__).resolve().parent.parent EXAMPLES = ROOT / "examples" +NOTEBOOKS = ROOT / "notebooks" SCRIPTED = ("tools_and_gates", "notes_and_context", "watch_the_artifact", "budgets_and_outcomes") LIVE = ("backends",) # tests/test_bench_agent.py drives this one: a fake store, and a child that writes a ledger @@ -34,22 +39,170 @@ def test_every_script_in_examples_is_listed_here(): assert sorted(p.stem for p in EXAMPLES.glob("*.py")) == sorted(SCRIPTED + LIVE + OWN_TESTS) -@pytest.mark.parametrize("script", SCRIPTED + LIVE) -def test_a_script_imports_without_running_its_main(script, tmp_path, monkeypatch): - monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) # a script's scratch directory lands here +def load(script: str): spec = importlib.util.spec_from_file_location(f"example_{script}", EXAMPLES / f"{script}.py") assert spec is not None and spec.loader is not None module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) - assert callable(module.main) + return module + + +@pytest.mark.parametrize("script", SCRIPTED + LIVE) +def test_a_script_imports_without_running_its_main(script, tmp_path, monkeypatch): + monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) # a script's scratch directory lands here + assert callable(load(script).main) @pytest.mark.parametrize("script", SCRIPTED) -def test_a_scripted_example_runs_to_its_end(script, tmp_path): +def test_a_scripted_example_runs_to_its_end_and_takes_its_scratch_with_it(script, tmp_path): done = run(script, env={**os.environ, "TMPDIR": str(tmp_path)}) assert done.returncode == 0, done.stderr assert done.stderr == "" assert done.stdout.count("\n") >= 5, done.stdout + assert sorted(tmp_path.iterdir()) == [] # a scratch directory made on import is gone at exit + + +def test_the_live_example_says_what_the_calls_cost_when_a_later_call_fails(tmp_path, monkeypatch, capsys): + """The table of calls is printed on the way out of a failure too: the calls before it were paid.""" + monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) + module = load("backends") + + class Backend: + error = module.OpenRouterError + + def __init__(self): + self.calls = [] + + async def __call__(self, messages): + if self.calls: + raise self.error("429 rate limited") + self.calls.append( + { + "usage": {"input_tokens": 12, "output_tokens": 3}, + "total_cost_usd": 0.00021, + "finish_reason": "stop", + "provider": "p", + "duration_ms": 40, + } + ) + return 'word_count(text="a b")' + + backend = Backend() + monkeypatch.setattr(module, "build_backend", lambda kind, model: backend) + with pytest.raises(module.OpenRouterError, match="429"): + asyncio.run(module.main("openrouter", None)) + out = capsys.readouterr().out + assert " 1 40 12 0 3 0 0.00021 stop / p" in out.splitlines() + assert "openrouter:" not in out # no run to summarize: the failure is the caller's to say + backend = Backend() + backend.error = asyncio.CancelledError # an interrupt after the paid call: the table still prints + with pytest.raises(asyncio.CancelledError): + asyncio.run(module.main("openrouter", None)) + assert ( + " 1 40 12 0 3 0 0.00021 stop / p" in capsys.readouterr().out.splitlines() + ) + + +async def run_cells(name: str, ns: dict, before=None) -> None: + """Every code cell of a notebook, in order, in one namespace, top-level awaits included. `before`, + given a cell's source, runs on the namespace ahead of that cell.""" + for i, cell in enumerate( + c for c in json.loads((NOTEBOOKS / name).read_text())["cells"] if c["cell_type"] == "code" + ): + source = "".join(cell["source"]) + if before is not None: + before(source, ns) + result = eval(compile(source, f"{name}:cell{i}", "exec", ast.PyCF_ALLOW_TOP_LEVEL_AWAIT), ns) # noqa: S307 + if inspect.iscoroutine(result): + await result + + +def test_the_call_table_marks_a_partial_cost_and_names_a_failed_calls_error(tmp_path, monkeypatch, capsys): + """A subtotal an attempt is missing from prints as at least that much, so a partial zero does not + read as a free call; a failed call's row carries its error where a finished one names its end.""" + monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) + module = load("backends") + + class Backend: + calls = [ + { + "usage": {"input_tokens": 1}, + "total_cost_usd": 0.0002, + "cost_complete": True, + "finish_reason": "stop", + "provider": "p", + }, + { + "usage": {"input_tokens": 1}, + "total_cost_usd": 0.0002, + "cost_complete": False, + "finish_reason": "stop", + "provider": "p", + }, + { + "usage": {}, + "total_cost_usd": 0.0, + "cost_complete": False, + "error": "cancelled\nafter 2 attempts", + }, + { + "usage": {}, + "total_cost_usd": 0.0, + "cost_complete": True, + "finish_reason": "stop", + "provider": "p", + }, + {"usage": {}, "total_cost_usd": None, "finish_reason": "stop", "provider": "p"}, + ] + + module.spent(Backend()) + rows = capsys.readouterr().out.splitlines() + assert [r.split(maxsplit=7)[6:] for r in rows if r[:5].strip().isdigit()] == [ + ["0.00020", "stop / p"], + [">=0.00020", "stop / p"], + [">=0.00000", "cancelled"], + ["0.00000", "stop / p"], + ["unknown", "stop / p"], + ] + assert rows[-1].endswith(">= marks a subtotal") + + +@pytest.mark.parametrize("name", sorted(p.name for p in NOTEBOOKS.glob("*.ipynb"))) +def test_a_notebook_runs_to_its_end(name, tmp_path, monkeypatch): + """The notebooks are scripted, so this is the whole notebook and it runs offline.""" + monkeypatch.chdir(tmp_path) + monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) + ns: dict = {"__name__": f"notebook_{name}"} + asyncio.run(run_cells(name, ns)) + if name.startswith("03"): # the peer's task is held and awaited, so its error would be this test's + assert ns["peer_tasks"] and all(t.done() and t.exception() is None for t in ns["peer_tasks"]) + + +def test_the_two_actor_notebook_raises_a_peers_error_on_the_page(tmp_path, monkeypatch): + """A researcher whose backend raises settles the reviewer's ask naming the error, and the + reviewer's run goes on; the peer's own task is awaited where the peer runs are shown, so the + error is raised there rather than left to the loop to report at exit.""" + monkeypatch.chdir(tmp_path) + ns: dict = {"__name__": "notebook_03_failing_peer"} + + async def failing(messages): + raise RuntimeError("the researcher's backend is down") + + seen: list[str] = [] + + def before(source, ns): + if source.startswith("run = await reviewer.run("): + seen.append("run") + ns["researcher_turns"] = failing + if source.startswith("await asyncio.gather"): + seen.append("gather") + assert ns["run"].output["verdict"].decision == "request_changes" # the reviewer finished + assert "RuntimeError" in str(ns["run"].resolve("q1")) # and its ask says why the peer failed + + with pytest.raises(RuntimeError, match="backend is down"): + asyncio.run(run_cells("03_two_actors.ipynb", ns, before)) + assert seen == ["run", "gather"] # both cells were met, so the checks inside ran + assert [t.exception().args for t in ns["peer_tasks"]] == [("the researcher's backend is down",)] @pytest.mark.parametrize(