Skip to content

Commit e34667e

Browse files
committed
feat: scan claim candidates
1 parent 4b818e3 commit e34667e

8 files changed

Lines changed: 128 additions & 53 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,8 @@
5656
and mailbox blockers.
5757
- Reject committed calls and message waits inside an ambient Solid Objects
5858
database transaction before they can self-deadlock.
59+
- Scan a bounded set of ready actor candidates so a lost lease race does not
60+
make a worker report idle while other actors are ready.
5961

6062
## 0.1.0 - 2026-08-13
6163

‎README.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -763,6 +763,8 @@ IDs and observable values are not authorization.
763763
- Different actor identities may execute concurrently.
764764
- A worker drains at most `maxMessagesPerActivationPass` turns from one actor,
765765
then yields its still-due work behind actors that were already waiting.
766+
- A global claim scans at most `claimScanLimit` ordered candidates, continuing
767+
to another ready actor after a lost lease race.
766768
- Long-running workers reuse hydrated actors for
767769
`idleDeactivationTimeoutMilliseconds` while renewing the same fenced lease.
768770
- State, completion, staged messages, effects, reminders, commit actions, and

‎docs/architecture.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,11 @@ moves only that actor's already-due ready memberships to current database time.
7070
Actors with older ready work therefore win the next global claim; delayed work
7171
keeps its original future availability.
7272

73+
The global claim reads at most `claimScanLimit` ordered candidates. If another
74+
worker acquires the first candidate's lease, the transaction continues through
75+
that bounded set instead of returning idle and sacrificing parallelism across
76+
independent actor identities.
77+
7378
When a pass becomes idle, a long-running worker keeps the hydrated actor and
7479
continues renewing the same fenced lease until its idle timeout. A later turn
7580
on that actor reuses both its persisted public fields and process-local private

‎docs/operations.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@ identities stay continuously busy; higher values reduce claim overhead for
2323
isolated backlogs. `solid_objects.activation.yielded` reports the actor
2424
identity, turns processed, and remaining due membership count.
2525

26+
`claimScanLimit` defaults to 100. Global claims inspect a bounded ordered set of
27+
actor identities and continue after a lost lease race, preserving worker
28+
parallelism without an unbounded scan.
29+
2630
Workers retain a hydrated actor and its fenced lease for
2731
`idleDeactivationTimeoutMilliseconds`, which defaults to 30 seconds. Idle
2832
leases renew at `leaseRenewalIntervalMilliseconds`; the worker polling cadence

‎docs/parity.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ Reference: Ruby `solid_objects` 0.12.0 at commit `a01b6f5`.
2626
| Ordered mailbox, sequence allocation, idempotency, retries, dead letters, leases, renewal, and fenced commits | Native | Relational ready/claimed membership tables, durable history, and adapter-appropriate sequence locking. |
2727
| Domain rejection and strict poison ordering | Native | Rejections roll back without retry; retryable failures block later operations until completion or dead-lettering. |
2828
| Bounded activation passes and hot-actor fairness | Native | Configurable turn-count and elapsed-time budgets bound each pass, then move only that actor's already-due memberships behind actors already waiting. |
29+
| Bounded claim candidate scan | Native | A configurable ordered scan continues to another ready actor when a worker loses the first candidate's lease race. |
2930
| Idle activation cache | Native | Long-running workers retain hydrated actors under renewable fenced leases, restore public state after failed turns, and release on timeout, fairness yield, lease loss, or shutdown. |
3031
| Transactional effects and outcome operations | Native | At-least-once effect handlers with stable IDs and success/failure actor operations. |
3132
| Actor-to-actor delivery | Native | `sendTo(reference).operation()` stages delivery in the source actor commit. |

‎src/configuration.ts‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ export interface SolidObjectsConfiguration {
5151
maxAttempts?: number
5252
maxMessagesPerActivationPass?: number
5353
maxActivationDurationMilliseconds?: number
54+
claimScanLimit?: number
5455
retryDelayMilliseconds?: (attempt: number) => number
5556
processHeartbeatIntervalMilliseconds?: number
5657
processAliveThresholdMilliseconds?: number
@@ -120,6 +121,7 @@ export function buildSettings(configuration: SolidObjectsConfiguration): Runtime
120121
maxAttempts: configuration.maxAttempts ?? 5,
121122
maxMessagesPerActivationPass: configuration.maxMessagesPerActivationPass ?? 50,
122123
maxActivationDurationMilliseconds: configuration.maxActivationDurationMilliseconds ?? 5_000,
124+
claimScanLimit: configuration.claimScanLimit ?? 100,
123125
retryDelayMilliseconds:
124126
configuration.retryDelayMilliseconds ??
125127
((attempt) => Math.min(2 ** (attempt - 1), 60) * 1_000),
@@ -210,11 +212,10 @@ function validateSettings(settings: RuntimeSettings): void {
210212
messageRetentionMilliseconds: settings.messageRetentionMilliseconds,
211213
processRetentionMilliseconds: settings.processRetentionMilliseconds,
212214
}
213-
if (
214-
!Number.isSafeInteger(settings.maxMessagesPerActivationPass) ||
215-
settings.maxMessagesPerActivationPass < 1
216-
) {
217-
throw new TypeError("maxMessagesPerActivationPass must be a positive safe integer")
215+
for (const name of ["maxMessagesPerActivationPass", "claimScanLimit"] as const) {
216+
if (!Number.isSafeInteger(settings[name]) || settings[name] < 1) {
217+
throw new TypeError(`${name} must be a positive safe integer`)
218+
}
218219
}
219220
for (const [name, value] of Object.entries(positive)) {
220221
if (!Number.isFinite(value) || value <= 0) throw new TypeError(`${name} must be positive`)

‎src/repository.ts‎

Lines changed: 63 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,7 @@ export class Repository {
320320
if (options.instanceId !== undefined) parameters.push(options.instanceId)
321321
parameters.push(now)
322322
}
323-
const candidate = await connection.get<{ message_id: string; instance_id: string }>(
323+
const candidates = await connection.all<{ message_id: string; instance_id: string }>(
324324
`SELECT ready.message_id, ready.instance_id
325325
FROM ${this.table("ready_messages")} ready
326326
JOIN ${this.table("instances")} instances ON instances.id = ready.instance_id
@@ -335,14 +335,18 @@ export class Repository {
335335
WHERE earlier.instance_id = ready.instance_id AND earlier.sequence < ready.sequence
336336
)
337337
${leaseCondition}
338-
ORDER BY ready.available_at_ms, ready.sequence
339-
LIMIT 1`,
340-
parameters,
338+
ORDER BY ready.available_at_ms, ready.sequence, ready.message_id
339+
LIMIT ?`,
340+
[
341+
...parameters,
342+
options.activation || options.instanceId ? 1 : this.settings.claimScanLimit,
343+
],
341344
)
342-
if (!candidate) return undefined
343-
344-
const token = options.activation?.activationToken ?? randomUUID()
345-
if (!options.activation) {
345+
for (const candidate of candidates) {
346+
const token = options.activation?.activationToken ?? randomUUID()
347+
if (options.activation) {
348+
return this.claimCandidate({ connection, candidate, processId, token, now })
349+
}
346350
const lease = await connection.run(
347351
`UPDATE ${this.table("instances")}
348352
SET activation_owner_id = ?, activation_token = ?, activation_generation = activation_generation + 1,
@@ -357,52 +361,63 @@ export class Repository {
357361
now,
358362
],
359363
)
360-
if (lease.changes !== 1) return undefined
364+
if (lease.changes !== 1) continue
365+
return this.claimCandidate({ connection, candidate, processId, token, now })
361366
}
367+
return undefined
368+
})
369+
}
362370

363-
const instance = await connection.get<InstanceRow>(
364-
`SELECT * FROM ${this.table("instances")} WHERE id = ?`,
365-
[candidate.instance_id],
366-
)
367-
const message = await connection.get<MessageRow>(
368-
`SELECT * FROM ${this.table("messages")} WHERE id = ?`,
369-
[candidate.message_id],
370-
)
371-
if (!instance || !message) return undefined
371+
private async claimCandidate(options: {
372+
connection: DatabaseConnection
373+
candidate: { message_id: string; instance_id: string }
374+
processId: string
375+
token: string
376+
now: number
377+
}): Promise<ClaimedTurn | undefined> {
378+
const { connection, candidate, processId, token, now } = options
379+
const instance = await connection.get<InstanceRow>(
380+
`SELECT * FROM ${this.table("instances")} WHERE id = ?`,
381+
[candidate.instance_id],
382+
)
383+
const message = await connection.get<MessageRow>(
384+
`SELECT * FROM ${this.table("messages")} WHERE id = ?`,
385+
[candidate.message_id],
386+
)
387+
if (!instance || !message) return undefined
372388

373-
const removed = await connection.run(
374-
`DELETE FROM ${this.table("ready_messages")} WHERE message_id = ?`,
375-
[message.id],
376-
)
377-
if (removed.changes !== 1) return undefined
378-
await connection.run(
379-
`INSERT INTO ${this.table("claimed_messages")}
389+
const removed = await connection.run(
390+
`DELETE FROM ${this.table("ready_messages")} WHERE message_id = ?`,
391+
[message.id],
392+
)
393+
if (removed.changes !== 1) return undefined
394+
await connection.run(
395+
`INSERT INTO ${this.table("claimed_messages")}
380396
(message_id, instance_id, sequence, process_id, activation_token, activation_generation, claimed_at_ms)
381397
VALUES (?, ?, ?, ?, ?, ?, ?)`,
382-
[
383-
message.id,
384-
instance.id,
385-
message.sequence,
386-
processId,
387-
token,
388-
instance.activation_generation,
389-
now,
390-
],
391-
)
392-
await connection.run(
393-
`UPDATE ${this.table("messages")} SET attempt_count = attempt_count + 1, updated_at_ms = ? WHERE id = ?`,
394-
[now, message.id],
395-
)
396-
message.attempt_count = Number(message.attempt_count) + 1
397-
return {
398-
instance,
399-
message,
398+
[
399+
message.id,
400+
instance.id,
401+
message.sequence,
400402
processId,
401-
activationToken: token,
402-
activationGeneration: BigInt(instance.activation_generation),
403-
nowMilliseconds: now,
404-
}
405-
})
403+
token,
404+
instance.activation_generation,
405+
now,
406+
],
407+
)
408+
await connection.run(
409+
`UPDATE ${this.table("messages")} SET attempt_count = attempt_count + 1, updated_at_ms = ? WHERE id = ?`,
410+
[now, message.id],
411+
)
412+
message.attempt_count = Number(message.attempt_count) + 1
413+
return {
414+
instance,
415+
message,
416+
processId,
417+
activationToken: token,
418+
activationGeneration: BigInt(instance.activation_generation),
419+
nowMilliseconds: now,
420+
}
406421
}
407422

408423
async renewTurn(turn: ClaimedTurn): Promise<void> {

‎test/lifecycle.test.ts‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -304,6 +304,51 @@ describe("runtime lifecycle", () => {
304304
expect(CachedActor.deactivations).toBe(1)
305305
})
306306

307+
it("claims another actor when an earlier candidate loses its lease race", async () => {
308+
runtime = configuredRuntime({ claimScanLimit: 2 })
309+
await runtime.install()
310+
const first = await LifecycleCounter.ref("first-candidate").send.increment()
311+
const second = await LifecycleCounter.ref("second-candidate").send.increment()
312+
await runtime.settings.database.connection(async (connection) => {
313+
await connection.run(
314+
`UPDATE ${runtime?.repository.table("ready_messages")} SET available_at_ms = 0 WHERE message_id = ?`,
315+
[first.id],
316+
)
317+
await connection.run(
318+
`UPDATE ${runtime?.repository.table("ready_messages")} SET available_at_ms = 1 WHERE message_id = ?`,
319+
[second.id],
320+
)
321+
})
322+
await runtime.repository.registerProcess("race-winner", "worker")
323+
const originalRun = runtime.settings.database.transaction.bind(runtime.settings.database)
324+
let leaseUpdates = 0
325+
const transaction = vi
326+
.spyOn(runtime.settings.database, "transaction")
327+
.mockImplementation((callback) =>
328+
originalRun((connection) =>
329+
callback({
330+
get: connection.get.bind(connection),
331+
all: connection.all.bind(connection),
332+
nowMilliseconds: connection.nowMilliseconds.bind(connection),
333+
run: async (sql, parameters) => {
334+
if (sql.includes("SET activation_owner_id")) {
335+
leaseUpdates += 1
336+
if (leaseUpdates === 1) return { changes: 0 }
337+
}
338+
return connection.run(sql, parameters)
339+
},
340+
}),
341+
),
342+
)
343+
344+
const turn = await runtime.repository.claim("race-winner")
345+
346+
expect(turn?.message.id).toBe(second.id)
347+
expect(await first.status()).toBe("ready")
348+
expect(await second.status()).toBe("claimed")
349+
transaction.mockRestore()
350+
})
351+
307352
it("restores cached state after a failed turn", async () => {
308353
runtime = configuredRuntime({ maxAttempts: 1 })
309354
runtime.register(CachedActor)

0 commit comments

Comments
 (0)