Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions docs/adr/010-delegation/ADR-0010-delegation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.

Expand All @@ -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,
Expand Down
51 changes: 34 additions & 17 deletions examples/backends.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,43 +90,60 @@ 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 {}
ms = c.get("duration_api_ms") or c.get("duration_ms") or 0
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__":
Expand Down
3 changes: 3 additions & 0 deletions examples/notes_and_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
from __future__ import annotations

import asyncio
import atexit
import shutil
import tempfile
from pathlib import Path

Expand All @@ -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)


Expand Down
3 changes: 3 additions & 0 deletions examples/tools_and_gates.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
from __future__ import annotations

import asyncio
import atexit
import shutil
import tempfile
from pathlib import Path

Expand All @@ -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")


Expand Down
4 changes: 3 additions & 1 deletion notebooks/03_two_actors.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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"
]
},
Expand Down Expand Up @@ -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)"
Expand Down
163 changes: 158 additions & 5 deletions tests/test_examples.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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
Expand All @@ -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 '<lact c1>word_count(text="a b")</lact>'

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(
Expand Down
Loading