Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,8 @@ db.prepare("SELECT * FROM streams WHERE sender = @sender").run({ sender: "G..."

## Pragmas and Performance

SQLite WAL mode is already enabled in `db.ts` (line 24). If issues #360 mentions additional pragmas, they should be added to the DB init:
SQLite WAL mode and related pragmas are applied in `backend/src/services/sqlite/` (see `docs/adr/0006-sqlite-wal-and-pool-tuning.md`):
- `PRAGMA journal_mode=WAL`: Concurrent reads during writes
- `PRAGMA synchronous=NORMAL`: Balance durability and speed
- `PRAGMA busy_timeout=5000`: Prevent SQLITE_BUSY on concurrent writes
- `PRAGMA cache_size=-64000`: 64MB page cache for read perf
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -424,6 +424,7 @@ Copy `backend/.env.example` to `backend/.env` and `frontend/.env.example` to `fr
| `NETWORK_PASSPHRASE` | Optional (Default: `Test SDF Network ; September 2015`) | Non-empty string passphrase matching Stellar network | `Test SDF Network ; September 2015` | Network identifier passphrase |
| `ALLOWED_ASSETS` | Optional (Default: `USDC,XLM`) | Comma-separated list of 1–12 alphanumeric asset codes | `USDC,XLM,EURC` | Supported token assets for payment streams |
| `DB_PATH` | Optional (Default: `backend/data/streams.db`) | Valid filesystem file path string | `backend/data/streams.db` | SQLite database file location |
| `SQLITE_READ_POOL_SIZE` | Optional (Default: `1`, max `4`) | Integer | `1` | Extra readonly SQLite connections for WAL reads (`0` uses the writer only) |
| `WEBHOOK_DESTINATION_URL` | Optional | Valid HTTP/HTTPS URL (`z.string().url()`) | `https://example.com/webhooks/stellar` | Destination URL for outbound stream event webhooks |
| `WEBHOOK_SIGNING_SECRET` | Optional (Recommended if URL set) | Secret string ($\ge 32$ random characters recommended) | `whsec_9a8b7c6d5e4f3a2b1c0d9e8f7a6b5c4d` | Secret used to sign webhook payloads via HMAC-SHA256 |
| `JWT_SECRET` | Optional (Required in Production) | Secret string ($\ge 32$ random characters recommended) | `jwt_sec_8f7e6d5c4b3a2f1e0d9c8b7a6f5e4d3c` | Secret used for signing JWT authentication tokens |
Expand Down
1 change: 1 addition & 0 deletions backend/eslint.config.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ export default tseslint.config(
'src/services/auth.ts',
'src/services/cache.ts',
'src/services/db.ts',
'src/services/sqlite/connection-pool.ts',
'src/services/eventHistory.ts',
'src/services/indexer.ts',
'src/services/metricsHistory.ts',
Expand Down
23 changes: 9 additions & 14 deletions backend/src/db/concurrency.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import path from "path";
import { Worker } from "worker_threads";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { runMigrations } from "../services/migrations";
import { applySqlitePragmas } from "../services/sqlite/apply-pragmas";

function createTempDbPath(): string {
return path.join(
Expand Down Expand Up @@ -36,22 +37,19 @@ describe("SQLite WAL mode and concurrent read/write safety", () => {

it("verifies PRAGMA journal_mode = wal on connection", () => {
const db = new Database(dbPath);
db.pragma("journal_mode = WAL");
const result = db.pragma("journal_mode");
const applied = applySqlitePragmas(db);
db.close();
expect(result).toEqual([{ journal_mode: "wal" }]);
expect(applied.journalMode).toBe("wal");
expect(applied.busyTimeoutMs).toBe(5000);
});

it("does not throw SQLITE_BUSY during concurrent read and write", () => {
const writerDb = new Database(dbPath);
writerDb.pragma("journal_mode = WAL");
writerDb.pragma("busy_timeout = 5000");
writerDb.pragma("synchronous = NORMAL");
applySqlitePragmas(writerDb);
runMigrations(writerDb);

const readerDb = new Database(dbPath);
readerDb.pragma("journal_mode = WAL");
readerDb.pragma("busy_timeout = 5000");
applySqlitePragmas(readerDb);

try {
writerDb
Expand Down Expand Up @@ -106,13 +104,11 @@ describe("SQLite WAL mode and concurrent read/write safety", () => {

it("returns consistent data when reader reads during active write transaction", () => {
const writerDb = new Database(dbPath);
writerDb.pragma("journal_mode = WAL");
writerDb.pragma("busy_timeout = 5000");
applySqlitePragmas(writerDb);
runMigrations(writerDb);

const readerDb = new Database(dbPath);
readerDb.pragma("journal_mode = WAL");
readerDb.pragma("busy_timeout = 5000");
applySqlitePragmas(readerDb);

try {
writerDb
Expand Down Expand Up @@ -154,8 +150,7 @@ describe("SQLite WAL mode and concurrent read/write safety", () => {

it("handles true concurrent read/write via worker threads without SQLITE_BUSY", async () => {
const setupDb = new Database(dbPath);
setupDb.pragma("journal_mode = WAL");
setupDb.pragma("busy_timeout = 5000");
applySqlitePragmas(setupDb);
runMigrations(setupDb);

setupDb
Expand Down
76 changes: 54 additions & 22 deletions backend/src/services/db.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
import Database from "better-sqlite3";
import path from "path";
import { runMigrations } from "./migrations";
import { SQLITE_DEFAULT_READ_POOL_SIZE } from "./sqlite/constants";
import { SqliteConnectionPool } from "./sqlite/connection-pool";

const DB_PATH =
process.env.DB_PATH || path.join(__dirname, "..", "..", "data", "streams.db");
function resolveDbPath(): string {
return process.env.DB_PATH || path.join(__dirname, "..", "..", "data", "streams.db");
}

let db: any;
let sqlitePool: SqliteConnectionPool | null = null;

export function getDb(): any {
if (!db) {
Expand All @@ -14,6 +17,26 @@ export function getDb(): any {
return db;
}

export function getReadDb(): any {
if (sqlitePool) {
return sqlitePool.getReader();
}
return getDb();
}

export function closeDb(): void {
if (sqlitePool) {
sqlitePool.close();
sqlitePool = null;
db = undefined;
return;
}
if (db) {
db.close();
db = undefined;
}
}

export function isPostgres(): boolean {
return !!process.env.DATABASE_URL;
}
Expand Down Expand Up @@ -259,39 +282,37 @@ class PostgresDatabase {
}

public prepare(sql: string): any {
const dbInstance = this;
return {
run(...params: any[]): any {
run: (...params: any[]): any => {
const paramObj = params.length === 1 && typeof params[0] === "object" && params[0] !== null && !Array.isArray(params[0]) ? params[0] : params;
const res = dbInstance.querySync(sql, paramObj);
const res = this.querySync(sql, paramObj);
return {
changes: res.rowCount,
lastInsertRowid: 0,
};
},
get(...params: any[]): any {
get: (...params: any[]): any => {
const paramObj = params.length === 1 && typeof params[0] === "object" && params[0] !== null && !Array.isArray(params[0]) ? params[0] : params;
const res = dbInstance.querySync(sql, paramObj);
const res = this.querySync(sql, paramObj);
return res.rows[0] || undefined;
},
all(...params: any[]): any {
all: (...params: any[]): any => {
const paramObj = params.length === 1 && typeof params[0] === "object" && params[0] !== null && !Array.isArray(params[0]) ? params[0] : params;
const res = dbInstance.querySync(sql, paramObj);
const res = this.querySync(sql, paramObj);
return res.rows;
},
};
}

public transaction(fn: Function): any {
const dbInstance = this;
public transaction(fn: (...args: any[]) => any): any {
return (...args: any[]) => {
dbInstance.exec("BEGIN");
this.exec("BEGIN");
try {
const result = fn(...args);
dbInstance.exec("COMMIT");
this.exec("COMMIT");
return result;
} catch (error) {
dbInstance.exec("ROLLBACK");
this.exec("ROLLBACK");
throw error;
}
};
Expand All @@ -308,22 +329,33 @@ class PostgresDatabase {
}
}

function resolveReadPoolSize(): number {
const raw = process.env.SQLITE_READ_POOL_SIZE;
if (!raw) {
return SQLITE_DEFAULT_READ_POOL_SIZE;
}
const parsed = Number.parseInt(raw, 10);
return Number.isFinite(parsed) ? parsed : SQLITE_DEFAULT_READ_POOL_SIZE;
}

export function initDb(): void {
closeDb();

if (isPostgres()) {
db = new PostgresDatabase(process.env.DATABASE_URL!);
} else {
const dir = path.dirname(DB_PATH);
const dbPath = resolveDbPath();
const dir = path.dirname(dbPath);
const fs = require("fs");
if (!fs.existsSync(dir)) {
fs.mkdirSync(dir, { recursive: true });
}

db = new Database(DB_PATH);
db.pragma("journal_mode = WAL");
db.pragma("foreign_keys = ON");
db.pragma("synchronous = NORMAL");
db.pragma("busy_timeout = 5000");
db.pragma("cache_size = -64000");
sqlitePool = new SqliteConnectionPool({
filePath: dbPath,
readPoolSize: resolveReadPoolSize(),
});
db = sqlitePool.getWriter();
}

runMigrations(db);
Expand Down
78 changes: 78 additions & 0 deletions backend/src/services/sqlite/apply-pragmas.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
import Database from "better-sqlite3";
import fs from "fs";
import os from "os";
import path from "path";
import { afterEach, describe, expect, it } from "vitest";
import {
applySqlitePragmas,
DEFAULT_SQLITE_PRAGMA_SETTINGS,
readSqlitePragmas,
} from "./apply-pragmas";
import {
SQLITE_BUSY_TIMEOUT_MS,
SQLITE_CACHE_SIZE_KIB,
SqliteConnectionRole,
} from "./constants";

function tempDbPath(): string {
return path.join(
os.tmpdir(),
`stellar-stream-pragmas-${Date.now()}-${Math.random().toString(36).slice(2)}.db`,
);
}

describe("applySqlitePragmas", () => {
const paths: string[] = [];

afterEach(() => {
for (const filePath of paths) {
for (const suffix of ["", "-wal", "-shm"]) {
try {
fs.unlinkSync(filePath + suffix);
} catch {}
}
}
paths.length = 0;
});

it("enables WAL, busy timeout, and tuned cache on a writer connection", () => {
const filePath = tempDbPath();
paths.push(filePath);
const db = new Database(filePath);
const applied = applySqlitePragmas(db);

expect(applied.journalMode).toBe("wal");
expect(applied.busyTimeoutMs).toBe(SQLITE_BUSY_TIMEOUT_MS);
expect(applied.cacheSize).toBe(SQLITE_CACHE_SIZE_KIB);

db.close();
});

it("applies busy timeout and cache size on a reader without requiring journal writes", () => {
const filePath = tempDbPath();
paths.push(filePath);
const writer = new Database(filePath);
applySqlitePragmas(writer, DEFAULT_SQLITE_PRAGMA_SETTINGS, SqliteConnectionRole.Writer);
const reader = new Database(filePath, { readonly: true });
const applied = applySqlitePragmas(
reader,
DEFAULT_SQLITE_PRAGMA_SETTINGS,
SqliteConnectionRole.Reader,
);

expect(applied.journalMode).toBe("wal");
expect(applied.busyTimeoutMs).toBe(SQLITE_BUSY_TIMEOUT_MS);
expect(applied.cacheSize).toBe(SQLITE_CACHE_SIZE_KIB);
reader.close();
writer.close();
});

it("reads back the same pragma snapshot after apply", () => {
const filePath = tempDbPath();
paths.push(filePath);
const db = new Database(filePath);
const applied = applySqlitePragmas(db);
expect(readSqlitePragmas(db)).toEqual(applied);
db.close();
});
});
95 changes: 95 additions & 0 deletions backend/src/services/sqlite/apply-pragmas.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import {
SQLITE_BUSY_TIMEOUT_MS,
SQLITE_CACHE_SIZE_KIB,
SQLITE_MMAP_SIZE_BYTES,
SQLITE_WAL_AUTOCHECKPOINT_PAGES,
SqliteConnectionRole,
SqliteForeignKeys,
SqliteJournalMode,
SqlitePragma,
SqliteSynchronous,
SqliteTempStore,
} from "./constants";

export interface SqliteQueryable {
pragma(source: string): unknown;
}

export interface SqlitePragmaSettings {
journalMode: SqliteJournalMode;
synchronous: SqliteSynchronous;
busyTimeoutMs: number;
cacheSizeKib: number;
foreignKeys: SqliteForeignKeys;
tempStore: SqliteTempStore;
walAutocheckpointPages: number;
mmapSizeBytes: number;
}

export const DEFAULT_SQLITE_PRAGMA_SETTINGS: SqlitePragmaSettings = {
journalMode: SqliteJournalMode.Wal,
synchronous: SqliteSynchronous.Normal,
busyTimeoutMs: SQLITE_BUSY_TIMEOUT_MS,
cacheSizeKib: SQLITE_CACHE_SIZE_KIB,
foreignKeys: SqliteForeignKeys.On,
tempStore: SqliteTempStore.Memory,
walAutocheckpointPages: SQLITE_WAL_AUTOCHECKPOINT_PAGES,
mmapSizeBytes: SQLITE_MMAP_SIZE_BYTES,
};

export interface AppliedSqlitePragmas {
journalMode: string;
synchronous: number | string;
busyTimeoutMs: number;
cacheSize: number;
foreignKeys: number;
tempStore: number | string;
walAutocheckpoint: number;
mmapSize: number;
}

function pragmaScalar(db: SqliteQueryable, name: SqlitePragma): string | number {
const result = db.pragma(name);
if (Array.isArray(result) && result.length > 0 && typeof result[0] === "object" && result[0] !== null) {
const values = Object.values(result[0] as Record<string, unknown>);
const value = values[0];
if (typeof value === "number" || typeof value === "string") {
return value;
}
}
if (typeof result === "number" || typeof result === "string") {
return result;
}
throw new Error(`Unexpected PRAGMA ${name} result`);
}

export function applySqlitePragmas(
db: SqliteQueryable,
settings: SqlitePragmaSettings = DEFAULT_SQLITE_PRAGMA_SETTINGS,
role: SqliteConnectionRole = SqliteConnectionRole.Writer,
): AppliedSqlitePragmas {
if (role === SqliteConnectionRole.Writer) {
db.pragma(`${SqlitePragma.JournalMode} = ${settings.journalMode}`);
db.pragma(`${SqlitePragma.Synchronous} = ${settings.synchronous}`);
db.pragma(`${SqlitePragma.ForeignKeys} = ${settings.foreignKeys}`);
db.pragma(`${SqlitePragma.TempStore} = ${settings.tempStore}`);
db.pragma(`${SqlitePragma.WalAutocheckpoint} = ${settings.walAutocheckpointPages}`);
db.pragma(`${SqlitePragma.MmapSize} = ${settings.mmapSizeBytes}`);
}
db.pragma(`${SqlitePragma.BusyTimeout} = ${settings.busyTimeoutMs}`);
db.pragma(`${SqlitePragma.CacheSize} = ${settings.cacheSizeKib}`);
return readSqlitePragmas(db);
}

export function readSqlitePragmas(db: SqliteQueryable): AppliedSqlitePragmas {
return {
journalMode: String(pragmaScalar(db, SqlitePragma.JournalMode)).toLowerCase(),
synchronous: pragmaScalar(db, SqlitePragma.Synchronous),
busyTimeoutMs: Number(pragmaScalar(db, SqlitePragma.BusyTimeout)),
cacheSize: Number(pragmaScalar(db, SqlitePragma.CacheSize)),
foreignKeys: Number(pragmaScalar(db, SqlitePragma.ForeignKeys)),
tempStore: pragmaScalar(db, SqlitePragma.TempStore),
walAutocheckpoint: Number(pragmaScalar(db, SqlitePragma.WalAutocheckpoint)),
mmapSize: Number(pragmaScalar(db, SqlitePragma.MmapSize)),
};
}
Loading