diff --git a/test/tus/tus_integration_test.go b/test/tus/tus_integration_test.go index 96f4bd505..0576caa98 100644 --- a/test/tus/tus_integration_test.go +++ b/test/tus/tus_integration_test.go @@ -6,6 +6,7 @@ import ( "encoding/base64" "fmt" "io" + "net" "net/http" "os" "os/exec" @@ -902,3 +903,96 @@ func TestTusResumeAfterInterruption(t *testing.T) { require.NoError(t, err) assert.Equal(t, testData, body, "Resumed upload should produce complete file") } + +// TestTusAbortedPatchKeepsStoredChunks checks that a PATCH cut off mid-body +// leaves the sub-chunks it already stored in place. The filer splits a PATCH +// into 4MB sub-chunks and records each one as it lands; the offset a resuming +// client reads back covers them, so their data has to survive the failure. +func TestTusAbortedPatchKeepsStoredChunks(t *testing.T) { + if testing.Short() { + t.Skip("Skipping integration test in short mode") + } + + ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second) + defer cancel() + + cluster, err := startTestCluster(t, ctx) + require.NoError(t, err) + defer func() { + cluster.Stop() + os.RemoveAll(cluster.dataDir) + }() + + const subChunkSize = 4 * 1024 * 1024 + testData := make([]byte, 3*subChunkSize) + for i := range testData { + testData[i] = byte(i % 251) + } + targetPath := "/aborted/interrupted.bin" + client := &http.Client{} + + createReq, err := http.NewRequest(http.MethodPost, cluster.TusURL()+targetPath, nil) + require.NoError(t, err) + createReq.Header.Set("Tus-Resumable", TusVersion) + createReq.Header.Set("Upload-Length", strconv.Itoa(len(testData))) + + createResp, err := client.Do(createReq) + require.NoError(t, err) + createResp.Body.Close() + require.Equal(t, http.StatusCreated, createResp.StatusCode) + uploadLocation := createResp.Header.Get("Location") + + // Promise the whole body, then reset the connection while a later + // sub-chunk is still being read. + conn, err := net.Dial("tcp", "127.0.0.1:"+testFilerPort) + require.NoError(t, err) + _, err = fmt.Fprintf(conn, "PATCH %s HTTP/1.1\r\nHost: 127.0.0.1:%s\r\nTus-Resumable: %s\r\nContent-Type: application/offset+octet-stream\r\nUpload-Offset: 0\r\nContent-Length: %d\r\n\r\n", + uploadLocation, testFilerPort, TusVersion, len(testData)) + require.NoError(t, err) + _, err = conn.Write(testData[:subChunkSize+1024*1024]) + require.NoError(t, err) + time.Sleep(3 * time.Second) + require.NoError(t, conn.(*net.TCPConn).SetLinger(0)) + require.NoError(t, conn.Close()) + t.Log("PATCH connection reset mid-body") + + time.Sleep(5 * time.Second) + + headReq, err := http.NewRequest(http.MethodHead, cluster.FullURL(uploadLocation), nil) + require.NoError(t, err) + headReq.Header.Set("Tus-Resumable", TusVersion) + + headResp, err := client.Do(headReq) + require.NoError(t, err) + headResp.Body.Close() + require.Equal(t, http.StatusOK, headResp.StatusCode) + currentOffset, err := strconv.Atoi(headResp.Header.Get("Upload-Offset")) + require.NoError(t, err) + require.Equal(t, subChunkSize, currentOffset, "the sub-chunk stored before the reset should count towards the offset") + + patchReq, err := http.NewRequest(http.MethodPatch, cluster.FullURL(uploadLocation), bytes.NewReader(testData[currentOffset:])) + require.NoError(t, err) + patchReq.Header.Set("Tus-Resumable", TusVersion) + patchReq.Header.Set("Upload-Offset", strconv.Itoa(currentOffset)) + patchReq.Header.Set("Content-Type", "application/offset+octet-stream") + + patchResp, err := client.Do(patchReq) + require.NoError(t, err) + patchResp.Body.Close() + require.Equal(t, http.StatusNoContent, patchResp.StatusCode) + + // A vacuum reclaims whatever the filer deleted, so the file survives this + // only if the chunks the session kept are still stored. + vacuumResp, err := client.Get(fmt.Sprintf("http://127.0.0.1:%s/vol/vacuum?garbageThreshold=0.001", testMasterPort)) + require.NoError(t, err) + vacuumResp.Body.Close() + require.Equal(t, http.StatusOK, vacuumResp.StatusCode) + + getResp, err := client.Get(cluster.FilerURL() + targetPath) + require.NoError(t, err) + defer getResp.Body.Close() + require.Equal(t, http.StatusOK, getResp.StatusCode) + body, err := io.ReadAll(getResp.Body) + require.NoError(t, err) + assert.Equal(t, testData, body, "the resumed upload should read back whole") +} diff --git a/weed/server/filer_server_tus_handlers.go b/weed/server/filer_server_tus_handlers.go index 598e45e0b..83f6e6d5c 100644 --- a/weed/server/filer_server_tus_handlers.go +++ b/weed/server/filer_server_tus_handlers.go @@ -632,7 +632,6 @@ func (fs *FilerServer) tusWriteData(ctx context.Context, session *TusSession, of // Upload in streaming chunks to avoid buffering entire content in memory var totalWritten int64 var uploadErr error - var uploadedChunks []*TusChunkInfo // Create one uploader for all sub-chunks to reuse HTTP client connections uploader, uploaderErr := operation.NewUploader() @@ -692,34 +691,19 @@ func (fs *FilerServer) tusWriteData(ctx context.Context, session *TusSession, of } if saveErr := fs.saveTusChunk(ctx, session.ID, chunk); saveErr != nil { - // Cleanup this chunk on failure - fs.filer.DeleteChunks(ctx, util.FullPath(session.TargetPath), []*filer_pb.FileChunk{ - {FileId: fileId}, - }) + fs.deleteTusChunk(ctx, session, chunk) uploadErr = fmt.Errorf("update session: %w", saveErr) break } - uploadedChunks = append(uploadedChunks, chunk) - totalWritten += int64(uploadResult.Size) currentOffset += int64(uploadResult.Size) stats.FilerHandlerCounter.WithLabelValues("tusUploadChunk").Inc() } - if uploadErr != nil { - // Cleanup all uploaded chunks on error - if len(uploadedChunks) > 0 { - var chunksToDelete []*filer_pb.FileChunk - for _, c := range uploadedChunks { - chunksToDelete = append(chunksToDelete, &filer_pb.FileChunk{FileId: c.FileId}) - } - fs.filer.DeleteChunks(ctx, util.FullPath(session.TargetPath), chunksToDelete) - } - return 0, uploadErr - } - - return totalWritten, nil + // Sub-chunks already recorded stay: the session offset a resuming client + // reads back covers them, and the completed entry is assembled from them. + return totalWritten, uploadErr } // parseTusMetadata parses the Upload-Metadata header diff --git a/weed/server/filer_server_tus_session.go b/weed/server/filer_server_tus_session.go index e530ed7d8..2e6a6d4d5 100644 --- a/weed/server/filer_server_tus_session.go +++ b/weed/server/filer_server_tus_session.go @@ -417,6 +417,18 @@ func (fs *FilerServer) saveTusChunk(ctx context.Context, uploadID string, chunk return nil } +// deleteTusChunk drops a chunk's record before freeing its data, so a save that +// failed after the record landed cannot leave the session pointing at a needle +// that is about to be deleted. Data outliving a lost record only leaks. +func (fs *FilerServer) deleteTusChunk(ctx context.Context, session *TusSession, chunk *TusChunkInfo) { + chunkPath := util.FullPath(fs.tusChunkPath(session.ID, chunk.Offset, chunk.Size, chunk.FileId)) + if err := fs.filer.DeleteEntryMetaAndData(ctx, chunkPath, false, false, false, false, nil, 0); err != nil && !errors.Is(err, filer_pb.ErrNotFound) { + glog.Errorf("TUS chunk %s record kept, its data leaks: %v", chunkPath, err) + return + } + fs.filer.DeleteChunks(ctx, util.FullPath(session.TargetPath), []*filer_pb.FileChunk{{FileId: chunk.FileId}}) +} + // deleteTusSession removes a TUS upload session and all its data func (fs *FilerServer) deleteTusSession(ctx context.Context, uploadID string) error { sessionPath := util.FullPath(fs.tusSessionPath(uploadID))