feat: stabilize multipart upload protocol (#537)

Preserve legacy upload clients while adding explicit resumable parts, safe completion recovery, idempotent cleanup, and downloader-bound authorization.
This commit is contained in:
agent-kanban[bot]
2026-07-29 10:37:31 -04:00
committed by GitHub
parent 1b1b1db772
commit f2aea1bedb
24 changed files with 1266 additions and 156 deletions
+41 -13
View File
@@ -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
}
+25 -4
View File
@@ -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) {
+22 -4
View File
@@ -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)
+7 -7
View File
@@ -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)
+21 -10
View File
@@ -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
+91 -16
View File
@@ -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
+62
View File
@@ -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.
@@ -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,
@@ -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)
+2 -2
View File
@@ -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(
+117 -6
View File
@@ -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<string, string> }> }
}
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
})
})
+40 -15
View File
@@ -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<Env>): string {
return c.get('userId') ?? 'system'
}
async function authorizeUploadSessionControl(
c: Context<Env>,
orgId: string,
objectId: string,
options: { uploadSessionId?: string } = {},
): Promise<void> {
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<Env>) => 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,
})
+523 -12
View File
@@ -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<S3Gateway>,
@@ -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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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<S3Gateway>,
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')
})
})
})
+183 -29
View File
@@ -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<string, string> {
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<number>()
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<string, string>
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<CompleteUploadOutcome> {
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<Deps, 'matter' | 'storages' | 's3' | 'objectUploadSessions'>,
params: { orgId: string; objectId: string; sessionId: string; actorId: string; strictStorageCleanup?: boolean },
): Promise<void> {
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<Deps, 'matter' | 'downloaders' | 'downloadTasks' | 'objectUploadSessions'>,
params: {
orgId: string
objectId: string
sessionId: string
taskId: string
downloaderId: string
targetFolder: string
},
): Promise<ConfirmAuthorizationOutcome> {
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 =
@@ -5,6 +5,7 @@ import type { ObjectUploadSession } from '@shared/types'
export type ObjectUploadSessionRecord = Omit<ObjectUploadSession, 'expiresAt' | 'createdAt' | 'updatedAt'> & {
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
+24 -12
View File
@@ -294,24 +294,36 @@ export const updateMatterSchema = z.object({
export type UpdateMatterInput = z.infer<typeof updateMatterSchema>
// 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),
})
+16 -2
View File
@@ -418,13 +418,27 @@ export interface ObjectUploadSession {
updatedAt: string
}
export interface ObjectUploadPartDescriptor {
partNumber: number
url: string
expiresAt: string
headers: Record<string, string>
}
// 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<string, string>
urls: string[]
parts: ObjectUploadPartDescriptor[]
}
export type BackgroundJobStatus = 'queued' | 'running' | 'completed' | 'failed' | 'canceled'
+17 -1
View File
@@ -27,7 +27,23 @@ function makeCtx(overrides: Partial<UploadRunnerContext> = {}): 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', () => {
+4 -6
View File
@@ -46,8 +46,7 @@ async function runPool<T>(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<Array<{ partNumber: number; etag: string }>> {
const { partSize, urls } = upload
const { partSize, parts: uploadParts } = upload
const completed: Array<{ partNumber: number; etag: string }> = []
const loadedByPart = new Map<number, number>()
@@ -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()
+11 -1
View File
@@ -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' }]
+12 -2
View File
@@ -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')
+9 -1
View File
@@ -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<string, string>
parts: ObjectUploadInstructions['parts']
}>(
objects[':id'].uploads[':uploadSessionId'].parts.$post({
param: { id, uploadSessionId },
json: data,
@@ -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),
}),
)
@@ -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<ReturnType<typeof createObject>> | 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 {