diff --git a/cmd/internal/client/client.go b/cmd/internal/client/client.go index f5eb6ea4..a8b02e7a 100644 --- a/cmd/internal/client/client.go +++ b/cmd/internal/client/client.go @@ -236,19 +236,27 @@ type ObjectDraft struct { Upload *ObjectUploadInstructions `json:"upload,omitempty"` } -// ObjectUploadInstructions is returned by CreateObject for a file draft: the -// server-decided part size and one presigned PUT URL per slice (1 URL = single -// PutObject, N URLs = multipart). The client PUTs each slice, reads the ETag, and -// posts them to CompleteObjectUpload. +// ObjectUploadInstructions is returned by CreateObject for a file draft. The +// client PUTs each explicit part descriptor, reads the ETag, and posts the +// partNumber+etag records to CompleteObjectUpload. type ObjectUploadInstructions struct { - SessionID string `json:"sessionId"` - PartSize int64 `json:"partSize"` - URLs []string `json:"urls"` + SessionID string `json:"sessionId"` + UploadID *string `json:"uploadId"` + Mode string `json:"mode"` + PartSize int64 `json:"partSize"` + PartCount int `json:"partCount"` + ExpiresAt string `json:"expiresAt"` + PresignedExpiresAt string `json:"presignedExpiresAt"` + RequiredHeaders map[string]string `json:"requiredHeaders"` + URLs []string `json:"urls"` + Parts []PresignedObjectUploadPart `json:"parts"` } type PresignedObjectUploadPart struct { - PartNumber int `json:"partNumber"` - URL string `json:"url"` + PartNumber int `json:"partNumber"` + URL string `json:"url"` + ExpiresAt string `json:"expiresAt"` + Headers map[string]string `json:"headers"` } type CompletedObjectUploadPart struct { @@ -612,9 +620,24 @@ func (c *Client) createMatter( draft := ObjectDraft{ID: res.JSON201.Id, Name: res.JSON201.Name} if u := res.JSON201.Upload; u != nil { draft.Upload = &ObjectUploadInstructions{ - SessionID: u.SessionId, - PartSize: int64(u.PartSize), - URLs: u.Urls, + SessionID: u.SessionId, + UploadID: u.UploadId, + Mode: string(u.Mode), + PartSize: int64(u.PartSize), + PartCount: u.PartCount, + ExpiresAt: u.ExpiresAt, + PresignedExpiresAt: u.PresignedExpiresAt, + RequiredHeaders: u.RequiredHeaders, + URLs: append([]string(nil), u.Urls...), + Parts: make([]PresignedObjectUploadPart, 0, len(u.Parts)), + } + for _, part := range u.Parts { + draft.Upload.Parts = append(draft.Upload.Parts, PresignedObjectUploadPart{ + PartNumber: part.PartNumber, + URL: part.Url, + ExpiresAt: part.ExpiresAt, + Headers: part.Headers, + }) } } return draft, nil @@ -638,7 +661,12 @@ func (c *Client) PresignObjectUploadParts(ctx context.Context, token string, id } parts := make([]PresignedObjectUploadPart, 0, len(res.JSON200.Parts)) for _, part := range res.JSON200.Parts { - parts = append(parts, PresignedObjectUploadPart{PartNumber: part.PartNumber, URL: part.Url}) + parts = append(parts, PresignedObjectUploadPart{ + PartNumber: part.PartNumber, + URL: part.Url, + ExpiresAt: part.ExpiresAt, + Headers: part.Headers, + }) } return parts, nil } diff --git a/cmd/internal/client/client_test.go b/cmd/internal/client/client_test.go index d9b20020..b4bf90bf 100644 --- a/cmd/internal/client/client_test.go +++ b/cmd/internal/client/client_test.go @@ -547,9 +547,23 @@ func TestCreateObjectMapsUploadInstructions(t *testing.T) { "id": "object-1", "name": "movie.mkv", "upload": map[string]any{ - "sessionId": "session-1", - "partSize": 1024, - "urls": []string{"https://s3/part-1"}, + "sessionId": "session-1", + "uploadId": nil, + "mode": "single", + "partSize": 1024, + "partCount": 1, + "expiresAt": "2026-01-01T01:00:00Z", + "presignedExpiresAt": "2026-01-01T00:15:00Z", + "requiredHeaders": map[string]string{}, + "urls": []string{"https://s3/part-1"}, + "parts": []map[string]any{ + { + "partNumber": 1, + "url": "https://s3/part-1", + "expiresAt": "2026-01-01T00:15:00Z", + "headers": map[string]string{}, + }, + }, }, }) })) @@ -562,9 +576,16 @@ func TestCreateObjectMapsUploadInstructions(t *testing.T) { if draft.Upload == nil { t.Fatalf("expected upload instructions: %#v", draft) } - if draft.Upload.SessionID != "session-1" || draft.Upload.PartSize != 1024 || !reflect.DeepEqual(draft.Upload.URLs, []string{"https://s3/part-1"}) { + if draft.Upload.SessionID != "session-1" || draft.Upload.PartSize != 1024 { t.Fatalf("unexpected upload instructions: %#v", draft.Upload) } + expectedParts := []PresignedObjectUploadPart{{PartNumber: 1, URL: "https://s3/part-1", ExpiresAt: "2026-01-01T00:15:00Z", Headers: map[string]string{}}} + if !reflect.DeepEqual(draft.Upload.Parts, expectedParts) { + t.Fatalf("unexpected upload instructions: %#v", draft.Upload) + } + if !reflect.DeepEqual(draft.Upload.URLs, []string{"https://s3/part-1"}) { + t.Fatalf("unexpected legacy urls: %#v", draft.Upload.URLs) + } } func TestCreateFolderUsesFolderMatterShape(t *testing.T) { diff --git a/cmd/internal/downloader/coverage_test.go b/cmd/internal/downloader/coverage_test.go index fc4db39d..45060414 100644 --- a/cmd/internal/downloader/coverage_test.go +++ b/cmd/internal/downloader/coverage_test.go @@ -19,6 +19,20 @@ import ( "github.com/saltbo/zpan/internal/config" ) +func testUploadInstructions(sessionID string, partSize int64, urls ...string) *client.ObjectUploadInstructions { + parts := make([]client.PresignedObjectUploadPart, 0, len(urls)) + for i, url := range urls { + parts = append(parts, client.PresignedObjectUploadPart{PartNumber: i + 1, URL: url}) + } + return &client.ObjectUploadInstructions{ + SessionID: sessionID, + Mode: "single", + PartSize: partSize, + PartCount: len(parts), + Parts: parts, + } +} + func TestDownloadTaskAccessors(t *testing.T) { runtime := &TaskRuntime{Engine: "aria2", Phase: "downloading"} task := DownloadTask{ @@ -310,7 +324,7 @@ func TestTaskRunnerRunHappyPathDownloadsUploadsAndCompletes(t *testing.T) { draft: client.ObjectDraft{ ID: "object-1", Name: "payload.bin", - Upload: &client.ObjectUploadInstructions{SessionID: "session-1", PartSize: int64(len(payload)), URLs: []string{uploadServer.URL}}, + Upload: testUploadInstructions("session-1", int64(len(payload)), uploadServer.URL), }, onCompleted: cancel, } @@ -505,7 +519,7 @@ func TestDirectoryUploadSuccessCreatesTree(t *testing.T) { createObjectDraft: client.ObjectDraft{ ID: "object-id", Name: "file", - Upload: &client.ObjectUploadInstructions{SessionID: "session", PartSize: 1, URLs: []string{uploadServer.URL}}, + Upload: testUploadInstructions("session", 1, uploadServer.URL), }, } uploader := NewUploader(api, nil) @@ -576,7 +590,11 @@ func TestUploadObjectSlicesTailPartAndMissingETag(t *testing.T) { Upload: &client.ObjectUploadInstructions{ SessionID: "session", PartSize: 3, - URLs: []string{uploadServer.URL, uploadServer.URL}, + PartCount: 2, + Parts: []client.PresignedObjectUploadPart{ + {PartNumber: 1, URL: uploadServer.URL}, + {PartNumber: 2, URL: uploadServer.URL}, + }, }, }, path, 5, &uploadProgress{totalBytes: 5, lastAt: time.Now()}) if err == nil || !strings.Contains(err.Error(), "missing ETag") { @@ -660,7 +678,7 @@ func TestUploadAndCompleteSuccessWithNilRuntimeCleansLocalFile(t *testing.T) { createObjectDraft: client.ObjectDraft{ ID: "object-1", Name: "payload.bin", - Upload: &client.ObjectUploadInstructions{SessionID: "session", PartSize: 1024, URLs: []string{uploadServer.URL}}, + Upload: testUploadInstructions("session", 1024, uploadServer.URL), }, } runner := NewTaskRunnerWithAPI(config.Config{}, api) diff --git a/cmd/internal/downloader/task_runner_test.go b/cmd/internal/downloader/task_runner_test.go index dc3bee01..a8b9c145 100644 --- a/cmd/internal/downloader/task_runner_test.go +++ b/cmd/internal/downloader/task_runner_test.go @@ -560,7 +560,7 @@ func TestUploadFilePartSendsContentLength(t *testing.T) { })) defer server.Close() - etag, err := uploadFilePart(context.Background(), server.URL, file, 0, 11, func(written int64) error { + etag, err := uploadFilePart(context.Background(), server.URL, nil, file, 0, 11, func(written int64) error { uploaded += written return nil }) @@ -594,7 +594,7 @@ func TestUploadFilePartIncludesErrorBody(t *testing.T) { })) defer server.Close() - _, err = uploadFilePart(context.Background(), server.URL, file, 0, 5, nil) + _, err = uploadFilePart(context.Background(), server.URL, nil, file, 0, 5, nil) if err == nil { t.Fatal("expected uploadFilePart error") } @@ -626,7 +626,7 @@ func TestUploadFilePartSendsSectionAndReturnsETag(t *testing.T) { })) defer server.Close() - etag, err := uploadFilePart(context.Background(), server.URL, file, 6, 9, func(written int64) error { + etag, err := uploadFilePart(context.Background(), server.URL, nil, file, 6, 9, func(written int64) error { uploaded += written return nil }) @@ -789,7 +789,7 @@ func TestWorkerLifecycleUploadFailurePreservesLocalResult(t *testing.T) { defer uploadServer.Close() api := &recordingAPI{ - createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: &client.ObjectUploadInstructions{SessionID: "session-1", PartSize: payloadSize, URLs: []string{uploadServer.URL}}}, + createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: testUploadInstructions("session-1", payloadSize, uploadServer.URL)}, completeErrs: []error{errors.New("unauthorized"), nil}, } eng := &recordingEngine{ @@ -842,7 +842,7 @@ func TestWorkerLifecycleHTTPUploadFailurePreservesLocalResult(t *testing.T) { payloadSize := int64(len(payload)) api := &recordingAPI{ - createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: &client.ObjectUploadInstructions{SessionID: "session-1", PartSize: payloadSize, URLs: []string{uploadServer.URL}}}, + createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: testUploadInstructions("session-1", payloadSize, uploadServer.URL)}, completeErrs: []error{errors.New("unauthorized"), nil}, } downloadDir := t.TempDir() @@ -1001,7 +1001,7 @@ func TestUploadShutdownMarksTaskInterrupted(t *testing.T) { payloadPath := writeTempFile(t, "downloaded payload") payloadSize := int64(len("downloaded payload")) api := &recordingAPI{ - createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: &client.ObjectUploadInstructions{SessionID: "session-1", PartSize: payloadSize, URLs: []string{"http://127.0.0.1:1"}}}, + createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: testUploadInstructions("session-1", payloadSize, "http://127.0.0.1:1")}, } w := NewTaskRunnerWithAPI(config.Config{}, api) @@ -1040,7 +1040,7 @@ func TestSuspendedUploadPreservesLocalResult(t *testing.T) { } payloadSize := int64(len(payload)) api := &recordingAPI{ - createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: &client.ObjectUploadInstructions{SessionID: "session-1", PartSize: payloadSize, URLs: []string{"http://127.0.0.1:1"}}}, + createObjectDraft: client.ObjectDraft{ID: "object-1", Name: "payload.bin", Upload: testUploadInstructions("session-1", payloadSize, "http://127.0.0.1:1")}, } w := NewTaskRunnerWithAPI(config.Config{}, api) diff --git a/cmd/internal/downloader/uploader.go b/cmd/internal/downloader/uploader.go index 7c85dc79..fca63339 100644 --- a/cmd/internal/downloader/uploader.go +++ b/cmd/internal/downloader/uploader.go @@ -101,16 +101,16 @@ func (u *Uploader) uploadSingleFile( if draft.Upload == nil { return "", fmt.Errorf("create remote object %s: missing upload instructions", draft.ID) } - log.Info("uploading file to object storage", "object_id", draft.ID, "path", path, "parts", len(draft.Upload.URLs)) + log.Info("uploading file to object storage", "object_id", draft.ID, "path", path, "parts", len(draft.Upload.Parts)) if err := u.uploadObjectSlices(ctx, log, task, draft, path, size, progress); err != nil { return "", fmt.Errorf("upload object %s: %w", draft.ID, err) } return draft.ID, nil } -// uploadObjectSlices runs the uniform upload: PUT each presigned slice (1 URL = -// single PutObject, N URLs = multipart), read each ETag, then finalize. On any -// failure it aborts the session, which also discards the draft. +// uploadObjectSlices runs the uniform upload: PUT each explicit presigned part, +// read each ETag, then finalize. On any failure it aborts the session, which also +// discards the draft. func (u *Uploader) uploadObjectSlices( ctx context.Context, log *slog.Logger, @@ -139,15 +139,15 @@ func (u *Uploader) uploadObjectSlices( } defer file.Close() - parts := make([]client.CompletedObjectUploadPart, 0, len(upload.URLs)) - for i, url := range upload.URLs { - partNumber := i + 1 - offset := int64(i) * upload.PartSize + parts := make([]client.CompletedObjectUploadPart, 0, len(upload.Parts)) + for _, part := range upload.Parts { + partNumber := part.PartNumber + offset := int64(partNumber-1) * upload.PartSize length := upload.PartSize if remaining := size - offset; remaining < length { length = remaining } - etag, err := uploadFilePart(ctx, url, file, offset, length, func(written int64) error { + etag, err := uploadFilePart(ctx, part.URL, part.Headers, file, offset, length, func(written int64) error { return u.reportUploadProgress(ctx, log, task, progress, written) }) if err != nil { @@ -326,7 +326,15 @@ func joinObjectPath(parent string, name string) string { return parent + "/" + name } -func uploadFilePart(ctx context.Context, url string, file *os.File, offset int64, length int64, progress func(written int64) error) (string, error) { +func uploadFilePart( + ctx context.Context, + url string, + headers map[string]string, + file *os.File, + offset int64, + length int64, + progress func(written int64) error, +) (string, error) { reader := io.NewSectionReader(file, offset, length) var body io.Reader = reader if progress != nil { @@ -337,6 +345,9 @@ func uploadFilePart(ctx context.Context, url string, file *os.File, offset int64 return "", err } req.ContentLength = length + for key, value := range headers { + req.Header.Set(key, value) + } res, err := http.DefaultClient.Do(req) if err != nil { return "", err diff --git a/cmd/internal/openapi/client.gen.go b/cmd/internal/openapi/client.gen.go index 3fe604cc..bb068bac 100644 --- a/cmd/internal/openapi/client.gen.go +++ b/cmd/internal/openapi/client.gen.go @@ -2509,6 +2509,24 @@ func (e CreateObjectJSONBodyOnConflict) Valid() bool { } } +// Defines values for CreateObject201JSONResponseBodyUploadMode. +const ( + CreateObject201JSONResponseBodyUploadModeMultipart CreateObject201JSONResponseBodyUploadMode = "multipart" + CreateObject201JSONResponseBodyUploadModeSingle CreateObject201JSONResponseBodyUploadMode = "single" +) + +// Valid indicates whether the value is a known member of the CreateObject201JSONResponseBodyUploadMode enum. +func (e CreateObject201JSONResponseBodyUploadMode) Valid() bool { + switch e { + case CreateObject201JSONResponseBodyUploadModeMultipart: + return true + case CreateObject201JSONResponseBodyUploadModeSingle: + return true + default: + return false + } +} + // Defines values for UpdateObjectJSONBodyOnConflict. const ( UpdateObjectJSONBodyOnConflictFail UpdateObjectJSONBodyOnConflict = "fail" @@ -2587,6 +2605,21 @@ func (e AbortObjectUploadParamsStrictStorageCleanup) Valid() bool { } } +// Defines values for PresignObjectUploadParts200JSONResponseBodyMode. +const ( + PresignObjectUploadParts200JSONResponseBodyModeMultipart PresignObjectUploadParts200JSONResponseBodyMode = "multipart" +) + +// Valid indicates whether the value is a known member of the PresignObjectUploadParts200JSONResponseBodyMode enum. +func (e PresignObjectUploadParts200JSONResponseBodyMode) Valid() bool { + switch e { + case PresignObjectUploadParts200JSONResponseBodyModeMultipart: + return true + default: + return false + } +} + // Defines values for ListSharesParamsStatus. const ( ListSharesParamsStatusActive ListSharesParamsStatus = "active" @@ -6700,6 +6733,9 @@ type CreateObjectJSONBody struct { // CreateObjectJSONBodyOnConflict defines parameters for CreateObject. type CreateObjectJSONBodyOnConflict string +// CreateObject201JSONResponseBodyUploadMode defines parameters for CreateObject. +type CreateObject201JSONResponseBodyUploadMode string + // UpdateObjectJSONBody defines parameters for UpdateObject. type UpdateObjectJSONBody struct { Name *string `json:"name,omitempty"` @@ -6750,6 +6786,9 @@ type PresignObjectUploadPartsJSONBody struct { PartNumbers []int `json:"partNumbers"` } +// PresignObjectUploadParts200JSONResponseBodyMode defines parameters for PresignObjectUploadParts. +type PresignObjectUploadParts200JSONResponseBodyMode string + // ListSharesParams defines parameters for ListShares. type ListSharesParams struct { PageSize *int `form:"pageSize,omitempty" json:"pageSize,omitempty"` @@ -30290,9 +30329,21 @@ type CreateObjectResponse struct { Type string `json:"type"` UpdatedAt string `json:"updatedAt"` Upload *struct { - PartSize int `json:"partSize"` - SessionId string `json:"sessionId"` - Urls []string `json:"urls"` + ExpiresAt string `json:"expiresAt"` + Mode CreateObject201JSONResponseBodyUploadMode `json:"mode"` + PartCount int `json:"partCount"` + PartSize int `json:"partSize"` + Parts []struct { + ExpiresAt string `json:"expiresAt"` + Headers map[string]string `json:"headers"` + PartNumber int `json:"partNumber"` + Url string `json:"url"` + } `json:"parts"` + PresignedExpiresAt string `json:"presignedExpiresAt"` + RequiredHeaders map[string]string `json:"requiredHeaders"` + SessionId string `json:"sessionId"` + UploadId *string `json:"uploadId"` + Urls []string `json:"urls"` } `json:"upload,omitempty"` } JSON400 *Error @@ -30576,12 +30627,18 @@ type PresignObjectUploadPartsResponse struct { Body []byte HTTPResponse *http.Response JSON200 *struct { - PartSize int `json:"partSize"` - Parts []struct { - PartNumber int `json:"partNumber"` - Url string `json:"url"` + Mode PresignObjectUploadParts200JSONResponseBodyMode `json:"mode"` + PartCount int `json:"partCount"` + PartSize int `json:"partSize"` + Parts []struct { + ExpiresAt string `json:"expiresAt"` + Headers map[string]string `json:"headers"` + PartNumber int `json:"partNumber"` + Url string `json:"url"` } `json:"parts"` - UploadId string `json:"uploadId"` + PresignedExpiresAt string `json:"presignedExpiresAt"` + RequiredHeaders map[string]string `json:"requiredHeaders"` + UploadId *string `json:"uploadId"` } JSON400 *Error JSON403 *Error @@ -45737,9 +45794,21 @@ func ParseCreateObjectResponse(rsp *http.Response) (*CreateObjectResponse, error Type string `json:"type"` UpdatedAt string `json:"updatedAt"` Upload *struct { - PartSize int `json:"partSize"` - SessionId string `json:"sessionId"` - Urls []string `json:"urls"` + ExpiresAt string `json:"expiresAt"` + Mode CreateObject201JSONResponseBodyUploadMode `json:"mode"` + PartCount int `json:"partCount"` + PartSize int `json:"partSize"` + Parts []struct { + ExpiresAt string `json:"expiresAt"` + Headers map[string]string `json:"headers"` + PartNumber int `json:"partNumber"` + Url string `json:"url"` + } `json:"parts"` + PresignedExpiresAt string `json:"presignedExpiresAt"` + RequiredHeaders map[string]string `json:"requiredHeaders"` + SessionId string `json:"sessionId"` + UploadId *string `json:"uploadId"` + Urls []string `json:"urls"` } `json:"upload,omitempty"` } if err := json.Unmarshal(bodyBytes, &dest); err != nil { @@ -46141,12 +46210,18 @@ func ParsePresignObjectUploadPartsResponse(rsp *http.Response) (*PresignObjectUp switch { case strings.Contains(rsp.Header.Get("Content-Type"), "json") && rsp.StatusCode == 200: var dest struct { - PartSize int `json:"partSize"` - Parts []struct { - PartNumber int `json:"partNumber"` - Url string `json:"url"` + Mode PresignObjectUploadParts200JSONResponseBodyMode `json:"mode"` + PartCount int `json:"partCount"` + PartSize int `json:"partSize"` + Parts []struct { + ExpiresAt string `json:"expiresAt"` + Headers map[string]string `json:"headers"` + PartNumber int `json:"partNumber"` + Url string `json:"url"` } `json:"parts"` - UploadId string `json:"uploadId"` + PresignedExpiresAt string `json:"presignedExpiresAt"` + RequiredHeaders map[string]string `json:"requiredHeaders"` + UploadId *string `json:"uploadId"` } if err := json.Unmarshal(bodyBytes, &dest); err != nil { return nil, err diff --git a/docs/design/object-upload-protocol.md b/docs/design/object-upload-protocol.md new file mode 100644 index 00000000..65d305b5 --- /dev/null +++ b/docs/design/object-upload-protocol.md @@ -0,0 +1,62 @@ +# Object Upload Protocol + +> Status: Implemented for v2.9 +> Product surface: `/api/objects` + +ZPan object uploads are control-plane-only through the server. File bytes are +PUT directly by the client to short-lived presigned S3 URLs and never traverse +ZPan server bandwidth. + +## Stable Operations + +Automation clients should discover and call these OpenAPI operation IDs: + +| Operation ID | Method and path | Purpose | +| --- | --- | --- | +| `createObject` | `POST /api/objects` | Create a folder, or create a file draft plus upload instructions. | +| `presignObjectUploadParts` | `POST /api/objects/{id}/uploads/{uploadSessionId}/parts` | Re-sign a bounded list of missing multipart part numbers. | +| `completeObjectUpload` | `POST /api/objects/{id}/uploads/{uploadSessionId}/completions` | Finalize a single or multipart upload with explicit part number + ETag records. | +| `abortObjectUpload` | `DELETE /api/objects/{id}/uploads/{uploadSessionId}` | Abort an active upload session and discard the draft. | + +## Upload Instructions + +`createObject` returns `upload` for file drafts: + +- `sessionId`: stable ZPan upload session ID. +- `uploadId`: S3 multipart upload ID for multipart sessions, otherwise `null`. +- `mode`: `single` or `multipart`. +- `partSize`: bytes per part. Single uploads use the file size, including `0` + for empty files. +- `partCount`: exact number of required part records for completion. +- `expiresAt`: ZPan upload session expiry. +- `presignedExpiresAt`: expiry for the returned presigned URLs. +- `requiredHeaders`: headers required by every returned part when applicable. +- `urls`: legacy positional presigned URL array retained for already-released + clients. New automation should prefer `parts`. +- `parts`: explicit descriptors with `partNumber`, `url`, `expiresAt`, and + `headers`. + +The v2.9 defaults are 64 MiB multipart parts and a 15 minute presign TTL. +Presigned upload URLs are short-lived. When a multipart URL expires, automation calls +`presignObjectUploadParts` with only the missing part numbers. Re-signing never +changes the workspace, object, session, storage key, or multipart upload +identity. + +## Completion And Abort + +`completeObjectUpload` accepts exactly the required `partCount` records. Each +record includes an explicit `partNumber` and the normalized S3 ETag returned by +that part PUT. Missing, duplicate, out-of-range, expired-session, or mismatched +single-PUT ETags are rejected before activation. Repeating completion after a +successful activation returns the active object. If S3 accepts multipart +completion but quota/conflict/database activation fails afterward, the upload +session records storage completion and a retry skips S3 completion while +retrying HEAD and draft activation. + +`abortObjectUpload` is idempotent for already-aborted sessions. Multipart aborts +call S3 `AbortMultipartUpload`; single PUT aborts delete the storage object on a +best-effort basis unless `strictStorageCleanup=1` is supplied. + +Every upload control-plane operation reauthorizes the caller's workspace and +object scope. Download-task upload tokens are also rechecked against their task, +downloader, and target folder before re-signing, completion, or abort. diff --git a/server/adapters/repos/object-upload-session.ts b/server/adapters/repos/object-upload-session.ts index 924c1a43..d0380d9d 100644 --- a/server/adapters/repos/object-upload-session.ts +++ b/server/adapters/repos/object-upload-session.ts @@ -17,6 +17,7 @@ function toRecord(row: SessionRow): ObjectUploadSessionRecord { onConflict: row.onConflict as ObjectUploadSessionRecord['onConflict'], status: row.status as ObjectUploadSessionRecord['status'], storageKey: row.storageKey, + createdBy: row.createdBy, expiresAt: row.expiresAt, createdAt: row.createdAt, updatedAt: row.updatedAt, diff --git a/server/http/downloads/download-tasks.integration.test.ts b/server/http/downloads/download-tasks.integration.test.ts index 727cbd34..7f7e829f 100644 --- a/server/http/downloads/download-tasks.integration.test.ts +++ b/server/http/downloads/download-tasks.integration.test.ts @@ -976,10 +976,10 @@ describe('Download tasks API integration', () => { const nestedObject = (await createNestedObjectRes.json()) as { id: string status: string - upload: { sessionId: string; urls: string[] } + upload: { sessionId: string; parts: Array<{ url: string }> } } expect(nestedObject.status).toBe('draft') - expect(nestedObject.upload.urls).toEqual(['https://presigned-upload.example.com']) + expect(nestedObject.upload.parts.map((part) => part.url)).toEqual(['https://presigned-upload.example.com']) const liveNestedParents = await db.all<{ count: number }>(sql` SELECT count(*) AS count FROM matters @@ -1028,11 +1028,12 @@ describe('Download tasks API integration', () => { const object = (await createObjectRes.json()) as { id: string status: string - upload: { sessionId: string; partSize: number; urls: string[] } + upload: { sessionId: string; partSize: number; mode: string; parts: Array<{ url: string }> } } expect(object.status).toBe('draft') - // 10 MiB ≤ 5 GiB → single PutObject: one presigned URL. - expect(object.upload.urls).toEqual(['https://presigned-upload.example.com']) + // 10 MiB stays single PutObject: one explicit presigned part. + expect(object.upload.mode).toBe('single') + expect(object.upload.parts.map((part) => part.url)).toEqual(['https://presigned-upload.example.com']) // Re-presign is multipart-only; this single-PUT session has no parts to re-presign. const partsRes = await app.request(`/api/objects/${object.id}/uploads/${object.upload.sessionId}/parts`, { @@ -1846,13 +1847,17 @@ describe('Download tasks API integration', () => { }), }) expect(createObjectRes.status).toBe(201) - const object = (await createObjectRes.json()) as { id: string; upload: { sessionId: string } } + const object = (await createObjectRes.json()) as { id: string; upload: { sessionId: string; partCount: number } } + const completeParts = Array.from({ length: object.upload.partCount }, (_, i) => ({ + partNumber: i + 1, + etag: `"etag-${i + 1}"`, + })) // Multipart completion calls CompleteMultipartUpload, which is mocked to reject. const completeRes = await app.request(`/api/objects/${object.id}/uploads/${object.upload.sessionId}/completions`, { method: 'POST', headers, - body: JSON.stringify({ parts: [{ partNumber: 1, etag: '"etag-1"' }] }), + body: JSON.stringify({ parts: completeParts }), }) expect(completeRes.status).toBe(502) diff --git a/server/http/objects.cf-test.ts b/server/http/objects.cf-test.ts index fef0cfe9..ed4a5dcf 100644 --- a/server/http/objects.cf-test.ts +++ b/server/http/objects.cf-test.ts @@ -181,10 +181,10 @@ describe('[CF] Objects upload + trash lifecycle (D1)', () => { const created = (await createRes.json()) as { id: string status: string - upload: { sessionId: string; urls: string[] } + upload: { sessionId: string; parts: Array<{ url: string }> } } expect(created.status).toBe('draft') - expect(created.upload.urls).toEqual(['https://cf-presigned-upload.example.com']) + expect(created.upload.parts.map((part) => part.url)).toEqual(['https://cf-presigned-upload.example.com']) // Finalize via completions (HEAD etag matches). const completeRes = await app.request( diff --git a/server/http/objects.integration.test.ts b/server/http/objects.integration.test.ts index 724e91fd..48950e99 100644 --- a/server/http/objects.integration.test.ts +++ b/server/http/objects.integration.test.ts @@ -4,10 +4,11 @@ import { eq, sql } from 'drizzle-orm' import { nanoid } from 'nanoid' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { S3Service } from '../adapters/gateways/s3.js' +import { createDownloadTokenGateway } from '../adapters/repos/download-tokens.js' import { createMatterRepo } from '../adapters/repos/matter.js' import { createQuotaRepo } from '../adapters/repos/quota.js' import { createStorageUsageRepo } from '../adapters/repos/storage-usage.js' -import { cloudTrafficReports, orgQuotaEntitlements, orgQuotas } from '../db/schema.js' +import { cloudTrafficReports, downloadTasks, orgQuotaEntitlements, orgQuotas } from '../db/schema.js' import { currentTrafficPeriod } from '../domain/quota.js' import { adminHeaders, authedHeaders, createTestApp, seedBusinessLicense, seedProLicense } from '../test/setup.js' import { type ConfirmUploadOptions, confirmUpload as confirmUploadUsecase } from '../usecases/object.js' @@ -506,6 +507,11 @@ describe('Objects API', () => { expect(S3Service.prototype.deleteObject).toHaveBeenCalled() const check = await app.request(`/api/objects/${created.id}`, { headers }) expect(check.status).toBe(404) + const repeated = await app.request(`/api/objects/${created.id}/uploads/${created.upload.sessionId}`, { + method: 'DELETE', + headers, + }) + expect(repeated.status).toBe(204) void db }) @@ -732,14 +738,24 @@ describe('Objects API', () => { const body = (await res.json()) as { status: string object: string - upload: { sessionId: string; partSize: number; urls: string[] } + upload: { + sessionId: string + partSize: number + partCount: number + mode: string + urls: string[] + parts: Array<{ url: string }> + } } expect(body.status).toBe('draft') expect(body.object).toBeTruthy() - // ≤5 GiB → single PutObject: one URL, partSize equals the file size. + // Small upload → single PutObject: one explicit part, partSize equals the file size. expect(body.upload.sessionId).toBeTruthy() + expect(body.upload.mode).toBe('single') expect(body.upload.partSize).toBe(2048) + expect(body.upload.partCount).toBe(1) expect(body.upload.urls).toEqual(['https://presigned-upload.example.com']) + expect(body.upload.parts.map((part) => part.url)).toEqual(['https://presigned-upload.example.com']) }) it('POST /api/objects with storageId uses that exact eligible storage', async () => { @@ -2502,13 +2518,17 @@ describe('object multipart upload API with S3-compatible storage', () => { expect(createRes.status).toBe(201) const object = (await createRes.json()) as { id: string - upload: { sessionId: string; partSize: number; urls: string[] } + upload: { sessionId: string; partSize: number; parts: Array<{ url: string; headers: Record }> } } - expect(object.upload.urls).toHaveLength(1) + expect(object.upload.parts).toHaveLength(1) expect(object.upload.partSize).toBe(11) // PUT the bytes directly to the presigned URL, read the ETag. - const putRes = await fetch(object.upload.urls[0], { method: 'PUT', body: 'hello world' }) + const putRes = await fetch(object.upload.parts[0].url, { + method: 'PUT', + headers: object.upload.parts[0].headers, + body: 'hello world', + }) expect(putRes.status).toBe(200) const etagHeader = putRes.headers.get('etag') expect(etagHeader).toBeTruthy() @@ -2825,4 +2845,95 @@ describe('Objects API — error branches', () => { const body = (await res.json()) as { error: { message: string } } expect(body.error.message).toBe('Forbidden') }) + + it('rejects download-task-upload re-sign outside the task target folder', async () => { + const { app, db } = await createTestApp({ DOWNLOAD_TOKEN_SECRET: 'test-download-token-secret' }) + await insertStorage(db) + const { uploadToken, orgId } = await mintTaskUploadContext(app, db, { targetFolder: 'Remote' }) + await insertFile(db, orgId, { + id: 'm-task-presign-outside', + name: 'file.txt', + parent: 'Elsewhere', + status: 'draft', + }) + + const res = await app.request('/api/objects/m-task-presign-outside/uploads/any-session/parts', { + method: 'POST', + headers: { Authorization: `Bearer ${uploadToken}`, 'Content-Type': 'application/json' }, + body: JSON.stringify({ partNumbers: [1] }), + }) + + expect(res.status).toBe(403) + const body = (await res.json()) as { error: { message: string } } + expect(body.error.message).toBe('Forbidden') + }) + + it('rejects download-task-upload abort outside the task target folder', async () => { + const { app, db } = await createTestApp({ DOWNLOAD_TOKEN_SECRET: 'test-download-token-secret' }) + await insertStorage(db) + const { uploadToken, orgId } = await mintTaskUploadContext(app, db, { targetFolder: 'Remote' }) + await insertFile(db, orgId, { id: 'm-task-abort-outside', name: 'file.txt', parent: 'Elsewhere', status: 'draft' }) + + const res = await app.request('/api/objects/m-task-abort-outside/uploads/any-session', { + method: 'DELETE', + headers: { Authorization: `Bearer ${uploadToken}` }, + }) + + expect(res.status).toBe(403) + const body = (await res.json()) as { error: { message: string } } + expect(body.error.message).toBe('Forbidden') + }) + + it('allows repeated download-task-upload abort after the draft is gone', async () => { + const { app, db, platform } = await createTestApp({ DOWNLOAD_TOKEN_SECRET: 'test-download-token-secret' }) + await insertStorage(db) + const targetFolder = 'Remote' + const { uploadToken } = await mintTaskUploadContext(app, db, { targetFolder }) + const createRes = await app.request('/api/objects', { + method: 'POST', + headers: { Authorization: `Bearer ${uploadToken}`, 'Content-Type': 'application/json' }, + body: JSON.stringify({ + name: 'task-cancel.txt', + type: 'text/plain', + size: 1, + parent: targetFolder, + }), + }) + expect(createRes.status).toBe(201) + const created = (await createRes.json()) as { id: string; upload: { sessionId: string } } + + const first = await app.request(`/api/objects/${created.id}/uploads/${created.upload.sessionId}`, { + method: 'DELETE', + headers: { Authorization: `Bearer ${uploadToken}` }, + }) + expect(first.status).toBe(204) + const tokenGateway = createDownloadTokenGateway() + const claims = await tokenGateway.verifyDownloadToken(platform, uploadToken) + if (!claims || claims.typ !== 'download-task-upload') throw new Error('task_upload_claims_missing') + const otherDownloaderId = 'other-downloader' + await db + .update(downloadTasks) + .set({ assignedDownloaderId: otherDownloaderId }) + .where(eq(downloadTasks.id, claims.taskId)) + const otherUploadToken = await tokenGateway.signDownloadToken(platform, { + ...claims, + downloaderId: otherDownloaderId, + jti: randomUUID(), + }) + const otherDownloader = await app.request(`/api/objects/${created.id}/uploads/${created.upload.sessionId}`, { + method: 'DELETE', + headers: { Authorization: `Bearer ${otherUploadToken}` }, + }) + expect(otherDownloader.status).toBe(403) + await db + .update(downloadTasks) + .set({ assignedDownloaderId: claims.downloaderId }) + .where(eq(downloadTasks.id, claims.taskId)) + const second = await app.request(`/api/objects/${created.id}/uploads/${created.upload.sessionId}`, { + method: 'DELETE', + headers: { Authorization: `Bearer ${uploadToken}` }, + }) + expect(second.status).toBe(204) + void db + }) }) diff --git a/server/http/objects.ts b/server/http/objects.ts index 5cbf0995..d043ffd1 100644 --- a/server/http/objects.ts +++ b/server/http/objects.ts @@ -18,6 +18,7 @@ import { transferAuditActor } from '../middleware/audit-transfers' import type { Env } from '../middleware/platform' import { abortUpload, + authorizeTaskUploadAbort, authorizeTaskUploadConfirm, completeUpload, copyObject, @@ -146,6 +147,36 @@ function actorId(c: Context): string { return c.get('userId') ?? 'system' } +async function authorizeUploadSessionControl( + c: Context, + orgId: string, + objectId: string, + options: { uploadSessionId?: string } = {}, +): Promise { + const principal = c.get('principal') + if (principal?.kind !== 'download-task-upload') return + if (options.uploadSessionId) { + const authorized = await authorizeTaskUploadAbort(c.get('deps'), { + orgId, + objectId, + sessionId: options.uploadSessionId, + taskId: principal.taskId, + downloaderId: principal.downloaderId, + targetFolder: principal.targetFolder, + }) + if (!authorized.ok) throw authorized.error + return + } + const authorized = await authorizeTaskUploadConfirm(c.get('deps'), { + orgId, + objectId, + taskId: principal.taskId, + downloaderId: principal.downloaderId, + targetFolder: principal.targetFolder, + }) + if (!authorized.ok) throw authorized.error +} + const cloudBaseUrl = (c: Context) => c.get('platform').getEnv('ZPAN_CLOUD_URL') ?? ZPAN_CLOUD_URL_DEFAULT const listRoute = authRoute( @@ -441,9 +472,11 @@ const objects = app .openapi(presignPartsRoute, async (c) => { const orgId = c.get('orgId') if (!orgId) throw new ObjectUploadSessionError('not_found') + const objectId = c.req.valid('param').id + await authorizeUploadSessionControl(c, orgId, objectId) const result = await presignUploadSessionParts(c.get('deps'), { orgId, - objectId: c.req.valid('param').id, + objectId, sessionId: c.req.valid('param').uploadSessionId, partNumbers: c.req.valid('json').partNumbers, }) @@ -457,18 +490,7 @@ const objects = app const orgId = c.get('orgId') if (!orgId) throw new ObjectUploadSessionError('not_found') const objectId = c.req.valid('param').id - - const principal = c.get('principal') - if (principal?.kind === 'download-task-upload') { - const authorized = await authorizeTaskUploadConfirm(c.get('deps'), { - orgId, - objectId, - taskId: principal.taskId, - downloaderId: principal.downloaderId, - targetFolder: principal.targetFolder, - }) - if (!authorized.ok) throw authorized.error - } + await authorizeUploadSessionControl(c, orgId, objectId) const result = await completeUpload(c.get('deps'), { orgId, @@ -486,10 +508,13 @@ const objects = app .openapi(abortUploadRoute, async (c) => { const orgId = c.get('orgId') if (!orgId) throw new ObjectUploadSessionError('not_found') + const objectId = c.req.valid('param').id + const uploadSessionId = c.req.valid('param').uploadSessionId + await authorizeUploadSessionControl(c, orgId, objectId, { uploadSessionId }) await abortUpload(c.get('deps'), { orgId, - objectId: c.req.valid('param').id, - sessionId: c.req.valid('param').uploadSessionId, + objectId, + sessionId: uploadSessionId, actorId: actorId(c), strictStorageCleanup: c.req.valid('query').strictStorageCleanup !== undefined, }) diff --git a/server/usecases/object.test.ts b/server/usecases/object.test.ts index 776b7cab..0eea3917 100644 --- a/server/usecases/object.test.ts +++ b/server/usecases/object.test.ts @@ -2,6 +2,7 @@ import { DirType } from '@shared/constants' import { beforeEach, describe, expect, it, vi } from 'vitest' import { abortUpload, + authorizeTaskUploadAbort, authorizeTaskUploadConfirm, completeUpload, copyObject, @@ -13,9 +14,11 @@ import { listObjects, listTrashedObjects, type ObjectActor, + presignUploadSessionParts, restoreObject, transferObject, trashObject, + UPLOAD_PRESIGNED_URL_TTL_SECONDS, updateObject, } from './object' import { @@ -320,12 +323,29 @@ describe('object usecase', () => { }) expect(out.ok).toBe(true) if (out.ok && 'upload' in out) { - // ≤5 GiB → single PutObject: one URL, partSize equals the file size. + // Small file → single PutObject: one explicit part, partSize equals the file size. expect(out.upload.sessionId).toBe('sess-1') - expect(out.upload.urls).toEqual(['https://up']) + expect(out.upload.mode).toBe('single') + expect(out.upload.uploadId).toBeNull() expect(out.upload.partSize).toBe(2048) + expect(out.upload.partCount).toBe(1) + expect(out.upload.requiredHeaders).toEqual({ 'content-type': 'image/jpeg' }) + expect(out.upload.urls).toEqual(['https://up']) + expect(out.upload.parts).toEqual([ + { + partNumber: 1, + url: 'https://up', + expiresAt: out.upload.presignedExpiresAt, + headers: { 'content-type': 'image/jpeg' }, + }, + ]) expect(out.matter.status).toBe('draft') - expect(presignUpload).toHaveBeenCalledWith(storage, expect.any(String), 'image/jpeg') + expect(presignUpload).toHaveBeenCalledWith( + storage, + expect.any(String), + 'image/jpeg', + UPLOAD_PRESIGNED_URL_TTL_SECONDS, + ) } else { throw new Error('expected upload outcome') } @@ -346,7 +366,12 @@ describe('object usecase', () => { expect(out.ok).toBe(true) expect(create).toHaveBeenCalledWith(expect.objectContaining({ type: '' })) - expect(presignUpload).toHaveBeenCalledWith(storage, expect.any(String), undefined) + expect(presignUpload).toHaveBeenCalledWith( + storage, + expect.any(String), + undefined, + UPLOAD_PRESIGNED_URL_TTL_SECONDS, + ) }) it('uses an eligible requested storage for a file draft', async () => { @@ -417,11 +442,13 @@ describe('object usecase', () => { }) it('creates a large file draft and returns multipart upload instructions', async () => { - const fiveGiB = 5 * 1024 * 1024 * 1024 - const size = fiveGiB + 1 // just over 5 GiB → two 5 GiB parts + const multipartPartSize = 64 * 1024 * 1024 + const size = multipartPartSize + 1 const draft = file('big', { status: 'draft', object: 'o1/u1/big.bin', name: 'big.bin', size }) const createMultipartUpload = vi.fn(async () => 'mp-1') - const presignUploadPart = vi.fn(async () => 'https://part') + const presignUploadPart = vi.fn( + async (_storage: unknown, _key: string, _uploadId: string, partNumber: number) => `https://part-${partNumber}`, + ) const { deps } = makeDeps({ matter: { create: async () => draft }, s3: { createMultipartUpload, presignUploadPart } as Partial, @@ -433,9 +460,24 @@ describe('object usecase', () => { }) expect(out.ok).toBe(true) if (out.ok && 'upload' in out) { - expect(out.upload.partSize).toBe(fiveGiB) - expect(out.upload.urls).toHaveLength(2) + expect(out.upload.mode).toBe('multipart') + expect(out.upload.uploadId).toBe('mp-1') + expect(out.upload.partSize).toBe(multipartPartSize) + expect(out.upload.partCount).toBe(2) + expect(out.upload.requiredHeaders).toEqual({}) + expect(out.upload.urls).toEqual(['https://part-1', 'https://part-2']) + expect(out.upload.parts).toEqual([ + { partNumber: 1, url: 'https://part-1', expiresAt: out.upload.presignedExpiresAt, headers: {} }, + { partNumber: 2, url: 'https://part-2', expiresAt: out.upload.presignedExpiresAt, headers: {} }, + ]) expect(createMultipartUpload).toHaveBeenCalled() + expect(presignUploadPart).toHaveBeenCalledWith( + storage, + expect.any(String), + 'mp-1', + 1, + UPLOAD_PRESIGNED_URL_TTL_SECONDS, + ) } else { throw new Error('expected upload outcome') } @@ -691,6 +733,193 @@ describe('object usecase', () => { expect(activateDraft).toHaveBeenCalledWith('d1', 'o1', 'd1.txt', 'audio/flac', expect.any(Date)) }) + it('retries activation after multipart storage completion without completing S3 again', async () => { + const draft = file('d1', { status: 'draft', size: 100 }) + let sessionStatus: ObjectUploadSessionRecord['status'] = 'active' + let quotaAttempts = 0 + const completeMultipartUpload = vi.fn(async () => {}) + const setStatus = vi.fn(async (_id: string, status: ObjectUploadSessionRecord['status']) => { + sessionStatus = status + }) + const headObject = vi.fn(async () => ({ size: 100, contentType: 'audio/flac', etag: 'multipart-etag' })) + const activateDraft = vi.fn(async () => true) + const { deps } = makeDeps({ + matter: { get: async () => draft, activateDraft }, + s3: { completeMultipartUpload, headObject } as Partial, + objectUploadSessions: { + get: async () => session({ status: sessionStatus, uploadId: 'mp-1' }), + setStatus, + }, + quota: { + incrementUsageIfEffectiveQuotaAllows: async () => { + quotaAttempts += 1 + return quotaAttempts > 1 + }, + }, + }) + + const first = await completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }) + expectError(first, 422, 'Quota exceeded', 'QUOTA_EXCEEDED') + + const second = await completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }) + + expect(second.ok).toBe(true) + expect(completeMultipartUpload).toHaveBeenCalledTimes(1) + expect(setStatus).toHaveBeenCalledTimes(1) + expect(setStatus).toHaveBeenCalledWith('sess-1', 'completed') + expect(headObject).toHaveBeenCalledTimes(2) + expect(activateDraft).toHaveBeenCalledTimes(1) + }) + + it('retries a multipart activation race without completing S3 again', async () => { + const draft = file('d1', { status: 'draft', size: 100 }) + let sessionStatus: ObjectUploadSessionRecord['status'] = 'active' + let activationAttempts = 0 + const completeMultipartUpload = vi.fn(async () => {}) + const setStatus = vi.fn(async (_id: string, status: ObjectUploadSessionRecord['status']) => { + sessionStatus = status + }) + const headObject = vi.fn(async () => ({ size: 100, contentType: 'audio/flac', etag: 'multipart-etag' })) + const activateDraft = vi.fn(async () => { + activationAttempts += 1 + return activationAttempts > 1 + }) + const { deps } = makeDeps({ + matter: { get: async () => draft, activateDraft }, + s3: { completeMultipartUpload, headObject } as Partial, + objectUploadSessions: { + get: async () => session({ status: sessionStatus, uploadId: 'mp-1' }), + setStatus, + }, + }) + + const first = await completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }) + expect(first).toEqual({ ok: false, reason: 'not_found' }) + + const second = await completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }) + + expect(second.ok).toBe(true) + expect(completeMultipartUpload).toHaveBeenCalledTimes(1) + expect(setStatus).toHaveBeenCalledTimes(1) + expect(setStatus).toHaveBeenCalledWith('sess-1', 'completed') + expect(headObject).toHaveBeenCalledTimes(2) + expect(activateDraft).toHaveBeenCalledTimes(2) + }) + + it('returns the active object when completion is repeated for an already-completed session', async () => { + const active = file('d1', { status: 'active' }) + const completeMultipartUpload = vi.fn() + const headObject = vi.fn() + const setStatus = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => active }, + s3: { completeMultipartUpload, headObject } as Partial, + objectUploadSessions: { get: async () => session({ status: 'completed', uploadId: 'mp-1' }), setStatus }, + }) + + const out = await completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }) + + expect(out).toEqual({ ok: true, matter: active }) + expect(completeMultipartUpload).not.toHaveBeenCalled() + expect(headObject).not.toHaveBeenCalled() + expect(setStatus).not.toHaveBeenCalled() + }) + + it('rejects duplicate multipart completion parts before completing storage', async () => { + const completeMultipartUpload = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 200 }) }, + s3: { completeMultipartUpload } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + await expect( + completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [ + { partNumber: 1, etag: 'e1' }, + { partNumber: 1, etag: 'e1-again' }, + ], + actorId: 'u1', + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(completeMultipartUpload).not.toHaveBeenCalled() + }) + + it('rejects missing multipart completion parts before completing storage', async () => { + const completeMultipartUpload = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 200 }) }, + s3: { completeMultipartUpload } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + await expect( + completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(completeMultipartUpload).not.toHaveBeenCalled() + }) + + it('rejects an expired upload session before completing storage', async () => { + const completeMultipartUpload = vi.fn() + const headObject = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft' }) }, + s3: { completeMultipartUpload, headObject } as Partial, + objectUploadSessions: { get: async () => session({ expiresAt: new Date(Date.now() - 1_000) }) }, + }) + + await expect( + completeUpload(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + parts: [{ partNumber: 1, etag: 'e1' }], + actorId: 'u1', + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(completeMultipartUpload).not.toHaveBeenCalled() + expect(headObject).not.toHaveBeenCalled() + }) + it('throws not_found when the upload session is missing', async () => { const draft = file('d1', { status: 'draft' }) const { deps } = makeDeps({ @@ -750,6 +979,62 @@ describe('object usecase', () => { }) }) + describe('presignUploadSessionParts', () => { + it('re-signs bounded explicit multipart part numbers with expiry metadata', async () => { + const presignUploadPart = vi.fn( + async (_storage: unknown, _key: string, _uploadId: string, partNumber: number) => `https://part-${partNumber}`, + ) + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 250 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + const out = await presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [3, 1], + }) + + expect(out.mode).toBe('multipart') + expect(out.uploadId).toBe('mp-1') + expect(out.partCount).toBe(3) + expect(out.parts).toEqual([ + { partNumber: 3, url: 'https://part-3', expiresAt: out.presignedExpiresAt, headers: {} }, + { partNumber: 1, url: 'https://part-1', expiresAt: out.presignedExpiresAt, headers: {} }, + ]) + expect(presignUploadPart).toHaveBeenCalledWith(storage, 'key/d1', 'mp-1', 3, UPLOAD_PRESIGNED_URL_TTL_SECONDS) + }) + + it('rejects duplicate or out-of-session re-sign part numbers', async () => { + const presignUploadPart = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 250 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + await expect( + presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [1, 1], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + await expect( + presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [4], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(presignUploadPart).not.toHaveBeenCalled() + }) + }) + describe('abortUpload', () => { it('aborts a single-PUT session: best-effort S3 delete + discards the draft', async () => { const setStatus = vi.fn(async () => {}) @@ -802,24 +1087,203 @@ describe('object usecase', () => { expect(abortMultipartUpload).toHaveBeenCalledWith(storage, 'key/d1', 'mp-1') }) - it('is idempotent for an already-aborted session', async () => { + it('abandons a completed multipart draft by deleting the finalized object', async () => { + const deleteObjectFn = vi.fn(async () => {}) + const abortMultipartUpload = vi.fn() + const setStatus = vi.fn(async () => {}) + const cancelDraft = vi.fn(async () => file('d1', { status: 'draft' })) + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft' }), cancelDraft }, + s3: { deleteObject: deleteObjectFn, abortMultipartUpload } as Partial, + objectUploadSessions: { get: async () => session({ status: 'completed', uploadId: 'mp-1' }), setStatus }, + }) + + await abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }) + + expect(deleteObjectFn).toHaveBeenCalledWith(storage, 'key/d1') + expect(abortMultipartUpload).not.toHaveBeenCalled() + expect(setStatus).toHaveBeenCalledWith('sess-1', 'aborted') + expect(cancelDraft).toHaveBeenCalledWith('d1', 'o1') + }) + + it('finishes draft cancellation when retrying an already-aborted session', async () => { const cancelDraft = vi.fn(async () => null) const { deps } = makeDeps({ matter: { get: async () => file('d1', { status: 'draft' }), cancelDraft }, objectUploadSessions: { get: async () => session({ status: 'aborted' }) }, }) await abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }) + expect(cancelDraft).toHaveBeenCalledWith('d1', 'o1') + }) + + it('returns successfully for an already-aborted session whose draft was already deleted', async () => { + const cancelDraft = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => null, cancelDraft }, + objectUploadSessions: { get: async () => session({ status: 'aborted' }) }, + }) + await abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }) expect(cancelDraft).not.toHaveBeenCalled() }) - it('rejects aborting an already-completed session', async () => { + it('is idempotent for an already-aborted non-draft session', async () => { + const cancelDraft = vi.fn(async () => null) const { deps } = makeDeps({ - matter: { get: async () => file('d1', { status: 'draft' }) }, + matter: { get: async () => file('d1', { status: 'cancelled' }), cancelDraft }, + objectUploadSessions: { get: async () => session({ status: 'aborted' }) }, + }) + await abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }) + expect(cancelDraft).not.toHaveBeenCalled() + }) + + it('recovers when completed draft abort marked aborted before draft cancellation failed', async () => { + let status: ObjectUploadSessionRecord['status'] = 'completed' + let cancelAttempts = 0 + const deleteObjectFn = vi.fn(async () => {}) + const setStatus = vi.fn(async (_id: string, next: ObjectUploadSessionRecord['status']) => { + status = next + }) + const cancelDraft = vi.fn(async () => { + cancelAttempts += 1 + if (cancelAttempts === 1) throw new Error('db write failed') + return file('d1', { status: 'draft' }) + }) + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft' }), cancelDraft }, + s3: { deleteObject: deleteObjectFn }, + objectUploadSessions: { get: async () => session({ status, uploadId: 'mp-1' }), setStatus }, + }) + + await expect( + abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }), + ).rejects.toThrow('db write failed') + await abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }) + + expect(deleteObjectFn).toHaveBeenCalledTimes(1) + expect(setStatus).toHaveBeenCalledWith('sess-1', 'aborted') + expect(cancelDraft).toHaveBeenCalledTimes(2) + expect(cancelDraft).toHaveBeenLastCalledWith('d1', 'o1') + }) + + it('rejects aborting an already-completed active session', async () => { + const deleteObjectFn = vi.fn() + const cancelDraft = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'active' }), cancelDraft }, + s3: { deleteObject: deleteObjectFn }, objectUploadSessions: { get: async () => session({ status: 'completed' }) }, }) await expect( abortUpload(deps, { orgId: 'o1', objectId: 'd1', sessionId: 'sess-1', actorId: 'u1' }), ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(deleteObjectFn).not.toHaveBeenCalled() + expect(cancelDraft).not.toHaveBeenCalled() + }) + }) + + describe('presignUploadSessionParts', () => { + it('re-signs bounded multipart parts with explicit descriptors', async () => { + const presignUploadPart = vi.fn( + async (_storage, _key, _uploadId, partNumber: number) => `https://part-${partNumber}`, + ) + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 300 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + const out = await presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [2, 3], + }) + + expect(out).toMatchObject({ + uploadId: 'mp-1', + mode: 'multipart', + partSize: 100, + partCount: 3, + requiredHeaders: {}, + parts: [ + { partNumber: 2, url: 'https://part-2', headers: {} }, + { partNumber: 3, url: 'https://part-3', headers: {} }, + ], + }) + expect(out.parts[0]?.expiresAt).toBe(out.presignedExpiresAt) + expect(presignUploadPart).toHaveBeenCalledWith(storage, 'key/d1', 'mp-1', 2, UPLOAD_PRESIGNED_URL_TTL_SECONDS) + }) + + it('rejects duplicate re-sign part numbers before signing', async () => { + const presignUploadPart = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 300 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + await expect( + presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [2, 2], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(presignUploadPart).not.toHaveBeenCalled() + }) + + it('rejects re-sign part numbers outside the session part count before signing', async () => { + const presignUploadPart = vi.fn() + const { deps } = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 300 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: 'mp-1', partSize: 100 }) }, + }) + + await expect( + presignUploadSessionParts(deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [4], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(presignUploadPart).not.toHaveBeenCalled() + }) + + it('rejects re-sign for expired sessions and single-PUT sessions before signing', async () => { + const presignUploadPart = vi.fn() + const expired = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 300 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { + get: async () => session({ uploadId: 'mp-1', partSize: 100, expiresAt: new Date(Date.now() - 1_000) }), + }, + }) + await expect( + presignUploadSessionParts(expired.deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [1], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + + const single = makeDeps({ + matter: { get: async () => file('d1', { status: 'draft', size: 100 }) }, + s3: { presignUploadPart } as Partial, + objectUploadSessions: { get: async () => session({ uploadId: null, partSize: 100 }) }, + }) + await expect( + presignUploadSessionParts(single.deps, { + orgId: 'o1', + objectId: 'd1', + sessionId: 'sess-1', + partNumbers: [1], + }), + ).rejects.toMatchObject({ code: 'invalid_state' }) + expect(presignUploadPart).not.toHaveBeenCalled() }) }) @@ -1267,4 +1731,51 @@ describe('object usecase', () => { expect(await authorizeTaskUploadConfirm(deps, taskParams)).toEqual({ ok: true }) }) }) + + describe('authorizeTaskUploadAbort', () => { + const taskParams = { + orgId: 'o1', + objectId: 'm1', + sessionId: 'sess-1', + taskId: 't1', + downloaderId: 'd1', + targetFolder: 'Inbox', + } + + it('authorizes an already-aborted upload session after the draft was deleted', async () => { + const matterGet = vi.fn(async () => null) + const { deps } = makeDeps({ + matter: { get: matterGet }, + objectUploadSessions: { + get: async () => session({ status: 'aborted', createdBy: 'downloader:d1' }), + }, + downloadTasks: { + findRecord: async () => + ({ id: 't1', assignedDownloaderId: 'd1', status: 'uploading' }) as unknown as DownloadTaskRecord, + }, + }) + expect(await authorizeTaskUploadAbort(deps, taskParams)).toEqual({ ok: true }) + expect(matterGet).not.toHaveBeenCalled() + }) + + it('forbids an already-aborted upload session created by another downloader', async () => { + const matterGet = vi.fn(async () => null) + const { deps } = makeDeps({ + matter: { get: matterGet }, + objectUploadSessions: { + get: async () => session({ status: 'aborted', createdBy: 'downloader:other' }), + }, + }) + expectError(await authorizeTaskUploadAbort(deps, taskParams), 403, 'Forbidden') + expect(matterGet).not.toHaveBeenCalled() + }) + + it('still forbids non-aborted abort outside the task target folder', async () => { + const { deps } = makeDeps({ + matter: { get: async () => file('m1', { parent: 'Other' }) }, + objectUploadSessions: { get: async () => session() }, + }) + expectError(await authorizeTaskUploadAbort(deps, taskParams), 403, 'Forbidden') + }) + }) }) diff --git a/server/usecases/object.ts b/server/usecases/object.ts index e12f51c0..f1597b71 100644 --- a/server/usecases/object.ts +++ b/server/usecases/object.ts @@ -143,12 +143,15 @@ export async function listObjects( // ─── Create (folder, or file draft + size-decided upload) ───────────────────── -// The server picks the S3 mechanism by size: ≤5 GiB → single PutObject (1 URL), -// >5 GiB → multipart with 5 GiB parts (N URLs), >5 TiB → rejected (S3's hard -// single-object ceiling). The client's flow is uniform: PUT each slice, read the -// ETag, then POST them to .../completions. -const PART_SIZE_BYTES = 5 * 1024 * 1024 * 1024 // 5 GiB +// The server picks the S3 mechanism by size: small objects use single PutObject, +// larger objects use S3 multipart with explicit part descriptors. Multipart +// defaults to automation-friendly 64 MiB chunks but grows when needed to respect +// S3's 10,000-part ceiling and 5 TiB single-object maximum. +const SINGLE_UPLOAD_MAX_BYTES = 64 * 1024 * 1024 // 64 MiB +const DEFAULT_MULTIPART_PART_SIZE_BYTES = 64 * 1024 * 1024 // 64 MiB +const MAX_UPLOAD_PARTS = 10_000 const MAX_OBJECT_BYTES = 5 * 1024 * 1024 * 1024 * 1024 // 5 TiB +export const UPLOAD_PRESIGNED_URL_TTL_SECONDS = 15 * 60 export type CreateObjectOutcome = | { ok: true; matter: Matter } @@ -246,16 +249,27 @@ async function prepareUpload( const { storage, storageKey, contentType, size } = params let uploadId: string | null = null let partSize: number - let urls: string[] + let partCount: number + let parts: ObjectUploadInstructions['parts'] + const headers = signedUploadHeaders(contentType) + const presignedExpiresAt = new Date(Date.now() + UPLOAD_PRESIGNED_URL_TTL_SECONDS * 1000).toISOString() - if (size <= PART_SIZE_BYTES) { + if (size <= SINGLE_UPLOAD_MAX_BYTES) { // Single PutObject — one presigned PUT, no S3 multipart overhead. When the // client supplied a content type it is signed and must be sent verbatim. partSize = size - urls = [await deps.s3.presignUpload(storage, storageKey, contentType)] + partCount = 1 + parts = [ + { + partNumber: 1, + url: await deps.s3.presignUpload(storage, storageKey, contentType, UPLOAD_PRESIGNED_URL_TTL_SECONDS), + expiresAt: presignedExpiresAt, + headers, + }, + ] } else { - partSize = PART_SIZE_BYTES - const partCount = Math.ceil(size / partSize) + partSize = Math.max(DEFAULT_MULTIPART_PART_SIZE_BYTES, Math.ceil(size / MAX_UPLOAD_PARTS)) + partCount = Math.ceil(size / partSize) try { uploadId = await deps.s3.createMultipartUpload(storage, storageKey, contentType) } catch (error) { @@ -265,8 +279,13 @@ async function prepareUpload( ) } const mpId = uploadId - urls = await Promise.all( - Array.from({ length: partCount }, (_, i) => deps.s3.presignUploadPart(storage, storageKey, mpId, i + 1)), + parts = await Promise.all( + Array.from({ length: partCount }, async (_, i) => ({ + partNumber: i + 1, + url: await deps.s3.presignUploadPart(storage, storageKey, mpId, i + 1, UPLOAD_PRESIGNED_URL_TTL_SECONDS), + expiresAt: presignedExpiresAt, + headers: {}, + })), ) } @@ -280,7 +299,49 @@ async function prepareUpload( onConflict: params.onConflict, actorId: params.actorId, }) - return { sessionId: record.id, partSize, urls } + return { + sessionId: record.id, + uploadId, + mode: uploadId == null ? 'single' : 'multipart', + partSize, + partCount, + expiresAt: record.expiresAt.toISOString(), + presignedExpiresAt, + requiredHeaders: uploadId == null ? headers : {}, + urls: parts.map((part) => part.url), + parts, + } +} + +function signedUploadHeaders(contentType?: string): Record { + return contentType ? { 'content-type': contentType } : {} +} + +function uploadPartCount(size: number, partSize: number): number { + if (partSize <= 0) return 1 + return Math.max(1, Math.ceil(size / partSize)) +} + +function normalizeETag(etag: string): string { + return etag.trim().replace(/^"+|"+$/g, '') +} + +function validateCompletionParts( + parts: CompleteObjectUploadInput['parts'], + requiredPartCount: number, +): Array<{ partNumber: number; etag: string }> { + const normalized = parts.map((part) => ({ partNumber: part.partNumber, etag: normalizeETag(part.etag) })) + const seen = new Set() + for (const part of normalized) { + if (part.partNumber < 1 || part.partNumber > requiredPartCount || seen.has(part.partNumber) || !part.etag) { + throw new ObjectUploadSessionError('invalid_state', 'Upload completion parts do not match the session') + } + seen.add(part.partNumber) + } + if (seen.size !== requiredPartCount) { + throw new ObjectUploadSessionError('invalid_state', 'Upload completion parts do not match the session') + } + return normalized.sort((a, b) => a.partNumber - b.partNumber) } // ─── Upload finalize / abort / re-presign ───────────────────────────────────── @@ -309,8 +370,16 @@ export async function presignUploadSessionParts( sessionId: string partNumbers: PresignObjectUploadPartsInput['partNumbers'] }, -): Promise<{ uploadId: string; partSize: number; parts: Array<{ partNumber: number; url: string }> }> { - const { storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) +): Promise<{ + uploadId: string | null + mode: 'multipart' + partSize: number + partCount: number + presignedExpiresAt: string + requiredHeaders: Record + parts: ObjectUploadInstructions['parts'] +}> { + const { matter, storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) const record = await deps.objectUploadSessions.get(params.orgId, params.objectId, params.sessionId) if (!record) throw new ObjectUploadSessionError('not_found') if (record.status !== 'active' || record.expiresAt.getTime() <= Date.now()) { @@ -319,13 +388,39 @@ export async function presignUploadSessionParts( // Only multipart sessions have re-presignable parts; a single PutObject does not. if (record.uploadId == null) throw new ObjectUploadSessionError('invalid_state') const uploadId = record.uploadId + const partCount = uploadPartCount(matter.size ?? 0, record.partSize) + if (new Set(params.partNumbers).size !== params.partNumbers.length) { + throw new ObjectUploadSessionError('invalid_state', 'Duplicate part numbers are not allowed') + } + for (const partNumber of params.partNumbers) { + if (partNumber < 1 || partNumber > partCount) { + throw new ObjectUploadSessionError('invalid_state', 'Part number is outside this upload session') + } + } + const presignedExpiresAt = new Date(Date.now() + UPLOAD_PRESIGNED_URL_TTL_SECONDS * 1000).toISOString() const parts = await Promise.all( params.partNumbers.map(async (partNumber) => ({ partNumber, - url: await deps.s3.presignUploadPart(storage, record.storageKey, uploadId, partNumber), + url: await deps.s3.presignUploadPart( + storage, + record.storageKey, + uploadId, + partNumber, + UPLOAD_PRESIGNED_URL_TTL_SECONDS, + ), + expiresAt: presignedExpiresAt, + headers: {}, })), ) - return { uploadId, partSize: record.partSize, parts } + return { + uploadId, + mode: 'multipart', + partSize: record.partSize, + partCount, + presignedExpiresAt, + requiredHeaders: {}, + parts, + } } export type CompleteUploadOutcome = @@ -349,20 +444,37 @@ export async function completeUpload( actorId: string }, ): Promise { - const { storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) + const { matter: sessionMatter, storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) const record = await deps.objectUploadSessions.get(params.orgId, params.objectId, params.sessionId) if (!record) throw new ObjectUploadSessionError('not_found') - if (record.status !== 'active') throw new ObjectUploadSessionError('invalid_state') + const retryingCompletedMultipartActivation = + record.status === 'completed' && record.uploadId != null && sessionMatter.status === 'draft' + if (record.status === 'completed') { + if (sessionMatter.status === 'active') return { ok: true, matter: sessionMatter } + if (!retryingCompletedMultipartActivation) throw new ObjectUploadSessionError('invalid_state') + } + if ( + (record.status !== 'active' && !retryingCompletedMultipartActivation) || + record.expiresAt.getTime() <= Date.now() + ) { + throw new ObjectUploadSessionError('invalid_state') + } - if (record.uploadId != null) { + const completedParts = validateCompletionParts( + params.parts, + uploadPartCount(sessionMatter.size ?? 0, record.partSize), + ) + + if (record.uploadId != null && !retryingCompletedMultipartActivation) { try { - await deps.s3.completeMultipartUpload(storage, record.storageKey, record.uploadId, params.parts) + await deps.s3.completeMultipartUpload(storage, record.storageKey, record.uploadId, completedParts) } catch (error) { throw new ObjectUploadSessionError( 'storage_failure', `Storage multipart upload complete failed: ${(error as Error).message}`, ) } + await deps.objectUploadSessions.setStatus(record.id, 'completed') } let head: { size: number; contentType?: string; etag: string } @@ -376,13 +488,11 @@ export async function completeUpload( ) } if (record.uploadId == null) { - const reported = params.parts[0]?.etag.replace(/"/g, '') + const reported = completedParts[0]?.etag if (!reported || reported !== head.etag) { throw new ObjectUploadSessionError('invalid_state', 'Uploaded object ETag does not match', 'etag_mismatch') } } - await deps.objectUploadSessions.setStatus(record.id, 'completed') - // Draft → live: reserve quota, apply the stored conflict strategy, activate. const { matter, quotaExceeded: exceeded } = await confirmUpload(deps, params.objectId, params.orgId, { onConflict: record.onConflict, @@ -395,20 +505,42 @@ export async function completeUpload( if (!matter) { return { ok: false, reason: 'not_found' } } + if (record.uploadId == null) { + await deps.objectUploadSessions.setStatus(record.id, 'completed') + } return { ok: true, matter } } -// Aborts an in-progress upload and discards the draft. Idempotent for an -// already-aborted session; rejects completing one. +// Aborts an in-progress upload and discards the draft. Retries finish draft +// cleanup if a previous abort marked the session first; rejects active completions. export async function abortUpload( deps: Pick, params: { orgId: string; objectId: string; sessionId: string; actorId: string; strictStorageCleanup?: boolean }, ): Promise { - const { storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) const record = await deps.objectUploadSessions.get(params.orgId, params.objectId, params.sessionId) if (!record) throw new ObjectUploadSessionError('not_found') - if (record.status === 'completed') throw new ObjectUploadSessionError('invalid_state', 'Upload already completed') - if (record.status === 'aborted') return + if (record.status === 'aborted') { + const sessionMatter = await deps.matter.get(params.objectId, params.orgId) + if (sessionMatter?.status === 'draft') { + await deps.matter.cancelDraft(params.objectId, params.orgId) + } + return + } + + const { matter: sessionMatter, storage } = await loadObjectForUploadSession(deps, params.orgId, params.objectId) + if (record.status === 'completed') { + if (sessionMatter.status !== 'draft' || record.uploadId == null) { + throw new ObjectUploadSessionError('invalid_state', 'Upload already completed') + } + try { + await deps.s3.deleteObject(storage, record.storageKey) + } catch (error) { + throw new ObjectUploadSessionError('storage_failure', `Storage cleanup failed: ${(error as Error).message}`) + } + await deps.objectUploadSessions.setStatus(record.id, 'aborted') + await deps.matter.cancelDraft(params.objectId, params.orgId) + return + } if (record.uploadId != null) { try { @@ -605,6 +737,28 @@ export async function authorizeTaskUploadConfirm( return { ok: true } } +export async function authorizeTaskUploadAbort( + deps: Pick, + params: { + orgId: string + objectId: string + sessionId: string + taskId: string + downloaderId: string + targetFolder: string + }, +): Promise { + const session = await deps.objectUploadSessions.get(params.orgId, params.objectId, params.sessionId) + if (session?.status === 'aborted') { + if (session.createdBy !== `downloader:${params.downloaderId}`) { + return { ok: false, error: forbidden() } + } + await assertTaskUploadAllowed(deps as Deps, { taskId: params.taskId, downloaderId: params.downloaderId }) + return { ok: true } + } + return authorizeTaskUploadConfirm(deps, params) +} + // ─── Permanent delete (purge) ───────────────────────────────────────────────── export type DeleteObjectOutcome = diff --git a/server/usecases/ports/object-upload-session.ts b/server/usecases/ports/object-upload-session.ts index 728a89cf..0b83adb4 100644 --- a/server/usecases/ports/object-upload-session.ts +++ b/server/usecases/ports/object-upload-session.ts @@ -5,6 +5,7 @@ import type { ObjectUploadSession } from '@shared/types' export type ObjectUploadSessionRecord = Omit & { storageKey: string + createdBy: string // The conflict strategy chosen at create time, applied when the upload is // finalized (a deferred 'replace' purges the incumbent only once bytes land). onConflict: ConflictStrategy diff --git a/shared/schemas/index.ts b/shared/schemas/index.ts index 3259ce28..3a7f9798 100644 --- a/shared/schemas/index.ts +++ b/shared/schemas/index.ts @@ -294,24 +294,36 @@ export const updateMatterSchema = z.object({ export type UpdateMatterInput = z.infer -// The upload instructions returned by POST /objects for a file draft. The server -// decides the S3 mechanism by size: 1 URL = single PutObject (≤5 GiB), N URLs = -// multipart 5 GiB parts (>5 GiB). The client PUTs each slice, reads the ETag, and -// posts them to .../completions. -export const objectUploadInstructionsSchema = z.object({ - sessionId: z.string(), - partSize: z.number().int(), - urls: z.array(z.string()), +export const presignedObjectUploadPartSchema = z.object({ + partNumber: z.number().int().min(1).max(10_000), + url: z.string(), + expiresAt: z.string(), + headers: z.record(z.string(), z.string()), }) -export const presignedObjectUploadPartSchema = z.object({ - partNumber: z.number().int(), - url: z.string(), +// The upload instructions returned by POST /objects for a file draft. File bytes +// still go directly to presigned S3 URLs; the server exposes stable part +// identities so automation never infers part numbers from URL positions. +export const objectUploadInstructionsSchema = z.object({ + sessionId: z.string(), + uploadId: z.string().nullable(), + mode: z.enum(['single', 'multipart']), + partSize: z.number().int(), + partCount: z.number().int().min(1).max(10_000), + expiresAt: z.string(), + presignedExpiresAt: z.string(), + requiredHeaders: z.record(z.string(), z.string()), + urls: z.array(z.string()), + parts: z.array(presignedObjectUploadPartSchema), }) export const presignObjectUploadPartsResponseSchema = z.object({ - uploadId: z.string(), + uploadId: z.string().nullable(), + mode: z.literal('multipart'), partSize: z.number().int(), + partCount: z.number().int().min(1).max(10_000), + presignedExpiresAt: z.string(), + requiredHeaders: z.record(z.string(), z.string()), parts: z.array(presignedObjectUploadPartSchema), }) diff --git a/shared/types/index.ts b/shared/types/index.ts index 2fb64986..38008f58 100644 --- a/shared/types/index.ts +++ b/shared/types/index.ts @@ -418,13 +418,27 @@ export interface ObjectUploadSession { updatedAt: string } +export interface ObjectUploadPartDescriptor { + partNumber: number + url: string + expiresAt: string + headers: Record +} + // The upload instructions returned by POST /objects for a file draft: the -// client PUTs each slice to urls[i] (slice i = bytes [i*partSize, …]), reads the -// ETag of each response, then POSTs them to .../completions. +// client PUTs each explicit part descriptor directly to S3, reads the ETag of +// each response, then POSTs the partNumber+etag records to .../completions. export interface ObjectUploadInstructions { sessionId: string + uploadId: string | null + mode: 'single' | 'multipart' partSize: number + partCount: number + expiresAt: string + presignedExpiresAt: string + requiredHeaders: Record urls: string[] + parts: ObjectUploadPartDescriptor[] } export type BackgroundJobStatus = 'queued' | 'running' | 'completed' | 'failed' | 'canceled' diff --git a/src/components/upload/multipart-upload.test.ts b/src/components/upload/multipart-upload.test.ts index 8bda4281..78414c5a 100644 --- a/src/components/upload/multipart-upload.test.ts +++ b/src/components/upload/multipart-upload.test.ts @@ -27,7 +27,23 @@ function makeCtx(overrides: Partial = {}): UploadRunnerCont } function makeUpload(urls: string[], partSize: number): ObjectUploadInstructions { - return { sessionId: 'sess-1', partSize, urls } + return { + sessionId: 'sess-1', + uploadId: urls.length > 1 ? 'upload-1' : null, + mode: urls.length > 1 ? 'multipart' : 'single', + partSize, + partCount: urls.length, + expiresAt: '2026-01-01T01:00:00.000Z', + presignedExpiresAt: '2026-01-01T00:15:00.000Z', + requiredHeaders: {}, + urls, + parts: urls.map((url, index) => ({ + partNumber: index + 1, + url, + expiresAt: '2026-01-01T00:15:00.000Z', + headers: {}, + })), + } } describe('uploadObjectSlices', () => { diff --git a/src/components/upload/multipart-upload.ts b/src/components/upload/multipart-upload.ts index 537ae0dd..fab80223 100644 --- a/src/components/upload/multipart-upload.ts +++ b/src/components/upload/multipart-upload.ts @@ -46,8 +46,7 @@ async function runPool(items: T[], concurrency: number, worker: (item: T) => } /** - * The uniform upload: PUT every presigned slice directly to S3 (1 URL = single - * PutObject, N URLs = 5 GiB-part multipart — the server already decided), read + * The uniform upload: PUT every explicit presigned part directly to S3, read * each slice's ETag, and return them sorted for the completions call. Bounded * concurrency, per-slice retry, aggregated progress. Bytes never touch our server. */ @@ -56,7 +55,7 @@ export async function uploadObjectSlices( file: File, ctx: UploadRunnerContext, ): Promise> { - const { partSize, urls } = upload + const { partSize, parts: uploadParts } = upload const completed: Array<{ partNumber: number; etag: string }> = [] const loadedByPart = new Map() @@ -66,14 +65,13 @@ export async function uploadObjectSlices( ctx.onProgress({ loaded, total: file.size }) } - const slices = urls.map((url, index) => ({ url, partNumber: index + 1 })) - await runPool(slices, PART_CONCURRENCY, async ({ url, partNumber }) => { + await runPool(uploadParts, PART_CONCURRENCY, async ({ url, partNumber, headers }) => { if (ctx.signal.aborted) throw new DOMException('Upload cancelled', 'AbortError') const start = (partNumber - 1) * partSize const slice = file.slice(start, Math.min(start + partSize, file.size)) const etag = await uploadPartWithRetry(url, slice, { signal: ctx.signal, - contentType: file.type || undefined, + contentType: headers['content-type'] ?? (file.type || undefined), onProgress: (loaded) => { loadedByPart.set(partNumber, loaded) reportProgress() diff --git a/src/components/upload/upload-dropzone.test.ts b/src/components/upload/upload-dropzone.test.ts index 54ad4d5d..65ea16bf 100644 --- a/src/components/upload/upload-dropzone.test.ts +++ b/src/components/upload/upload-dropzone.test.ts @@ -252,7 +252,17 @@ describe('UploadDropzone — uploadFn path', () => { describe('UploadDropzone — default upload path (no uploadFn)', () => { const draft = { id: 'obj-1', - upload: { sessionId: 'sess-1', partSize: 5 * 1024 * 1024, urls: ['https://s3/part-1'] }, + upload: { + sessionId: 'sess-1', + uploadId: null, + mode: 'single', + partSize: 5 * 1024 * 1024, + partCount: 1, + expiresAt: '2026-01-01T01:00:00.000Z', + presignedExpiresAt: '2026-01-01T00:15:00.000Z', + requiredHeaders: {}, + parts: [{ partNumber: 1, url: 'https://s3/part-1', expiresAt: '2026-01-01T00:15:00.000Z', headers: {} }], + }, } const parts = [{ partNumber: 1, etag: 'etag-1' }] diff --git a/src/lib/api.test.ts b/src/lib/api.test.ts index b5f6b36d..e0c28cac 100644 --- a/src/lib/api.test.ts +++ b/src/lib/api.test.ts @@ -508,7 +508,17 @@ describe('api', () => { const created = { id: 'new1', name: 'doc.pdf', - upload: { sessionId: 'sess-1', partSize: 5 * 1024 * 1024, urls: ['https://s3/part-1'] }, + upload: { + sessionId: 'sess-1', + uploadId: null, + mode: 'single', + partSize: 5 * 1024 * 1024, + partCount: 1, + expiresAt: '2026-01-01T01:00:00.000Z', + presignedExpiresAt: '2026-01-01T00:15:00.000Z', + requiredHeaders: {}, + parts: [{ partNumber: 1, url: 'https://s3/part-1', expiresAt: '2026-01-01T00:15:00.000Z', headers: {} }], + }, } vi.mocked(fetch).mockResolvedValueOnce(makeResponse(created)) @@ -521,7 +531,7 @@ describe('api', () => { }) expect(result).toEqual(created) - expect(result.upload).toEqual({ sessionId: 'sess-1', partSize: 5 * 1024 * 1024, urls: ['https://s3/part-1'] }) + expect(result.upload).toEqual(created.upload) const [url, init] = vi.mocked(fetch).mock.calls[0] as [string, RequestInit] expect(url).toContain('/api/objects') expect(init.method).toBe('POST') diff --git a/src/lib/api.ts b/src/lib/api.ts index 0d41ad6e..b9329f6d 100644 --- a/src/lib/api.ts +++ b/src/lib/api.ts @@ -352,7 +352,15 @@ export function abortObjectUpload(id: string, uploadSessionId: string, opts: { s // Re-presign expired part URLs mid-upload (multipart only); the happy path uses // the URLs from createObject. export function presignObjectUploadParts(id: string, uploadSessionId: string, data: PresignObjectUploadPartsInput) { - return unwrap<{ uploadId: string; partSize: number; parts: Array<{ partNumber: number; url: string }> }>( + return unwrap<{ + uploadId: string | null + mode: 'multipart' + partSize: number + partCount: number + presignedExpiresAt: string + requiredHeaders: Record + parts: ObjectUploadInstructions['parts'] + }>( objects[':id'].uploads[':uploadSessionId'].parts.$post({ param: { id, uploadSessionId }, json: data, diff --git a/src/routes/_authenticated/admin/storages/index.test.tsx b/src/routes/_authenticated/admin/storages/index.test.tsx index 235ab3d0..96ff9824 100644 --- a/src/routes/_authenticated/admin/storages/index.test.tsx +++ b/src/routes/_authenticated/admin/storages/index.test.tsx @@ -101,7 +101,25 @@ const uploadDraft: CreateObjectResult = { trashedAt: null, createdAt: '2026-01-01T00:00:00.000Z', updatedAt: '2026-01-01T00:00:00.000Z', - upload: { sessionId: 'session-1', urls: ['https://uploads.example.com/object-1'], partSize: 5_242_880 }, + upload: { + sessionId: 'session-1', + uploadId: null, + mode: 'single', + partSize: 5_242_880, + partCount: 1, + expiresAt: '2026-01-01T01:00:00.000Z', + presignedExpiresAt: '2026-01-01T00:15:00.000Z', + requiredHeaders: { 'content-type': 'text/plain' }, + urls: ['https://uploads.example.com/object-1'], + parts: [ + { + partNumber: 1, + url: 'https://uploads.example.com/object-1', + expiresAt: '2026-01-01T00:15:00.000Z', + headers: { 'content-type': 'text/plain' }, + }, + ], + }, } function renderStoragesPage() { @@ -236,7 +254,7 @@ describe('StoragesPage connection test action', () => { 'https://uploads.example.com/object-1', expect.objectContaining({ method: 'PUT', - headers: { 'Content-Type': 'text/plain' }, + headers: { 'content-type': 'text/plain' }, body: expect.any(Blob), }), ) diff --git a/src/routes/_authenticated/admin/storages/index.tsx b/src/routes/_authenticated/admin/storages/index.tsx index dfe5c75f..f441bc0e 100644 --- a/src/routes/_authenticated/admin/storages/index.tsx +++ b/src/routes/_authenticated/admin/storages/index.tsx @@ -234,7 +234,7 @@ export function StoragesPage() { setTestHealth({ status: 'testing' }) setTestStep('creating') setTestFailedStep(null) - let draft: { id: string; upload?: { sessionId: string; urls: string[] } } | null = null + let draft: Awaited> | null = null let result: StorageHealth | null = null let currentStep: StorageTestStep = 'creating' let failedStep: StorageTestStep | null = null @@ -256,13 +256,14 @@ export function StoragesPage() { const upload = draft.upload setCurrentStep('uploading') - if (!upload?.urls[0]) throw new Error(t('admin.storages.testNoUploadUrl')) + const uploadPart = upload?.parts[0] + if (!uploadPart) throw new Error(t('admin.storages.testNoUploadUrl')) let uploadResponse: Response try { - uploadResponse = await fetch(upload.urls[0], { + uploadResponse = await fetch(uploadPart.url, { method: 'PUT', - headers: { 'Content-Type': 'text/plain' }, + headers: uploadPart.headers, body: blob, }) } catch {