Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 26 additions & 2 deletions packages/storage/src/__tests__/context-offload-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,8 @@ import {

after(removeTrackedControlDirectories);

const SINGLE_FLIGHT_RACE_ROUNDS = 500;

test('requires authentic Storage Root leases and writer facades', async () => {
await assert.rejects(
() =>
Expand Down Expand Up @@ -81,9 +83,8 @@ test('single-flights one limit-bound writer and snapshots admitted inputs', asyn
);
(mutableLimits.ownerMaxBytes as { tool_result_archive: number }).tool_result_archive = 0;
(mutableLimits as { sessionLogicalBytes: number }).sessionLogicalBytes = 0;
const [first, second] = await Promise.all([opening, concurrentOpening]);
const [first, second] = await Promise.all([opening, concurrentOpening, conflictingOpening]);
try {
await conflictingOpening;
assert.strictEqual(second, first);
assert.strictEqual(authenticateInteractiveContextOffloadWriter(first), first);
const reader = createInteractiveContextOffloadReader(first);
Expand Down Expand Up @@ -150,6 +151,29 @@ test('single-flights one limit-bound writer and snapshots admitted inputs', asyn
});
});

test('binds the first caller limits whatever order concurrent lease checks finish in', async () => {
// Each open checks the root identity on disk before joining the single flight, and
// those checks finish in no fixed order, so replay the race enough times to observe it.
await withInteractiveOwner(async (owner) => {
for (let round = 0; round < SINGLE_FLIGHT_RACE_ROUNDS; round += 1) {
const [first, conflicting] = await Promise.allSettled([
openInteractiveContextOffloadStoreForWrite(owner.lease, { limits: testLimits() }),
openInteractiveContextOffloadStoreForWrite(owner.lease, {
limits: { ...testLimits(), workspacePhysicalBytes: 63 },
}),
]);
await Promise.all(
[first, conflicting].map((outcome) =>
outcome.status === 'fulfilled' ? outcome.value.close() : undefined,
),
);
assert.equal(first.status, 'fulfilled', `round ${round}: the first caller was rejected`);
assert.equal(conflicting.status, 'rejected', `round ${round}: the conflict was admitted`);
assert.match(String(conflicting.reason), /different limits/u);
}
});
});

test('close drains admitted work, revokes the facade, and permits a clean reopen', async () => {
await withInteractiveOwner(async (owner) => {
const writer = await openInteractiveContextOffloadStoreForWrite(owner.lease, {
Expand Down
9 changes: 7 additions & 2 deletions packages/storage/src/context-offload-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import type {
} from '@maka/core/context-offload';
import {
assertStorageRootLease,
assertStorageRootLeaseActive,
runWithStorageRootLease,
StorageRootAuthorityError,
type StorageRootLease,
Expand Down Expand Up @@ -123,11 +124,15 @@ export async function openInteractiveContextOffloadStoreForWrite(
): Promise<InteractiveContextOffloadWriter> {
const limits = snapshotLimits(options.limits);
const limitsKey = serializeLimits(limits);
await assertStorageRootLease(lease, 'interactive', 'write');
// Claim or join the per-lease slot before the first await so the earliest caller binds
// the limits; the on-disk root identity checks of concurrent callers finish in no fixed order.
assertStorageRootLeaseActive(lease, 'interactive', 'write');
const existing = writerByLease.get(lease);
if (existing) {
assertSameLimits(existing.limitsKey, limitsKey);
return existing.writer;
await assertStorageRootLease(lease, 'interactive', 'write');
if (writerByLease.get(lease) === existing) return existing.writer;
return openInteractiveContextOffloadStoreForWrite(lease, { limits });
}
const opening = writerOpeningByLease.get(lease);
if (opening) {
Expand Down