Skip to content

Commit 6fcd528

Browse files
committed
feat: read the schedule from an actor
reminder returns one armed alarm as a ScheduledReminder and reminders lists every key of one operation, matching what the Ruby port already offers. Both are async, because a TypeScript actor holds no rows and reads its own through a reader the runtime supplies at hydration. That is one injection point, so both engines get it: the SQL runtime reads the reminders table for the instance, and the Durable Objects engine reads the object's own store. A read starts from the committed rows and applies the intents staged so far, so an actor that schedules and then reads sees what the commit will write, and one that cancels and then reads sees the alarm gone. key and intervalMilliseconds are null rather than undefined. The first draft used undefined and a test caught it: returning a ScheduledReminder straight from an operation failed with InvalidPayload, because undefined does not serialise. Returning one is the obvious thing to do, so the type makes it work. A projection has no reader and throws rather than reporting an armed alarm as absent, which is the failure this whole feature exists to prevent. ScheduledReminder rather than ReminderStatus, because the administration API already exports that name for the status string.
1 parent ab5fd6a commit 6fcd528

13 files changed

Lines changed: 286 additions & 17 deletions

File tree

‎CHANGELOG.md‎

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

33
## Unreleased
44

5+
- Add reminder reading. `reminder()` returns one armed alarm as a
6+
`ScheduledReminder`, and `reminders()` lists every key of one operation. Both
7+
apply the intents staged so far in the turn, so a read agrees with what the
8+
commit will write. Reading works on the SQL backends and on Durable Objects.
59
- Refuse an unknown operation in `unschedule()` and `unscheduleAll()`.
610
`schedule()` already threw `UnknownOperation` for one, so a typo cancelled
711
nothing quietly and left a recurring reminder running.

‎docs/api.md‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,9 @@ createdAtMs }` shape returned by `SolidObjectsRuntime.snapshotWithIncarnation`.
7070
`DestroyOptions`: the options for authorization, idempotency, time, and
7171
schedule that the reference methods use.
7272

73+
`ScheduledReminder` is one armed reminder as an actor reads it, and
74+
`ReminderReader` is how a runtime supplies them.
75+
7376
`ActorIntents`, `EffectIntent`, `CommitActionIntent`, `ReminderIntent`,
7477
`UnscheduleIntent`, `UnscheduleAllIntent`, `ReminderMutation`,
7578
`OutboundMessageIntent`, `ReminderOptions`, `OutboundMessageOptions`,
@@ -245,6 +248,32 @@ declare, with the `UnknownOperation` that `schedule()` already throws, so a typo
245248
fails the turn rather than cancelling nothing. A handle skips that check, because
246249
the `schedule()` call that produced it was already checked.
247250

251+
#### Reading the schedule
252+
253+
`reminder()` returns the armed alarm as a `ScheduledReminder`, or `undefined`.
254+
`reminders()` returns every key of one operation. Both are async, because an
255+
actor reads its own rows rather than holding them in memory:
256+
257+
```typescript
258+
async nextChargeAt(): Promise<number | null> {
259+
return (await this.reminder("chargeRenewal"))?.runAtMilliseconds ?? null
260+
}
261+
262+
async pendingCarriers(): Promise<(string | null)[]> {
263+
return (await this.reminders("chaseCarrier")).map((reminder) => reminder.key)
264+
}
265+
```
266+
267+
A read applies the intents staged so far in the turn, so an actor that schedules
268+
and then reads sees what the commit will write, and one that cancels and then
269+
reads sees the alarm gone.
270+
271+
`key` and `intervalMilliseconds` are `null` rather than `undefined` when absent,
272+
so a `ScheduledReminder` returns from an operation without a serialization error.
273+
274+
Reading is available during a turn. A projection has no reader and throws, rather
275+
than reporting an armed alarm as absent.
276+
248277
A cancellation is staged like a schedule, so it commits with the state change
249278
that decided it and a turn that throws cancels nothing. Both apply in the order
250279
the turn called them, so cancelling and then scheduling the same name leaves it

‎docs/correctness.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@
55
- Delivery is ordered per actor identity and at least once.
66
- Different identities may execute concurrently.
77
- Sequence allocation and durable enqueue are one transaction.
8+
- An actor reads its own schedule. A read applies the intents staged so far in
9+
the turn, so it agrees with what the commit will write rather than with what
10+
the turn began with.
811
- A reminder can be cancelled. A cancellation commits with the state change that
912
decided it, and applies in the order the turn called it. A cancellation cannot
1013
recall an occurrence the scheduler already turned into a message. It does

‎src/actor.ts‎

Lines changed: 105 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { currentMessage, currentRuntime } from "./context.js"
22
import { getDefaultRuntime } from "./default-runtime.js"
33
import type { StateMigration } from "./definition.js"
44
import { InvalidRejectionCode, Rejected, UnknownOperation } from "./errors.js"
5+
import type { ReminderStatus } from "./reminder-administration.js"
56
import { TRANSMIT_EFFECT } from "./transmit-effect.js"
67
import { randomUUID } from "./platform/uuid.js"
78
import {
@@ -20,6 +21,8 @@ import type {
2021
EffectHandle,
2122
JsonObject,
2223
ReminderHandle,
24+
ReminderReader,
25+
ScheduledReminder,
2326
JsonValue,
2427
MessageContext,
2528
} from "./types.js"
@@ -179,6 +182,66 @@ function validatedReminderKey(key: string | number | undefined): string | undefi
179182
* reverse, and it is refused here rather than at the insert, once the turn is
180183
* already doing work.
181184
*/
185+
function handleName(handle: ReminderHandle, key: string | number | undefined): string {
186+
if (key !== undefined) throw new TypeError("a reminder handle already names its key")
187+
188+
const name = handle?.name
189+
if (typeof name !== "string" || name.length === 0) {
190+
throw new TypeError("a reminder handle returned by schedule is required")
191+
}
192+
193+
return name
194+
}
195+
196+
function reminderKeyOf(name: string, operation: string): string | null {
197+
return name === operation ? null : name.slice(operation.length + 1)
198+
}
199+
200+
function reminderStatusOf(options: {
201+
name: string
202+
operation: string
203+
runAtMilliseconds: number
204+
intervalMilliseconds: number | null
205+
missedPolicy: "all" | "latest"
206+
status: ReminderStatus
207+
}): ScheduledReminder {
208+
return {
209+
name: options.name,
210+
operation: options.operation,
211+
key: reminderKeyOf(options.name, options.operation),
212+
runAtMilliseconds: options.runAtMilliseconds,
213+
intervalMilliseconds: options.intervalMilliseconds,
214+
missedPolicy: options.missedPolicy,
215+
status: options.status,
216+
handle: { name: options.name },
217+
}
218+
}
219+
220+
function applyReminderIntent(view: Map<string, ScheduledReminder>, intent: ReminderMutation): void {
221+
if (intent.cancel === "all") {
222+
for (const [name, status] of view) {
223+
if (status.operation === intent.operation) view.delete(name)
224+
}
225+
return
226+
}
227+
if (intent.cancel === "one") {
228+
view.delete(intent.name)
229+
return
230+
}
231+
232+
view.set(
233+
intent.name,
234+
reminderStatusOf({
235+
name: intent.name,
236+
operation: intent.operation,
237+
runAtMilliseconds: intent.atMilliseconds,
238+
intervalMilliseconds: intent.intervalMilliseconds ?? null,
239+
missedPolicy: intent.missedPolicy,
240+
status: "scheduled",
241+
}),
242+
)
243+
}
244+
182245
function reminderName(operation: string, key: string | undefined): string {
183246
if (key === undefined) return operation
184247

@@ -211,6 +274,8 @@ export abstract class Actor {
211274
}
212275

213276
readonly #actorId: string
277+
#readReminders: ReminderReader | undefined
278+
214279
readonly #intents: ActorIntents = {
215280
effects: [],
216281
commitActions: [],
@@ -363,25 +428,50 @@ export abstract class Actor {
363428
}
364429

365430
unschedule(operationOrHandle: string | ReminderHandle, options: { key?: string | number } = {}) {
431+
this.#intents.reminders.push({
432+
cancel: "one",
433+
name: this.#reminderNameOf(operationOrHandle, options),
434+
})
435+
}
436+
437+
unscheduleAll(operation: string) {
438+
this.#assertOperation(operation)
439+
this.#intents.reminders.push({ cancel: "all", operation })
440+
}
441+
442+
async reminder(
443+
operationOrHandle: string | ReminderHandle,
444+
options: { key?: string | number } = {},
445+
): Promise<ScheduledReminder | undefined> {
446+
return (await this.#reminderView()).get(this.#reminderNameOf(operationOrHandle, options))
447+
}
448+
449+
async reminders(operation: string): Promise<ScheduledReminder[]> {
450+
this.#assertOperation(operation)
451+
const view = await this.#reminderView()
452+
return [...view.values()].filter((status) => status.operation === operation)
453+
}
454+
455+
#reminderNameOf(
456+
operationOrHandle: string | ReminderHandle,
457+
options: { key?: string | number },
458+
): string {
366459
if (typeof operationOrHandle === "string") {
367460
this.#assertOperation(operationOrHandle)
368-
const key = validatedReminderKey(options.key)
369-
this.#intents.reminders.push({ cancel: "one", name: reminderName(operationOrHandle, key) })
370-
return
371-
}
372-
if (options.key !== undefined) throw new TypeError("a reminder handle already names its key")
373-
374-
const name = operationOrHandle?.name
375-
if (typeof name !== "string" || name.length === 0) {
376-
throw new TypeError("unschedule requires a reminder handle returned by schedule")
461+
return reminderName(operationOrHandle, validatedReminderKey(options.key))
377462
}
378463

379-
this.#intents.reminders.push({ cancel: "one", name })
464+
return handleName(operationOrHandle, options.key)
380465
}
381466

382-
unscheduleAll(operation: string) {
383-
this.#assertOperation(operation)
384-
this.#intents.reminders.push({ cancel: "all", operation })
467+
async #reminderView(): Promise<Map<string, ScheduledReminder>> {
468+
if (!this.#readReminders) {
469+
throw new TypeError("reading reminders is not available outside an actor turn")
470+
}
471+
const view = new Map<string, ScheduledReminder>()
472+
for (const status of await this.#readReminders()) view.set(status.name, status)
473+
for (const intent of this.#intents.reminders) applyReminderIntent(view, intent)
474+
return view
385475
}
386476

387477
#assertOperation(operation: string): void {
@@ -413,8 +503,9 @@ export abstract class Actor {
413503
}
414504

415505
/** @internal */
416-
prepare(operations: ReadonlySet<string>): void {
506+
prepare(operations: ReadonlySet<string>, readReminders?: ReminderReader): void {
417507
this.#operations = operations
508+
this.#readReminders = readReminders
418509
}
419510

420511
/** @internal */

‎src/cloudflare/engine.ts‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ import type {
3434
JsonObject,
3535
JsonValue,
3636
SerializedError,
37+
ScheduledReminder,
3738
} from "../types.js"
3839
import type { CloudflareSettings } from "./configuration.js"
3940
import { actorName, callHost, type ActorIdentity, type HostRequest } from "./protocol.js"
@@ -357,6 +358,21 @@ export class ActorEngine {
357358
return { definition, instance, state }
358359
}
359360

361+
private readReminders = async (): Promise<ScheduledReminder[]> =>
362+
this.store.rows<Reminder>("SELECT record FROM reminders ORDER BY name").map((reminder) => ({
363+
name: reminder.name,
364+
operation: reminder.operation,
365+
key:
366+
reminder.name === reminder.operation
367+
? null
368+
: reminder.name.slice(reminder.operation.length + 1),
369+
runAtMilliseconds: reminder.at,
370+
intervalMilliseconds: reminder.interval,
371+
missedPolicy: reminder.missed,
372+
status: reminder.status,
373+
handle: { name: reminder.name },
374+
}))
375+
360376
private async snapshot(identity: ActorIdentity): Promise<JsonObject> {
361377
const { definition, instance, state } = this.committed(identity)
362378
const actor = hydrateActor({ definition, actorId: identity.actorId, state })
@@ -387,7 +403,11 @@ export class ActorEngine {
387403
}): Promise<JsonObject> {
388404
const { input, payloadNames } = options
389405
const { definition, instance, state } = this.committed(input)
390-
const actor = hydrateActor({ definition, actorId: input.actorId, state })
406+
const actor = hydrateActor({
407+
definition,
408+
actorId: input.actorId,
409+
state,
410+
})
391411
const identity = {
392412
actorType: input.actorType,
393413
actorId: input.actorId,
@@ -486,6 +506,7 @@ export class ActorEngine {
486506
storedVersion: instance.stateVersion,
487507
storedState: instance.state,
488508
}),
509+
readReminders: this.readReminders,
489510
})
490511
if (this.cached?.actor !== actor) {
491512
await withActorContext({ actor, runtime: this.runtime }, () => actor.activate())

‎src/definition.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { Actor, type ActorClass } from "./actor.js"
22
import { withApplicationWritesForbidden } from "./context.js"
33
import { ApplicationWriteForbidden, InvalidActor, StateMigrationError } from "./errors.js"
44
import { deepCopy, jsonObject, normalizeJson } from "./serialization.js"
5+
import type { ReminderReader } from "./types.js"
56
import type { JsonObject } from "./types.js"
67

78
export type PayloadBroadcastHandler = (
@@ -146,10 +147,11 @@ export function hydrateActor<ActorType extends Actor>(options: {
146147
definition: ValidatedActorDefinition<ActorType>
147148
actorId: string
148149
state: JsonObject
150+
readReminders?: ReminderReader
149151
}): ActorType {
150152
const { definition, actorId, state } = options
151153
const actor = new definition.actorClass(actorId)
152-
actor.prepare(new Set(definition.operations))
154+
actor.prepare(new Set(definition.operations), options.readReminders)
153155
const target = actor as unknown as Record<string, unknown>
154156
for (const key of definition.stateKeys) target[key] = deepCopy(state[key])
155157
return actor

‎src/index.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,8 @@ export type {
132132
EffectFailurePayload,
133133
EffectHandle,
134134
ReminderHandle,
135+
ReminderReader,
136+
ScheduledReminder,
135137
EffectSuccessPayload,
136138
InvocationOptions,
137139
JsonObject,

‎src/repository.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1169,6 +1169,15 @@ export class Repository {
11691169
})
11701170
}
11711171

1172+
async remindersForInstance(instanceId: string): Promise<ReminderRow[]> {
1173+
return this.settings.database.connection((connection) =>
1174+
connection.all<ReminderRow>(
1175+
`SELECT * FROM ${this.table("reminders")} WHERE instance_id = ? ORDER BY operation`,
1176+
[instanceId],
1177+
),
1178+
)
1179+
}
1180+
11721181
async findInstanceByIdentity(
11731182
actorType: string,
11741183
actorId: string,

‎src/runtime.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,7 @@ import type {
117117
MessageContext,
118118
MessageStatus,
119119
SnapshotOptions,
120+
ScheduledReminder,
120121
} from "./types.js"
121122
import { SolidObjectsTestHelper } from "./test-helper.js"
122123
import { waitFor, Worker } from "./worker.js"
@@ -165,6 +166,21 @@ type CommitActionHandler = (
165166
context: CommitActionContext,
166167
) => unknown | Promise<unknown>
167168

169+
function scheduledReminderOf(row: ReminderRow): ScheduledReminder {
170+
const operation = row.message_operation ?? row.operation
171+
const name = row.operation
172+
return {
173+
name,
174+
operation,
175+
key: name === operation ? null : name.slice(operation.length + 1),
176+
runAtMilliseconds: Number(row.run_at_ms),
177+
intervalMilliseconds: row.interval_ms === null ? null : Number(row.interval_ms),
178+
missedPolicy: row.missed_policy,
179+
status: row.status,
180+
handle: { name },
181+
}
182+
}
183+
168184
export class SolidObjectsRuntime {
169185
readonly settings
170186
readonly repository
@@ -1117,6 +1133,8 @@ export class SolidObjectsRuntime {
11171133
definition,
11181134
actorId: turn.message.actor_id,
11191135
state: deepCopy(state),
1136+
readReminders: async () =>
1137+
(await this.repository.remindersForInstance(turn.instance.id)).map(scheduledReminderOf),
11201138
})
11211139
} catch (error) {
11221140
throw new ActorSetupFailed(error)

‎src/types.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import type { ReminderStatus } from "./reminder-administration.js"
2+
13
export type JsonPrimitive = null | boolean | number | string
24
export type JsonValue = JsonPrimitive | JsonValue[] | { [key: string]: JsonValue }
35
export type JsonObject = { [key: string]: JsonValue }
@@ -6,6 +8,19 @@ export type EffectHandle = { readonly id: string }
68

79
export type ReminderHandle = { readonly name: string }
810

11+
export interface ScheduledReminder {
12+
readonly name: string
13+
readonly operation: string
14+
readonly key: string | null
15+
readonly runAtMilliseconds: number
16+
readonly intervalMilliseconds: number | null
17+
readonly missedPolicy: "all" | "latest"
18+
readonly status: ReminderStatus
19+
readonly handle: ReminderHandle
20+
}
21+
22+
export type ReminderReader = () => Promise<readonly ScheduledReminder[]>
23+
924
export type SerializedError = {
1025
name: string
1126
message: string

0 commit comments

Comments
 (0)