返回 DeepSeek-Reasonix
history_message_index.go
根目录 / internal / session / history_message_index.go
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
88 lines GO