diff --git a/src/browser/js/Scheduler.zig b/src/browser/js/Scheduler.zig index 34844fb16..7d26bf5b6 100644 --- a/src/browser/js/Scheduler.zig +++ b/src/browser/js/Scheduler.zig @@ -36,12 +36,18 @@ const Queue = std.PriorityQueue(Task, void, struct { const Scheduler = @This(); _sequence: u64, +// Some things (e.g. IndexedDB) can have operations that are only valid for a +// specific task boundary. So every time we start a task, we increment the +// scheduler's generation. Code can snapshot this version and then compare it +// later to see if we're still in the same task. +generation: u64, low_priority: Queue, high_priority: Queue, pub fn init(allocator: std.mem.Allocator) Scheduler { return .{ ._sequence = 0, + .generation = 0, .low_priority = Queue.init(allocator, {}), .high_priority = Queue.init(allocator, {}), }; @@ -116,6 +122,8 @@ fn runQueue(self: *Scheduler, queue: *Queue) !void { log.debug(.scheduler, "scheduler.runTask", .{ .name = task.name }); } + self.generation +%= 1; + const repeat_in_ms = task.callback(task.ctx) catch |err| { log.warn(.scheduler, "task.callback", .{ .name = task.name, .err = err }); continue; diff --git a/src/browser/tests/indexeddb.html b/src/browser/tests/indexeddb.html index fe9e9b0e5..aee2f92ea 100644 --- a/src/browser/tests/indexeddb.html +++ b/src/browser/tests/indexeddb.html @@ -1120,4 +1120,37 @@ }); } + + + diff --git a/src/browser/webapi/storage/idb/Engine.zig b/src/browser/webapi/storage/idb/Engine.zig index 6638ca5ca..8b8d73b0c 100644 --- a/src/browser/webapi/storage/idb/Engine.zig +++ b/src/browser/webapi/storage/idb/Engine.zig @@ -28,6 +28,60 @@ const Engine = @This(); conn: Sqlite.Conn, +// A single sqlite connection backs every database of an origin, so only one +// transaction (or open/delete) may hold it in a `begin`/`commit` bracket at a +// time. `_gate_owner` is whoever currently holds it; contenders park their +// (intrusive, caller-owned) node on `_gate_waiters` and are handed ownership in +// FIFO order as it's released. +// +// LIFETIME GAP (to resolve in the IDB memory rework): the Engine is +// session-scoped but the waiters (IDBTransaction / Open+DeleteContext) are +// page-scoped, so these are the only session->page pointers in IDB. A waiter +// freed on navigation while still owning or parked here leaves a dangling +// pointer/node that a later same-origin page can deref or deadlock on; a parked +// waiter has no scheduler task and so no finalizer to unlink it on teardown. The +// rework must session-scope these objects, add a page-teardown unlink, or drop +// the wait-list for a bool gate. +_gate_owner: ?*GateWaiter = null, +_gate_waiters: std.DoublyLinkedList = .{}, + +pub const GateWaiter = struct { + node: std.DoublyLinkedList.Node = .{}, + // Called when this waiter is handed the gate; typically reschedules the + // owner's task so it re-runs and finds itself the owner. + wake: *const fn (waiter: *GateWaiter) void, +}; + +pub fn acquireGate(self: *Engine, waiter: *GateWaiter) bool { + if (self._gate_owner == null) { + self._gate_owner = waiter; + return true; + } + + if (self._gate_owner == waiter) { + return true; + } + + self._gate_waiters.append(&waiter.node); + return false; +} + +// Release the gate held by `waiter`, handing it directly to the next parked +// waiter (no window where an unrelated contender can grab it) and waking it. +// No-op if `waiter` isn't the current owner. +pub fn releaseGate(self: *Engine, waiter: *GateWaiter) void { + if (self._gate_owner != waiter) { + return; + } + if (self._gate_waiters.popFirst()) |node| { + const next: *GateWaiter = @fieldParentPtr("node", node); + self._gate_owner = next; + next.wake(next); + return; + } + self._gate_owner = null; +} + pub fn open(path: [:0]const u8) !Engine { const conn = try Sqlite.Conn.open(path); errdefer conn.close(); diff --git a/src/browser/webapi/storage/idb/IDBFactory.zig b/src/browser/webapi/storage/idb/IDBFactory.zig index 52a9dca45..4b2fffc8b 100644 --- a/src/browser/webapi/storage/idb/IDBFactory.zig +++ b/src/browser/webapi/storage/idb/IDBFactory.zig @@ -52,6 +52,7 @@ pub fn open(_: *IDBFactory, name: []const u8, version: ?u64, exec: *Execution) ! .name = try exec.dupeString(name), .version = version, .exec = exec, + ._gate_waiter = .{ .wake = OpenContext.wakeUp }, }); try exec.js.scheduler.add(ctx, OpenContext.run, 0, .{ @@ -66,6 +67,8 @@ const OpenContext = struct { name: []const u8, version: ?u64, exec: *Execution, + // Our node in the engine's connection gate wait-list. See Engine.acquireGate. + _gate_waiter: Engine.GateWaiter, fn cancelled(ctx: *anyopaque) void { const self: *OpenContext = @ptrCast(@alignCast(ctx)); @@ -74,9 +77,24 @@ const OpenContext = struct { fn run(ctx: *anyopaque) !?u32 { const self: *OpenContext = @ptrCast(@alignCast(ctx)); - defer self.exec._factory.destroy(self); - self.runOpen() catch |err| { + const engine = self.resolveEngine() catch |err| { + self.exec._factory.destroy(self); + self.request.setError(err); + self.request.deliver(self.exec) catch {}; + return null; + }; + + // An open that upgrades runs a versionchange transaction on the shared + // connection, so it must serialize with other transactions/opens. Park + // on the gate if it's held; wakeUp re-runs us when it's handed over. + if (!engine.acquireGate(&self._gate_waiter)) { + return null; // parked; not destroyed + } + defer self.exec._factory.destroy(self); + defer engine.releaseGate(&self._gate_waiter); + + self.runOpen(engine) catch |err| { log.warn(.storage, "idb open", .{ .err = err, .name = self.name }); self.request.setError(err); self.request.deliver(self.exec) catch {}; @@ -84,14 +102,29 @@ const OpenContext = struct { return null; } - fn runOpen(self: *OpenContext) !void { - const exec = self.exec; + // Scheduler wake-up: the connection gate was handed to us, so re-run. + fn wakeUp(waiter: *Engine.GateWaiter) void { + const self: *OpenContext = @fieldParentPtr("_gate_waiter", waiter); + self.exec.js.scheduler.add(self, run, 0, .{ + .name = "IDBFactory.open", + .finalizer = cancelled, + }) catch |err| { + // We were handed the gate; if we can't reschedule, hand it off so the + // waiters behind us aren't stranded. + if (self.resolveEngine()) |engine| engine.releaseGate(&self._gate_waiter) else |_| {} + log.warn(.storage, "idb resume open", .{ .err = err }); + }; + } + fn resolveEngine(self: *OpenContext) !*Engine { // origin being null was already guarded against, so this should be // unreachable, but this is safer. - const origin = exec.origin() orelse return error.SecurityError; + const origin = self.exec.origin() orelse return error.SecurityError; + return self.exec.session.idb.engineForOrigin(origin); + } - const engine = try exec.session.idb.engineForOrigin(origin); + fn runOpen(self: *OpenContext, engine: *Engine) !void { + const exec = self.exec; const existing = try engine.databaseVersion(self.name); // No explicit version means "open at the current version" (or 1 for a @@ -166,6 +199,7 @@ pub fn deleteDatabase(_: *IDBFactory, name: []const u8, exec: *Execution) !*IDBR .request = request, .name = try exec.dupeString(name), .exec = exec, + ._gate_waiter = .{ .wake = DeleteContext.wakeUp }, }); try exec.js.scheduler.add(ctx, DeleteContext.run, 0, .{ @@ -179,6 +213,7 @@ const DeleteContext = struct { request: *IDBRequest, name: []const u8, exec: *Execution, + _gate_waiter: Engine.GateWaiter, fn cancelled(ctx: *anyopaque) void { const self: *DeleteContext = @ptrCast(@alignCast(ctx)); @@ -187,9 +222,21 @@ const DeleteContext = struct { fn run(ctx: *anyopaque) !?u32 { const self: *DeleteContext = @ptrCast(@alignCast(ctx)); - defer self.exec._factory.destroy(self); - self.runDelete() catch |err| { + const engine = self.resolveEngine() catch |err| { + self.exec._factory.destroy(self); + self.request.setError(err); + self.request.deliver(self.exec) catch {}; + return null; + }; + + if (!engine.acquireGate(&self._gate_waiter)) { + return null; // parked; not destroyed + } + defer self.exec._factory.destroy(self); + defer engine.releaseGate(&self._gate_waiter); + + self.runDelete(engine) catch |err| { log.warn(.storage, "idb deleteDatabase", .{ .err = err, .name = self.name }); self.request.setError(err); self.request.deliver(self.exec) catch {}; @@ -197,12 +244,28 @@ const DeleteContext = struct { return null; } - fn runDelete(self: *DeleteContext) !void { - const exec = self.exec; - const origin = exec.origin() orelse return error.SecurityError; - const engine = try exec.session.idb.engineForOrigin(origin); + // Scheduler wake-up: the connection gate was handed to us, so re-run. + fn wakeUp(waiter: *Engine.GateWaiter) void { + const self: *DeleteContext = @fieldParentPtr("_gate_waiter", waiter); + self.exec.js.scheduler.add(self, run, 0, .{ + .name = "IDBFactory.deleteDatabase", + .finalizer = cancelled, + }) catch |err| { + // We were handed the gate; if we can't reschedule, hand it off so the + // waiters behind us aren't stranded. + if (self.resolveEngine()) |engine| engine.releaseGate(&self._gate_waiter) else |_| {} + log.warn(.storage, "idb resume delete", .{ .err = err }); + }; + } + + fn resolveEngine(self: *DeleteContext) !*Engine { + const origin = self.exec.origin() orelse return error.SecurityError; + return self.exec.session.idb.engineForOrigin(origin); + } + + fn runDelete(self: *DeleteContext, engine: *Engine) !void { try engine.deleteDatabase(self.name); - return self.request.fireSuccess(exec); + return self.request.fireSuccess(self.exec); } }; diff --git a/src/browser/webapi/storage/idb/IDBTransaction.zig b/src/browser/webapi/storage/idb/IDBTransaction.zig index ab0cd0985..549a66239 100644 --- a/src/browser/webapi/storage/idb/IDBTransaction.zig +++ b/src/browser/webapi/storage/idb/IDBTransaction.zig @@ -34,6 +34,7 @@ const DOMStringList = @import("../../collections.zig").DOMStringList; const log = lp.log; const Execution = js.Execution; const FunctionSetter = idb.FunctionSetter; +const IS_DEBUG = @import("builtin").mode == .Debug; const IDBTransaction = @This(); @@ -48,11 +49,22 @@ _mode: Mode, _scope: []const []const u8 = &.{}, // Advisory only: we always commit through sqlite. Stored to expose the property. _durability: Durability = .default, -_requests: std.ArrayList(*IDBRequest) = .empty, + +// request queue, swaps between &_queue_a and &_queue_b so that, as we drain, new +// requests are queued in the new queue and will be processed on the next drain +_queue: *std.ArrayList(*IDBRequest), +_queue_a: std.ArrayList(*IDBRequest) = .empty, +_queue_b: std.ArrayList(*IDBRequest) = .empty, + _begun: bool = false, _settled: bool = false, _aborted: bool = false, _committing: bool = false, +_gate_waiter: Engine.GateWaiter, +// A transaction is only active for one execution of a Scheduler's task. We +// capture the scheduler's generation here and reject any request made in a +// later generation (see assertActive). +_active_turn: u64 = 0, _on_complete: ?js.Function.Global = null, _on_error: ?js.Function.Global = null, @@ -89,7 +101,12 @@ pub fn init(exec: *Execution, db: *IDBDatabase, mode: Mode, durability: Durabili ._mode = mode, ._scope = scope, ._durability = durability, + ._active_turn = exec.js.scheduler.generation, + ._queue = undefined, + ._gate_waiter = undefined, }); + self._queue = &self._queue_a; + self._gate_waiter = .{ .wake = resumeDrain }; // Schedule the drain even for an empty transaction so it still `complete`s. try exec.js.scheduler.add(self, drain, 0, .{ @@ -101,14 +118,21 @@ pub fn init(exec: *Execution, db: *IDBDatabase, mode: Mode, durability: Durabili // We need a "special" transaction for upgradeneeded pub fn initVersionChange(exec: *Execution, db: *IDBDatabase) !*IDBTransaction { - return exec._factory.eventTarget(IDBTransaction{ + const self = try exec._factory.eventTarget(IDBTransaction{ ._proto = undefined, ._exec = exec, ._db = db, ._engine = db._engine, ._mode = .versionchange, ._begun = true, + ._queue = undefined, + ._gate_waiter = undefined, }); + self._queue = &self._queue_a; + // A versionchange transaction never contends for the gate (the open path + // holds it); keep the node well-formed so releaseGate's owner check no-ops. + self._gate_waiter = .{ .wake = resumeDrain }; + return self; } pub fn asEventTarget(self: *IDBTransaction) *EventTarget { @@ -144,50 +168,70 @@ pub fn abort(self: *IDBTransaction, exec: *Execution) !void { if (self._begun) { self._engine.rollback(); } + self._engine.releaseGate(&self._gate_waiter); - // request whose op already ran (op == .none) was delivered and is skipped. - for (self._requests.items, 0..) |request, i| { - if (i != request._txn_index or request._op == .none) { - continue; + for ([_]*std.ArrayList(*IDBRequest){ &self._queue_a, &self._queue_b }) |queue| { + for (queue.items, 0..) |request, i| { + if (i != request._txn_index or request._op == .none) { + continue; + } + request._op = .none; + request.setError(error.AbortError); + request.deliver(exec) catch |err| { + log.warn(.storage, "idb abort deliver", .{ .err = err }); + }; } - request._op = .none; - request.setError(error.AbortError); - request.deliver(exec) catch |err| { - log.warn(.storage, "idb abort deliver", .{ .err = err }); - }; } self.fire(exec, comptime .wrap("abort"), self._on_abort); } -// Deliver pending requests, commit, then fire `complete` (or `abort` if the -// commit fails). Drives a transaction to a successful close — shared by the -// drain task and the open path's synchronous versionchange transaction. pub fn settle(self: *IDBTransaction, exec: *Execution) void { + if (comptime IS_DEBUG) { + // non versionchange mode goes through the scheduler + drain + std.debug.assert(self._mode == .versionchange); + } + if (self._settled) { return; } - self.deliverPending(exec); - if (self._settled) { - // a request handler might have settled this (e.g. called abort) - return; + // Deliver batches until the queue stays empty — a handler may enqueue more. + while (self._queue.items.len > 0) { + self.deliverBatch(exec); + if (self._settled) { + // a request handler might have settled this (e.g. called abort) + return; + } } + self.commitAndComplete(exec); +} + +// Commit the underlying sqlite transaction (if begun), release the connection +// gate, then fire `complete` — or `abort` if the commit fails. +fn commitAndComplete(self: *IDBTransaction, exec: *Execution) void { if (self._begun) { self._engine.commit() catch |err| { log.warn(.storage, "idb commit", .{ .err = err, .sqlite = self._engine.conn.lastError() }); self._engine.rollback(); + self._engine.releaseGate(&self._gate_waiter); self.fire(exec, comptime .wrap("abort"), self._on_abort); return; }; } + self._engine.releaseGate(&self._gate_waiter); self.fire(exec, comptime .wrap("complete"), self._on_complete); } -// "is this transaction still usable". Once settled or explicitly committing, it no -// longer accepts new requests. +// "is this transaction still usable". Once settled or explicitly committing, it +// no longer accepts new requests; nor does it outside its active turn (a request +// made from an unrelated task). A versionchange transaction runs synchronously +// during upgradeneeded and stays active until settled, so it skips the turn check. pub fn assertActive(self: *const IDBTransaction) !void { if (self._settled or self._committing) { return error.TransactionInactiveError; } + if (self._mode != .versionchange and self._active_turn != self._exec.js.scheduler.generation) { + return error.TransactionInactiveError; + } } pub fn ensureBegun(self: *IDBTransaction) !void { @@ -209,8 +253,8 @@ pub fn newRequest(self: *IDBTransaction) !*IDBRequest { } pub fn enqueue(self: *IDBTransaction, request: *IDBRequest) !void { - request._txn_index = self._requests.items.len; - try self._requests.append(self._exec.arena, request); + request._txn_index = self._queue.items.len; + try self._queue.append(self._exec.arena, request); } pub fn objectStore(self: *IDBTransaction, name: []const u8, exec: *Execution) !*IDBObjectStore { @@ -297,25 +341,70 @@ fn fire(self: *IDBTransaction, exec: *Execution, typ: lp.String, handler: ?js.Fu fn drain(ctx: *anyopaque) !?u32 { const self: *IDBTransaction = @ptrCast(@alignCast(ctx)); if (self._settled) { - // could have already been settled, e.g. via an abort) + // Already settled (e.g. via an abort). + self._engine.releaseGate(&self._gate_waiter); return null; } - self.settle(self._exec); + + const exec = self._exec; + + if (self._queue.items.len > 0) { + if (self._engine.acquireGate(&self._gate_waiter) == false) { + return null; // parked; resumeDrain reschedules us + } + + self.deliverBatch(exec); + if (self._settled) { + // a handler aborted us mid-delivery; abort() released the gate. + return null; + } + if (self._queue.items.len > 0) { + // handlers enqueued more — keep the gate and resume next turn. + return 1; + } + } + + // Nothing left to deliver: commit and fire `complete` (releases the gate). + self.commitAndComplete(exec); return null; } +// Scheduler wake-up: the gate was handed to us, so run the drain again. +fn resumeDrain(waiter: *Engine.GateWaiter) void { + const self: *IDBTransaction = @fieldParentPtr("_gate_waiter", waiter); + if (comptime IS_DEBUG) { + std.debug.assert(self._mode != .versionchange); + } + + self._exec.js.scheduler.add(self, drain, 0, .{ + .name = "IDBTransaction.drain", + .finalizer = finalize, + }) catch |err| { + self._engine.releaseGate(&self._gate_waiter); + log.warn(.storage, "idb resume drain", .{ .err = err }); + }; +} + fn finalize(ctx: *anyopaque) void { const self: *IDBTransaction = @ptrCast(@alignCast(ctx)); if (self._begun and !self._settled) { self._engine.rollback(); } + self._engine.releaseGate(&self._gate_waiter); } -fn deliverPending(self: *IDBTransaction, exec: *Execution) void { +// Deliver the requests queued as of now, one queue's worth. Requests enqueued +// by handlers during delivery go to the other queue and are handled on a later +// turn. +fn deliverBatch(self: *IDBTransaction, exec: *Execution) void { + const batch = self._queue; + // New requests now accumulate in the other queue. + self._queue = if (batch == &self._queue_a) &self._queue_b else &self._queue_a; + defer batch.clearRetainingCapacity(); + // exec.js.local must be non-null. In some cases it is (e.g. tx.commit() // called directly from JS), in others is isn't (drain schedule tasks). // Easier to explicitly create one and then restore whatever was there before. - const prev_local = exec.js.local; defer exec.js.local = prev_local; @@ -325,14 +414,17 @@ fn deliverPending(self: *IDBTransaction, exec: *Execution) void { exec.js.local = &ls.local; - var i: usize = 0; - while (i < self._requests.items.len) : (i += 1) { + // The transaction is active for this batch's dispatch and the rest of this + // task's turn (including the microtasks it queues). generation is constant + // within a task, so stamp it once. + self._active_turn = exec.js.scheduler.generation; + + for (batch.items) |request| { // A handler may have aborted the transaction mid-delivery; abort() already // delivered AbortError to the remaining requests, so stop here. if (self._settled) { return; } - const request = self._requests.items[i]; request.execute(exec) catch |err| { log.warn(.storage, "idb request execute", .{ .err = err }); request.setError(err);