Skip to content

Commit d1c0d02

Browse files
authored
feat(runtime): publish model retry lifecycle (#80)
1 parent 52fe762 commit d1c0d02

21 files changed

Lines changed: 1273 additions & 60 deletions

docs/architecture/domain-model.md

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -54,12 +54,14 @@ The public session status is a projection:
5454

5555
- `idle` — waiting for input or a required client action;
5656
- `running` — accepted work is executing or queued;
57-
- `rescheduling` — reserved for automatic retry behavior;
57+
- `rescheduling` — an active model request is waiting for its next bounded
58+
provider retry;
5859
- `terminated` — terminal failure or completion.
5960

60-
Temporal may retry failed Activities internally without changing the public
61-
status. The public `rescheduling` transition is not implemented; an exhausted
62-
or terminal turn failure projects directly to `terminated`.
61+
Temporal may retry infrastructure-failed Activities internally without changing
62+
the public status. Retryable provider responses instead use the public
63+
`running -> rescheduling -> running` lifecycle. Exhausting that bounded budget
64+
returns the Session to `idle`; permanent failures project to `terminated`.
6365

6466
## Event
6567

docs/architecture/session-lifecycle.md

Lines changed: 16 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,9 @@ stateDiagram-v2
1515
[*] --> idle
1616
idle --> running: claimable input admitted
1717
running --> idle: turn completed or requires action
18+
running --> rescheduling: retryable model failure
19+
rescheduling --> running: retry delay elapsed
20+
rescheduling --> idle: interrupted
1821
running --> terminated: terminal failure
1922
idle --> terminated: deleted or terminal operation
2023
terminated --> [*]
@@ -23,8 +26,8 @@ stateDiagram-v2
2326
- `idle` means the Session is waiting for ordinary input or a required client
2427
action.
2528
- `running` means accepted, claimable work is queued or executing.
26-
- `rescheduling` is reserved for a future public retry policy. Temporal
27-
Activity retries remain private and do not use it.
29+
- `rescheduling` means an immutable model request is waiting for its next
30+
bounded provider retry. Later input remains queued behind the active turn.
2831
- `terminated` is final.
2932

3033
The status is a projection of committed events. It is not derived from worker
@@ -129,16 +132,21 @@ not supported.
129132

130133
## Failure and delivery semantics
131134

132-
- Permanent model failures commit `session.error` and
133-
`session.status_terminated`.
134-
- Retryable model and infrastructure failures remain Activity failures so
135-
Temporal can apply its retry policy.
135+
- Permanent model failures commit `session.error` with `retry_status: terminal`
136+
and `session.status_terminated`.
137+
- Retryable provider responses publish `session.error` with
138+
`retry_status: retrying`, atomically move through
139+
`session.status_rescheduled`, and publish `session.status_running` before the
140+
next attempt. Exhaustion returns the Session to idle with
141+
`stop_reason: retries_exhausted` and flushes later queued messages.
142+
- Infrastructure failures remain Activity failures so Temporal can recover
143+
without spending the public provider retry budget.
136144
- Authoritative output is visible only after the PostgreSQL completion
137145
transaction commits.
138146
- NATS previews and wakeups are best-effort. Clients recover complete events
139147
from the PostgreSQL-backed event cursor.
140-
- Workflow and Activity retries are private implementation details and never
141-
create duplicate public events.
148+
- Activity retries for infrastructure recovery remain private. Workflow-owned
149+
provider retries use deterministic event IDs and never create duplicates.
142150

143151
## Session deletion
144152

docs/compatibility.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ compatibility.
4646
| Tool confirmations | Supported | Interceptable `always_ask` built-ins and MCP tools park durably. Allow executes the original server-owned call through the tool journal; deny returns an error tool result with the optional denial message. Provider-native Web Search/Fetch cannot be intercepted and reject `always_ask`. |
4747
| User interrupt | Limited | An untargeted interrupt durably cancels an active model, outcome grader, or tool Activity across API and worker processes. PostgreSQL defines finish-vs-interrupt ordering, closes model spans, emits one idle `end_turn`, and fences uncertain started tool steps as ambiguous. Targeted multi-agent interrupts are not supported. |
4848
| Terminal execution errors | Supported | A deterministic turn failure emits `session.error` with the documented `unknown_error` variant and `retry_status: terminal`, followed by `session.status_terminated`. MCP setup failures retain their more specific documented event variant. Older persisted histories may contain the legacy `api_error` spelling. |
49+
| Model retry lifecycle | Supported | Retryable provider responses emit typed `model_overloaded_error`, `model_rate_limited_error`, or `model_request_failed_error` events. Each bounded retry publishes `retrying`, `session.status_rescheduled`, then `session.status_running`; exhaustion emits `exhausted`, returns idle with `retries_exhausted`, and flushes later queued messages. Provider `Retry-After` is honored up to the server cap. Infrastructure recovery remains private to Temporal. |
4950
| MCP execution | Limited | Remote tools are discovered over unauthenticated Streamable HTTP, pinned per Session, permission-checked, journaled, and executed with large or binary results materialized in the Session sandbox. Calls publish `agent.mcp_tool_use` with the required `mcp_server_name` and the bare server-side tool name, answered by `agent.mcp_tool_result` carrying `mcp_tool_use_id`; an `always_ask` MCP call still parks for a `user.tool_confirmation` referencing `tool_use_id`. Authentication, private-network access, deprecated-SSE fallback, resources, and prompts are not supported. |
5051
| Files, skills, memory, and vaults | Not supported | These product surfaces are outside the core harness scope. File-backed outcome rubrics and file-sourced message content are rejected explicitly; inline or URL image/document content still participates in the provider transcript and context projection. |
5152
| Multi-agent orchestration | Not supported | `multiagent` configuration can be stored, but rosters, threads, delegation, and orchestration are not executed. |

internal/domain/session.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ const (
1717
var allowed = map[Status]map[Status]bool{
1818
StatusIdle: {StatusRunning: true, StatusTerminated: true},
1919
StatusRunning: {StatusIdle: true, StatusRescheduling: true, StatusTerminated: true},
20-
StatusRescheduling: {StatusRunning: true, StatusTerminated: true},
20+
StatusRescheduling: {StatusIdle: true, StatusRunning: true, StatusTerminated: true},
2121
StatusTerminated: {},
2222
}
2323

internal/domain/session_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ func TestStatusTransitions(t *testing.T) {
99
ok := [][2]Status{
1010
{StatusIdle, StatusRunning}, {StatusRunning, StatusIdle},
1111
{StatusRunning, StatusRescheduling}, {StatusRescheduling, StatusRunning},
12+
{StatusRescheduling, StatusIdle},
1213
{StatusIdle, StatusTerminated}, {StatusRunning, StatusTerminated},
1314
}
1415
for _, p := range ok {

internal/httpapi/openapi.yaml

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2195,7 +2195,16 @@ components:
21952195
type:
21962196
type: string
21972197
description: api_error may appear only in histories written before the terminal Session error contract.
2198-
enum: [unknown_error, mcp_connection_failed_error, api_error]
2198+
enum:
2199+
- unknown_error
2200+
- model_overloaded_error
2201+
- model_rate_limited_error
2202+
- model_request_failed_error
2203+
- mcp_connection_failed_error
2204+
- mcp_authentication_failed_error
2205+
- billing_error
2206+
- credential_host_unreachable_error
2207+
- api_error
21992208
message: {type: string}
22002209
retry_status:
22012210
type: object

internal/httpapi/sdk_test.go

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,66 @@ func TestSDK_TerminalSessionErrorEvent(t *testing.T) {
8888
}
8989
}
9090

91+
func TestSDK_ModelRetryLifecycleEvents(t *testing.T) {
92+
client, ts, sessions := sdkClientServerAndSessions(t)
93+
ctx := context.Background()
94+
95+
agent := mustAgent(t, client, "opus", "sys")
96+
environmentID := mustEnv(t, ts.URL)
97+
session, err := client.Beta.Sessions.New(ctx, anthropic.BetaSessionNewParams{
98+
Agent: anthropic.BetaSessionNewParamsAgentUnion{OfString: anthropic.String(agent.ID)},
99+
EnvironmentID: environmentID,
100+
})
101+
if err != nil {
102+
t.Fatalf("create session: %v", err)
103+
}
104+
105+
sessions.mu.Lock()
106+
sessions.appendEventLocked(session.ID, domain.EventDraft{
107+
Type: domain.EvSessionError,
108+
Payload: map[string]any{"error": map[string]any{
109+
"type": "model_rate_limited_error", "message": "slow down",
110+
"retry_status": map[string]any{"type": "retrying"},
111+
}},
112+
})
113+
sessions.appendEventLocked(session.ID, domain.EventDraft{
114+
Type: domain.EvSessionStatusRescheduling, Payload: map[string]any{},
115+
})
116+
sessions.appendEventLocked(session.ID, domain.EventDraft{
117+
Type: domain.EvSessionStatusRunning, Payload: map[string]any{},
118+
})
119+
sessions.mu.Unlock()
120+
121+
page, err := client.Beta.Sessions.Events.List(
122+
ctx,
123+
session.ID,
124+
anthropic.BetaSessionEventListParams{Types: []string{
125+
domain.EvSessionError,
126+
domain.EvSessionStatusRescheduling,
127+
domain.EvSessionStatusRunning,
128+
}},
129+
)
130+
if err != nil {
131+
t.Fatalf("list retry lifecycle: %v", err)
132+
}
133+
if len(page.Data) != 3 {
134+
t.Fatalf("retry lifecycle event count = %d, want 3", len(page.Data))
135+
}
136+
retry := page.Data[0].AsSessionError().Error.AsModelRateLimitedError()
137+
if retry.Type != anthropic.BetaManagedAgentsModelRateLimitedErrorTypeModelRateLimitedError ||
138+
retry.Message != "slow down" || retry.RetryStatus.AsRetrying().Type != "retrying" {
139+
t.Fatalf("retry event = %s", page.Data[0].RawJSON())
140+
}
141+
if page.Data[1].AsSessionStatusRescheduled().Type !=
142+
anthropic.BetaManagedAgentsSessionStatusRescheduledEventTypeSessionStatusRescheduled {
143+
t.Fatalf("rescheduled event = %s", page.Data[1].RawJSON())
144+
}
145+
if page.Data[2].AsSessionStatusRunning().Type !=
146+
anthropic.BetaManagedAgentsSessionStatusRunningEventTypeSessionStatusRunning {
147+
t.Fatalf("running event = %s", page.Data[2].RawJSON())
148+
}
149+
}
150+
91151
func TestSDK_AgentLifecycle(t *testing.T) {
92152
client, _ := sdkClientAndServer(t)
93153
ctx := context.Background()

internal/pg/completion_test.go

Lines changed: 217 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -171,3 +171,220 @@ func TestAppendWorkflowEvents_IdempotentAndIncludedInCompletionReplay(t *testing
171171
t.Fatal("mismatched idempotency payload unexpectedly succeeded")
172172
}
173173
}
174+
175+
func TestWorkflowRetryTransitionsAreAtomicIdempotentAndKeepQueuedInputBehindTurn(t *testing.T) {
176+
store := testStore(t)
177+
ctx := context.Background()
178+
session := newSession("sess_retry_transitions")
179+
if _, err := store.CreateSession(ctx, session, nil); err != nil {
180+
t.Fatalf("create: %v", err)
181+
}
182+
admission, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{{
183+
Type: domain.EvUserMessage, Payload: map[string]any{"content": "first"},
184+
}})
185+
if err != nil {
186+
t.Fatalf("admit trigger: %v", err)
187+
}
188+
trigger := admission.Events[0]
189+
errorPayload := map[string]any{
190+
"type": "model_overloaded_error",
191+
"message": "provider overloaded",
192+
"retry_status": map[string]any{
193+
"type": "retrying",
194+
},
195+
}
196+
for attempt := 0; attempt < 2; attempt++ {
197+
if err := store.RecordWorkflowRetry(
198+
ctx,
199+
session.ID,
200+
trigger.ID,
201+
"sevt_retry_error",
202+
"sevt_retry_rescheduled",
203+
errorPayload,
204+
); err != nil {
205+
t.Fatalf("record retry %d: %v", attempt+1, err)
206+
}
207+
}
208+
projected, err := store.GetSession(ctx, session.ID)
209+
if err != nil {
210+
t.Fatalf("get rescheduling session: %v", err)
211+
}
212+
if projected.Status != domain.StatusRescheduling {
213+
t.Fatalf("status = %s, want rescheduling", projected.Status)
214+
}
215+
216+
queued, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{{
217+
Type: domain.EvUserMessage, Payload: map[string]any{"content": "queued"},
218+
}})
219+
if err != nil {
220+
t.Fatalf("queue later message: %v", err)
221+
}
222+
if queued.Session.Status != domain.StatusRescheduling {
223+
t.Fatalf("queued admission status = %s, want rescheduling", queued.Session.Status)
224+
}
225+
if len(queued.Events) != 1 || queued.Events[0].Type != domain.EvUserMessage {
226+
t.Fatalf("queued admission events = %v", eventTypes(queued.Events))
227+
}
228+
229+
for attempt := 0; attempt < 2; attempt++ {
230+
if err := store.ResumeWorkflowRetry(
231+
ctx, session.ID, trigger.ID, "sevt_retry_running",
232+
); err != nil {
233+
t.Fatalf("resume retry %d: %v", attempt+1, err)
234+
}
235+
}
236+
projected, err = store.GetSession(ctx, session.ID)
237+
if err != nil {
238+
t.Fatalf("get resumed session: %v", err)
239+
}
240+
if projected.Status != domain.StatusRunning {
241+
t.Fatalf("status = %s, want running", projected.Status)
242+
}
243+
events, err := store.EventsAfter(ctx, session.ID, 0, 100)
244+
if err != nil {
245+
t.Fatalf("list events: %v", err)
246+
}
247+
if got := eventTypes(events); !slices.Equal(got, []string{
248+
domain.EvUserMessage,
249+
domain.EvSessionStatusRunning,
250+
domain.EvSessionError,
251+
domain.EvSessionStatusRescheduling,
252+
domain.EvUserMessage,
253+
domain.EvSessionStatusRunning,
254+
}) {
255+
t.Fatalf("event order = %v", got)
256+
}
257+
}
258+
259+
func TestRetriesExhaustedReturnsIdleAndFlushesQueuedMessages(t *testing.T) {
260+
store := testStore(t)
261+
ctx := context.Background()
262+
session := newSession("sess_retry_flush")
263+
if _, err := store.CreateSession(ctx, session, nil); err != nil {
264+
t.Fatalf("create: %v", err)
265+
}
266+
first, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{{
267+
Type: domain.EvUserMessage, Payload: map[string]any{"content": "first"},
268+
}})
269+
if err != nil {
270+
t.Fatalf("admit first: %v", err)
271+
}
272+
queued, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{
273+
{Type: domain.EvUserMessage, Payload: map[string]any{"content": "second"}},
274+
{Type: domain.EvSystemMessage, Payload: map[string]any{"content": "context"}},
275+
})
276+
if err != nil {
277+
t.Fatalf("admit queued pair: %v", err)
278+
}
279+
completion, err := store.CompleteWorkflowTurn(
280+
ctx,
281+
session.ID,
282+
first.Events[0].ID,
283+
[]domain.EventDraft{
284+
{Type: domain.EvSessionError, Payload: map[string]any{"error": map[string]any{
285+
"type": "model_request_failed_error",
286+
"message": "retry budget exhausted",
287+
"retry_status": map[string]any{
288+
"type": "exhausted",
289+
},
290+
}}},
291+
{Type: domain.EvSessionStatusIdle, Payload: map[string]any{
292+
"stop_reason": map[string]any{"type": "retries_exhausted"},
293+
}},
294+
},
295+
domain.StatusIdle,
296+
"",
297+
"",
298+
nil,
299+
nil,
300+
nil,
301+
)
302+
if err != nil {
303+
t.Fatalf("complete exhausted turn: %v", err)
304+
}
305+
if completion.Session.Status != domain.StatusIdle {
306+
t.Fatalf("status = %s, want idle", completion.Session.Status)
307+
}
308+
queuedMessage, err := store.GetEvent(ctx, session.ID, queued.Events[0].ID)
309+
if err != nil {
310+
t.Fatalf("get queued message: %v", err)
311+
}
312+
if queuedMessage.ProcessedAt == nil {
313+
t.Fatal("queued message was not flushed")
314+
}
315+
queuedSystem, err := store.GetEvent(ctx, session.ID, queued.Events[1].ID)
316+
if err != nil {
317+
t.Fatalf("get queued system message: %v", err)
318+
}
319+
if queuedSystem.ProcessedAt == nil {
320+
t.Fatal("queued companion system message was not flushed")
321+
}
322+
}
323+
324+
func TestInterruptCanReturnAReschedulingSessionToIdle(t *testing.T) {
325+
store := testStore(t)
326+
ctx := context.Background()
327+
session := newSession("sess_retry_interrupt")
328+
if _, err := store.CreateSession(ctx, session, nil); err != nil {
329+
t.Fatalf("create: %v", err)
330+
}
331+
admission, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{{
332+
Type: domain.EvUserMessage, Payload: map[string]any{"content": "retry me"},
333+
}})
334+
if err != nil {
335+
t.Fatalf("admit trigger: %v", err)
336+
}
337+
trigger := admission.Events[0]
338+
if err := store.RecordWorkflowRetry(
339+
ctx,
340+
session.ID,
341+
trigger.ID,
342+
"sevt_interrupt_retry_error",
343+
"sevt_interrupt_rescheduled",
344+
map[string]any{
345+
"type": "model_request_failed_error",
346+
"message": "temporary failure",
347+
"retry_status": map[string]any{
348+
"type": "retrying",
349+
},
350+
},
351+
); err != nil {
352+
t.Fatalf("record retry: %v", err)
353+
}
354+
interrupt, err := store.AdmitEvents(ctx, session.ID, []domain.EventDraft{{
355+
Type: domain.EvUserInterrupt, Payload: map[string]any{},
356+
}})
357+
if err != nil {
358+
t.Fatalf("admit interrupt: %v", err)
359+
}
360+
completion, err := store.CompleteWorkflowTurn(
361+
ctx,
362+
session.ID,
363+
trigger.ID,
364+
[]domain.EventDraft{{
365+
Type: domain.EvSessionStatusIdle,
366+
Payload: map[string]any{
367+
"stop_reason": map[string]any{"type": "end_turn"},
368+
},
369+
}},
370+
domain.StatusIdle,
371+
"",
372+
"",
373+
nil,
374+
nil,
375+
nil,
376+
)
377+
if err != nil {
378+
t.Fatalf("complete interrupted retry: %v", err)
379+
}
380+
if completion.Session.Status != domain.StatusIdle {
381+
t.Fatalf("status = %s, want idle", completion.Session.Status)
382+
}
383+
storedInterrupt, err := store.GetEvent(ctx, session.ID, interrupt.Events[0].ID)
384+
if err != nil {
385+
t.Fatalf("get interrupt: %v", err)
386+
}
387+
if storedInterrupt.ProcessedAt == nil {
388+
t.Fatal("interrupt was not acknowledged")
389+
}
390+
}

0 commit comments

Comments
 (0)