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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
80 changes: 80 additions & 0 deletions packages/shared/src/nodeSqliteClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,8 +61,88 @@ layer("NodeSqliteClient", (it) => {
assert.equal(error.reason.operation, "prepare");
}),
);

it.effect("classifies constraint failures by their SQLite result code", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* sql`CREATE TABLE constrained(name TEXT NOT NULL UNIQUE)`;
yield* sql`INSERT INTO constrained VALUES ('taken')`;

const duplicate = yield* sql`INSERT INTO constrained VALUES ('taken')`.pipe(Effect.flip);
assert(duplicate.reason._tag === "UniqueViolation");
assert.equal(duplicate.reason.constraint, "constrained.name");

const missing = yield* sql`INSERT INTO constrained VALUES (NULL)`.pipe(Effect.flip);
assert.equal(missing.reason._tag, "ConstraintError");
}),
);
});

const makeTempDatabase = Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const directory = yield* fs.makeTempDirectoryScoped({ prefix: "t3-sqlite-transaction-" });
const filename = path.join(directory, "state.sqlite");
// node:sqlite connections fail a busy statement at once unless given a timeout.
const other = yield* Effect.acquireRelease(
Effect.sync(() => new NodeSqlite.DatabaseSync(filename)),
(database) => Effect.sync(() => database.close()),
);
yield* Effect.sync(() => {
other.exec(`
PRAGMA journal_mode = WAL;
CREATE TABLE counters(id INTEGER PRIMARY KEY, value INTEGER NOT NULL);
INSERT INTO counters VALUES (1, 0);
`);
});
return { filename, other };
});

it.effect("keeps another connection from committing between a transaction's read and write", () =>
Effect.gen(function* () {
const { filename, other } = yield* makeTempDatabase;
yield* Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* sql.withTransaction(
Effect.gen(function* () {
const [row] = yield* sql<{ readonly value: number }>`
SELECT value FROM counters WHERE id = 1
`;
// With a deferred BEGIN this commit lands and the write below fails
// with SQLITE_BUSY_SNAPSHOT, which no busy timeout can wait out.
yield* Effect.sync(() =>
assert.throws(
() => other.exec("UPDATE counters SET value = value + 1 WHERE id = 1"),
/database is locked/,
),
);
yield* sql`UPDATE counters SET value = ${(row?.value ?? 0) + 10} WHERE id = 1`;
}),
);
yield* Effect.sync(() => other.exec("UPDATE counters SET value = value + 1 WHERE id = 1"));
assert.deepEqual(yield* sql`SELECT value FROM counters WHERE id = 1`.values, [[11]]);
}).pipe(Effect.provide(SqliteClient.layer({ filename })));
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("reports a transaction blocked by another writer as a lock timeout", () =>
Effect.gen(function* () {
const { filename, other } = yield* makeTempDatabase;
yield* Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* sql`PRAGMA busy_timeout = 0`;
const read = sql.withTransaction(sql`SELECT value FROM counters WHERE id = 1`.values);

yield* Effect.sync(() => other.exec("BEGIN IMMEDIATE"));
const error = yield* read.pipe(Effect.flip);
assert.equal(error.reason._tag, "LockTimeoutError");

yield* Effect.sync(() => other.exec("COMMIT"));
assert.deepEqual(yield* read, [[0]]);
}).pipe(Effect.provide(SqliteClient.layer({ filename })));
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("returns a typed failure when the database cannot be opened", () =>
Effect.gen(function* () {
const error = yield* Effect.flip(
Expand Down
58 changes: 34 additions & 24 deletions packages/shared/src/nodeSqliteClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
import { identity } from "effect/Function";
import * as Layer from "effect/Layer";
import * as Predicate from "effect/Predicate";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
Expand Down Expand Up @@ -82,6 +83,22 @@ const checkNodeSqliteCompat = () => {
return Effect.void;
};

/**
* `node:sqlite` reports the SQLite result code as `errcode`, while
* `classifySqliteError` reads `errno`. Copy it across so busy, locked and
* constraint failures get their own reasons instead of `UnknownError`.
*/
const classifyError = (cause: unknown, message: string, operation: string) => {
if (
Predicate.hasProperty(cause, "errcode") &&
typeof cause.errcode === "number" &&
!Predicate.hasProperty(cause, "errno")
) {
Object.assign(cause, { errno: cause.errcode });
}
return classifySqliteError(cause, { message, operation });
};

const make = Effect.fn("makeWithDatabase")(function* (
options: SqliteClientConfig,
): Effect.fn.Return<Client.SqlClient, SqlError, Scope.Scope | Reactivity.Reactivity> {
Expand All @@ -102,10 +119,7 @@ const make = Effect.fn("makeWithDatabase")(function* (
}),
catch: (cause) =>
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to open database",
operation: "open",
}),
reason: classifyError(cause, "Failed to open database", "open"),
}),
});
yield* Scope.addFinalizer(
Expand All @@ -114,10 +128,7 @@ const make = Effect.fn("makeWithDatabase")(function* (
try: () => db.close(),
catch: (cause) =>
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to close database",
operation: "close",
}),
reason: classifyError(cause, "Failed to close database", "close"),
}),
}).pipe(Effect.orDie),
);
Expand All @@ -138,10 +149,7 @@ const make = Effect.fn("makeWithDatabase")(function* (
try: () => db.prepare(sql),
catch: (cause) =>
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to prepare statement",
operation: "prepare",
}),
reason: classifyError(cause, "Failed to prepare statement", "prepare"),
}),
});

Expand All @@ -168,10 +176,7 @@ const make = Effect.fn("makeWithDatabase")(function* (
} catch (cause) {
return Effect.fail(
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to execute statement",
operation: "execute",
}),
reason: classifyError(cause, "Failed to execute statement", "execute"),
}),
);
}
Expand Down Expand Up @@ -201,10 +206,7 @@ const make = Effect.fn("makeWithDatabase")(function* (
},
catch: (cause) =>
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to execute statement",
operation: "execute",
}),
reason: classifyError(cause, "Failed to execute statement", "execute"),
}),
}),
(statement) =>
Expand All @@ -216,10 +218,11 @@ const make = Effect.fn("makeWithDatabase")(function* (
},
catch: (cause) =>
new SqlError({
reason: classifySqliteError(cause, {
message: "Failed to reset statement result mode",
operation: "resetResultMode",
}),
reason: classifyError(
cause,
"Failed to reset statement result mode",
"resetResultMode",
),
}),
}).pipe(Effect.orDie),
);
Expand Down Expand Up @@ -273,6 +276,13 @@ const make = Effect.fn("makeWithDatabase")(function* (
acquirer,
compiler,
transactionAcquirer,
// A deferred BEGIN only takes the write lock at the first write. If another
// process commits after this transaction's first read, that write fails at
// once with SQLITE_BUSY_SNAPSHOT, which busy_timeout cannot wait out. Taking
// the lock up front makes it wait instead, at the cost of serializing
// read-only transactions behind other processes' writers. Read-only
// connections cannot write, so they keep the deferred BEGIN.
beginTransaction: options.readonly === true ? "BEGIN" : "BEGIN IMMEDIATE",
Comment thread
coderabbitai[bot] marked this conversation as resolved.
spanAttributes: [
...(options.spanAttributes ? Object.entries(options.spanAttributes) : []),
[ATTR_DB_SYSTEM_NAME, "sqlite"],
Expand Down
Loading