diff --git a/Dockerfile b/Dockerfile index 9750682..3881e5f 100644 --- a/Dockerfile +++ b/Dockerfile @@ -3,6 +3,5 @@ RUN apk add --no-cache ca-certificates COPY zig-out/bin/zlay /usr/local/bin/zlay RUN mkdir -p /data/events ENV RELAY_DATA_DIR=/data/events -ENV RELAY_DB_PATH=/data/relay.sqlite EXPOSE 3000 3001 ENTRYPOINT ["/usr/local/bin/zlay"] diff --git a/build.zig b/build.zig index b802445..f0a7b73 100644 --- a/build.zig +++ b/build.zig @@ -12,7 +12,7 @@ pub fn build(b: *std.Build) void { .target = target, .optimize = optimize, }); - const zqlite = b.dependency("zqlite", .{ + const pg = b.dependency("pg", .{ .target = target, .optimize = optimize, }); @@ -20,7 +20,7 @@ pub fn build(b: *std.Build) void { const imports: []const std.Build.Module.Import = &.{ .{ .name = "zat", .module = zat.module("zat") }, .{ .name = "websocket", .module = websocket.module("websocket") }, - .{ .name = "zqlite", .module = zqlite.module("zqlite") }, + .{ .name = "pg", .module = pg.module("pg") }, }; // relay executable diff --git a/build.zig.zon b/build.zig.zon index 9f1d452..7abc3d2 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -12,9 +12,9 @@ .url = "https://github.com/karlseguin/websocket.zig/archive/97fefafa59cc78ce177cff540b8685cd7f699276.tar.gz", .hash = "websocket-0.1.0-ZPISdRlzAwBB_Bz2UMMqxYqF6YEVTIBoFsbzwPUJTHIc", }, - .zqlite = .{ - .url = "https://github.com/karlseguin/zqlite.zig/archive/05a88d6758753e1c63fdd45b211dde2057094b0c.tar.gz", - .hash = "zqlite-0.0.1-RWLaYz6bmAAT7E_jxopXf-j5Ea8VQldnxsd6TU8sa0Bb", + .pg = .{ + .url = "git+https://github.com/karlseguin/pg.zig?ref=master#e58b318b7867ef065b3135983f829219c5eef891", + .hash = "pg-0.0.0-Wp_7gXFoBgD0fQ72WICKa-bxLga03AXXQ3BbIsjjohQ3", }, }, .paths = .{ diff --git a/src/event_log.zig b/src/event_log.zig index 50a46ea..8a96748 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -1,7 +1,7 @@ //! disk persistence matching indigo's diskpersist format //! //! append-only log files with relay-assigned sequence numbers. -//! SQLite metadata index for fast cursor→file lookup. +//! Postgres metadata index for fast cursor→file lookup. //! //! on-disk entry format (28-byte LE header + CBOR payload): //! [4B flags LE] [4B kind LE] [4B payload_len LE] [8B uid LE] [8B seq LE] [payload] @@ -11,7 +11,7 @@ //! see: indigo cmd/relay/stream/persist/diskpersist/diskpersist.go const std = @import("std"); -const zqlite = @import("zqlite"); +const pg = @import("pg"); const Allocator = std.mem.Allocator; const log = std.log.scoped(.relay); @@ -38,17 +38,6 @@ const default_events_per_file: u32 = 10_000; const default_flush_interval_ms: u64 = 100; const default_flush_threshold: usize = 400; -/// convert u64 seq to i64 for SQLite storage. -/// XOR with sign bit to preserve ordering across the full u64 range. -fn seqToSqlite(seq: u64) i64 { - return @bitCast(seq ^ (1 << 63)); -} - -/// convert i64 from SQLite back to u64 seq. -fn sqliteToSeq(val: i64) u64 { - return @as(u64, @bitCast(val)) ^ (1 << 63); -} - // --- header --- pub const EvtHeader = struct { @@ -90,7 +79,7 @@ pub const DiskPersist = struct { allocator: Allocator, dir_path: []const u8, dir: std.fs.Dir, - db: zqlite.Conn, + db: *pg.Pool, current_file: ?std.fs.File = null, current_file_path: ?[]const u8 = null, @@ -116,7 +105,7 @@ pub const DiskPersist = struct { alive: std.atomic.Value(bool) = .{ .raw = true }, flush_cond: std.Thread.Condition = .{}, - pub fn init(allocator: Allocator, dir_path: []const u8, db_path: []const u8) !DiskPersist { + pub fn init(allocator: Allocator, dir_path: []const u8, database_url: []const u8) !DiskPersist { // ensure directory exists std.fs.cwd().makePath(dir_path) catch |err| switch (err) { error.PathAlreadyExists => {}, @@ -126,51 +115,52 @@ pub const DiskPersist = struct { var dir = try std.fs.cwd().openDir(dir_path, .{ .iterate = true }); errdefer dir.close(); - // ensure db parent directory exists - if (std.fs.path.dirname(db_path)) |parent| { - std.fs.cwd().makePath(parent) catch |err| switch (err) { - error.PathAlreadyExists => {}, - else => return err, - }; - } + // connect to Postgres + const uri = std.Uri.parse(database_url) catch return error.InvalidDatabaseUrl; + const pool = try pg.Pool.initUri(allocator, uri, .{ .size = 5 }); + errdefer pool.deinit(); - // open SQLite - const db_path_z = try allocator.dupeZ(u8, db_path); - defer allocator.free(db_path_z); - var db = try zqlite.open(db_path_z, zqlite.OpenFlags.Create | zqlite.OpenFlags.ReadWrite); - errdefer db.close(); - - // pragmas - try db.execNoArgs("PRAGMA journal_mode=WAL"); - try db.execNoArgs("PRAGMA busy_timeout=5000"); - try db.execNoArgs("PRAGMA synchronous=NORMAL"); - - // create tables - try db.execNoArgs( + // create tables (matching indigo's Go relay schema) + _ = try pool.exec( \\CREATE TABLE IF NOT EXISTS log_file_refs ( - \\ id INTEGER PRIMARY KEY AUTOINCREMENT, + \\ id BIGSERIAL PRIMARY KEY, \\ path TEXT NOT NULL, - \\ archived INTEGER NOT NULL DEFAULT 0, - \\ seq_start INTEGER NOT NULL, - \\ created_at TEXT NOT NULL DEFAULT (datetime('now')) + \\ archived BOOLEAN NOT NULL DEFAULT false, + \\ seq_start BIGINT NOT NULL, + \\ created_at TIMESTAMPTZ NOT NULL DEFAULT now() \\) - ); - try db.execNoArgs( + , .{}); + + _ = try pool.exec( \\CREATE TABLE IF NOT EXISTS account ( - \\ uid INTEGER PRIMARY KEY AUTOINCREMENT, + \\ uid BIGSERIAL PRIMARY KEY, \\ did TEXT NOT NULL UNIQUE, \\ status TEXT NOT NULL DEFAULT 'active', - \\ rev TEXT, - \\ commit_data_cid BLOB, - \\ created_at TEXT NOT NULL DEFAULT (datetime('now')) + \\ created_at TIMESTAMPTZ NOT NULL DEFAULT now() \\) - ); + , .{}); + + _ = try pool.exec( + \\CREATE TABLE IF NOT EXISTS account_repo ( + \\ uid BIGINT PRIMARY KEY REFERENCES account(uid), + \\ rev TEXT NOT NULL, + \\ commit_data_cid TEXT NOT NULL + \\) + , .{}); + + _ = try pool.exec( + \\CREATE TABLE IF NOT EXISTS domain_ban ( + \\ id BIGSERIAL PRIMARY KEY, + \\ domain TEXT NOT NULL UNIQUE, + \\ created_at TIMESTAMPTZ NOT NULL DEFAULT now() + \\) + , .{}); var self = DiskPersist{ .allocator = allocator, .dir_path = try allocator.dupe(u8, dir_path), .dir = dir, - .db = db, + .db = pool, }; // recover from existing log files @@ -205,7 +195,7 @@ pub const DiskPersist = struct { if (self.current_file) |f| f.close(); if (self.current_file_path) |p| self.allocator.free(p); self.dir.close(); - self.db.close(); + self.db.deinit(); self.allocator.free(self.dir_path); } @@ -225,27 +215,26 @@ pub const DiskPersist = struct { } // check database - if (self.db.row( - "SELECT uid FROM account WHERE did = ?", + if (try self.db.rowUnsafe( + "SELECT uid FROM account WHERE did = $1", .{did}, - )) |maybe_row| { - if (maybe_row) |r| { - defer r.deinit(); - const uid: u64 = @intCast(r.int(0)); - // populate cache - const did_duped = try self.allocator.dupe(u8, did); - self.did_cache_mutex.lock(); - defer self.did_cache_mutex.unlock(); - self.did_cache.put(self.allocator, did_duped, uid) catch { - self.allocator.free(did_duped); - }; - return uid; - } - } else |_| {} + )) |row| { + var r = row; + defer r.deinit() catch {}; + const uid: u64 = @intCast(r.get(i64, 0)); + // populate cache + const did_duped = try self.allocator.dupe(u8, did); + self.did_cache_mutex.lock(); + defer self.did_cache_mutex.unlock(); + self.did_cache.put(self.allocator, did_duped, uid) catch { + self.allocator.free(did_duped); + }; + return uid; + } // create new account row (ignore if already exists from concurrent insert) - self.db.exec( - "INSERT OR IGNORE INTO account (did) VALUES (?)", + _ = self.db.exec( + "INSERT INTO account (did) VALUES ($1) ON CONFLICT (did) DO NOTHING", .{did}, ) catch |err| { log.warn("failed to create account for {s}: {s}", .{ did, @errorName(err) }); @@ -253,12 +242,12 @@ pub const DiskPersist = struct { }; // read back the UID (whether we just created it or it already existed) - const row = try self.db.row( - "SELECT uid FROM account WHERE did = ?", + var row = try self.db.rowUnsafe( + "SELECT uid FROM account WHERE did = $1", .{did}, ) orelse return error.AccountCreationFailed; - defer row.deinit(); - const uid: u64 = @intCast(row.int(0)); + defer row.deinit() catch {}; + const uid: u64 = @intCast(row.get(i64, 0)); // populate cache const did_duped = try self.allocator.dupe(u8, did); @@ -277,15 +266,15 @@ pub const DiskPersist = struct { data_cid: []const u8, }; - /// get stored sync state for a user + /// get stored sync state for a user (from account_repo table) pub fn getAccountState(self: *DiskPersist, uid: u64, allocator: Allocator) !?AccountState { - const row = (try self.db.row( - "SELECT rev, commit_data_cid FROM account WHERE uid = ? AND rev IS NOT NULL", + var row = (try self.db.rowUnsafe( + "SELECT rev, commit_data_cid FROM account_repo WHERE uid = $1", .{@as(i64, @intCast(uid))}, )) orelse return null; - defer row.deinit(); - const rev = row.text(0); - const data_cid = row.blob(1); + defer row.deinit() catch {}; + const rev = row.get([]const u8, 0); + const data_cid = row.get([]const u8, 1); if (rev.len == 0 or data_cid.len == 0) return null; return .{ .rev = try allocator.dupe(u8, rev), @@ -293,11 +282,11 @@ pub const DiskPersist = struct { }; } - /// update stored sync state after a verified commit + /// update stored sync state after a verified commit (upsert into account_repo) pub fn updateAccountState(self: *DiskPersist, uid: u64, rev: []const u8, data_cid: []const u8) !void { - try self.db.exec( - "UPDATE account SET rev = ?, commit_data_cid = ? WHERE uid = ?", - .{ rev, data_cid, @as(i64, @intCast(uid)) }, + _ = try self.db.exec( + "INSERT INTO account_repo (uid, rev, commit_data_cid) VALUES ($1, $2, $3) ON CONFLICT (uid) DO UPDATE SET rev = EXCLUDED.rev, commit_data_cid = EXCLUDED.commit_data_cid", + .{ @as(i64, @intCast(uid)), rev, data_cid }, ); } @@ -341,35 +330,42 @@ pub const DiskPersist = struct { } /// playback events with seq > since. calls cb for each event. - pub fn playback(self: *DiskPersist, since: u64, allocator: Allocator, result: *std.ArrayListUnmanaged(PlaybackEntry)) !void { + pub fn playback(self: *DiskPersist, since: u64, allocator: Allocator, entries: *std.ArrayListUnmanaged(PlaybackEntry)) !void { self.mutex.lock(); defer self.mutex.unlock(); + const since_i: i64 = @intCast(since); + // find the log file containing `since` var start_files: std.ArrayListUnmanaged(LogFileRef) = .{}; defer start_files.deinit(allocator); if (since > 0) { // find file whose seq_start is just before `since` - if (self.db.row("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start <= ? ORDER BY seq_start DESC LIMIT 1", .{seqToSqlite(since)})) |row| { - if (row) |r| { - defer r.deinit(); - try start_files.append(allocator, .{ - .path = try allocator.dupe(u8, r.text(1)), - .seq_start = sqliteToSeq(r.int(2)), - }); - } - } else |_| {} + if (try self.db.rowUnsafe( + "SELECT id, path, seq_start FROM log_file_refs WHERE seq_start <= $1 ORDER BY seq_start DESC LIMIT 1", + .{since_i}, + )) |row| { + var r = row; + defer r.deinit() catch {}; + try start_files.append(allocator, .{ + .path = try allocator.dupe(u8, r.get([]const u8, 1)), + .seq_start = @intCast(r.get(i64, 2)), + }); + } } // find all subsequent files { - var rows = try self.db.rows("SELECT id, path, seq_start FROM log_file_refs WHERE seq_start > ? ORDER BY seq_start ASC", .{seqToSqlite(since)}); - defer rows.deinit(); - while (rows.next()) |r| { + var result = try self.db.query( + "SELECT id, path, seq_start FROM log_file_refs WHERE seq_start > $1 ORDER BY seq_start ASC", + .{since_i}, + ); + defer result.deinit(); + while (result.nextUnsafe() catch null) |r| { try start_files.append(allocator, .{ - .path = try allocator.dupe(u8, r.text(1)), - .seq_start = sqliteToSeq(r.int(2)), + .path = try allocator.dupe(u8, r.get([]const u8, 1)), + .seq_start = @intCast(r.get(i64, 2)), }); } } @@ -380,7 +376,7 @@ pub const DiskPersist = struct { for (start_files.items) |ref| { var file = self.dir.openFile(ref.path, .{}) catch continue; defer file.close(); - try readEventsFrom(allocator, file, since, result); + try readEventsFrom(allocator, file, since, entries); } } @@ -395,9 +391,8 @@ pub const DiskPersist = struct { self.mutex.lock(); defer self.mutex.unlock(); - const cutoff_hours = self.retention_hours; - const cutoff_sql = try std.fmt.allocPrint(self.allocator, "-{d} hours", .{cutoff_hours}); - defer self.allocator.free(cutoff_sql); + const cutoff_interval = try std.fmt.allocPrint(self.allocator, "{d} hours", .{self.retention_hours}); + defer self.allocator.free(cutoff_interval); // find expired refs var expired: std.ArrayListUnmanaged(GcRef) = .{}; @@ -407,15 +402,15 @@ pub const DiskPersist = struct { } { - var rows = try self.db.rows( - "SELECT id, path FROM log_file_refs WHERE created_at < datetime('now', ?)", - .{cutoff_sql}, + var result = try self.db.query( + "SELECT id, path FROM log_file_refs WHERE created_at < now() - $1::interval", + .{cutoff_interval}, ); - defer rows.deinit(); - while (rows.next()) |r| { + defer result.deinit(); + while (result.nextUnsafe() catch null) |r| { try expired.append(self.allocator, .{ - .id = r.int(0), - .path = try self.allocator.dupe(u8, r.text(1)), + .id = r.get(i64, 0), + .path = try self.allocator.dupe(u8, r.get([]const u8, 1)), }); } } @@ -427,7 +422,7 @@ pub const DiskPersist = struct { } // delete db record first (prevents playback from finding it) - self.db.exec("DELETE FROM log_file_refs WHERE id = ?", .{ref.id}) catch |err| { + _ = self.db.exec("DELETE FROM log_file_refs WHERE id = $1", .{ref.id}) catch |err| { log.warn("gc: failed to delete db record {d}: {s}", .{ ref.id, @errorName(err) }); continue; }; @@ -456,10 +451,10 @@ pub const DiskPersist = struct { } { - var rows = try self.db.rows("SELECT path FROM log_file_refs ORDER BY seq_start DESC", .{}); - defer rows.deinit(); - while (rows.next()) |r| { - try refs.append(self.allocator, try self.allocator.dupe(u8, r.text(0))); + var result = try self.db.query("SELECT path FROM log_file_refs ORDER BY seq_start DESC", .{}); + defer result.deinit(); + while (result.nextUnsafe() catch null) |r| { + try refs.append(self.allocator, try self.allocator.dupe(u8, r.get([]const u8, 0))); } } @@ -476,11 +471,14 @@ pub const DiskPersist = struct { fn resumeLog(self: *DiskPersist) !void { // find most recent log file - const r = try self.db.row("SELECT id, path, seq_start FROM log_file_refs ORDER BY seq_start DESC LIMIT 1", .{}); - if (r) |row| { - defer row.deinit(); - const path = row.text(1); - const seq_start: u64 = sqliteToSeq(row.int(2)); + if (try self.db.rowUnsafe( + "SELECT id, path, seq_start FROM log_file_refs ORDER BY seq_start DESC LIMIT 1", + .{}, + )) |row| { + var r = row; + defer r.deinit() catch {}; + const path = r.get([]const u8, 1); + const seq_start: u64 = @intCast(r.get(i64, 2)); var file = self.dir.openFile(path, .{ .mode = .read_write }) catch { // file missing, start fresh @@ -527,10 +525,10 @@ pub const DiskPersist = struct { self.current_file = try self.dir.createFile(name, .{ .truncate = false }); self.current_file_path = try self.allocator.dupe(u8, name); - // register in SQLite - try self.db.exec( - "INSERT INTO log_file_refs (path, seq_start) VALUES (?, ?)", - .{ name, seqToSqlite(start_seq) }, + // register in Postgres + _ = try self.db.exec( + "INSERT INTO log_file_refs (path, seq_start) VALUES ($1, $2)", + .{ name, @as(i64, @intCast(start_seq)) }, ); self.event_counter = 0; @@ -753,17 +751,20 @@ test "header is little-endian" { try std.testing.expectEqual(@as(u8, 0x01), buf[9]); } +fn requireDatabaseUrl() ![]const u8 { + return std.posix.getenv("DATABASE_URL") orelse return error.SkipZigTest; +} + test "persist and playback" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); // persist some events (sync flush, no background thread) @@ -798,16 +799,15 @@ test "persist and playback" { } test "playback with cursor" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); _ = try dp.persist(.commit, 1, "a"); @@ -833,18 +833,17 @@ test "playback with cursor" { } test "seq recovery after reinit" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - // write some events { - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); _ = try dp.persist(.commit, 1, "x"); _ = try dp.persist(.commit, 2, "y"); @@ -856,7 +855,7 @@ test "seq recovery after reinit" { // reinit — should recover seq { - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); try std.testing.expectEqual(@as(u64, 3), dp.lastSeq().?); const seq4 = try dp.persist(.commit, 1, "w"); @@ -865,16 +864,15 @@ test "seq recovery after reinit" { } test "takedown zeros payload" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); _ = try dp.persist(.commit, 42, "secret-data"); @@ -901,16 +899,15 @@ test "takedown zeros payload" { } test "uidForDid assigns and caches UIDs" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); // first call creates the account @@ -928,25 +925,24 @@ test "uidForDid assigns and caches UIDs" { } test "uidForDid survives reinit" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - var uid1: u64 = undefined; { - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); uid1 = try dp.uidForDid("did:plc:carol"); } // reinit — UID should be the same from database { - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); const uid1_again = try dp.uidForDid("did:plc:carol"); try std.testing.expectEqual(uid1, uid1_again); @@ -954,16 +950,15 @@ test "uidForDid survives reinit" { } test "takedown with real UIDs" { + const database_url = try requireDatabaseUrl(); + var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); defer std.testing.allocator.free(dir_path); - const db_path = try std.fmt.allocPrint(std.testing.allocator, "{s}/relay.sqlite", .{dir_path}); - defer std.testing.allocator.free(db_path); - - var dp = try DiskPersist.init(std.testing.allocator, dir_path, db_path); + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url); defer dp.deinit(); const alice_uid = try dp.uidForDid("did:plc:alice"); diff --git a/src/main.zig b/src/main.zig index 9a195d9..c5fc007 100644 --- a/src/main.zig +++ b/src/main.zig @@ -62,9 +62,9 @@ pub fn main() !void { defer val.deinit(); try val.start(); - // init disk persistence (indigo-compatible diskpersist format + SQLite index) - const db_path = std.posix.getenv("RELAY_DB_PATH") orelse "data/relay.sqlite"; - var dp = event_log_mod.DiskPersist.init(allocator, data_dir, db_path) catch |err| { + // init disk persistence (indigo-compatible diskpersist format + Postgres index) + const database_url = std.posix.getenv("DATABASE_URL") orelse "postgres://relay:relay@localhost:5432/relay"; + var dp = event_log_mod.DiskPersist.init(allocator, data_dir, database_url) catch |err| { log.err("failed to init disk persist at {s}: {s}", .{ data_dir, @errorName(err) }); return err; };