From d858b63391b87337c94158394723883cc772050d Mon Sep 17 00:00:00 2001 From: Lucas Carlson Date: Mon, 21 Sep 2026 12:23:20 -0700 Subject: [PATCH] fix: lock the actor instance by primary key Concurrent callers that created the same actor from inside their own transaction deadlocked on MySQL. The portable conflict clause becomes INSERT IGNORE, which leaves a shared lock on the identity index when the row already exists. The mailbox then read the same row FOR UPDATE, so two callers each held shared and each waited for exclusive on one index record. Six of eight concurrent callers failed with ER_LOCK_DEADLOCK. Top level enqueue hid this, because it retries ER_LOCK_DEADLOCK eight times. Nested callers get no retry, and the effect executor, the reminder scheduler, the retry path, and effect recovery all enqueue inside a transaction that already wrote. Read the instance without a lock first, on every database, then lock it by primary key. After an ignored insert on MySQL, read the winning row in shared mode and select only its id, so the read stays inside the identity index and never locks the clustered record. A shared read that selects every column locks the clustered record too, which only moves the same upgrade deadlock onto the primary key. PostgreSQL and SQLite take no shared lock. A share lock on PostgreSQL creates the upgrade deadlock it prevents on MySQL, and its regression test fails when that lock is applied there. Validate with the default suite, the MySQL suite, the PostgreSQL suite, format:check, check, build, test:package, and test:recovery. --- CHANGELOG.md | 11 +++++++++++ docs/correctness.md | 3 +++ package.json | 2 +- src/repository.ts | 34 ++++++++++++++++++++++---------- src/version.ts | 2 +- test/mysql.test.ts | 43 +++++++++++++++++++++++++++++++++++++++++ test/postgresql.test.ts | 43 +++++++++++++++++++++++++++++++++++++++++ 7 files changed, 126 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d04a688..237f7d5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,16 @@ # Changelog +## 0.15.2 - 2026-09-21 + +- Stop the deadlock between concurrent callers that create the same actor from + a transaction of their own. MySQL turns a portable conflict clause into + `INSERT IGNORE`, which keeps a shared lock on the identity index, and the + mailbox then asked to upgrade that lock. It now reads the winning row in + shared mode, selects only its id so the read stays inside the index, and + locks the row by its primary key. +- Read the instance without a lock before the insert, on every database, and + take the row lock by primary key rather than by actor type and actor id. + ## 0.15.1 - 2026-09-16 - Index opt-in instance expiration by actor type and update time so pruning diff --git a/docs/correctness.md b/docs/correctness.md index dd686d7..17858a3 100644 --- a/docs/correctness.md +++ b/docs/correctness.md @@ -5,6 +5,9 @@ - Delivery is ordered per actor identity and at least once. - Different identities may execute concurrently. - Sequence allocation and durable enqueue are one transaction. +- Concurrent callers that create the same actor produce one instance row and + distinct sequences. The mailbox locks that row by its primary key, so MySQL + does not upgrade a shared lock and the enqueue does not deadlock. - A retryable failure rolls state and staged intents back and blocks later work. - A stale activation may finish JavaScript but cannot commit. - Effects can execute more than once and must deduplicate by their stable id. diff --git a/package.json b/package.json index cb45151..0527e08 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "solid-objects", - "version": "0.15.1", + "version": "0.15.2", "description": "Race-free realtime state per application identity, backed by your SQL database", "type": "module", "license": "MIT", diff --git a/src/repository.ts b/src/repository.ts index 93b9e24..9d0611b 100644 --- a/src/repository.ts +++ b/src/repository.ts @@ -298,13 +298,12 @@ export class Repository { input: EnqueueInput, ): Promise { const now = await connection.nowMilliseconds() - let instance = await this.findInstance({ + let identified = await this.findInstanceId({ connection, actorType: input.actorType, actorId: input.actorId, - ...(this.settings.database.family === "postgresql" ? { lock: "update" } : {}), }) - if (!instance) { + if (!identified) { const instanceId = randomUUID() await connection.run( `INSERT INTO ${this.table("instances")} @@ -321,15 +320,18 @@ export class Repository { now, ], ) - } - if (this.settings.database.family === "mysql" || !instance) { - instance = await this.findInstance({ + identified = await this.findInstanceId({ connection, actorType: input.actorType, actorId: input.actorId, - ...(this.settings.database.family === "sqlite" ? {} : { lock: "update" }), + ...(this.settings.database.family === "mysql" ? { lock: "share" as const } : {}), }) } + if (!identified) throw new ActorDestroyed("actor disappeared during enqueue") + const instance = await connection.get( + `SELECT * FROM ${this.table("instances")} WHERE id = ?${this.rowLockClause()}`, + [identified.id], + ) if (!instance) throw new ActorDestroyed("actor disappeared during enqueue") if (input.idempotencyKey !== undefined) { @@ -2071,12 +2073,24 @@ export class Repository { connection: DatabaseConnection actorType: string actorId: string - lock?: "update" }): Promise { const { connection, actorType, actorId } = options return connection.get( - `SELECT * FROM ${this.table("instances")} WHERE actor_type = ? AND actor_id = ?${ - options.lock === "update" ? " FOR UPDATE" : "" + `SELECT * FROM ${this.table("instances")} WHERE actor_type = ? AND actor_id = ?`, + [actorType, actorId], + ) + } + + private findInstanceId(options: { + connection: DatabaseConnection + actorType: string + actorId: string + lock?: "share" + }): Promise<{ id: string } | undefined> { + const { connection, actorType, actorId } = options + return connection.get<{ id: string }>( + `SELECT id FROM ${this.table("instances")} WHERE actor_type = ? AND actor_id = ?${ + options.lock === "share" ? " FOR SHARE" : "" }`, [actorType, actorId], ) diff --git a/src/version.ts b/src/version.ts index f8ec2c0..e25e508 100644 --- a/src/version.ts +++ b/src/version.ts @@ -1 +1 @@ -export const VERSION = "0.15.1" +export const VERSION = "0.15.2" diff --git a/test/mysql.test.ts b/test/mysql.test.ts index 01f06f5..c0c8e2c 100644 --- a/test/mysql.test.ts +++ b/test/mysql.test.ts @@ -146,6 +146,49 @@ describe("MySQL SQL compatibility", () => { }) describeMySQL("MySQL adapter", () => { + it("creates one instance when nested callers race to enqueue to a new actor", async () => { + if (!connectionString) throw new Error("MySQL connection string is required") + database = mysql({ connectionString, maximumConnections: 16 }) + runtime = configure({ + database, + tableNamePrefix: "mysql_test_", + authorizeMessage: () => true, + authorizeQuery: () => true, + logger: quietLogger, + }) + runtime.register(MySQLWorkflow) + await runtime.install() + const repository = runtime.repository + const actorId = `create-race-${crypto.randomUUID()}` + + const outcomes = await Promise.all( + Array.from({ length: 8 }, (_, index) => + database! + .transaction((connection) => + repository.enqueueInTransaction(connection, { + actorType: MySQLWorkflow.actorType, + actorId, + operation: "increment", + deliveryMode: "async", + arguments: {}, + idempotencyKey: `race-${index}`, + }), + ) + .then((message) => ({ sequence: Number(message.sequence) })) + .catch((error: { code?: string }) => ({ code: error.code ?? String(error) })), + ), + ) + + const failures = outcomes.filter((outcome) => "code" in outcome) + expect(failures).toEqual([]) + const sequences = outcomes.flatMap((outcome) => + "sequence" in outcome ? [outcome.sequence] : [], + ) + expect(sequences.slice().sort((left, right) => left - right)).toEqual([1, 2, 3, 4, 5, 6, 7, 8]) + const instance = await repository.findInstanceByIdentity(MySQLWorkflow.actorType, actorId) + expect(instance).toBeDefined() + }, 60_000) + it("keeps a fenced commit exclusive after its lease expires", async () => { if (!connectionString) throw new Error("MySQL connection string is required") database = mysql({ connectionString, maximumConnections: 5 }) diff --git a/test/postgresql.test.ts b/test/postgresql.test.ts index 99bc2a3..4eb1b96 100644 --- a/test/postgresql.test.ts +++ b/test/postgresql.test.ts @@ -180,6 +180,49 @@ describe("PostgreSQL SQL parameters", () => { }) describePostgreSQL("PostgreSQL adapter", () => { + it("creates one instance when nested callers race to enqueue to a new actor", async () => { + if (!connectionString) throw new Error("PostgreSQL connection string is required") + database = postgresql({ connectionString, maximumConnections: 16 }) + runtime = configure({ + database, + tableNamePrefix: "postgresql_test_", + authorizeMessage: () => true, + authorizeQuery: () => true, + logger: quietLogger, + }) + runtime.register(PostgreSQLCounter) + await runtime.install() + const repository = runtime.repository + const actorId = `create-race-${crypto.randomUUID()}` + + const outcomes = await Promise.all( + Array.from({ length: 8 }, (_, index) => + database! + .transaction((connection) => + repository.enqueueInTransaction(connection, { + actorType: PostgreSQLCounter.actorType, + actorId, + operation: "increment", + deliveryMode: "async", + arguments: {}, + idempotencyKey: `race-${index}`, + }), + ) + .then((message) => ({ sequence: Number(message.sequence) })) + .catch((error: { code?: string }) => ({ code: error.code ?? String(error) })), + ), + ) + + const failures = outcomes.filter((outcome) => "code" in outcome) + expect(failures).toEqual([]) + const sequences = outcomes.flatMap((outcome) => + "sequence" in outcome ? [outcome.sequence] : [], + ) + expect(sequences.slice().sort((left, right) => left - right)).toEqual([1, 2, 3, 4, 5, 6, 7, 8]) + const instance = await repository.findInstanceByIdentity(PostgreSQLCounter.actorType, actorId) + expect(instance).toBeDefined() + }, 60_000) + it("keeps a fenced commit exclusive after its lease expires", async () => { if (!connectionString) throw new Error("PostgreSQL connection string is required") database = postgresql({ connectionString, maximumConnections: 5 })