返回 DeepSeek-Reasonix
metadata_read_cache.go
根目录 / internal / session / metadata_read_cache.go
1 package session
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/binary"
7 "encoding/json"
8 "fmt"
9 "io"
10 "os"
11 "path/filepath"
12 "strings"
13 "sync"
14 "time"
15
16 "reasonix/internal/historywork"
17 )
18
19 type sessionMetadataRead struct {
20 fingerprint [32]byte
21 info SessionInfo
22 used uint64
23 }
24
25 type sessionMetadataReads struct {
26 mu sync.Mutex
27 clock uint64
28 entries map[string]sessionMetadataRead
29 }
30
31 // Keep list metadata separate from execution state. Validation reads actual
32 // metadata bytes and the selected log revision, including external writers;
33 // no body replay, TTL staleness or writer ownership is involved.
34 func (p *FilesystemPersistence) cachedSessionInfo(ctx context.Context, id string) (SessionInfo, error) {
35 if err := ctx.Err(); err != nil {
36 return SessionInfo{}, err
37 }
38 id = strings.TrimSpace(id)
39 input, key, err := p.readSessionMetadata(ctx, id)
40 if err != nil {
41 return SessionInfo{}, err
42 }
43 c := &p.metadataReads
44 c.mu.Lock()
45 c.clock++
46 if entry, ok := c.entries[id]; ok && entry.fingerprint == key {
47 entry.used = c.clock
48 c.entries[id] = entry
49 c.mu.Unlock()
50 return entry.info, nil
51 }
52 c.mu.Unlock()
53 info, err := input.info(id)
54 if err != nil {
55 return info, err
56 }
57 c.mu.Lock()
58 defer c.mu.Unlock()
59 if c.entries == nil {
60 c.entries = map[string]sessionMetadataRead{}
61 }
62 if len(c.entries) >= 2048 {
63 oldest := ""
64 var age uint64
65 for name, entry := range c.entries {
66 if oldest == "" || entry.used < age {
67 oldest, age = name, entry.used
68 }
69 }
70 delete(c.entries, oldest)
71 }
72 c.clock++
73 c.entries[id] = sessionMetadataRead{key, info, c.clock}
74 return info, nil
75 }
76
77 type sessionMetadataInput struct {
78 dir string
79 manifest Manifest
80 header, catalog []byte
81 hasHeader bool
82 revision logRevision
83 }
84
85 func (p *FilesystemPersistence) readSessionMetadata(ctx context.Context, id string) (sessionMetadataInput, [32]byte, error) {
86 var input sessionMetadataInput
87 var zero [32]byte
88 dir, err := p.sessionDir(id, true)
89 if err != nil {
90 return input, zero, err
91 }
92 input.dir = dir
93 h := sha256.New()
94 _, _ = fmt.Fprintf(h, "root:%s\x00", dir)
95 var manifest Manifest
96 // Derive cache identity from the confined directory, as the store writer does.
97 // Keep protocol input out of filesystem paths after sessionDir resolves it.
98 cacheDir := filepath.Join(filepath.Dir(dir), ".query-cache", filepath.Base(dir))
99 paths := []string{filepath.Join(dir, "manifest.json"), filepath.Join(dir, sessionHeaderName), catalogMetadataPath(cacheDir)}
100 for index, path := range paths {
101 body, readErr := readListingMetadata(ctx, path)
102 if index == 0 {
103 if readErr != nil {
104 if os.IsNotExist(readErr) {
105 readErr = fmt.Errorf("%w: %s", ErrSessionNotFound, id)
106 }
107 return input, zero, readErr
108 }
109 if err := json.Unmarshal(body, &manifest); err != nil {
110 return input, zero, err
111 }
112 if !supportedStoredManifest(manifest) {
113 return input, zero, ErrUnsupportedVersion
114 }
115 if manifest.SessionID != id {
116 return input, zero, fmt.Errorf("%w: manifest belongs to another session", ErrDamagedStore)
117 }
118 input.manifest = manifest
119 } else if index == 1 {
120 if readErr != nil && !os.IsNotExist(readErr) {
121 return input, zero, readErr
122 }
123 input.header, input.hasHeader = body, readErr == nil
124 } else if readErr == nil {
125 input.catalog = body
126 }
127 // Distinguish absent/unreadable fields from a readable empty file.
128 _, _ = fmt.Fprintf(h, "%d:%v:", index, readErr)
129 _ = binary.Write(h, binary.LittleEndian, uint64(len(body)))
130 _, _ = h.Write(body)
131 }
132 stat, err := os.Stat(logPathForManifest(dir, manifest))
133 if err != nil && !os.IsNotExist(err) {
134 return input, zero, err
135 }
136 if stat != nil {
137 _, _ = fmt.Fprintf(h, "log:%d:%d", stat.Size(), stat.ModTime().UnixNano())
138 input.revision = logRevision{Size: stat.Size(), ModTimeNS: stat.ModTime().UnixNano(), Exists: true}
139 }
140 copy(zero[:], h.Sum(nil))
141 return input, zero, nil
142 }
143
144 // Listing metadata is optional display information. An oversized sidecar must
145 // not turn discovery into an unbounded read; its caller keeps the source as an
146 // unknown/degraded entry and explicit session recovery uses its native reader.
147 func readListingMetadata(ctx context.Context, path string) ([]byte, error) {
148 f, err := os.Open(path)
149 if err != nil {
150 return nil, err
151 }
152 defer f.Close()
153 body, err := io.ReadAll(io.LimitReader(&historywork.Reader{Context: ctx, Source: f}, historywork.ReadChunk+1))
154 if err == nil && len(body) > historywork.ReadChunk {
155 return nil, fmt.Errorf("listing metadata exceeds read budget")
156 }
157 return body, err
158 }
159
160 func (in sessionMetadataInput) info(id string) (SessionInfo, error) {
161 m := in.manifest
162 info := SessionInfo{SessionID: id, Codec: m.Codec, Kind: m.Kind, CreatedAt: m.CreatedAt, UpdatedAt: m.CreatedAt, MetadataStatus: MetadataPending, Path: in.dir}
163 if in.revision.Exists && in.revision.ModTimeNS > info.UpdatedAt.UnixNano() {
164 info.UpdatedAt = time.Unix(0, in.revision.ModTimeNS)
165 }
166 if in.hasHeader {
167 header, _, err := decodeSessionHeader(in.header, id)
168 if err != nil {
169 return SessionInfo{}, err
170 }
171 info.CWD, info.ParentSessionID, info.Origin = header.CWD, header.ParentSessionID, header.Origin
172 }
173 if metadata, err := decodeCatalogMetadata(in.catalog, m, in.revision); err == nil {
174 info.Title, info.TitleSequence = metadata.Title, metadata.TitleSequence
175 info.ModelRef, info.ModelIdentity = metadata.ModelRef, metadata.ModelIdentity
176 info.Turns, info.Preview, info.MetadataStatus = metadata.Turns, metadata.Preview, MetadataReady
177 info.EventSequence, info.ResultSequence = metadata.Sequence, metadata.ResultSequence
178 }
179 return info, nil
180 }
181
181 lines GO