返回 DeepSeek-Reasonix
heads.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "database/sql"
6 "fmt"
7
8 "reasonix/internal/agent"
9 )
10
11 // HeadRecord is one head of a schema-2 session as the catalog projects it.
12 // Rows come from the kernel's head index sidecar, never from replaying a log.
13 type HeadRecord struct {
14 Path string `json:"path"`
15 ID string `json:"id"`
16 ParentHeadID string `json:"parentHeadId,omitempty"`
17 Kind string `json:"kind"`
18 Name string `json:"name,omitempty"`
19 LeafMessageID string `json:"leafMessageId,omitempty"`
20 WriterID string `json:"writerId,omitempty"`
21 LastActivityAt int64 `json:"lastActivityAt,omitempty"`
22 Turns int `json:"turns"`
23 Preview string `json:"preview,omitempty"`
24 Retired bool `json:"retired,omitempty"`
25 Selected bool `json:"selected,omitempty"`
26 }
27
28 // sessionHeadProjection is what recordFromOrder learns about a schema-2
29 // session without opening its log: the sidecar mirror is always available,
30 // the per-head rows only while the head index still matches the log.
31 type sessionHeadProjection struct {
32 logFormat int
33 headCount int
34 selected string
35 fingerprint string
36 heads []HeadRecord
37 stale bool
38 }
39
40 func projectSessionHeads(info agent.SessionOrderInfo) sessionHeadProjection {
41 if info.LogSchema < 2 {
42 return sessionHeadProjection{logFormat: 1}
43 }
44 out := sessionHeadProjection{logFormat: info.LogSchema, headCount: info.HeadCount, selected: info.HeadID}
45 idx, err := agent.ReadSessionHeadIndex(info.Path)
46 if err != nil || idx == nil || !idx.Current(info.Path) {
47 out.stale = true
48 return out
49 }
50 out.selected = idx.SelectedHead
51 out.headCount = len(idx.Heads)
52 out.heads, out.fingerprint = headRecordsFromIndex(info.Path, idx)
53 return out
54 }
55
56 func headRecordsFromIndex(path string, idx *agent.SessionHeadIndex) ([]HeadRecord, string) {
57 heads := make([]HeadRecord, 0, len(idx.Heads))
58 fingerprint := ""
59 for _, h := range idx.Heads {
60 heads = append(heads, HeadRecord{
61 Path: path, ID: h.ID, ParentHeadID: h.ParentHead, Kind: h.Kind, Name: h.Name,
62 LeafMessageID: h.LeafID, WriterID: h.Writer, LastActivityAt: unixMilli(h.LastActivity),
63 Turns: h.Turns, Preview: h.Preview, Retired: h.Retired, Selected: h.ID == idx.SelectedHead,
64 })
65 if h.ID == idx.SelectedHead {
66 fingerprint = "|h:" + h.ID + ":" + h.LeafID
67 }
68 }
69 return heads, fingerprint
70 }
71
72 // projectRepairedHeads lands the head rows of a session whose repair just
73 // completed inside that transaction: repair rebuilds a stale head index, and
74 // the row must not turn valid while its head_count names rows it lacks.
75 func projectRepairedHeads(ctx context.Context, tx *sql.Tx, path, pathKey string) error {
76 idx, err := agent.ReadSessionHeadIndex(path)
77 if err != nil || idx == nil || !idx.Current(path) {
78 return nil
79 }
80 heads, _ := headRecordsFromIndex(path, idx)
81 if _, err := tx.ExecContext(ctx, `UPDATE catalog_sessions SET log_format=?,head_count=?,selected_head_id=?
82 WHERE path_key=?`, idx.SchemaVersion, len(heads), idx.SelectedHead, pathKey); err != nil {
83 return err
84 }
85 return upsertHeadRows(ctx, tx, pathKey, heads)
86 }
87
88 // writeDirectoryRow lands one scanned session row and its head rows inside
89 // the directory projection transaction.
90 func (c *Catalog) writeDirectoryRow(ctx context.Context, tx *sql.Tx, stmt *sql.Stmt, record SessionRecord, pathKey, directoryKey string, generation int64) error {
91 if _, err := stmt.ExecContext(ctx, c.sessionRowValues(record, pathKey, directoryKey, generation)...); err != nil {
92 return err
93 }
94 return upsertHeadRows(ctx, tx, pathKey, record.heads)
95 }
96
97 // upsertHeadRows replaces the head rows of one session inside the projection
98 // transaction. Stale projections keep the previous rows until the head index
99 // is current again.
100 func upsertHeadRows(ctx context.Context, tx *sql.Tx, pathKey string, heads []HeadRecord) error {
101 if heads == nil {
102 return nil
103 }
104 if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_heads WHERE path_key=?`, pathKey); err != nil {
105 return err
106 }
107 for _, h := range heads {
108 if _, err := tx.ExecContext(ctx, `INSERT INTO catalog_heads(path_key,head_id,parent_head_id,kind,name,leaf_message_id,
109 writer_id,last_activity_at,turns,preview,retired,selected) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)`,
110 pathKey, h.ID, h.ParentHeadID, h.Kind, h.Name, h.LeafMessageID, h.WriterID, h.LastActivityAt, h.Turns, h.Preview,
111 boolToInt(h.Retired), boolToInt(h.Selected)); err != nil {
112 return err
113 }
114 }
115 return nil
116 }
117
118 // ListHeads returns the projected heads of one session in creation order. A
119 // schema-1 session has none and returns an empty slice.
120 func (c *Catalog) ListHeads(ctx context.Context, path string) ([]HeadRecord, error) {
121 out := []HeadRecord{}
122 if c == nil || path == "" {
123 return out, nil
124 }
125 rows, err := c.readDB(ctx).QueryContext(ctx, `SELECT head_id,parent_head_id,kind,name,leaf_message_id,writer_id,last_activity_at,
126 turns,preview,retired,selected FROM catalog_heads WHERE path_key=? ORDER BY rowid`, c.pathKey(path))
127 if err != nil {
128 return out, fmt.Errorf("list session heads: %w", err)
129 }
130 defer rows.Close()
131 for rows.Next() {
132 var h HeadRecord
133 var retired, selected int
134 if err := rows.Scan(&h.ID, &h.ParentHeadID, &h.Kind, &h.Name, &h.LeafMessageID, &h.WriterID, &h.LastActivityAt,
135 &h.Turns, &h.Preview, &retired, &selected); err != nil {
136 return out, err
137 }
138 h.Path = path
139 h.Retired, h.Selected = retired != 0, selected != 0
140 out = append(out, h)
141 }
142 return out, rows.Err()
143 }
144
144 lines GO