add IDBIndex

This commit is contained in:
Karl Seguin committed 2026-07-08 10:19:38 +08:00
1 parent edbd24bbac
commit 437dac68ef
9 files changed
+1006 -114

No files matched your search

+25 -14
View File
@@ -256,22 +256,13 @@ pub fn dispatchDirect(
// Per spec, currentTarget is only set while listeners are being invoked
defer event._current_target = null;
// Call the property handler (e.g., onmessage) if present
if (getFunction(handler, &ls.local)) |func| {
event._current_target = target;
_ = func.callWithThis(void, target, .{event}) catch |err| {
log.warn(.event, opts.context, .{ .err = err });
};
}
// Call listeners registered via addEventListener
const list = self.getListeners(target, event._type_string) orelse return;
// This is a slightly simplified version of what you'll find in EventManager.
// dispatchPhase. It is simpler because, for direct dispatching, we know
// there's no ancestors and only the single target phase.
// Track dispatch depth for deferred removal
// Track dispatch depth for deferred removal. Bump it *before* the property
// handler runs so any listener it removes is deferred (keeping our sentinel
// node alive) rather than freed mid-dispatch.
self.dispatch_depth += 1;
defer {
self.dispatch_depth -= 1;
@@ -285,8 +276,28 @@ pub fn dispatchDirect(
}
}
// Use the last listener in the list as sentinel - listeners added during dispatch will be after it
const last_node = list.last orelse return;
// Snapshot the listener list *before* invoking the property handler. Per
// spec the set of listeners is collected at the start of dispatch, so a
// listener added while we're dispatching — including one added by the
// property handler itself (e.g. onupgradeneeded calling addEventListener) —
// must not be invoked for this event.
const maybe_list = self.getListeners(target, event._type_string);
const sentinel = if (maybe_list) |list| list.last else null;
// Call the property handler (e.g., onmessage) if present
if (getFunction(handler, &ls.local)) |func| {
event._current_target = target;
_ = func.callWithThis(void, target, .{event}) catch |err| {
log.warn(.event, opts.context, .{ .err = err });
};
}
// No listeners were registered via addEventListener at dispatch start.
const last_node = sentinel orelse return;
const list = maybe_list.?;
// Use the last listener present at dispatch start as sentinel - listeners
// added during dispatch will be after it
const last_listener: *Listener = @alignCast(@fieldParentPtr("node", last_node));
// Iterate through the list, stopping after we've encountered the last_listener
+144 -2
View File
@@ -647,18 +647,21 @@
open.onsuccess = (e) => {
const store = e.target.result.transaction("s", "readonly").objectStore("s");
const seen = [];
let sourceIsStore = null;
const req = store.openCursor();
req.onsuccess = () => {
const cursor = req.result;
if (cursor) {
if (sourceIsStore === null) sourceIsStore = cursor.source === store;
seen.push([cursor.key, cursor.primaryKey, cursor.value, cursor.direction]);
cursor.continue();
} else {
state.resolve(seen);
state.resolve({ seen, sourceIsStore });
}
};
};
await state.done((seen) => {
await state.done((got) => {
const seen = got.seen;
testing.expectEqual(3, seen.length);
testing.expectEqual(1, seen[0][0]); // key
testing.expectEqual(1, seen[0][1]); // primaryKey == key for object store
@@ -666,6 +669,7 @@
testing.expectEqual("next", seen[0][3]); // direction
testing.expectEqual(3, seen[2][0]);
testing.expectEqual("c", seen[2][2]);
testing.expectEqual(true, got.sourceIsStore); // cursor.source === the store handle
});
}
</script>
@@ -773,6 +777,144 @@
}
</script>
<!-- createIndex during upgrade, populated from existing data; index get/getKey/
getAll/count look records up by a non-primary property. -->
<script id="index_basic" type=module>
{
const state = await testing.async();
const open = indexedDB.open("idx-db", 1);
open.onupgradeneeded = (e) => {
const s = e.target.result.createObjectStore("people", { keyPath: "id" });
s.add({ id: 1, age: 30, name: "Ann" }); // seeded before the index exists
const idx = s.createIndex("byAge", "age");
testing.expectEqual("byAge", idx.name);
testing.expectEqual("age", idx.keyPath);
testing.expectEqual(false, idx.unique);
s.add({ id: 2, age: 25, name: "Bob" }); // indexed on insert
s.add({ id: 3, age: 25, name: "Cy" });
};
open.onsuccess = (e) => {
const store = e.target.result.transaction("people").objectStore("people");
testing.expectEqual("byAge", store.indexNames[0]);
const idx = store.index("byAge");
const g = idx.get(25); // first record (by primaryKey) with age 25
const gk = idx.getKey(25); // its primary key
const all = idx.getAll(25); // all records with age 25
const cnt = idx.count(); // total index entries
cnt.onsuccess = () => state.resolve({
got: g.result, key: gk.result, all: all.result, count: cnt.result,
});
};
await state.done((r) => {
testing.expectEqual("Bob", r.got.name); // lowest primaryKey among age 25
testing.expectEqual(2, r.key);
testing.expectEqual(2, r.all.length);
testing.expectEqual(3, r.count);
});
}
</script>
<!-- A unique index rejects a duplicate key with ConstraintError, and the failed
write leaves the store unchanged (savepoint rollback). -->
<script id="index_unique" type=module>
{
const state = await testing.async();
const open = indexedDB.open("idxu-db", 1);
open.onupgradeneeded = (e) => {
const s = e.target.result.createObjectStore("u", { keyPath: "id" });
s.createIndex("email", "email", { unique: true });
s.add({ id: 1, email: "a@x.com" });
};
open.onsuccess = (e) => {
const store = e.target.result.transaction("u", "readwrite").objectStore("u");
const dup = store.add({ id: 2, email: "a@x.com" }); // collides on the unique index
dup.onerror = (ev) => {
ev.preventDefault();
const c = store.count();
c.onsuccess = () => state.resolve({ err: dup.error.name, count: c.result });
};
};
await state.done((r) => {
testing.expectEqual("ConstraintError", r.err);
testing.expectEqual(1, r.count); // id:2 was rolled back
});
}
</script>
<!-- An index cursor exposes key (index key), primaryKey (store key) and value,
in index-key order; multiEntry indexes one entry per array element. -->
<script id="index_cursor_multientry" type=module>
{
const state = await testing.async();
const open = indexedDB.open("idxc-db", 1);
open.onupgradeneeded = (e) => {
const s = e.target.result.createObjectStore("posts", { keyPath: "id" });
s.createIndex("tags", "tags", { multiEntry: true });
s.add({ id: 10, tags: ["a", "b"] });
s.add({ id: 20, tags: ["b", "c"] });
};
open.onsuccess = (e) => {
const store = e.target.result.transaction("posts").objectStore("posts");
const idx = store.index("tags");
// Cursor over the whole index: (indexKey, primaryKey) pairs.
const seen = [];
const cur = idx.openCursor();
cur.onsuccess = () => {
const c = cur.result;
if (c) { seen.push([c.key, c.primaryKey, c.value.id]); c.continue(); return; }
// "b" appears in both posts.
const bPosts = idx.getAll("b");
bPosts.onsuccess = () => state.resolve({ seen, bCount: bPosts.result.length });
};
};
await state.done((r) => {
// tags expand to: a->10, b->10, b->20, c->20 (index-key then primaryKey order)
testing.expectEqual(4, r.seen.length);
testing.expectEqual("a", r.seen[0][0]);
testing.expectEqual(10, r.seen[0][1]);
testing.expectEqual("b", r.seen[1][0]);
testing.expectEqual(10, r.seen[1][1]);
testing.expectEqual(20, r.seen[2][1]);
testing.expectEqual(2, r.bCount);
});
}
</script>
<!-- deleteDatabase removes index data too: re-creating a unique index after a
delete must not collide with orphaned index records (sqlite reuses ids). -->
<script id="index_delete_db_isolation" type=module>
{
const state = await testing.async();
const open1 = indexedDB.open("idx-iso-db", 1);
open1.onupgradeneeded = (e) => {
const s = e.target.result.createObjectStore("s", { keyPath: "id" });
s.createIndex("email", "email", { unique: true });
s.add({ id: 1, email: "a@x.com" });
};
open1.onsuccess = (e) => {
e.target.result.close();
const del = indexedDB.deleteDatabase("idx-iso-db");
del.onsuccess = () => {
const open2 = indexedDB.open("idx-iso-db", 1);
open2.onupgradeneeded = (e2) => {
const s = e2.target.result.createObjectStore("s", { keyPath: "id" });
s.createIndex("email", "email", { unique: true });
};
open2.onsuccess = (e2) => {
const store = e2.target.result.transaction("s", "readwrite").objectStore("s");
// Same email as the deleted db — must succeed, not hit a stale index row.
const req = store.add({ id: 1, email: "a@x.com" });
req.onsuccess = () => state.resolve({ ok: true });
req.onerror = () => state.resolve({ ok: false, err: req.error.name });
};
};
};
await state.done((r) => testing.expectEqual(true, r.ok));
}
</script>
<!-- clear() empties the store; commit() settles the transaction explicitly. -->
<script id="clear_and_commit" type=module>
{
+312 -19
View File
@@ -34,27 +34,51 @@ pub fn open(path: [:0]const u8) !Engine {
try conn.busyTimeout(1000);
try conn.exec("pragma journal_mode=wal", .{});
try conn.exec("pragma foreign_keys = on", .{});
try conn.exec(
\\ create table if not exists idb_databases (
\\ id integer primary key,
\\ name text not null unique,
\\ version integer not null
\\ );
\\
\\ create table if not exists idb_object_stores (
\\ id integer primary key,
\\ database_id integer not null references idb_databases(id),
\\ database_id integer not null references idb_databases(id) on delete cascade,
\\ name text not null,
\\ key_path text,
\\ auto_increment integer not null default 0,
\\ key_generator integer not null default 1,
\\ unique(database_id, name)
\\ );
\\
\\ create table if not exists idb_records (
\\ object_store_id integer not null references idb_object_stores(id),
\\ object_store_id integer not null references idb_object_stores(id) on delete cascade,
\\ key blob not null,
\\ value blob not null,
\\ primary key (object_store_id, key)
\\ ) without rowid;
\\
\\ create table if not exists idb_indexes (
\\ id integer primary key,
\\ object_store_id integer not null references idb_object_stores(id) on delete cascade,
\\ name text not null,
\\ key_path text not null,
\\ is_unique integer not null default 0,
\\ multi_entry integer not null default 0,
\\ unique(object_store_id, name)
\\ );
\\
\\ create table if not exists idb_index_records (
\\ index_id integer not null references idb_indexes(id) on delete cascade,
\\ key blob not null,
\\ primary_key blob not null,
\\ is_unique integer not null,
\\ primary key (index_id, key, primary_key)
\\ ) without rowid;
\\ create unique index if not exists idb_index_unique
\\ on idb_index_records(index_id, key) where is_unique = 1;
, .{});
return .{ .conn = conn };
}
@@ -94,17 +118,7 @@ pub fn upsertDatabase(self: *Engine, name: []const u8, version: i64) !i64 {
}
pub fn deleteDatabase(self: *Engine, name: []const u8) !void {
const database_id = (try self.databaseId(name)) orelse return;
try self.begin();
errdefer self.rollback();
try self.conn.exec(
"delete from idb_records where object_store_id in (select id from idb_object_stores where database_id = ?1)",
.{database_id},
);
try self.conn.exec("delete from idb_object_stores where database_id = ?1", .{database_id});
try self.conn.exec("delete from idb_databases where id = ?1", .{database_id});
try self.commit();
return self.conn.exec("delete from idb_databases where name = ?1", .{name});
}
pub fn objectStoreId(self: *const Engine, database_id: i64, name: []const u8) !?i64 {
@@ -177,10 +191,15 @@ pub fn createObjectStore(
}
pub fn deleteObjectStore(self: *Engine, database_id: i64, name: []const u8) !void {
const store_id = (try self.objectStoreId(database_id, name)) orelse return error.NotFound;
// caller has a transaction open
try self.conn.exec("delete from idb_records where object_store_id = ?1", .{store_id});
try self.conn.exec("delete from idb_object_stores where id = ?1", .{store_id});
// caller has a transaction open; cascade drops records, indexes and index
// records.
const deleted = try self.conn.scalar(i64,
\\ delete from idb_object_stores where database_id = ?1 and name = ?2
\\ returning id
, .{ database_id, name });
if (deleted == null) {
return error.NotFound;
}
}
pub fn add(self: *Engine, object_store_id: i64, key: []const u8, value: []const u8) !void {
@@ -210,6 +229,265 @@ pub fn clear(self: *Engine, object_store_id: i64) !void {
return self.conn.exec("delete from idb_records where object_store_id = ?1", .{object_store_id});
}
pub const IndexInfo = struct {
id: i64,
key_path: []const u8,
unique: bool,
multi_entry: bool,
};
pub fn createIndexRow(self: *Engine, object_store_id: i64, name: []const u8, key_path: []const u8, unique: bool, multi_entry: bool) !i64 {
return (try self.conn.scalar(
i64,
\\ insert into idb_indexes (object_store_id, name, key_path, is_unique, multi_entry)
\\ values (?1, ?2, ?3, ?4, ?5) returning id
,
.{ object_store_id, name, key_path, unique, multi_entry },
)) orelse error.UnknownError;
}
pub fn deleteIndexRow(self: *Engine, object_store_id: i64, name: []const u8) !void {
const deleted = try self.conn.scalar(
i64,
"delete from idb_indexes where object_store_id = ?1 and name = ?2 returning id",
.{ object_store_id, name },
);
if (deleted == null) {
return error.NotFound;
}
}
pub fn indexInfo(self: *const Engine, arena: Allocator, object_store_id: i64, name: []const u8) !?IndexInfo {
var row = (try self.conn.row(
"select id, key_path, is_unique, multi_entry from idb_indexes where object_store_id = ?1 and name = ?2",
.{ object_store_id, name },
)) orelse return null;
defer row.deinit();
return .{
.id = row.get(i64, 0),
.key_path = try arena.dupe(u8, row.get([]const u8, 1)),
.unique = row.get(bool, 2),
.multi_entry = row.get(bool, 3),
};
}
pub fn indexesForStore(self: *const Engine, arena: Allocator, object_store_id: i64) ![]IndexInfo {
var rows = try self.conn.rows(
"select id, key_path, is_unique, multi_entry from idb_indexes where object_store_id = ?1",
.{object_store_id},
);
defer rows.deinit();
var list: std.ArrayList(IndexInfo) = .empty;
while (try rows.next()) |row| {
try list.append(arena, .{
.id = row.get(i64, 0),
.key_path = try arena.dupe(u8, row.get([]const u8, 1)),
.unique = row.get(bool, 2),
.multi_entry = row.get(bool, 3),
});
}
return list.items;
}
pub fn indexNames(self: *const Engine, arena: Allocator, object_store_id: i64) ![]const []const u8 {
var rows = try self.conn.rows(
"select name from idb_indexes where object_store_id = ?1 order by name",
.{object_store_id},
);
defer rows.deinit();
var list: std.ArrayList([]const u8) = .empty;
while (try rows.next()) |row| {
try list.append(arena, try arena.dupe(u8, row.get([]const u8, 0)));
}
return list.items;
}
pub fn addIndexRecord(self: *Engine, index_id: i64, key: []const u8, primary_key: []const u8, unique: bool) !void {
return self.conn.exec(
"insert into idb_index_records (index_id, key, primary_key, is_unique) values (?1, ?2, ?3, ?4)",
.{ index_id, key, primary_key, unique },
);
}
pub fn deleteIndexRecordsForKey(self: *Engine, object_store_id: i64, primary_key: []const u8) !void {
return self.conn.exec(
\\ delete from idb_index_records
\\ where primary_key = ?2 and index_id in (
\\ select id from idb_indexes where object_store_id = ?1
\\ )
, .{ object_store_id, primary_key });
}
pub fn clearIndexRecordsForStore(self: *Engine, object_store_id: i64) !void {
return self.conn.exec(
"delete from idb_index_records where index_id in (select id from idb_indexes where object_store_id = ?1)",
.{object_store_id},
);
}
// Drop index entries for every record about to be deleted by a ranged delete.
pub fn deleteIndexRecordsForRange(self: *Engine, object_store_id: i64, b: Bounds) !void {
const ops = rangeOps(b);
var buf: [400]u8 = undefined;
const sql = try std.fmt.bufPrintZ(&buf,
\\ delete from idb_index_records where index_id in (
\\ select id from idb_indexes where object_store_id = ?1
\\ ) and primary_key in (
\\ select key from idb_records where object_store_id = ?1 and key {s} ?2 and key {s} ?3
\\)
, .{ ops.lo, ops.hi });
return self.conn.exec(sql, .{ object_store_id, b.lower, b.upper });
}
pub fn savepoint(self: *Engine) !void {
return self.conn.exec("savepoint idb_op", .{});
}
pub fn releaseSavepoint(self: *Engine) !void {
return self.conn.exec("release idb_op", .{});
}
pub fn rollbackSavepoint(self: *Engine) void {
self.conn.exec("rollback to idb_op", .{}) catch |err| {
log.warn(.storage, "idb savepoint rollback", .{ .err = err, .sqlite = self.conn.lastError() });
};
self.conn.exec("release idb_op", .{}) catch |err| {
log.warn(.storage, "idb savepoint release", .{ .err = err });
};
}
// First record value in an index range (joins back to the store by primary key).
pub fn indexGetRange(self: *const Engine, arena: Allocator, object_store_id: i64, index_id: i64, b: Bounds) !?[]u8 {
const ops = rangeOps(b);
var buf: [512]u8 = undefined;
const sql = try std.fmt.bufPrint(&buf,
\\ select r.value
\\ from idb_index_records ir
\\ join idb_records r on r.object_store_id = ?1 and r.key = ir.primary_key
\\ where ir.index_id = ?2 and ir.key {s} ?3 and ir.key {s} ?4 order by ir.key, ir.primary_key
\\ limit 1
, .{ ops.lo, ops.hi });
var row = (try self.conn.row(sql, .{ object_store_id, index_id, b.lower, b.upper })) orelse return null;
defer row.deinit();
return try arena.dupe(u8, row.get([]const u8, 0));
}
// First primary key in an index range.
pub fn indexGetKeyRange(self: *const Engine, arena: Allocator, index_id: i64, b: Bounds) !?[]u8 {
const ops = rangeOps(b);
var buf: [320]u8 = undefined;
const sql = try std.fmt.bufPrint(&buf,
\\ select primary_key
\\ from idb_index_records
\\ where index_id = ?1 and key {s} ?2 and key {s} ?3
\\ order by key, primary_key limit 1
, .{ ops.lo, ops.hi });
var row = (try self.conn.row(sql, .{ index_id, b.lower, b.upper })) orelse return null;
defer row.deinit();
return try arena.dupe(u8, row.get([]const u8, 0));
}
pub fn indexCountRange(self: *const Engine, index_id: i64, b: Bounds) !i64 {
const ops = rangeOps(b);
var buf: [320]u8 = undefined;
const sql = try std.fmt.bufPrint(&buf,
\\ select count(*)
\\ from idb_index_records
\\ where index_id = ?1 and key {s} ?2 and key {s} ?3
, .{ ops.lo, ops.hi });
return (try self.conn.scalar(i64, sql, .{ index_id, b.lower, b.upper })) orelse 0;
}
// Open an index getAll/getAllKeys cursor (.value joins the store, .key is the
// primary key). The JS layer streams rows straight into a JS array.
pub fn indexGetAllRangeRows(self: *const Engine, object_store_id: i64, index_id: i64, b: Bounds, column: Column, limit_: ?u32) !Sqlite.Rows {
const ops = rangeOps(b);
const limit: i64 = if (limit_) |c| @intCast(c) else -1;
var buf: [512]u8 = undefined;
if (column == .value) {
const sql = try std.fmt.bufPrint(&buf,
\\ select r.value
\\ from idb_index_records ir
\\ join idb_records r on r.object_store_id = ?1 and r.key = ir.primary_key
\\ where ir.index_id = ?2 and ir.key {s} ?3 and ir.key {s} ?4
\\ order by ir.key, ir.primary_key
\\ limit ?5
, .{ ops.lo, ops.hi });
return self.conn.rows(sql, .{ object_store_id, index_id, b.lower, b.upper, limit });
}
const sql = try std.fmt.bufPrint(&buf,
\\ select primary_key
\\ from idb_index_records
\\ where index_id = ?1 and key {s} ?2 and key {s} ?3
\\ order by key, primary_key
\\ limit ?4
, .{ ops.lo, ops.hi });
return self.conn.rows(sql, .{ index_id, b.lower, b.upper, limit });
}
pub const IndexCursorRecord = struct {
key: []u8,
primary_key: []u8,
value: ?[]u8,
};
pub fn indexCursorSeek(
self: *const Engine,
arena: Allocator,
object_store_id: i64,
index_id: i64,
b: Bounds,
reverse: bool,
from_key: []const u8,
from_pk: []const u8,
pk_inclusive: bool,
with_value: bool,
offset: u32,
) !?IndexCursorRecord {
const ops = rangeOps(b);
const order = if (reverse) "desc" else "asc";
// Position past the current (key, primary_key): a strictly-greater key, or an
// equal key with a greater (or, for continuePrimaryKey, >=) primary key.
const key_op = if (reverse) "< " else "> ";
const pk_op = if (reverse) (if (pk_inclusive) "<= " else "< ") else (if (pk_inclusive) ">= " else "> ");
var buf: [640]u8 = undefined;
const select = if (with_value)
"select ir.key, ir.primary_key, r.value from idb_index_records ir join idb_records r on r.object_store_id = ?1 and r.key = ir.primary_key"
else
"select ir.key, ir.primary_key from idb_index_records ir";
const sql = try std.fmt.bufPrint(
&buf,
\\ {s} where ir.index_id = ?2 and ir.key {s} ?3 and ir.key {s} ?4 and (ir.key {s}?5 or (ir.key = ?5 and ir.primary_key {s}?6))
\\ order by ir.key {s}, ir.primary_key {s}
\\ limit 1 offset ?7
,
.{ select, ops.lo, ops.hi, key_op, pk_op, order, order },
);
var row = (try self.conn.row(sql, .{ object_store_id, index_id, b.lower, b.upper, from_key, from_pk, @as(i64, offset) })) orelse return null;
defer row.deinit();
return .{
.key = try arena.dupe(u8, row.get([]const u8, 0)),
.primary_key = try arena.dupe(u8, row.get([]const u8, 1)),
.value = if (with_value) try arena.dupe(u8, row.get([]const u8, 2)) else null,
};
}
pub const Bounds = struct {
lower: []const u8,
upper: []const u8,
@@ -284,20 +562,35 @@ fn rangeSql(buf: []u8, head: []const u8, b: Bounds, tail: []const u8) ![:0]u8 {
);
}
// A point range stores empty operators; index queries want a closed range.
fn rangeOps(b: Bounds) struct { lo: []const u8, hi: []const u8 } {
return .{
.lo = if (b.is_point) ">= " else b.lower_op,
.hi = if (b.is_point) "<= " else b.upper_op,
};
}
// What a ranged getAll/getAllKeys returns: the value or key column.
pub const Column = enum {
value,
key,
};
pub fn getAllRange(self: *const Engine, arena: Allocator, object_store_id: i64, b: Bounds, column: Column, limit_: ?u32) ![]const []u8 {
// Open a getAll/getAllKeys cursor. The JS layer streams rows straight into a JS
// array, avoiding a copy of the whole result set out of sqlite. (The SQL text is
// copied into the prepared statement, so the stack `buf` can be discarded.)
pub fn getAllRangeRows(self: *const Engine, object_store_id: i64, b: Bounds, column: Column, limit_: ?u32) !Sqlite.Rows {
var buf: [256]u8 = undefined;
const head = if (column == .value) "select value" else "select key";
const sql = try rangeSql(&buf, head, b, " order by key limit ?4");
// SQLite treats a negative LIMIT as "no limit".
const limit: i64 = if (limit_) |c| @intCast(c) else -1;
return self.conn.rows(sql, .{ object_store_id, b.lower, b.upper, limit });
}
var rows = try self.conn.rows(sql, .{ object_store_id, b.lower, b.upper, limit });
pub fn getAllRange(self: *const Engine, arena: Allocator, object_store_id: i64, b: Bounds, column: Column, limit_: ?u32) ![]const []u8 {
var rows = try self.getAllRangeRows(object_store_id, b, column, limit_);
defer rows.deinit();
var list: std.ArrayList([]u8) = .empty;
+125 -37
View File
@@ -23,6 +23,7 @@ const js = @import("../../../js/js.zig");
const Key = @import("Key.zig");
const Engine = @import("Engine.zig");
const IDBIndex = @import("IDBIndex.zig");
const IDBRequest = @import("IDBRequest.zig");
const IDBObjectStore = @import("IDBObjectStore.zig");
const IDBTransaction = @import("IDBTransaction.zig");
@@ -40,18 +41,24 @@ _request: *IDBRequest,
_bounds: Engine.Bounds,
_direction: Direction,
_key_only: bool,
_source: Source,
_got_value: bool = false,
// null for an object-store cursor
_index_id: ?i64 = null,
// the JS value of this Cursor, pre-converted and cached as an optimization
// since this cursor will be the request value on every iteration.
_js: js.Value.Global,
// Encoded current key; null before iteration and at the end
// Encoded current key; null before iteration and at the end. For an index cursor
// this is the index key; for an object store it equals the primary key.
_key: ?[]const u8 = null,
// Encoded primary (store) key. Equals _key for an object-store cursor.
_primary_key: ?[]const u8 = null,
// Current record's serialized value bytes (null when key-only or exhausted).
_value: ?[]const u8 = null,
// The spec's "got value flag": gated true only while a success handler can read
// the cursor, set just before each delivery and cleared by continue/advance.
_got_value: bool = false,
pub const Direction = enum {
next,
@@ -70,10 +77,38 @@ pub const Direction = enum {
}
};
// Create a cursor over `store` and run the first seek. Returns the request that
// delivers the cursor (or null) via repeated `success` events.
// What this cursor iterates — used only by the `source` accessor.
const Source = union(enum) {
store: *IDBObjectStore,
index: *IDBIndex,
};
// How the next iterate() positions the cursor.
const Seek = union(enum) {
first, // start of the range, in the iteration direction
next, // strictly past the current position
to: []const u8, // continue(key): first record at/after an index/store key
to_primary: struct { key: []const u8, primary_key: []const u8 }, // continuePrimaryKey
};
fn startSentinel(reverse: bool) []const u8 {
return if (reverse) Engine.Bounds.max_sentinel else Engine.Bounds.min_sentinel;
}
// Cursor over an object store.
pub fn open(store: *IDBObjectStore, bounds: Engine.Bounds, direction: Direction, key_only: bool, exec: *Execution) !*IDBRequest {
const txn = store._txn orelse return error.TransactionInactiveError;
return create(store, txn, null, .{ .store = store }, bounds, direction, key_only, exec);
}
// Cursor over an index (key = index key, primaryKey = store key).
pub fn openIndex(index: *IDBIndex, bounds: Engine.Bounds, direction: Direction, key_only: bool, exec: *Execution) !*IDBRequest {
const store = index._store;
const txn = store._txn orelse return error.TransactionInactiveError;
return create(store, txn, index._index_id, .{ .index = index }, bounds, direction, key_only, exec);
}
fn create(store: *IDBObjectStore, txn: *IDBTransaction, index_id: ?i64, source: Source, bounds: Engine.Bounds, direction: Direction, key_only: bool, exec: *Execution) !*IDBRequest {
try txn.ensureBegun();
const request = try txn.newRequest();
@@ -85,24 +120,23 @@ pub fn open(store: *IDBObjectStore, bounds: Engine.Bounds, direction: Direction,
._bounds = bounds,
._direction = direction,
._key_only = key_only,
._index_id = index_id,
._js = undefined,
._source = source,
});
request._cursor = self;
const local = exec.js.local.?;
const public: js.Value = if (key_only)
try local.zigValueToJs(self, .{})
else
try local.zigValueToJs(try IDBCursorWithValue.init(self, exec), .{});
// An optimization. Potentially looked up _a lot_, so calculating upfront
// storing it, and setting it in the IDBRequest (which is already js.Value.Global
// aware), avoids the lookup we'd normally have to do in the bridge.
// Pre-converted and cached because it's the request value on every iteration,
// avoiding the bridge's pointer -> js.Value lookup each time.
self._js = try public.persist();
const reverse = direction.reverse();
try self.iterate(if (reverse) "<= " else ">= ", if (reverse) Engine.Bounds.max_sentinel else Engine.Bounds.min_sentinel, 0, exec);
try self.iterate(.first, 0, exec);
return request;
}
@@ -113,31 +147,55 @@ pub fn beforeDeliver(self: *IDBCursor) void {
pub fn @"continue"(self: *IDBCursor, key_arg: ?js.Value, exec: *Execution) !void {
try self.prepareIterate();
const reverse = self._direction.reverse();
if (key_arg) |k| {
const encoded = try Key.encodeValue(exec.arena, k);
// The target must move past the current key in the iteration direction.
const order = std.mem.order(u8, encoded, self._key.?);
if (if (reverse) order != .lt else order != .gt) {
if (if (self._direction.reverse()) order != .lt else order != .gt) {
return error.DataError;
}
try self.iterate(if (reverse) "<= " else ">= ", encoded, 0, exec);
try self.iterate(.{ .to = encoded }, 0, exec);
} else {
try self.iterate(if (reverse) "< " else "> ", self._key.?, 0, exec);
try self.iterate(.next, 0, exec);
}
try self._txn.enqueue(self._request);
}
pub fn continuePrimaryKey(self: *IDBCursor, key_arg: js.Value, primary_key_arg: js.Value, exec: *Execution) !void {
// Only meaningful on an index cursor with a directed (non-unique) direction.
if (self._index_id == null or self._direction == .nextunique or self._direction == .prevunique) {
return error.InvalidAccessError;
}
try self.prepareIterate();
const reverse = self._direction.reverse();
const key = try Key.encodeValue(exec.arena, key_arg);
const primary_key = try Key.encodeValue(exec.arena, primary_key_arg);
// The (key, primaryKey) pair must move past the current position.
const ok = switch (std.mem.order(u8, key, self._key.?)) {
.gt => !reverse,
.lt => reverse,
.eq => blk: {
const pk_order = std.mem.order(u8, primary_key, self._primary_key.?);
break :blk if (reverse) pk_order == .lt else pk_order == .gt;
},
};
if (!ok) return error.DataError;
try self.iterate(.{ .to_primary = .{ .key = key, .primary_key = primary_key } }, 0, exec);
try self._txn.enqueue(self._request);
}
pub fn advance(self: *IDBCursor, count: u32, exec: *Execution) !void {
if (count == 0) {
return error.TypeError;
}
try self.prepareIterate();
const reverse = self._direction.reverse();
// `count` records forward = skip count-1 past the immediate next.
try self.iterate(if (reverse) "< " else "> ", self._key.?, count - 1, exec);
try self.iterate(.next, count - 1, exec);
try self._txn.enqueue(self._request);
}
@@ -154,9 +212,10 @@ pub fn update(self: *IDBCursor, value: js.Value, exec: *Execution) !*IDBRequest
return error.InvalidStateError;
}
const key = self._key orelse return error.InvalidStateError;
// The record sits at the primary (store) key, even for an index cursor.
const key = self._primary_key orelse return error.InvalidStateError;
// For an in-line store, the value's own key must match the cursor's key.
// For an in-line store, the value's own key must match the record's key.
if (self._store._key_path) |kp| {
const extracted = Key.evaluatePath(value, kp) orelse return error.DataError;
const encoded = try Key.encodeValue(exec.call_arena, extracted);
@@ -169,7 +228,7 @@ pub fn update(self: *IDBCursor, value: js.Value, exec: *Execution) !*IDBRequest
defer serialized.deinit();
const request = try self._txn.newRequest();
self._engine.put(self._store._store_id, key, serialized.bytes()) catch |err| {
self._store.writeAt(key, value, serialized.bytes(), exec) catch |err| {
log.warn(.storage, "idb cursor update", .{ .err = err });
request.setError(err);
return request;
@@ -191,10 +250,10 @@ pub fn delete(self: *IDBCursor) !*IDBRequest {
return error.InvalidStateError;
}
const key = self._key orelse return error.InvalidStateError;
const key = self._primary_key orelse return error.InvalidStateError;
const request = try self._txn.newRequest();
self._engine.deleteRange(self._store._store_id, Engine.Bounds.point(key)) catch |err| {
self._store.deleteAt(key) catch |err| {
log.warn(.storage, "idb cursor delete", .{ .err = err });
request.setError(err);
};
@@ -206,35 +265,63 @@ pub fn getKey(self: *const IDBCursor, exec: *Execution) !?js.Value {
return try Key.decodeToJs(exec.call_arena, exec.js.local.?, encoded);
}
// For an object store the primary key is the key.
pub fn getPrimaryKey(self: *const IDBCursor, exec: *Execution) !?js.Value {
return self.getKey(exec);
const encoded = self._primary_key orelse return null;
return try Key.decodeToJs(exec.call_arena, exec.js.local.?, encoded);
}
pub fn getDirection(self: *const IDBCursor) Direction {
return self._direction;
}
pub fn getSource(self: *const IDBCursor) *IDBObjectStore {
return self._store;
// The bridge converts the active union variant (the IDBObjectStore/IDBIndex).
pub fn getSource(self: *const IDBCursor) Source {
return self._source;
}
// Seek the next record and stage it as the request result. The key/value live on
// Seek the next record and stage it as the request result. The keys/value live on
// the page arena (they must survive across event-loop turns within the txn).
fn iterate(self: *IDBCursor, from_op: []const u8, from_key: []const u8, offset: u32, exec: *Execution) !void {
fn iterate(self: *IDBCursor, seek: Seek, offset: u32, exec: *Execution) !void {
const reverse = self._direction.reverse();
const rec = try self._engine.cursorSeek(exec.arena, self._store._store_id, self._bounds, reverse, from_op, from_key, offset, !self._key_only);
if (rec) |r| {
self._key = r.key;
self._value = r.value; // already null for a key-only cursor
self._request.setValueGlobal(self._js);
const arena = exec.arena;
const store_id = self._store._store_id;
if (self._index_id) |index_id| {
const from_key, const from_pk, const pk_inclusive = switch (seek) {
.first => .{ startSentinel(reverse), startSentinel(reverse), false },
.next => .{ self._key.?, self._primary_key.?, false },
.to => |t| .{ t, startSentinel(reverse), false },
.to_primary => |tp| .{ tp.key, tp.primary_key, true },
};
const rec = try self._engine.indexCursorSeek(arena, store_id, index_id, self._bounds, reverse, from_key, from_pk, pk_inclusive, !self._key_only, offset);
if (rec) |r| self.position(r.key, r.primary_key, r.value) else self.exhaust();
} else {
self._key = null;
self._value = null;
self._request.setNull();
const from_op, const from_key = switch (seek) {
.first => .{ if (reverse) "<= " else ">= ", startSentinel(reverse) },
.next => .{ if (reverse) "< " else "> ", self._key.? },
.to => |t| .{ if (reverse) "<= " else ">= ", t },
.to_primary => return error.InvalidAccessError,
};
const rec = try self._engine.cursorSeek(arena, store_id, self._bounds, reverse, from_op, from_key, offset, !self._key_only);
// For an object store the key is the primary key.
if (rec) |r| self.position(r.key, r.key, r.value) else self.exhaust();
}
}
fn position(self: *IDBCursor, key: []const u8, primary_key: []const u8, value: ?[]const u8) void {
self._key = key;
self._primary_key = primary_key;
self._value = value;
self._request.setValueGlobal(self._js);
}
fn exhaust(self: *IDBCursor) void {
self._key = null;
self._primary_key = null;
self._value = null;
self._request.setNull();
}
// validate the state before we can advance/continue
fn prepareIterate(self: *IDBCursor) !void {
if (self._txn._settled == true) {
@@ -262,6 +349,7 @@ pub const JsApi = struct {
pub const direction = bridge.accessor(IDBCursor.getDirection, null, .{});
pub const source = bridge.accessor(IDBCursor.getSource, null, .{});
pub const @"continue" = bridge.function(IDBCursor.@"continue", .{ });
pub const continuePrimaryKey = bridge.function(IDBCursor.continuePrimaryKey, .{ });
pub const advance = bridge.function(IDBCursor.advance, .{ });
pub const update = bridge.function(IDBCursor.update, .{ });
pub const delete = bridge.function(IDBCursor.delete, .{ });
+203
View File
@@ -0,0 +1,203 @@
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
//
// Francis Bouvier <francis@lightpanda.io>
// Pierre Tachoire <pierre@lightpanda.io>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const lp = @import("lightpanda");
const js = @import("../../../js/js.zig");
const Key = @import("Key.zig");
const Engine = @import("Engine.zig");
const IDBCursor = @import("IDBCursor.zig");
const IDBRequest = @import("IDBRequest.zig");
const IDBKeyRange = @import("IDBKeyRange.zig");
const IDBTransaction = @import("IDBTransaction.zig");
const IDBObjectStore = @import("IDBObjectStore.zig");
const log = lp.log;
const Execution = js.Execution;
const IDBIndex = @This();
_store: *IDBObjectStore,
_engine: *Engine,
_index_id: i64,
_name: []const u8,
_key_path: []const u8,
_unique: bool,
_multi_entry: bool,
pub fn init(obj_store: *IDBObjectStore, info: Engine.IndexInfo, name: []const u8, exec: *Execution) !*IDBIndex {
return exec._factory.create(IDBIndex{
._store = obj_store,
._engine = obj_store._engine,
._index_id = info.id,
._name = name,
._key_path = info.key_path,
._unique = info.unique,
._multi_entry = info.multi_entry,
});
}
fn txn(self: *IDBIndex) !*IDBTransaction {
const t = self._store._txn orelse return error.TransactionInactiveError;
try t.ensureBegun();
return t;
}
pub fn get(self: *IDBIndex, query: js.Value, exec: *Execution) !*IDBRequest {
const t = try self.txn();
const arena = exec.call_arena;
const bounds = try IDBKeyRange.resolveQuery(arena, query, exec);
const request = try t.newRequest();
const bytes = self._engine.indexGetRange(arena, self._store._store_id, self._index_id, bounds) catch |err| {
log.warn(.storage, "idb index get", .{ .err = err });
request.setError(err);
return request;
};
const b = bytes orelse return request;
try request.setValue(try js.Value.deserialize(exec.js.local.?, b));
return request;
}
pub fn getKey(self: *IDBIndex, query: js.Value, exec: *Execution) !*IDBRequest {
const t = try self.txn();
const arena = exec.call_arena;
const bounds = try IDBKeyRange.resolveQuery(arena, query, exec);
const request = try t.newRequest();
const bytes = self._engine.indexGetKeyRange(arena, self._index_id, bounds) catch |err| {
log.warn(.storage, "idb index getKey", .{ .err = err });
request.setError(err);
return request;
};
const b = bytes orelse return request;
try request.setValue(try Key.decodeToJs(arena, exec.js.local.?, b));
return request;
}
pub fn getAll(self: *IDBIndex, query: ?js.Value, count_: ?u32, exec: *Execution) !*IDBRequest {
return self._getAll(query, count_, .value, exec);
}
pub fn getAllKeys(self: *IDBIndex, query: ?js.Value, count_: ?u32, exec: *Execution) !*IDBRequest {
return self._getAll(query, count_, .key, exec);
}
fn _getAll(self: *IDBIndex, query: ?js.Value, count_: ?u32, column: Engine.Column, exec: *Execution) !*IDBRequest {
const t = try self.txn();
const bounds = try IDBKeyRange.resolveQuery(exec.call_arena, query, exec);
const request = try t.newRequest();
const arr = self.collectAll(bounds, count_, column, exec) catch |err| {
log.warn(.storage, "idb index getAll", .{ .err = err });
request.setError(err);
return request;
};
try request.setValue(arr);
return request;
}
// Stream an index getAll/getAllKeys straight into a JS array: .value rows are
// the joined records (deserialized), .key rows are primary keys (decoded).
fn collectAll(self: *IDBIndex, bounds: Engine.Bounds, count_: ?u32, column: Engine.Column, exec: *Execution) !js.Value {
const local = exec.js.local.?;
const arena = exec.call_arena;
var rows = try self._engine.indexGetAllRangeRows(self._store._store_id, self._index_id, bounds, column, count_);
defer rows.deinit();
const arr = local.newArray(0);
var i: u32 = 0;
while (try rows.next()) |row| {
const bytes = row.get([]const u8, 0);
const value = if (column == .value) try js.Value.deserialize(local, bytes) else try Key.decodeToJs(arena, local, bytes);
_ = try arr.set(i, value, .{});
i += 1;
}
return arr.toValue();
}
pub fn count(self: *IDBIndex, query: ?js.Value, exec: *Execution) !*IDBRequest {
const t = try self.txn();
const bounds = try IDBKeyRange.resolveQuery(exec.call_arena, query, exec);
const request = try t.newRequest();
const n = self._engine.indexCountRange(self._index_id, bounds) catch |err| {
log.warn(.storage, "idb index count", .{ .err = err });
request.setError(err);
return request;
};
try request.setValue(try exec.js.local.?.zigValueToJs(n, .{}));
return request;
}
pub fn openCursor(self: *IDBIndex, query: ?js.Value, direction: ?IDBCursor.Direction, exec: *Execution) !*IDBRequest {
const bounds = try IDBKeyRange.resolveQuery(exec.arena, query, exec);
return IDBCursor.openIndex(self, bounds, direction orelse .next, false, exec);
}
pub fn openKeyCursor(self: *IDBIndex, query: ?js.Value, direction: ?IDBCursor.Direction, exec: *Execution) !*IDBRequest {
const bounds = try IDBKeyRange.resolveQuery(exec.arena, query, exec);
return IDBCursor.openIndex(self, bounds, direction orelse .next, true, exec);
}
pub fn getName(self: *const IDBIndex) []const u8 {
return self._name;
}
pub fn getKeyPath(self: *const IDBIndex) []const u8 {
return self._key_path;
}
pub fn getUnique(self: *const IDBIndex) bool {
return self._unique;
}
pub fn getMultiEntry(self: *const IDBIndex) bool {
return self._multi_entry;
}
pub fn getObjectStore(self: *IDBIndex) *IDBObjectStore {
return self._store;
}
pub const JsApi = struct {
pub const bridge = js.Bridge(IDBIndex);
pub const Meta = struct {
pub const name = "IDBIndex";
pub const prototype_chain = bridge.prototypeChain();
pub var class_id: bridge.ClassId = undefined;
};
pub const name = bridge.accessor(IDBIndex.getName, null, .{});
pub const keyPath = bridge.accessor(IDBIndex.getKeyPath, null, .{});
pub const unique = bridge.accessor(IDBIndex.getUnique, null, .{});
pub const multiEntry = bridge.accessor(IDBIndex.getMultiEntry, null, .{});
pub const objectStore = bridge.accessor(IDBIndex.getObjectStore, null, .{});
pub const get = bridge.function(IDBIndex.get, .{ .dom_exception = true });
pub const getKey = bridge.function(IDBIndex.getKey, .{ .dom_exception = true });
pub const getAll = bridge.function(IDBIndex.getAll, .{ .dom_exception = true });
pub const getAllKeys = bridge.function(IDBIndex.getAllKeys, .{ .dom_exception = true });
pub const count = bridge.function(IDBIndex.count, .{ .dom_exception = true });
pub const openCursor = bridge.function(IDBIndex.openCursor, .{ .dom_exception = true });
pub const openKeyCursor = bridge.function(IDBIndex.openKeyCursor, .{ .dom_exception = true });
};
+191 -38
View File
@@ -16,12 +16,14 @@
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
const std = @import("std");
const lp = @import("lightpanda");
const js = @import("../../../js/js.zig");
const Key = @import("Key.zig");
const Engine = @import("Engine.zig");
const IDBIndex = @import("IDBIndex.zig");
const IDBCursor = @import("IDBCursor.zig");
const IDBRequest = @import("IDBRequest.zig");
const IDBKeyRange = @import("IDBKeyRange.zig");
@@ -29,6 +31,7 @@ const IDBTransaction = @import("IDBTransaction.zig");
const log = lp.log;
const Execution = js.Execution;
const Allocator = std.mem.Allocator;
const IDBObjectStore = @This();
@@ -100,13 +103,18 @@ pub fn delete(self: *IDBObjectStore, query: js.Value, exec: *Execution) !*IDBReq
const bounds = try IDBKeyRange.resolveQuery(exec.call_arena, query, exec);
const request = try txn.newRequest();
self._engine.deleteRange(self._store_id, bounds) catch |err| {
self.deleteBounds(bounds) catch |err| {
log.warn(.storage, "idb delete", .{ .err = err });
request.setError(err);
};
return request;
}
fn deleteBounds(self: *IDBObjectStore, bounds: Engine.Bounds) !void {
try self._engine.deleteIndexRecordsForRange(self._store_id, bounds);
try self._engine.deleteRange(self._store_id, bounds);
}
pub fn clear(self: *IDBObjectStore, _: *Execution) !*IDBRequest {
const txn = self._txn orelse return error.TransactionInactiveError;
if (txn._mode == .readonly) {
@@ -115,13 +123,18 @@ pub fn clear(self: *IDBObjectStore, _: *Execution) !*IDBRequest {
try txn.ensureBegun();
const request = try txn.newRequest();
self._engine.clear(self._store_id) catch |err| {
self.clearAll() catch |err| {
log.warn(.storage, "idb clear", .{ .err = err });
request.setError(err);
};
return request;
}
fn clearAll(self: *IDBObjectStore) !void {
try self._engine.clearIndexRecordsForStore(self._store_id);
try self._engine.clear(self._store_id);
}
pub fn count(self: *IDBObjectStore, query: ?js.Value, exec: *Execution) !*IDBRequest {
const txn = self._txn orelse return error.TransactionInactiveError;
try txn.ensureBegun();
@@ -138,29 +151,45 @@ pub fn count(self: *IDBObjectStore, query: ?js.Value, exec: *Execution) !*IDBReq
}
pub fn getAll(self: *IDBObjectStore, query: ?js.Value, count_: ?u32, exec: *Execution) !*IDBRequest {
return self._getAll(query, count_, .value, exec);
}
fn _getAll(self: *IDBObjectStore, query: ?js.Value, count_: ?u32, column: Engine.Column, exec: *Execution) !*IDBRequest {
const txn = self._txn orelse return error.TransactionInactiveError;
try txn.ensureBegun();
const local = exec.js.local.?;
const arena = exec.call_arena;
const bounds = try IDBKeyRange.resolveQuery(arena, query, exec);
const bounds = try IDBKeyRange.resolveQuery(exec.call_arena, query, exec);
const request = try txn.newRequest();
const values = self._engine.getAllRange(arena, self._store_id, bounds, .value, count_) catch |err| {
const arr = self.collectAll(exec, bounds, column, count_) catch |err| {
log.warn(.storage, "idb getAll", .{ .err = err });
request.setError(err);
return request;
};
const arr = local.newArray(@intCast(values.len));
for (values, 0..) |bytes, i| {
const value = try js.Value.deserialize(local, bytes);
_ = try arr.set(@intCast(i), value, .{});
}
try request.setValue(arr.toValue());
try request.setValue(arr);
return request;
}
// Stream a getAll/getAllKeys result straight into a JS array: .value rows
// deserialize, .key rows decode — nothing is copied out of sqlite first.
fn collectAll(self: *IDBObjectStore, exec: *Execution, bounds: Engine.Bounds, column: Engine.Column, count_: ?u32) !js.Value {
const local = exec.js.local.?;
const arena = exec.call_arena;
var rows = try self._engine.getAllRangeRows(self._store_id, bounds, column, count_);
defer rows.deinit();
const arr = local.newArray(0);
var i: u32 = 0;
while (try rows.next()) |row| {
const bytes = row.get([]const u8, 0);
const value = if (column == .value) try js.Value.deserialize(local, bytes) else try Key.decodeToJs(arena, local, bytes);
_ = try arr.set(i, value, .{});
i += 1;
}
return arr.toValue();
}
pub fn getKey(self: *IDBObjectStore, query: js.Value, exec: *Execution) !*IDBRequest {
const txn = self._txn orelse return error.TransactionInactiveError;
try txn.ensureBegun();
@@ -181,26 +210,7 @@ pub fn getKey(self: *IDBObjectStore, query: js.Value, exec: *Execution) !*IDBReq
}
pub fn getAllKeys(self: *IDBObjectStore, query: ?js.Value, count_: ?u32, exec: *Execution) !*IDBRequest {
const txn = self._txn orelse return error.TransactionInactiveError;
try txn.ensureBegun();
const arena = exec.call_arena;
const bounds = try IDBKeyRange.resolveQuery(arena, query, exec);
const request = try txn.newRequest();
const keys = self._engine.getAllRange(arena, self._store_id, bounds, .key, count_) catch |err| {
log.warn(.storage, "idb getAllKeys", .{ .err = err });
request.setError(err);
return request;
};
const local = exec.js.local.?;
const arr = local.newArray(@intCast(keys.len));
for (keys, 0..) |bytes, i| {
_ = try arr.set(@intCast(i), try Key.decodeToJs(arena, local, bytes), .{});
}
try request.setValue(arr.toValue());
return request;
return self._getAll(query, count_, .key, exec);
}
pub fn openCursor(self: *IDBObjectStore, query: ?js.Value, direction: ?IDBCursor.Direction, exec: *Execution) !*IDBRequest {
@@ -296,20 +306,159 @@ fn write(self: *IDBObjectStore, value: js.Value, key_arg: ?js.Value, kind: Write
const request = try txn.newRequest();
const result = switch (kind) {
.add => self._engine.add(self._store_id, encoded, serialized.bytes()),
.put => self._engine.put(self._store_id, encoded, serialized.bytes()),
};
result catch |err| {
// Record + index rows are atomic: a unique-index violation rolls the record
// write back too.
try self._engine.savepoint();
self.writeRecord(kind, encoded, serialized.bytes(), value, exec) catch |err| {
self._engine.rollbackSavepoint();
log.warn(.storage, "idb write", .{ .err = err, .kind = kind });
request.setError(err);
return request;
};
try self._engine.releaseSavepoint();
try request.setValue(key_value);
return request;
}
fn writeRecord(self: *IDBObjectStore, kind: WriteKind, key: []const u8, bytes: []const u8, value: js.Value, exec: *Execution) !void {
switch (kind) {
.add => try self._engine.add(self._store_id, key, bytes),
.put => try self._engine.put(self._store_id, key, bytes),
}
try self.reindex(value, key, exec);
}
// Used by IDBCursor.update: overwrite the record at `key` and re-index, atomically.
pub fn writeAt(self: *IDBObjectStore, key: []const u8, value: js.Value, bytes: []const u8, exec: *Execution) !void {
try self._engine.savepoint();
errdefer self._engine.rollbackSavepoint();
try self._engine.put(self._store_id, key, bytes);
try self.reindex(value, key, exec);
try self._engine.releaseSavepoint();
}
// Used by IDBCursor.delete: drop the record at `key` and its index entries.
pub fn deleteAt(self: *IDBObjectStore, key: []const u8) !void {
try self._engine.deleteIndexRecordsForKey(self._store_id, key);
try self._engine.deleteRange(self._store_id, Engine.Bounds.point(key));
}
// Drop a record's old index entries and add fresh ones from `value`.
fn reindex(self: *IDBObjectStore, value: js.Value, primary_key: []const u8, exec: *Execution) !void {
const arena = exec.call_arena;
const indexes = try self._engine.indexesForStore(arena, self._store_id);
if (indexes.len == 0) {
return;
}
var seen: std.ArrayList([]const u8) = .empty;
try self._engine.deleteIndexRecordsForKey(self._store_id, primary_key);
for (indexes) |idx| {
try self.addIndexEntries(arena, &seen, idx.id, idx.unique, idx.multi_entry, idx.key_path, value, primary_key);
}
}
fn addIndexEntries(self: *IDBObjectStore, arena: Allocator, seen: *std.ArrayList([]const u8), index_id: i64, unique: bool, multi_entry: bool, key_path: []const u8, value: js.Value, primary_key: []const u8) !void {
const extracted = Key.evaluatePath(value, key_path) orelse return;
if (multi_entry and extracted.isArray()) {
seen.clearRetainingCapacity();
// every value of an array is added, but only once (e.g. deduplicated)
const arr = extracted.toArray();
for (0..arr.len()) |i| {
const element = try arr.get(@intCast(i));
const encoded = Key.encodeValue(arena, element) catch continue; // skip invalid elements
for (seen.items) |s| {
if (std.mem.eql(u8, s, encoded)) {
continue;
}
}
try seen.append(arena, encoded);
try self._engine.addIndexRecord(index_id, encoded, primary_key, unique);
}
} else {
const encoded = Key.encodeValue(arena, extracted) catch return; // not a valid key -> not indexed
try self._engine.addIndexRecord(index_id, encoded, primary_key, unique);
}
}
const CreateIndexOptions = struct {
unique: bool = false,
multiEntry: bool = false,
};
// Only callable during an upgrade (versionchange transaction).
pub fn createIndex(self: *IDBObjectStore, name: []const u8, key_path: []const u8, options: ?CreateIndexOptions, exec: *Execution) !*IDBIndex {
const txn = self._txn orelse return error.InvalidStateError;
if (txn._mode != .versionchange) {
return error.InvalidStateError;
}
const opts = options orelse CreateIndexOptions{};
try self._engine.savepoint();
errdefer self._engine.rollbackSavepoint();
const index_id = self._engine.createIndexRow(self._store_id, name, key_path, opts.unique, opts.multiEntry) catch |err| switch (err) {
error.Constraint => return error.ConstraintError, // duplicate index name
else => return err,
};
const arena = exec.call_arena;
const local = exec.js.local.?;
{
// we reach directly in to _engine.conn here to avoid copying the values out
// sqlite
var rows = try self._engine.conn.rows("select key, value from idb_records where object_store_id = ?1", .{self._store_id});
defer rows.deinit();
var seen: std.ArrayList([]const u8) = .empty;
while (try rows.next()) |row| {
const value = try js.Value.deserialize(local, row.get([]const u8, 1));
try self.addIndexEntries(arena, &seen, index_id, opts.unique, opts.multiEntry, key_path, value, row.get([]const u8, 0));
}
}
const owned_name = try exec.dupeString(name);
const owned_key_path = try exec.dupeString(key_path);
const idb_index = try IDBIndex.init(self, .{
.id = index_id,
.key_path = owned_key_path,
.unique = opts.unique,
.multi_entry = opts.multiEntry,
}, owned_name, exec);
try self._engine.releaseSavepoint();
return idb_index;
}
// Only callable during an upgrade (versionchange transaction).
pub fn deleteIndex(self: *IDBObjectStore, name: []const u8, _: *Execution) !void {
const txn = self._txn orelse return error.InvalidStateError;
if (txn._mode != .versionchange) {
return error.InvalidStateError;
}
self._engine.deleteIndexRow(self._store_id, name) catch |err| switch (err) {
error.NotFound => return error.NotFoundError,
else => return err,
};
}
pub fn index(self: *IDBObjectStore, name: []const u8, exec: *Execution) !*IDBIndex {
if (self._txn == null) {
return error.InvalidStateError;
}
const info = (try self._engine.indexInfo(exec.arena, self._store_id, name)) orelse return error.NotFound;
const owned_name = try exec.dupeString(name);
return IDBIndex.init(self, info, owned_name, exec);
}
pub fn getIndexNames(self: *IDBObjectStore, exec: *Execution) ![]const []const u8 {
return self._engine.indexNames(exec.arena, self._store_id);
}
pub const JsApi = struct {
pub const bridge = js.Bridge(IDBObjectStore);
@@ -334,4 +483,8 @@ pub const JsApi = struct {
pub const getAllKeys = bridge.function(IDBObjectStore.getAllKeys, .{ });
pub const openCursor = bridge.function(IDBObjectStore.openCursor, .{ });
pub const openKeyCursor = bridge.function(IDBObjectStore.openKeyCursor, .{ });
pub const indexNames = bridge.accessor(IDBObjectStore.getIndexNames, null, .{});
pub const createIndex = bridge.function(IDBObjectStore.createIndex, .{ });
pub const deleteIndex = bridge.function(IDBObjectStore.deleteIndex, .{ });
pub const index = bridge.function(IDBObjectStore.index, .{ });
};
@@ -37,7 +37,7 @@ const FunctionSetter = idb.FunctionSetter;
const IDBRequest = @This();
_proto: *EventTarget,
_result: Result = .{.none = js.Undefined{}},
_result: Result = .{ .none = js.Undefined{} },
_error: ?anyerror = null,
_txn: ?*IDBTransaction = null,
_cursor: ?*IDBCursor = null,
@@ -57,7 +57,7 @@ const ReadyState = enum {
};
const Result = union(enum) {
none: ?js.Undefined, // null or undefined (different APIs return different values)
none: ?js.Undefined, // null or undefined (different APIs return different values)
value: js.Value.Global, // the result of a get/add/put, or a positioned cursor
database: *IDBDatabase, // the result of an open
};
@@ -82,7 +82,7 @@ pub fn setValueGlobal(self: *IDBRequest, global: js.Value.Global) void {
// Not exposed to JS, called internally. Result becomes JS `null` (not undefined).
pub fn setNull(self: *IDBRequest) void {
self._result = .{.none = null};
self._result = .{ .none = null };
}
// Not exposed to JS, called internally
+2
View File
@@ -25,6 +25,7 @@ pub const Manager = @import("Manager.zig");
pub const IDBFactory = @import("IDBFactory.zig");
pub const IDBRequest = @import("IDBRequest.zig");
pub const IDBCursor = @import("IDBCursor.zig");
pub const IDBIndex = @import("IDBIndex.zig");
pub const IDBDatabase = @import("IDBDatabase.zig");
pub const IDBKeyRange = @import("IDBKeyRange.zig");
pub const IDBTransaction = @import("IDBTransaction.zig");
@@ -37,6 +38,7 @@ pub fn registerTypes() []const type {
IDBFactory,
IDBRequest,
IDBCursor,
IDBIndex,
IDBDatabase,
IDBKeyRange,
IDBTransaction,
+1 -1
View File
@@ -306,7 +306,7 @@ const Row = struct {
}
};
const Rows = struct {
pub const Rows = struct {
stmt: Statement,
pub fn deinit(self: Rows) void {