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
29 changes: 27 additions & 2 deletions internal/nzbfilesystem/metadata_remote_file.go
Original file line number Diff line number Diff line change
Expand Up @@ -1075,6 +1075,24 @@ func (mvf *MetadataVirtualFile) classifyReadError(readErr error) error {
return readErr
}

// truncatedTailError describes a file whose advertised size extends past the
// bytes its articles actually hold: a reader built at `at` produced nothing
// and rebuilding it cannot help. The bytes from `at` to FileSize are gone for
// good, so this is a permanent data corruption (NoRetry) that the health
// pipeline must record and hand to repair. Without that verdict every client
// that wants the tail (MKV cues live there) re-fetches the final article on
// each attempt and retries indefinitely — the 15-requests-per-second storm
// seen in the field. The error still unwraps to io.ErrUnexpectedEOF.
func (mvf *MetadataVirtualFile) truncatedTailError(at int64) error {
return &usenet.DataCorruptionError{
UnderlyingErr: fmt.Errorf("%w: reader ended at offset %d before the requested end (file advertises %d bytes)",
io.ErrUnexpectedEOF, at, mvf.meta.FileSize),
BytesRead: at,
NoRetry: true,
FileOffset: at,
}
}

// segmentOffsetIndex provides O(1) lookup for offset→segment mapping using binary search
type segmentOffsetIndex struct {
offsets []int64 // Cumulative start offset of each segment in file coordinates
Expand Down Expand Up @@ -1212,7 +1230,7 @@ func (mvf *MetadataVirtualFile) Read(p []byte) (n int, err error) {
// the range end; rotating once more would spin here holding mvf.mu.
if totalRead == 0 && mvf.position == stalledAt {
mvf.closeCurrentReader()
return n, fmt.Errorf("%w: reader ended at offset %d before the requested end", io.ErrUnexpectedEOF, mvf.position)
return n, mvf.classifyReadError(mvf.truncatedTailError(mvf.position))
}
stalledAt = mvf.position
// Close current reader and try to get a new one for the next range in next iteration
Expand Down Expand Up @@ -1363,7 +1381,7 @@ func (mvf *MetadataVirtualFile) ReadAtContext(readCtx context.Context, p []byte,
// same offset cannot reach the range end.
at := off + int64(n)
if rn == 0 && at == stalledAt {
sharedErr = fmt.Errorf("%w: reader ended at offset %d before the requested end", io.ErrUnexpectedEOF, at)
sharedErr = mvf.truncatedTailError(at)
break
}
stalledAt = at
Expand Down Expand Up @@ -1452,6 +1470,13 @@ ephemeral:
if err == io.ErrUnexpectedEOF {
err = nil
}
// The window lies inside the advertised file, yet a fresh reader at its
// start had nothing at all: the data behind this offset does not exist
// (a truncated final article). Report it as corruption so the client
// gets a definitive answer instead of a short read it will retry forever.
if n == 0 && errors.Is(err, io.EOF) && off < mvf.meta.FileSize {
err = mvf.truncatedTailError(off)
}

// Only update the shared cursor when the shared reader was torn down.
// If it is still alive, readAtSharedNext already points to the reader's
Expand Down
142 changes: 142 additions & 0 deletions internal/nzbfilesystem/truncated_tail_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
package nzbfilesystem

import (
"context"
"errors"
"fmt"
"io"
"runtime"
"testing"

"github.com/kipsilabs/altmount/internal/config"
"github.com/kipsilabs/altmount/internal/database"
metapb "github.com/kipsilabs/altmount/internal/metadata/proto"
"github.com/kipsilabs/altmount/internal/testsupport/fakepool"
"github.com/kipsilabs/altmount/internal/testsupport/segments"
"github.com/kipsilabs/altmount/internal/utils"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// newTruncatedTailPool serves n segments where the final article carries
// `short` bytes fewer than its SegmentSize claims — a truncated last part.
// The metadata still advertises the full n*segSize, so the last `short`
// bytes of the file can never be produced.
func newTruncatedTailPool(n, segSize, short int) *fakepool.Client {
fp := fakepool.New()
all := segments.FileBytes(n, segSize)
for i := range n {
b := all[i*segSize : (i+1)*segSize]
if i == n-1 {
b = b[:segSize-short]
}
fp.SetBehavior(segments.MessageID(i), fakepool.SegmentBehavior{Bytes: b})
}
return fp
}

func rangeCtx(start, end int64) context.Context {
return context.WithValue(context.Background(), utils.RangeKey, fmt.Sprintf("bytes=%d-%d", start, end))
}

// A client asking for the unreachable tail (the retry loop seen in the field:
// the same 4760-byte range every 60 ms) must get a definitive corruption
// error, not a bare unexpected EOF that reads as a transient failure.
func TestTruncatedFinalArticleReportsCorruption(t *testing.T) {
const n, segSize, short = 8, 64 << 10, 4760
fileSize := int64(n * segSize)
realEnd := fileSize - short

for _, tc := range []struct {
name string
start int64
}{
{"range starts where the data ends", realEnd},
{"range starts before the data ends", realEnd - 2000},
} {
t.Run(tc.name, func(t *testing.T) {
fp := newTruncatedTailPool(n, segSize, short)
mvf := newTestMVF(t, rangeCtx(tc.start, fileSize-1), fp, n, segSize, 4)
mvf.originalRangeEnd = 0 // parse the Range header on first read

_, err := mvf.Seek(tc.start, io.SeekStart)
require.NoError(t, err)

got, err := readAllOrHang(t, mvf)
require.Error(t, err)
var corrupted *CorruptedFileError
assert.True(t, errors.As(err, &corrupted), "want *CorruptedFileError, got %T: %v", err, err)
assert.True(t, errors.Is(err, io.ErrUnexpectedEOF), "must still unwrap to io.ErrUnexpectedEOF: %v", err)
assert.Equal(t, int(realEnd-tc.start), len(got), "everything before the truncation point is delivered")
})
}
}

// ReadAt takes the shared-cursor path; it must reach the same verdict.
func TestTruncatedFinalArticleReportsCorruptionOnReadAt(t *testing.T) {
const n, segSize, short = 8, 64 << 10, 4760
fileSize := int64(n * segSize)
realEnd := fileSize - short

fp := newTruncatedTailPool(n, segSize, short)
mvf := newTestMVF(t, context.Background(), fp, n, segSize, 4)

buf := make([]byte, short)
_, err := mvf.ReadAtContext(context.Background(), buf, realEnd)
require.Error(t, err)
var corrupted *CorruptedFileError
assert.True(t, errors.As(err, &corrupted), "want *CorruptedFileError, got %T: %v", err, err)
assert.True(t, errors.Is(err, io.ErrUnexpectedEOF), "must still unwrap to io.ErrUnexpectedEOF: %v", err)
}

// With the health system wired, the truncation is recorded and the file
// handed to repair like any other corruption, so the metadata is moved and
// the next client request stops at a 404 instead of re-fetching the article.
func TestTruncatedFinalArticleTriggersRepair(t *testing.T) {
if runtime.GOOS == "windows" {
t.Skip("symlinks not supported on Windows")
}
const n, segSize, short = 4, 64 << 10, 4760
fileSize := int64(n * segSize)
realEnd := fileSize - short

repo, db, ms := setupStreamHealthEnv(t)
ctx := context.Background()
filePath := "series/truncated.s01e01.mkv"

fp := newTruncatedTailPool(n, segSize, short)
mvf := newTestMVF(t, rangeCtx(realEnd, fileSize-1), fp, n, segSize, 4)
mvf.originalRangeEnd = 0
mvf.name = filePath

meta := ms.CreateFileMetadata(fileSize, "test.nzb", metapb.FileStatus_FILE_STATUS_HEALTHY,
mvf.meta.SegmentData, metapb.Encryption_NONE, "", "", nil, nil, 0, nil, "")
require.NoError(t, ms.WriteFileMetadata(filePath, meta))
_, err := db.Exec(
`INSERT INTO file_health (file_path, library_path, status, scheduled_check_at) VALUES (?, ?, 'healthy', datetime('now'))`,
filePath, "/media/library/truncated.s01e01.mkv",
)
require.NoError(t, err)

enabled := true
cfg := config.DefaultConfig()
cfg.Health.Enabled = &enabled
cfg.MountPath = ""
mvf.metadataService = ms
mvf.healthRepository = repo
mvf.configGetter = func() *config.Config { return cfg }

_, err = mvf.Seek(realEnd, io.SeekStart)
require.NoError(t, err)
_, err = readAllOrHang(t, mvf)
require.Error(t, err)

fh, err := repo.GetFileHealth(ctx, filePath)
require.NoError(t, err)
require.NotNil(t, fh)
assert.Equal(t, database.HealthStatusRepairTriggered, fh.Status)

moved, err := ms.ReadFileMetadata(filePath)
require.NoError(t, err)
assert.Nil(t, moved, "metadata must be moved to the corrupted folder")
}
Loading