返回 DeepSeek-Reasonix
metadata_scan.go
根目录 / internal / sessioncatalog / metadata_scan.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "io"
8 "os"
9 "path/filepath"
10 "strings"
11 "sync"
12 "time"
13
14 "reasonix/internal/agent"
15 "reasonix/internal/historywork"
16 "reasonix/internal/store"
17 )
18
19 // metadataRecord never opens the transcript, display index, or event log.
20 // Listing fields without a valid sidecar remain explicitly unknown.
21 func metadataRecord(ctx context.Context, target DirectoryTarget, path string) (SessionRecord, error) {
22 st, err := os.Stat(path)
23 if err != nil {
24 return SessionRecord{}, err
25 }
26 if !st.Mode().IsRegular() {
27 return SessionRecord{}, fmt.Errorf("not a session file")
28 }
29 meta, known, metaErr := agent.LoadBranchMetaBounded(ctx, path)
30 if err := ctx.Err(); err != nil {
31 return SessionRecord{}, err
32 }
33 order := agent.SessionOrderInfo{Path: path, Scope: target.Scope, WorkspaceRoot: target.WorkspaceRoot,
34 CreatedAt: st.ModTime(), LastActivityAt: st.ModTime(), ModTime: st.ModTime()}
35 if known && metaErr == nil {
36 if strings.TrimSpace(meta.Scope) != "" {
37 order.Scope, order.WorkspaceRoot = meta.DefaultScope(), meta.WorkspaceRoot
38 }
39 order.TopicID, order.TopicTitle, order.CustomTitle = meta.TopicID, meta.TopicTitle, meta.CustomTitle
40 if !meta.CreatedAt.IsZero() {
41 order.CreatedAt = meta.CreatedAt
42 }
43 if !meta.UpdatedAt.IsZero() {
44 order.LastActivityAt = meta.UpdatedAt
45 }
46 order.Recovered, order.ParentID = meta.Recovered, meta.ParentID
47 order.RecoveryReason, order.RecoveryDigest = meta.RecoveryReason, meta.RecoveryDigest
48 order.Turns, order.Preview, order.SchemaVersion = meta.Turns, meta.Preview, meta.SchemaVersion
49 order.Revision, order.ContentDigest = meta.Revision, meta.ContentDigest
50 order.ListingRevision, order.ListingContentDigest = meta.ListingRevision, meta.ListingContentDigest
51 }
52 // Every physical session remains reachable without rewriting old sidecars
53 // merely to assign a topic. A later authoritative topic assignment replaces it.
54 if order.TopicID == "" {
55 order.TopicID = agent.LegacySessionTopicID(path)
56 }
57 if order.TopicTitle == "" {
58 order.TopicTitle = strings.TrimSuffix(filepath.Base(path), ".jsonl")
59 }
60 record := recordFromOrder(target, order) // LogSchema is deliberately unset: no whole head-index load.
61 record.LogFormat, record.HeadCount, record.SelectedHeadID = max(1, meta.LogSchema), meta.HeadCount, meta.HeadID
62 record.LogicalTopicID, record.OrdinaryVisible = record.TopicID, true
63 if metaErr != nil {
64 record.Health, record.TurnsState = HealthDegraded, TurnsUnknown
65 }
66 return record, nil
67 }
68
69 func (c *Catalog) indexMetadataPath(ctx context.Context, target DirectoryTarget, path string, sequence uint64) error {
70 path = cleanCatalogAccessPath(path)
71 if target.Path == "" {
72 target.Path = filepath.Dir(path)
73 }
74 target.Path = cleanCatalogAccessPath(target.Path)
75 target.Scope, target.WorkspaceRoot = normalizeScope(target.Scope, target.WorkspaceRoot)
76 lock := c.directoryLock(target.Path)
77 lock.Lock()
78 defer lock.Unlock()
79 record, err := metadataRecord(ctx, target, path)
80 if os.IsNotExist(err) {
81 return nil
82 }
83 if err != nil {
84 return err
85 }
86 record.enqueueSequence = sequence
87 _, err = c.upsertExactPathSession(ctx, record)
88 return err
89 }
90
91 // reconcileMetadata is the explicit, synchronous entry point. The ordinary
92 // queue drives the same scan one slice at a time, preserving its directory
93 // iterator while other roots receive their own slices.
94 func (c *Catalog) reconcileMetadata(ctx context.Context, target DirectoryTarget, sequence uint64) error {
95 scan, err := c.startMetadataScan(ctx, target, sequence, false)
96 if err != nil {
97 return err
98 }
99 defer func() { scan.close(ctx, err) }()
100 for {
101 var done bool
102 var bytes int64
103 done, bytes, err = scan.step(ctx)
104 if err != nil || done {
105 return err
106 }
107 if c.opts.Maintenance == nil {
108 err = historywork.Pause(ctx, bytes)
109 }
110 if err != nil {
111 return err
112 }
113 }
114 }
115
116 type metadataScan struct {
117 c *Catalog
118 target DirectoryTarget
119 sequence uint64
120 generation int64
121 signature string
122 started int64
123 total int
124 file *os.File
125 lock *sync.Mutex
126 scanLock *sync.Mutex
127 yield bool
128 }
129
130 var errMetadataScanBusy = errors.New("directory scan already active")
131
132 func (c *Catalog) startMetadataScan(ctx context.Context, target DirectoryTarget, sequence uint64, try bool) (_ *metadataScan, result error) {
133 target.Path = cleanCatalogAccessPath(target.Path)
134 target.Scope, target.WorkspaceRoot = normalizeScope(target.Scope, target.WorkspaceRoot)
135 value, _ := c.metadataScans.LoadOrStore(c.pathKey(target.Path), &sync.Mutex{})
136 scanLock := value.(*sync.Mutex)
137 if try {
138 if !scanLock.TryLock() {
139 return nil, errMetadataScanBusy
140 }
141 } else {
142 scanLock.Lock()
143 }
144 scan := &metadataScan{c: c, target: target, sequence: sequence, scanLock: scanLock, lock: c.directoryLock(target.Path), yield: try}
145 defer func() {
146 if result != nil {
147 scan.close(ctx, result)
148 }
149 }()
150 if err := ctx.Err(); err != nil {
151 return nil, err
152 }
153 scan.started = c.opts.Now().UnixMilli()
154 scan.signature = fmt.Sprintf("metadata:%d:%d", scan.started, sequence)
155 scan.lock.Lock()
156 generation, _, err := c.beginDirectoryScan(ctx, target, scan.signature, scan.started)
157 scan.lock.Unlock()
158 if err != nil {
159 return nil, err
160 }
161 scan.generation = generation
162 scan.file, err = os.Open(target.Path)
163 if os.IsNotExist(err) {
164 // Optional legacy roots do not exist on a fresh installation. They are
165 // an empty discovery only while the catalog has never retained a row
166 // there. An unavailable root with known history may be an unplugged
167 // device; never turn that into missing-row confirmation.
168 var known bool
169 queryErr := c.readDB(ctx).QueryRowContext(ctx, `SELECT EXISTS(
170 SELECT 1 FROM catalog_sessions WHERE directory_key=? LIMIT 1)`, c.pathKey(target.Path)).Scan(&known)
171 if queryErr != nil {
172 return nil, queryErr
173 }
174 if !known {
175 return scan, nil
176 }
177 }
178 if err != nil {
179 return nil, err
180 }
181 st, err := scan.file.Stat()
182 if err != nil {
183 return nil, err
184 }
185 if !st.IsDir() {
186 return nil, fmt.Errorf("not a directory")
187 }
188 return scan, nil
189 }
190
191 func (s *metadataScan) close(ctx context.Context, result error) {
192 if s.file != nil {
193 _ = s.file.Close()
194 s.file = nil
195 }
196 if result != nil && s.generation != 0 {
197 s.c.failDirectoryScan(context.WithoutCancel(ctx), s.target.Path, result)
198 }
199 if s.scanLock != nil {
200 s.scanLock.Unlock()
201 s.scanLock = nil
202 }
203 }
204
205 // step commits only the visited prefix. Only a successful EOF can mark
206 // missing rows, including when shutdown, cancellation or device loss occurs.
207 func (s *metadataScan) step(ctx context.Context) (done bool, bytes int64, result error) {
208 if err := ctx.Err(); err != nil {
209 return false, 0, err
210 }
211 c, target := s.c, s.target
212 release := func(int64) {}
213 if c.opts.Maintenance != nil {
214 ctx = c.opts.Maintenance.Context(ctx)
215 var err error
216 if s.yield {
217 release, err = c.opts.Maintenance.BackgroundSliceYielding(ctx, c.isPriorityDirectory(target))
218 } else {
219 release, err = c.opts.Maintenance.BackgroundSlice(ctx, c.isPriorityDirectory(target))
220 }
221 if err != nil {
222 return false, 0, err
223 }
224 }
225 defer func() { release(bytes) }()
226 start := time.Now()
227 records := make([]SessionRecord, 0, historywork.BatchEntries)
228 s.lock.Lock()
229 defer s.lock.Unlock()
230 for count := 0; count < historywork.BatchEntries && bytes+historywork.ReadChunk+1 <= historywork.BatchBytes && time.Since(start) < historywork.SliceDuration; count++ {
231 if err := ctx.Err(); err != nil {
232 return false, bytes, err
233 }
234 if s.file == nil {
235 done = true
236 break
237 }
238 entries, readErr := s.file.ReadDir(1)
239 if errors.Is(readErr, io.EOF) {
240 done = true
241 break
242 }
243 if readErr != nil {
244 return false, bytes, readErr
245 }
246 entry := entries[0]
247 if entry.IsDir() || entry.Type()&os.ModeSymlink != 0 || !store.IsSessionTranscriptName(entry.Name()) {
248 continue
249 }
250 record, readErr := metadataRecord(ctx, target, filepath.Join(target.Path, entry.Name()))
251 if os.IsNotExist(readErr) {
252 continue
253 }
254 if readErr != nil {
255 return false, bytes, readErr
256 }
257 if metaInfo, statErr := os.Stat(agent.BranchMetaPath(record.Path)); statErr == nil {
258 bytes += min(metaInfo.Size()+1, int64(historywork.ReadChunk+1))
259 }
260 if existing, found, loadErr := c.GetSession(ctx, record.Path); loadErr != nil {
261 return false, bytes, loadErr
262 } else if found && existing.ContentFingerprint == record.ContentFingerprint && existing.MetaFingerprint == record.MetaFingerprint {
263 preserveDirectoryProjection(existing, &record)
264 if existing.TurnsState != TurnsUnknown {
265 record.Turns, record.Preview, record.TurnsState = existing.Turns, existing.Preview, existing.TurnsState
266 }
267 record.metadataUnchanged = sameMetadataProjection(existing, record)
268 }
269 record.enqueueSequence = s.sequence
270 records = append(records, record)
271 }
272 if c.testReconcileBatchHook != nil {
273 c.testReconcileBatchHook(len(records))
274 }
275 _, err := c.upsertSessionsWithNotification(ctx, records, nil, "metadata_batch", true, upsertDirectoryProjection)
276 if err != nil {
277 return false, bytes, err
278 }
279 s.total += len(records)
280 err = c.updateDirectoryScanProgress(ctx, target.Path, s.generation, s.total)
281 if err == nil && done {
282 err = c.finishDirectoryScan(ctx, target, s.signature, s.generation, s.started, s.total)
283 }
284 return done, bytes, err
285 }
286
286 lines GO