Skip to content

Commit f5c667a

Browse files
committed
feat: schedule runtime housekeeping
1 parent f627578 commit f5c667a

8 files changed

Lines changed: 319 additions & 6 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,9 @@
4747
timeout, fairness yield, lease loss, and shutdown.
4848
- Bound activation passes by configurable elapsed time as well as message
4949
count so slow hot actors yield workers fairly.
50+
- Add independent supervised schedulers for expired message and process
51+
retention and stale process recovery, with disable switches and bounded
52+
failure backoff.
5053

5154
## 0.1.0 - 2026-08-13
5255

‎docs/operations.md‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -61,12 +61,21 @@ every repair through the actor's typed `send` dispatcher, optionally with a
6161
future `availableAt` to spread large repairs. Never update persisted actor
6262
state from reconciliation code.
6363

64-
Retention is explicit through `runtime.retention.preview()` and
64+
The runtime automatically prunes expired message and stopped-process history
65+
once at startup and every `retentionIntervalMilliseconds`, which defaults to
66+
one hour. Stale process ownership is recovered independently every
67+
`deadProcessCleanupIntervalMilliseconds`, which defaults to one minute. Set an
68+
interval to zero to disable its scheduler. Failed passes emit metadata-only
69+
events and retry with bounded exponential backoff without stopping other
70+
runtime roles.
71+
72+
Operators can also use `runtime.retention.preview()` and
6573
`runtime.retention.prune()`. Both calls require administration authorization;
6674
use preview first and alert on an unexpected count before executing deletion.
6775
Message history defaults to 30 days with optional per-actor overrides. Stopped
6876
process history defaults to 7 days. Instance expiration is disabled unless an
69-
actor type appears in `instanceRetentionByActorType`.
77+
actor type appears in `instanceRetentionByActorType`, and remains an explicit
78+
operator action because it deletes the entire actor incarnation.
7079

7180
Pruning selects and rechecks at most `pruneBatchSize` rows per transaction. It
7281
preserves ready and claimed messages, dead-letter originals and replacements,

‎docs/parity.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,12 +41,12 @@ Reference: Ruby `solid_objects` 0.12.0 at commit `a01b6f5`.
4141

4242
| Capability | Status | TypeScript shape or remaining work |
4343
| ----------------------------------------------------------------------------- | ------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------- |
44-
| Process registration, heartbeats, stale claim recovery, and graceful shutdown | Native | Runtime roles persist process state and release ownership on shutdown; authorized inspection and cleanup recover every stale claim type. |
44+
| Process registration, heartbeats, stale claim recovery, and graceful shutdown | Native | Runtime roles persist process state and release ownership on shutdown; scheduled and authorized manual cleanup recover every stale claim type. |
4545
| Failed-role replacement | Native | Built-in and registered roles are rebuilt through their factories with capped backoff; shutdown is the terminal replacement boundary. |
4646
| Additional supervised components | Native | `registerComponent()` builds, validates, runs, and stops application components with the runtime. |
4747
| Dead-letter inspection and retry | Native | `runtime.deadLetters` provides deny-by-default immutable inspection and idempotent durable retry linkage. |
4848
| Reconciliation reads | Native | Authorized cursor pages cover active, quiet, and orphaned instances; bounded state batches are migrated and deeply frozen. |
49-
| Message, process, and opt-in instance retention | Native | Authorized preview and bounded prune APIs preserve live work and support default, per-actor, and opt-in policies. |
49+
| Message, process, and opt-in instance retention | Native | Supervised scheduling bounds message and process growth; authorized manual APIs add preview and keep destructive instance expiration explicit. |
5050
| Doctor and schema verification | Native | Structured checks cover configuration, schema/version shape, adapter server versions, policy posture, live roles, and a targeted round trip. |
5151
| CLI | Native | The packaged executable loads an application runtime and exposes start, diagnostics, processes, dead letters, reminders, and explicit retention pruning as JSON. |
5252
| Structured instrumentation | Native | An isolated transport-neutral sink emits immutable lifecycle metadata and structurally excludes application payloads. |

‎src/configuration.ts‎

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,8 @@ export interface SolidObjectsConfiguration {
6060
effectWorkerCount?: number
6161
broadcastWorkerCount?: number
6262
reminderSchedulerCount?: number
63+
retentionIntervalMilliseconds?: number
64+
deadProcessCleanupIntervalMilliseconds?: number
6365
messageRetentionMilliseconds?: number
6466
messageRetentionByActorType?: Readonly<Record<string, number>>
6567
instanceRetentionByActorType?: Readonly<Record<string, number>>
@@ -131,6 +133,9 @@ export function buildSettings(configuration: SolidObjectsConfiguration): Runtime
131133
effectWorkerCount: configuration.effectWorkerCount ?? 1,
132134
broadcastWorkerCount: configuration.broadcastWorkerCount ?? 1,
133135
reminderSchedulerCount: configuration.reminderSchedulerCount ?? 1,
136+
retentionIntervalMilliseconds: configuration.retentionIntervalMilliseconds ?? 3_600_000,
137+
deadProcessCleanupIntervalMilliseconds:
138+
configuration.deadProcessCleanupIntervalMilliseconds ?? 60_000,
134139
messageRetentionMilliseconds: configuration.messageRetentionMilliseconds ?? 30 * 86_400_000,
135140
messageRetentionByActorType: Object.freeze({
136141
...(configuration.messageRetentionByActorType ?? {}),
@@ -224,6 +229,12 @@ function validateSettings(settings: RuntimeSettings): void {
224229
) {
225230
throw new TypeError("idleDeactivationTimeoutMilliseconds must be non-negative")
226231
}
232+
for (const [name, value] of Object.entries({
233+
retentionIntervalMilliseconds: settings.retentionIntervalMilliseconds,
234+
deadProcessCleanupIntervalMilliseconds: settings.deadProcessCleanupIntervalMilliseconds,
235+
})) {
236+
if (!Number.isFinite(value) || value < 0) throw new TypeError(`${name} must be non-negative`)
237+
}
227238
if (
228239
settings.supervisorMaximumRestartDelayMilliseconds < settings.supervisorRestartDelayMilliseconds
229240
) {
@@ -253,7 +264,9 @@ function validateSettings(settings: RuntimeSettings): void {
253264
settings.workerCount +
254265
settings.effectWorkerCount +
255266
settings.reminderSchedulerCount +
256-
(broadcastsEnabled(settings) ? settings.broadcastWorkerCount : 0)
267+
(broadcastsEnabled(settings) ? settings.broadcastWorkerCount : 0) +
268+
(settings.retentionIntervalMilliseconds > 0 ? 1 : 0) +
269+
(settings.deadProcessCleanupIntervalMilliseconds > 0 ? 1 : 0)
257270
if (roleCount === 0) throw new TypeError("at least one runtime role must be configured")
258271
}
259272

‎src/maintenance-scheduler.ts‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
import type { SolidObjectsRuntime } from "./runtime.js"
2+
import { waitFor } from "./worker.js"
3+
4+
export class MaintenanceScheduler {
5+
private stopping = false
6+
7+
constructor(
8+
private readonly options: {
9+
runtime: SolidObjectsRuntime
10+
intervalMilliseconds: number
11+
failureEvent: string
12+
operation: () => Promise<void>
13+
},
14+
) {}
15+
16+
async run(signal: AbortSignal): Promise<void> {
17+
let failureCount = 0
18+
while (!signal.aborted && !this.stopping) {
19+
try {
20+
await this.options.operation()
21+
failureCount = 0
22+
} catch (error) {
23+
failureCount += 1
24+
this.options.runtime.emitInstrumentation(this.options.failureEvent, {
25+
errorName: error instanceof Error ? error.name : "Error",
26+
failureCount,
27+
})
28+
}
29+
await waitFor(this.pauseMilliseconds(failureCount), signal)
30+
}
31+
}
32+
33+
requestShutdown(): void {
34+
this.stopping = true
35+
}
36+
37+
stopped(): boolean {
38+
return this.stopping
39+
}
40+
41+
stop(): void {
42+
this.stopping = true
43+
}
44+
45+
private pauseMilliseconds(failureCount: number): number {
46+
if (failureCount === 0) return this.options.intervalMilliseconds
47+
const exponent = Math.min(failureCount - 1, 16)
48+
return Math.min(
49+
this.options.runtime.settings.supervisorRestartDelayMilliseconds * 2 ** exponent,
50+
this.options.runtime.settings.supervisorMaximumRestartDelayMilliseconds,
51+
this.options.intervalMilliseconds,
52+
)
53+
}
54+
}

‎src/runtime.ts‎

Lines changed: 38 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,7 @@ import {
6666
} from "./process-administration.js"
6767
import { RealtimeManager } from "./realtime.js"
6868
import type { PayloadEnvelope } from "./browser/index.js"
69+
import { MaintenanceScheduler } from "./maintenance-scheduler.js"
6970
import type {
7071
BroadcastRow,
7172
ClaimedTurn,
@@ -253,6 +254,8 @@ export class SolidObjectsRuntime {
253254
broadcastWorkerCount: broadcastsEnabled(this.settings)
254255
? this.settings.broadcastWorkerCount
255256
: 0,
257+
retentionEnabled: this.settings.retentionIntervalMilliseconds > 0,
258+
deadProcessCleanupEnabled: this.settings.deadProcessCleanupIntervalMilliseconds > 0,
256259
})
257260
const controller = new AbortController()
258261
const abort = () => controller.abort(signal.reason)
@@ -607,6 +610,11 @@ export class SolidObjectsRuntime {
607610
resource: "processes",
608611
authorizationContext: options.authorizationContext,
609612
})
613+
const cleaned = await this.cleanupStaleProcesses()
614+
return Object.freeze({ cleaned })
615+
}
616+
617+
private async cleanupStaleProcesses(): Promise<number> {
610618
const cleaned = await this.repository.cleanupStaleProcesses()
611619
if (cleaned > 0) {
612620
this.wakeUp("actors")
@@ -615,7 +623,7 @@ export class SolidObjectsRuntime {
615623
this.wakeUp("broadcasts")
616624
this.emitInstrumentation("processes.cleaned", { count: cleaned })
617625
}
618-
return Object.freeze({ cleaned })
626+
return cleaned
619627
}
620628

621629
async instancesWithoutPendingWork(
@@ -1158,6 +1166,28 @@ export class SolidObjectsRuntime {
11581166
() => () => this.broadcastWorker(),
11591167
)
11601168
: []),
1169+
...(this.settings.retentionIntervalMilliseconds > 0
1170+
? [
1171+
() =>
1172+
new MaintenanceScheduler({
1173+
runtime: this,
1174+
intervalMilliseconds: this.settings.retentionIntervalMilliseconds,
1175+
failureEvent: "supervisor.retention_failed",
1176+
operation: () => this.pruneExpiredRecords(),
1177+
}),
1178+
]
1179+
: []),
1180+
...(this.settings.deadProcessCleanupIntervalMilliseconds > 0
1181+
? [
1182+
() =>
1183+
new MaintenanceScheduler({
1184+
runtime: this,
1185+
intervalMilliseconds: this.settings.deadProcessCleanupIntervalMilliseconds,
1186+
failureEvent: "supervisor.process_cleanup_failed",
1187+
operation: () => this.cleanupStaleProcesses().then(() => undefined),
1188+
}),
1189+
]
1190+
: []),
11611191
]
11621192

11631193
for (const { count, factory } of this.additionalComponents) {
@@ -1179,6 +1209,13 @@ export class SolidObjectsRuntime {
11791209
}
11801210
}
11811211

1212+
private async pruneExpiredRecords(): Promise<void> {
1213+
for (const target of ["messages", "processes"] as const) {
1214+
const count = await this.repository.pruneRetention(target)
1215+
this.emitInstrumentation(`${target}.pruned`, { count })
1216+
}
1217+
}
1218+
11821219
private async supervise(slot: ComponentSlot, signal: AbortSignal): Promise<void> {
11831220
let failureCount = 0
11841221
let replacementErrorName: string | null = null

‎test/definition.test.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,8 +219,16 @@ describe("runtime configuration", () => {
219219
workerCount: 0,
220220
effectWorkerCount: 0,
221221
reminderSchedulerCount: 0,
222+
retentionIntervalMilliseconds: 0,
223+
deadProcessCleanupIntervalMilliseconds: 0,
222224
}),
223225
).toThrow("at least one runtime role")
226+
expect(() => buildSettings({ database, retentionIntervalMilliseconds: -1 })).toThrow(
227+
"retentionIntervalMilliseconds must be non-negative",
228+
)
229+
expect(() => buildSettings({ database, deadProcessCleanupIntervalMilliseconds: -1 })).toThrow(
230+
"deadProcessCleanupIntervalMilliseconds must be non-negative",
231+
)
224232
expect(() => validateComponent({} as never)).toThrow("must implement run")
225233

226234
await database.close()

0 commit comments

Comments
 (0)