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 })