Skip to content
Merged
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
205 changes: 205 additions & 0 deletions src/wasix/sdks/ts/src/__tests__/opfs-pool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import { restoreGlobals, stubGlobal } from './test-globals.js';
const OP = {
metadata: 1,
createDirectory: 3,
removeDirectory: 4,
rename: 5,
open: 7,
close: 8,
Expand Down Expand Up @@ -197,6 +198,210 @@ describe('WASIX pooled OPFS storage', () => {
await pool.close(false);
});

for (const boundary of ['operation', 'full', 'close'] as const) {
for (const explicitFlush of [false, true]) {
it(`skips net-zero namespace publication at ${boundary} with guest file sync ${explicitFlush}`, async () => {
const io: FakeIo = { writes: [], flushes: [] };
const root = installOpfs(io);
const name = `net-zero-${boundary}-${explicitFlush}`;
const pool = await DirectOpfsPool.open(name, clusterSeed(), compatible());
const relation = open(pool, 'base/value', FLAG_WRITE | FLAG_CREATE);
const wal = open(pool, 'pg_wal/0001', FLAG_WRITE | FLAG_CREATE);
const control = open(pool, 'global/pg_control', FLAG_WRITE);
write(pool, relation, Uint8Array.of(1));
write(pool, wal, Uint8Array.of(1));
await pool.sync('full');
const database = await databaseDirectory(root, name);
const published = await database.file('state.json').text();
const state = JSON.parse(published) as {
entries: Array<{ path: string; backing?: string }>;
};
const backing = (path: string): string => {
const value = state.entries.find((entry) => entry.path === path)?.backing;
if (value === undefined) throw new Error(`missing backing for ${path}`);
return value;
};
io.writes!.length = 0;
io.flushes!.length = 0;
write(pool, relation, Uint8Array.of(2));
write(pool, wal, Uint8Array.of(2));
write(pool, control, Uint8Array.of(2));
if (explicitFlush) requestOk(pool, OP.flush, '', new Uint8Array(), relation);
requestOk(pool, OP.createDirectory, 'base/temporary');
// Exhaust the 32 spares too: unlinked staged overflow needs no backing.
for (let index = 0; index < 40; index += 1) {
const temporary = open(pool, `base/temporary/${index}`, FLAG_WRITE | FLAG_CREATE);
write(pool, temporary, new Uint8Array(8192));
requestOk(pool, OP.unlink, '', new Uint8Array(), temporary);
requestOk(pool, OP.close, '', new Uint8Array(), temporary);
}
requestOk(pool, OP.removeDirectory, 'base/temporary');
await pool.sync(boundary);
expect(io.writes).toEqual([]);
expect(await database.file('state.json').text()).toBe(published);
const flushed = io.flushes!.map(lastPathSegment);
expect(flushed).toContain(backing('pg_wal/0001'));
expect(flushed.includes(backing('base/value'))).toBe(
explicitFlush || boundary !== 'operation',
);
expect(flushed.includes(backing('global/pg_control'))).toBe(boundary !== 'operation');
if (boundary === 'operation') {
expect(flushed).toEqual([
...(explicitFlush ? [backing('base/value')] : []),
backing('pg_wal/0001'),
]);
}
for (const descriptor of [relation, wal, control]) {
requestOk(pool, OP.close, '', new Uint8Array(), descriptor);
}
await pool.close(false);
const reopened = await DirectOpfsPool.open(name, clusterSeed(), compatible());
expect(reopened.state).toBe('existing');
expect(reopened.request(OP.metadata, 'base/temporary', new Uint8Array(), 0, 0, 0)[0]).toBe(
1,
);
await reopened.close(false);
});
}
}

it.each([
'wal',
'relation',
'state',
] as const)('preserves the published namespace and retries after a %s publication failure', async (failure) => {
const io: FakeIo = { writes: [], flushes: [] };
const root = installOpfs(io);
const pool = await DirectOpfsPool.open(`retry-${failure}`, clusterSeed(), compatible());
const wal = open(pool, 'pg_wal/0001', FLAG_WRITE | FLAG_CREATE);
write(pool, wal, Uint8Array.of(1));
await pool.sync('full');
const database = await databaseDirectory(root, `retry-${failure}`);
const published = await database.file('state.json').text();
const relation = open(pool, 'base/persistent', FLAG_WRITE | FLAG_CREATE);
write(pool, relation, Uint8Array.of(3));
requestOk(pool, OP.close, '', new Uint8Array(), relation);
if (failure === 'wal') write(pool, wal, Uint8Array.of(2));
if (failure === 'state') io.failNextStateCommit = true;
else io.failNextFlush = true;
await expect(pool.sync('operation')).rejects.toThrow(
failure === 'state' ? 'injected state commit failure' : 'injected OPFS flush failure',
);
expect(await database.file('state.json').text()).toBe(published);
io.writes!.length = 0;
await pool.sync('operation');
expect(io.writes).toEqual([`${OPFS_POOL_ROOT}/retry-${failure}/state.json`]);
expect(await database.file('state.json').text()).not.toBe(published);
io.writes!.length = 0;
const temporary = open(pool, 'base/temporary', FLAG_WRITE | FLAG_CREATE);
requestOk(pool, OP.unlink, '', new Uint8Array(), temporary);
requestOk(pool, OP.close, '', new Uint8Array(), temporary);
await pool.sync('operation');
expect(io.writes).toEqual([]);
requestOk(pool, OP.close, '', new Uint8Array(), wal);
await pool.close(false);
const reopened = await DirectOpfsPool.open(`retry-${failure}`, clusterSeed(), compatible());
const recovered = open(reopened, 'base/persistent', FLAG_READ);
expect(read(reopened, recovered, 1)).toEqual(Uint8Array.of(3));
await reopened.close(false);
});

it('keeps open-unlinked bytes isolated when transient slots are reused and handed to the portable path', async () => {
const io: FakeIo = { writes: [] };
const root = installOpfs(io);
const pool = await DirectOpfsPool.open('reuse', clusterSeed(), compatible());
await pool.sync('full');
const database = await databaseDirectory(root, 'reuse');
const published = await database.file('state.json').text();
const unlinked = open(pool, 'base/transient', FLAG_READ | FLAG_WRITE | FLAG_CREATE);
write(pool, unlinked, Uint8Array.of(7, 8));
requestOk(pool, OP.unlink, '', new Uint8Array(), unlinked);
await pool.sync('operation');
expect(await database.file('state.json').text()).toBe(published);
for (let index = 0; index < 40; index += 1) {
const temporary = open(pool, `base/temporary-${index}`, FLAG_WRITE | FLAG_CREATE);
write(pool, temporary, Uint8Array.of(index));
requestOk(pool, OP.unlink, '', new Uint8Array(), temporary);
requestOk(pool, OP.close, '', new Uint8Array(), temporary);
}
await pool.sync('operation');
expect(read(pool, unlinked, 8)).toEqual(Uint8Array.of(7, 8));
write(pool, unlinked, Uint8Array.of(9), 1);
requestOk(pool, OP.flush, '', new Uint8Array(), unlinked);
expect(read(pool, unlinked, 8)).toEqual(Uint8Array.of(7, 9));
requestOk(pool, OP.close, '', new Uint8Array(), unlinked);
await pool.sync('operation');
const persistent = open(pool, 'base/transient', FLAG_WRITE | FLAG_CREATE);
write(pool, persistent, Uint8Array.of(4));
await pool.sync('operation');
await pool.close(false);
const snapshot = await readSnapshot(database, 'reuse', compatible());
expect(snapshot?.files.find(({ path }) => path === 'base/transient')?.bytes).toEqual(
Uint8Array.of(4),
);
expect(snapshot?.files.some(({ path }) => path.startsWith('base/temporary-'))).toBe(false);
await applyPooledOpfsDelta(database.asHandle(), 'reuse', compatible(), {
directories: [],
files: [{ path: 'base/transient', bytes: Uint8Array.of(5) }],
deleted: [],
});
const reopened = await DirectOpfsPool.open('reuse', clusterSeed(), compatible());
expect(read(reopened, open(reopened, 'base/transient', FLAG_READ), 8)).toEqual(
Uint8Array.of(5),
);
await reopened.close(false);
});

it('distinguishes unchanged rename mappings from deletion, replacement and stale process-file cleanup', async () => {
const io: FakeIo = { writes: [] };
const root = installOpfs(io);
const pool = await DirectOpfsPool.open('mapping', clusterSeed(), compatible());
const descriptor = open(pool, 'base/value', FLAG_WRITE | FLAG_CREATE);
write(pool, descriptor, Uint8Array.of(1));
requestOk(pool, OP.close, '', new Uint8Array(), descriptor);
await pool.sync('full');
const database = await databaseDirectory(root, 'mapping');
const published = await database.file('state.json').text();
io.writes!.length = 0;
requestOk(pool, OP.rename, 'base/value', encoder.encode('base/renamed'));
requestOk(pool, OP.rename, 'base/renamed', encoder.encode('base/value'));
await pool.sync('operation');
expect(io.writes).toEqual([]);
expect(await database.file('state.json').text()).toBe(published);
const removed = open(pool, 'base/value', FLAG_WRITE);
requestOk(pool, OP.unlink, '', new Uint8Array(), removed);
requestOk(pool, OP.close, '', new Uint8Array(), removed);
const replacement = open(pool, 'base/value', FLAG_WRITE | FLAG_CREATE);
write(pool, replacement, Uint8Array.of(2));
requestOk(pool, OP.close, '', new Uint8Array(), replacement);
io.failNextStateCommit = true;
await expect(pool.sync('operation')).rejects.toThrow('injected state commit failure');
expect(await database.file('state.json').text()).toBe(published);
await pool.close(false);
const recovered = await DirectOpfsPool.open('mapping', clusterSeed(), compatible());
expect(read(recovered, open(recovered, 'base/value', FLAG_READ), 1)).toEqual(Uint8Array.of(1));
const stale = open(recovered, 'postmaster.pid', FLAG_WRITE | FLAG_CREATE);
write(recovered, stale, encoder.encode('stale'));
requestOk(recovered, OP.close, '', new Uint8Array(), stale);
await recovered.sync('operation');
await recovered.close(false);
io.writes!.length = 0;
const scrubbed = await DirectOpfsPool.open('mapping', clusterSeed(), compatible());
expect(io.writes).toEqual([`${OPFS_POOL_ROOT}/mapping/state.json`]);
expect(scrubbed.request(OP.metadata, 'postmaster.pid', new Uint8Array(), 0, 0, 0)[0]).toBe(1);
const changed = open(scrubbed, 'base/value', FLAG_WRITE);
requestOk(scrubbed, OP.unlink, '', new Uint8Array(), changed);
requestOk(scrubbed, OP.close, '', new Uint8Array(), changed);
const newFile = open(scrubbed, 'base/value', FLAG_WRITE | FLAG_CREATE);
write(scrubbed, newFile, Uint8Array.of(3));
await scrubbed.sync('operation');
expect(await database.file('state.json').text()).not.toBe(published);
await scrubbed.close(false);
const final = await DirectOpfsPool.open('mapping', clusterSeed(), compatible());
expect(read(final, open(final, 'base/value', FLAG_READ), 1)).toEqual(Uint8Array.of(3));
await final.close(false);
});

it('stages a creation burst beyond the preopened fast path and publishes it durably', async () => {
installOpfs();
const pool = await DirectOpfsPool.open('burst', clusterSeed(), compatible());
Expand Down
28 changes: 17 additions & 11 deletions src/wasix/sdks/ts/src/storage/opfs-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,7 @@ export class DirectOpfsPool {
#initializationComplete: boolean;
#nextDescriptor = 1;
#namespaceDirty = false;
#publishedNamespace = '';
#closed = false;

private constructor(
Expand Down Expand Up @@ -200,6 +201,9 @@ export class DirectOpfsPool {
const pool = new DirectOpfsPool(name, database, data, selectedIdentity, state);
try {
await pool.#hydrate(stored);
pool.#publishedNamespace = JSON.stringify(
serializePoolState(name, selectedIdentity, pool.#entries, stored.phase),
);
await pool.#restoreSpares();
await pool.#removeVolatileEntries();
await pool.#maintainSpares();
Expand All @@ -223,18 +227,20 @@ export class DirectOpfsPool {
await this.#materializeStagedFiles();
const completesInitialization = !this.#initializationComplete && boundary === 'full';
if (this.#namespaceDirty || completesInitialization) {
// A namespace publication makes newly mapped bytes reachable. Reuse the
// complete PostgreSQL order instead of depending on map insertion order.
this.#flushAll();
await writePoolState(
this.#database,
serializePoolState(
this.#name,
this.#physicalIdentity,
this.#entries,
this.#initializationComplete || completesInitialization ? 'ready' : 'initializing',
),
const state = serializePoolState(
this.#name,
this.#physicalIdentity,
this.#entries,
this.#initializationComplete || completesInitialization ? 'ready' : 'initializing',
);
const namespace = JSON.stringify(state);
// Creation and deletion within one operation can leave the published
// mapping unchanged. Only a real publication needs the complete barrier.
if (namespace !== this.#publishedNamespace) {
this.#flushAll();
await writePoolState(this.#database, state);
this.#publishedNamespace = namespace;
}
this.#namespaceDirty = false;
if (completesInitialization) this.#initializationComplete = true;
}
Expand Down
Loading