| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| 5 | "context" |
| 6 | "encoding/json" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | "strings" |
| 10 | |
| 11 | "reasonix/internal/agent" |
| 12 | "reasonix/internal/attachment" |
| 13 | "reasonix/internal/provider" |
| 14 | "reasonix/internal/sessioncontent" |
| 15 | ) |
| 16 | |
| 17 | func indexOneMessage(ctx context.Context, content *sessioncontent.Store, state *historyBuildState, message provider.Message, sequence uint64, upsert bool) error { |
| 18 | body, err := json.Marshal(message) |
| 19 | if err != nil { |
| 20 | return err |
| 21 | } |
| 22 | return indexOneMessageBody(ctx, content, state, message, sequence, upsert, body) |
| 23 | } |
| 24 | |
| 25 | func indexOneMessageBody(ctx context.Context, content *sessioncontent.Store, state *historyBuildState, message provider.Message, sequence uint64, upsert bool, body json.RawMessage) error { |
| 26 | id := strings.TrimSpace(message.ID) |
| 27 | if id == "" { |
| 28 | return errors.New("session: indexed message has no stable id") |
| 29 | } |
| 30 | position, exists := state.positions[id] |
| 31 | visibleTurn := state.turns[id] |
| 32 | if !exists { |
| 33 | state.nextPosition++ |
| 34 | position = state.nextPosition |
| 35 | state.positions[id] = position |
| 36 | if agent.IsUserAuthoredTurnMessage(message) { |
| 37 | state.visibleTurn++ |
| 38 | } |
| 39 | visibleTurn = state.visibleTurn |
| 40 | state.turns[id] = visibleTurn |
| 41 | } else if !upsert { |
| 42 | return fmt.Errorf("session: duplicate indexed message id %q", id) |
| 43 | } |
| 44 | version := state.versions[id] + 1 |
| 45 | state.versions[id] = version |
| 46 | if exists { |
| 47 | if err := flushHistoryBuildRows(ctx, state.tx, state); err != nil { |
| 48 | return err |
| 49 | } |
| 50 | if _, err := state.statements.expire.ExecContext(ctx, sequence, id); err != nil { |
| 51 | return err |
| 52 | } |
| 53 | } |
| 54 | ref, err := content.Put(ctx, bytes.NewReader(body), sessioncontent.Metadata{MediaType: "application/json"}) |
| 55 | if err != nil { |
| 56 | return err |
| 57 | } |
| 58 | if err := insertContentRef(ctx, state, ref); err != nil { |
| 59 | return err |
| 60 | } |
| 61 | for _, extra := range attachment.CollectContentRefs(message.ImageInputs) { |
| 62 | if err := insertContentRef(ctx, state, extra); err != nil { |
| 63 | return err |
| 64 | } |
| 65 | } |
| 66 | if message.Role == provider.RoleTool && message.ToolCallID != "" { |
| 67 | if _, err := state.tx.ExecContext(ctx, `INSERT OR IGNORE INTO tool_links(message_id,digest,call_id,is_result,state) VALUES(?,?,?,1,?)`, id, ref.Digest, message.ToolCallID, string(provider.ToolResultRunState(message))); err != nil { |
| 68 | return err |
| 69 | } |
| 70 | } |
| 71 | for _, call := range message.ToolCalls { |
| 72 | if call.ID != "" { |
| 73 | if _, err := state.tx.ExecContext(ctx, `INSERT OR IGNORE INTO tool_links(message_id,digest,call_id,is_result,state) VALUES(?,?,?,0,'unknown')`, id, ref.Digest, call.ID); err != nil { |
| 74 | return err |
| 75 | } |
| 76 | } |
| 77 | } |
| 78 | visibleUser := 0 |
| 79 | if agent.IsUserAuthoredTurnMessage(message) { |
| 80 | visibleUser = 1 |
| 81 | } |
| 82 | state.messages = append(state.messages, []any{id, version, position, sequence, 0, string(message.Role), messagePreview(message), nil, ref.Digest, ref.Bytes, ref.IndexDigest, 1, "", visibleTurn, visibleUser}) |
| 83 | return nil |
| 84 | } |
| 85 | |
| 86 | // Re-number only rows whose visible user boundary changed. New versions keep |
| 87 | // previous numbering available to fixed-snapshot cursors. |
| 88 |