| 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 |