diff --git a/packages/mcp/server.integration.test.ts b/packages/mcp/server.integration.test.ts index e2e51f6..8fd216a 100644 --- a/packages/mcp/server.integration.test.ts +++ b/packages/mcp/server.integration.test.ts @@ -33,6 +33,7 @@ import { decodeMessageBodySync, decodeMessageIdSync, decodeTimestampSync, + IdentityError, InboxError, MessagePermalinkSchema, userMentions, @@ -173,6 +174,11 @@ interface AdapterOverrides { * replay already covered the backlog — onAcquire must then SKIP catch-up. */ readonly resumeQueueReplayed?: boolean + /** + * Whether the realm answers a sign-in, asked on every acquire. While it + * returns false the acquire fails as a realm that cannot be reached. + */ + readonly realmReachable?: () => boolean /** * Capture stderr-shaped log output (default: route to the runner's * STDERR via the production logger). Pass an array to collect the @@ -272,7 +278,19 @@ const buildHarness = async (overrides: AdapterOverrides = {}): Promise acquire: (name) => Effect.sync(() => { acquires.push(name) - }).pipe(Effect.flatMap(() => base.identity.acquire(name))), + }).pipe( + Effect.flatMap(() => + (overrides.realmReachable ?? (() => true))() + ? base.identity.acquire(name) + : Effect.fail( + new IdentityError({ + operation: 'acquire', + cause: new Error('Transport error (GET https://zulip.example.com/api/v1/users)'), + transient: true, + }), + ), + ), + ), release: () => Effect.sync(() => { releases += 1 @@ -2470,6 +2488,28 @@ test('boot feeder is a no-op when CLAUDE_CODE_SESSION_ID is absent — deferred * first `listTools()` and the test would pass without ever seeing the tool * present. */ +test('a seat with a bot name completes the handshake while the realm is down, then signs in once it answers', async () => { + let reachable = false + const cap = captureSubscribes() + const h = await buildHarness({ + env: { COMMY_PROJECT: 'myproject' }, + seedChannels: ['myproject'], + inboxOverrides: cap.inboxOverrides, + realmReachable: () => reachable, + }) + try { + await waitFor(() => h.identityCalls.acquires.length >= 2, 2000) + expect((await h.client.listTools()).tools.map((t) => t.name)).toContain('post') + expect(cap.tokens).toEqual([]) + + reachable = true + await waitFor(() => cap.tokens.length === 2, 5000) + expect(new Set(cap.tokens)).toEqual(new Set(['new-topics:myproject', 'myproject/general'])) + } finally { + await h.cleanup() + } +}) + test('an administrator turning message editing off withdraws edit_message from a connected seat', async () => { const flip = Deferred.unsafeMake(FiberId.none) const h = await buildHarness({ diff --git a/packages/mcp/server.test.ts b/packages/mcp/server.test.ts index 3ff2c8e..37d9c2a 100644 --- a/packages/mcp/server.test.ts +++ b/packages/mcp/server.test.ts @@ -185,6 +185,8 @@ const buildFakeAdapter = ( readonly identityOrigin?: IdentityOrigin /** Reject every substrate-side subscribe, for the part-way-failure paths. */ readonly subscribeError?: InboxError + /** How the realm answers the editing-setting probe. */ + readonly editingAvailable?: MessagePublisher['editingAvailable'] } = {}, ): { readonly adapter: ZulipAdapter; readonly calls: FakeAdapterCalls } => { const acquired: string[] = [] @@ -222,7 +224,7 @@ const buildFakeAdapter = ( const publisher: MessagePublisher = { post: () => Effect.die(new Error('unused fake')), edit: () => Effect.void, - editingAvailable: () => Effect.succeed(true), + editingAvailable: options.editingAvailable ?? (() => Effect.succeed(true)), react: () => Effect.void, unreact: () => Effect.void, resolveThread: () => Effect.void, @@ -329,7 +331,6 @@ test('main writes acquire failure to stderr in the canonical format, fails boot, // A seat that boots while the realm is unreachable retries rather than dying, // because the host caches a failed MCP connection for the rest of its session. -// The retry has to finish inside the host's startup timeout, so it is bounded. const unreachable = () => new IdentityError({ operation: 'acquire', @@ -369,24 +370,34 @@ test('boot retries an unreachable realm and comes up once it answers', async () expect(logs.filter((line) => line.includes('could not reach the realm'))).toHaveLength(2) }) -test('boot gives up on a realm that stays unreachable, inside the host startup timeout', async () => { +test('boot keeps retrying a realm that stays unreachable, long after the host startup timeout', async () => { const fake = buildFakeAdapter({ - acquireFailures: Array.from({ length: 100 }, unreachable), + acquireFailures: Array.from({ length: 1000 }, unreachable), }) - const stderr: string[] = [] + const stillRunning = await Effect.runPromise( + Effect.gen(function* () { + const fiber = yield* Effect.fork(provideProgram(validEnv, fake.adapter)) + yield* TestClock.adjust('10 minutes') + const exit = yield* Fiber.poll(fiber) + yield* Fiber.interrupt(fiber) + return Option.isNone(exit) + }).pipe(Effect.provide(TestContext.TestContext)), + ) + expect(stillRunning).toBe(true) + expect(fake.calls.acquired.length).toBeGreaterThan(20) + expect(fake.calls.closes.count).toBe(1) +}) + +test('boot does not wait on a realm that never answers the editing-setting probe', async () => { + const fake = buildFakeAdapter({ editingAvailable: () => Effect.never }) const exit = await runProgramOnTestClock( validEnv, fake.adapter, - { loggerLayer: captureLogger(stderr) }, - '25 seconds', - ) - expect(Exit.isFailure(exit)).toBe(true) - expect(fake.calls.acquired.length).toBeGreaterThan(2) - expect(fake.calls.acquired.length).toBeLessThan(100) - expect(stderr.at(-1)).toBe( - 'commy plugin: acquire("myproject-concierge") failed: Transport error (GET https://zulip.example.com/api/v1/users)', + { readGitContext: () => Effect.succeed(NotInRepo()) }, + '1 minute', ) - expect(fake.calls.closes.count).toBe(1) + expect(Exit.isSuccess(exit)).toBe(true) + expect(fake.calls.acquired).toEqual(['myproject-concierge']) }) test('boot does not retry a realm that refuses the sign-in', async () => { diff --git a/packages/mcp/server.ts b/packages/mcp/server.ts index 0c86730..e9ada3a 100644 --- a/packages/mcp/server.ts +++ b/packages/mcp/server.ts @@ -334,17 +334,16 @@ const buildIdentityCache = ( } /** - * How long a seat with a bot name keeps retrying a realm it cannot reach at - * boot. The acquire runs before the MCP handshake, so the retry has to finish - * inside the host's startup timeout — Claude Code's `MCP_TIMEOUT`, 30 seconds - * by default — with room left for the rest of boot. A host that gives up + * How long boot waits for the realm's editing setting. The probe runs before + * the MCP handshake, which has to finish inside the host's startup timeout — + * Claude Code's `MCP_TIMEOUT`, 30 seconds by default. A host that gives up * first caches the failure for the rest of its session. */ -const BOOT_ACQUIRE_RETRY_BUDGET = Duration.seconds(20) +const EDITING_PROBE_TIMEOUT = Duration.seconds(5) -const bootAcquireRetrySchedule = Schedule.exponential(Duration.millis(250)).pipe( - Schedule.union(Schedule.spaced(Duration.seconds(4))), - Schedule.upTo(BOOT_ACQUIRE_RETRY_BUDGET), +/** A seat with a bot name retries a realm it cannot reach for as long as it runs. */ +const signInRetrySchedule = Schedule.exponential(Duration.millis(250)).pipe( + Schedule.union(Schedule.spaced(Duration.seconds(30))), ) const isTransientIdentityFailure = (err: EnsureBoundError): boolean => @@ -375,27 +374,28 @@ const isTransientIdentityFailure = (err: EnsureBoundError): boolean => * 5. buildIdentityCache — single (persistent) or ephemeral (1-slot, * release-then-acquire across session_id * transitions). - * 6. env-driven branch on `parsed.botName`: - * • Persistent (botName set): eager acquire now. A realm it cannot - * reach is retried for up to BOOT_ACQUIRE_RETRY_BUDGET. Rejection, - * or a realm still unreachable after that, writes the canonical - * diagnostic and FAILS with BootError — runMain's teardown maps that - * to exit 1. - * • Ephemeral (botName unset): skip the boot acquire; the first - * attribution-producing call mints `cc-[-]` lazily. - * 7. release finalizer — registered after acquire, gated on - * `boundIdentityIds().size > 0`. + * 6. release finalizer — gated on `boundIdentityIds().size > 0`. * LIFO ordering: it runs AFTER pump-cancel * (registered later) and BEFORE the substrate * `close()` finalizer (the outer layer scope), * giving cancel-pump → release → close. - * 8. subscribeFromEnv — apply COMMY_SUBSCRIBE tokens. - * 9. registerTools — wire the cache + projectForCwd into the MCP + * 7. ephemeral COMMY_SUBSCRIBE — trigger the lazy mint when the seat has + * tokens to hold; the mint's onAcquire applies them. + * 8. registerTools — wire the cache + projectForCwd into the MCP * tool surface. - * 10. mcp.connect — bind to the supplied transport (stdio in prod). - * 11. persistent boot-time mentions + channels catch-up. No-op in - * ephemeral mode (onAcquire owns that path). - * 12. startEventPump — filter inbound events through the narrowSet, + * 9. mcp.connect — bind to the supplied transport (stdio in prod). + * Nothing before it waits on the realm for longer + * than EDITING_PROBE_TIMEOUT. + * 10. persistent sign-in (botName set) — acquire, retrying a realm it + * cannot reach for as long as the seat runs, then + * Type-1 defaults, the COMMY_SUBSCRIBE seed, and + * the mentions + channels catch-up. A rejection + * writes the canonical diagnostic and FAILS with + * BootError — runMain's teardown maps that to + * exit 1. Ephemeral seats skip this: the first + * attribution-producing call mints + * `cc-[-]` lazily. + * 11. startEventPump — filter inbound events through the narrowSet, * dispatch via the notifier. The pump is a * daemon, so a `pump.cancel` finalizer stops it * on scope unwind (signal interrupt under @@ -575,15 +575,14 @@ export const makeProgram = ( // worse failure of the two: the tool still refuses legibly if editing // really is off, whereas a wrongly-hidden tool is simply gone with no // way for the caller to discover why. Do not make this fatal. - const canEditMessages = yield* adapter.publisher - .editingAvailable() - .pipe( - Effect.catchAll((error) => - Effect.logWarning( - `commy plugin: could not read the realm's message-editing setting (${error.message}); offering edit_message and letting the substrate refuse if it is off`, - ).pipe(Effect.as(true)), - ), - ) + const canEditMessages = yield* adapter.publisher.editingAvailable().pipe( + Effect.timeout(EDITING_PROBE_TIMEOUT), + Effect.catchAll((error) => + Effect.logWarning( + `commy plugin: could not read the realm's message-editing setting (${error.message}); offering edit_message and letting the substrate refuse if it is off`, + ).pipe(Effect.as(true)), + ), + ) const narrowSet = createNarrowSet() const mcp = buildMcpServer() @@ -797,57 +796,7 @@ export const makeProgram = ( // identity to arrive. yield* installBinder(yield* SessionBinderTag, binderFor(identityCache)) - // Persistent mode (COMMY_BOT_NAME set): eager acquire so a - // misconfigured concierge dies at boot rather than on first message. - // The single-identity cache ignores the session_id. - // Ephemeral mode: skip — the first tool call mints lazily. - let type1Intents: ReadonlyArray = [] - // The persistent bot's boot acquire, held so the COMMY_SUBSCRIBE - // bootstrap below can read the origin of the bind that just happened. - // Persistent mode has no post-acquire hook to hang the seeding off, so - // this is where its mint is observable. - let persistentIdentity: AcquiredIdentity | undefined - if (parsed.botName !== undefined) { - const botName = parsed.botName - const ensureBound = yield* identityCache.ensureBoundFor(undefined) - type1Intents = yield* ensureBound().pipe( - Effect.tapError((err) => - isTransientIdentityFailure(err) - ? Effect.logWarning( - `commy plugin: acquire("${botName}") could not reach the realm: ${err.message}`, - ) - : Effect.void, - ), - Effect.retry({ schedule: bootAcquireRetrySchedule, while: isTransientIdentityFailure }), - // `ensureBound()` is the acquire Effect (ensure-bound.ts). Squash the - // whole Cause — typed acquire failure or an adapter defect — into the - // BootError the operator must fix. - Effect.catchAllCause((cause) => { - const err = Cause.squash(cause) - return Effect.fail( - new BootError({ message: Predicate.isError(err) ? err.message : String(err) }), - ) - }), - Effect.matchEffect({ - onFailure: (bootErr) => - Effect.logError( - `commy plugin: acquire("${botName}") failed: ${bootErr.message}`, - ).pipe(Effect.zipRight(Effect.fail(bootErr))), - // Type-1 defaults: post-acquire register the - // universal `mentions` narrow plus project-specific subs. - onSuccess: (acquired) => - Effect.sync(() => { - persistentIdentity = acquired - }).pipe( - Effect.zipRight( - registerType1DefaultsOnBoot(adapter.inbox, narrowSet, parsed.project), - ), - ), - }), - ) - } - - // Release-on-shutdown finalizer. Registered after acquire and + // Release-on-shutdown finalizer. Registered // before the pump-cancel finalizer, so LIFO teardown runs // pump-cancel → release; the substrate `close()` (outer layer scope) // then runs last. Gated on the acquisition state: an @@ -889,10 +838,10 @@ export const makeProgram = ( ) // Two mint paths, one bootstrap, and boot's job differs between them. // - // PERSISTENT: the eager acquire above IS the mint, and there is no - // post-acquire hook, so the seeding happens here — after the Type-1 - // defaults, which are a separate, per-boot concern that this bead does not - // touch (they are computed from the project slug, not read from + // PERSISTENT: the sign-in after the handshake IS the mint, and there is + // no post-acquire hook, so `signInPersistentSeat` seeds it — + // after the Type-1 defaults, which are a separate, per-boot concern + // (they are computed from the project slug, not read from // COMMY_SUBSCRIBE, and a pane depends on them being re-registered every // boot). // @@ -902,35 +851,42 @@ export const makeProgram = ( // principal, so a seat asked to hold subscriptions needs an identity at // boot. A seat asked to hold nothing must NOT mint here; it stays // identity-free until it acts. - const subscribedIntents = yield* withSessionContext( - persistentIdentity !== undefined - ? seedSubscriptionsOnMint(persistentIdentity, seedDeps) - : parsed.subscribe === undefined - ? Effect.succeed>([]) - : binderFor(identityCache).pipe(Effect.as>([])), - { sessionId: Option.getOrUndefined(bootSessionId), project: parsed.project }, - ).pipe( - // A seat with no way to bind cannot hold subscriptions at all — the - // ephemeral bot name is derived from the session id, so there is no name - // to mint under. Log what was lost and carry on serving: deaf is the - // accepted outcome here (the residual gap Graeme's 2026-07-05 ruling - // names), a dead MCP child is not. The seeding is not lost with it: an - // ephemeral seat that could not bind at boot still seeds when its first - // action mints it. - Effect.catchIf(isBindError, (cause) => - Effect.logWarning( - `commy plugin: boot-time subscribe could not bind an identity, so no ` + - `subscriptions were applied yet — this seat will not receive channel traffic ` + - `until it acts. ${Predicate.isError(cause) ? cause.message : String(cause)}`, - ).pipe(Effect.provide(loggerLayer), Effect.as>([])), - ), - ) + const withBootSubscribeContext = ( + seed: Effect.Effect< + ReadonlyArray, + SubscribeTokenError | InboxError | BindError + >, + ) => + withSessionContext(seed, { + sessionId: Option.getOrUndefined(bootSessionId), + project: parsed.project, + }).pipe( + // A seat with no way to bind cannot hold subscriptions at all — the + // ephemeral bot name is derived from the session id, so there is no name + // to mint under. Log what was lost and carry on serving: deaf is the + // accepted outcome here (the residual gap Graeme's 2026-07-05 ruling + // names), a dead MCP child is not. The seeding is not lost with it: an + // ephemeral seat that could not bind at boot still seeds when its first + // action mints it. + Effect.catchIf(isBindError, (cause) => + Effect.logWarning( + `commy plugin: boot-time subscribe could not bind an identity, so no ` + + `subscriptions were applied yet — this seat will not receive channel traffic ` + + `until it acts. ${Predicate.isError(cause) ? cause.message : String(cause)}`, + ).pipe(Effect.provide(loggerLayer), Effect.as>([])), + ), + ) + if (parsed.botName === undefined && parsed.subscribe !== undefined) { + yield* withBootSubscribeContext( + binderFor(identityCache).pipe(Effect.as>([])), + ) + } - // Start journaling runtime subscribe/unsubscribe deltas now — after the - // boot-time defaults and any COMMY_SUBSCRIBE bootstrap have landed, but - // before a tool call can mutate — so a delta racing the boot-forked - // rebuild below is replayed onto its result rather than clobbered by it. - // Both modes buffer, because both now rebuild. + // Start journaling subscribe/unsubscribe deltas now, before a tool call + // can mutate, so a delta racing the rebuild below is replayed onto its + // result rather than clobbered by it. A persistent seat's Type-1 + // defaults and COMMY_SUBSCRIBE seed land after this and are replayed the + // same way. Both modes buffer, because both now rebuild. narrowSet.beginBuffering() const toolsCache = registerTools(mcp, { @@ -981,104 +937,155 @@ export const makeProgram = ( return ids.size === 0 ? undefined : ids.values().next().value } - // Missed-mentions catch-up on persistent-mode resume + - // boot-time channel/thread catch-up. Both non-fatal — - // log + continue. Ephemeral mode runs the equivalent via onAcquire. - const persistentBotId = parsed.botName !== undefined ? getBotIdentityId() : undefined - if (persistentBotId !== undefined) { - yield* catchUpMentions({ - cursorStore, - inbox: adapter.inbox, - identityId: persistentBotId, - notifier, - }).pipe(Effect.catchAllCause(logCatchUpFailure('mentions'))) + // Persistent mode (COMMY_BOT_NAME set): sign in after the handshake, so + // the host has its tool list whatever state the realm is in — a host + // that times out the handshake caches the failure for its whole session. + // A realm it cannot reach is retried for as long as the seat runs; a + // rejection still ends the server, so a misconfigured concierge dies + // loudly rather than sitting deaf. The single-identity cache ignores the + // session_id. Ephemeral mode skips this: the first tool call mints lazily. + const signInPersistentSeat = (botName: BotName) => + Effect.gen(function* () { + const ensureBound = yield* identityCache.ensureBoundFor(undefined) + const acquired = yield* ensureBound().pipe( + Effect.tapError((err) => + isTransientIdentityFailure(err) + ? Effect.logWarning( + `commy plugin: acquire("${botName}") could not reach the realm: ${err.message}`, + ) + : Effect.void, + ), + Effect.retry({ schedule: signInRetrySchedule, while: isTransientIdentityFailure }), + // `ensureBound()` is the acquire Effect (ensure-bound.ts). Squash the + // whole Cause — typed acquire failure or an adapter defect — into the + // BootError the operator must fix. + Effect.catchAllCause((cause) => { + const err = Cause.squash(cause) + return Effect.fail( + new BootError({ message: Predicate.isError(err) ? err.message : String(err) }), + ) + }), + Effect.tapError((bootErr) => + Effect.logError(`commy plugin: acquire("${botName}") failed: ${bootErr.message}`), + ), + ) + const type1Intents = yield* registerType1DefaultsOnBoot( + adapter.inbox, + narrowSet, + parsed.project, + ) + const subscribedIntents = yield* withBootSubscribeContext( + seedSubscriptionsOnMint(acquired, seedDeps), + ) - const windowSeconds = parsed.catchupWindowSeconds ?? DEFAULT_CATCHUP_WINDOW_SECONDS - const catchUpIntents = [...type1Intents, ...subscribedIntents] - if (windowSeconds > 0 && catchUpIntents.length > 0) { - yield* catchUpChannels({ - intents: catchUpIntents, - history: adapter.history, + // Missed-mentions catch-up on persistent-mode resume + + // boot-time channel/thread catch-up. Both non-fatal — + // log + continue. Ephemeral mode runs the equivalent via onAcquire. + const persistentBotId = acquired.identity.id + yield* catchUpMentions({ + cursorStore, + inbox: adapter.inbox, + identityId: persistentBotId, notifier, - botIdentityId: persistentBotId, - windowSeconds, - }).pipe(Effect.catchAllCause(logCatchUpFailure('channels'))) - } - } + }).pipe(Effect.catchAllCause(logCatchUpFailure('mentions'))) - const pump = yield* startEventPump({ - inbox: adapter.inbox, - notifier, - getBotIdentityId, - // No binding means the seat owns no events queue, so an event in hand - // is not this seat's to filter — the same statement the producer makes - // one layer down by idling instead of polling when `ownerHttp` refuses - // (zulip/events.ts, `UnboundEphemeralSession`). Reached only if a - // release lands between an event being produced and dispatched. - match: (event) => { - const botIdentityId = getBotIdentityId() - return botIdentityId !== undefined && narrowSet.matches(event, botIdentityId) - }, - // Populate the tools-side identity cache from inbound events so - // `presence` / `post` mentions / `react` can resolve ids only ever - // seen via a notification. - rememberIdentity: toolsCache.rememberIdentity, - // Advance the per-identity cursor on every observed mention so the - // next resume's catch-up has an accurate "have-seen-up-to" mark. - // Returns the write Effect for the pump to sequence; - // its failures are swallowed because the advance is best-effort and - // cursor writes are monotonic. - onMention: (ts) => { - const id = getBotIdentityId() - if (id === undefined) return Effect.void - return cursorStore.write(id, ts).pipe(Effect.catchAllCause(() => Effect.void)) - }, - // The boot-time `canEditMessages` sample above is only good until an - // administrator moves it. Rebuild the tool list against the new value - // and tell the client its list changed — without the notification the - // seat holds a stale list until it reconnects, which for a persistent - // seat can be days. - // - // This does not replace the `editing-disabled` arm of - // `MessageEditRefused`: the two are complements. The arm still covers - // the window between the flip and this signal arriving, and the - // fail-open path above where the boot probe never answered. - onRealmSettings: (settings) => - Effect.sync(() => toolsCache.setEditingAvailable(settings.editingAvailable)).pipe( - Effect.zipRight( - Effect.promise(() => mcp.sendToolListChanged()).pipe( - Effect.catchAllCause( - logSettingsFailure( - `tools/list_changed after editing=${settings.editingAvailable}`, + const windowSeconds = parsed.catchupWindowSeconds ?? DEFAULT_CATCHUP_WINDOW_SECONDS + const catchUpIntents = [...type1Intents, ...subscribedIntents] + if (windowSeconds > 0 && catchUpIntents.length > 0) { + yield* catchUpChannels({ + intents: catchUpIntents, + history: adapter.history, + notifier, + botIdentityId: persistentBotId, + windowSeconds, + }).pipe(Effect.catchAllCause(logCatchUpFailure('channels'))) + } + }) + + // Started after a persistent seat signs in: the stream decides at start + // whether to synthesise mention events, from the identity bound then. + const serve = Effect.gen(function* () { + const pump = yield* startEventPump({ + inbox: adapter.inbox, + notifier, + getBotIdentityId, + // No binding means the seat owns no events queue, so an event in hand + // is not this seat's to filter — the same statement the producer makes + // one layer down by idling instead of polling when `ownerHttp` refuses + // (zulip/events.ts, `UnboundEphemeralSession`). Reached only if a + // release lands between an event being produced and dispatched. + match: (event) => { + const botIdentityId = getBotIdentityId() + return botIdentityId !== undefined && narrowSet.matches(event, botIdentityId) + }, + // Populate the tools-side identity cache from inbound events so + // `presence` / `post` mentions / `react` can resolve ids only ever + // seen via a notification. + rememberIdentity: toolsCache.rememberIdentity, + // Advance the per-identity cursor on every observed mention so the + // next resume's catch-up has an accurate "have-seen-up-to" mark. + // Returns the write Effect for the pump to sequence; + // its failures are swallowed because the advance is best-effort and + // cursor writes are monotonic. + onMention: (ts) => { + const id = getBotIdentityId() + if (id === undefined) return Effect.void + return cursorStore.write(id, ts).pipe(Effect.catchAllCause(() => Effect.void)) + }, + // The boot-time `canEditMessages` sample above is only good until an + // administrator moves it. Rebuild the tool list against the new value + // and tell the client its list changed — without the notification the + // seat holds a stale list until it reconnects, which for a persistent + // seat can be days. + // + // This does not replace the `editing-disabled` arm of + // `MessageEditRefused`: the two are complements. The arm still covers + // the window between the flip and this signal arriving, and the + // fail-open path above where the boot probe never answered. + onRealmSettings: (settings) => + Effect.sync(() => toolsCache.setEditingAvailable(settings.editingAvailable)).pipe( + Effect.zipRight( + Effect.promise(() => mcp.sendToolListChanged()).pipe( + Effect.catchAllCause( + logSettingsFailure( + `tools/list_changed after editing=${settings.editingAvailable}`, + ), ), ), ), ), - ), - }) - // The pump is a daemon (forkDaemon) — not scope-tied — so an explicit - // finalizer interrupts it on scope unwind (signal interrupt under - // runMain, or natural stream end). Registered after the release - // finalizer so LIFO teardown runs pump-cancel first. - yield* Effect.addFinalizer(() => pump.cancel) + }) + // The pump is a daemon (forkDaemon) — not scope-tied — so an explicit + // finalizer interrupts it on scope unwind (signal interrupt under + // runMain, or natural stream end). Registered after the release + // finalizer so LIFO teardown runs pump-cancel first. + yield* Effect.addFinalizer(() => pump.cancel) + + // Boot-forked narrow-set rebuild, scope-tied. Asks the + // realm what this seat is subscribed to and narrows the answer with its + // recorded topic intents, then replays the deltas journaled since + // `beginBuffering` so a subscribe racing the load is never lost. Never + // awaited inline: an ephemeral seat's record parks on the shared session + // id, and an unfilled id on a listen-only seat would otherwise block the + // scope forever. + yield* Effect.forkScoped(rebuildNarrowSet) - // Boot-forked narrow-set rebuild, scope-tied. Asks the - // realm what this seat is subscribed to and narrows the answer with its - // recorded topic intents, then replays the deltas journaled since - // `beginBuffering` so a subscribe racing the load is never lost. Never - // awaited inline: an ephemeral seat's record parks on the shared session - // id, and an unfilled id on a listen-only seat would otherwise block the - // scope forever. - yield* Effect.forkScoped(rebuildNarrowSet) + yield* pump.done + }) // Block until either the event stream ends / the pump fatally parks - // and is interrupted (the SIGINT/SIGTERM path), OR the MCP client - // disconnects. The pump is a daemon that + // and is interrupted (the SIGINT/SIGTERM path), the persistent sign-in + // is rejected, OR the MCP client disconnects. The pump is a daemon that // long-polls Zulip forever, so without the disconnect race a plain // client exit — which sends no signal — would block here and orphan // the server child. `raceFirst` interrupts the losing wait and the // scope then unwinds the finalizers (pump-cancel → release → close). - yield* Effect.raceFirst(pump.done, params.shutdownSignal ?? Effect.never) + yield* Effect.raceFirst( + parsed.botName === undefined + ? serve + : signInPersistentSeat(parsed.botName).pipe(Effect.zipRight(serve)), + params.shutdownSignal ?? Effect.never, + ) }), )