From 1885db3d8b27265631a83d3d3fe6b3f40f29f395 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 14 Feb 2026 10:37:47 -0500 Subject: [PATCH] discord: simplify ConversationLedger to offsets + in-memory dedup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the full inbox state machine (admit/start/setTarget/setPrompt/ setResponse/complete/retry/prune/replayPending) with three methods: dedup (bounded in-memory Set), getOffset, and setOffset. Discord is already the durable inbox — on startup we just resume from persisted offsets per source. Within a session, in-memory dedup prevents double-processing. This removes ~565 lines of inbox lifecycle code, the conversation_inbox SQLite table, MessageState, ReliabilityError, and the prune schedule. --- packages/discord/src/conversation/README.md | 5 +- .../implementations/discord/index.ts | 26 +- .../discord/src/conversation/model/errors.ts | 10 - .../services/conversation.test.ts | 141 +--------- .../src/conversation/services/conversation.ts | 248 +++++++----------- .../src/conversation/services/ledger.test.ts | 133 +++------- .../src/conversation/services/ledger.ts | 245 ++--------------- packages/discord/src/db/init.test.ts | 39 --- .../db/migrations/0001_discord_sessions.ts | 48 ---- 9 files changed, 165 insertions(+), 730 deletions(-) diff --git a/packages/discord/src/conversation/README.md b/packages/discord/src/conversation/README.md index cde6733619..73d2522dc2 100644 --- a/packages/discord/src/conversation/README.md +++ b/packages/discord/src/conversation/README.md @@ -21,10 +21,9 @@ This module is wired into `src/index.ts`. Reliability semantics: -- inbound events are durably admitted by `message_id` before processing -- pending events replay on startup +- in-memory dedup by `message_id` prevents double-processing within a session - startup catch-up fetches missed Discord messages from tracked thread sources and allowed channels using persisted offsets -- response text is cached before Discord delivery so retries can re-publish without re-calling the model +- Discord is the durable inbox — on startup we resume from where we left off per source Local CLI notes (`bun run conversation:cli`): diff --git a/packages/discord/src/conversation/implementations/discord/index.ts b/packages/discord/src/conversation/implementations/discord/index.ts index 15d22e0604..2eda47fe1f 100644 --- a/packages/discord/src/conversation/implementations/discord/index.ts +++ b/packages/discord/src/conversation/implementations/discord/index.ts @@ -193,7 +193,7 @@ export class DiscordConversationServices { const uniq = (values: ReadonlyArray): Array => [...new Set(values)] const offer = (event: Inbound, onFresh: Effect.Effect) => - ledger.admit(event).pipe( + ledger.dedup(event.message_id).pipe( Effect.flatMap((fresh) => { if (!fresh) { return Effect.logDebug("Message deduped (already seen)").pipe( @@ -543,7 +543,7 @@ export class DiscordConversationServices { }), content: text, }) - const ingest = ledger.admit(event).pipe( + const ingest = ledger.dedup(event.message_id).pipe( Effect.flatMap((fresh) => { if (!fresh) return Effect.void return Effect.sync(() => { @@ -617,28 +617,6 @@ export class DiscordConversationServices { )), ) - yield* ledger.replayPending().pipe( - Effect.flatMap((events) => - Effect.forEach( - events, - (event) => - Effect.sync(() => { - input.unsafeOffer(event) - }), - { discard: true }, - ).pipe( - Effect.zipRight( - Effect.logInfo("Replayed pending conversation events").pipe( - Effect.annotateLogs({ event: "conversation.ledger.replay", count: events.length }), - ), - ), - )), - Effect.catchAll((error) => - Effect.logError("Failed replaying pending conversation events").pipe( - Effect.annotateLogs({ event: "conversation.ledger.replay.failed", error: messageOf(error) }), - )), - ) - const inbox = Inbox.of({ events: Stream.fromQueue(input, { shutdown: false }), }) diff --git a/packages/discord/src/conversation/model/errors.ts b/packages/discord/src/conversation/model/errors.ts index f2debaf029..0fe36c87c5 100644 --- a/packages/discord/src/conversation/model/errors.ts +++ b/packages/discord/src/conversation/model/errors.ts @@ -48,22 +48,12 @@ export class SandboxSendError extends Schema.TaggedError()( }, ) {} -export class ReliabilityError extends Schema.TaggedError()( - "ReliabilityError", - { - message_id: Schema.String, - message: Schema.String, - retriable: Schema.Boolean, - }, -) {} - export const ConversationError = Schema.Union( ThreadEnsureError, HistoryError, DeliveryError, RoutingError, SandboxSendError, - ReliabilityError, ) export type ConversationError = typeof ConversationError.Type diff --git a/packages/discord/src/conversation/services/conversation.test.ts b/packages/discord/src/conversation/services/conversation.test.ts index 9f0cf321e6..3081c8c847 100644 --- a/packages/discord/src/conversation/services/conversation.test.ts +++ b/packages/discord/src/conversation/services/conversation.test.ts @@ -8,7 +8,7 @@ import { ChannelId, GuildId, SandboxId, SessionId, SessionInfo, ThreadId } from import { Mention, ThreadMessage, ThreadRef, Typing, type Action, type Inbound } from "../model/schema" import { History } from "./history" import { Inbox } from "./inbox" -import { ConversationLedger, MessageState } from "./ledger" +import { ConversationLedger } from "./ledger" import { Outbox } from "./outbox" import { Threads } from "./threads" import { Conversation } from "./conversation" @@ -447,55 +447,24 @@ describe("Conversation", () => { // --- Duplicate processing tests --- -/** A ledger that tracks admit/start calls and enforces dedup like the real one */ +/** A ledger that tracks dedup calls */ const makeTrackingLedger = () => { - const admitted = new Set() - const started = new Set() - const completed = new Set() - const admitCalls: Array = [] - const startCalls: Array = [] + const seen = new Set() + const dedupCalls: Array = [] const service: ConversationLedger.Service = { - admit: (event) => + dedup: (message_id) => Effect.sync(() => { - admitCalls.push(event.message_id) - if (admitted.has(event.message_id)) return false - admitted.add(event.message_id) + dedupCalls.push(message_id) + if (seen.has(message_id)) return false + seen.add(message_id) return true }), - replayPending: () => Effect.succeed([]), - start: (message_id) => - Effect.sync(() => { - startCalls.push(message_id) - if (started.has(message_id) || completed.has(message_id)) return Option.none() - started.add(message_id) - return Option.some(MessageState.make({ - thread_id: null, - channel_id: null, - response_text: null, - prompt_text: null, - session_id: null, - })) - }), - setTarget: () => Effect.void, - setPrompt: () => Effect.void, - setResponse: () => Effect.void, - complete: (message_id) => - Effect.sync(() => { - started.delete(message_id) - completed.add(message_id) - }), - retry: (message_id) => - Effect.sync(() => { - started.delete(message_id) - // Back to pending — NOT completed, so start will work again - }), - prune: () => Effect.void, getOffset: () => Effect.succeed(Option.none()), setOffset: () => Effect.void, } - return { service, admitted, started, completed, admitCalls, startCalls } + return { service, seen, dedupCalls } } const makeConversationLayerWithLedger = (props: { @@ -622,10 +591,8 @@ describe("Conversation duplicate processing", () => { const conversation = yield* Conversation yield* conversation.run - // The ledger should have been called twice with admit - expect(ledger.admitCalls).toEqual(["m1", "m1"]) - // But only one start should have succeeded - expect(ledger.startCalls.length).toBeLessThanOrEqual(2) + // The ledger should have been called twice with dedup + expect(ledger.dedupCalls).toEqual(["m1", "m1"]) // The agent should have received the prompt only once expect(prompts).toEqual(["hello"]) // Only one typing + one send @@ -633,12 +600,12 @@ describe("Conversation duplicate processing", () => { }).pipe(Effect.provide(live)) }) - effectTest("noop ledger now deduplicates (bug is fixed)", () => { + effectTest("noop ledger deduplicates", () => { const actions: Array = [] const prompts: Array = [] const event = makeEvent("hello") - // The noop ledger now tracks seen message_ids + // The noop ledger tracks seen message_ids const live = makeConversationLayer({ events: [event, event], tracked: Option.none(), @@ -687,88 +654,6 @@ describe("Conversation duplicate processing", () => { }).pipe(Effect.provide(live)) }) - effectTest("retry resets to pending and second run processes again", () => { - const actions: Array = [] - const prompts: Array = [] - const sendCount = { value: 0 } - - // A ledger that allows retry -> re-process flow - const admitted = new Set() - const state = new Map() - const ledgerService: ConversationLedger.Service = { - admit: (event) => - Effect.sync(() => { - if (admitted.has(event.message_id)) return false - admitted.add(event.message_id) - state.set(event.message_id, "pending") - return true - }), - replayPending: () => Effect.succeed([]), - start: (message_id) => - Effect.sync(() => { - if (state.get(message_id) !== "pending") return Option.none() - state.set(message_id, "processing") - return Option.some(MessageState.make({ - thread_id: null, - channel_id: null, - response_text: null, - prompt_text: null, - session_id: null, - })) - }), - setTarget: () => Effect.void, - setPrompt: () => Effect.void, - setResponse: () => Effect.void, - complete: (message_id) => - Effect.sync(() => { state.set(message_id, "completed") }), - retry: (message_id) => - Effect.sync(() => { state.set(message_id, "pending") }), - prune: () => Effect.void, - getOffset: () => Effect.succeed(Option.none()), - setOffset: () => Effect.void, - } - - const event = makeEvent("build it") - - const live = makeConversationLayerWithLedger({ - events: [], - tracked: Option.none(), - resolves: [makeSession("s1"), makeSession("s2")], - send: (_session, text) => { - sendCount.value += 1 - if (sendCount.value === 1) { - return Effect.fail( - SandboxDeadError.make({ - threadId: ThreadId.make("t1"), - reason: "dead", - }), - ) - } - return Effect.succeed(`ok:${text}`) - }, - rehydrate: (_threadId, latest) => Effect.succeed(`rehydrated:${latest}`), - actions, - prompts, - ledger: ledgerService, - }) - - return Effect.gen(function* () { - const conversation = yield* Conversation - // First turn fails (SandboxDeadError caught + recovery), so it processes once - yield* conversation.turn(event) - - // After the first turn, the ledger state should be completed (recovery succeeded inline) - // A second turn with the same event should be blocked by admit - yield* conversation.turn(event) - - // The second turn should NOT have sent the prompt again - // (admit returns false because the message was already admitted) - const sendActions = actions.filter((x) => x.kind === "send") - // The first turn should have: recovery message + successful reply - expect(sendActions.length).toBeGreaterThanOrEqual(1) - }).pipe(Effect.provide(live)) - }) - effectTest("two different messages on same thread are processed sequentially (not lost)", () => { const actions: Array = [] const prompts: Array = [] diff --git a/packages/discord/src/conversation/services/conversation.ts b/packages/discord/src/conversation/services/conversation.ts index d188664c99..b635d47066 100644 --- a/packages/discord/src/conversation/services/conversation.ts +++ b/packages/discord/src/conversation/services/conversation.ts @@ -4,7 +4,7 @@ import { TurnRouter } from "../../discord/turn-routing" import { ActorMap } from "../../lib/actors/keyed" import { ThreadAgentPool } from "../../sandbox/pool" import type { ChannelId, ThreadId } from "../../types" -import { type ConversationError, messageOf, ReliabilityError, RoutingError, SandboxSendError } from "../model/errors" +import { type ConversationError, messageOf, RoutingError, SandboxSendError } from "../model/errors" import { Send, type Inbound } from "../model/schema" import { History } from "./history" import { Inbox } from "./inbox" @@ -53,15 +53,7 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio SandboxSendError.make({ thread_id, message: messageOf(cause), - retriable: RETRIABLE_TAGS.has(cause._tag), - }) - - const asReliabilityError = (message_id: string) => - (cause: unknown): ReliabilityError => - ReliabilityError.make({ - message_id, - message: messageOf(cause), - retriable: true, + retriable: RETRIABLE_TAGS.has(cause._tag), }) const publishText = (threadId: ThreadId, text: string) => @@ -86,29 +78,31 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio if (event.mentions_everyone) return false if (!event.content.trim()) return false - const mentioned = event.mentions.user_ids.includes(event.bot_user_id) - || (event.bot_role_id.length > 0 && event.mentions.role_ids.includes(event.bot_role_id)) + const mentioned = + event.mentions.user_ids.includes(event.bot_user_id) || + (event.bot_role_id.length > 0 && event.mentions.role_ids.includes(event.bot_role_id)) if (event.kind === "channel_message") return mentioned if (mentioned) return true - const owned = yield* pool.hasTrackedThread(event.thread_id).pipe( - Effect.mapError(asSendError(event.thread_id)), - ) + const owned = yield* pool.hasTrackedThread(event.thread_id).pipe(Effect.mapError(asSendError(event.thread_id))) if (!owned) return false - const decision = yield* router.shouldRespond({ - content: event.content, - botUserId: event.bot_user_id, - botRoleId: event.bot_role_id, - mentionedUserIds: event.mentions.user_ids, - mentionedRoleIds: event.mentions.role_ids, - }).pipe( - Effect.mapError((cause) => - RoutingError.make({ - message: messageOf(cause), - retriable: false, - })), - ) + const decision = yield* router + .shouldRespond({ + content: event.content, + botUserId: event.bot_user_id, + botRoleId: event.bot_role_id, + mentionedUserIds: event.mentions.user_ids, + mentionedRoleIds: event.mentions.role_ids, + }) + .pipe( + Effect.mapError((cause) => + RoutingError.make({ + message: messageOf(cause), + retriable: false, + }), + ), + ) return decision.shouldRespond }) @@ -125,18 +119,15 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio target: { thread_id: ThreadId; channel_id: ChannelId }, ) { const toSendError = asSendError(target.thread_id) - const tracked = yield* pool.getTrackedSession(target.thread_id).pipe( - Effect.mapError(toSendError), - ) - const agent = yield* pool.getOrCreate(target.thread_id, target.channel_id, event.guild_id).pipe( - Effect.mapError(toSendError), - ) - const current = yield* agent.current().pipe( - Effect.mapError(toSendError), - ) - const prompt = Option.isSome(tracked) && tracked.value.sessionId !== current.sessionId - ? yield* history.rehydrate(target.thread_id, event.content) - : event.content + const tracked = yield* pool.getTrackedSession(target.thread_id).pipe(Effect.mapError(toSendError)) + const agent = yield* pool + .getOrCreate(target.thread_id, target.channel_id, event.guild_id) + .pipe(Effect.mapError(toSendError)) + const current = yield* agent.current().pipe(Effect.mapError(toSendError)) + const prompt = + Option.isSome(tracked) && tracked.value.sessionId !== current.sessionId + ? yield* history.rehydrate(target.thread_id, event.content) + : event.content return { target, agent, prompt, session: current } }) @@ -149,7 +140,9 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio return true } if (text === "!status") { - const tracked = yield* pool.getTrackedSession(target.thread_id).pipe(Effect.catchAll(() => Effect.succeed(Option.none()))) + const tracked = yield* pool + .getTrackedSession(target.thread_id) + .pipe(Effect.catchAll(() => Effect.succeed(Option.none()))) if (Option.isNone(tracked)) { yield* publishText(target.thread_id, "*No active session for this thread.*") } else { @@ -170,30 +163,10 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio return false }) - const turnRaw = Effect.fn("Conversation.turnRaw")(function* ( - event: Inbound, - state: { - thread_id: ThreadId | null - channel_id: ChannelId | null - response_text: string | null - prompt_text: string | null - session_id: string | null - }, - ) { + const turnRaw = Effect.fn("Conversation.turnRaw")(function* (event: Inbound) { if (!(yield* route(event))) return - const target = state.thread_id && state.channel_id - ? { thread_id: state.thread_id, channel_id: state.channel_id } - : yield* resolve(event) - - yield* ledger.setTarget(event.message_id, target.thread_id, target.channel_id).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - ) - - if (state.response_text) { - yield* publishText(target.thread_id, state.response_text) - return - } + const target = yield* resolve(event) if (yield* command(event, target)) return @@ -206,116 +179,75 @@ export class Conversation extends Context.Tag("@discord/conversation/Conversatio }), ) - yield* outbox.withTyping( - target.thread_id, - Effect.gen(function* () { - const input = yield* buildInput(event, target) - const reuse = state.prompt_text !== null - && state.session_id !== null - && state.session_id === input.session.sessionId - const prompt = reuse ? (state.prompt_text ?? input.prompt) : input.prompt - if (!reuse) { - yield* ledger.setPrompt(event.message_id, prompt, input.session.sessionId).pipe( - Effect.mapError(asReliabilityError(event.message_id)), + yield* outbox + .withTyping( + target.thread_id, + Effect.gen(function* () { + const input = yield* buildInput(event, target) + + const reply = yield* input.agent.send(input.prompt).pipe( + Effect.catchTag("SandboxDeadError", () => + Effect.gen(function* () { + yield* publishText(input.target.thread_id, "*Session changed state, recovering...*") + const toErr = asSendError(input.target.thread_id) + const next = yield* pool + .getOrCreate(input.target.thread_id, input.target.channel_id, event.guild_id) + .pipe(Effect.mapError(toErr)) + const nextSession = yield* next.current().pipe(Effect.mapError(toErr)) + const prompt = + nextSession.sessionId !== input.session.sessionId + ? yield* history.rehydrate(input.target.thread_id, event.content) + : event.content + return yield* next.send(prompt) + }), + ), + Effect.mapError(asSendError(input.target.thread_id)), ) - } else { - yield* Effect.logInfo("Recovered in-flight prompt from ledger").pipe( + + yield* Effect.logInfo("Bot reply").pipe( Effect.annotateLogs({ - event: "conversation.ledger.prompt.reused", - message_id: event.message_id, - thread_id: target.thread_id, + event: "conversation.bot.reply", + thread_id: input.target.thread_id, + content: reply.slice(0, 200), }), ) - } - - const reply = yield* input.agent.send(prompt).pipe( - Effect.catchTag("SandboxDeadError", () => - Effect.gen(function* () { - yield* publishText(input.target.thread_id, "*Session changed state, recovering...*") - const toErr = asSendError(input.target.thread_id) - const next = yield* pool.getOrCreate( - input.target.thread_id, - input.target.channel_id, - event.guild_id, - ).pipe(Effect.mapError(toErr)) - const nextSession = yield* next.current().pipe( - Effect.mapError(toErr), - ) - const prompt = nextSession.sessionId !== input.session.sessionId - ? yield* history.rehydrate(input.target.thread_id, event.content) - : event.content - return yield* next.send(prompt) - }), - ), - Effect.mapError(asSendError(input.target.thread_id)), - ) - - yield* Effect.logInfo("Bot reply").pipe( - Effect.annotateLogs({ - event: "conversation.bot.reply", - thread_id: input.target.thread_id, - content: reply.slice(0, 200), - }), - ) - yield* ledger.setResponse(event.message_id, reply).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - ) - yield* publishText(input.target.thread_id, reply) - }), - ).pipe( - Effect.catchAll(reportFailure(target.thread_id)), - ) + yield* publishText(input.target.thread_id, reply) + }), + ) + .pipe(Effect.catchAll(reportFailure(target.thread_id))) }) const keyOf = (event: Inbound) => - event.kind === "thread_message" - ? `thread:${event.thread_id}` - : `channel:${event.channel_id}` + event.kind === "thread_message" ? `thread:${event.thread_id}` : `channel:${event.channel_id}` - const runEvent = Effect.fn("Conversation.runEvent")(function* (event: Inbound) { - yield* ledger.admit(event).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - ) - const state = yield* ledger.start(event.message_id).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - ) - if (Option.isNone(state)) return - - yield* turnRaw(event, state.value).pipe( - Effect.tap(() => - ledger.complete(event.message_id).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - )), - Effect.catchAll((error) => - ledger.retry(event.message_id, messageOf(error).slice(0, 500)).pipe( - Effect.mapError(asReliabilityError(event.message_id)), - Effect.zipRight(Effect.fail(error)), - )), - ) - }) + const processEvent = (event: Inbound) => actors.run(keyOf(event), turnRaw(event), { touch: false }) const turn = Effect.fn("Conversation.turn")(function* (event: Inbound) { - yield* actors.run( - keyOf(event), - runEvent(event), - { touch: false }, - ) + const fresh = yield* ledger.dedup(event.message_id) + if (!fresh) return + yield* processEvent(event) }) const run = inbox.events.pipe( Stream.mapEffect( (event) => - turn(event).pipe( - Effect.retry(turnRetry), - Effect.catchAll((error) => - Effect.logError("Conversation turn failed").pipe( - Effect.annotateLogs({ - event: "conversation.turn.failed", - tag: error._tag, - retriable: error.retriable, - message: error.message, - }), - )), + ledger.dedup(event.message_id).pipe( + Effect.flatMap((fresh) => { + if (!fresh) return Effect.void + return processEvent(event).pipe( + Effect.retry(turnRetry), + Effect.catchAll((error) => + Effect.logError("Conversation turn failed").pipe( + Effect.annotateLogs({ + event: "conversation.turn.failed", + tag: error._tag, + retriable: error.retriable, + message: error.message, + }), + ), + ), + ) + }), ), { concurrency: "unbounded", unordered: true }, ), diff --git a/packages/discord/src/conversation/services/ledger.test.ts b/packages/discord/src/conversation/services/ledger.test.ts index 305a3073f5..be029f3fdd 100644 --- a/packages/discord/src/conversation/services/ledger.test.ts +++ b/packages/discord/src/conversation/services/ledger.test.ts @@ -1,12 +1,8 @@ -import * as Client from "@effect/sql/SqlClient" import { describe, expect } from "bun:test" import { Duration, Effect, Layer, Option, Redacted } from "effect" import { AppConfig } from "../../config" import { SqliteDb } from "../../db/client" -import { initializeSchema } from "../../db/init" import { effectTest, withTempSqliteFile } from "../../test/effect" -import { ChannelId, GuildId, SessionId, ThreadId } from "../../types" -import { Mention, ThreadMessage, type Inbound } from "../model/schema" import { ConversationLedger } from "./ledger" const makeConfig = (databasePath: string) => @@ -38,118 +34,30 @@ const makeConfig = (databasePath: string) => openCodeModel: "opencode/claude-sonnet-4-5", }) -const event = (message_id: string, content: string): Inbound => - ThreadMessage.make({ - kind: "thread_message", - thread_id: ThreadId.make("t1"), - channel_id: ChannelId.make("c1"), - message_id, - guild_id: GuildId.make("g1"), - bot_user_id: "bot-1", - bot_role_id: "role-1", - author_id: "u1", - author_is_bot: false, - mentions_everyone: false, - mentions: Mention.make({ user_ids: ["bot-1"], role_ids: [] }), - content, - }) - const withLedger = ( - run: (ledger: ConversationLedger.Service, sql: Client.SqlClient) => Effect.Effect, + run: (ledger: ConversationLedger.Service) => Effect.Effect, ) => withTempSqliteFile((databasePath) => Effect.gen(function* () { const config = Layer.succeed(AppConfig, makeConfig(databasePath)) const sqlite = SqliteDb.layer.pipe(Layer.provide(config)) const deps = Layer.merge(sqlite, config) - const live = Layer.merge( - ConversationLedger.layer.pipe(Layer.provide(deps)), - sqlite, - ) - const program = Effect.all([ConversationLedger, SqliteDb]).pipe( - Effect.flatMap(([ledger, sql]) => - initializeSchema.pipe( - Effect.provideService(Client.SqlClient, sql), - Effect.zipRight(run(ledger, sql)), - )), - ) + const live = ConversationLedger.layer.pipe(Layer.provide(deps)) + const program = Effect.flatMap(ConversationLedger, (ledger) => run(ledger)) return yield* program.pipe(Effect.provide(live)) }), "discord-ledger-", ) describe("ConversationLedger", () => { - effectTest("deduplicates by message id and tracks completion", () => - withLedger((ledger) => - Effect.gen(function* () { - const m = event("m1", "hello") - expect(yield* ledger.admit(m)).toBe(true) - expect(yield* ledger.admit(m)).toBe(false) - - const started = yield* ledger.start(m.message_id) - expect(Option.isSome(started)).toBe(true) - if (Option.isNone(started)) return - - yield* ledger.setTarget(m.message_id, ThreadId.make("t1"), ChannelId.make("c1")) - yield* ledger.setPrompt(m.message_id, "prompt:hello", SessionId.make("s1")) - yield* ledger.setResponse(m.message_id, "reply:hello") - yield* ledger.complete(m.message_id) - - const next = yield* ledger.start(m.message_id) - expect(Option.isNone(next)).toBe(true) - }), - ), - ) - - effectTest("replays pending rows and recovers processing rows", () => - withLedger((ledger) => - Effect.gen(function* () { - const a = event("m-a", "one") - const b = event("m-b", "two") - yield* ledger.admit(a) - yield* ledger.admit(b) - - const started = yield* ledger.start(a.message_id) - expect(Option.isSome(started)).toBe(true) - - const replay = yield* ledger.replayPending() - expect(replay.map((x) => x.message_id)).toEqual(["m-a", "m-b"]) - - const again = yield* ledger.start(a.message_id) - expect(Option.isSome(again)).toBe(true) - }), - ), - ) - - effectTest("retains cached response across retry and prunes old completed rows", () => - withLedger((ledger, sql) => - Effect.gen(function* () { - const m = event("m-cache", "cache") - yield* ledger.admit(m) - yield* ledger.start(m.message_id) - yield* ledger.setTarget(m.message_id, ThreadId.make("t1"), ChannelId.make("c1")) - yield* ledger.setPrompt(m.message_id, "prompt:cache", SessionId.make("s1")) - yield* ledger.setResponse(m.message_id, "reply:cache") - yield* ledger.retry(m.message_id, "send failed") - - const resumed = yield* ledger.start(m.message_id) - expect(Option.isSome(resumed)).toBe(true) - if (Option.isSome(resumed)) { - expect(resumed.value.response_text).toBe("reply:cache") - expect(resumed.value.thread_id).toBe(ThreadId.make("t1")) - expect(resumed.value.channel_id).toBe(ChannelId.make("c1")) - } - - yield* ledger.complete(m.message_id) - yield* sql`UPDATE conversation_inbox - SET completed_at = datetime('now', '-10 minutes') - WHERE message_id = ${m.message_id}` - yield* ledger.prune() - - const rows = yield* sql<{ n: number }>`SELECT COUNT(*) AS n FROM conversation_inbox WHERE message_id = ${m.message_id}` - expect(rows[0]?.n ?? 0).toBe(0) - }), - ), + effectTest("dedup returns true first time, false second time", () => + Effect.gen(function* () { + const ledger = yield* ConversationLedger + expect(yield* ledger.dedup("m1")).toBe(true) + expect(yield* ledger.dedup("m1")).toBe(false) + expect(yield* ledger.dedup("m2")).toBe(true) + expect(yield* ledger.dedup("m2")).toBe(false) + }).pipe(Effect.provide(ConversationLedger.noop)), ) effectTest("stores and updates source offsets", () => @@ -163,4 +71,23 @@ describe("ConversationLedger", () => { }), ), ) + + effectTest("dedup works in layer mode", () => + withLedger((ledger) => + Effect.gen(function* () { + expect(yield* ledger.dedup("m1")).toBe(true) + expect(yield* ledger.dedup("m1")).toBe(false) + expect(yield* ledger.dedup("m2")).toBe(true) + }), + ), + ) + + effectTest("noop offsets always return none", () => + Effect.gen(function* () { + const ledger = yield* ConversationLedger + expect(Option.isNone(yield* ledger.getOffset("thread:t1"))).toBe(true) + yield* ledger.setOffset("thread:t1", "m1") + expect(Option.isNone(yield* ledger.getOffset("thread:t1"))).toBe(true) + }).pipe(Effect.provide(ConversationLedger.noop)), + ) }) diff --git a/packages/discord/src/conversation/services/ledger.ts b/packages/discord/src/conversation/services/ledger.ts index 6d5142854f..603b302db8 100644 --- a/packages/discord/src/conversation/services/ledger.ts +++ b/packages/discord/src/conversation/services/ledger.ts @@ -1,65 +1,32 @@ import * as Client from "@effect/sql/SqlClient" -import { Context, Effect, Layer, Option, Schedule, Schema } from "effect" -import { AppConfig } from "../../config" +import { Context, Effect, Layer, Option } from "effect" import { SqliteDb } from "../../db/client" import { initializeSchema } from "../../db/init" import { DatabaseError } from "../../errors" -import { ChannelId, SessionId, ThreadId } from "../../types" -import { Inbound } from "../model/schema" -const DEDUP_TTL_MINUTES = 5 -const PRUNE_BATCH_SIZE = 500 - -type Snapshot = { - thread_id: ThreadId | null - channel_id: ChannelId | null - response_text: string | null - prompt_text: string | null - session_id: SessionId | null -} - -export class MessageState extends Schema.Class("MessageState")({ - thread_id: Schema.NullOr(ThreadId), - channel_id: Schema.NullOr(ChannelId), - response_text: Schema.NullOr(Schema.String), - prompt_text: Schema.NullOr(Schema.String), - session_id: Schema.NullOr(SessionId), -}) {} - -const InboundJson = Schema.parseJson(Inbound) -const decode = Schema.decodeUnknown(InboundJson) -const encode = Schema.encode(InboundJson) +const DEDUP_LIMIT = 4_000 const db = (effect: Effect.Effect) => effect.pipe(Effect.mapError((cause) => new DatabaseError({ cause }))) -const changes = (sql: Client.SqlClient) => - db( - sql<{ n: number }>`SELECT changes() AS n`.pipe( - Effect.map((rows) => rows[0]?.n ?? 0), - ), - ) - -const toState = (row: Snapshot) => - MessageState.make({ - thread_id: row.thread_id, - channel_id: row.channel_id, - response_text: row.response_text, - prompt_text: row.prompt_text, - session_id: row.session_id, - }) +const makeDedupSet = () => { + const seen = new Set() + const order: Array = [] + return (message_id: string): boolean => { + if (seen.has(message_id)) return false + seen.add(message_id) + order.push(message_id) + if (order.length > DEDUP_LIMIT) { + const oldest = order.shift() + if (oldest) seen.delete(oldest) + } + return true + } +} export declare namespace ConversationLedger { export interface Service { - readonly admit: (event: Inbound) => Effect.Effect - readonly replayPending: () => Effect.Effect, DatabaseError> - readonly start: (message_id: string) => Effect.Effect, DatabaseError> - readonly setTarget: (message_id: string, thread_id: ThreadId, channel_id: ChannelId) => Effect.Effect - readonly setPrompt: (message_id: string, prompt: string, session_id: SessionId) => Effect.Effect - readonly setResponse: (message_id: string, response: string) => Effect.Effect - readonly complete: (message_id: string) => Effect.Effect - readonly retry: (message_id: string, error: string) => Effect.Effect - readonly prune: () => Effect.Effect + readonly dedup: (message_id: string) => Effect.Effect readonly getOffset: (source_id: string) => Effect.Effect, DatabaseError> readonly setOffset: (source_id: string, message_id: string) => Effect.Effect } @@ -69,161 +36,22 @@ export class ConversationLedger extends Context.Tag("@discord/conversation/Conve ConversationLedger, ConversationLedger.Service >() { - static readonly noop = Layer.effect( - ConversationLedger, - Effect.sync(() => { - const pending = new Set() - const completed = new Set() - return ConversationLedger.of({ - admit: (event) => - Effect.sync(() => { - if (pending.has(event.message_id) || completed.has(event.message_id)) return false - pending.add(event.message_id) - return true - }), - replayPending: () => Effect.succeed([]), - start: (message_id) => - Effect.sync(() => { - if (!pending.has(message_id)) return Option.none() - pending.delete(message_id) - return Option.some(MessageState.make({ - thread_id: null, - channel_id: null, - response_text: null, - prompt_text: null, - session_id: null, - })) - }), - setTarget: () => Effect.void, - setPrompt: () => Effect.void, - setResponse: () => Effect.void, - complete: (message_id) => Effect.sync(() => { completed.add(message_id) }), - retry: (message_id) => Effect.sync(() => { pending.add(message_id) }), - prune: () => Effect.void, - getOffset: () => Effect.succeed(Option.none()), - setOffset: () => Effect.void, - }) - }), - ) + static readonly noop = Layer.sync(ConversationLedger, () => { + const check = makeDedupSet() + return ConversationLedger.of({ + dedup: (message_id) => Effect.sync(() => check(message_id)), + getOffset: () => Effect.succeed(Option.none()), + setOffset: () => Effect.void, + }) + }) - static readonly layer = Layer.scoped( + static readonly layer = Layer.effect( ConversationLedger, Effect.gen(function* () { const sql = yield* SqliteDb - const config = yield* AppConfig yield* db(initializeSchema.pipe(Effect.provideService(Client.SqlClient, sql))) - const admit = Effect.fn("ConversationLedger.admit")(function* (event: Inbound) { - const payload = yield* encode(event).pipe( - Effect.mapError((cause) => new DatabaseError({ cause })), - ) - yield* db( - sql`INSERT OR IGNORE INTO conversation_inbox (message_id, kind, payload_json, status, created_at, updated_at) - VALUES (${event.message_id}, ${event.kind}, ${payload}, 'pending', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`, - ) - return (yield* changes(sql)) > 0 - }) - - const replayPending = Effect.fn("ConversationLedger.replayPending")(function* () { - yield* db( - sql`UPDATE conversation_inbox - SET status = 'pending', updated_at = CURRENT_TIMESTAMP - WHERE status = 'processing'`, - ) - const rows = yield* db( - sql<{ payload_json: string }>`SELECT payload_json - FROM conversation_inbox - WHERE status = 'pending' - ORDER BY created_at ASC`, - ) - return yield* Effect.forEach(rows, (row) => - decode(row.payload_json).pipe(Effect.mapError((cause) => new DatabaseError({ cause }))), - ) - }) - - const start = Effect.fn("ConversationLedger.start")(function* (message_id: string) { - yield* db( - sql`UPDATE conversation_inbox - SET status = 'processing', attempts = attempts + 1, - processing_started_at = CURRENT_TIMESTAMP, - updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id} AND status = 'pending'`, - ) - if ((yield* changes(sql)) === 0) return Option.none() - const rows = yield* db( - sql`SELECT thread_id, channel_id, response_text, prompt_text, session_id - FROM conversation_inbox - WHERE message_id = ${message_id} - LIMIT 1`, - ) - const row = rows[0] - if (!row) return Option.none() - return Option.some(toState(row)) - }) - - const setTarget = Effect.fn("ConversationLedger.setTarget")(function* ( - message_id: string, - thread_id: ThreadId, - channel_id: ChannelId, - ) { - yield* db( - sql`UPDATE conversation_inbox - SET thread_id = ${thread_id}, channel_id = ${channel_id}, updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id}`, - ) - }) - - const setPrompt = Effect.fn("ConversationLedger.setPrompt")(function* ( - message_id: string, - prompt: string, - session_id: SessionId, - ) { - yield* db( - sql`UPDATE conversation_inbox - SET prompt_text = ${prompt}, session_id = ${session_id}, updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id}`, - ) - }) - - const setResponse = Effect.fn("ConversationLedger.setResponse")(function* (message_id: string, response: string) { - yield* db( - sql`UPDATE conversation_inbox - SET response_text = ${response}, updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id}`, - ) - }) - - const complete = Effect.fn("ConversationLedger.complete")(function* (message_id: string) { - yield* db( - sql`UPDATE conversation_inbox - SET status = 'completed', completed_at = CURRENT_TIMESTAMP, - processing_started_at = NULL, last_error = NULL, updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id}`, - ) - }) - - const retry = Effect.fn("ConversationLedger.retry")(function* (message_id: string, error: string) { - yield* db( - sql`UPDATE conversation_inbox - SET status = 'pending', last_error = ${error}, updated_at = CURRENT_TIMESTAMP - WHERE message_id = ${message_id}`, - ) - }) - - const prune = Effect.fn("ConversationLedger.prune")(function* () { - yield* db( - sql`DELETE FROM conversation_inbox - WHERE message_id IN ( - SELECT message_id - FROM conversation_inbox - WHERE status = 'completed' - AND completed_at IS NOT NULL - AND completed_at < datetime('now', '-' || ${DEDUP_TTL_MINUTES} || ' minutes') - ORDER BY completed_at ASC - LIMIT ${PRUNE_BATCH_SIZE} - )`, - ) - }) + const check = makeDedupSet() const getOffset = Effect.fn("ConversationLedger.getOffset")(function* (source_id: string) { const rows = yield* db( @@ -247,25 +75,8 @@ export class ConversationLedger extends Context.Tag("@discord/conversation/Conve ) }) - yield* prune().pipe( - Effect.catchAll((error) => - Effect.logError("Conversation ledger prune failed").pipe( - Effect.annotateLogs({ event: "conversation.ledger.prune.failed", error: String(error) }), - )), - Effect.repeat(Schedule.spaced(config.cleanupInterval)), - Effect.forkScoped, - ) - return ConversationLedger.of({ - admit, - replayPending, - start, - setTarget, - setPrompt, - setResponse, - complete, - retry, - prune, + dedup: (message_id) => Effect.sync(() => check(message_id)), getOffset, setOffset, }) diff --git a/packages/discord/src/db/init.test.ts b/packages/discord/src/db/init.test.ts index d88a7d1cbd..f23cd5a5b8 100644 --- a/packages/discord/src/db/init.test.ts +++ b/packages/discord/src/db/init.test.ts @@ -31,29 +31,6 @@ const indexes = [ "discord_sessions_status_updated_at_idx", ] -const inboxColumns = [ - "message_id", - "kind", - "payload_json", - "status", - "thread_id", - "channel_id", - "prompt_text", - "session_id", - "response_text", - "attempts", - "processing_started_at", - "completed_at", - "last_error", - "created_at", - "updated_at", -] - -const inboxIndexes = [ - "conversation_inbox_status_created_at_idx", - "conversation_inbox_completed_at_idx", -] - const offsetColumns = [ "source_id", "last_message_id", @@ -69,21 +46,11 @@ const getColumns = (db: Client.SqlClient) => Effect.map((rows) => rows.map((row: { name: string }) => row.name)), ) -const getInboxColumns = (db: Client.SqlClient) => - db<{ name: string }>`PRAGMA table_info(conversation_inbox)`.pipe( - Effect.map((rows) => rows.map((row: { name: string }) => row.name)), - ) - const getIndexes = (db: Client.SqlClient) => db<{ name: string }>`PRAGMA index_list(discord_sessions)`.pipe( Effect.map((rows) => rows.map((row: { name: string }) => row.name)), ) -const getInboxIndexes = (db: Client.SqlClient) => - db<{ name: string }>`PRAGMA index_list(conversation_inbox)`.pipe( - Effect.map((rows) => rows.map((row: { name: string }) => row.name)), - ) - const getOffsetColumns = (db: Client.SqlClient) => db<{ name: string }>`PRAGMA table_info(conversation_offsets)`.pipe( Effect.map((rows) => rows.map((row: { name: string }) => row.name)), @@ -100,27 +67,21 @@ describe("initializeSchema", () => { Effect.gen(function* () { yield* withSqlite(filename, (db) => initializeSchema.pipe(Effect.provideService(Client.SqlClient, db))) const one = yield* withSqlite(filename, getColumns) - const inboxOne = yield* withSqlite(filename, getInboxColumns) const offsetOne = yield* withSqlite(filename, getOffsetColumns) expect(one).toEqual(columns) - expect(inboxOne).toEqual(inboxColumns) expect(offsetOne).toEqual(offsetColumns) yield* withSqlite(filename, (db) => initializeSchema.pipe(Effect.provideService(Client.SqlClient, db))) const two = yield* withSqlite(filename, getColumns) - const inboxTwo = yield* withSqlite(filename, getInboxColumns) const offsetTwo = yield* withSqlite(filename, getOffsetColumns) expect(two).toEqual(one) - expect(inboxTwo).toEqual(inboxOne) expect(offsetTwo).toEqual(offsetOne) const seen = new Set(two) expect(seen.size).toBe(two.length) const actual = (yield* withSqlite(filename, getIndexes)).filter((name) => !name.startsWith("sqlite_")) - const inboxActual = (yield* withSqlite(filename, getInboxIndexes)).filter((name) => !name.startsWith("sqlite_")) const offsetActual = (yield* withSqlite(filename, getOffsetIndexes)).filter((name) => !name.startsWith("sqlite_")) expect(new Set(actual)).toEqual(new Set(indexes)) - expect(new Set(inboxActual)).toEqual(new Set(inboxIndexes)) expect(new Set(offsetActual)).toEqual(new Set(offsetIndexes)) }), "discord-sessions-", diff --git a/packages/discord/src/db/migrations/0001_discord_sessions.ts b/packages/discord/src/db/migrations/0001_discord_sessions.ts index 3ac7b84142..9a34091d5c 100644 --- a/packages/discord/src/db/migrations/0001_discord_sessions.ts +++ b/packages/discord/src/db/migrations/0001_discord_sessions.ts @@ -23,24 +23,6 @@ const TABLE = `CREATE TABLE IF NOT EXISTS discord_sessions ( updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )` -const INBOX_TABLE = `CREATE TABLE IF NOT EXISTS conversation_inbox ( - message_id TEXT PRIMARY KEY, - kind TEXT NOT NULL CHECK (kind IN ('thread_message', 'channel_message')), - payload_json TEXT NOT NULL, - status TEXT NOT NULL CHECK (status IN ('pending', 'processing', 'completed')), - thread_id TEXT, - channel_id TEXT, - prompt_text TEXT, - session_id TEXT, - response_text TEXT, - attempts INTEGER NOT NULL DEFAULT 0, - processing_started_at TEXT, - completed_at TEXT, - last_error TEXT, - created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, - updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP -)` - const OFFSETS_TABLE = `CREATE TABLE IF NOT EXISTS conversation_offsets ( source_id TEXT PRIMARY KEY, last_message_id TEXT NOT NULL, @@ -62,21 +44,6 @@ const COLUMNS = [ ["resume_fail_count", "INTEGER NOT NULL DEFAULT 0"], ] as const -const INBOX_COLUMNS = [ - ["status", "TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'processing', 'completed'))"], - ["thread_id", "TEXT"], - ["channel_id", "TEXT"], - ["prompt_text", "TEXT"], - ["session_id", "TEXT"], - ["response_text", "TEXT"], - ["attempts", "INTEGER NOT NULL DEFAULT 0"], - ["processing_started_at", "TEXT"], - ["completed_at", "TEXT"], - ["last_error", "TEXT"], - ["created_at", "TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP"], - ["updated_at", "TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP"], -] as const - const OFFSET_COLUMNS = [ ["last_message_id", "TEXT NOT NULL"], ["updated_at", "TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP"], @@ -89,13 +56,6 @@ const INDEXES = [ ON discord_sessions (status, updated_at)`, ] as const -const INBOX_INDEXES = [ - `CREATE INDEX IF NOT EXISTS conversation_inbox_status_created_at_idx - ON conversation_inbox (status, created_at)`, - `CREATE INDEX IF NOT EXISTS conversation_inbox_completed_at_idx - ON conversation_inbox (completed_at)`, -] as const - const OFFSET_INDEXES = [ `CREATE INDEX IF NOT EXISTS conversation_offsets_updated_at_idx ON conversation_offsets (updated_at)`, @@ -104,7 +64,6 @@ const OFFSET_INDEXES = [ export default Effect.gen(function* () { const db = yield* Client.SqlClient yield* db.unsafe(TABLE) - yield* db.unsafe(INBOX_TABLE) yield* db.unsafe(OFFSETS_TABLE) const names = new Set((yield* db<{ name: string }>`PRAGMA table_info(discord_sessions)`).map((row) => row.name)) @@ -113,12 +72,6 @@ export default Effect.gen(function* () { discard: true, }) - const inboxNames = new Set((yield* db<{ name: string }>`PRAGMA table_info(conversation_inbox)`).map((row) => row.name)) - const inboxMissing = INBOX_COLUMNS.filter(([name]) => !inboxNames.has(name)) - yield* Effect.forEach(inboxMissing, ([name, definition]) => db.unsafe(`ALTER TABLE conversation_inbox ADD COLUMN ${name} ${definition}`), { - discard: true, - }) - const offsetNames = new Set((yield* db<{ name: string }>`PRAGMA table_info(conversation_offsets)`).map((row) => row.name)) const offsetMissing = OFFSET_COLUMNS.filter(([name]) => !offsetNames.has(name)) yield* Effect.forEach(offsetMissing, ([name, definition]) => db.unsafe(`ALTER TABLE conversation_offsets ADD COLUMN ${name} ${definition}`), { @@ -126,6 +79,5 @@ export default Effect.gen(function* () { }) yield* Effect.forEach(INDEXES, (index) => db.unsafe(index), { discard: true }) - yield* Effect.forEach(INBOX_INDEXES, (index) => db.unsafe(index), { discard: true }) yield* Effect.forEach(OFFSET_INDEXES, (index) => db.unsafe(index), { discard: true }) })