Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/cluster-message-storage-next-deliver-at.md
Original file line number Diff line number Diff line change
@@ -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.
54 changes: 52 additions & 2 deletions packages/effect/src/unstable/cluster/MessageStorage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -143,6 +144,17 @@ export class MessageStorage extends Context.Service<MessageStorage, {
messageIds: Iterable<Snowflake.Snowflake>
) => Effect.Effect<Array<Message.Incoming<R>>, 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<ShardId.ShardId>
) => Effect.Effect<Option.Option<Duration.Duration>, PersistenceError>

/**
* Reset the mailbox state for the provided shards.
*/
Expand Down Expand Up @@ -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<string>,
now: number
) => Effect.Effect<Option.Option<number>, PersistenceError>

/**
* Reset the mailbox state for the provided address.
*/
Expand Down Expand Up @@ -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) => {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -849,7 +886,7 @@ export class MemoryDriver extends Context.Service<MemoryDriver>()("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({
Expand Down Expand Up @@ -967,7 +1004,7 @@ export class MemoryDriver extends Context.Service<MemoryDriver>()("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({
Expand All @@ -992,6 +1029,19 @@ export class MemoryDriver extends Context.Service<MemoryDriver>()("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(() => {
Expand Down
31 changes: 31 additions & 0 deletions packages/effect/src/unstable/cluster/Sharding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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<Message.Incoming<any>> = []
const removableNotifications = new Set<PendingNotification>()
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
89 changes: 89 additions & 0 deletions packages/effect/src/unstable/cluster/SqlMessageStorage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,43 @@ export const make: (options?: {
)
})

const getNextDeliverAt = sql.onDialectOrElse({
mysql: () => (shardIds: ReadonlyArray<string>, 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<string>, 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<string>, 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))}
`
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.

return yield* MessageStorage.makeEncoded({
saveEnvelope: ({ deliverAt, envelope, primaryKey }) =>
Effect.suspend(() => {
Expand Down Expand Up @@ -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}
Expand Down Expand Up @@ -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
`
})
})
})
}
Expand Down
Loading
Loading