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..4a81d9bcfb3 --- /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 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 40a6ad37e2c..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" @@ -143,6 +144,17 @@ export class MessageStorage extends Context.Service ) => Effect.Effect>, PersistenceError> + /** + * 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> + /** * Reset the mailbox state for the provided shards. */ @@ -375,6 +387,18 @@ export type Encoded = { PersistenceError > + /** + * 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, + now: number + ) => Effect.Effect, PersistenceError> + /** * Reset the mailbox state for the provided address. */ @@ -651,6 +675,18 @@ 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())).pipe( + Effect.map((next) => + Option.map( + next, + (deliverAt) => Duration.max(Duration.millis(deliverAt - clock.currentTimeMillisUnsafe()), Duration.zero) + ) + ) + ) + }, resetAddress: encoded.resetAddress, clearAddress: encoded.clearAddress, resetShards: (shardIds) => { @@ -779,6 +815,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 +886,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 +1004,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 +1029,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/Sharding.ts b/packages/effect/src/unstable/cluster/Sharding.ts index a202047cf87..7197790faff 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,15 @@ 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() + // 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> = [] const removableNotifications = new Set() @@ -555,6 +565,26 @@ const make = Effect.gen(function*() { } ) + const rediscoverStorageWake = Effect.suspend(() => storage.nextDeliverAt(acquiredShards)).pipe( + Effect.flatMap(Option.match({ + onNone: () => FiberHandle.clear(storageWakeHandle), + 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(delay), + 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 @@ -573,6 +603,7 @@ const make = Effect.gen(function*() { pendingNotifications.forEach((entry) => removableNotifications.add(entry)) } + yield* FiberHandle.run(storageWakeDiscoveryHandle, rediscoverStorageWake) messages = yield* storage.unprocessedMessages(acquiredShards) index = 0 yield* processMessages diff --git a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts index 63c1c38c480..46e7b5733d4 100644 --- a/packages/effect/src/unstable/cluster/SqlMessageStorage.ts +++ b/packages/effect/src/unstable/cluster/SqlMessageStorage.ts @@ -403,6 +403,43 @@ export const make: (options?: { ) }) + const getNextDeliverAt = sql.onDialectOrElse({ + mysql: () => (shardIds: ReadonlyArray, now: number) => { + 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 ( + 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 + ` + }, + // 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.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} + 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 ${sql.in("shard_id", shardIds)} + AND processed = ${sqlFalse} + AND deliver_at > ${sql.literal(String(now))} + ` + }) + return yield* MessageStorage.makeEncoded({ saveEnvelope: ({ deliverAt, envelope, primaryKey }) => Effect.suspend(() => { @@ -577,6 +614,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 +1025,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..fd1abf2ab87 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, Duration, Effect, Exit, Fiber, Latch, Layer, Option, Schema } from "effect" import { TestClock } from "effect/testing" import { + DeliverAt, EntityAddress, EntityId, EntityType, @@ -61,6 +62,159 @@ 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(Duration.millis(60_000))) + }).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(Duration.millis(1))) + 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 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 + 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 +255,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 +263,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 +291,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/effect/test/cluster/Sharding.test.ts b/packages/effect/test/cluster/Sharding.test.ts index 2a526bbc692..ef8a231acaa 100644 --- a/packages/effect/test/cluster/Sharding.test.ts +++ b/packages/effect/test/cluster/Sharding.test.ts @@ -1,7 +1,22 @@ 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, + Deferred, + Effect, + Exit, + Fiber, + Layer, + MutableRef, + Option, + Queue, + Stream +} from "effect" import { TestClock } from "effect/testing" import { + ClusterError, MessageStorage, RunnerAddress, RunnerHealth, @@ -183,6 +198,274 @@ 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("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() + 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 + 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, 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..e50682ee99e 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, Duration, 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,62 @@ 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(Duration.minutes(1))) + + next = yield* storage.nextDeliverAt([shard2]) + 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(Duration.minutes(2))) + })) + + 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]) + assert.deepStrictEqual(next, Option.some(Duration.minutes(1))) + })) + it.effect("repliesFor", () => Effect.gen(function*() { yield* truncate