Skip to content

Commit b020540

Browse files
committed
fix(platform-cloudflare): preserve delayed request waiters
1 parent 0d9b328 commit b020540

8 files changed

Lines changed: 268 additions & 55 deletions

File tree

packages/platform/cloudflare/src/CloudflareCluster.ts

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -119,16 +119,16 @@ const notImplemented = (method: string) =>
119119

120120
interface EntityStub {
121121
readonly invoke: (envelope: string, discard: boolean, delivery?: {
122-
readonly deliverAt: number
123-
readonly primaryKey: string | null
122+
readonly deliverAt?: number | undefined
123+
readonly primaryKey?: string | null | undefined
124124
readonly replyTo?: string | undefined
125125
}) => Promise<{
126126
readonly requestId: string
127127
readonly replies: ReadonlyArray<string>
128128
readonly error?: "MailboxFull" | "EncodedMessageTooLarge" | undefined
129129
}>
130130
readonly acknowledge: (requestId: string, replyId: string) => Promise<ReadonlyArray<string>>
131-
readonly interrupt?: (requestId: string) => Promise<void>
131+
readonly interrupt?: (storageRequestId: string, clientRequestId?: string) => Promise<void>
132132
readonly reset?: (requestId: string) => Promise<void>
133133
}
134134

@@ -192,6 +192,7 @@ const make = Effect.fnUntraced(function*(options: LayerOptions) {
192192
readonly clientRequestId: string
193193
storageRequestId: string
194194
lastChunkId?: string
195+
replyHandler?: (reply: string) => Promise<void>
195196
}
196197
const entries = new Map<string, ClientEntry>()
197198
let client!: Effect.Success<ReturnType<typeof RpcClient.makeNoSerialization<any, MailboxFull | PersistenceError>>>
@@ -292,14 +293,17 @@ const make = Effect.fnUntraced(function*(options: LayerOptions) {
292293
}
293294
const delivery = delayed
294295
? { deliverAt: deliverAt!, primaryKey, ...(replyTo === undefined ? undefined : { replyTo }) }
295-
: undefined
296+
: replyTo === undefined
297+
? undefined
298+
: { replyTo }
296299
let replyHandler: ((reply: string) => Promise<void>) | undefined
297300
if (delivery?.replyTo !== undefined) {
298301
replyHandler = async (reply) => {
299-
unregisterReplyHandler(clientRequestId)
300-
unregisterReplyHandler(entry.storageRequestId)
302+
unregisterReplyHandler(clientRequestId, replyHandler)
303+
unregisterReplyHandler(entry.storageRequestId, replyHandler)
301304
await Effect.runPromise(deliverReplies(entry, [reply]))
302305
}
306+
entry.replyHandler = replyHandler
303307
registerReplyHandler(clientRequestId, replyHandler)
304308
}
305309
return Effect.promise(() => target.stub.invoke(envelope, discard, delivery)).pipe(
@@ -328,15 +332,15 @@ const make = Effect.fnUntraced(function*(options: LayerOptions) {
328332
return Effect.ensuring(
329333
deliver,
330334
Effect.sync(() => {
331-
unregisterReplyHandler(clientRequestId)
332-
unregisterReplyHandler(entry.storageRequestId)
335+
unregisterReplyHandler(clientRequestId, replyHandler)
336+
unregisterReplyHandler(entry.storageRequestId, replyHandler)
333337
})
334338
)
335339
}),
336340
Effect.tapCause(() =>
337341
Effect.sync(() => {
338-
unregisterReplyHandler(clientRequestId)
339-
unregisterReplyHandler(entry.storageRequestId)
342+
unregisterReplyHandler(clientRequestId, replyHandler)
343+
unregisterReplyHandler(entry.storageRequestId, replyHandler)
340344
})
341345
)
342346
)
@@ -354,10 +358,12 @@ const make = Effect.fnUntraced(function*(options: LayerOptions) {
354358
const entry = entries.get(clientRequestId)
355359
entries.delete(clientRequestId)
356360
requestTargets.delete(clientRequestId)
357-
unregisterReplyHandler(clientRequestId)
358-
if (entry !== undefined) unregisterReplyHandler(entry.storageRequestId)
361+
unregisterReplyHandler(clientRequestId, entry?.replyHandler)
362+
if (entry?.replyHandler !== undefined) {
363+
unregisterReplyHandler(entry.storageRequestId, entry.replyHandler)
364+
}
359365
if (entry === undefined || target.stub.interrupt === undefined) return Effect.void
360-
return Effect.promise(() => target.stub.interrupt!(entry.storageRequestId))
366+
return Effect.promise(() => target.stub.interrupt!(entry.storageRequestId, clientRequestId))
361367
}
362368
default:
363369
return Effect.void

packages/platform/cloudflare/src/CloudflareDurableObjects.ts

Lines changed: 52 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,7 @@ interface ReplayMessage {
7474
readonly lastSentChunk: string | undefined
7575
readonly discard: boolean
7676
readonly deliverAt?: number | undefined
77-
readonly replyTo?: string | undefined
77+
readonly replyTos?: ReadonlyArray<string> | undefined
7878
}
7979

8080
interface InvokeResult {
@@ -89,8 +89,8 @@ interface InvokeOutcome {
8989
}
9090

9191
interface DeliveryOptions {
92-
readonly deliverAt: number
93-
readonly primaryKey: string | null
92+
readonly deliverAt?: number | undefined
93+
readonly primaryKey?: string | null | undefined
9494
readonly replyTo?: string | undefined
9595
}
9696

@@ -121,7 +121,6 @@ export class ClusterEntity extends DurableObject<unknown> {
121121
#serial: Promise<void> = Promise.resolve()
122122
readonly #sessions = new Map<string, ReplySession>()
123123
readonly #workerWaiters = new Map<string, Array<WorkerWaiter>>()
124-
readonly #workerWaiterTargets = new Map<string, string>()
125124

126125
constructor(ctx: DurableObjectState, env: unknown) {
127126
super(ctx, env)
@@ -176,7 +175,7 @@ export class ClusterEntity extends DurableObject<unknown> {
176175
? undefined
177176
: {
178177
scheduled: true,
179-
...(row.replyTo === undefined ? undefined : { replyTo: row.replyTo })
178+
...(row.replyTos === undefined ? undefined : { replyTos: row.replyTos })
180179
}
181180
).pipe(
182181
Effect.catchCause((cause) => this.#completeReplayFailure(registration, row, cause))
@@ -223,7 +222,7 @@ export class ClusterEntity extends DurableObject<unknown> {
223222
if (persisted.processed) {
224223
return { result: { requestId: persisted.originalId, replies: [] } }
225224
}
226-
if (delivery !== undefined) {
225+
if (delivery?.deliverAt !== undefined) {
227226
yield* this.#armEarliestAlarm()
228227
return this.#delayedOutcome(
229228
persisted.originalId,
@@ -234,6 +233,15 @@ export class ClusterEntity extends DurableObject<unknown> {
234233
}
235234
const original = loadMessage(storage.sql, persisted.originalId)
236235
if (original === undefined) return yield* Effect.die("Duplicate mailbox row disappeared")
236+
if (original.deliverAt !== undefined && original.deliverAt > Date.now()) {
237+
yield* this.#armEarliestAlarm()
238+
return this.#delayedOutcome(
239+
persisted.originalId,
240+
discard,
241+
delivery?.replyTo,
242+
String(envelope.requestId)
243+
)
244+
}
237245
const replies = yield* this.#runStored(
238246
registration,
239247
runtime,
@@ -243,7 +251,7 @@ export class ClusterEntity extends DurableObject<unknown> {
243251
)
244252
return { result: { requestId: persisted.originalId, replies } }
245253
}
246-
if (delivery !== undefined) {
254+
if (delivery?.deliverAt !== undefined) {
247255
yield* this.#armEarliestAlarm()
248256
return this.#delayedOutcome(String(envelope.requestId), discard, delivery.replyTo)
249257
}
@@ -264,7 +272,6 @@ export class ClusterEntity extends DurableObject<unknown> {
264272
const waiters = this.#workerWaiters.get(requestId) ?? []
265273
waiters.push({ clientRequestId, resolve, reject })
266274
this.#workerWaiters.set(requestId, waiters)
267-
this.#workerWaiterTargets.set(clientRequestId, requestId)
268275
})
269276
return { result, deferred }
270277
}
@@ -287,22 +294,20 @@ export class ClusterEntity extends DurableObject<unknown> {
287294
}
288295

289296
/** @internal Interrupts an in-memory handler execution. Persisted rows remain replayable. */
290-
interrupt(requestId: string): Promise<void> {
291-
const storageRequestId = this.#workerWaiterTargets.get(requestId) ?? requestId
297+
interrupt(storageRequestId: string, clientRequestId = storageRequestId): Promise<void> {
292298
const waiters = this.#workerWaiters.get(storageRequestId)
293299
if (waiters !== undefined) {
294300
const remaining = waiters.filter((waiter) => {
295-
if (waiter.clientRequestId !== requestId) return true
301+
if (waiter.clientRequestId !== clientRequestId) return true
296302
waiter.reject(new Error("Delayed entity request interrupted"))
297303
return false
298304
})
299-
this.#workerWaiterTargets.delete(requestId)
300305
if (remaining.length === 0) this.#workerWaiters.delete(storageRequestId)
301306
else this.#workerWaiters.set(storageRequestId, remaining)
302307
}
303-
const session = this.#sessions.get(requestId)
308+
const session = this.#sessions.get(storageRequestId)
304309
if (session === undefined) return Promise.resolve()
305-
this.#sessions.delete(requestId)
310+
this.#sessions.delete(storageRequestId)
306311
session.ack?.resolve()
307312
session.ack = undefined
308313
session.done = true
@@ -338,7 +343,7 @@ export class ClusterEntity extends DurableObject<unknown> {
338343
envelopeText: string,
339344
lastSentChunk: string | undefined,
340345
discard: boolean,
341-
options?: { readonly scheduled?: boolean; readonly replyTo?: string | undefined }
346+
options?: { readonly scheduled?: boolean; readonly replyTos?: ReadonlyArray<string> | undefined }
342347
) {
343348
return Effect.flatMap(
344349
decodeRequest(registration, envelopeText),
@@ -358,7 +363,7 @@ export class ClusterEntity extends DurableObject<unknown> {
358363
(row) =>
359364
this.#runStored(registration, runtime, row.envelope, row.lastSentChunk, row.discard, {
360365
scheduled: true,
361-
...(row.replyTo === undefined ? undefined : { replyTo: row.replyTo })
366+
...(row.replyTos === undefined ? undefined : { replyTos: row.replyTos })
362367
}).pipe(
363368
Effect.catchCause((cause) => this.#completeReplayFailure(registration, row, cause))
364369
),
@@ -403,7 +408,7 @@ export class ClusterEntity extends DurableObject<unknown> {
403408
),
404409
(reply) =>
405410
Effect.sync(() => storage.transactionSync(() => saveReply(storage.sql, reply))).pipe(
406-
Effect.andThen(this.#deliverScheduledReply(encoded.requestId as string, reply, row.replyTo))
411+
Effect.andThen(this.#deliverScheduledReply(encoded.requestId as string, reply, row.replyTos))
407412
)
408413
).pipe(
409414
Effect.catchCause(() =>
@@ -422,7 +427,7 @@ export class ClusterEntity extends DurableObject<unknown> {
422427
lastSentChunkText: string | undefined,
423428
discard: boolean,
424429
persisted: boolean,
425-
options?: { readonly scheduled?: boolean; readonly replyTo?: string | undefined }
430+
options?: { readonly scheduled?: boolean; readonly replyTos?: ReadonlyArray<string> | undefined }
426431
): Effect.Effect<ReadonlyArray<string>> {
427432
const rpc = registration.entity.protocol.requests.get(envelope.tag) as Rpc.AnyWithProps
428433
const storage = this.#state.storage
@@ -457,7 +462,7 @@ export class ClusterEntity extends DurableObject<unknown> {
457462
yield* Effect.promise(() => this.#offerReply(session, encoded))
458463
}
459464
if (scheduled && reply._tag === "WithExit") {
460-
yield* this.#deliverScheduledReply(requestId, encoded, options?.replyTo)
465+
yield* this.#deliverScheduledReply(requestId, encoded, options?.replyTos)
461466
}
462467
})
463468
)
@@ -477,25 +482,48 @@ export class ClusterEntity extends DurableObject<unknown> {
477482
})
478483
}
479484

480-
#deliverScheduledReply(requestId: string, reply: string, replyTo: string | undefined): Effect.Effect<void> {
485+
#deliverScheduledReply(
486+
requestId: string,
487+
reply: string,
488+
replyTos: ReadonlyArray<string> | undefined
489+
): Effect.Effect<void> {
481490
const waiters = this.#workerWaiters.get(requestId)
482491
if (waiters !== undefined) {
483492
this.#workerWaiters.delete(requestId)
484493
for (const waiter of waiters) {
485-
this.#workerWaiterTargets.delete(waiter.clientRequestId)
486494
waiter.resolve({ requestId, replies: [reply] })
487495
}
488496
}
489-
if (replyTo === undefined) return Effect.void
497+
if (replyTos === undefined) return Effect.void
490498
const namespace = (this.#state.exports as Record<string, unknown>).ClusterEntity as
491499
| {
492500
readonly getByName: (
493501
name: string
494502
) => { readonly deliverReply: (requestId: string, reply: string) => Promise<boolean> }
495503
}
496504
| undefined
497-
if (namespace === undefined) return Effect.void
498-
return Effect.promise(() => namespace.getByName(replyTo).deliverReply(requestId, reply)).pipe(Effect.ignore)
505+
if (namespace === undefined) {
506+
return Effect.logError(
507+
"Scheduled entity reply delivery failed",
508+
new Error("CloudflareCluster: ClusterEntity export is unavailable for scheduled reply delivery")
509+
)
510+
}
511+
return Effect.forEach(
512+
replyTos,
513+
(replyTo) =>
514+
Effect.promise(() => namespace.getByName(replyTo).deliverReply(requestId, reply)).pipe(
515+
Effect.flatMap((delivered) =>
516+
delivered
517+
? Effect.void
518+
: Effect.logError(
519+
"Scheduled entity reply delivery failed",
520+
new Error(`Scheduled entity reply target is unavailable: ${replyTo}`)
521+
)
522+
),
523+
Effect.catchCause((cause) => Effect.logError("Scheduled entity reply delivery failed", cause))
524+
),
525+
{ discard: true }
526+
)
499527
}
500528

501529
/** @internal Completes an in-memory delayed ask owned by this entity object. */

packages/platform/cloudflare/src/internal/entityMailbox.ts

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ export const persistRequest = (
5353
}
5454

5555
const existing = sql.exec(
56-
`SELECT m.request_id, m.processed, r.reply AS last_reply
56+
`SELECT m.request_id, m.processed, m.reply_to, r.reply AS last_reply
5757
FROM cluster_messages m
5858
LEFT JOIN cluster_replies r ON r.reply_id = m.last_reply_id
5959
WHERE m.request_id = ? OR (? IS NOT NULL AND m.message_id = ?)
@@ -64,9 +64,11 @@ export const persistRequest = (
6464
).toArray()[0]
6565
if (existing !== undefined) {
6666
if (replyTo !== null && Number(existing.processed) === 0) {
67+
const replyTos = decodeReplyTargets(existing.reply_to)
68+
if (!replyTos.includes(replyTo)) replyTos.push(replyTo)
6769
sql.exec(
6870
"UPDATE cluster_messages SET reply_to = ? WHERE request_id = ?",
69-
replyTo,
71+
JSON.stringify(replyTos),
7072
String(existing.request_id)
7173
)
7274
}
@@ -99,7 +101,7 @@ export const persistRequest = (
99101
envelopeText,
100102
discard ? 1 : 0,
101103
deliverAt,
102-
replyTo
104+
replyTo === null ? null : JSON.stringify([replyTo])
103105
)
104106
return { _tag: "Success" }
105107
}
@@ -148,16 +150,28 @@ export interface StoredMessage {
148150
readonly lastSentChunk: string | undefined
149151
readonly discard: boolean
150152
readonly deliverAt?: number | undefined
151-
readonly replyTo?: string | undefined
153+
readonly replyTos?: ReadonlyArray<string> | undefined
154+
}
155+
156+
const decodeReplyTargets = (value: unknown): Array<string> => {
157+
if (typeof value !== "string") return []
158+
try {
159+
const decoded = JSON.parse(value)
160+
if (Array.isArray(decoded) && decoded.every((item) => typeof item === "string")) return decoded
161+
} catch {
162+
// Rows written before reply targets became a collection contain one plain name.
163+
}
164+
return [value]
152165
}
153166

154167
const rowToMessage = (row: Record<string, unknown>): StoredMessage => {
168+
const replyTos = decodeReplyTargets(row.reply_to)
155169
const message: StoredMessage = {
156170
envelope: String(row.envelope),
157171
lastSentChunk: typeof row.last_reply === "string" ? row.last_reply : undefined,
158172
discard: Number(row.discard) === 1,
159173
...(typeof row.deliver_at === "number" ? { deliverAt: row.deliver_at } : undefined),
160-
...(typeof row.reply_to === "string" ? { replyTo: row.reply_to } : undefined)
174+
...(replyTos.length === 0 ? undefined : { replyTos })
161175
}
162176
return message
163177
}

packages/platform/cloudflare/src/internal/entityReply.ts

Lines changed: 16 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -9,23 +9,32 @@ export const CurrentEntityName = Context.Reference<string | undefined>(
99

1010
type ReplyHandler = (reply: string) => Promise<void>
1111

12-
const handlers = new Map<string, ReplyHandler>()
12+
const handlers = new Map<string, Set<ReplyHandler>>()
1313

1414
/** @internal */
1515
export const registerReplyHandler = (requestId: string, handler: ReplyHandler): void => {
16-
handlers.set(requestId, handler)
16+
const registered = handlers.get(requestId) ?? new Set()
17+
registered.add(handler)
18+
handlers.set(requestId, registered)
1719
}
1820

1921
/** @internal */
20-
export const unregisterReplyHandler = (requestId: string): void => {
21-
handlers.delete(requestId)
22+
export const unregisterReplyHandler = (requestId: string, handler?: ReplyHandler): void => {
23+
if (handler === undefined) {
24+
handlers.delete(requestId)
25+
return
26+
}
27+
const registered = handlers.get(requestId)
28+
if (registered === undefined) return
29+
registered.delete(handler)
30+
if (registered.size === 0) handlers.delete(requestId)
2231
}
2332

2433
/** @internal */
2534
export const deliverReply = async (requestId: string, reply: string): Promise<boolean> => {
26-
const handler = handlers.get(requestId)
27-
if (handler === undefined) return false
35+
const registered = handlers.get(requestId)
36+
if (registered === undefined) return false
2837
handlers.delete(requestId)
29-
await handler(reply)
38+
await Promise.all(Array.from(registered, (handler) => handler(reply)))
3039
return true
3140
}

0 commit comments

Comments
 (0)