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
204 changes: 144 additions & 60 deletions src/examples/IDBMirrorVFS.js
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ class File {
/** @type {{write?: function, reserved?: function, hint?: function}} */ locks;

/** @type {AbortController} */ abortController;
/** @type {number} */ commitsPending;

/** @type {Transaction?} */ txActive;
/** @type {boolean} */ txWriteHint;
Expand All @@ -56,6 +57,8 @@ class File {
this.broadcastReceived = [];
this.lockState = VFS.SQLITE_LOCK_NONE;
this.locks = {};
this.abortController = new AbortController();
this.commitsPending = 0;
Comment thread
rhashimoto marked this conversation as resolved.
this.txActive = null;
this.txWriteHint = false;
this.txOverwrite = false;
Expand Down Expand Up @@ -141,39 +144,7 @@ export class IDBMirrorVFS extends FacadeVFS {
}
}

// Load pages into memory from IndexedDB.
await new Promise((resolve, reject) => {
const range = IDBKeyRange.bound([path, 0], [path, Infinity]);
const request = blocks.openCursor(range);
request.onsuccess = () => {
const cursor = request.result;
if (cursor) {
const { offset, data } = cursor.value;
file.blocks.set(offset, data);
cursor.continue();
} else {
resolve();
}
};
request.onerror = () => reject(request.error);
});
file.blockSize = file.blocks.get(0)?.byteLength ?? 0;

// Get the last transaction id.
const transactions = idbTx.objectStore('tx');
file.viewTx = await new Promise((resolve, reject) => {
const range = IDBKeyRange.bound([path, 0], [path, Infinity]);
const request = transactions.openCursor(range, 'prev');
request.onsuccess = () => {
const cursor = request.result;
if (cursor) {
resolve(cursor.value);
} else {
resolve({ txId: 0 });
}
};
request.onerror = () => reject(request.error);
});
await this.#loadFile(file, idbTx);

// Publish our view of the database. This prevents other connections
// from overwriting file data we still need.
Expand Down Expand Up @@ -263,6 +234,10 @@ export class IDBMirrorVFS extends FacadeVFS {
this.#mapIdToFile.delete(fileId);

if (file?.flags & VFS.SQLITE_OPEN_MAIN_DB) {
if (file.abortController.signal.aborted) {
// The journal belongs to a view that was never stored.
this.#mapPathToFile.delete(file.path + '-journal');
}
file.broadcastChannel.close();
file.viewReleaser?.();
}
Expand Down Expand Up @@ -418,6 +393,11 @@ export class IDBMirrorVFS extends FacadeVFS {
if (lockType <= file.lockState) return VFS.SQLITE_OK;
switch (lockType) {
case VFS.SQLITE_LOCK_SHARED:
if (file.abortController.signal.aborted) {
// SQLite validates its cache against the file here, so this is
// where a view including an aborted commit can be replaced.
await this.#reloadFile(file);
}
if (file.txWriteHint) {
// xFileControl() has hinted that this transaction will
// write. Acquire the hint lock, which is required to reach
Expand Down Expand Up @@ -471,6 +451,12 @@ export class IDBMirrorVFS extends FacadeVFS {
}

console.assert(entries[0]?.txId === file.viewTx.txId || !file.viewTx.txId);
if (file.abortController.signal.aborted) {
// Our view includes a commit that was never stored. Make the
// transaction start over, which reloads the file.
file.locks.reserved();
return VFS.SQLITE_BUSY;
}
break;
case VFS.SQLITE_LOCK_EXCLUSIVE:
await this.#lock(file, 'write');
Expand Down Expand Up @@ -670,54 +656,152 @@ export class IDBMirrorVFS extends FacadeVFS {
* @param {File} file
*/
async #commitTx(file) {
if (file.abortController.signal.aborted) {
// This transaction was built on a commit that was never stored.
// SQLite discards its cache after this error, so the view can be
// reloaded, unless SQLite rolls back from a journal written on it.
this.#dropTx(file);
if (!this.#mapPathToFile.has(file.path + '-journal')) {
await this.#reloadFile(file);
}
throw new Error('an earlier commit was aborted');
}

// Advance our own view. Even if we received our own broadcasts (we
// don't), we want our view to be updated synchronously.
this.#acceptTx(file, file.txActive);
this.#setView(file, file.txActive);

const oldestTxId = await this.#getOldestTxInUse(file);

// Update IndexedDB page data.
const tx = file.txActive;
const idbTx = this.#idb.transaction(['blocks', 'tx'], 'readwrite');
const blocks = idbTx.objectStore('blocks');
for (const [offset, data] of file.txActive.blocks) {
blocks.put({ path: file.path, offset, data });
}
const complete = new Promise(async (resolve, reject) => {
idbTx.oncomplete = () => {
file.commitsPending--;
file.broadcastChannel.postMessage(tx);
resolve();
};
idbTx.onabort = () => {
// Our view includes this transaction, so it must be reloaded.
file.commitsPending--;
file.abortController.abort();
reject(idbTx.error);
};

// Delete blocks past the end of the file.
blocks.delete(IDBKeyRange.bound(
[file.path, file.txActive.fileSize], [file.path, Infinity]));
if (file.commitsPending++ > 0) {
// There are other write transactions still in flight. Use
// a read request purely for the side effect of flushing
// the IDB pipeline. Then we can ensure we are not building
// the new transaction on a failed transaction.
const fence = idbTx.objectStore('tx').get([file.path, -1]);
await idbX(fence);
}

// Delete obsolete transactions no longer needed.
const oldRange = IDBKeyRange.bound(
[file.path, -Infinity], [file.path, oldestTxId],
false, true);
idbTx.objectStore('tx').delete(oldRange);
try {
if (file.abortController.signal.aborted) {
idbTx.abort();
return;
}

// Save transaction object. Omit page data as an optimization.
const txSansData = Object.assign({}, file.txActive);
txSansData.blocks = new Map(Array.from(file.txActive.blocks, ([k]) => [k, null]));
idbTx.objectStore('tx').put(txSansData);
// Update IndexedDB page data.
const blocks = idbTx.objectStore('blocks');
for (const [offset, data] of tx.blocks) {
blocks.put({ path: file.path, offset, data });
}

// Broadcast transaction once it commits.
const complete = new Promise((resolve, reject) => {
const message = file.txActive;
idbTx.oncomplete = () => {
file.broadcastChannel.postMessage(message);
resolve();
};
idbTx.onabort = () => reject(idbTx.error);
idbTx.commit();
// Delete blocks past the end of the file.
blocks.delete(IDBKeyRange.bound(
[file.path, tx.fileSize], [file.path, Infinity]));

// Delete obsolete transactions no longer needed.
const oldRange = IDBKeyRange.bound(
[file.path, -Infinity], [file.path, oldestTxId],
false, true);
idbTx.objectStore('tx').delete(oldRange);

// Save transaction object. Omit page data as an optimization.
const txSansData = Object.assign({}, tx);
txSansData.blocks = new Map(Array.from(tx.blocks, ([k]) => [k, null]));
idbTx.objectStore('tx').put(txSansData);
idbTx.commit();
} catch (e) {
// This code no longer runs inside a request event handler, so an
// exception would not abort the transaction on its own.
idbTx.abort();
throw e;
}
});

if (file.synchronous === 'full') {
await complete;
try {
await complete;
} catch (e) {
// SQLite discards its cache after this error, so the view can be
// reloaded now, even in exclusive locking mode.
await this.#reloadFile(file);
throw e;
}
} else {
// A failure is handled by the next transaction.
complete.catch(() => {});
}

file.txActive = null;
file.txWriteHint = false;
}

/**
* Load the stored blocks and the last transaction into a main database.
* @param {File} file
* @param {IDBTransaction} idbTx
*/
async #loadFile(file, idbTx) {
const blocks = new Map();
await new Promise((resolve, reject) => {
// TODO: Reading the database in batches with `getAll()` and a
// batch size count would be faster than using a cursor.
const range = IDBKeyRange.bound([file.path, 0], [file.path, Infinity]);
Comment thread
rhashimoto marked this conversation as resolved.
const request = idbTx.objectStore('blocks').openCursor(range);
request.onsuccess = () => {
const cursor = request.result;
if (cursor) {
const { offset, data } = cursor.value;
blocks.set(offset, data);
cursor.continue();
} else {
resolve();
}
};
request.onerror = () => reject(request.error);
});

const viewTx = await new Promise((resolve, reject) => {
const range = IDBKeyRange.bound([file.path, 0], [file.path, Infinity]);
const request = idbTx.objectStore('tx').openCursor(range, 'prev');
request.onsuccess = () => resolve(request.result?.value ?? { txId: 0 });
request.onerror = () => reject(request.error);
});

file.blocks = blocks;
file.blockSize = blocks.get(0)?.byteLength ?? 0;
file.viewTx = viewTx;
}

/**
* Replace a view that includes an aborted commit with the stored one.
* @param {File} file
*/
async #reloadFile(file) {
await this.#loadFile(file, this.#idb.transaction(['blocks', 'tx']));
file.txActive = null;
file.abortController = new AbortController();
await this.#setView(file, file.viewTx);
if (file.lockState === VFS.SQLITE_LOCK_NONE) {
this.#processBroadcasts(file);
}
}

/**
* @param {File} file
*/
Expand Down
2 changes: 2 additions & 0 deletions test/IDBMirrorVFS.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import { vfs_xRead } from "./vfs_xRead.js";
import { vfs_xWrite } from "./vfs_xWrite.js";
import { vfs_leak } from "./vfs_leak.js";
import { vfs_rollback } from "./vfs_rollback.js";
import { vfs_commit_abort } from "./vfs_commit_abort.js";

const CONFIG = 'IDBMirrorVFS';
const BUILDS = ['asyncify', 'jspi'];
Expand All @@ -26,6 +27,7 @@ describe(CONFIG, function() {
vfs_xWrite(context);
vfs_leak(context);
vfs_rollback(context);
vfs_commit_abort({ build });
});
}
});
108 changes: 108 additions & 0 deletions test/vfs_commit_abort-worker.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
// A worker for vfs_commit_abort.js, holding one connection on
// IDBMirrorVFS. It can make the next IndexedDB commit abort, as a full
// quota would, optionally after keeping the transaction alive.
import * as SQLite from '../src/sqlite-api.js';
import { IDBMirrorVFS } from '../src/examples/IDBMirrorVFS.js';

const BUILDS = new Map([
['default', '../dist/wa-sqlite.mjs'],
['asyncify', '../dist/wa-sqlite-async.mjs'],
['jspi', '../dist/wa-sqlite-jspi.mjs'],
]);

const searchParams = new URLSearchParams(location.search);

let abortNextCommit = false;
let abortDelay = 0;
const commit = IDBTransaction.prototype.commit;
IDBTransaction.prototype.commit = function() {
if (abortNextCommit && this.mode === 'readwrite') {
abortNextCommit = false;
if (!abortDelay) return this.abort();

// Keep the transaction active with requests, then abort it.
const deadline = Date.now() + abortDelay;
const store = this.objectStore('tx');
const keepAlive = () => {
if (Date.now() < deadline) {
store.get(['', 0]).onsuccess = keepAlive;
} else {
this.abort();
}
};
return keepAlive();
}
return commit.call(this);
};

const ready = (async () => {
const { default: moduleFactory } = await import(BUILDS.get(searchParams.get('build')));
const module = await moduleFactory();
const sqlite3 = SQLite.Factory(module);
const vfs = await IDBMirrorVFS.create(searchParams.get('idb'), module);
sqlite3.vfs_register(vfs, true);
return { sqlite3, db: await sqlite3.open_v2('commit-abort') };
})();

// A read-only transaction completes after the read-write transactions
// created before it, and after their completion or abort events. Closing
// before a commit completes makes its broadcast throw, which is not what
// these tests are about.
async function commitsFinished() {
const idb = await new Promise((resolve, reject) => {
const request = indexedDB.open(searchParams.get('idb'));
request.onsuccess = () => resolve(request.result);
request.onerror = () => reject(request.error);
});
await new Promise(resolve => {
const tx = idb.transaction(['blocks', 'tx']);
tx.objectStore('tx').count();
tx.oncomplete = tx.onabort = resolve;
});
idb.close();
}

addEventListener('message', async ({ data }) => {
try {
const state = await ready;
const { sqlite3 } = state;
switch (data.type) {
case 'exec': {
const rows = [];
try {
await sqlite3.exec(state.db, data.sql, row => rows.push(row));
postMessage({ rows });
} catch (e) {
postMessage({ rows, error: e.message });
}
break;
}
case 'abort-next-commit':
abortNextCommit = true;
abortDelay = data.delay ?? 0;
postMessage({});
break;
case 'sleep':
await new Promise(resolve => setTimeout(resolve, data.ms));
postMessage({});
break;
case 'settle':
await commitsFinished();
postMessage({});
break;
case 'reopen':
await commitsFinished();
await sqlite3.close(state.db);
state.db = await sqlite3.open_v2('commit-abort');
postMessage({});
break;
case 'close':
await commitsFinished();
await sqlite3.close(state.db);
postMessage({});
break;
}
} catch (e) {
postMessage({ error: e.message });
}
});
Loading
Loading