Skip to content

Commit b893a72

Browse files
committed
feat: support orchestrated agent sessions
1 parent 54482ff commit b893a72

13 files changed

Lines changed: 357 additions & 29 deletions

File tree

‎README.md‎

Lines changed: 31 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,40 @@
11
# Agent Ledger
22

33
Agent Ledger is a framework-neutral specification and a set of polyglot adapters for durable agent
4-
sessions. It records model and tool attempts before execution, preserves causal timelines across
5-
distributed agents, and lets each framework rebuild its own native session after interruption.
4+
execution records. Agent loops and orchestrators append to the same session history, producing a
5+
causal account of model calls, tool calls, delegation, framework-native state, and outcomes.
66

77
The specification is the stable product. Language SDKs are deliberately small; most project code
88
lives in adapters that understand a framework's hooks, messages, checkpoints, and resume APIs.
99

10+
## Architecture position
11+
12+
Agent Ledger is not another agent loop or workflow engine. It is a shared evidence layer across
13+
both:
14+
15+
```text
16+
Orchestrator ── decisions, delegation, approvals ──┐
17+
├── Agent Ledger session
18+
Agent loops ── steps, attempts, native state ──────┘ │
19+
├── recovery
20+
├── global timeline
21+
└── analysis / evaluation
22+
```
23+
24+
An orchestrator owns desired state, scheduling, and run ownership. Each agent framework owns its
25+
native context and resume API. Agent Ledger owns the immutable facts that let those systems explain
26+
and reconstruct what happened.
27+
1028
## Model
1129

12-
- `Session` groups one end-to-end task across processes, languages, and agents.
13-
- `Run` identifies one semantic agent execution and participates in the causal DAG.
30+
- `Session` groups one end-to-end task across processes, languages, agents, and orchestration runs.
31+
- `Run` identifies one semantic execution by an agent or orchestrator and participates in the
32+
causal DAG.
1433
- `EventStream` is an optimistic-concurrency partition. It may contain one run's execution events or
1534
framework-native state that survives several runtime runs.
1635
- `Step` is logical work that survives retries; `Attempt` is one physical model or tool invocation.
1736
- Normalized events are the source for timelines and trajectories. Framework-native records are the
18-
source for resume.
37+
lossless input to framework-owned resume.
1938

2039
Requested events are committed before an external call. A requested event without a terminal event
2140
is unresolved after a crash. It is input to the adapter's reconciliation policy; an adapter must
@@ -31,7 +50,7 @@ not silently replay a side-effecting tool.
3150
| `typescript/` | TypeScript core SDK and Pi adapter |
3251
| `go/` | Go core SDK and AgentGo adapter |
3352

34-
Current framework profiles:
53+
Current framework profiles are integration examples, not definitions of the core session model:
3554

3655
| Adapter | Recording | Recovery |
3756
| --- | --- | --- |
@@ -42,6 +61,10 @@ Current framework profiles:
4261
Every adapter publishes machine-readable capabilities such as `strict`, `best_effort`, and
4362
`unsupported`; installing a telemetry-only hook never silently claims durable recovery.
4463

64+
Pi's append-only session tree is preserved in a dedicated framework stream because Pi needs it for
65+
lossless reconstruction. Its entry types, active leaf, and branching rules remain Pi-owned rather
66+
than becoming requirements for other agents or orchestrators.
67+
4568
## Store contract
4669

4770
Applications inject an `EventStore`. V1 has no mandatory collector or `/agent-session` service:
@@ -69,3 +92,5 @@ make build
6992

7093
See [RFC 0001](spec/rfcs/0001-agent-ledger.md) for the ledger contract and
7194
[RFC 0002](spec/rfcs/0002-polyglot-adapters.md) for framework recording and recovery boundaries.
95+
The [orchestrated agents example](python/examples/orchestrated_agents.py) shows an orchestrator and
96+
multiple agent loops contributing to one causal session.

‎go/recorder.go‎

Lines changed: 45 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ type RecorderOptions struct {
1313
RunID string
1414
StreamID string
1515
Actor Actor
16+
Parent *CausalParent
1617
ExpectedVersion *int64
1718
}
1819

@@ -21,6 +22,7 @@ type SessionRecorder struct {
2122
stream EventStream
2223
runID string
2324
actor Actor
25+
parent *CausalParent
2426
expectedVersion int64
2527
mu sync.Mutex
2628
}
@@ -34,11 +36,17 @@ func NewSessionRecorder(options RecorderOptions) *SessionRecorder {
3436
if options.ExpectedVersion != nil {
3537
expectedVersion = *options.ExpectedVersion
3638
}
39+
var parent *CausalParent
40+
if options.Parent != nil {
41+
copy := *options.Parent
42+
parent = &copy
43+
}
3744
return &SessionRecorder{
3845
store: options.Store,
3946
stream: EventStream{SessionID: options.SessionID, StreamID: streamID},
4047
runID: options.RunID,
4148
actor: options.Actor,
49+
parent: parent,
4250
expectedVersion: expectedVersion,
4351
}
4452
}
@@ -64,12 +72,39 @@ func (r *SessionRecorder) RunID() string { return r.runID }
6472
func (r *SessionRecorder) Store() EventStore { return r.store }
6573

6674
func (r *SessionRecorder) Record(ctx context.Context, eventType string, payload map[string]any, stepID, attemptID string) (StoredEvent, error) {
67-
r.mu.Lock()
68-
defer r.mu.Unlock()
6975
event := NewEvent(eventType, r.stream.SessionID, r.runID, r.actor)
70-
event.Payload = payload
76+
event.Payload = payloadOrEmpty(payload)
7177
event.StepID = stepID
7278
event.AttemptID = attemptID
79+
return r.appendEvent(ctx, event)
80+
}
81+
82+
func (r *SessionRecorder) StartRun(ctx context.Context, payload map[string]any) (StoredEvent, error) {
83+
event := NewEvent("run.started", r.stream.SessionID, r.runID, r.actor)
84+
event.Payload = payloadOrEmpty(payload)
85+
if r.parent != nil {
86+
event.ParentRunID = r.parent.RunID
87+
event.CausedByEventID = r.parent.CausedByEventID
88+
}
89+
return r.appendEvent(ctx, event)
90+
}
91+
92+
func (r *SessionRecorder) Child(runID string, actor Actor, causedByEventID string) *SessionRecorder {
93+
return NewSessionRecorder(RecorderOptions{
94+
Store: r.store,
95+
SessionID: r.stream.SessionID,
96+
RunID: runID,
97+
Actor: actor,
98+
Parent: &CausalParent{
99+
RunID: r.runID,
100+
CausedByEventID: causedByEventID,
101+
},
102+
})
103+
}
104+
105+
func (r *SessionRecorder) appendEvent(ctx context.Context, event ProposedEvent) (StoredEvent, error) {
106+
r.mu.Lock()
107+
defer r.mu.Unlock()
73108
receipt, err := r.store.Append(ctx, r.stream, r.expectedVersion, NewID(), event)
74109
if err != nil {
75110
return StoredEvent{}, err
@@ -127,3 +162,10 @@ func errorPayload(err error) map[string]any {
127162
}
128163
return map[string]any{"error": err.Error()}
129164
}
165+
166+
func payloadOrEmpty(payload map[string]any) map[string]any {
167+
if payload == nil {
168+
return map[string]any{}
169+
}
170+
return payload
171+
}

‎go/store_test.go‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,3 +96,44 @@ func TestResumeRecorderRejectsExpectedVersion(t *testing.T) {
9696
t.Fatal("resume accepted an explicit expected version")
9797
}
9898
}
99+
100+
func TestOrchestratorLinksMultipleAgentRuns(t *testing.T) {
101+
ctx := context.Background()
102+
store := NewMemoryEventStore()
103+
orchestrator := NewSessionRecorder(RecorderOptions{
104+
Store: store, SessionID: "session", RunID: "orchestrator-run",
105+
Actor: Actor{Type: "orchestrator", ID: "planner"},
106+
})
107+
if _, err := orchestrator.StartRun(ctx, nil); err != nil {
108+
t.Fatalf("start orchestrator: %v", err)
109+
}
110+
for _, role := range []string{"researcher", "reviewer"} {
111+
dispatch, err := orchestrator.Record(
112+
ctx, "orchestration.agent.dispatched", map[string]any{"role": role}, "", "",
113+
)
114+
if err != nil {
115+
t.Fatalf("record %s dispatch: %v", role, err)
116+
}
117+
child := orchestrator.Child(role+"-run", Actor{Type: "agent", ID: role}, dispatch.EventID)
118+
if _, err := child.StartRun(ctx, nil); err != nil {
119+
t.Fatalf("start %s: %v", role, err)
120+
}
121+
}
122+
123+
var childRuns []string
124+
for event, err := range store.ScanSession(ctx, "session", "") {
125+
if err != nil {
126+
t.Fatalf("scan session: %v", err)
127+
}
128+
if event.ParentRunID == "" {
129+
continue
130+
}
131+
if event.ParentRunID != "orchestrator-run" || event.CausedByEventID == "" {
132+
t.Fatalf("invalid causal edge: %#v", event.ProposedEvent)
133+
}
134+
childRuns = append(childRuns, event.RunID)
135+
}
136+
if len(childRuns) != 2 || childRuns[0] != "researcher-run" || childRuns[1] != "reviewer-run" {
137+
t.Fatalf("child runs = %v", childRuns)
138+
}
139+
}

‎go/types.go‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,11 @@ type EventStream struct {
1313
StreamID string `json:"stream_id"`
1414
}
1515

16+
type CausalParent struct {
17+
RunID string
18+
CausedByEventID string
19+
}
20+
1621
type ProposedEvent struct {
1722
SchemaVersion string `json:"schema_version"`
1823
EventID string `json:"event_id"`
Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
import asyncio
2+
from uuid import uuid4
3+
4+
from agent_ledger import Actor, SessionRecorder, inspect_session
5+
from agent_ledger.stores.memory import MemoryEventStore
6+
7+
8+
async def main() -> None:
9+
store = MemoryEventStore()
10+
session_id = str(uuid4())
11+
orchestrator = SessionRecorder(
12+
store=store,
13+
session_id=session_id,
14+
run_id="orchestrator-1",
15+
actor=Actor(type="orchestrator", id="planner"),
16+
)
17+
await orchestrator.start_run(payload={"goal": "research and review a proposal"})
18+
19+
for role in ("researcher", "reviewer"):
20+
dispatch = await orchestrator.record(
21+
"orchestration.agent.dispatched",
22+
payload={"role": role},
23+
)
24+
agent = orchestrator.child(
25+
run_id=f"{role}-1",
26+
actor=Actor(type="agent", id=role, framework="plain-loop"),
27+
caused_by_event_id=dispatch.event_id,
28+
)
29+
await agent.start_run(
30+
payload={
31+
"agent": {"id": role, "version": "1"},
32+
"framework": {"name": "plain-loop"},
33+
"code": {"revision": "example"},
34+
}
35+
)
36+
step_id = str(uuid4())
37+
await agent.start_step(step_id, payload={"role": role})
38+
await agent.complete_step(step_id)
39+
await agent.complete_run()
40+
41+
events = [event async for event in store.scan_session(session_id)]
42+
inspection = inspect_session(events)
43+
print(f"events={len(inspection.timeline)} run_edges={len(inspection.run_edges)}")
44+
45+
46+
if __name__ == "__main__":
47+
asyncio.run(main())

‎python/tests/test_models.py‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,9 @@
33
import json
44
from pathlib import Path
55

6+
import pytest
67
from jsonschema import Draft202012Validator, FormatChecker
8+
from jsonschema.exceptions import ValidationError
79

810
from agent_ledger import Actor, MemoryArtifactStore, ProposedEvent, StoredEvent
911
from agent_ledger.frameworks.plain_loop import PlainLoopProfile
@@ -32,6 +34,25 @@ def test_event_matches_normative_json_schema() -> None:
3234
validator.validate(stored.model_dump(mode="json", exclude_none=True))
3335

3436

37+
def test_event_schema_requires_complete_causal_parent() -> None:
38+
schema_path = Path(__file__).parents[2] / "spec" / "schemas" / "event.schema.json"
39+
schema = json.loads(schema_path.read_text())
40+
event = {
41+
"schema_version": "1.0",
42+
"event_id": "event",
43+
"event_type": "run.started",
44+
"session_id": "session",
45+
"run_id": "run",
46+
"actor": {"type": "agent", "id": "child"},
47+
"occurred_at": "2026-08-09T00:00:00Z",
48+
"parent_run_id": "parent",
49+
"payload": {},
50+
}
51+
52+
with pytest.raises(ValidationError):
53+
Draft202012Validator(schema).validate(event)
54+
55+
3556
async def test_memory_artifact_round_trip() -> None:
3657
store = MemoryArtifactStore()
3758
ref = await store.put("session", b"large model output", "text/plain")

‎python/tests/test_recorder.py‎

Lines changed: 25 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -54,24 +54,37 @@ async def test_retry_keeps_step_and_gets_new_attempt() -> None:
5454
assert inspection.unresolved_attempts[0].step_id == "step-1"
5555

5656

57-
async def test_child_run_forms_causal_edge() -> None:
57+
async def test_orchestrator_links_multiple_agent_runs() -> None:
5858
store = MemoryEventStore()
59-
parent = _recorder(store)
60-
trigger = await parent.start_step("delegate")
61-
child = parent.child(
62-
run_id=str(uuid4()),
63-
actor=Actor(type="agent", id="child"),
64-
caused_by_event_id=trigger.event_id,
59+
parent = SessionRecorder(
60+
store=store,
61+
session_id=str(uuid4()),
62+
run_id="orchestrator-run",
63+
actor=Actor(type="orchestrator", id="planner"),
6564
)
66-
await child.start_run()
65+
await parent.start_run()
66+
children: list[SessionRecorder] = []
67+
for role in ("researcher", "reviewer"):
68+
trigger = await parent.record(
69+
"orchestration.agent.dispatched",
70+
payload={"role": role},
71+
)
72+
child = parent.child(
73+
run_id=f"{role}-run",
74+
actor=Actor(type="agent", id=role),
75+
caused_by_event_id=trigger.event_id,
76+
)
77+
await child.start_run()
78+
children.append(child)
6779

6880
events = [event async for event in store.scan_session(parent.stream.session_id)]
6981
inspection = inspect_session(events)
7082

71-
assert len(inspection.run_edges) == 1
72-
assert inspection.run_edges[0].parent_run_id == parent.run_id
73-
assert inspection.run_edges[0].child_run_id == child.run_id
74-
assert inspection.run_edges[0].caused_by_event_id == trigger.event_id
83+
assert len(inspection.run_edges) == 2
84+
assert {edge.parent_run_id for edge in inspection.run_edges} == {parent.run_id}
85+
assert {edge.child_run_id for edge in inspection.run_edges} == {
86+
child.run_id for child in children
87+
}
7588

7689

7790
async def test_plain_loop_profile_restores_snapshot_and_tail() -> None:

0 commit comments

Comments
 (0)