diff --git a/packages/testing/stub-http-client.test.ts b/packages/testing/stub-http-client.test.ts index 75935a5..13fb0a3 100644 --- a/packages/testing/stub-http-client.test.ts +++ b/packages/testing/stub-http-client.test.ts @@ -1,6 +1,6 @@ import { expect, test } from 'bun:test' import { HttpClient, HttpClientRequest } from '@effect/platform' -import { Effect } from 'effect' +import { Duration, Effect, Fiber, Option } from 'effect' import { makeStubHttpClient } from './stub-http-client.ts' const REALM = 'https://zulip.example.com/api/v1' @@ -202,6 +202,29 @@ test('opens no socket — a request to an unroutable host still resolves', () => }), )) +test('a hang response captures the request, then never resolves and stays interruptible', () => + Effect.runPromise( + Effect.gen(function* () { + const stub = yield* makeStubHttpClient + // The long-poll hold: the stub answers instantly for everything else, so + // a terminal hang is what stops an eager consumer draining the sequence. + yield* stub.respondSequence('GET', '/api/v1/events', [{ hang: true }]) + const fiber = yield* Effect.fork( + stub.client.execute(HttpClientRequest.get(`${REALM}/events`)), + ) + // Let the forked request issue and park on the hang. + yield* Effect.sleep(Duration.millis(10)) + const captured = yield* stub.captured + expect(captured).toHaveLength(1) + expect(captured[0]?.url.pathname).toBe('/api/v1/events') + // Still parked — a hang resolves to no Exit. + expect(Option.isNone(yield* fiber.poll)).toBe(true) + // Interrupting unwinds it cleanly (the scope-close path). + const exit = yield* Fiber.interrupt(fiber) + expect(exit._tag).toBe('Failure') + }), + )) + test('is a drop-in for the HttpClient.HttpClient service', () => Effect.runPromise( Effect.gen(function* () { diff --git a/packages/testing/stub-http-client.ts b/packages/testing/stub-http-client.ts index 7a37f20..af72f09 100644 --- a/packages/testing/stub-http-client.ts +++ b/packages/testing/stub-http-client.ts @@ -23,9 +23,22 @@ import { type HttpClientRequest, HttpClientResponse, } from '@effect/platform' -import { Data, Effect, HashMap, Option, Ref } from 'effect' +import { Data, Effect, HashMap, Option, Predicate, Ref } from 'effect' -export type StubResponse = { +/** + * A response that never resolves — the request is captured, then the effect + * parks on `Effect.never`. Models the real Zulip long-poll *holding* the + * connection open: an in-memory stub answers instantly, so without a hold an + * eager consumer (the event-pump's `Stream.runDrain`) would burn through the + * whole response sequence in a hot loop. A fiber blocked on a hang is + * *interrupted* (not errored) when its scope closes — which is exactly the + * scope-close-interrupt path the event-pump tests exercise. + */ +export type StubHang = { + readonly hang: true +} + +export type StubBody = { /** Object → JSON-encoded; string → verbatim; `Uint8Array` → raw bytes. */ readonly body: unknown /** HTTP status; defaults to 200. */ @@ -34,6 +47,8 @@ export type StubResponse = { readonly headers?: Readonly> } +export type StubResponse = StubBody | StubHang + export type CapturedHttpRequest = { readonly method: string readonly url: URL @@ -87,19 +102,19 @@ const requestBodyInit = (body: HttpBody.HttpBody): string | Uint8Array | FormDat } } -const responseBodyInit = (response: StubResponse): string | Uint8Array => { +const responseBodyInit = (response: StubBody): string | Uint8Array => { if (response.body instanceof Uint8Array) return response.body if (typeof response.body === 'string') return response.body return JSON.stringify(response.body) } -const responseHeaders = (response: StubResponse): Record => { +const responseHeaders = (response: StubBody): Record => { const base: Record = response.body instanceof Uint8Array ? {} : { 'content-type': 'application/json' } return { ...base, ...response.headers } } -const notFound = (method: string, path: string): StubResponse => ({ +const notFound = (method: string, path: string): StubBody => ({ body: { result: 'error', code: 'NO_STUB_HANDLER', msg: `no stub handler for ${method} ${path}` }, status: 404, }) @@ -155,14 +170,18 @@ export const makeStubHttpClient: Effect.Effect = Effect.gen(func const client = HttpClient.make((request, url) => capture(request, url).pipe( Effect.zipRight(nextResponse(request.method, url.pathname)), - Effect.map((response) => - HttpClientResponse.fromWeb( - request, - new Response(responseBodyInit(response), { - status: response.status ?? 200, - headers: responseHeaders(response), - }), - ), + Effect.flatMap((response) => + Predicate.hasProperty(response, 'hang') + ? Effect.never + : Effect.succeed( + HttpClientResponse.fromWeb( + request, + new Response(responseBodyInit(response), { + status: response.status ?? 200, + headers: responseHeaders(response), + }), + ), + ), ), ), ) diff --git a/packages/zulip/adapter-events.test.ts b/packages/zulip/adapter-events.test.ts new file mode 100644 index 0000000..68ea314 --- /dev/null +++ b/packages/zulip/adapter-events.test.ts @@ -0,0 +1,511 @@ +/** + * Event-pump LOGIC, exercised through the full adapter stack + * (adapter → ZulipHttp → HttpClient) on the **owned-fake stub HttpClient + + * TestClock** — no `Bun.serve`, no real socket. + * + * These are the long-poll / reconnect tests that used to drive a real + * `FetchHttpClient` against an in-process `Bun.serve` realm in `adapter.test.ts`. + * Moving the LOGIC onto the stub makes them deterministic and dissolves the + * real-socket contention flake (comms-hbm9 / comms-xwqm): the infinite + * long-poll hold becomes a stub `{ hang: true }` response that parks on + * `Effect.never`, and scope-close interruption of `Effect.never` is + * deterministic. The retry/backoff sleeps run on the virtual clock. + * + * SCOPE EDGE: the stub proves the *Effect fiber-interrupt* logic. The genuine + * `AbortSignal → fetch → TCP teardown` of an in-flight `FetchHttpClient` + * long-poll on scope close stays a real-socket Tier-3 test (comms-4lz5) — that + * integration cannot move off the socket. + * + * The producer-level pump logic (transient-retry breadcrumbs, the exponential + * backoff cap, per-event parsing) is covered separately at the `inboxEvents` + * unit level in `events.test.ts` with a `ZulipHttp`-level fake. These tests add + * the *adapter wiring*: the adapter-scoped gap-replay watermark surviving across + * `events()` iterator instances, `inbox.replay()` wired into the producer's + * BAD_EVENT_QUEUE_ID recovery, and the forked-drain scope-close path. + */ + +import { expect } from 'bun:test' +import type { InboundEvent } from '@commy/core/ports' +import { + decodeBotNameSync, + decodeMessageBodySync, + decodeTimestampSync, + type IdentityError, + type UnknownIdentity, +} from '@commy/core/ports' +import { effectTest } from '@commy/testing/effect-test' +import { + type CapturedHttpRequest, + makeStubHttpClient, + type StubHttpClient, +} from '@commy/testing/stub-http-client' +import { HttpClient } from '@effect/platform' +import { + Duration, + Effect, + Queue, + Redacted, + type Scope, + Stream, + TestClock, + TestContext, +} from 'effect' +import type { ZulipAdapter } from './adapter.ts' +import { zulipAdapter } from './adapter.ts' +import { ApiKey, BotEmail, RealmUrl } from './http.ts' + +const HERMES = { + user_id: 9, + email: 'hermes-agent-bot@example.com', + full_name: 'hermes-agent', + is_bot: true, + is_active: true, + role: 400, +} as const + +const GRAEME = { + user_id: 5, + email: 'graeme@example.com', + full_name: 'Graeme Foster', + is_bot: false, + is_active: true, + role: 100, +} as const + +const REALM_URL = 'https://zulip.example.com' + +const seedUsers = (stub: StubHttpClient, members: ReadonlyArray): Effect.Effect => + stub.respond('GET', '/api/v1/users', { body: { result: 'success', members } }) + +const seedRegenerate = (stub: StubHttpClient, userId: number): Effect.Effect => + stub.respond('POST', `/api/v1/bots/${userId}/api_key/regenerate`, { + body: { result: 'success', api_key: 'fresh-key' }, + }) + +const seedRegister = (stub: StubHttpClient): Effect.Effect => + stub.respond('POST', '/api/v1/register', { + body: { result: 'success', queue_id: 'queue-1', last_event_id: 0 }, + }) + +const buildAdapter = ( + stub: StubHttpClient, +): Effect.Effect => + Effect.gen(function* () { + yield* seedUsers(stub, [HERMES, GRAEME]) + yield* seedRegenerate(stub, HERMES.user_id) + const config = { + realmUrl: yield* RealmUrl(REALM_URL).pipe(Effect.orDie), + minterEmail: yield* BotEmail('minter@example.com').pipe(Effect.orDie), + minterApiKey: Redacted.make(yield* ApiKey('minter-key').pipe(Effect.orDie)), + } + const adapter = yield* zulipAdapter(config).pipe( + Effect.provideService(HttpClient.HttpClient, stub.client), + ) + yield* adapter.identity.acquire(decodeBotNameSync('hermes-agent')) + return adapter + }) + +// Mirror inbox.events() into an unbounded Queue under the caller's Scope. The +// forked Stream.runDrain fiber lives for the scope's lifetime, so scope close +// interrupts it (and any in-flight long-poll). The stub answers instantly, so +// every event sequence ends with a `{ hang: true }` entry — the long-poll hold +// that keeps the eager drain from spinning past the canned responses. +const eventQueue = ( + adapter: ZulipAdapter, +): Effect.Effect, never, Scope.Scope> => + Effect.gen(function* () { + const queue = yield* Queue.unbounded() + yield* Effect.forkScoped( + adapter.inbox.events().pipe( + Stream.tap((event) => Queue.offer(queue, event)), + Stream.runDrain, + ), + ) + return queue + }) + +const isEventsPoll = (r: CapturedHttpRequest): boolean => + r.method === 'GET' && r.url.pathname === '/api/v1/events' + +const isRegisterPost = (r: CapturedHttpRequest): boolean => + r.method === 'POST' && r.url.pathname === '/api/v1/register' + +const eventPolls = (stub: StubHttpClient): Effect.Effect> => + stub.captured.pipe(Effect.map((reqs) => reqs.filter(isEventsPoll))) + +const registerPosts = (stub: StubHttpClient): Effect.Effect> => + stub.captured.pipe(Effect.map((reqs) => reqs.filter(isRegisterPost))) + +// Condition-gate (not a sleep): yield to the forked drain until it has issued +// `n` GET /events polls. Replaces the original real `Effect.sleep` that waited +// for a long-poll "to get into flight" — here the drain parks on a stub hang, +// and we wait for that poll to have been captured before unwinding the scope. +const awaitEventPolls = (stub: StubHttpClient, n: number): Effect.Effect => + Effect.yieldNow().pipe( + Effect.zipRight(eventPolls(stub)), + Effect.map((polls) => polls.length), + Effect.repeat({ until: (count) => count >= n }), + Effect.asVoid, + ) + +const aZulipMessage = ( + overrides: Partial<{ + id: number + sender_id: number + sender_full_name: string + stream_id: number + display_recipient: string + subject: string + content: string + timestamp: number + }> = {}, +): Record => ({ + id: 100, + sender_id: GRAEME.user_id, + sender_full_name: GRAEME.full_name, + stream_id: 1234, + display_recipient: 'general', + subject: 'lobby', + content: 'hello', + timestamp: 1715000000, + ...overrides, +}) + +const messageEvent = (id: number, message: Record): Record => ({ + id, + type: 'message', + message, + flags: [], +}) + +const gapMessagesBody = (content: string): Record => ({ + result: 'success', + messages: [ + { + id: 150, + sender_id: GRAEME.user_id, + sender_full_name: GRAEME.full_name, + stream_id: 1234, + display_recipient: 'general', + subject: 'lobby', + content, + timestamp: 1500, + flags: [], + }, + ], + anchor: 0, + found_anchor: false, + found_newest: true, + found_oldest: false, + history_limited: false, +}) + +effectTest( + 'inbox.events advances last_event_id between long-poll calls', + () => + Effect.gen(function* () { + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { + body: { + result: 'success', + events: [messageEvent(5, aZulipMessage({ content: 'first' }))], + }, + }, + { + body: { + result: 'success', + events: [messageEvent(11, aZulipMessage({ id: 200, content: 'second' }))], + }, + }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + yield* Queue.take(queue) + yield* Queue.take(queue) + const polls = yield* eventPolls(stub) + expect(polls.length).toBeGreaterThanOrEqual(2) + expect(polls[0]?.url.searchParams.get('last_event_id')).toBe('0') + expect(polls[1]?.url.searchParams.get('last_event_id')).toBe('5') + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events ignores Zulip heartbeat events but advances last_event_id', + () => + Effect.gen(function* () { + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { body: { result: 'success', events: [{ id: 17, type: 'heartbeat' }] } }, + { body: { result: 'success', events: [messageEvent(18, aZulipMessage())] } }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + const event = yield* Queue.take(queue) + expect(event.kind).toBe('message-posted') + const polls = yield* eventPolls(stub) + expect(polls.length).toBeGreaterThanOrEqual(2) + expect(polls[1]?.url.searchParams.get('last_event_id')).toBe('17') + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events re-registers and resumes when /events returns BAD_EVENT_QUEUE_ID', + () => + Effect.gen(function* () { + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { + body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, + status: 400, + }, + { body: { result: 'success', events: [messageEvent(1, aZulipMessage())] } }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + const event = yield* Queue.take(queue) + expect(event.kind).toBe('message-posted') + // First poll's dead queue forces a fresh /register before the retry. + expect(yield* registerPosts(stub)).toHaveLength(2) + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events backfills the gap via inbox.replay() on BAD_EVENT_QUEUE_ID and marks events replayed=true (comms-jnn)', + () => + Effect.gen(function* () { + // Live message at ts=1000 sets the watermark; BAD_EVENT_QUEUE_ID dies the + // queue; the iterator calls replay(since=1000) which returns a message + // posted at ts=1500 during the dead window; a fresh live message at + // ts=2000 follows on the re-registered queue. The middle event must + // surface with replayed=true; the live ones must not. + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respond('GET', '/api/v1/messages', { + body: gapMessagesBody('posted during the gap'), + }) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { + body: { + result: 'success', + events: [messageEvent(1, aZulipMessage({ id: 100, timestamp: 1000 }))], + }, + }, + { + body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, + status: 400, + }, + { + body: { + result: 'success', + events: [messageEvent(1, aZulipMessage({ id: 200, timestamp: 2000 }))], + }, + }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + const collected: InboundEvent[] = [] + for (let i = 0; i < 3; i += 1) { + collected.push(yield* Queue.take(queue)) + } + expect(collected).toHaveLength(3) + expect(collected[0]?.kind).toBe('message-posted') + if (collected[0]?.kind === 'message-posted') { + expect(collected[0].replayed).toBeUndefined() + expect(String(collected[0].message.ref.id)).toBe('100') + } + expect(collected[1]?.kind).toBe('message-posted') + if (collected[1]?.kind === 'message-posted') { + expect(collected[1].replayed).toBe(true) + expect(String(collected[1].message.ref.id)).toBe('150') + expect(collected[1].message.body).toBe(decodeMessageBodySync('posted during the gap')) + } + expect(collected[2]?.kind).toBe('message-posted') + if (collected[2]?.kind === 'message-posted') { + expect(collected[2].replayed).toBeUndefined() + expect(String(collected[2].message.ref.id)).toBe('200') + } + expect(yield* registerPosts(stub)).toHaveLength(2) + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events persists the gap-replay watermark across iterator instances so reconnect-then-BAD_EVENT_QUEUE_ID still backfills (comms-4au)', + () => + Effect.gen(function* () { + // The gap-replay watermark lives at the adapter, not in the iterator's + // closure, so a second events() iterator inherits the first's last-seen + // ts and a BAD_EVENT_QUEUE_ID on the new iterator fires the replay. + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respond('GET', '/api/v1/messages', { + body: gapMessagesBody('posted during the reconnect gap'), + }) + // Sequence across both iterators: + // poll1 (iter1) — live ts=1000, advances the watermark + // poll2 (iter1) — HANG: the long-poll holds, so iter1's eager drain + // does not race ahead and consume the BAD_QUEUE meant + // for iter2; iter1's scope close interrupts it. + // poll3 (iter2) — BAD_EVENT_QUEUE_ID; iter2's first poll hits replay. + // poll4 (iter2) — live ts=2000 on the re-registered queue. + yield* stub.respondSequence('GET', '/api/v1/events', [ + { + body: { + result: 'success', + events: [messageEvent(1, aZulipMessage({ id: 100, timestamp: 1000 }))], + }, + }, + { hang: true }, + { + body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, + status: 400, + }, + { + body: { + result: 'success', + events: [messageEvent(1, aZulipMessage({ id: 200, timestamp: 2000 }))], + }, + }, + { hang: true }, + ]) + + yield* Effect.scoped( + Effect.gen(function* () { + const queue1 = yield* eventQueue(adapter) + const first = yield* Queue.take(queue1) + expect(first.kind).toBe('message-posted') + if (first.kind === 'message-posted') { + expect(first.message.ts).toBe(decodeTimestampSync(1000)) + expect(first.replayed).toBeUndefined() + } + // Hold the scope open until poll2 (the hang) is in flight, so iter1 + // owns exactly poll1+poll2 and iter2 starts at the BAD_QUEUE poll. + yield* awaitEventPolls(stub, 2) + }), + ) + + yield* Effect.scoped( + Effect.gen(function* () { + const queue2 = yield* eventQueue(adapter) + const collected: InboundEvent[] = [] + for (let i = 0; i < 2; i += 1) { + collected.push(yield* Queue.take(queue2)) + } + expect(collected[0]?.kind).toBe('message-posted') + if (collected[0]?.kind === 'message-posted') { + expect(collected[0].replayed).toBe(true) + expect(collected[0].message.ts).toBe(decodeTimestampSync(1500)) + expect(collected[0].message.body).toBe( + decodeMessageBodySync('posted during the reconnect gap'), + ) + } + expect(collected[1]?.kind).toBe('message-posted') + if (collected[1]?.kind === 'message-posted') { + expect(collected[1].replayed).toBeUndefined() + expect(collected[1].message.ts).toBe(decodeTimestampSync(2000)) + } + }), + ) + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events recovers when /events returns 429 RATE_LIMIT_HIT (comms-9wi)', + () => + Effect.gen(function* () { + // A 429 from /events is backpressure: the send path waits out the + // retry-after and retries. The wait runs on the virtual clock. + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { + body: { + result: 'error', + code: 'RATE_LIMIT_HIT', + msg: 'API usage exceeded rate limit', + 'retry-after': 0, + }, + status: 429, + }, + { body: { result: 'success', events: [messageEvent(1, aZulipMessage())] } }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + // The rate-limit backoff (min 100ms) sleeps on the virtual clock. + yield* TestClock.adjust(Duration.millis(100)) + const event = yield* Queue.take(queue) + expect(event.kind).toBe('message-posted') + expect((yield* eventPolls(stub)).length).toBeGreaterThanOrEqual(2) + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events recovers when /register returns 429 RATE_LIMIT_HIT (comms-9wi)', + () => + Effect.gen(function* () { + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* stub.respondSequence('POST', '/api/v1/register', [ + { + body: { + result: 'error', + code: 'RATE_LIMIT_HIT', + msg: 'API usage exceeded rate limit', + 'retry-after': 0, + }, + status: 429, + }, + { body: { result: 'success', queue_id: 'queue-1', last_event_id: 0 } }, + ]) + yield* stub.respondSequence('GET', '/api/v1/events', [ + { body: { result: 'success', events: [messageEvent(1, aZulipMessage())] } }, + { hang: true }, + ]) + const queue = yield* eventQueue(adapter) + yield* TestClock.adjust(Duration.millis(100)) + const event = yield* Queue.take(queue) + expect(event.kind).toBe('message-posted') + expect(yield* registerPosts(stub)).toHaveLength(2) + }), + { layer: TestContext.TestContext }, +) + +effectTest( + 'inbox.events long-poll fiber interrupts cleanly when the consumer scope closes (comms-spj3.8)', + () => + Effect.gen(function* () { + // A hung long-poll must not pin the pump on shutdown. The stub holds the + // /events poll open (Effect.never); scope close interrupts the forked + // Stream fiber. The test passes iff the outer scope close completes — the + // genuine AbortSignal→fetch→TCP teardown is the Tier-3 residue (comms-4lz5). + const stub = yield* makeStubHttpClient + const adapter = yield* buildAdapter(stub) + yield* seedRegister(stub) + yield* stub.respondSequence('GET', '/api/v1/events', [{ hang: true }]) + yield* Effect.scoped( + Effect.gen(function* () { + const queue = yield* eventQueue(adapter) + // Let the long-poll get into flight before the scope unwinds. + yield* awaitEventPolls(stub, 1) + void queue + }), + ) + expect((yield* eventPolls(stub)).length).toBeGreaterThanOrEqual(1) + }), + { layer: TestContext.TestContext }, +) diff --git a/packages/zulip/adapter.test.ts b/packages/zulip/adapter.test.ts index 39760d1..934fde3 100644 --- a/packages/zulip/adapter.test.ts +++ b/packages/zulip/adapter.test.ts @@ -1,5 +1,5 @@ import { expect, test } from 'bun:test' -import type { ChannelRef, Identity, InboundEvent, MessageRef } from '@commy/core/ports' +import type { ChannelRef, Identity, MessageRef } from '@commy/core/ports' import { DirectoryError, decodeBotNameSync, @@ -19,7 +19,7 @@ import { } from '@commy/core/ports' import { registerRealmHooks } from '@commy/testing/realm-hooks' import { FetchHttpClient, HttpClient } from '@effect/platform' -import { Cause, Duration, Effect, Exit, Option, Queue, Redacted, type Scope, Stream } from 'effect' +import { Cause, Effect, Exit, Option, Redacted } from 'effect' import type { ZulipAdapter, ZulipAdapterConfig } from './adapter.ts' import { attachmentReference, zulipAdapter as zulipAdapterRaw } from './adapter.ts' import { ApiKey, BotEmail, decodeUserUploadPathSync, RealmUrl, ZulipApiError } from './http.ts' @@ -139,25 +139,6 @@ const findRequest = (method: string, pathname: string) => { return req } -// Mirror inbox.events() into an unbounded Queue under the caller's Scope. -// The forked Stream.runDrain fiber lives for the scope's lifetime, so -// scope close interrupts it — which aborts any in-flight long-poll via the -// HttpClient's AbortSignal. Tests that need iter-shaped semantics use -// `yield* Queue.take(queue)` per substrate event. -const eventQueue = ( - adapter: ZulipAdapter, -): Effect.Effect, never, Scope.Scope> => - Effect.gen(function* () { - const queue = yield* Queue.unbounded() - yield* Effect.forkScoped( - adapter.inbox.events().pipe( - Stream.tap((event) => Queue.offer(queue, event)), - Stream.runDrain, - ), - ) - return queue - }) - test('identity.acquire on an existing bot regenerates its API key and binds', () => Effect.runPromise( Effect.gen(function* () { @@ -1484,585 +1465,6 @@ test('inbox.unsubscribe with mentions target does not call /users/me/subscriptio }), )) -interface RawZulipEvent { - readonly id: number - readonly type: string - readonly [key: string]: unknown -} - -const seedEventBatches = (batches: ReadonlyArray>): void => { - let idx = 0 - realm.handle('GET', '/api/v1/events', () => { - const events = batches[idx] ?? [] - if (idx < batches.length) idx += 1 - return { body: { result: 'success', events } } - }) -} - -const aZulipMessage = ( - overrides: Partial<{ - id: number - sender_id: number - sender_full_name: string - stream_id: number - display_recipient: string - subject: string - content: string - timestamp: number - }> = {}, -): Record => ({ - id: 100, - sender_id: 5, - sender_full_name: 'Graeme Foster', - stream_id: 1234, - display_recipient: 'general', - subject: 'lobby', - content: 'hello', - timestamp: 1715000000, - ...overrides, -}) - -test('inbox.events first .next() yields a message-posted InboundEvent for a Zulip message event', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - seedRegisterOk('queue-1', 0) - seedEventBatches([[{ id: 1, type: 'message', message: aZulipMessage(), flags: [] }]]) - const queue = yield* eventQueue(adapter) - const event = yield* Queue.take(queue) - expect(event).toEqual({ - kind: 'message-posted', - message: { - ref: { - id: decodeMessageIdSync('100'), - channel: { id: decodeChannelIdSync('1234'), name: decodeChannelNameSync('general') }, - thread: { name: decodeThreadNameSync('lobby') }, - }, - sender: { - id: decodeIdentityIdSync('5'), - name: decodeDisplayNameSync('Graeme Foster'), - kind: 'human', - }, - body: decodeMessageBodySync('hello'), - ts: decodeTimestampSync(1715000000), - mentions: [], - reactions: [], - }, - }) - }), - ), - )) - -test('inbox.events surfaces mention-received in addition to message-posted when the bound bot is mentioned in the content', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - seedRegisterOk('queue-1', 0) - seedEventBatches([ - [ - { - id: 7, - type: 'message', - message: aZulipMessage({ content: '@**hermes-agent** wake up' }), - flags: [], - }, - ], - ]) - const queue = yield* eventQueue(adapter) - const first = yield* Queue.take(queue) - const second = yield* Queue.take(queue) - expect(first.kind).toBe('message-posted') - expect(second.kind).toBe('mention-received') - if (second.kind === 'mention-received') { - expect(second.mentions.map((i: Identity) => i.name)).toEqual([ - decodeDisplayNameSync('hermes-agent'), - ]) - } - }), - ), - )) - -test('inbox.events advances last_event_id between long-poll calls', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - seedRegisterOk('queue-1', 0) - seedEventBatches([ - [ - { - id: 5, - type: 'message', - message: aZulipMessage({ id: 100, content: 'first' }), - flags: [], - }, - ], - [ - { - id: 11, - type: 'message', - message: aZulipMessage({ id: 200, content: 'second' }), - flags: [], - }, - ], - ]) - const queue = yield* eventQueue(adapter) - yield* Queue.take(queue) - yield* Queue.take(queue) - const eventCalls = realm.captured.filter( - (r) => r.method === 'GET' && r.url.pathname === '/api/v1/events', - ) - expect(eventCalls.length).toBeGreaterThanOrEqual(2) - const [first, second] = eventCalls - if (first === undefined || second === undefined) throw new Error('expected 2 event calls') - expect(first.url.searchParams.get('last_event_id')).toBe('0') - expect(second.url.searchParams.get('last_event_id')).toBe('5') - }), - ), - )) - -test('inbox.events ignores Zulip heartbeat events but advances last_event_id', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - seedRegisterOk('queue-1', 0) - seedEventBatches([ - [{ id: 17, type: 'heartbeat' }], - [{ id: 18, type: 'message', message: aZulipMessage(), flags: [] }], - ]) - const queue = yield* eventQueue(adapter) - const event = yield* Queue.take(queue) - expect(event.kind).toBe('message-posted') - const eventCalls = realm.captured.filter( - (r) => r.method === 'GET' && r.url.pathname === '/api/v1/events', - ) - expect(eventCalls.length).toBeGreaterThanOrEqual(2) - const secondCall = eventCalls[1] - if (secondCall === undefined) throw new Error('expected 2 event calls') - expect(secondCall.url.searchParams.get('last_event_id')).toBe('17') - }), - ), - )) - -test('inbox.events re-registers and resumes when /events returns BAD_EVENT_QUEUE_ID', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - let registerCalls = 0 - realm.handle('POST', '/api/v1/register', () => { - registerCalls += 1 - return { - body: { result: 'success', queue_id: `q${registerCalls}`, last_event_id: 0 }, - } - }) - let eventsCalls = 0 - realm.handle('GET', '/api/v1/events', () => { - eventsCalls += 1 - if (eventsCalls === 1) { - return { - body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, - init: { status: 400 }, - } - } - return { - body: { - result: 'success', - events: [{ id: 1, type: 'message', message: aZulipMessage(), flags: [] }], - }, - } - }) - const queue = yield* eventQueue(adapter) - const event = yield* Queue.take(queue) - expect(event.kind).toBe('message-posted') - expect(registerCalls).toBe(2) - }), - ), - )) - -test('inbox.events backfills the gap via inbox.replay() on BAD_EVENT_QUEUE_ID and marks events replayed=true (comms-jnn)', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - // End-to-end through the real adapter: live message at ts=1000 sets the - // watermark, BAD_EVENT_QUEUE_ID dies the queue, the iterator calls - // adapter.replay(since=1000) which hits /messages and returns a message - // posted at ts=1500 during the dead window, and a fresh live message at - // ts=2000 follows on the new queue. The middle event must surface with - // replayed=true; the live ones must not. - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - let registerCalls = 0 - realm.handle('POST', '/api/v1/register', () => { - registerCalls += 1 - return { - body: { result: 'success', queue_id: `q${registerCalls}`, last_event_id: 0 }, - } - }) - realm.handle('GET', '/api/v1/messages', () => ({ - body: { - result: 'success', - messages: [ - { - id: 150, - sender_id: GRAEME.user_id, - sender_full_name: GRAEME.full_name, - stream_id: 1234, - display_recipient: 'general', - subject: 'lobby', - content: 'posted during the gap', - timestamp: 1500, - flags: [], - }, - ], - anchor: 0, - found_anchor: false, - found_newest: true, - found_oldest: false, - history_limited: false, - }, - })) - let eventsCalls = 0 - realm.handle('GET', '/api/v1/events', () => { - eventsCalls += 1 - if (eventsCalls === 1) { - return { - body: { - result: 'success', - events: [ - { - id: 1, - type: 'message', - message: aZulipMessage({ id: 100, timestamp: 1000 }), - flags: [], - }, - ], - }, - } - } - if (eventsCalls === 2) { - return { - body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, - init: { status: 400 }, - } - } - return { - body: { - result: 'success', - events: [ - { - id: 1, - type: 'message', - message: aZulipMessage({ id: 200, timestamp: 2000 }), - flags: [], - }, - ], - }, - } - }) - const queue = yield* eventQueue(adapter) - const collected: InboundEvent[] = [] - for (let i = 0; i < 3; i += 1) { - collected.push(yield* Queue.take(queue)) - } - expect(collected).toHaveLength(3) - expect(collected[0]?.kind).toBe('message-posted') - if (collected[0]?.kind === 'message-posted') { - expect(collected[0].replayed).toBeUndefined() - expect(String(collected[0].message.ref.id)).toBe('100') - } - expect(collected[1]?.kind).toBe('message-posted') - if (collected[1]?.kind === 'message-posted') { - expect(collected[1].replayed).toBe(true) - expect(String(collected[1].message.ref.id)).toBe('150') - expect(collected[1].message.body).toBe(decodeMessageBodySync('posted during the gap')) - } - expect(collected[2]?.kind).toBe('message-posted') - if (collected[2]?.kind === 'message-posted') { - expect(collected[2].replayed).toBeUndefined() - expect(String(collected[2].message.ref.id)).toBe('200') - } - expect(registerCalls).toBe(2) - }), - ), - )) - -// Per-test wall clock (comms-9tik). Unlike the fakeHttp iterator tests in -// events.test.ts — which use a 2s fail-fast wall because their fake resolves -// instantly — this test drives the REAL in-process Bun.serve realm across -// ~7 sequential round-trips (users, two register/events pairs, a replay, a -// final live poll). Every wait here is condition-gated (`Queue.take` blocks on -// the forked drain producing an event); there is NO artificial sleep or -// Schedule backoff in the path — BAD_EVENT_QUEUE_ID is handled inline without -// retry — so TestClock has nothing to virtualise. The cost is irreducible real -// local I/O plus fiber scheduling, all sharing one starved event loop. In -// isolation the body runs ~60ms, but under a full-parallel `bun run check` -// (34 test files) that real work stretched to 5002ms and tripped bun's default -// 5s ceiling by 2ms. 30s clears pathological contention with wide margin while -// still failing honestly if the pump ever genuinely hangs (infinite long-poll -// loop, or scope-close abort never firing). -const GAP_REPLAY_TEST_TIMEOUT_MS = 30_000 - -test( - 'inbox.events persists the gap-replay watermark across iterator instances so reconnect-then-BAD_EVENT_QUEUE_ID still backfills (comms-4au)', - () => - Effect.runPromise( - Effect.gen(function* () { - // event-pump auto-reconnect (comms-ynb) restarts the iterator on - // transient errors, but the gap-replay watermark used to live in - // the iterator's closure — iter2's first poll had no watermark to - // anchor the replay (comms-jnn), and the very poll where replay - // mattered most was the one that skipped it. The watermark now - // lives at the adapter so iter2 inherits iter1's last-seen ts and - // BAD_EVENT_QUEUE_ID on the new iterator fires the replay. - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - let registerCalls = 0 - realm.handle('POST', '/api/v1/register', () => { - registerCalls += 1 - return { - body: { result: 'success', queue_id: `q${registerCalls}`, last_event_id: 0 }, - } - }) - realm.handle('GET', '/api/v1/messages', () => ({ - body: { - result: 'success', - messages: [ - { - id: 150, - sender_id: GRAEME.user_id, - sender_full_name: GRAEME.full_name, - stream_id: 1234, - display_recipient: 'general', - subject: 'lobby', - content: 'posted during the reconnect gap', - timestamp: 1500, - flags: [], - }, - ], - anchor: 0, - found_anchor: false, - found_newest: true, - found_oldest: false, - history_limited: false, - }, - })) - // Sequence: - // call 1 (queue1) — one live message at ts=1000 (advances watermark) - // call 2 (queue1) — Stream.runDrain pulls eagerly past the consumer's - // first take; long-poll forever so the BAD_QUEUE - // meant for queue2 is not consumed by queue1's race. - // queue1's scope close interrupts the consumer fiber - // and unwinds this fetch via AbortSignal. - // call 3 (queue2) — BAD_EVENT_QUEUE_ID — queue2's first poll hits the - // replay path. - // call 4+ (queue2) — live ts=2000 on the re-registered queue - let eventsCalls = 0 - realm.handle('GET', '/api/v1/events', async () => { - eventsCalls += 1 - if (eventsCalls === 1) { - return { - body: { - result: 'success', - events: [ - { - id: 1, - type: 'message', - message: aZulipMessage({ id: 100, timestamp: 1000 }), - flags: [], - }, - ], - }, - } - } - if (eventsCalls === 2) { - await new Promise(() => {}) - return { body: { result: 'success', events: [] } } - } - if (eventsCalls === 3) { - return { - body: { result: 'error', code: 'BAD_EVENT_QUEUE_ID', msg: 'queue expired' }, - init: { status: 400 }, - } - } - return { - body: { - result: 'success', - events: [ - { - id: 1, - type: 'message', - message: aZulipMessage({ id: 200, timestamp: 2000 }), - flags: [], - }, - ], - }, - } - }) - - // queue1: pull a live message at ts=1000 (advances the watermark), - // then close the scope to mimic the pump dropping the consumer after - // a transient error. - yield* Effect.scoped( - Effect.gen(function* () { - const queue1 = yield* eventQueue(adapter) - const first = yield* Queue.take(queue1) - expect(first.kind).toBe('message-posted') - if (first.kind === 'message-posted') { - expect(first.message.ts).toBe(decodeTimestampSync(1000)) - expect(first.replayed).toBeUndefined() - } - }), - ) - - // queue2: the pump-side reconnect. First poll hits BAD_EVENT_QUEUE_ID - // because the queue was GC'd during the backoff window. The persisted - // watermark anchors the replay at since=1000, so the message posted - // during the gap (ts=1500) surfaces as replayed=true before the next - // live event (ts=2000) on the freshly-registered queue. - yield* Effect.scoped( - Effect.gen(function* () { - const queue2 = yield* eventQueue(adapter) - const collected: InboundEvent[] = [] - for (let i = 0; i < 2; i += 1) { - collected.push(yield* Queue.take(queue2)) - } - expect(collected).toHaveLength(2) - expect(collected[0]?.kind).toBe('message-posted') - if (collected[0]?.kind === 'message-posted') { - expect(collected[0].replayed).toBe(true) - expect(collected[0].message.ts).toBe(decodeTimestampSync(1500)) - expect(collected[0].message.body).toBe( - decodeMessageBodySync('posted during the reconnect gap'), - ) - } - expect(collected[1]?.kind).toBe('message-posted') - if (collected[1]?.kind === 'message-posted') { - expect(collected[1].replayed).toBeUndefined() - expect(collected[1].message.ts).toBe(decodeTimestampSync(2000)) - } - }), - ) - }), - ), - GAP_REPLAY_TEST_TIMEOUT_MS, -) - -test('inbox.events sleeps and retries when /events returns 429 RATE_LIMIT_HIT (comms-9wi)', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - // Pre-comms-9wi: a single 429 from /events bubbled out of the iterator, - // was caught by the event-pump, and process.exit(1)'d the whole MCP - // server — disconnecting every subscribed bot in lock-step. Recovery - // must mirror the BAD_EVENT_QUEUE_ID pattern: honour the retry-after - // value Zulip emits and let the outer loop try again. - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - seedRegisterOk('queue-1', 0) - let eventsCalls = 0 - realm.handle('GET', '/api/v1/events', () => { - eventsCalls += 1 - if (eventsCalls === 1) { - return { - body: { - result: 'error', - code: 'RATE_LIMIT_HIT', - msg: 'API usage exceeded rate limit', - 'retry-after': 0, - }, - init: { status: 429 }, - } - } - return { - body: { - result: 'success', - events: [{ id: 1, type: 'message', message: aZulipMessage(), flags: [] }], - }, - } - }) - const queue = yield* eventQueue(adapter) - const event = yield* Queue.take(queue) - expect(event.kind).toBe('message-posted') - expect(eventsCalls).toBe(2) - }), - ), - )) - -test('inbox.events sleeps and retries when /register returns 429 RATE_LIMIT_HIT (comms-9wi)', () => - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - // Same recovery semantics on the registration round-trip — the - // initial /register is fired lazily inside the iterator when - // subscribe() has not yet ensured a queue. A 429 there must not - // escape the iterator either. - const adapter = yield* buildAdapter() - seedUsers([HERMES, GRAEME]) - let registerCalls = 0 - realm.handle('POST', '/api/v1/register', () => { - registerCalls += 1 - if (registerCalls === 1) { - return { - body: { - result: 'error', - code: 'RATE_LIMIT_HIT', - msg: 'API usage exceeded rate limit', - 'retry-after': 0, - }, - init: { status: 429 }, - } - } - return { body: { result: 'success', queue_id: 'q1', last_event_id: 0 } } - }) - seedEventBatches([[{ id: 1, type: 'message', message: aZulipMessage(), flags: [] }]]) - const queue = yield* eventQueue(adapter) - const event = yield* Queue.take(queue) - expect(event.kind).toBe('message-posted') - expect(registerCalls).toBe(2) - }), - ), - )) - -test('inbox.events long-poll fiber interrupts cleanly when the consumer scope closes (comms-spj3.8)', () => - Effect.runPromise( - Effect.gen(function* () { - // A hung Zulip long-poll must not pin the pump on shutdown — fiber - // interrupt has to abort the in-flight fetch and let the Effect.scoped - // boundary return. The handler hangs forever; scope close interrupts - // the underlying Stream fiber; the @effect/platform HttpClient signals - // abort via AbortController so the fetch unwinds. Test passes iff the - // outer scope close completes — if interrupt leaks, bun's per-test - // timeout fires. - const adapter = yield* buildAdapter() - seedUsers([HERMES]) - seedRegisterOk('queue-1', 0) - realm.handle('GET', '/api/v1/events', async () => { - await new Promise(() => {}) - return { body: { result: 'success', events: [] } } - }) - yield* Effect.scoped( - Effect.gen(function* () { - const queue = yield* eventQueue(adapter) - // Let the long-poll get into flight before the scope unwinds. - yield* Effect.sleep(Duration.millis(20)) - void queue - }), - ) - expect( - realm.captured.some((r) => r.method === 'GET' && r.url.pathname === '/api/v1/events'), - ).toBe(true) - }), - )) - test('inbox.replay(since) returns message-posted events for messages with ts >= since', () => Effect.runPromise( Effect.gen(function* () {