Skip to content

Commit 54c3ce1

Browse files
fix(service-storage): chunks sent in parallel to one chunked upload each record their part (#22332) (#22348)
Fixes #22332 Clause-②: no ## What was wrong The chunk door (`PUT /storage/upload/chunked/:uploadId/chunk/:chunkIndex`) read the session row, merged its part into the parts list in memory (`recordChunk`) and wrote the whole record back with an unconditional by-id update. Two PUTs to one upload at once both read the same record, and the later write erased the earlier part: both answered `200`, `sys_upload_session.parts` kept one part, `uploaded_chunks` / `uploaded_size` counted one chunk, and the completion door refused the upload `409` naming the lost chunk. ## The atomic route taken, and why **A compare-and-set on the row's three progress columns, with a bounded re-read-and-merge retry.** This is the route the store already supports; it needs no new object, no new field and no `packages/spec` change. - **On a wired engine** it is the engine's own conditional update: `engine.update('sys_upload_session', progress, { where: { parts, uploaded_chunks, uploaded_size, id }, multi: true, context })`. The update-dispatch module (`resolveEngineUpdateDispatch`, `@objectstack/metadata-core`) names exactly this shape, an id in `where` beside further keys with `multi: true`, as "the compare-and-set spelling". It routes to `driver.updateMany`, which evaluates the guard in the same statement that writes and answers the matched-row count. driver-sql issues one `UPDATE … WHERE` with the tenant scope applied; driver-memory matches and writes in one synchronous step. Every driver's `updateMany` (sql, memory, turso, mongodb) resolves a count. - **On the engine-absent stand-in** it is one synchronous compare-and-set on the Map. - **Why these columns are the guard:** they are exactly the values the merge read, so the comparison is exact. `updated_at` as a version has millisecond grain, so two writes in one millisecond would still match. A per-part row would be a new object, a schema change, and was not needed. - ⛔ No process-local lock: a second server process would defeat it. ⛔ The completion door's guard is untouched. The store side is `updateSessionProgressIfUnchanged` in `metadata-store.ts`. It is a module export that is not re-exported from the package entry, following the `organizationOutOfWriteReach` pattern, so `StorageMetadataStore`'s public face does not grow. The door merges and writes up to `CHUNK_RECORD_ATTEMPTS` (16) times. A conditional write loses only to a write that landed between its read and itself, so each lost attempt is another chunk recorded. - **On exhaustion** the door answers `409 RESOURCE_CONFLICT` with `error.details: { chunkIndex, attempts }`. The chunk's bytes are stored but the record does not hold them, and re-sending the chunk replaces its slot. It never answers a silent `200`. - **When the conditional write cannot reach the row** (it misses, and the row still holds the progress the write was conditioned on), the door answers `500` at once. It does not loop into a `409` that would tell the uploader to retry something that cannot succeed. ## Measured: what the store supports (H2) A probe against the real `ObjectQL` over `SqlDriver` (better-sqlite3, `:memory:`) at `e36ee5351b`: | conditional update | answer | | :--- | ---: | | guard as read | 1 | | stale guard | 0 | | guard on an 8-part `parts` string | 1 | | guard of NULL columns (lowered to `IS NULL`) | 1 | | row stamped org_B, acting organization org_A | 0, no throw | | the same row, acting organization org_B | 1 | | missing row | 0 | ## Measured: the reproduction, before and after `pnpm --filter @objectstack/service-storage exec vitest run --maxWorkers=2 src/chunk-part-record-concurrency.test.ts` **Before**, the pins on `main` at `e36ee5351b` (which carries the completion guard): 8 failed and 2 passed; the 2 that passed are the sequential controls. - The stand-in and the real engine alike lost parts: `every concurrently sent chunk is in the record: expected [ 1 ] to deeply equal [ +0, 1 ]` (2 concurrent), and `expected [ +0 ] to deeply equal [ +0, 1, 2, 3, 4, 5, 6, 7 ]` (8 concurrent). - The parallel upload's completion was refused `409` with `missingChunks: [1]`. - With a competing write landing between the read and the write, the door answered `{"success":true,…}` (`expected 200 to be 409`), and the record held `[ +0 ]`: the competitor's part was erased. **After**, at `006344eb50`: 24 passed (24). Both stores pass the 2 and 8 concurrent pins (every part, every eTag and size, `uploaded_chunks`, `uploaded_size`, and the progress door), the parallel completion `200` with the whole file, the sequential control (the same answers and record, one session write per chunk), and the exhaustion pin (`409 RESOURCE_CONFLICT` after exactly 16 reads, with every competitor's part kept). The store-level pins also pass on both stores: the write lands on an unchanged row, misses once any one progress column moved, and misses on a gone row. On the real engine, the call shape is pinned (dispatch verdict `multi`, the same `{ tenantId, isSystem }` context), another organization's row is not reached, a non-count answer is refused loudly, and an unreachable row answers `500` once. **Over HTTP**, a real boot (`bootStack`, sqlite-wasm): `packages/qa/dogfood/test/storage-chunked-parallel-parts.dogfood.test.ts` passes 2 of 2 (2 concurrent PUTs, then progress and completion `200`; 8 concurrent PUTs, then progress). With the sibling `storage-chunked-resume-integrity.dogfood.test.ts`, 6 passed (6) at `006344eb50`. ## Ablations Each mutation went through `scripts/ablation-replace.mjs`: the anchor hit 1 to 0, the blob changed, and the restore was proven with blob == HEAD and an empty `git diff HEAD`. The runner also carried its own `EXIT`/`INT`/`TERM` restore over absolute paths. The control run with no mutation passed 24 of 24. | mutation | pins turned red | | :--- | :--- | | M1: the stand-in's comparison removed | 7, all stand-in: 2 and 8 concurrent, the parallel completion, exhaustion, and 3 column-moved pins | | M2: the engine guard replaced by an unchanging column | 8, all engine: 2 and 8 concurrent, the parallel completion, exhaustion, 3 column-moved pins, and the where-shape pin | | M3: the door treats a lost write as landed | 9: on both stores, 2 and 8 concurrent, the parallel completion and exhaustion; plus the unreachable-row `500` | | M4: exhaustion falls through to `200` | 2: exhaustion, on both stores | | M5: a non-count engine answer accepted | 1 | | M6: the unreachable-row check removed | 1 | | M7: the stand-in lands on a missing row | 1 | | M8: the organization scope dropped from the conditional write | 2 | The dogfood pin reads `@objectstack/service-storage` from `dist/`, so its ablation included a rebuild. - **Mutation leg:** M3, then a rebuild. `ablation-dist-preflight` found the marker in `dist/index.js` and `dist/index.cjs`, and 2 of 2 failed (the progress fell short of `uploadedChunks`). - **Restore leg:** a rebuild, after which the marker was absent from all 6 built files and the tree was clean against HEAD. 2 of 2 passed. ## Tests and gates (HEAD `006344eb50`) - `@objectstack/service-storage`: 46 files and 773 tests passed. `typecheck` passed (`tsc`, the scripts project, and `check:test-typecheck: OK`). `--listFiles` shows the new test file is in the test-layer program. - `@objectstack/dogfood`: `typecheck` passed. - Two fake engines now answer a declared predicate update with its matched-row count, as the real engine does: `tenant-audit-update-delete-half-repairs.test.ts` and `storage-routes.metadata-outage.test.ts`. The chunk-door pin in the first now expects the conditional `where` with `multi: true`, under the same context. - `check:tenant-audit-census`: the new conditional write is one more engine write call site (236 to 237). The census was regenerated (`node scripts/tenant-audit-census.mjs --write`), and the page's hand-written prose figures were moved with it. The gate and its self-test pass. - Lint, narrowed and stated as a measurement. ① Population: the 6 touched TypeScript files, all inside the `eslint.config.mjs` glob `**/*.{ts,tsx,mts,cts,js,jsx,mjs,cjs}`. ② `--format json` linted 6 files with 0 errors and 0 warnings. ③ Invariance: the config sets no `parserOptions.project` (`--print-config` gives `{"ecmaVersion":"latest","sourceType":"module"}`), so linting is not type-aware and this diff cannot move a verdict on an untouched file. The repo-wide `pnpm lint` is CI's. - The derived gate families (`dispatch-gates --commands`, derived at `006344eb50`) were run after the last commit, with exit codes captured before any pipe. `dispatch-gates --ran` reads `97 derived, 97 run, 0 NOT-MEASURED, 0 UNRUN` (a derived zero: every family recorded `exit 0`). `pnpm check:error-status-conformance` was run by hand and exited 0: `✓ every derivable runtime status is documented, and every documented status is reachable.` A few verdict lines: - `check-engine-double-contract: OK — 983 pinned, 129 in the DEBT ledger, 3 exempt.` - `check-nul-bytes: OK (scanned 10310 text file(s) …; no raw ASCII control bytes).` - `✓ check-tenant-audit-census: OK -- 237 write call sites certified …` - `check-test-source-alias OK — 73 packages with tests scanned …` - `✓ check:dual-build-cjs-loads — 106 published require entry point(s) across 66 package(s) load …`. Its first run answered `PREREQUISITE NOT MET` (no `dist/` for 8 packages), which is not a measurement. It passed after a full cache-backed `turbo run build`. - ✓ This diff introduces no `major` bump. (`check-changeset-no-major`) ## Acceptance notes - **Clause-② stays `no`, as the claim declared.** No accept set, export or public signature moves. The one new answer, `409` after 16 lost record writes, replaces a `200` that had recorded nothing. - **The progress write is now a predicate update.** On an engine-backed store, the chunk door's write publishes the bulk `data.records.updated` event (a count) where it published `data.record.updated`, and its after-hooks go through the bulk per-row path. Nothing in the tree registers an update hook on `sys_upload_session` or subscribes to its record events, and the config-change audit excludes the object. - **Observation, read-only and not measured, not filed** (承接者:无). The chunk door's upload-limit check judges the record as it was read. Concurrent PUTs can therefore each pass it against the same stale total, and the backend can briefly hold more bytes than `maxUploadBytes`. The completion door still refuses any upload whose held bytes differ from its declared size, and the init door bounds that size. --- _Generated by [Claude Code](https://claude.ai/code/session_01WkL6Eijt432S1Y7ekb6ovQ)_ --------- Co-authored-by: Claude <noreply@anthropic.com>
1 parent c6fc938 commit 54c3ce1

9 files changed

Lines changed: 912 additions & 46 deletions
Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,13 @@
1+
---
2+
'@objectstack/service-storage': patch
3+
---
4+
5+
fix(service-storage): chunks sent in parallel to one chunked upload each record their part
6+
7+
Clause-②: no
8+
9+
`PUT /storage/upload/chunked/:uploadId/chunk/:chunkIndex` merges its chunk into the upload session's record of the chunks it holds (`parts`, `uploaded_chunks`, `uploaded_size` on `sys_upload_session`). It read that record, merged in memory and wrote the whole record back, so two chunk PUTs to one upload at the same time both answered `200` while the record kept only one of them: `GET …/progress` undercounted, and the completion was refused `409 RESOURCE_CONFLICT` naming the chunk the record had lost until the client sent it again.
10+
11+
The record is now written with a compare-and-set: the write lands only while the row still holds the progress the chunk door read, and when another chunk's write landed first the door reads the record again and merges again. On a wired data engine this is the engine's own conditional update, evaluated in the same statement that writes, so it holds across server processes. Every chunk sent in parallel is recorded, and a parallel upload completes on its first completion. A sequential upload is unchanged.
12+
13+
A chunk whose record write loses to another write on 16 attempts in a row is refused `409 RESOURCE_CONFLICT`, with `error.details` `{ chunkIndex, attempts }`: its bytes are stored, but the upload does not hold it. Send that chunk again.

‎content/docs/permissions/tenant-audit-census.mdx‎

Lines changed: 23 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,7 @@ what moved this page's population from 225 to 227; nothing about the two sites
8383
changed, only whether this instrument could see them.
8484

8585
**The expensive failure direction is a keyword.** Sites whose receiver the author
86-
typed `any` have no type to read, and there are 50 of them — just over a fifth
86+
typed `any` have no type to read, and there are 51 of them — just over a fifth
8787
of the population, concentrated in exactly the seed and bootstrap paths this
8888
control exists for. Scoring an unreadable receiver as "not an engine" would have
8989
dropped every one of them silently, with a clean exit and a smaller number that
@@ -122,7 +122,7 @@ are reported as `undecidable` rather than assumed either way.
122122

123123
The same holds twice over for the context. An options argument spelled as a
124124
literal can be read; one spelled `options`, `{ ...opts }`, or handed through a
125-
forwarding shim cannot, and **52 of the 236 sites are spelled that way**. A
125+
forwarding shim cannot, and **53 of the 237 sites are spelled that way**. A
126126
context resolved from an inline literal or a local `const` can be tested for
127127
`isSystem`; one arriving from a helper call cannot.
128128

@@ -150,8 +150,8 @@ now **0**: nothing on this surface threads a context that provably lacks the fla
150150

151151
**"No tenant context" counted sites it had not read.** An options argument the
152152
walker could not parse was folded into the same bucket as one it had read and
153-
found empty. That published **60 sites "carrying no tenant context at all"**
154-
when 8 said so and 52 were simply unread — an over-claim in the *alarming*
153+
found empty. That published **61 sites "carrying no tenant context at all"**
154+
when 8 said so and 53 were simply unread — an over-claim in the *alarming*
155155
direction, on the very figure this page tells other cards to cite. `carries` is
156156
now three-valued, and an unreadable argument can never contribute to the
157157
provable count.
@@ -187,10 +187,10 @@ reproduce them. Where it disagrees, it disagrees on the page:
187187

188188
| carried figure | where it survives | this census |
189189
| :--- | :--- | ---: |
190-
| 175 write call sites | quoted in the merged changeset | **236** |
191-
| 24 carrying no tenant context | quoted in the merged changeset | **2** provable and tenancy-enabled; **31** more whose options argument is unreadable |
192-
| 127 of 175 statically decidable, 48 runtime-parameter-name sites | restated on the `isSystem`-scoping card | **156 of 236** decidable, **80** undecidable |
193-
| 135 (77%) silenced by the `isSystem` guard before the posture gate | the lost issue body — **no surviving corroboration** | **not reproduced**: 125 decidably elevated, 0 decidably not, 103 undecidable |
190+
| 175 write call sites | quoted in the merged changeset | **237** |
191+
| 24 carrying no tenant context | quoted in the merged changeset | **2** provable and tenancy-enabled; **32** more whose options argument is unreadable |
192+
| 127 of 175 statically decidable, 48 runtime-parameter-name sites | restated on the `isSystem`-scoping card | **157 of 237** decidable, **80** undecidable |
193+
| 135 (77%) silenced by the `isSystem` guard before the posture gate | the lost issue body — **no surviving corroboration** | **not reproduced**: 125 decidably elevated, 0 decidably not, 104 undecidable |
194194
| 141 and 132, two independent re-derivations | the card that filed this work | — |
195195

196196
**The differences are not reconciled, and deliberately so.** The old census's
@@ -200,21 +200,21 @@ be stated is what this instrument counts, which is written above and re-runnable
200200
at any commit.
201201

202202
Two structural facts do plausibly widen this reading against any hand or regex
203-
one, and both are counted in the generated tables below: the 50 sites reached
203+
one, and both are counted in the generated tables below: the 51 sites reached
204204
through an erased (`any`) receiver, and the 50 that name their object through a
205205
`const` rather than inline. An instrument that read either the way a person does
206206
would report a smaller number and would not say so.
207207

208208
The fourth row is the one worth flagging to anyone citing it. **The 135 / 77%
209209
figure has no surviving corroboration anywhere in the tree.** This census reads
210-
125 of 236 (53%) as decidably elevated, with 103 more whose elevation is a
210+
125 of 237 (53%) as decidably elevated, with 104 more whose elevation is a
211211
run-time fact — so the claim is neither confirmed nor refuted, and the honest
212212
answer is that a static reading cannot settle it.
213213

214-
⇒ **Cite `2 / 236`, and say what it is**: the sites whose options argument was
214+
⇒ **Cite `2 / 237`, and say what it is**: the sites whose options argument was
215215
READ and holds no tenant context, against a decidably tenancy-enabled object.
216216
That is the control's provable yield surface. ⛔ Do not cite it as "the sites
217-
without tenant context" — **31 further sites** have an options argument this
217+
without tenant context" — **32 further sites** have an options argument this
218218
cannot read, and they are neither in nor out.
219219

220220
{/* BEGIN GENERATED: tenant-audit-census (scripts/tenant-audit-census.mjs) — DO NOT EDIT */}
@@ -223,28 +223,28 @@ cannot read, and they are neither in nor out.
223223

224224
| what | count |
225225
| :--- | ---: |
226-
| write call sites on the application surface | **236** |
227-
| …whose object name is statically decidable | 156 |
226+
| write call sites on the application surface | **237** |
227+
| …whose object name is statically decidable | 157 |
228228
| …whose object name is chosen at run time | 80 |
229-
| …against an object with tenancy ENABLED | 155 |
229+
| …against an object with tenancy ENABLED | 156 |
230230
| …against an object that declares tenancy off | 1 |
231231
| threading a tenant context | 176 |
232232
| PROVABLY carrying none (options read, no context key) | **8** |
233233
| …of those, against a decidably tenancy-enabled object | **2** |
234-
| options argument UNREADABLE — may or may not carry one | 52 |
235-
| …of those, against a decidably tenancy-enabled object | 31 |
234+
| options argument UNREADABLE — may or may not carry one | 53 |
235+
| …of those, against a decidably tenancy-enabled object | 32 |
236236
| threading a decidably ELEVATED (`isSystem`) context | 125 |
237237
| threading a context that is decidably NOT elevated | 0 |
238-
| threading a context whose elevation is a run-time fact | 103 |
238+
| threading a context whose elevation is a run-time fact | 104 |
239239

240240
| how the instrument reached the site | count |
241241
| :--- | ---: |
242242
| receiver carried a readable engine type | 186 |
243-
| receiver erased, placed by the object NAME | 30 |
243+
| receiver erased, placed by the object NAME | 31 |
244244
| receiver erased, placed by an `object: string` PARAMETER | 15 |
245245
| receiver erased, placed by an `UNTYPED_RECEIVERS` row | 5 |
246246

247-
| object name spelled inline | 106 |
247+
| object name spelled inline | 107 |
248248
| object name spelled through a `const` | 50 |
249249
| object name is an `object: string` parameter | 19 |
250250
| object name is some other run-time expression | 61 |
@@ -297,13 +297,13 @@ holds still. They are required to be HERE and to say WHEN they were true;
297297
their values are not compared. The reasoning, and the measurement behind it,
298298
are in `scripts/check-tenant-audit-census.mjs`.
299299

300-
Measured on 2026-10-08 at `4b40ca31f`.
300+
Measured on 2026-10-08 at `8e432893f`.
301301

302302
| corpus scale (not enforced) | count |
303303
| :--- | ---: |
304-
| tracked non-test sources scanned | 618 |
304+
| tracked non-test sources scanned | 621 |
305305
| engine-shaped types recognised | 70 |
306306
| declared objects in the registry | 117 |
307-
| same-named calls subtracted as non-engine | 161 |
307+
| same-named calls subtracted as non-engine | 162 |
308308

309309
{/* END GENERATED: tenant-audit-census */}

‎docs/audits/2026-08-tenant-audit-write-call-sites.counts.md‎

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -33,19 +33,19 @@ silent, and `node scripts/tenant-audit-census.mjs --write` is the resolution.
3333

3434
| Measure | Value |
3535
|---|---:|
36-
| Write call sites | 236 |
37-
| Object name statically decidable | 156 |
36+
| Write call sites | 237 |
37+
| Object name statically decidable | 157 |
3838
| Object name chosen at run time | 80 |
39-
| Against a tenancy-enabled object | 155 |
39+
| Against a tenancy-enabled object | 156 |
4040
| Against an object declaring tenancy off | 1 |
4141
| Threading a tenant context | 176 |
4242
| Provably carrying none | 8 |
4343
| …and decidably tenancy-enabled | 2 |
44-
| Options argument unreadable | 52 |
45-
| …and decidably tenancy-enabled | 31 |
44+
| Options argument unreadable | 53 |
45+
| …and decidably tenancy-enabled | 32 |
4646
| Threading a decidably elevated context | 125 |
4747
| Threading a decidably non-elevated context | 0 |
48-
| Threading a context of undecidable elevation | 103 |
48+
| Threading a context of undecidable elevation | 104 |
4949

5050
## Subtractions the census could NOT defend — enforced
5151

@@ -90,14 +90,14 @@ holds still. They are required to be HERE and to say WHEN they were true;
9090
their values are not compared. The reasoning, and the measurement behind it,
9191
are in `scripts/check-tenant-audit-census.mjs`.
9292

93-
Measured on 2026-10-08 at `4b40ca31f`.
93+
Measured on 2026-10-08 at `8e432893f`.
9494

9595
| corpus scale (not enforced) | count |
9696
| :--- | ---: |
97-
| tracked non-test sources scanned | 618 |
97+
| tracked non-test sources scanned | 621 |
9898
| engine-shaped types recognised | 70 |
9999
| declared objects in the registry | 117 |
100-
| same-named calls subtracted as non-engine | 161 |
100+
| same-named calls subtracted as non-engine | 162 |
101101

102102
## Every site
103103

@@ -255,4 +255,4 @@ Measured on 2026-10-08 at `4b40ca31f`.
255255
| `packages/services/service-storage/src/metadata-store.ts` | `update` | `sys_file` | enabled | options unreadable | 1 |
256256
| `packages/services/service-storage/src/metadata-store.ts` | `delete` | `sys_upload_session` | enabled | options unreadable | 1 |
257257
| `packages/services/service-storage/src/metadata-store.ts` | `insert` | `sys_upload_session` | enabled | options unreadable | 1 |
258-
| `packages/services/service-storage/src/metadata-store.ts` | `update` | `sys_upload_session` | enabled | options unreadable | 1 |
258+
| `packages/services/service-storage/src/metadata-store.ts` | `update` | `sys_upload_session` | enabled | options unreadable | 2 |
Lines changed: 117 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,117 @@
1+
// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license.
2+
//
3+
// [#22332] Chunk PUTs sent in parallel to one upload each record their part —
4+
// over the HTTP doors of a REAL boot.
5+
//
6+
// The package suite (`service-storage/src/chunk-part-record-concurrency.test.ts`)
7+
// pins the rule at the route handlers, over the engine-absent stand-in and a
8+
// real ObjectQL on SqlDriver. This file sends the chunks the way a parallel
9+
// (S3-style) uploader does — every PUT in flight at once — through the booted
10+
// stack's own engine and driver, then reads the progress back and completes the
11+
// upload listing every part. Before the fix both PUTs answered `200`, the
12+
// progress counted one chunk, and the completion was refused `409` naming the
13+
// chunk the record had lost.
14+
//
15+
// Not eligible for the shared showcase project: it boots its own storage root.
16+
17+
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
18+
import { mkdtempSync } from 'node:fs';
19+
import { join } from 'node:path';
20+
import { tmpdir } from 'node:os';
21+
import { bootStack, type VerifyStack } from '@objectstack/verify';
22+
import { StorageServicePlugin } from '@objectstack/service-storage';
23+
import { attachmentsFixtureStack, attachmentsFixtureSecurity } from './fixtures/attachments-fixture.js';
24+
25+
const MIB = 1024 * 1024;
26+
/** The chunk size the init door floors every upload at. */
27+
const MIN_CHUNK = 5 * MIB;
28+
29+
describe('[#22332] parallel chunk PUTs each record their part, over a real boot', () => {
30+
let stack: VerifyStack;
31+
let token: string;
32+
33+
const json = { 'Content-Type': 'application/json' };
34+
const auth = () => ({ Authorization: `Bearer ${token}` });
35+
36+
const init = async (totalSize: number) => {
37+
const res = await stack.api('/storage/upload/chunked', {
38+
method: 'POST',
39+
headers: { ...json, ...auth() },
40+
body: JSON.stringify({ filename: 'parallel.bin', mimeType: 'application/octet-stream', totalSize }),
41+
});
42+
expect(res.status, await res.clone().text()).toBe(200);
43+
return ((await res.json()) as any).data as { uploadId: string; resumeToken: string; fileId: string; totalChunks: number };
44+
};
45+
46+
const putChunk = async (uploadId: string, resumeToken: string, chunkIndex: number, bytes: Uint8Array<ArrayBuffer>) => {
47+
const res = await stack.api(`/storage/upload/chunked/${uploadId}/chunk/${chunkIndex}`, {
48+
method: 'PUT',
49+
headers: { ...auth(), 'Content-Type': 'application/octet-stream', 'x-resume-token': resumeToken },
50+
body: bytes,
51+
});
52+
expect(res.status, await res.clone().text()).toBe(200);
53+
return ((await res.json()) as any).data.eTag as string;
54+
};
55+
56+
const progress = async (uploadId: string) => {
57+
const res = await stack.api(`/storage/upload/chunked/${uploadId}/progress`, { headers: auth() });
58+
expect(res.status, await res.clone().text()).toBe(200);
59+
return ((await res.json()) as any).data as Record<string, unknown>;
60+
};
61+
62+
const fill = (byte: number, n: number) => new Uint8Array(n).fill(byte);
63+
64+
beforeAll(async () => {
65+
const rootDir = mkdtempSync(join(tmpdir(), 'chunked-parallel-dogfood-'));
66+
stack = await bootStack(attachmentsFixtureStack as never, {
67+
security: attachmentsFixtureSecurity(),
68+
// `bindToSettings:false` keeps the constructor rootDir; the size limit is
69+
// not what this file is about.
70+
extraPlugins: [new StorageServicePlugin({ adapter: 'local', local: { rootDir }, bindToSettings: false })],
71+
});
72+
token = await stack.signIn();
73+
}, 120_000);
74+
75+
afterAll(async () => {
76+
await stack?.stop();
77+
});
78+
79+
it('two chunks sent at once are both recorded, and the upload completes with the whole file', async () => {
80+
const chunk0 = fill(0x61, MIN_CHUNK);
81+
const chunk1 = fill(0x62, 10);
82+
const { uploadId, resumeToken, fileId } = await init(MIN_CHUNK + 10);
83+
84+
const [eTag0, eTag1] = await Promise.all([
85+
putChunk(uploadId, resumeToken, 0, chunk0),
86+
putChunk(uploadId, resumeToken, 1, chunk1),
87+
]);
88+
89+
expect(await progress(uploadId)).toMatchObject({ uploadedChunks: 2, uploadedSize: MIN_CHUNK + 10 });
90+
const res = await stack.api(`/storage/upload/chunked/${uploadId}/complete`, {
91+
method: 'POST',
92+
headers: { ...json, ...auth() },
93+
body: JSON.stringify({
94+
uploadId,
95+
parts: [
96+
{ chunkIndex: 0, eTag: eTag0 },
97+
{ chunkIndex: 1, eTag: eTag1 },
98+
],
99+
}),
100+
});
101+
expect(res.status, await res.clone().text()).toBe(200);
102+
expect(((await res.json()) as any).data).toMatchObject({ fileId, size: MIN_CHUNK + 10 });
103+
});
104+
105+
it('eight chunks sent at once are all recorded', async () => {
106+
const { uploadId, resumeToken, totalChunks } = await init(8 * MIN_CHUNK);
107+
expect(totalChunks).toBe(8);
108+
// Small, distinct sizes: the chunk door records what it is sent, and a lost
109+
// part shows in both counts.
110+
const chunks = Array.from({ length: 8 }, (_, i) => fill(0x61 + i, 10 + i));
111+
112+
await Promise.all(chunks.map((bytes, i) => putChunk(uploadId, resumeToken, i, bytes)));
113+
114+
const total = chunks.reduce((sum, c) => sum + c.length, 0);
115+
expect(await progress(uploadId)).toMatchObject({ uploadedChunks: 8, uploadedSize: total });
116+
});
117+
});

0 commit comments

Comments
 (0)