返回 DeepSeek-Reasonix
stream_checkpoint.go
根目录 / internal / session / stream_checkpoint.go
1 package session
2
3 import (
4 "encoding/json"
5
6 "reasonix/internal/provider"
7 )
8
9 // StreamCheckpointKind carries the output a still-open stream has produced,
10 // as the local-only record its interruption would leave. It is optional so a
11 // reader that does not know it skips it; it enters history only when restart
12 // recovery closes the turn it belongs to.
13 const StreamCheckpointKind = "stream/checkpoint"
14
15 // StreamCheckpointEvent builds the event for one open-stream checkpoint.
16 func StreamCheckpointEvent(record provider.Message) (Event, error) {
17 payload, err := json.Marshal(map[string]any{"message": record})
18 if err != nil {
19 return Event{}, err
20 }
21 return Event{Kind: StreamCheckpointKind, Optional: true, Payload: payload}, nil
22 }
23
24 // projectStreamCheckpoint keeps the newest checkpoint of the open turn. A
25 // payload this reader cannot decode is skipped like any unknown optional event.
26 func projectStreamCheckpoint(projection *Projection, ev Event) {
27 var body struct {
28 Message *provider.Message `json:"message"`
29 }
30 if projection.TurnID == "" || json.Unmarshal(ev.Payload, &body) != nil || body.Message == nil || body.Message.ID == "" || !body.Message.LocalOnly {
31 return
32 }
33 projection.StreamCheckpoint = body.Message
34 }
35
36 // supersedeStreamCheckpoint drops the checkpoint once anything after it
37 // commits a message or closes the turn: the stream it described has either
38 // become a durable message or been recorded as interrupted.
39 func supersedeStreamCheckpoint(projection *Projection, kind string) {
40 switch kind {
41 case "message/complete", "message/upsert", "message/retract", "history/replace", "legacy/import", "turn/start", "turn/end":
42 projection.StreamCheckpoint = nil
43 }
44 }
45
46 // streamCheckpointRecoveryEvents materialises the open turn's checkpoint as
47 // the interrupted record a cancelled stream would have committed.
48 func streamCheckpointRecoveryEvents(projection Projection) []Event {
49 if projection.StreamCheckpoint == nil {
50 return nil
51 }
52 payload, err := json.Marshal(map[string]any{"message": projection.StreamCheckpoint})
53 if err != nil {
54 return nil
55 }
56 return []Event{{Kind: "message/upsert", Payload: payload}}
57 }
58
58 lines GO