Skip to content

Commit e92ceb9

Browse files
committed
Release all process ownership
1 parent a9da0ce commit e92ceb9

6 files changed

Lines changed: 163 additions & 108 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@
22

33
## Unreleased
44

5+
- Release actor, message, effect, reminder, and broadcast ownership atomically
6+
during graceful shutdown and recover stale draining processes as well.
57
- Persist hostname, host PID, Node version, and Solid Objects version for every
68
runtime role and expose them through immutable process administration reads.
79
- Make typed snapshots return persisted fields and getter values from one

‎docs/correctness.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,11 @@
99
- Destruction creates an incarnation boundary; old leases cannot address a
1010
recreated actor. A caller authorized before destruction receives
1111
`ActorDestroyed`; an unknown or forged reference remains unauthorized.
12+
- Graceful process shutdown and stale-process cleanup use the same atomic
13+
ownership release: claimed messages return to ready membership, activations
14+
are unfenced, and processing effect, reminder, and broadcast claims become
15+
available again. A stale draining process is recoverable like a stale running
16+
process.
1217
- Permanent operation failure raises `MessageFailed` with the durable message
1318
ID and persisted error details instead of treating actor code text as the
1419
public exception contract.

‎docs/operations.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,9 @@ failure is logged without preventing lease release.
4848
metadata with hostname, host process ID, Node and Solid Objects versions, and a
4949
current `stale` flag. `cleanup()` reauthorizes separately and atomically fences
5050
stale processes out of every owned role claim before waking the affected
51-
runtime roles. The cleanup count is instrumented; application payloads are not.
51+
runtime roles. Graceful role shutdown performs the same atomic release instead
52+
of leaving claims to wait for lease expiry. Stale `draining` records are also
53+
recovered. The cleanup count is instrumented; application payloads are not.
5254

5355
Committed calls and `message.wait()` apply `timeoutMilliseconds` to the entire
5456
durable wait, beginning before enqueue or message lookup. Adapter deadlines

‎src/repository.ts‎

Lines changed: 59 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -77,15 +77,7 @@ export class Repository {
7777
async stopProcess(processId: string): Promise<void> {
7878
await this.settings.database.transaction(async (connection) => {
7979
const now = await connection.nowMilliseconds()
80-
await connection.run(
81-
`UPDATE ${this.table("processes")} SET shutdown_state = 'stopped', stopped_at_ms = ?, heartbeat_at_ms = ? WHERE id = ?`,
82-
[now, now, processId],
83-
)
84-
await connection.run(
85-
`UPDATE ${this.table("instances")} SET activation_owner_id = NULL, activation_token = NULL,
86-
activation_expires_at_ms = NULL WHERE activation_owner_id = ?`,
87-
[processId],
88-
)
80+
await this.releaseProcessOwnership({ connection, processId, now })
8981
})
9082
}
9183

@@ -119,65 +111,73 @@ export class Repository {
119111
const staleAt = now - this.settings.processAliveThresholdMilliseconds
120112
const processes = await connection.all<{ id: string }>(
121113
`SELECT id FROM ${this.table("processes")}
122-
WHERE shutdown_state = 'running' AND heartbeat_at_ms <= ?
114+
WHERE shutdown_state <> 'stopped' AND heartbeat_at_ms <= ?
123115
ORDER BY id`,
124116
[staleAt],
125117
)
126118
for (const process of processes) {
127-
const claims = await connection.all<{
128-
message_id: string
129-
instance_id: string
130-
sequence: number | bigint
131-
}>(
132-
`SELECT message_id, instance_id, sequence
133-
FROM ${this.table("claimed_messages")} WHERE process_id = ?`,
134-
[process.id],
135-
)
136-
for (const claim of claims) {
137-
await connection.run(
138-
`DELETE FROM ${this.table("claimed_messages")} WHERE message_id = ?`,
139-
[claim.message_id],
140-
)
141-
await connection.run(
142-
`INSERT INTO ${this.table("ready_messages")}
143-
(message_id, instance_id, sequence, available_at_ms)
144-
VALUES (?, ?, ?, ?) ON CONFLICT(message_id) DO NOTHING`,
145-
[claim.message_id, claim.instance_id, claim.sequence, now],
146-
)
147-
}
148-
await connection.run(
149-
`UPDATE ${this.table("instances")}
150-
SET activation_owner_id = NULL, activation_token = NULL,
151-
activation_expires_at_ms = NULL, updated_at_ms = ?
152-
WHERE activation_owner_id = ?`,
153-
[now, process.id],
154-
)
155-
await connection.run(
156-
`UPDATE ${this.table("effects")} SET status = 'pending', claimed_by = NULL
157-
WHERE status = 'processing' AND claimed_by = ?`,
158-
[process.id],
159-
)
160-
await connection.run(
161-
`UPDATE ${this.table("reminders")} SET claimed_by = NULL, claimed_at_ms = NULL
162-
WHERE claimed_by = ?`,
163-
[process.id],
164-
)
165-
await connection.run(
166-
`UPDATE ${this.table("broadcasts")} SET status = 'pending', claimed_by = NULL
167-
WHERE status = 'processing' AND claimed_by = ?`,
168-
[process.id],
169-
)
170-
await connection.run(
171-
`UPDATE ${this.table("processes")}
172-
SET shutdown_state = 'stopped', stopped_at_ms = ?, heartbeat_at_ms = ?
173-
WHERE id = ? AND shutdown_state = 'running'`,
174-
[now, now, process.id],
175-
)
119+
await this.releaseProcessOwnership({ connection, processId: process.id, now })
176120
}
177121
return processes.length
178122
})
179123
}
180124

125+
private async releaseProcessOwnership(options: {
126+
connection: DatabaseConnection
127+
processId: string
128+
now: number
129+
}): Promise<void> {
130+
const { connection, processId, now } = options
131+
const claims = await connection.all<{
132+
message_id: string
133+
instance_id: string
134+
sequence: number | bigint
135+
}>(
136+
`SELECT message_id, instance_id, sequence
137+
FROM ${this.table("claimed_messages")} WHERE process_id = ?`,
138+
[processId],
139+
)
140+
for (const claim of claims) {
141+
await connection.run(`DELETE FROM ${this.table("claimed_messages")} WHERE message_id = ?`, [
142+
claim.message_id,
143+
])
144+
await connection.run(
145+
`INSERT INTO ${this.table("ready_messages")}
146+
(message_id, instance_id, sequence, available_at_ms)
147+
VALUES (?, ?, ?, ?) ON CONFLICT(message_id) DO NOTHING`,
148+
[claim.message_id, claim.instance_id, claim.sequence, now],
149+
)
150+
}
151+
await connection.run(
152+
`UPDATE ${this.table("instances")}
153+
SET activation_owner_id = NULL, activation_token = NULL,
154+
activation_expires_at_ms = NULL, updated_at_ms = ?
155+
WHERE activation_owner_id = ?`,
156+
[now, processId],
157+
)
158+
await connection.run(
159+
`UPDATE ${this.table("effects")} SET status = 'pending', claimed_by = NULL
160+
WHERE status = 'processing' AND claimed_by = ?`,
161+
[processId],
162+
)
163+
await connection.run(
164+
`UPDATE ${this.table("reminders")} SET claimed_by = NULL, claimed_at_ms = NULL
165+
WHERE claimed_by = ?`,
166+
[processId],
167+
)
168+
await connection.run(
169+
`UPDATE ${this.table("broadcasts")} SET status = 'pending', claimed_by = NULL
170+
WHERE status = 'processing' AND claimed_by = ?`,
171+
[processId],
172+
)
173+
await connection.run(
174+
`UPDATE ${this.table("processes")}
175+
SET shutdown_state = 'stopped', stopped_at_ms = ?, heartbeat_at_ms = ?
176+
WHERE id = ? AND shutdown_state <> 'stopped'`,
177+
[now, now, processId],
178+
)
179+
}
180+
181181
async enqueue(input: EnqueueInput): Promise<MessageRow> {
182182
for (let attempt = 1; attempt <= 8; attempt += 1) {
183183
try {

‎src/runtime.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -872,7 +872,7 @@ export class SolidObjectsRuntime {
872872
metadata: processMetadata(row.metadata),
873873
shutdownState: row.shutdown_state,
874874
stale:
875-
row.shutdown_state === "running" &&
875+
row.shutdown_state !== "stopped" &&
876876
Number(row.heartbeat_at_ms) <= result.staleAtMilliseconds,
877877
startedAt: new Date(Number(row.started_at_ms)),
878878
heartbeatAt: new Date(Number(row.heartbeat_at_ms)),

‎test/process-administration.test.ts‎

Lines changed: 93 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -60,59 +60,41 @@ describe("process administration", () => {
6060
it("atomically releases every claim owned by a stale process", async () => {
6161
runtime = configuredRuntime()
6262
await runtime.install()
63-
const message = await ProcessActor.ref("claimed").send.run()
64-
await runtime.repository.registerProcess("stale", "worker")
65-
const turn = await runtime.repository.claim("stale")
66-
expect(turn?.message.id).toBe(message.id)
67-
await runtime.settings.database.connection(async (connection) => {
68-
await connection.run(
69-
`INSERT INTO ${runtime?.repository.table("effects")}
70-
(id, message_id, instance_id, name, arguments, status, max_attempts, available_at_ms, claimed_by)
71-
VALUES ('effect', ?, ?, 'test', '{}', 'processing', 1, 0, 'stale')`,
72-
[message.id, turn?.instance.id],
73-
)
74-
await connection.run(
75-
`INSERT INTO ${runtime?.repository.table("reminders")}
76-
(id, instance_id, operation, run_at_ms, arguments, missed_policy, status, claimed_by, claimed_at_ms)
77-
VALUES ('reminder', ?, 'run', 0, '{}', 'latest', 'scheduled', 'stale', 0)`,
78-
[turn?.instance.id],
79-
)
80-
await connection.run(
81-
`INSERT INTO ${runtime?.repository.table("broadcasts")}
82-
(id, message_id, instance_id, actor_type, actor_id, state_revision, observables,
83-
status, available_at_ms, claimed_by)
84-
VALUES ('broadcast', ?, ?, ?, 'claimed', 1, '{}', 'processing', 0, 'stale')`,
85-
[message.id, turn?.instance.id, ProcessActor.actorType],
86-
)
87-
})
63+
const message = await claimEveryRole("stale")
8864
await staleProcess("stale")
8965

9066
const result = await runtime.processes.cleanup()
9167

9268
expect(result.cleaned).toBe(1)
93-
expect(await message.status()).toBe("ready")
94-
const rows = await runtime.settings.database.connection(async (connection) => ({
95-
process: await connection.get<{ shutdown_state: string }>(
96-
`SELECT shutdown_state FROM ${runtime?.repository.table("processes")} WHERE id = 'stale'`,
97-
),
98-
instance: await connection.get<{ activation_owner_id: string | null }>(
99-
`SELECT activation_owner_id FROM ${runtime?.repository.table("instances")}`,
100-
),
101-
effect: await connection.get<{ status: string; claimed_by: string | null }>(
102-
`SELECT status, claimed_by FROM ${runtime?.repository.table("effects")}`,
103-
),
104-
reminder: await connection.get<{ claimed_by: string | null }>(
105-
`SELECT claimed_by FROM ${runtime?.repository.table("reminders")}`,
106-
),
107-
broadcast: await connection.get<{ status: string; claimed_by: string | null }>(
108-
`SELECT status, claimed_by FROM ${runtime?.repository.table("broadcasts")}`,
69+
await expectOwnershipReleased("stale", message)
70+
})
71+
72+
it("atomically releases every claim during graceful shutdown", async () => {
73+
runtime = configuredRuntime()
74+
await runtime.install()
75+
const message = await claimEveryRole("stopping")
76+
77+
await runtime.repository.stopProcess("stopping")
78+
79+
await expectOwnershipReleased("stopping", message)
80+
})
81+
82+
it("recovers a stale draining process", async () => {
83+
runtime = configuredRuntime()
84+
await runtime.install()
85+
const message = await claimEveryRole("draining")
86+
await runtime.settings.database.connection((connection) =>
87+
connection.run(
88+
`UPDATE ${runtime?.repository.table("processes")}
89+
SET shutdown_state = 'draining', heartbeat_at_ms = 0 WHERE id = 'draining'`,
10990
),
110-
}))
111-
expect(rows.process?.shutdown_state).toBe("stopped")
112-
expect(rows.instance?.activation_owner_id).toBeNull()
113-
expect(rows.effect).toEqual({ status: "pending", claimed_by: null })
114-
expect(rows.reminder?.claimed_by).toBeNull()
115-
expect(rows.broadcast).toEqual({ status: "pending", claimed_by: null })
91+
)
92+
93+
expect(await runtime.processes.all()).toEqual([
94+
expect.objectContaining({ id: "draining", shutdownState: "draining", stale: true }),
95+
])
96+
expect(await runtime.processes.cleanup()).toEqual({ cleaned: 1 })
97+
await expectOwnershipReleased("draining", message)
11698
})
11799

118100
it("leaves live process ownership unchanged", async () => {
@@ -151,3 +133,67 @@ async function staleProcess(id: string): Promise<void> {
151133
),
152134
)
153135
}
136+
137+
async function claimEveryRole(processId: string) {
138+
const message = await ProcessActor.ref(`claimed-${processId}`).send.run()
139+
await runtime?.repository.registerProcess(processId, "worker")
140+
const turn = await runtime?.repository.claim(processId)
141+
expect(turn?.message.id).toBe(message.id)
142+
await runtime?.settings.database.connection(async (connection) => {
143+
await connection.run(
144+
`INSERT INTO ${runtime?.repository.table("effects")}
145+
(id, message_id, instance_id, name, arguments, status, max_attempts, available_at_ms, claimed_by)
146+
VALUES (?, ?, ?, 'test', '{}', 'processing', 1, 0, ?)`,
147+
[`effect-${processId}`, message.id, turn?.instance.id, processId],
148+
)
149+
await connection.run(
150+
`INSERT INTO ${runtime?.repository.table("reminders")}
151+
(id, instance_id, operation, run_at_ms, arguments, missed_policy, status, claimed_by,
152+
claimed_at_ms)
153+
VALUES (?, ?, 'run', 0, '{}', 'latest', 'scheduled', ?, 0)`,
154+
[`reminder-${processId}`, turn?.instance.id, processId],
155+
)
156+
await connection.run(
157+
`INSERT INTO ${runtime?.repository.table("broadcasts")}
158+
(id, message_id, instance_id, actor_type, actor_id, state_revision, observables,
159+
status, available_at_ms, claimed_by)
160+
VALUES (?, ?, ?, ?, 'claimed', 1, '{}', 'processing', 0, ?)`,
161+
[`broadcast-${processId}`, message.id, turn?.instance.id, ProcessActor.actorType, processId],
162+
)
163+
})
164+
return message
165+
}
166+
167+
async function expectOwnershipReleased(
168+
processId: string,
169+
message: Awaited<ReturnType<typeof claimEveryRole>>,
170+
) {
171+
expect(await message.status()).toBe("ready")
172+
const rows = await runtime?.settings.database.connection(async (connection) => ({
173+
process: await connection.get<{ shutdown_state: string }>(
174+
`SELECT shutdown_state FROM ${runtime?.repository.table("processes")} WHERE id = ?`,
175+
[processId],
176+
),
177+
instance: await connection.get<{ activation_owner_id: string | null }>(
178+
`SELECT activation_owner_id FROM ${runtime?.repository.table("instances")} WHERE actor_id = ?`,
179+
[`claimed-${processId}`],
180+
),
181+
effect: await connection.get<{ status: string; claimed_by: string | null }>(
182+
`SELECT status, claimed_by FROM ${runtime?.repository.table("effects")} WHERE id = ?`,
183+
[`effect-${processId}`],
184+
),
185+
reminder: await connection.get<{ claimed_by: string | null }>(
186+
`SELECT claimed_by FROM ${runtime?.repository.table("reminders")} WHERE id = ?`,
187+
[`reminder-${processId}`],
188+
),
189+
broadcast: await connection.get<{ status: string; claimed_by: string | null }>(
190+
`SELECT status, claimed_by FROM ${runtime?.repository.table("broadcasts")} WHERE id = ?`,
191+
[`broadcast-${processId}`],
192+
),
193+
}))
194+
expect(rows?.process?.shutdown_state).toBe("stopped")
195+
expect(rows?.instance?.activation_owner_id).toBeNull()
196+
expect(rows?.effect).toEqual({ status: "pending", claimed_by: null })
197+
expect(rows?.reminder?.claimed_by).toBeNull()
198+
expect(rows?.broadcast).toEqual({ status: "pending", claimed_by: null })
199+
}

0 commit comments

Comments
 (0)