import { afterEach, beforeEach, describe, expect, it } from 'bun:test'; import { mkdtempSync, rmSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { defineEvents, openBus } from '@celilo/event-bus'; import { type ConfigReply, type ConfigRequiredPayload, EVENT_TYPES, busInterview, } from './bus-interview'; const NO_SCHEMAS = defineEvents({}); describe('busInterview', () => { let dir: string; let dbPath: string; let origEnv: string | undefined; beforeEach(() => { dir = mkdtempSync(join(tmpdir(), 'bus-interview-test-')); dbPath = join(dir, 'events.db'); origEnv = process.env.EVENT_BUS_DB; process.env.EVENT_BUS_DB = dbPath; }); afterEach(() => { if (origEnv === undefined) delete process.env.EVENT_BUS_DB; else process.env.EVENT_BUS_DB = origEnv; try { rmSync(dir, { recursive: true, force: true }); } catch { /* ignore */ } }); it('emits a query and returns the responder reply', async () => { // Stand up a responder in a separate Bus instance against the same // DB. The deploy's busInterview opens its own short-lived Bus and // queries; the responder's Bus emits the reply. The cross-process // poll inside bus.query is what makes this work. const responderBus = openBus({ dbPath, events: NO_SCHEMAS }); const watch = responderBus.watch('config.required.lunacycle.domain', (event) => { responderBus.emitRaw( `${event.type}.reply`, { value: 'example.net' }, { replyFor: event.id, emittedBy: 'test-responder' }, ); }); const payload: ConfigRequiredPayload = { module: 'lunacycle', key: 'domain', type: 'string', required: true, }; const reply = await busInterview( EVENT_TYPES.configRequired('lunacycle', 'domain'), payload, ); expect(reply.value).toBe('example.net'); watch.close(); responderBus.close(); }); it('honors first-reply-wins when multiple responders compete', async () => { const fastResponder = openBus({ dbPath, events: NO_SCHEMAS }); const slowResponder = openBus({ dbPath, events: NO_SCHEMAS }); fastResponder.watch('config.required.foo.bar', (event) => { fastResponder.emitRaw( `${event.type}.reply`, { value: 'fast' }, { replyFor: event.id, emittedBy: 'fast' }, ); }); let slowFired = false; slowResponder.watch('config.required.foo.bar', (event) => { // Reply after a delay; the fast responder should win. setTimeout(() => { // Bus may have been closed by the time this fires (test won // already). Skip the emit if so — checking via the // testFinished flag set after assertions. if (slowFired) return; slowFired = true; try { slowResponder.emitRaw( `${event.type}.reply`, { value: 'slow' }, { replyFor: event.id, emittedBy: 'slow' }, ); } catch { // Bus closed; benign in this test. } }, 200); }); const reply = await busInterview(EVENT_TYPES.configRequired('foo', 'bar'), { module: 'foo', key: 'bar', type: 'string', required: true, }); expect(reply.value).toBe('fast'); slowFired = true; // suppress the slow timer if it hasn't fired yet fastResponder.close(); slowResponder.close(); }); it('uses the existing Bus instance when one is passed in', async () => { const sharedBus = openBus({ dbPath, events: NO_SCHEMAS }); sharedBus.watch('config.required.shared.x', (event) => { sharedBus.emitRaw( `${event.type}.reply`, { value: 'shared' }, { replyFor: event.id, emittedBy: 'in-process' }, ); }); const reply = await busInterview( EVENT_TYPES.configRequired('shared', 'x'), { module: 'shared', key: 'x', type: 'string', required: true }, sharedBus, ); expect(reply.value).toBe('shared'); sharedBus.close(); }); }); describe('EVENT_TYPES', () => { it('builds dotted event names for each interview shape', () => { expect(EVENT_TYPES.configRequired('lunacycle', 'domain')).toBe( 'config.required.lunacycle.domain', ); expect(EVENT_TYPES.secretRequired('authentik', 'admin_password')).toBe( 'secret.required.authentik.admin_password', ); expect(EVENT_TYPES.ensureRequired('namecheap', 'add_domain')).toBe( 'ensure.required.namecheap.add_domain', ); }); }); describe('multi-select payload roundtrip', () => { let dir: string; let dbPath: string; let origEnv: string | undefined; beforeEach(() => { dir = mkdtempSync(join(tmpdir(), 'bus-interview-multi-')); dbPath = join(dir, 'events.db'); origEnv = process.env.EVENT_BUS_DB; process.env.EVENT_BUS_DB = dbPath; }); afterEach(() => { if (origEnv === undefined) delete process.env.EVENT_BUS_DB; else process.env.EVENT_BUS_DB = origEnv; try { rmSync(dir, { recursive: true, force: true }); } catch { /* ignore */ } }); it('options-bearing payload can be answered with a string[] reply', async () => { // Stand up a responder that ignores text input and replies with a // selection. Mimics what the terminal-responder does when it sees // options[] and uses p.multiselect instead of promptText. const responder = openBus({ dbPath, events: NO_SCHEMAS }); responder.watch('config.required.greenwave.zones', (event) => { const payload = event.payload as ConfigRequiredPayload; expect(payload.options).toBeDefined(); expect(payload.options?.map((o) => o.value).sort()).toEqual(['app', 'dmz', 'internal']); responder.emitRaw( `${event.type}.reply`, { value: ['dmz', 'internal'] }, { replyFor: event.id, emittedBy: 'test' }, ); }); const reply = await busInterview( EVENT_TYPES.configRequired('greenwave', 'zones'), { module: 'greenwave', key: 'zones', type: 'array', required: true, options: [ { value: 'dmz', label: 'DMZ' }, { value: 'app', label: 'App tier' }, { value: 'internal', label: 'Internal' }, ], } satisfies ConfigRequiredPayload, ); expect(reply.value).toEqual(['dmz', 'internal']); responder.close(); }); });