From df99a84838fe1368f199125d234d50d641ed0768 Mon Sep 17 00:00:00 2001 From: Andrew Levine Date: Thu, 3 Sep 2026 23:21:42 -0400 Subject: [PATCH 1/2] feat(storage): add zero-drawdown reservation release --- modules/jarvos-storage-janitor/README.md | 11 ++ modules/jarvos-storage-janitor/src/ports.js | 4 +- .../src/reservation-store.js | 57 ++++++- .../jarvos-storage-janitor/test/ports.test.js | 7 +- .../test/reservation-store.test.js | 158 +++++++++++++++++- 5 files changed, 230 insertions(+), 7 deletions(-) diff --git a/modules/jarvos-storage-janitor/README.md b/modules/jarvos-storage-janitor/README.md index 23b734b7..f2a67350 100644 --- a/modules/jarvos-storage-janitor/README.md +++ b/modules/jarvos-storage-janitor/README.md @@ -56,6 +56,17 @@ contract. In short: inputs; reserving against a pool subtracts every already-active reservation held for that `poolId` across all fence generations, so two distinct reservations cannot double-commit the same headroom. + `release({ reservationId, idempotencyKey, now })` ends an active, + unexpired hold without any drawdown: `consumedBytes` stays zero and the + reservation immediately stops counting toward its pool's active headroom. + A same-key replay is a stable success; a different-key replay against an + already-released reservation, or a release attempt against a consumed, + expired, or missing reservation, is a typed blocked result rather than a + mutation. `released` is terminal -- a `reserve` idempotency replay against + a released reservation is rejected, never reopened, and a `consume` + attempt against a released reservation is rejected as `already_released` + with no drawdown, so no transition can ever reopen or rewrite a released + reservation. - **Ports** (`ports.js`) define the three explicit boundaries this package depends on -- capacity observation, external reclaim provider, and reservation persistence -- as typed method-shape contracts. None receives a diff --git a/modules/jarvos-storage-janitor/src/ports.js b/modules/jarvos-storage-janitor/src/ports.js index b7f7d435..b473d810 100644 --- a/modules/jarvos-storage-janitor/src/ports.js +++ b/modules/jarvos-storage-janitor/src/ports.js @@ -31,11 +31,11 @@ function assertExternalReclaimPort(port) { assertPortShape(port, ['proposeDryRun', 'execute'], 'external reclaim port'); } -// ReservationPersistencePort: reserve/consume/reap/get, matching the +// ReservationPersistencePort: reserve/consume/release/reap/get, matching the // reservation-store.js contract. A conforming implementation must satisfy // checkReservationStoreConformance from reservation-store.js. function assertReservationPort(port) { - assertPortShape(port, ['reserve', 'consume', 'reap', 'get'], 'reservation-persistence port'); + assertPortShape(port, ['reserve', 'consume', 'release', 'reap', 'get'], 'reservation-persistence port'); } module.exports = { diff --git a/modules/jarvos-storage-janitor/src/reservation-store.js b/modules/jarvos-storage-janitor/src/reservation-store.js index 081fe86d..6195b14c 100644 --- a/modules/jarvos-storage-janitor/src/reservation-store.js +++ b/modules/jarvos-storage-janitor/src/reservation-store.js @@ -3,7 +3,7 @@ const { isObject, clone, isOpaqueId, isSafeNonNegativeInt, isSafePositiveInt, isValidClockValue, normalizeTime, digestOf } = require('./primitives'); const RESERVATION_STORE_SCHEMA_VERSION = 'jarvos-storage-janitor.reservation-store.v1'; -const RESERVATION_STATES = Object.freeze(['active', 'consumed', 'expired']); +const RESERVATION_STATES = Object.freeze(['active', 'consumed', 'expired', 'released']); const MAX_MUTATE_ATTEMPTS = 8; function emptyState() { @@ -56,6 +56,7 @@ function publicReservation(record) { createdAt: record.createdAt, expiresAt: record.expiresAt, consumedAt: record.consumedAt, + releasedAt: record.releasedAt, }; } @@ -160,6 +161,9 @@ function createReservationStore(options = {}) { } if (existing.status === 'expired') return noCommit({ ok: false, reason: 'expired' }); if (existing.status === 'consumed') return noCommit({ ok: false, reason: 'already_consumed' }); + // A released reservation is terminal: a reserve replay must never + // reopen it, regardless of matching parameters or fence. + if (existing.status === 'released') return noCommit({ ok: false, reason: 'already_released' }); const matches = existing.amountBytes === amountBytes && existing.fenceGeneration === fenceGeneration @@ -222,9 +226,11 @@ function createReservationStore(options = {}) { consumedBytes: 0, status: 'active', consumeIdempotencyKey: null, + releaseIdempotencyKey: null, createdAt: effectiveNow, expiresAt: expiresIso, consumedAt: null, + releasedAt: null, }; state.reservations[id] = record; state.idempotencyIndex[idempotencyKey] = id; @@ -258,6 +264,10 @@ function createReservationStore(options = {}) { return noCommit({ ok: true, replayed: true, reservation: publicReservation(record) }); } if (record.status === 'consumed') return noCommit({ ok: false, reason: 'already_consumed' }); + // A released reservation is terminal: a consume attempt must never + // reopen it and rewrite it to `consumed`, regardless of idempotency + // key or requested amount. + if (record.status === 'released') return noCommit({ ok: false, reason: 'already_released' }); if (record.status === 'expired') return noCommit({ ok: false, reason: 'expired' }); if (record.status === 'active' && isExpired(record, effectiveNow)) { // Commit the expiry transition rather than reporting `expired` @@ -276,6 +286,49 @@ function createReservationStore(options = {}) { }); } + function release({ reservationId: id, idempotencyKey, now } = {}) { + return guarded(async () => { + if (!isOpaqueId(id) || !isOpaqueId(idempotencyKey)) { + return { ok: false, reason: 'invalid_request', errors: ['reservationId and idempotencyKey are required and must be well-formed'] }; + } + const resolvedNow = resolveNow(now, clock); + if (!resolvedNow.ok) return { ok: false, reason: 'invalid_request', errors: ['now must be a valid UTC ISO-8601 timestamp'] }; + const effectiveNow = resolvedNow.value; + + return mutate((state) => { + const record = state.reservations[id]; + if (!record) return noCommit({ ok: false, reason: 'not_found' }); + + // A repeat release with the same idempotency key replays the + // original result; a released reservation is terminal, so a + // same-key replay is the only way a second release call can + // succeed. A different key against an already-released + // reservation is a typed conflict, never a silent reuse. + if (record.status === 'released' && record.releaseIdempotencyKey === idempotencyKey) { + return noCommit({ ok: true, replayed: true, reservation: publicReservation(record) }); + } + if (record.status === 'released') return noCommit({ ok: false, reason: 'already_released' }); + if (record.status === 'consumed') return noCommit({ ok: false, reason: 'already_consumed' }); + if (record.status === 'expired') return noCommit({ ok: false, reason: 'expired' }); + if (record.status === 'active' && isExpired(record, effectiveNow)) { + // Commit the expiry transition rather than reporting `expired` + // for a mutation that was never actually persisted. + record.status = 'expired'; + return { ok: false, reason: 'expired' }; + } + + // Freeing an active, unexpired reservation drops its status out of + // `active`, so it stops counting toward reserve()'s aggregate + // active-headroom sum for its pool immediately, with no + // consumedBytes drawdown. + record.status = 'released'; + record.releaseIdempotencyKey = idempotencyKey; + record.releasedAt = effectiveNow; + return { ok: true, replayed: false, reservation: publicReservation(record) }; + }); + }); + } + function reap({ now } = {}) { return guarded(async () => { const resolvedNow = resolveNow(now, clock); @@ -308,7 +361,7 @@ function createReservationStore(options = {}) { }); } - return { reserve, consume, reap, get }; + return { reserve, consume, release, reap, get }; } function createMemoryReservationBackend() { diff --git a/modules/jarvos-storage-janitor/test/ports.test.js b/modules/jarvos-storage-janitor/test/ports.test.js index 8cd699b2..226daec5 100644 --- a/modules/jarvos-storage-janitor/test/ports.test.js +++ b/modules/jarvos-storage-janitor/test/ports.test.js @@ -20,9 +20,12 @@ test('an external reclaim port must expose proposeDryRun() and execute()', () => assert.doesNotThrow(() => assertExternalReclaimPort({ proposeDryRun: () => {}, execute: () => {} })); }); -test('a reservation port must expose reserve(), consume(), reap(), and get()', () => { +test('a reservation port must expose reserve(), consume(), release(), reap(), and get()', () => { assert.throws(() => assertReservationPort({ reserve: () => {} })); - assert.doesNotThrow(() => assertReservationPort({ + assert.throws(() => assertReservationPort({ reserve: () => {}, consume: () => {}, reap: () => {}, get: () => {}, })); + assert.doesNotThrow(() => assertReservationPort({ + reserve: () => {}, consume: () => {}, release: () => {}, reap: () => {}, get: () => {}, + })); }); diff --git a/modules/jarvos-storage-janitor/test/reservation-store.test.js b/modules/jarvos-storage-janitor/test/reservation-store.test.js index 55661177..91275d07 100644 --- a/modules/jarvos-storage-janitor/test/reservation-store.test.js +++ b/modules/jarvos-storage-janitor/test/reservation-store.test.js @@ -292,6 +292,159 @@ test('idempotent reserve replay against an expired reservation is rejected', asy assert.equal(replay.reason, 'expired'); }); +test('release frees an active reservation with zero consumed bytes', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + const result = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.replayed, false); + assert.equal(result.reservation.status, 'released'); + assert.equal(result.reservation.consumedBytes, 0); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.ok, true); + assert.equal(fetched.reservation.status, 'released'); + assert.equal(fetched.reservation.consumedBytes, 0); +}); + +test('release is idempotent for a repeated idempotencyKey', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + const first = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + const second = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + assert.equal(first.ok, true); + assert.equal(second.ok, true); + assert.equal(second.replayed, true); + assert.equal(second.reservation.status, 'released'); +}); + +test('a release replay with a different idempotencyKey against an already-released reservation is rejected without mutation', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + const different = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:different', now: req().now }); + assert.equal(different.ok, false); + assert.equal(different.reason, 'already_released'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'released'); +}); + +test('idempotent reserve replay against a released reservation is rejected as terminal, never reopened', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + const replay = await store.reserve(req()); + assert.equal(replay.ok, false); + assert.equal(replay.reason, 'already_released'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'released'); +}); + +test('release rejects a consumed reservation without mutation', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + await store.consume({ reservationId: reservation.reservationId, idempotencyKey: 'consume:once', amountBytes: reservation.amountBytes, now: req().now }); + const result = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + assert.equal(result.ok, false); + assert.equal(result.reason, 'already_consumed'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'consumed'); +}); + +test('release rejects an expired reservation and persists the expiry transition', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req({ expiresAt: '2026-09-03T12:04:00.000Z' })); + const result = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: '2026-09-03T12:05:00.000Z' }); + assert.equal(result.ok, false); + assert.equal(result.reason, 'expired'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'expired'); +}); + +test('release rejects a missing, invalid, or malformed request without mutation', async () => { + const store = createMemoryReservationStore(); + const missing = await store.release({ reservationId: 'reservation_does_not_exist', idempotencyKey: 'release:once', now: req().now }); + assert.equal(missing.ok, false); + assert.equal(missing.reason, 'not_found'); + + const invalid = await store.release({ reservationId: 'not an opaque id!', idempotencyKey: 'release:once', now: req().now }); + assert.equal(invalid.ok, false); + assert.equal(invalid.reason, 'invalid_request'); + + const { reservation } = await store.reserve(req()); + const malformedClock = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: '2026-09-03T12:03:00' }); + assert.equal(malformedClock.ok, false); + assert.equal(malformedClock.reason, 'invalid_request'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'active'); +}); + +test('release frees a reservation from aggregate active headroom so a new reservation can reuse it', async () => { + const store = createMemoryReservationStore(); + const first = await store.reserve(req({ idempotencyKey: 'reserve:a', amountBytes: 1200000000 })); + assert.equal(first.ok, true, JSON.stringify(first)); + const blocked = await store.reserve(req({ idempotencyKey: 'reserve:b', amountBytes: 900000000 })); + assert.equal(blocked.ok, false); + assert.equal(blocked.reason, 'capacity_exceeded'); + + const released = await store.release({ reservationId: first.reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + assert.equal(released.ok, true, JSON.stringify(released)); + + const second = await store.reserve(req({ idempotencyKey: 'reserve:b', amountBytes: 900000000 })); + assert.equal(second.ok, true, JSON.stringify(second)); +}); + +test('concurrent release calls for the same reservation and idempotencyKey cannot double-mutate', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + const [a, b] = await Promise.all([ + store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }), + store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }), + ]); + assert.equal(a.ok, true); + assert.equal(b.ok, true); + assert.equal([a.replayed, b.replayed].filter(Boolean).length, 1); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.reservation.status, 'released'); +}); + +test('consume rejects a released reservation without mutation', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + const result = await store.consume({ reservationId: reservation.reservationId, idempotencyKey: 'consume:once', amountBytes: reservation.amountBytes, now: req().now }); + assert.equal(result.ok, false); + assert.equal(result.reason, 'already_released'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.ok, true); + assert.equal(fetched.reservation.status, 'released'); + assert.equal(fetched.reservation.consumedBytes, 0); +}); + +test('a release racing a consume for the same reservation lets only one transition win, and the loser is rejected as already_released', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req()); + const [released, consumed] = await Promise.all([ + store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }), + store.consume({ reservationId: reservation.reservationId, idempotencyKey: 'consume:once', amountBytes: reservation.amountBytes, now: req().now }), + ]); + assert.equal(released.ok, true, JSON.stringify(released)); + assert.equal(consumed.ok, false, JSON.stringify(consumed)); + assert.equal(consumed.reason, 'already_released'); + + const fetched = await store.get(reservation.reservationId); + assert.equal(fetched.ok, true); + assert.equal(fetched.reservation.status, 'released'); + assert.equal(fetched.reservation.consumedBytes, 0); +}); + test('get validates its identifier', async () => { const store = createMemoryReservationStore(); const result = await store.get('not an opaque id!'); @@ -317,7 +470,7 @@ test('fail-closed recovery: a backend that cannot load denies rather than procee assert.equal(result.reason, 'store_unavailable'); }); -test('fail-closed recovery covers get, consume, and reap as well as reserve', async () => { +test('fail-closed recovery covers get, consume, release, and reap as well as reserve', async () => { const backend = { load() { throw new Error('backend unavailable'); }, save() { throw new Error('backend unavailable'); }, @@ -325,11 +478,14 @@ test('fail-closed recovery covers get, consume, and reap as well as reserve', as const store = createReservationStore({ backend }); const getResult = await store.get('reservation_abc'); const consumeResult = await store.consume({ reservationId: 'reservation_abc', idempotencyKey: 'consume:x', amountBytes: 1 }); + const releaseResult = await store.release({ reservationId: 'reservation_abc', idempotencyKey: 'release:x' }); const reapResult = await store.reap({}); assert.equal(getResult.ok, false); assert.equal(getResult.reason, 'store_unavailable'); assert.equal(consumeResult.ok, false); assert.equal(consumeResult.reason, 'store_unavailable'); + assert.equal(releaseResult.ok, false); + assert.equal(releaseResult.reason, 'store_unavailable'); assert.equal(reapResult.ok, false); assert.equal(reapResult.reason, 'store_unavailable'); }); From f1ee2e181ea66d69e5d36610b566497f2fdd2344 Mon Sep 17 00:00:00 2001 From: Andrew Levine Date: Thu, 3 Sep 2026 23:35:19 -0400 Subject: [PATCH 2/2] fix(storage): version released reservation state --- modules/jarvos-storage-janitor/README.md | 27 +- .../src/reservation-store.js | 80 +++++- .../test/reservation-store.test.js | 236 +++++++++++++++++- 3 files changed, 326 insertions(+), 17 deletions(-) diff --git a/modules/jarvos-storage-janitor/README.md b/modules/jarvos-storage-janitor/README.md index f2a67350..4fce9ccc 100644 --- a/modules/jarvos-storage-janitor/README.md +++ b/modules/jarvos-storage-janitor/README.md @@ -66,7 +66,32 @@ contract. In short: a released reservation is rejected, never reopened, and a `consume` attempt against a released reservation is rejected as `already_released` with no drawdown, so no transition can ever reopen or rewrite a released - reservation. + reservation. `reap` never reverts a released reservation back to `expired` + either: terminal means terminal against every transition, not just + `reserve` and `consume`. + + `release()` and the `released` status are a **breaking addition** to the + reservation-persistence port and its persisted schema: the store's + persisted `schemaVersion` is `jarvos-storage-janitor.reservation-store.v2`. + Any host adapter implementing this port directly (not through + `createMemoryReservationStore`) must add `release()` -- see + `assertReservationPort` in `ports.js`. A **genuine** v1 store (one written + before `release` existed) contains only `active`, `consumed`, or `expired` + records; this package still reads it, normalizes its missing + `releasedAt`/`releaseIdempotencyKey` fields to `null` in memory, and + upgrades it to v2 in storage the next time it is mutated -- `get()` alone + never writes anything. A v1 store that impossibly contains a `released` + record (which genuine v1 code could never have written) is rejected as + invalid rather than silently accepted. This upgrade path is one-way: a + v1-only rollback boundary. Once any node has upgraded a store to v2 (by + mutating it, or simply by starting from an already-v2 empty state), old + v1-only code reading that store sees a `schemaVersion` it does not + recognize and fails closed with a schema-mismatch error -- it does not + crash, and it does not silently misinterpret a `released` record as some + other status. Rolling back to v1-only code against a store any v2 code has + touched is therefore unsafe and unsupported; a rollback requires either a + v1-only store that has never been touched by v2 code, or restoring a v1 + snapshot taken before the first v2 write. - **Ports** (`ports.js`) define the three explicit boundaries this package depends on -- capacity observation, external reclaim provider, and reservation persistence -- as typed method-shape contracts. None receives a diff --git a/modules/jarvos-storage-janitor/src/reservation-store.js b/modules/jarvos-storage-janitor/src/reservation-store.js index 6195b14c..739dd965 100644 --- a/modules/jarvos-storage-janitor/src/reservation-store.js +++ b/modules/jarvos-storage-janitor/src/reservation-store.js @@ -2,7 +2,15 @@ const { isObject, clone, isOpaqueId, isSafeNonNegativeInt, isSafePositiveInt, isValidClockValue, normalizeTime, digestOf } = require('./primitives'); -const RESERVATION_STORE_SCHEMA_VERSION = 'jarvos-storage-janitor.reservation-store.v1'; +// v1 shipped before `release`/`released` existed: it is the pre-release +// contract, and a genuine v1 store can therefore never contain a `released` +// record. v2 is a breaking addition (a new terminal status and its +// `releasedAt`/`releaseIdempotencyKey` fields); old v1-only code must fail +// closed on a v2 store via this schema-version mismatch rather than silently +// misreading a status it does not know about. See README.md for the +// documented rollback boundary this implies. +const RESERVATION_STORE_SCHEMA_VERSION_V1 = 'jarvos-storage-janitor.reservation-store.v1'; +const RESERVATION_STORE_SCHEMA_VERSION = 'jarvos-storage-janitor.reservation-store.v2'; const RESERVATION_STATES = Object.freeze(['active', 'consumed', 'expired', 'released']); const MAX_MUTATE_ATTEMPTS = 8; @@ -21,6 +29,55 @@ function validateState(state) { return { ok: errors.length === 0, errors }; } +// A legitimate v1 store only ever wrote `active`, `consumed`, or `expired`; +// a `released` record under a v1 schemaVersion is impossible for genuine v1 +// data, so it is rejected as invalid rather than silently normalized. +function validateV1State(state) { + const errors = []; + if (!isObject(state)) return { ok: false, errors: ['reservation-store state must be an object'] }; + if (state.schemaVersion !== RESERVATION_STORE_SCHEMA_VERSION_V1) errors.push(`state.schemaVersion must be ${RESERVATION_STORE_SCHEMA_VERSION_V1}`); + if (!Number.isInteger(state.revision) || state.revision < 0) errors.push('state.revision must be a non-negative integer'); + if (!Number.isInteger(state.currentFence) || state.currentFence < 0) errors.push('state.currentFence must be a non-negative integer'); + if (!isObject(state.reservations)) errors.push('state.reservations must be an object'); + if (!isObject(state.idempotencyIndex)) errors.push('state.idempotencyIndex must be an object'); + if (isObject(state.reservations)) { + for (const record of Object.values(state.reservations)) { + if (isObject(record) && record.status === 'released') { + errors.push(`reservation ${record.reservationId} has status "released" under schemaVersion ${RESERVATION_STORE_SCHEMA_VERSION_V1}, which is impossible for a genuine v1 store`); + } + } + } + return { ok: errors.length === 0, errors }; +} + +// Upgrades an already-validated v1 state to the v2 shape in memory, adding +// the `released`-status fields a v1 record never had. This does not persist +// anything by itself: a mutation persists the upgrade via its normal save, +// while a read-only `get` returns the upgraded shape without writing it back. +function upgradeV1State(state) { + const upgraded = clone(state); + upgraded.schemaVersion = RESERVATION_STORE_SCHEMA_VERSION; + for (const record of Object.values(upgraded.reservations)) { + if (record.releaseIdempotencyKey === undefined) record.releaseIdempotencyKey = null; + if (record.releasedAt === undefined) record.releasedAt = null; + } + return upgraded; +} + +// Every read path funnels through here so a legitimate v1 store loads and +// upgrades exactly once, a v2 store loads as-is, and anything else -- +// including a v1 store impossibly marked `released` -- fails closed. +function loadAndUpgradeState(loaded) { + if (isObject(loaded) && loaded.schemaVersion === RESERVATION_STORE_SCHEMA_VERSION_V1) { + const v1Validation = validateV1State(loaded); + if (!v1Validation.ok) return { ok: false, errors: v1Validation.errors }; + return { ok: true, state: upgradeV1State(loaded) }; + } + const validation = validateState(loaded); + if (!validation.ok) return { ok: false, errors: validation.errors }; + return { ok: true, state: loaded }; +} + // A conflict here is the store's declared atomic primitive speaking: a // conforming backend detects a stale compare-and-set precondition and raises // exactly this shape rather than silently overwriting the loser's read. @@ -84,9 +141,9 @@ function createReservationStore(options = {}) { async function mutate(mutator) { for (let attempt = 0; attempt < MAX_MUTATE_ATTEMPTS; attempt += 1) { const loaded = await backend.load(); - const validation = validateState(loaded); - if (!validation.ok) throw new Error(`invalid reservation-store state: ${validation.errors.join('; ')}`); - const state = clone(loaded); + const normalized = loadAndUpgradeState(loaded); + if (!normalized.ok) throw new Error(`invalid reservation-store state: ${normalized.errors.join('; ')}`); + const state = clone(normalized.state); const expectedRevision = state.revision; const outcome = mutator(state); if (outcome && outcome.__noCommit) return outcome.value; @@ -352,10 +409,10 @@ function createReservationStore(options = {}) { function get(id) { return guarded(async () => { if (!isOpaqueId(id)) return { ok: false, reason: 'invalid_request', errors: ['reservationId must be an opaque identifier'] }; - const state = await backend.load(); - const validation = validateState(state); - if (!validation.ok) throw new Error(`invalid reservation-store state: ${validation.errors.join('; ')}`); - const record = state.reservations[id]; + const loaded = await backend.load(); + const normalized = loadAndUpgradeState(loaded); + if (!normalized.ok) throw new Error(`invalid reservation-store state: ${normalized.errors.join('; ')}`); + const record = normalized.state.reservations[id]; if (!record) return { ok: false, reason: 'not_found' }; return { ok: true, reservation: publicReservation(record) }; }); @@ -390,7 +447,11 @@ function createMemoryReservationStore(options = {}) { // equivalent serialization) this store depends on for no-double-spend: a // second save() against a precondition its own first save() already // invalidated must be rejected with a reservationConflict error, not -// silently accepted as a last-writer-wins overwrite. +// silently accepted as a last-writer-wins overwrite. This is a mechanical +// property of a *fresh* backend, not a business-semantics check, so it +// validates against the current schema directly rather than through +// loadAndUpgradeState's legacy-store business rules (a fresh backend has no +// legacy store to upgrade from). async function checkReservationStoreConformance(createBackend) { if (typeof createBackend !== 'function') return { ok: false, errors: ['createBackend must be a factory function returning a fresh backend'] }; const backend = createBackend(); @@ -424,6 +485,7 @@ async function checkReservationStoreConformance(createBackend) { module.exports = { RESERVATION_STORE_SCHEMA_VERSION, + RESERVATION_STORE_SCHEMA_VERSION_V1, RESERVATION_STATES, emptyState, validateState, diff --git a/modules/jarvos-storage-janitor/test/reservation-store.test.js b/modules/jarvos-storage-janitor/test/reservation-store.test.js index 91275d07..99503acf 100644 --- a/modules/jarvos-storage-janitor/test/reservation-store.test.js +++ b/modules/jarvos-storage-janitor/test/reservation-store.test.js @@ -7,10 +7,14 @@ const { createMemoryReservationStore, createReservationStore, createMemoryReservationBackend, + createReservationConflictError, checkReservationStoreConformance, emptyState, RESERVATION_STATES, + RESERVATION_STORE_SCHEMA_VERSION, + RESERVATION_STORE_SCHEMA_VERSION_V1, } = require('../src/reservation-store'); +const { clone, isObject } = require('../src/primitives'); function countingBackend(backend) { const calls = { save: 0 }; @@ -34,6 +38,103 @@ function req(overrides = {}) { }; } +function createDeferred() { + let resolve; + const promise = new Promise((res) => { resolve = res; }); + return { promise, resolve }; +} + +// A single unsynchronized core shared by two independently-gated backend +// handles, so a test -- not incidental microtask/array ordering -- controls +// which of two racing operations' save() the CAS accepts first, while both +// are guaranteed to have loaded the same pre-mutation revision (neither's +// save can complete, and so neither can change the loadable state, until the +// test explicitly releases its gate). +function createSharedReservationCore() { + let state = emptyState(); + return { + async rawLoad() { return clone(state); }, + async rawSave(next, expectedRevision) { + if (state.revision !== expectedRevision) throw createReservationConflictError(); + state = clone(next); + }, + }; +} + +function createPlainReservationBackend(core) { + return { load: () => core.rawLoad(), save: (next, rev) => core.rawSave(next, rev) }; +} + +function createGatedReservationBackend(core) { + let saveGate = null; + let loadCalls = 0; + return { + get loadCalls() { return loadCalls; }, + holdSave() { saveGate = createDeferred(); }, + releaseSave() { + if (!saveGate) return; + const gate = saveGate; + saveGate = null; + gate.resolve(); + }, + async load() { + loadCalls += 1; + return core.rawLoad(); + }, + async save(next, expectedRevision) { + if (saveGate) await saveGate.promise; + return core.rawSave(next, expectedRevision); + }, + }; +} + +function legacyV1Record(overrides = {}) { + return { + reservationId: 'reservation_legacy001', + idempotencyKey: 'reserve:legacy-001', + poolId: 'pool:legacy-001', + capacityLimitBytes: 2000000000, + fenceGeneration: 1, + amountBytes: 900000000, + consumedBytes: 0, + status: 'active', + consumeIdempotencyKey: null, + createdAt: '2026-09-03T12:00:00.000Z', + expiresAt: '2026-09-03T13:00:00.000Z', + consumedAt: null, + ...overrides, + }; +} + +// Deliberately has no releaseIdempotencyKey/releasedAt fields at all: a +// genuine v1 record was written before those fields existed. +function legacyV1State(recordOverrides = {}) { + const record = legacyV1Record(recordOverrides); + return { + schemaVersion: RESERVATION_STORE_SCHEMA_VERSION_V1, + revision: 3, + currentFence: 1, + reservations: { [record.reservationId]: record }, + idempotencyIndex: { [record.idempotencyKey]: record.reservationId }, + }; +} + +// Replicates the pre-release validateState exactly as it shipped before +// `release`/`released` existed (schemaVersion pinned to v1, no released-field +// awareness at all), to prove old v1-only code fails closed -- via a schema +// mismatch, not a crash or a silent misread -- on a store a v2-aware node has +// since upgraded and released against. +function legacyV1OnlyValidateState(state) { + const errors = []; + if (!isObject(state)) return { ok: false, errors: ['reservation-store state must be an object'] }; + if (state.schemaVersion !== RESERVATION_STORE_SCHEMA_VERSION_V1) errors.push(`state.schemaVersion must be ${RESERVATION_STORE_SCHEMA_VERSION_V1}`); + if (!Number.isInteger(state.revision) || state.revision < 0) errors.push('state.revision must be a non-negative integer'); + if (!Number.isInteger(state.currentFence) || state.currentFence < 0) errors.push('state.currentFence must be a non-negative integer'); + if (!isObject(state.reservations)) errors.push('state.reservations must be an object'); + if (!isObject(state.idempotencyIndex)) errors.push('state.idempotencyIndex must be an object'); + return { ok: errors.length === 0, errors }; +} + test('reserve creates a reservation with a typed active state', async () => { const store = createMemoryReservationStore(); const result = await store.reserve(req()); @@ -428,16 +529,79 @@ test('consume rejects a released reservation without mutation', async () => { assert.equal(fetched.reservation.consumedBytes, 0); }); -test('a release racing a consume for the same reservation lets only one transition win, and the loser is rejected as already_released', async () => { - const store = createMemoryReservationStore(); - const { reservation } = await store.reserve(req()); - const [released, consumed] = await Promise.all([ - store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }), - store.consume({ reservationId: reservation.reservationId, idempotencyKey: 'consume:once', amountBytes: reservation.amountBytes, now: req().now }), - ]); +// A barrier/CAS backend forces release and consume to both load the exact +// same pre-mutation revision (neither's save can complete until the test +// explicitly releases its gate), then the test -- not incidental +// Promise.all/array ordering -- picks which save the CAS accepts first. The +// loser's save() is rejected by the real CAS check and mutate() retries: it +// re-loads the now-updated state and re-reads it as terminal, proving actual +// loser retry/re-read rather than a result baked in by call order. Either +// operation may win, so both orderings are exercised. +async function raceReleaseAndConsume({ releaseWinsFirst }) { + const core = createSharedReservationCore(); + const setupStore = createReservationStore({ backend: createPlainReservationBackend(core) }); + const { reservation } = await setupStore.reserve(req()); + + const releaseBackend = createGatedReservationBackend(core); + const consumeBackend = createGatedReservationBackend(core); + releaseBackend.holdSave(); + consumeBackend.holdSave(); + const releaseStore = createReservationStore({ backend: releaseBackend }); + const consumeStore = createReservationStore({ backend: consumeBackend }); + + const releasePromise = releaseStore.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + const consumePromise = consumeStore.consume({ reservationId: reservation.reservationId, idempotencyKey: 'consume:once', amountBytes: reservation.amountBytes, now: req().now }); + + if (releaseWinsFirst) { + releaseBackend.releaseSave(); + const released = await releasePromise; + consumeBackend.releaseSave(); + const consumed = await consumePromise; + return { released, consumed, releaseBackend, consumeBackend, reservationId: reservation.reservationId, setupStore }; + } + consumeBackend.releaseSave(); + const consumed = await consumePromise; + releaseBackend.releaseSave(); + const released = await releasePromise; + return { released, consumed, releaseBackend, consumeBackend, reservationId: reservation.reservationId, setupStore }; +} + +test('release winning a barrier-forced race against consume commits released, and the loser retries, re-reads, and reports already_released', async () => { + const { released, consumed, consumeBackend, reservationId, setupStore } = await raceReleaseAndConsume({ releaseWinsFirst: true }); assert.equal(released.ok, true, JSON.stringify(released)); + assert.equal(released.reservation.status, 'released'); assert.equal(consumed.ok, false, JSON.stringify(consumed)); assert.equal(consumed.reason, 'already_released'); + assert.ok(consumeBackend.loadCalls >= 2, 'the loser must retry: re-load after its first save is rejected by CAS'); + + const fetched = await setupStore.get(reservationId); + assert.equal(fetched.ok, true); + assert.equal(fetched.reservation.status, 'released'); + assert.equal(fetched.reservation.consumedBytes, 0); +}); + +test('consume winning a barrier-forced race against release commits consumed, and the loser retries, re-reads, and reports already_consumed', async () => { + const { released, consumed, releaseBackend, reservationId, setupStore } = await raceReleaseAndConsume({ releaseWinsFirst: false }); + assert.equal(consumed.ok, true, JSON.stringify(consumed)); + assert.equal(consumed.reservation.status, 'consumed'); + assert.equal(released.ok, false, JSON.stringify(released)); + assert.equal(released.reason, 'already_consumed'); + assert.ok(releaseBackend.loadCalls >= 2, 'the loser must retry: re-load after its first save is rejected by CAS'); + + const fetched = await setupStore.get(reservationId); + assert.equal(fetched.ok, true); + assert.equal(fetched.reservation.status, 'consumed'); +}); + +test('reap never reverts an already-released reservation back to expired: released remains terminal', async () => { + const store = createMemoryReservationStore(); + const { reservation } = await store.reserve(req({ expiresAt: '2026-09-03T12:04:00.000Z' })); + const released = await store.release({ reservationId: reservation.reservationId, idempotencyKey: 'release:once', now: req().now }); + assert.equal(released.ok, true, JSON.stringify(released)); + + const reap = await store.reap({ now: '2026-09-03T13:00:00.000Z' }); + assert.equal(reap.ok, true); + assert.deepEqual(reap.expired, []); const fetched = await store.get(reservation.reservationId); assert.equal(fetched.ok, true); @@ -535,3 +699,61 @@ test('conformance check fails a backend lacking atomic compare-and-set, not a ma assert.ok(result.errors.length > 0); assert.ok(result.errors.some((e) => /compare-and-set/i.test(e))); }); + +test('a legitimate v1 store normalizes a missing releasedAt to null and upgrades to v2 on the next mutation', async () => { + let persisted = null; + const v1State = legacyV1State(); + const backend = { + async load() { return clone(persisted || v1State); }, + async save(next, expectedRevision) { + const current = persisted || v1State; + if (current.revision !== expectedRevision) throw createReservationConflictError(); + persisted = clone(next); + }, + }; + const store = createReservationStore({ backend }); + + const result = await store.release({ reservationId: 'reservation_legacy001', idempotencyKey: 'release:once', now: '2026-09-03T12:05:00.000Z' }); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.reservation.status, 'released'); + assert.equal(result.reservation.releasedAt, '2026-09-03T12:05:00.000Z'); + + assert.equal(persisted.schemaVersion, RESERVATION_STORE_SCHEMA_VERSION); + assert.equal(persisted.reservations.reservation_legacy001.releasedAt, '2026-09-03T12:05:00.000Z'); + assert.equal(persisted.reservations.reservation_legacy001.releaseIdempotencyKey, 'release:once'); +}); + +test('get() of a legacy v1 record returns a stable public shape with releasedAt: null, without persisting the upgrade', async () => { + const v1State = legacyV1State(); + const backend = { + async load() { return clone(v1State); }, + async save() { throw new Error('a read-only get() must never save'); }, + }; + const store = createReservationStore({ backend }); + + const result = await store.get('reservation_legacy001'); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.reservation.status, 'active'); + assert.equal(result.reservation.releasedAt, null); + assert.equal(result.reservation.version, RESERVATION_STORE_SCHEMA_VERSION); +}); + +test('a v1 store impossibly containing status: released is rejected as invalid rather than silently normalized', async () => { + const v1State = legacyV1State({ status: 'released' }); + const backend = { + async load() { return clone(v1State); }, + async save() { throw new Error('an invalid state must never be saved'); }, + }; + const store = createReservationStore({ backend }); + + const result = await store.get('reservation_legacy001'); + assert.equal(result.ok, false); + assert.equal(result.reason, 'store_unavailable'); +}); + +test('old v1-only code fails closed on a v2 store via schema-version mismatch, the documented rollback boundary', () => { + const v2State = emptyState(); + const legacyValidation = legacyV1OnlyValidateState(v2State); + assert.equal(legacyValidation.ok, false); + assert.ok(legacyValidation.errors.some((e) => e.includes(RESERVATION_STORE_SCHEMA_VERSION_V1))); +});