| 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 |