From 4d6f0f3a7466c19c15a94671b7e3282f327b3864 Mon Sep 17 00:00:00 2001 From: Stephen Brian King Date: Fri, 17 Jul 2026 00:00:18 -0600 Subject: [PATCH 1/5] Track the next scheduled delivery time in message storage --- ...cluster-message-storage-next-deliver-at.md | 5 + .../src/unstable/cluster/MessageStorage.ts | 40 ++++- .../src/unstable/cluster/SqlMessageStorage.ts | 90 +++++++++++ .../test/cluster/MessageStorage.test.ts | 150 +++++++++++++++++- .../test/cluster/MessageStorageTest.ts | 17 +- .../test/cluster/SqlMessageStorage.test.ts | 43 ++++- 6 files changed, 337 insertions(+), 8 deletions(-) create mode 100644 .changeset/cluster-message-storage-next-deliver-at.md diff --git a/.changeset/cluster-message-storage-next-deliver-at.md b/.changeset/cluster-message-storage-next-deliver-at.md new file mode 100644 index 00000000000..1e25a94082b --- /dev/null +++ b/.changeset/cluster-message-storage-next-deliver-at.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Add a required `MessageStorage.nextDeliverAt` operation that returns the earliest future scheduled delivery time for a set of shards, with memory and SQL driver implementations and a supporting index migration for all SQL dialects. diff --git a/packages/effect/src/unstable/cluster/MessageStorage.ts b/packages/effect/src/unstable/cluster/MessageStorage.ts index 40a6ad37e2c..47897e8ecf5 100644 --- a/packages/effect/src/unstable/cluster/MessageStorage.ts +++ b/packages/effect/src/unstable/cluster/MessageStorage.ts @@ -143,6 +143,14 @@ export class MessageStorage extends Context.Service ) => Effect.Effect>, PersistenceError> + /** + * Retrieves the earliest `deliverAt` of the scheduled messages for the + * specified shards that are not yet deliverable. + */ + readonly nextDeliverAt: ( + shardIds: Iterable + ) => Effect.Effect, PersistenceError> + /** * Reset the mailbox state for the provided shards. */ @@ -375,6 +383,15 @@ export type Encoded = { PersistenceError > + /** + * Retrieves the earliest `deliverAt` of the scheduled messages for the + * given shards that is later than `now`. + */ + readonly nextDeliverAt: ( + shardIds: Arr.NonEmptyArray, + now: number + ) => Effect.Effect, PersistenceError> + /** * Reset the mailbox state for the provided address. */ @@ -651,6 +668,11 @@ export const makeEncoded: (encoded: Encoded) => Effect.Effect< (messages) => decodeMessages(storage, messages) ) }, + nextDeliverAt(shardIds) { + const shards = Array.from(shardIds, (id) => id.toString()) + if (!Arr.isArrayNonEmpty(shards)) return Effect.succeedNone + return Effect.suspend(() => encoded.nextDeliverAt(shards, clock.currentTimeMillisUnsafe())) + }, resetAddress: encoded.resetAddress, clearAddress: encoded.clearAddress, resetShards: (shardIds) => { @@ -779,6 +801,7 @@ export const noop: MessageStorage["Service"] = Effect.runSync(make({ requestIdForPrimaryKey: () => Effect.succeedNone, unprocessedMessages: () => Effect.succeed([]), unprocessedMessagesById: () => Effect.succeed([]), + nextDeliverAt: () => Effect.succeedNone, resetAddress: () => Effect.void, clearAddress: () => Effect.void, resetShards: () => Effect.void, @@ -849,7 +872,7 @@ export class MemoryDriver extends Context.Service()("effect/cluste } if (envelope._tag === "Request") { const entry = requests.get(envelope.requestId) - if (entry?.deliverAt && entry.deliverAt > now) { + if (entry?.deliverAt != null && entry.deliverAt > now) { continue } messages.push({ @@ -967,7 +990,7 @@ export class MemoryDriver extends Context.Service()("effect/cluste } if (envelope._tag === "Request") { const entry = requests.get(envelope.requestId)! - if (entry.deliverAt && entry.deliverAt > now) { + if (entry.deliverAt != null && entry.deliverAt > now) { continue } messages.push({ @@ -992,6 +1015,19 @@ export class MemoryDriver extends Context.Service()("effect/cluste } return unprocessedWith((envelope) => envelopeIds.has(envelope.requestId)) }), + nextDeliverAt: (shardIds, now) => + Effect.sync(() => { + let min: number | undefined + for (const envelope of unprocessed) { + if (envelope._tag !== "Request") continue + const shardId = ShardId.make(envelope.address.shardId.group, envelope.address.shardId.id) + if (!shardIds.includes(shardId.toString())) continue + const entry = requests.get(envelope.requestId) + if (entry?.deliverAt == null || entry.deliverAt <= now) continue + if (min === undefined || entry.deliverAt < min) min = entry.deliverAt + } + return min === undefined ? Option.none() : Option.some(min) + }), resetAddress: () => Effect.void, clearAddress: (address) => Effect.sync(() => { diff --git a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts index 63c1c38c480..5c0eea9517e 100644 --- a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts +++ b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts @@ -403,6 +403,44 @@ export const make: (options?: { ) }) + const getNextDeliverAt = sql.onDialectOrElse({ + mysql: () => (shardIds: ReadonlyArray, now: number) => { + const shards = sql.literal( + shardIds.map((id, i) => i === 0 ? `SELECT ${wrapString(id)} AS shard_id` : `SELECT ${wrapString(id)}`).join( + " UNION ALL " + ) + ) + return sql<{ readonly next_deliver_at: bigint | null }>` + SELECT MIN(t.next_at) AS next_deliver_at FROM ( + SELECT ( + SELECT MIN(deliver_at) FROM ${messagesTableSql} + WHERE shard_id = s.shard_id AND processed = ${sqlFalse} AND deliver_at > ${sql.literal(String(now))} + ) AS next_at + FROM (${shards}) AS s + ) AS t + ` + }, + mssql: () => (shardIds: ReadonlyArray, now: number) => + sql<{ readonly next_deliver_at: bigint | null }>` + SELECT MIN(l.next_at) AS next_deliver_at + FROM (VALUES ${sql.literal(shardIds.map((id) => `(${wrapString(id)})`).join(","))}) AS s(shard_id) + CROSS APPLY ( + SELECT TOP 1 deliver_at AS next_at + FROM ${messagesTableSql} + WHERE shard_id = s.shard_id AND processed = ${sqlFalse} AND deliver_at > ${sql.literal(String(now))} + ORDER BY deliver_at + ) AS l + `, + orElse: () => (shardIds: ReadonlyArray, now: number) => + sql<{ readonly next_deliver_at: bigint | null }>` + SELECT MIN(deliver_at) AS next_deliver_at + FROM ${messagesTableSql} + WHERE shard_id IN (${sql.literal(shardIds.map(wrapString).join(","))}) + AND processed = ${sqlFalse} + AND deliver_at > ${sql.literal(String(now))} + ` + }) + return yield* MessageStorage.makeEncoded({ saveEnvelope: ({ deliverAt, envelope, primaryKey }) => Effect.suspend(() => { @@ -577,6 +615,17 @@ export const make: (options?: { ) }, + nextDeliverAt: (shardIds, now) => + getNextDeliverAt(shardIds, now).unprepared.pipe( + Effect.map((rows) => { + const next = rows[0]?.next_deliver_at + return next == null ? Option.none() : Option.some(Number(next)) + }), + Effect.provideService(SqlClient.SafeIntegers, true), + PersistenceError.refail, + withTracerDisabled + ), + resetAddress: (address) => sql` UPDATE ${messagesTableSql} @@ -977,6 +1026,47 @@ const migrations = (options?: { // sqlite Effect.void }) + }), + "0003_deliver_at_index": Effect.gen(function*() { + const sql = (yield* SqlClient.SqlClient).withoutTransforms() + const messagesTableSql = sql(messagesTable) + const deliverAtIndex = `${messagesTable}_deliver_at_idx` + + yield* sql.onDialectOrElse({ + mssql: () => + sql` + IF NOT EXISTS (SELECT * FROM sys.indexes WHERE name = ${deliverAtIndex}) + CREATE INDEX ${sql(deliverAtIndex)} + ON ${messagesTableSql} (shard_id, processed, deliver_at, last_read) + `, + mysql: () => + sql` + CREATE INDEX ${sql(deliverAtIndex)} + ON ${messagesTableSql} (shard_id, processed, deliver_at, last_read) + `.unprepared.pipe(Effect.ignore), + pg: () => + sql` + CREATE INDEX IF NOT EXISTS ${sql(deliverAtIndex)} + ON ${messagesTableSql} (shard_id, deliver_at) + WHERE processed = FALSE AND deliver_at IS NOT NULL + `.pipe( + Effect.tapDefect((error) => + Effect.annotateLogs(Effect.logDebug("Failed to create indexes", error), { + package: "@effect/cluster", + module: "SqlMessageStorage" + }) + ), + Effect.retry({ + schedule: Schedule.spaced(1000) + }) + ), + orElse: () => + sql` + CREATE INDEX IF NOT EXISTS ${sql(deliverAtIndex)} + ON ${messagesTableSql} (shard_id, deliver_at) + WHERE processed = FALSE AND deliver_at IS NOT NULL + ` + }) }) }) } diff --git a/packages/effect/test/cluster/MessageStorage.test.ts b/packages/effect/test/cluster/MessageStorage.test.ts index f580519bc78..33a2326a9a9 100644 --- a/packages/effect/test/cluster/MessageStorage.test.ts +++ b/packages/effect/test/cluster/MessageStorage.test.ts @@ -1,7 +1,8 @@ -import { describe, expect, it } from "@effect/vitest" -import { Context, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" +import { assert, describe, expect, it } from "@effect/vitest" +import { Clock, Context, DateTime, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" import { TestClock } from "effect/testing" import { + DeliverAt, EntityAddress, EntityId, EntityType, @@ -61,6 +62,135 @@ describe("MessageStorage", () => { expect(messages).toHaveLength(0) }).pipe(Effect.provide(MemoryLive))) + it.effect("nextDeliverAt returns none when nothing is scheduled", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const next = yield* storage.nextDeliverAt([ShardId.make("default", 1)]) + assert.isTrue(Option.isNone(next)) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt returns the minimum future delivery time", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const now = yield* Clock.currentTimeMillis + const shard1 = ShardId.make("default", 1) + const shard2 = ShardId.make("default", 2) + const earlier = now + 60_000 + const later = now + 120_000 + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(later) }), + shardId: shard1 + }) + ) + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 2, deliverAt: DateTime.makeUnsafe(earlier) }), + shardId: shard2 + }) + ) + + const next = yield* storage.nextDeliverAt([shard1, shard2]) + assert.deepStrictEqual(next, Option.some(earlier)) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt handles epoch zero across the due boundary", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const shard = ShardId.make("default", 1) + yield* TestClock.setTime(-1) + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(0) }), + shardId: shard + }) + ) + + let next = yield* storage.nextDeliverAt([shard]) + assert.deepStrictEqual(next, Option.some(0)) + let messages = yield* storage.unprocessedMessages([shard]) + assert.strictEqual(messages.length, 0) + + yield* TestClock.setTime(0) + next = yield* storage.nextDeliverAt([shard]) + assert.isTrue(Option.isNone(next)) + messages = yield* storage.unprocessedMessages([shard]) + assert.strictEqual(messages.length, 1) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt ignores messages on other shards", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const now = yield* Clock.currentTimeMillis + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(now + 60_000) }), + shardId: ShardId.make("default", 2) + }) + ) + + const next = yield* storage.nextDeliverAt([ShardId.make("default", 1)]) + assert.isTrue(Option.isNone(next)) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt ignores already-due messages", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const now = yield* Clock.currentTimeMillis + const shard = ShardId.make("default", 1) + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(now) }), + shardId: shard + }) + ) + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 2, deliverAt: DateTime.makeUnsafe(now - 1) }), + shardId: shard + }) + ) + + const next = yield* storage.nextDeliverAt([shard]) + assert.isTrue(Option.isNone(next)) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt ignores completed requests", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const now = yield* Clock.currentTimeMillis + const request = yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(now + 60_000) }) + }) + yield* storage.saveRequest(request) + yield* storage.saveReply(yield* makeReply(request)) + + const next = yield* storage.nextDeliverAt([request.envelope.address.shardId]) + assert.isTrue(Option.isNone(next)) + }).pipe(Effect.provide(MemoryLive))) + + it.effect("nextDeliverAt returns none for an empty shard iterable", () => + Effect.gen(function*() { + const storage = yield* MessageStorage.MessageStorage + const now = yield* Clock.currentTimeMillis + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(now + 60_000) }) + }) + ) + + const next = yield* storage.nextDeliverAt([]) + assert.isTrue(Option.isNone(next)) + }).pipe(Effect.provide(MemoryLive))) + it.effect("repliesFor", () => Effect.gen(function*() { const storage = yield* MessageStorage.MessageStorage @@ -101,6 +231,7 @@ export const GetUserRpc = Rpc.make("GetUser", { export const makeRequest = Effect.fnUntraced(function*(options?: { readonly rpc?: Rpc.AnyWithProps readonly payload?: any + readonly shardId?: ShardId.ShardId }) { const snowflake = yield* Snowflake.Generator const rpc = options?.rpc ?? GetUserRpc @@ -108,7 +239,7 @@ export const makeRequest = Effect.fnUntraced(function*(options?: { envelope: Envelope.makeRequest({ requestId: snowflake.nextUnsafe(), address: EntityAddress.make({ - shardId: ShardId.make("default", 1), + shardId: options?.shardId ?? ShardId.make("default", 1), entityType: EntityType.make("test"), entityId: EntityId.make("1") }), @@ -136,6 +267,19 @@ export class PrimaryKeyTest extends Rpc.make("PrimaryKeyTest", { primaryKey: (value) => value.id.toString() }) {} +export class ScheduledPayload extends Schema.Class("ScheduledPayload")({ + id: Schema.Number, + deliverAt: Schema.DateTimeUtcFromMillis +}) { + [DeliverAt.symbol]() { + return this.deliverAt + } +} + +export class ScheduledRpc extends Rpc.make("ScheduledRpc", { + payload: ScheduledPayload +}) {} + export class StreamRpc extends Rpc.make("StreamTest", { success: RpcSchema.Stream(Schema.Void, Schema.Never), payload: { diff --git a/packages/platform-node/test/cluster/MessageStorageTest.ts b/packages/platform-node/test/cluster/MessageStorageTest.ts index f580519bc78..4021d769b4a 100644 --- a/packages/platform-node/test/cluster/MessageStorageTest.ts +++ b/packages/platform-node/test/cluster/MessageStorageTest.ts @@ -2,6 +2,7 @@ import { describe, expect, it } from "@effect/vitest" import { Context, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" import { TestClock } from "effect/testing" import { + DeliverAt, EntityAddress, EntityId, EntityType, @@ -101,6 +102,7 @@ export const GetUserRpc = Rpc.make("GetUser", { export const makeRequest = Effect.fnUntraced(function*(options?: { readonly rpc?: Rpc.AnyWithProps readonly payload?: any + readonly shardId?: ShardId.ShardId }) { const snowflake = yield* Snowflake.Generator const rpc = options?.rpc ?? GetUserRpc @@ -108,7 +110,7 @@ export const makeRequest = Effect.fnUntraced(function*(options?: { envelope: Envelope.makeRequest({ requestId: snowflake.nextUnsafe(), address: EntityAddress.make({ - shardId: ShardId.make("default", 1), + shardId: options?.shardId ?? ShardId.make("default", 1), entityType: EntityType.make("test"), entityId: EntityId.make("1") }), @@ -136,6 +138,19 @@ export class PrimaryKeyTest extends Rpc.make("PrimaryKeyTest", { primaryKey: (value) => value.id.toString() }) {} +export class ScheduledPayload extends Schema.Class("ScheduledPayload")({ + id: Schema.Number, + deliverAt: Schema.DateTimeUtcFromMillis +}) { + [DeliverAt.symbol]() { + return this.deliverAt + } +} + +export class ScheduledRpc extends Rpc.make("ScheduledRpc", { + payload: ScheduledPayload +}) {} + export class StreamRpc extends Rpc.make("StreamTest", { success: RpcSchema.Stream(Schema.Void, Schema.Never), payload: { diff --git a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts index 16deaaa863d..b7421db1d39 100644 --- a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts +++ b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts @@ -1,9 +1,9 @@ import { NodeFileSystem } from "@effect/platform-node" import { SqliteClient } from "@effect/sql-sqlite-node" import { assert, describe, expect, it } from "@effect/vitest" -import { Effect, Fiber, FileSystem, Latch, Layer, Option } from "effect" +import { Clock, DateTime, Effect, Fiber, FileSystem, Latch, Layer, Option } from "effect" import { TestClock } from "effect/testing" -import { Message, MessageStorage, ShardingConfig, Snowflake, SqlMessageStorage } from "effect/unstable/cluster" +import { Message, MessageStorage, ShardId, ShardingConfig, Snowflake, SqlMessageStorage } from "effect/unstable/cluster" import { SqlClient } from "effect/unstable/sql" import { MysqlContainer } from "../fixtures/mysql2-utils.ts" import { PgContainer } from "../fixtures/pg-utils.ts" @@ -13,6 +13,8 @@ import { makeReply, makeRequest, PrimaryKeyTest, + ScheduledPayload, + ScheduledRpc, StreamRpc } from "./MessageStorageTest.ts" @@ -163,6 +165,43 @@ describe("SqlMessageStorage", () => { expect(messages).toHaveLength(0) })) + it.effect("nextDeliverAt", () => + Effect.gen(function*() { + yield* truncate + + const storage = yield* MessageStorage.MessageStorage + const shard1 = ShardId.make("default", 1) + const shard2 = ShardId.make("default", 2) + let next = yield* storage.nextDeliverAt([shard1, shard2]) + assert.isTrue(Option.isNone(next)) + + const now = yield* Clock.currentTimeMillis + const earlier = now + 60_000 + const later = now + 120_000 + const earlierRequest = yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(earlier) }), + shardId: shard1 + }) + const laterRequest = yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 2, deliverAt: DateTime.makeUnsafe(later) }), + shardId: shard2 + }) + yield* storage.saveRequest(earlierRequest) + yield* storage.saveRequest(laterRequest) + + next = yield* storage.nextDeliverAt([shard1, shard2]) + assert.deepStrictEqual(next, Option.some(earlier)) + + next = yield* storage.nextDeliverAt([shard2]) + assert.deepStrictEqual(next, Option.some(later)) + + yield* storage.saveReply(yield* makeReply(earlierRequest)) + next = yield* storage.nextDeliverAt([shard1, shard2]) + assert.deepStrictEqual(next, Option.some(later)) + })) + it.effect("repliesFor", () => Effect.gen(function*() { yield* truncate From 6e63cbbd17c84fab6678d0849c010f2a87e24850 Mon Sep 17 00:00:00 2001 From: Stephen Brian King Date: Fri, 17 Jul 2026 00:00:46 -0600 Subject: [PATCH 2/5] Wake the storage read loop when scheduled messages become deliverable --- .../effect/src/unstable/cluster/Sharding.ts | 26 +++ packages/effect/test/cluster/Sharding.test.ts | 186 +++++++++++++++++- packages/effect/test/cluster/TestEntity.ts | 20 +- 3 files changed, 230 insertions(+), 2 deletions(-) diff --git a/packages/effect/src/unstable/cluster/Sharding.ts b/packages/effect/src/unstable/cluster/Sharding.ts index a202047cf87..f0f4d852a66 100644 --- a/packages/effect/src/unstable/cluster/Sharding.ts +++ b/packages/effect/src/unstable/cluster/Sharding.ts @@ -20,6 +20,7 @@ import * as Effect from "../../Effect.ts" import * as Equal from "../../Equal.ts" import type * as Exit from "../../Exit.ts" import * as Fiber from "../../Fiber.ts" +import * as FiberHandle from "../../FiberHandle.ts" import * as FiberMap from "../../FiberMap.ts" import { constant, flow } from "../../Function.ts" import * as HashRing from "../../HashRing.ts" @@ -472,6 +473,12 @@ const make = Effect.gen(function*() { yield* Effect.logDebug("Starting") yield* Effect.addFinalizer(() => Effect.logDebug("Shutting down")) + // A single replaceable timer that re-opens the storage read latch when the + // next scheduled message becomes deliverable. Re-derived from storage after + // every read, so a missed or failed wake only delays delivery until the + // next poll interval. + const storageWakeHandle = yield* FiberHandle.make() + let index = 0 let messages: Array> = [] const removableNotifications = new Set() @@ -555,6 +562,24 @@ const make = Effect.gen(function*() { } ) + const rediscoverStorageWake = Effect.suspend(() => storage.nextDeliverAt(acquiredShards)).pipe( + Effect.flatMap(Option.match({ + onNone: () => FiberHandle.clear(storageWakeHandle), + onSome: (deliverAt) => { + const remainingMillis = deliverAt - clock.currentTimeMillisUnsafe() + if (remainingMillis <= 0) { + return Effect.asVoid(storageReadLatch.open) + } + return storageReadLatch.open.pipe( + Effect.delay(Duration.millis(remainingMillis)), + FiberHandle.run(storageWakeHandle), + Effect.asVoid + ) + } + })), + Effect.catchCause((cause) => Effect.logDebug("Could not read the next scheduled delivery time", cause)) + ) + while (true) { // wait for the next poll interval, or if we get notified of a change yield* storageReadLatch.await @@ -576,6 +601,7 @@ const make = Effect.gen(function*() { messages = yield* storage.unprocessedMessages(acquiredShards) index = 0 yield* processMessages + yield* rediscoverStorageWake if (removableNotifications.size > 0) { removableNotifications.forEach(({ message, resume }) => { diff --git a/packages/effect/test/cluster/Sharding.test.ts b/packages/effect/test/cluster/Sharding.test.ts index 2a526bbc692..4ca68006d7c 100644 --- a/packages/effect/test/cluster/Sharding.test.ts +++ b/packages/effect/test/cluster/Sharding.test.ts @@ -1,7 +1,8 @@ import { assert, describe, expect, it } from "@effect/vitest" -import { Array, Cause, Clock, Effect, Exit, Fiber, Layer, MutableRef, Option, Queue, Stream } from "effect" +import { Array, Cause, Clock, DateTime, Effect, Exit, Fiber, Layer, MutableRef, Option, Queue, Stream } from "effect" import { TestClock } from "effect/testing" import { + ClusterError, MessageStorage, RunnerAddress, RunnerHealth, @@ -183,6 +184,189 @@ describe.concurrent("Sharding", () => { assert(reply._tag === "WithExit" && reply.exit._tag === "Failure" && reply.exit.cause[0]._tag === "Die") }).pipe(Effect.provide(TestSharding))) + it.effect("delivers scheduled messages at their deadline instead of the next poll", () => + Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + const fiber = yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 2000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + + // the message is persisted, but not yet deliverable + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + assert.isUndefined(fiber.pollUnsafe()) + + // just before the deadline nothing is delivered + yield* TestClock.adjust(1997) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + // the deadline wake delivers the message before the 5000ms poll interval + yield* TestClock.adjust(2) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + const exit = fiber.pollUnsafe() + assert(exit && Exit.isSuccess(exit)) + }).pipe(Effect.provide(TestSharding))) + + it.effect("wakes for each scheduled deadline independently", () => + Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 3500) + }).pipe(Effect.forkChild) + yield* client.GetUserScheduled({ + id: 2, + deliverAt: DateTime.makeUnsafe(now + 2000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + // the earlier deadline fires without cancelling the later one + yield* TestClock.adjust(2000) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + + yield* TestClock.adjust(1500) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 2) + }).pipe(Effect.provide(TestSharding))) + + it.effect("delivers scheduled messages with a past deadline immediately", () => + Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + const fiber = yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now - 1000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + const exit = fiber.pollUnsafe() + assert(exit && Exit.isSuccess(exit)) + }).pipe(Effect.provide(TestSharding))) + + it.effect("re-arms the wake when an earlier deadline arrives", () => + Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 4000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + + yield* client.GetUserScheduled({ + id: 2, + deliverAt: DateTime.makeUnsafe(now + 1000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1000) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + + yield* TestClock.adjust(3000) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 2) + }).pipe(Effect.provide(TestSharding))) + + it.effect("discovers scheduled messages persisted before the runner starts", () => + Effect.gen(function*() { + const driver = yield* MessageStorage.MemoryDriver + const state = yield* TestEntityState + const now = yield* Clock.currentTimeMillis + const EnvLayer = TestShardingWithoutState.pipe( + Layer.provide(Runners.layerNoop), + Layer.provide(TestShardingConfig) + ) + + yield* Effect.gen(function*() { + yield* TestClock.adjust(1) + const makeClient = yield* TestEntity.client + const client = makeClient("1") + yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 3000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + }).pipe( + Effect.provide(EnvLayer), + Effect.scoped + ) + + assert.strictEqual(driver.unprocessed.size, 1) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + yield* Effect.gen(function*() { + yield* TestClock.adjust(1) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + yield* TestClock.adjust(2996) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + yield* TestClock.adjust(1) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + }).pipe( + Effect.provide(EnvLayer), + Effect.scoped + ) + }).pipe(Effect.provide(MessageStorage.layerMemory.pipe( + Layer.provide(TestShardingConfig), + Layer.merge(TestEntityState.layer) + )))) + + it.effect("falls back to polling when nextDeliverAt fails", () => + Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + const scheduledFiber = yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 2000) + }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + + yield* TestClock.adjust(1999) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + assert.isUndefined(scheduledFiber.pollUnsafe()) + + yield* TestClock.adjust(3000) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + const scheduledExit = scheduledFiber.pollUnsafe() + assert(scheduledExit && Exit.isSuccess(scheduledExit)) + + const immediateFiber = yield* client.GetUser({ id: 2 }).pipe(Effect.forkChild) + yield* TestClock.adjust(1) + const user = yield* Fiber.join(immediateFiber) + assert.deepStrictEqual(user, new User({ id: 2, name: "User 2" })) + }).pipe(Effect.provide(TestShardingWithoutStorage.pipe( + Layer.updateService(MessageStorage.MessageStorage, (storage) => ({ + ...storage, + nextDeliverAt: () => + Effect.fail(new ClusterError.PersistenceError({ cause: new Error("nextDeliverAt failed") })) + })), + Layer.provide(MessageStorage.layerMemory), + Layer.provide(TestShardingConfig) + )))) + it.effect("fails volatile requests immediately when the mailbox is full", () => Effect.gen(function*() { const makeClient = yield* TestEntity.client diff --git a/packages/effect/test/cluster/TestEntity.ts b/packages/effect/test/cluster/TestEntity.ts index 8341f60f344..3f8e9a22072 100644 --- a/packages/effect/test/cluster/TestEntity.ts +++ b/packages/effect/test/cluster/TestEntity.ts @@ -1,6 +1,6 @@ import { type Cause, Context, Effect, Layer, MutableRef, Option, Queue, Schedule, Schema, Stream } from "effect" import type { Envelope } from "effect/unstable/cluster" -import { ClusterSchema, Entity } from "effect/unstable/cluster" +import { ClusterSchema, DeliverAt, Entity } from "effect/unstable/cluster" import { MemoryTransaction } from "effect/unstable/cluster/MessageStorage" import type { RpcGroup } from "effect/unstable/rpc" import { Rpc, RpcSchema } from "effect/unstable/rpc" @@ -10,6 +10,15 @@ export class User extends Schema.Class("User")({ name: Schema.String }) {} +export class ScheduledPayload extends Schema.Class("ScheduledPayload")({ + id: Schema.Number, + deliverAt: Schema.DateTimeUtcFromMillis +}) { + [DeliverAt.symbol]() { + return this.deliverAt + } +} + export class StreamWithKey extends Rpc.make("StreamWithKey", { success: RpcSchema.Stream(Schema.Number, Schema.Never), payload: { key: Schema.String }, @@ -25,6 +34,10 @@ export const TestEntity = Entity.make("TestEntity", [ success: User, payload: { id: Schema.Number } }).annotate(ClusterSchema.Persisted, false), + Rpc.make("GetUserScheduled", { + success: User, + payload: ScheduledPayload + }), Rpc.make("Never"), Rpc.make("NeverFork"), Rpc.make("NeverVolatile").annotate(ClusterSchema.Persisted, false), @@ -101,6 +114,11 @@ export const TestEntityNoState = TestEntity.toLayer( Queue.offerUnsafe(state.envelopes, envelope) return new User({ id: envelope.payload.id, name: `User ${envelope.payload.id}` }) }), + GetUserScheduled: (envelope) => + Effect.sync(() => { + Queue.offerUnsafe(state.envelopes, envelope) + return new User({ id: envelope.payload.id, name: `User ${envelope.payload.id}` }) + }), Never: never, NeverFork: (envelope) => Rpc.fork(never(envelope)), NeverVolatile: never, From 461d01443d2b6ed74a95e063f15c83b765aa9f7b Mon Sep 17 00:00:00 2001 From: Stephen Brian King Date: Fri, 17 Jul 2026 02:10:27 -0600 Subject: [PATCH 3/5] Cast SQL Server shard values to VARCHAR in the deliver-at lookup The NVARCHAR literals cause SQL Server to convert the indexed VARCHAR `shard_id` column, preventing the per-shard index seek. --- packages/effect/src/unstable/cluster/SqlMessageStorage.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts index 5c0eea9517e..80893ae89df 100644 --- a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts +++ b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts @@ -420,10 +420,13 @@ export const make: (options?: { ) AS t ` }, + // Keep the VALUES column as VARCHAR so SQL Server can seek the shard_id index. mssql: () => (shardIds: ReadonlyArray, now: number) => sql<{ readonly next_deliver_at: bigint | null }>` SELECT MIN(l.next_at) AS next_deliver_at - FROM (VALUES ${sql.literal(shardIds.map((id) => `(${wrapString(id)})`).join(","))}) AS s(shard_id) + FROM (VALUES ${ + sql.literal(shardIds.map((id) => `(CAST(${wrapString(id)} AS VARCHAR(50)))`).join(",")) + }) AS s(shard_id) CROSS APPLY ( SELECT TOP 1 deliver_at AS next_at FROM ${messagesTableSql} From bb0e8802b131000b3dfd5d9944afc21073eb06d5 Mon Sep 17 00:00:00 2001 From: Stephen Brian King Date: Fri, 17 Jul 2026 20:22:06 -0600 Subject: [PATCH 4/5] Bind shard values in deliver-at lookups and arm the wake before reading Shard identifiers now reach the nextDeliverAt queries as bound parameters in all three dialect branches, with SQL Server keeping its VARCHAR cast. The storage read loop discovers the next scheduled deadline before reading due messages, so a deadline crossing between the two reads arms a wake instead of waiting for the next poll interval. Regression tests cover a quote-bearing shard group and a deadline that lands between the reads. --- .../effect/src/unstable/cluster/Sharding.ts | 2 +- .../src/unstable/cluster/SqlMessageStorage.ts | 12 ++-- packages/effect/test/cluster/Sharding.test.ts | 67 ++++++++++++++++++- .../test/cluster/SqlMessageStorage.test.ts | 19 ++++++ 4 files changed, 90 insertions(+), 10 deletions(-) diff --git a/packages/effect/src/unstable/cluster/Sharding.ts b/packages/effect/src/unstable/cluster/Sharding.ts index f0f4d852a66..fc86031f617 100644 --- a/packages/effect/src/unstable/cluster/Sharding.ts +++ b/packages/effect/src/unstable/cluster/Sharding.ts @@ -598,10 +598,10 @@ const make = Effect.gen(function*() { pendingNotifications.forEach((entry) => removableNotifications.add(entry)) } + yield* rediscoverStorageWake messages = yield* storage.unprocessedMessages(acquiredShards) index = 0 yield* processMessages - yield* rediscoverStorageWake if (removableNotifications.size > 0) { removableNotifications.forEach(({ message, resume }) => { diff --git a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts index 80893ae89df..46e7b5733d4 100644 --- a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts +++ b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts @@ -405,10 +405,8 @@ export const make: (options?: { const getNextDeliverAt = sql.onDialectOrElse({ mysql: () => (shardIds: ReadonlyArray, now: number) => { - const shards = sql.literal( - shardIds.map((id, i) => i === 0 ? `SELECT ${wrapString(id)} AS shard_id` : `SELECT ${wrapString(id)}`).join( - " UNION ALL " - ) + const shards = sql.join(" UNION ALL ", false)( + shardIds.map((id, i) => i === 0 ? sql`SELECT ${id} AS shard_id` : sql`SELECT ${id}`) ) return sql<{ readonly next_deliver_at: bigint | null }>` SELECT MIN(t.next_at) AS next_deliver_at FROM ( @@ -424,9 +422,7 @@ export const make: (options?: { mssql: () => (shardIds: ReadonlyArray, now: number) => sql<{ readonly next_deliver_at: bigint | null }>` SELECT MIN(l.next_at) AS next_deliver_at - FROM (VALUES ${ - sql.literal(shardIds.map((id) => `(CAST(${wrapString(id)} AS VARCHAR(50)))`).join(",")) - }) AS s(shard_id) + FROM (VALUES ${sql.csv(shardIds.map((id) => sql`(CAST(${id} AS VARCHAR(50)))`))}) AS s(shard_id) CROSS APPLY ( SELECT TOP 1 deliver_at AS next_at FROM ${messagesTableSql} @@ -438,7 +434,7 @@ export const make: (options?: { sql<{ readonly next_deliver_at: bigint | null }>` SELECT MIN(deliver_at) AS next_deliver_at FROM ${messagesTableSql} - WHERE shard_id IN (${sql.literal(shardIds.map(wrapString).join(","))}) + WHERE ${sql.in("shard_id", shardIds)} AND processed = ${sqlFalse} AND deliver_at > ${sql.literal(String(now))} ` diff --git a/packages/effect/test/cluster/Sharding.test.ts b/packages/effect/test/cluster/Sharding.test.ts index 4ca68006d7c..2ae29e2a4ae 100644 --- a/packages/effect/test/cluster/Sharding.test.ts +++ b/packages/effect/test/cluster/Sharding.test.ts @@ -1,5 +1,19 @@ import { assert, describe, expect, it } from "@effect/vitest" -import { Array, Cause, Clock, DateTime, Effect, Exit, Fiber, Layer, MutableRef, Option, Queue, Stream } from "effect" +import { + Array, + Cause, + Clock, + DateTime, + Deferred, + Effect, + Exit, + Fiber, + Layer, + MutableRef, + Option, + Queue, + Stream +} from "effect" import { TestClock } from "effect/testing" import { ClusterError, @@ -213,6 +227,57 @@ describe.concurrent("Sharding", () => { assert(exit && Exit.isSuccess(exit)) }).pipe(Effect.provide(TestSharding))) + it.effect("arms a deadline wake before reading due messages", () => + Effect.gen(function*() { + const blockedReadComplete = yield* Deferred.make() + const releaseBlockedRead = yield* Deferred.make() + let blockNextReadAfterQuery = false + + yield* Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const now = yield* Clock.currentTimeMillis + blockNextReadAfterQuery = true + const fiber = yield* client.GetUserScheduled({ + id: 1, + deliverAt: DateTime.makeUnsafe(now + 200) + }).pipe(Effect.forkChild) + + yield* Deferred.await(blockedReadComplete) + yield* TestClock.adjust(200) + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 0) + + yield* Deferred.succeed(releaseBlockedRead, undefined) + yield* TestClock.adjust(1) + + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + const exit = fiber.pollUnsafe() + assert(exit && Exit.isSuccess(exit)) + }).pipe(Effect.provide(TestShardingWithoutStorage.pipe( + Layer.updateService(MessageStorage.MessageStorage, (storage) => ({ + ...storage, + unprocessedMessages(shardIds) { + return storage.unprocessedMessages(shardIds).pipe( + Effect.tap(() => { + if (!blockNextReadAfterQuery) { + return Effect.void + } + blockNextReadAfterQuery = false + return Deferred.succeed(blockedReadComplete, undefined).pipe( + Effect.andThen(Deferred.await(releaseBlockedRead)) + ) + }) + ) + } + })), + Layer.provide(MessageStorage.layerMemory), + Layer.provide(TestShardingConfig) + ))) + })) + it.effect("wakes for each scheduled deadline independently", () => Effect.gen(function*() { const state = yield* TestEntityState diff --git a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts index b7421db1d39..c32362695f7 100644 --- a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts +++ b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts @@ -202,6 +202,25 @@ describe("SqlMessageStorage", () => { assert.deepStrictEqual(next, Option.some(later)) })) + it.effect("nextDeliverAt binds quote-bearing shard groups", () => + Effect.gen(function*() { + yield* truncate + + const storage = yield* MessageStorage.MessageStorage + const shard = ShardId.make("quote's-group", 1) + const now = yield* Clock.currentTimeMillis + const deadline = now + 60_000 + const request = yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(deadline) }), + shardId: shard + }) + yield* storage.saveRequest(request) + + const next = yield* storage.nextDeliverAt([shard]) + expect(next).toEqual(Option.some(deadline)) + })) + it.effect("repliesFor", () => Effect.gen(function*() { yield* truncate From 536f197eff015dd7fae224c4326940bc1f88e0b3 Mon Sep 17 00:00:00 2001 From: Stephen Brian King <3913213+sbking@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:04:22 -0600 Subject: [PATCH 5/5] Address scheduled delivery review feedback --- ...cluster-message-storage-next-deliver-at.md | 2 +- .../src/unstable/cluster/MessageStorage.ts | 22 +++++++++--- .../effect/src/unstable/cluster/Sharding.ts | 17 ++++++---- .../test/cluster/MessageStorage.test.ts | 30 ++++++++++++++-- packages/effect/test/cluster/Sharding.test.ts | 34 +++++++++++++++++++ .../test/cluster/SqlMessageStorage.test.ts | 10 +++--- 6 files changed, 96 insertions(+), 19 deletions(-) diff --git a/.changeset/cluster-message-storage-next-deliver-at.md b/.changeset/cluster-message-storage-next-deliver-at.md index 1e25a94082b..4a81d9bcfb3 100644 --- a/.changeset/cluster-message-storage-next-deliver-at.md +++ b/.changeset/cluster-message-storage-next-deliver-at.md @@ -2,4 +2,4 @@ "effect": patch --- -Add a required `MessageStorage.nextDeliverAt` operation that returns the earliest future scheduled delivery time for a set of shards, with memory and SQL driver implementations and a supporting index migration for all SQL dialects. +Add a required `MessageStorage.nextDeliverAt` operation that returns the delay until the earliest future scheduled delivery for a set of shards, with memory and SQL driver implementations, background wake discovery, and a supporting index migration for all SQL dialects. diff --git a/packages/effect/src/unstable/cluster/MessageStorage.ts b/packages/effect/src/unstable/cluster/MessageStorage.ts index 47897e8ecf5..4b18fe9d92e 100644 --- a/packages/effect/src/unstable/cluster/MessageStorage.ts +++ b/packages/effect/src/unstable/cluster/MessageStorage.ts @@ -14,6 +14,7 @@ import * as Arr from "../../Array.ts" import { Clock } from "../../Clock.ts" import * as Context from "../../Context.ts" import * as Data from "../../Data.ts" +import * as Duration from "../../Duration.ts" import * as Effect from "../../Effect.ts" import * as Exit from "../../Exit.ts" import { constFalse, identity } from "../../Function.ts" @@ -144,12 +145,15 @@ export class MessageStorage extends Context.Service Effect.Effect>, PersistenceError> /** - * Retrieves the earliest `deliverAt` of the scheduled messages for the - * specified shards that are not yet deliverable. + * Retrieves the delay until the earliest scheduled message for the + * specified shards becomes deliverable. + * + * The delay must measure when this implementation's `unprocessedMessages` + * predicate will consider the message due. */ readonly nextDeliverAt: ( shardIds: Iterable - ) => Effect.Effect, PersistenceError> + ) => Effect.Effect, PersistenceError> /** * Reset the mailbox state for the provided shards. @@ -386,6 +390,9 @@ export type Encoded = { /** * Retrieves the earliest `deliverAt` of the scheduled messages for the * given shards that is later than `now`. + * + * The returned time must use the same clock domain as this implementation's + * `unprocessedMessages` eligibility predicate. */ readonly nextDeliverAt: ( shardIds: Arr.NonEmptyArray, @@ -671,7 +678,14 @@ export const makeEncoded: (encoded: Encoded) => Effect.Effect< nextDeliverAt(shardIds) { const shards = Array.from(shardIds, (id) => id.toString()) if (!Arr.isArrayNonEmpty(shards)) return Effect.succeedNone - return Effect.suspend(() => encoded.nextDeliverAt(shards, clock.currentTimeMillisUnsafe())) + return Effect.suspend(() => encoded.nextDeliverAt(shards, clock.currentTimeMillisUnsafe())).pipe( + Effect.map((next) => + Option.map( + next, + (deliverAt) => Duration.max(Duration.millis(deliverAt - clock.currentTimeMillisUnsafe()), Duration.zero) + ) + ) + ) }, resetAddress: encoded.resetAddress, clearAddress: encoded.clearAddress, diff --git a/packages/effect/src/unstable/cluster/Sharding.ts b/packages/effect/src/unstable/cluster/Sharding.ts index fc86031f617..7197790faff 100644 --- a/packages/effect/src/unstable/cluster/Sharding.ts +++ b/packages/effect/src/unstable/cluster/Sharding.ts @@ -478,6 +478,9 @@ const make = Effect.gen(function*() { // every read, so a missed or failed wake only delays delivery until the // next poll interval. const storageWakeHandle = yield* FiberHandle.make() + // Keep discovery separate so replacing a stale lookup does not cancel the + // last successfully armed deadline. + const storageWakeDiscoveryHandle = yield* FiberHandle.make() let index = 0 let messages: Array> = [] @@ -565,13 +568,15 @@ const make = Effect.gen(function*() { const rediscoverStorageWake = Effect.suspend(() => storage.nextDeliverAt(acquiredShards)).pipe( Effect.flatMap(Option.match({ onNone: () => FiberHandle.clear(storageWakeHandle), - onSome: (deliverAt) => { - const remainingMillis = deliverAt - clock.currentTimeMillisUnsafe() - if (remainingMillis <= 0) { - return Effect.asVoid(storageReadLatch.open) + onSome: (delay) => { + if (Duration.isLessThanOrEqualTo(delay, Duration.zero)) { + return FiberHandle.clear(storageWakeHandle).pipe( + Effect.andThen(storageReadLatch.open), + Effect.asVoid + ) } return storageReadLatch.open.pipe( - Effect.delay(Duration.millis(remainingMillis)), + Effect.delay(delay), FiberHandle.run(storageWakeHandle), Effect.asVoid ) @@ -598,7 +603,7 @@ const make = Effect.gen(function*() { pendingNotifications.forEach((entry) => removableNotifications.add(entry)) } - yield* rediscoverStorageWake + yield* FiberHandle.run(storageWakeDiscoveryHandle, rediscoverStorageWake) messages = yield* storage.unprocessedMessages(acquiredShards) index = 0 yield* processMessages diff --git a/packages/effect/test/cluster/MessageStorage.test.ts b/packages/effect/test/cluster/MessageStorage.test.ts index 33a2326a9a9..fd1abf2ab87 100644 --- a/packages/effect/test/cluster/MessageStorage.test.ts +++ b/packages/effect/test/cluster/MessageStorage.test.ts @@ -1,5 +1,5 @@ import { assert, describe, expect, it } from "@effect/vitest" -import { Clock, Context, DateTime, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" +import { Clock, Context, DateTime, Duration, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" import { TestClock } from "effect/testing" import { DeliverAt, @@ -93,7 +93,7 @@ describe("MessageStorage", () => { ) const next = yield* storage.nextDeliverAt([shard1, shard2]) - assert.deepStrictEqual(next, Option.some(earlier)) + assert.deepStrictEqual(next, Option.some(Duration.millis(60_000))) }).pipe(Effect.provide(MemoryLive))) it.effect("nextDeliverAt handles epoch zero across the due boundary", () => @@ -110,7 +110,7 @@ describe("MessageStorage", () => { ) let next = yield* storage.nextDeliverAt([shard]) - assert.deepStrictEqual(next, Option.some(0)) + assert.deepStrictEqual(next, Option.some(Duration.millis(1))) let messages = yield* storage.unprocessedMessages([shard]) assert.strictEqual(messages.length, 0) @@ -121,6 +121,30 @@ describe("MessageStorage", () => { assert.strictEqual(messages.length, 1) }).pipe(Effect.provide(MemoryLive))) + it.effect("nextDeliverAt accounts for lookup latency", () => + Effect.gen(function*() { + yield* TestClock.setTime(0) + const driver = yield* MessageStorage.MemoryDriver + const storage = yield* MessageStorage.makeEncoded({ + ...driver.encoded, + nextDeliverAt: (shardIds, now) => + driver.encoded.nextDeliverAt(shardIds, now).pipe(Effect.delay(Duration.seconds(2))) + }) + const shard = ShardId.make("default", 1) + yield* storage.saveRequest( + yield* makeRequest({ + rpc: ScheduledRpc, + payload: new ScheduledPayload({ id: 1, deliverAt: DateTime.makeUnsafe(1000) }), + shardId: shard + }) + ) + + const fiber = yield* storage.nextDeliverAt([shard]).pipe(Effect.forkChild) + yield* TestClock.adjust(2000) + + assert.deepStrictEqual(yield* Fiber.join(fiber), Option.some(Duration.zero)) + }).pipe(Effect.provide(MemoryLive))) + it.effect("nextDeliverAt ignores messages on other shards", () => Effect.gen(function*() { const storage = yield* MessageStorage.MessageStorage diff --git a/packages/effect/test/cluster/Sharding.test.ts b/packages/effect/test/cluster/Sharding.test.ts index 2ae29e2a4ae..ef8a231acaa 100644 --- a/packages/effect/test/cluster/Sharding.test.ts +++ b/packages/effect/test/cluster/Sharding.test.ts @@ -227,6 +227,40 @@ describe.concurrent("Sharding", () => { assert(exit && Exit.isSuccess(exit)) }).pipe(Effect.provide(TestSharding))) + it.effect("reads due messages while discovering the next deadline in the background", () => + Effect.gen(function*() { + const queryStarted = yield* Deferred.make() + const releaseQuery = yield* Deferred.make() + + yield* Effect.gen(function*() { + const state = yield* TestEntityState + const makeClient = yield* TestEntity.client + yield* TestClock.adjust(1) + const client = makeClient("1") + + const fiber = yield* client.GetUser({ id: 1 }).pipe(Effect.forkChild) + yield* Deferred.await(queryStarted) + yield* TestClock.adjust(1) + + assert.strictEqual(Queue.sizeUnsafe(state.envelopes), 1) + const exit = fiber.pollUnsafe() + assert(exit && Exit.isSuccess(exit)) + yield* Deferred.succeed(releaseQuery, undefined) + }).pipe(Effect.provide(TestShardingWithoutStorage.pipe( + Layer.updateService(MessageStorage.MessageStorage, (storage) => ({ + ...storage, + nextDeliverAt(shardIds) { + return Deferred.succeed(queryStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseQuery)), + Effect.andThen(storage.nextDeliverAt(shardIds)) + ) + } + })), + Layer.provide(MessageStorage.layerMemory), + Layer.provide(TestShardingConfig) + ))) + })) + it.effect("arms a deadline wake before reading due messages", () => Effect.gen(function*() { const blockedReadComplete = yield* Deferred.make() diff --git a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts index c32362695f7..e50682ee99e 100644 --- a/packages/platform-node/test/cluster/SqlMessageStorage.test.ts +++ b/packages/platform-node/test/cluster/SqlMessageStorage.test.ts @@ -1,7 +1,7 @@ import { NodeFileSystem } from "@effect/platform-node" import { SqliteClient } from "@effect/sql-sqlite-node" import { assert, describe, expect, it } from "@effect/vitest" -import { Clock, DateTime, Effect, Fiber, FileSystem, Latch, Layer, Option } from "effect" +import { Clock, DateTime, Duration, Effect, Fiber, FileSystem, Latch, Layer, Option } from "effect" import { TestClock } from "effect/testing" import { Message, MessageStorage, ShardId, ShardingConfig, Snowflake, SqlMessageStorage } from "effect/unstable/cluster" import { SqlClient } from "effect/unstable/sql" @@ -192,14 +192,14 @@ describe("SqlMessageStorage", () => { yield* storage.saveRequest(laterRequest) next = yield* storage.nextDeliverAt([shard1, shard2]) - assert.deepStrictEqual(next, Option.some(earlier)) + assert.deepStrictEqual(next, Option.some(Duration.minutes(1))) next = yield* storage.nextDeliverAt([shard2]) - assert.deepStrictEqual(next, Option.some(later)) + assert.deepStrictEqual(next, Option.some(Duration.minutes(2))) yield* storage.saveReply(yield* makeReply(earlierRequest)) next = yield* storage.nextDeliverAt([shard1, shard2]) - assert.deepStrictEqual(next, Option.some(later)) + assert.deepStrictEqual(next, Option.some(Duration.minutes(2))) })) it.effect("nextDeliverAt binds quote-bearing shard groups", () => @@ -218,7 +218,7 @@ describe("SqlMessageStorage", () => { yield* storage.saveRequest(request) const next = yield* storage.nextDeliverAt([shard]) - expect(next).toEqual(Option.some(deadline)) + assert.deepStrictEqual(next, Option.some(Duration.minutes(1))) })) it.effect("repliesFor", () =>