Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion sdk/experimental/tdf/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
54 changes: 44 additions & 10 deletions sdk/internal/zipstream/segment_writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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}
}
Comment thread
dmihalcik-virtru marked this conversation as resolved.

// 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))
Expand All @@ -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}
}
Expand Down Expand Up @@ -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
}

Expand Down
172 changes: 172 additions & 0 deletions sdk/internal/zipstream/segment_writer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
18 changes: 16 additions & 2 deletions sdk/internal/zipstream/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand All @@ -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)")
)
Expand Down
Loading