返回 DeepSeek-Reasonix
recovery.go
根目录 / internal / session / recovery.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "sort"
8
9 "reasonix/internal/event"
10 "reasonix/internal/provider"
11 )
12
13 // RecoverInterrupted closes persisted runtime authority that cannot survive a
14 // process restart. It never reruns a tool or restores an approval. Callers must
15 // hold the exclusive write handle returned by Open.
16 //
17 // This is a Session operation because it derives a closure batch from the
18 // projection; the physical handle only writes the resulting commit.
19 func (s *Session) RecoverInterrupted(ctx context.Context) (Commit, bool, error) {
20 if s == nil {
21 return Commit{}, false, fmt.Errorf("session: nil session")
22 }
23 if s.readOnly {
24 return Commit{}, false, ErrReadOnly
25 }
26 // Restart recovery needs only live authority, never historical message
27 // bodies or the provider workset. Using the full compatibility Snapshot here
28 // would replay the entire durable transcript on every cold open.
29 snapshot := s.StateSnapshot()
30 turnID := snapshot.Projection.TurnID
31 if turnID == "" {
32 return Commit{}, false, nil
33 }
34 events := streamCheckpointRecoveryEvents(snapshot.Projection)
35 events = append(events, ClosureEvents(snapshot.Projection, "unavailable", "previous runtime exited before recording a result")...)
36 terminal, _ := json.Marshal(map[string]any{"status": event.TurnInterrupted})
37 events = append(events, Event{Kind: "turn/end", Payload: terminal})
38 commit, err := s.Append(ctx, Batch{OperationID: "turn-finalize:" + turnID, TurnID: turnID, Events: events})
39 if err != nil {
40 return Commit{}, false, err
41 }
42 if _, err := s.Flush(ctx); err != nil {
43 return Commit{}, false, err
44 }
45 return commit, true, nil
46 }
47
48 // ClosureEvents derives deterministic terminal facts without executing tools
49 // or changing the provider workset. Live termination and restart share it.
50 func ClosureEvents(projection Projection, interactionState, detail string) []Event {
51 events := make([]Event, 0, len(projection.ActiveTools)+len(projection.Interactions)+len(projection.ActiveSteps))
52 toolIDs := make([]string, 0, len(projection.ActiveTools))
53 for id := range projection.ActiveTools {
54 toolIDs = append(toolIDs, id)
55 }
56 sort.Strings(toolIDs)
57 for _, id := range toolIDs {
58 state := provider.ToolRunNotStarted
59 legacyState := "not_started"
60 if projection.StartedTools[id] {
61 state = provider.ToolRunUnknown
62 legacyState = "result_unknown"
63 }
64 payload, _ := json.Marshal(map[string]any{"id": id, "name": projection.ActiveTools[id], "state": legacyState, "runState": state, "error": detail})
65 events = append(events, Event{Kind: "tool/result", Payload: payload})
66 }
67 requestIDs := make([]string, 0, len(projection.Interactions))
68 for id := range projection.Interactions {
69 requestIDs = append(requestIDs, id)
70 }
71 sort.Strings(requestIDs)
72 for _, id := range requestIDs {
73 payload, _ := json.Marshal(map[string]any{"id": id, "state": interactionState})
74 events = append(events, Event{Kind: "interaction/resolved", Payload: payload})
75 }
76 stepIDs := make([]string, 0, len(projection.ActiveSteps))
77 for id := range projection.ActiveSteps {
78 stepIDs = append(stepIDs, id)
79 }
80 sort.Strings(stepIDs)
81 for _, id := range stepIDs {
82 payload, _ := json.Marshal(map[string]any{"id": id})
83 events = append(events, Event{Kind: "step/end", Payload: payload})
84 }
85 return events
86 }
87
87 lines GO