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
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
3 changes: 3 additions & 0 deletions docs/correctness.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -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",
Expand Down
34 changes: 24 additions & 10 deletions src/repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -298,13 +298,12 @@ export class Repository {
input: EnqueueInput,
): Promise<MessageRow> {
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")}
Expand All @@ -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<InstanceRow>(
`SELECT * FROM ${this.table("instances")} WHERE id = ?${this.rowLockClause()}`,
[identified.id],
)
if (!instance) throw new ActorDestroyed("actor disappeared during enqueue")

if (input.idempotencyKey !== undefined) {
Expand Down Expand Up @@ -2071,12 +2073,24 @@ export class Repository {
connection: DatabaseConnection
actorType: string
actorId: string
lock?: "update"
}): Promise<InstanceRow | undefined> {
const { connection, actorType, actorId } = options
return connection.get<InstanceRow>(
`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],
)
Expand Down
2 changes: 1 addition & 1 deletion src/version.ts
Original file line number Diff line number Diff line change
@@ -1 +1 @@
export const VERSION = "0.15.1"
export const VERSION = "0.15.2"
43 changes: 43 additions & 0 deletions test/mysql.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 })
Expand Down
43 changes: 43 additions & 0 deletions test/postgresql.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 })
Expand Down
Loading