diff --git a/sdk/experimental/tdf/writer.go b/sdk/experimental/tdf/writer.go index 02d2f8af53..2a58af6c87 100644 --- a/sdk/experimental/tdf/writer.go +++ b/sdk/experimental/tdf/writer.go @@ -341,7 +341,10 @@ func (w *Writer) WriteSegment(ctx context.Context, index int, data []byte) (*Seg // // Error conditions: // - ErrAlreadyFinalized: Finalize already called -// - Missing segments: Gaps in segment indices (e.g., segments 0,1,3 written but 2 missing) +// - Missing segment 0: Index 0 carries the payload's ZIP local file header, which +// every recorded offset is measured from, so a write set that omits it is +// rejected. Gaps between the remaining indices are legal (e.g., segments 0,1,3 +// with 2 missing); order is inferred by sorting whichever indices are present. // - Key splitting failures: Invalid attributes or KAS configuration // - Manifest generation errors: JSON marshaling failures // - Archive finalization errors: ZIP structure generation failures diff --git a/sdk/internal/zipstream/segment_writer.go b/sdk/internal/zipstream/segment_writer.go index b6f1d7191f..18b35b62ce 100644 --- a/sdk/internal/zipstream/segment_writer.go +++ b/sdk/internal/zipstream/segment_writer.go @@ -126,6 +126,29 @@ func (sw *segmentWriter) Finalize(ctx context.Context, manifest []byte) ([]byte, default: } + // Nothing arrived at all: report the general incomplete-input error + // rather than the segment-0-specific one below. + if len(sw.metadata.Segments) == 0 { + return nil, &Error{Op: "finalize", Type: "segment", Err: ErrSegmentMissing} + } + + // Only segment 0 emits the payload's local file header, and every offset + // recorded below is measured from it: without it the manifest entry, the + // central directory, and the EOCD all overshoot by headerSize. The result + // is a corrupt archive rather than a clean failure -- under Zip64Always + // the trailer points past the end of the buffer, while under Zip64Auto it + // lands mid-archive, where some readers parse the manifest happily and + // only choke on the payload. + // + // This has to run before the order derivation below: Order is derived + // once and kept, so a caller that supplies segment 0 and retries must not + // inherit an order that already excluded it. IsComplete cannot catch the + // absence either, since that derived order is self-consistent by + // construction. + if _, ok := sw.metadata.Segments[0]; !ok { + return nil, &Error{Op: "finalize", Type: "segment", Err: ErrNoSegmentZero} + } + // If no explicit order was provided, derive order from present indices (sorted). if len(sw.metadata.Order) == 0 { order := make([]int, 0, len(sw.metadata.Segments)) @@ -139,7 +162,10 @@ func (sw *segmentWriter) Finalize(ctx context.Context, manifest []byte) ([]byte, } } - // Verify all segments are present + // Verify all segments are present. Unreachable with an order derived + // above -- that order is built from the present indices, so it is + // complete by construction, and the empty set already returned. Kept for + // a future caller that supplies an explicit order. if !sw.metadata.IsComplete() { return nil, &Error{Op: "finalize", Type: "segment", Err: ErrSegmentMissing} } @@ -222,21 +248,29 @@ func (sw *segmentWriter) Finalize(ctx context.Context, manifest []byte) ([]byte, return buffer.Bytes(), nil } -// CleanupSegment removes the presence marker for a segment index. Since payload -// bytes are not retained, this only affects metadata tracking. Calling this -// before Finalize will cause IsComplete() to fail for that index. +// CleanupSegment implements SegmentWriter. Payload bytes are never retained, so +// this only rolls back the metadata and size accounting the segment contributed. func (sw *segmentWriter) CleanupSegment(index int) error { sw.mu.Lock() defer sw.mu.Unlock() - // Remove segment from unprocessed map (no-op if already processed or not found) - if _, ok := sw.metadata.Segments[index]; ok { - delete(sw.metadata.Segments, index) - if sw.metadata.presentCount > 0 { - sw.metadata.presentCount-- - } + // No-op if the index was never written or was already cleaned up. + seg, ok := sw.metadata.Segments[index] + if !ok { + return nil } + delete(sw.metadata.Segments, index) + sw.metadata.presentCount-- + + // Undo everything the segment contributed, so that a cleaned-up index is + // indistinguishable from one that was never written. Leaving the sizes + // behind would make Finalize describe a payload larger than the one the + // caller can assemble, and the offsets it records would overshoot. + sw.metadata.TotalSize -= seg.Size + sw.payloadEntry.Size -= seg.Size + sw.payloadEntry.CompressedSize -= seg.Size + return nil } diff --git a/sdk/internal/zipstream/segment_writer_test.go b/sdk/internal/zipstream/segment_writer_test.go index c57830d858..a4319e66f5 100644 --- a/sdk/internal/zipstream/segment_writer_test.go +++ b/sdk/internal/zipstream/segment_writer_test.go @@ -362,6 +362,178 @@ func TestSegmentWriter_AllowsGapsOnFinalize(t *testing.T) { writer.Close() } +func TestSegmentWriter_FinalizeRequiresSegmentZero(t *testing.T) { + // Gaps are fine, but the set has to start at 0: only segment 0 emits + // the payload local file header, and Finalize sizes every offset it + // records as though that header were at the front of the stream. + writer := NewSegmentTDFWriter(1) + ctx := t.Context() + + _, err := writer.WriteSegment(ctx, 1, 5, crc32.ChecksumIEEE([]byte("first"))) + require.NoError(t, err) + + _, err = writer.WriteSegment(ctx, 2, 6, crc32.ChecksumIEEE([]byte("second"))) + require.NoError(t, err) + + _, err = writer.Finalize(ctx, []byte("manifest")) + require.ErrorIs(t, err, ErrNoSegmentZero) + + writer.Close() +} + +func TestSegmentWriter_CleanupSegmentZeroBlocksFinalize(t *testing.T) { + // Dropping segment 0 after the fact is invisible to IsComplete -- + // the inferred order becomes [1], which is internally consistent -- + // so the header check has to stand on its own. + writer := NewSegmentTDFWriter(2) + ctx := t.Context() + + _, err := writer.WriteSegment(ctx, 0, 5, crc32.ChecksumIEEE([]byte("first"))) + require.NoError(t, err) + + _, err = writer.WriteSegment(ctx, 1, 6, crc32.ChecksumIEEE([]byte("second"))) + require.NoError(t, err) + + require.NoError(t, writer.CleanupSegment(0)) + + _, err = writer.Finalize(ctx, []byte("manifest")) + require.ErrorIs(t, err, ErrNoSegmentZero) + require.NotErrorIs(t, err, ErrSegmentMissing, "the remaining segments are complete; only the header is gone") + + writer.Close() +} + +func TestSegmentWriter_FinalizeAfterSupplyingSegmentZero(t *testing.T) { + // ErrNoSegmentZero invites the caller to write segment 0 and try again, + // so the retry has to produce a correct archive. Finalize derives the + // segment order once and keeps it: if the check ran after that + // derivation, this second Finalize would succeed against an order that + // still omitted 0, combining a CRC over two of the three segments. + writer := NewSegmentTDFWriter(3) + ctx := t.Context() + + segments := [][]byte{[]byte("first"), []byte("second"), []byte("third")} + + for _, index := range []int{1, 2} { + data := segments[index] + _, err := writer.WriteSegment(ctx, index, uint64(len(data)), crc32.ChecksumIEEE(data)) + require.NoError(t, err) + } + + _, err := writer.Finalize(ctx, []byte("manifest")) + require.ErrorIs(t, err, ErrNoSegmentZero) + + headerBytes, err := writer.WriteSegment(ctx, 0, uint64(len(segments[0])), crc32.ChecksumIEEE(segments[0])) + require.NoError(t, err, "the writer stays usable after ErrNoSegmentZero") + require.NotEmpty(t, headerBytes, "segment 0 carries the payload local file header") + + var archive []byte + archive = append(archive, headerBytes...) + for _, data := range segments { + archive = append(archive, data...) + } + + finalBytes, err := writer.Finalize(ctx, []byte("manifest")) + require.NoError(t, err, "Finalize should succeed once segment 0 arrives") + archive = append(archive, finalBytes...) + + zipReader, err := zip.NewReader(bytes.NewReader(archive), int64(len(archive))) + require.NoError(t, err, "retry should produce a readable ZIP") + + payloadFile := findFileByName(zipReader, TDFPayloadFileName) + require.NotNil(t, payloadFile) + + payloadReader, err := payloadFile.Open() + require.NoError(t, err) + defer payloadReader.Close() + + // Reading through archive/zip validates the recorded CRC against the + // bytes actually present, which is what a stale order would break. + content, err := io.ReadAll(payloadReader) + require.NoError(t, err, "payload CRC must cover every segment, including 0") + assert.Equal(t, bytes.Join(segments, nil), content) + + writer.Close() +} + +func TestSegmentWriter_FinalizeWithoutAnySegments(t *testing.T) { + // No segments at all is incomplete input, not a missing-header problem: + // the general error stays reachable and keeps its distinct meaning. + writer := NewSegmentTDFWriter(2) + + _, err := writer.Finalize(t.Context(), []byte("manifest")) + require.ErrorIs(t, err, ErrSegmentMissing) + require.NotErrorIs(t, err, ErrNoSegmentZero) + + writer.Close() +} + +func TestSegmentWriter_CleanupOnlySegmentZero(t *testing.T) { + // Cleaning up the last remaining segment leaves nothing to finalize, so + // this reports the general incomplete-input error rather than the + // segment-0-specific one -- "nothing to assemble" is the more useful + // diagnosis than "the header is gone". + writer := NewSegmentTDFWriter(1) + ctx := t.Context() + + _, err := writer.WriteSegment(ctx, 0, 5, crc32.ChecksumIEEE([]byte("first"))) + require.NoError(t, err) + + require.NoError(t, writer.CleanupSegment(0)) + + _, err = writer.Finalize(ctx, []byte("manifest")) + require.ErrorIs(t, err, ErrSegmentMissing) + require.NotErrorIs(t, err, ErrNoSegmentZero) + + writer.Close() +} + +func TestSegmentWriter_CleanupSegmentRollsBackSizeAccounting(t *testing.T) { + // A cleaned-up index has to become indistinguishable from one that was + // never written, which sparse write sets already allow. Without the + // rollback the recorded sizes still cover the removed segment while the + // CRC covers only the survivors, and Finalize emits a trailer describing + // a payload the caller cannot produce. + writer := NewSegmentTDFWriter(3) + ctx := t.Context() + + segments := [][]byte{[]byte("first"), []byte("second"), []byte("third")} + + var archive []byte + for index, data := range segments { + headerBytes, err := writer.WriteSegment(ctx, index, uint64(len(data)), crc32.ChecksumIEEE(data)) + require.NoError(t, err) + archive = append(archive, headerBytes...) + if index != 1 { + archive = append(archive, data...) + } + } + + require.NoError(t, writer.CleanupSegment(1)) + + finalBytes, err := writer.Finalize(ctx, []byte("manifest")) + require.NoError(t, err) + archive = append(archive, finalBytes...) + + zipReader, err := zip.NewReader(bytes.NewReader(archive), int64(len(archive))) + require.NoError(t, err, "offsets must describe the payload the caller actually assembled") + + payloadFile := findFileByName(zipReader, TDFPayloadFileName) + require.NotNil(t, payloadFile) + + payloadReader, err := payloadFile.Open() + require.NoError(t, err) + defer payloadReader.Close() + + // Reading through archive/zip validates the recorded CRC against the + // bytes present, which is what stale size accounting would break. + content, err := io.ReadAll(payloadReader) + require.NoError(t, err, "payload CRC must cover exactly the surviving segments") + assert.Equal(t, []byte("firstthird"), content) + + writer.Close() +} + func TestSegmentWriter_CleanupSegment(t *testing.T) { // Test memory cleanup functionality writer := NewSegmentTDFWriter(3) diff --git a/sdk/internal/zipstream/writer.go b/sdk/internal/zipstream/writer.go index 11eea8a5b5..774e1c9488 100644 --- a/sdk/internal/zipstream/writer.go +++ b/sdk/internal/zipstream/writer.go @@ -24,9 +24,22 @@ type Writer interface { type SegmentWriter interface { Writer WriteSegment(ctx context.Context, index int, size uint64, crc32 uint32) ([]byte, error) + // Finalize writes the trailer (data descriptor, manifest, central + // directory) for the segments recorded so far. Segment 0 must be among + // them: it carries the payload local file header that every recorded + // offset is measured from. Finalize returns ErrNoSegmentZero when index 0 + // was never written or was cleaned up, and ErrSegmentMissing when no + // segments remain at all -- none were written, or every one was cleaned + // up. Gaps between the remaining indices are accepted; order is inferred + // by sorting whichever indices are present. Finalize(ctx context.Context, manifest []byte) ([]byte, error) - // CleanupSegment removes the presence marker for a segment index. - // Calling this before Finalize will cause IsComplete() to fail for that index. + // CleanupSegment drops a segment index, rolling back both its presence + // marker and the payload size it contributed. A cleaned-up index becomes + // indistinguishable from one that was never written: it falls out of the + // order Finalize infers, exactly as a gap in the write set would, and the + // caller must leave its bytes out of the assembled archive. Index 0 is the + // exception -- Finalize rejects its absence with ErrNoSegmentZero. + // Cleaning up an index that was never written is a no-op. CleanupSegment(index int) error } @@ -52,6 +65,7 @@ var ( ErrOutOfOrder = errors.New("segment out of order") ErrDuplicateSegment = errors.New("duplicate segment already written") ErrSegmentMissing = errors.New("segment missing") + ErrNoSegmentZero = errors.New("segment 0 missing; it carries the payload local file header") ErrInvalidSize = errors.New("invalid size") ErrZip64Required = errors.New("ZIP64 required but disabled (Zip64Never)") )