返回 DeepSeek-Reasonix
store.go
1 // Package sessioncontent stores immutable session payloads by content digest.
2 // It deliberately owns bytes only: session ordering, authorization and model
3 // projection remain responsibilities of their respective services.
4 package sessioncontent
5
6 import (
7 "bytes"
8 "context"
9 "crypto/sha256"
10 "encoding/hex"
11 "encoding/json"
12 "errors"
13 "fmt"
14 "hash"
15 "io"
16 "os"
17 "path/filepath"
18 "runtime"
19 "syscall"
20
21 filelock "reasonix/internal/identitylock"
22 )
23
24 const (
25 copyBufferBytes = 1 << 20
26 // IntegrityBlockBytes bounds range verification work independently of the
27 // total object size.
28 IntegrityBlockBytes = 1 << 20
29 // MaxReadRange bounds one allocation, not the size of an object or session.
30 MaxReadRange = 8 << 20
31 )
32
33 // Metadata describes how callers display or interpret content. None of these
34 // fields participate in the storage path; identical bytes are deduplicated.
35 type Metadata struct {
36 MediaType string
37 Name string
38 }
39
40 // Ref is the durable, path-independent identity of one immutable object.
41 type Ref struct {
42 Digest string `json:"digest"`
43 Bytes int64 `json:"bytes"`
44 MediaType string `json:"mediaType,omitempty"`
45 Name string `json:"name,omitempty"`
46 IndexDigest string `json:"indexDigest,omitempty"`
47 IntegrityBlock int64 `json:"integrityBlockBytes,omitempty"`
48 }
49
50 type integrityIndex struct {
51 Version int `json:"version"`
52 Object string `json:"object"`
53 Bytes int64 `json:"bytes"`
54 BlockBytes int64 `json:"blockBytes"`
55 Blocks []string `json:"blocks"`
56 }
57
58 // Store is a process-independent content-addressed object store. Publishing is
59 // no-overwrite: concurrent writers of the same digest converge on one object.
60 type Store struct {
61 root string
62 }
63
64 // ObjectWitness is a process-local observation of an immutable object. It is
65 // suitable only for deciding whether a previously verified cache entry may be
66 // reused; authorization remains the caller's responsibility.
67 type ObjectWitness struct {
68 info os.FileInfo
69 bytes int64
70 modTimeNano int64
71 indexDigest string
72 }
73
74 func (w ObjectWitness) Same(other ObjectWitness) bool {
75 return w.info != nil && other.info != nil && w.bytes == other.bytes &&
76 w.modTimeNano == other.modTimeNano && w.indexDigest == other.indexDigest &&
77 os.SameFile(w.info, other.info)
78 }
79
80 func (w ObjectWitness) Valid() bool { return w.info != nil }
81
82 func New(root string) *Store { return &Store{root: root} }
83
84 func (s *Store) Root() string {
85 if s == nil {
86 return ""
87 }
88 return s.root
89 }
90
91 // Put streams r into a private temporary file, fsyncs it, then atomically
92 // publishes the file under its SHA-256 digest. A returned Ref always names a
93 // complete object. Cancellation never publishes the partial temporary file.
94 func (s *Store) Put(ctx context.Context, r io.Reader, meta Metadata) (Ref, error) {
95 if s == nil || s.root == "" {
96 return Ref{}, errors.New("session content store unavailable")
97 }
98 if r == nil {
99 return Ref{}, errors.New("session content reader is nil")
100 }
101 if err := ctx.Err(); err != nil {
102 return Ref{}, err
103 }
104 tmpDir := filepath.Join(s.root, ".tmp")
105 if err := os.MkdirAll(tmpDir, 0o700); err != nil {
106 return Ref{}, fmt.Errorf("create session content temp directory: %w", err)
107 }
108 tmp, err := os.CreateTemp(tmpDir, "content-*.tmp")
109 if err != nil {
110 return Ref{}, fmt.Errorf("create session content temporary file: %w", err)
111 }
112 tmpPath := tmp.Name()
113 closed := false
114 defer func() {
115 if !closed {
116 _ = tmp.Close()
117 }
118 _ = os.Remove(tmpPath)
119 }()
120
121 digests := newBlockDigestWriter(tmp)
122 n, err := copyWithContext(ctx, digests, r)
123 if err != nil {
124 return Ref{}, fmt.Errorf("stage session content: %w", err)
125 }
126 if err := tmp.Sync(); err != nil {
127 return Ref{}, fmt.Errorf("fsync session content: %w", err)
128 }
129 if err := tmp.Chmod(0o600); err != nil {
130 return Ref{}, fmt.Errorf("protect session content: %w", err)
131 }
132 if err := tmp.Close(); err != nil {
133 return Ref{}, fmt.Errorf("close session content: %w", err)
134 }
135 closed = true
136
137 ref := Ref{
138 Digest: hex.EncodeToString(digests.full.Sum(nil)),
139 Bytes: n,
140 MediaType: meta.MediaType,
141 Name: meta.Name,
142 IntegrityBlock: IntegrityBlockBytes,
143 }
144 index := integrityIndex{Version: 1, Object: ref.Digest, Bytes: ref.Bytes, BlockBytes: IntegrityBlockBytes, Blocks: digests.finish()}
145 indexBytes, err := json.Marshal(index)
146 if err != nil {
147 return Ref{}, fmt.Errorf("encode session content integrity index: %w", err)
148 }
149 indexSum := sha256.Sum256(indexBytes)
150 ref.IndexDigest = hex.EncodeToString(indexSum[:])
151 dest := s.objectPath(ref.Digest)
152 if err := os.MkdirAll(filepath.Dir(dest), 0o700); err != nil {
153 return Ref{}, fmt.Errorf("create session content object directory: %w", err)
154 }
155 if err := s.publishObject(ctx, tmpPath, dest, ref); err != nil {
156 return Ref{}, err
157 }
158 if err := s.publishIndex(ctx, ref, indexBytes); err != nil {
159 return Ref{}, err
160 }
161 return ref, nil
162 }
163
164 // Open validates the complete immutable object before returning it positioned
165 // at byte zero. This favors integrity over trusting a mutable local filesystem.
166 func (s *Store) Open(ctx context.Context, ref Ref) (*os.File, error) {
167 if err := validateRef(ref); err != nil {
168 return nil, err
169 }
170 f, err := s.openRaw(ref)
171 if err != nil {
172 return nil, err
173 }
174 if err := s.verifyOpenFile(ctx, f, ref); err != nil {
175 _ = f.Close()
176 return nil, err
177 }
178 if _, err := f.Seek(0, io.SeekStart); err != nil {
179 _ = f.Close()
180 return nil, fmt.Errorf("rewind session content %s: %w", ref.Digest, err)
181 }
182 return f, nil
183 }
184
185 // Verify checks the immutable object and its block index without materializing
186 // the object. It is the explicit integrity boundary used by import/export and
187 // diagnostics.
188 func (s *Store) Verify(ctx context.Context, ref Ref) error {
189 if err := validateRef(ref); err != nil {
190 return err
191 }
192 f, err := s.openRaw(ref)
193 if err != nil {
194 return err
195 }
196 defer f.Close()
197 return s.verifyOpenFile(ctx, f, ref)
198 }
199
200 // Stat validates the object and returns its caller-owned display metadata.
201 func (s *Store) Stat(ctx context.Context, ref Ref) (Ref, error) {
202 if err := s.Verify(ctx, ref); err != nil {
203 return Ref{}, err
204 }
205 return ref, nil
206 }
207
208 // Probe opens the object through the bounded content root and validates its
209 // immutable integrity index without hashing the complete original.
210 func (s *Store) Probe(ctx context.Context, ref Ref) (ObjectWitness, error) {
211 if err := validateRef(ref); err != nil {
212 return ObjectWitness{}, err
213 }
214 if err := ctx.Err(); err != nil {
215 return ObjectWitness{}, err
216 }
217 f, err := s.openRaw(ref)
218 if err != nil {
219 return ObjectWitness{}, err
220 }
221 defer f.Close()
222 info, err := f.Stat()
223 if err != nil {
224 return ObjectWitness{}, err
225 }
226 if ref.IndexDigest != "" {
227 if _, err := s.readIndex(ctx, ref); err != nil {
228 return ObjectWitness{}, err
229 }
230 }
231 return ObjectWitness{info: info, bytes: info.Size(), modTimeNano: info.ModTime().UnixNano(), indexDigest: ref.IndexDigest}, nil
232 }
233
234 // ReadRange reads exactly length bytes starting at offset after validating the
235 // object. The per-call allocation is bounded independently of object size.
236 func (s *Store) ReadRange(ctx context.Context, ref Ref, offset, length int64) ([]byte, error) {
237 if offset < 0 || length < 0 {
238 return nil, errors.New("session content range must be non-negative")
239 }
240 if length > MaxReadRange {
241 return nil, fmt.Errorf("session content range %d exceeds per-read budget %d", length, MaxReadRange)
242 }
243 if offset > ref.Bytes || length > ref.Bytes-offset {
244 return nil, fmt.Errorf("session content range [%d,%d) exceeds object size %d", offset, offset+length, ref.Bytes)
245 }
246 if err := validateRef(ref); err != nil {
247 return nil, err
248 }
249 f, err := s.openRaw(ref)
250 if err != nil {
251 return nil, err
252 }
253 defer f.Close()
254 if ref.IndexDigest == "" {
255 if err := s.verifyOpenFile(ctx, f, ref); err != nil {
256 return nil, err
257 }
258 } else if err := s.verifyRange(ctx, f, ref, offset, length); err != nil {
259 return nil, err
260 }
261 if _, err := f.Seek(offset, io.SeekStart); err != nil {
262 return nil, fmt.Errorf("seek session content %s: %w", ref.Digest, err)
263 }
264 buf := make([]byte, int(length))
265 if _, err := io.ReadFull(f, buf); err != nil {
266 return nil, fmt.Errorf("read session content %s range: %w", ref.Digest, err)
267 }
268 return buf, nil
269 }
270
271 func (s *Store) openRaw(ref Ref) (*os.File, error) {
272 if s == nil || s.root == "" {
273 return nil, errors.New("session content store unavailable")
274 }
275 root, err := os.OpenRoot(s.root)
276 if err != nil {
277 return nil, fmt.Errorf("open session content root: %w", err)
278 }
279 defer root.Close()
280 f, err := root.Open(s.objectRelativePath(ref.Digest))
281 if err != nil {
282 return nil, fmt.Errorf("open session content %s: %w", ref.Digest, err)
283 }
284 info, err := f.Stat()
285 if err != nil {
286 _ = f.Close()
287 return nil, fmt.Errorf("stat session content %s: %w", ref.Digest, err)
288 }
289 if !info.Mode().IsRegular() {
290 _ = f.Close()
291 return nil, fmt.Errorf("session content %s is not a regular file", ref.Digest)
292 }
293 if info.Size() != ref.Bytes {
294 _ = f.Close()
295 return nil, fmt.Errorf("session content %s size is %d, expected %d", ref.Digest, info.Size(), ref.Bytes)
296 }
297 return f, nil
298 }
299
300 func (s *Store) verifyOpenFile(ctx context.Context, f *os.File, ref Ref) error {
301 if ref.IndexDigest != "" {
302 if _, err := s.readIndex(ctx, ref); err != nil {
303 return err
304 }
305 }
306 if _, err := f.Seek(0, io.SeekStart); err != nil {
307 return err
308 }
309 digest := sha256.New()
310 if _, err := copyWithContext(ctx, digest, f); err != nil {
311 return fmt.Errorf("verify session content %s: %w", ref.Digest, err)
312 }
313 got := hex.EncodeToString(digest.Sum(nil))
314 if got != ref.Digest {
315 return fmt.Errorf("session content %s failed SHA-256 verification: got %s", ref.Digest, got)
316 }
317 return nil
318 }
319
320 func validateRef(ref Ref) error {
321 if ref.Bytes < 0 {
322 return errors.New("session content byte size must be non-negative")
323 }
324 if len(ref.Digest) != sha256.Size*2 {
325 return fmt.Errorf("invalid session content digest %q", ref.Digest)
326 }
327 for _, c := range ref.Digest {
328 if (c < '0' || c > '9') && (c < 'a' || c > 'f') {
329 return fmt.Errorf("invalid session content digest %q", ref.Digest)
330 }
331 }
332 if ref.IndexDigest != "" {
333 if len(ref.IndexDigest) != sha256.Size*2 || ref.IntegrityBlock != IntegrityBlockBytes {
334 return fmt.Errorf("invalid session content integrity index for %q", ref.Digest)
335 }
336 for _, c := range ref.IndexDigest {
337 if (c < '0' || c > '9') && (c < 'a' || c > 'f') {
338 return fmt.Errorf("invalid session content index digest %q", ref.IndexDigest)
339 }
340 }
341 }
342 return nil
343 }
344
345 func (s *Store) objectPath(digest string) string {
346 if len(digest) < 4 {
347 return filepath.Join(s.root, "objects", digest)
348 }
349 return filepath.Join(s.root, "objects", digest[:2], digest[2:4], digest)
350 }
351
352 func (s *Store) objectRelativePath(digest string) string {
353 if len(digest) < 4 {
354 return filepath.Join("objects", digest)
355 }
356 return filepath.Join("objects", digest[:2], digest[2:4], digest)
357 }
358
359 func (s *Store) indexPath(digest string) string {
360 if len(digest) < 4 {
361 return filepath.Join(s.root, "indexes", digest+".json")
362 }
363 return filepath.Join(s.root, "indexes", digest[:2], digest[2:4], digest+".json")
364 }
365
366 func (s *Store) indexRelativePath(digest string) string {
367 if len(digest) < 4 {
368 return filepath.Join("indexes", digest+".json")
369 }
370 return filepath.Join("indexes", digest[:2], digest[2:4], digest+".json")
371 }
372
373 type blockDigestWriter struct {
374 dst io.Writer
375 full hash.Hash
376 block hash.Hash
377 blockN int
378 blocks []string
379 }
380
381 func newBlockDigestWriter(dst io.Writer) *blockDigestWriter {
382 return &blockDigestWriter{dst: dst, full: sha256.New(), block: sha256.New()}
383 }
384
385 func (w *blockDigestWriter) Write(p []byte) (int, error) {
386 written := 0
387 for len(p) > 0 {
388 part := min(len(p), IntegrityBlockBytes-w.blockN)
389 chunk := p[:part]
390 n, err := w.dst.Write(chunk)
391 if n > 0 {
392 _, _ = w.full.Write(chunk[:n])
393 _, _ = w.block.Write(chunk[:n])
394 w.blockN += n
395 written += n
396 p = p[n:]
397 }
398 if err != nil {
399 return written, err
400 }
401 if n != part {
402 return written, io.ErrShortWrite
403 }
404 if w.blockN == IntegrityBlockBytes {
405 w.blocks = append(w.blocks, hex.EncodeToString(w.block.Sum(nil)))
406 w.block.Reset()
407 w.blockN = 0
408 }
409 }
410 return written, nil
411 }
412
413 func (w *blockDigestWriter) finish() []string {
414 if w.blockN > 0 {
415 w.blocks = append(w.blocks, hex.EncodeToString(w.block.Sum(nil)))
416 w.block.Reset()
417 w.blockN = 0
418 }
419 return append([]string(nil), w.blocks...)
420 }
421
422 func (s *Store) publishObject(ctx context.Context, tmpPath, dest string, ref Ref) error {
423 if err := os.Link(tmpPath, dest); err == nil {
424 _ = syncParent(filepath.Dir(dest))
425 return nil
426 } else if os.IsExist(err) {
427 if verifyErr := s.Verify(ctx, Ref{Digest: ref.Digest, Bytes: ref.Bytes}); verifyErr != nil {
428 return fmt.Errorf("existing session content %s is invalid: %w", ref.Digest, verifyErr)
429 }
430 return nil
431 }
432 release, err := filelock.Acquire(ctx, dest+".publish.lock")
433 if err != nil {
434 return fmt.Errorf("lock session content %s publication: %w", ref.Digest, err)
435 }
436 defer release()
437 if _, err := os.Stat(dest); err == nil {
438 return s.Verify(ctx, Ref{Digest: ref.Digest, Bytes: ref.Bytes})
439 } else if !os.IsNotExist(err) {
440 return err
441 }
442 if err := os.Rename(tmpPath, dest); err != nil {
443 return fmt.Errorf("publish session content %s: %w", ref.Digest, err)
444 }
445 return syncParent(filepath.Dir(dest))
446 }
447
448 func (s *Store) publishIndex(ctx context.Context, ref Ref, data []byte) error {
449 dest := s.indexPath(ref.Digest)
450 if err := os.MkdirAll(filepath.Dir(dest), 0o700); err != nil {
451 return err
452 }
453 tmp, err := os.CreateTemp(filepath.Join(s.root, ".tmp"), "index-*.tmp")
454 if err != nil {
455 return err
456 }
457 tmpPath := tmp.Name()
458 defer os.Remove(tmpPath)
459 if err := writeAll(tmp, data); err != nil {
460 _ = tmp.Close()
461 return err
462 }
463 if err := tmp.Sync(); err != nil {
464 _ = tmp.Close()
465 return err
466 }
467 if err := tmp.Close(); err != nil {
468 return err
469 }
470 release, err := filelock.Acquire(ctx, dest+".publish.lock")
471 if err != nil {
472 return err
473 }
474 defer release()
475 if existing, err := os.ReadFile(dest); err == nil {
476 if !bytes.Equal(existing, data) {
477 return fmt.Errorf("session content %s integrity index conflicts", ref.Digest)
478 }
479 return nil
480 } else if !os.IsNotExist(err) {
481 return err
482 }
483 if err := os.Rename(tmpPath, dest); err != nil {
484 return err
485 }
486 return syncParent(filepath.Dir(dest))
487 }
488
489 func (s *Store) readIndex(ctx context.Context, ref Ref) (integrityIndex, error) {
490 if err := ctx.Err(); err != nil {
491 return integrityIndex{}, err
492 }
493 if err := validateRef(ref); err != nil {
494 return integrityIndex{}, err
495 }
496 if s == nil || s.root == "" {
497 return integrityIndex{}, errors.New("session content store unavailable")
498 }
499 root, err := os.OpenRoot(s.root)
500 if err != nil {
501 return integrityIndex{}, fmt.Errorf("open session content root: %w", err)
502 }
503 defer root.Close()
504 data, err := root.ReadFile(s.indexRelativePath(ref.Digest))
505 if err != nil {
506 return integrityIndex{}, fmt.Errorf("read session content %s integrity index: %w", ref.Digest, err)
507 }
508 sum := sha256.Sum256(data)
509 if hex.EncodeToString(sum[:]) != ref.IndexDigest {
510 return integrityIndex{}, fmt.Errorf("session content %s integrity index failed SHA-256 verification", ref.Digest)
511 }
512 var index integrityIndex
513 if err := json.Unmarshal(data, &index); err != nil {
514 return integrityIndex{}, err
515 }
516 wantBlocks := int((ref.Bytes + IntegrityBlockBytes - 1) / IntegrityBlockBytes)
517 if index.Version != 1 || index.Object != ref.Digest || index.Bytes != ref.Bytes || index.BlockBytes != IntegrityBlockBytes || len(index.Blocks) != wantBlocks {
518 return integrityIndex{}, fmt.Errorf("session content %s integrity index metadata mismatch", ref.Digest)
519 }
520 return index, nil
521 }
522
523 func (s *Store) verifyRange(ctx context.Context, f *os.File, ref Ref, offset, length int64) error {
524 index, err := s.readIndex(ctx, ref)
525 if err != nil {
526 return err
527 }
528 if length == 0 {
529 return nil
530 }
531 first := offset / IntegrityBlockBytes
532 last := (offset + length - 1) / IntegrityBlockBytes
533 buf := make([]byte, IntegrityBlockBytes)
534 for block := first; block <= last; block++ {
535 if err := ctx.Err(); err != nil {
536 return err
537 }
538 start := block * IntegrityBlockBytes
539 size := min(IntegrityBlockBytes, ref.Bytes-start)
540 if _, err := f.ReadAt(buf[:size], start); err != nil {
541 return err
542 }
543 sum := sha256.Sum256(buf[:size])
544 if hex.EncodeToString(sum[:]) != index.Blocks[block] {
545 return fmt.Errorf("session content %s block %d failed SHA-256 verification", ref.Digest, block)
546 }
547 }
548 return nil
549 }
550
551 func writeAll(w io.Writer, data []byte) error {
552 for len(data) > 0 {
553 n, err := w.Write(data)
554 if err != nil {
555 return err
556 }
557 if n == 0 {
558 return io.ErrShortWrite
559 }
560 data = data[n:]
561 }
562 return nil
563 }
564
565 func copyWithContext(ctx context.Context, dst io.Writer, src io.Reader) (int64, error) {
566 buf := make([]byte, copyBufferBytes)
567 var written int64
568 for {
569 if err := ctx.Err(); err != nil {
570 return written, err
571 }
572 n, readErr := src.Read(buf)
573 if n > 0 {
574 wn, writeErr := dst.Write(buf[:n])
575 written += int64(wn)
576 if writeErr != nil {
577 return written, writeErr
578 }
579 if wn != n {
580 return written, io.ErrShortWrite
581 }
582 }
583 if readErr != nil {
584 if errors.Is(readErr, io.EOF) {
585 return written, nil
586 }
587 return written, readErr
588 }
589 }
590 }
591
592 func syncParent(path string) error {
593 f, err := os.Open(path)
594 if err != nil {
595 return err
596 }
597 defer f.Close()
598 if err := f.Sync(); err != nil && runtime.GOOS != "windows" && !errors.Is(err, syscall.EINVAL) && !errors.Is(err, syscall.ENOTSUP) && !errors.Is(err, syscall.ENOSYS) {
599 return err
600 }
601 return nil
602 }
603
603 lines GO