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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
306 changes: 306 additions & 0 deletions src/services/__tests__/omni-bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -705,3 +705,309 @@ describe('OmniBridge β€” session lifecycle (Group 2)', () => {
}
});
});

// ============================================================================
// Session reset subscription β€” issue #1089
// ============================================================================

describe('OmniBridge β€” session reset (#1089)', () => {
/** Pre-populate a live session entry on the bridge for reset tests. */
function injectSession(
bridge: OmniBridge,
key: string,
instanceId: string,
session: OmniSession,
): { entry: { idleTimer: ReturnType<typeof setTimeout> | null } } {
const idleTimer = setTimeout(() => {}, 60_000);
const entry = {
session,
instanceId,
spawning: false,
buffer: [],
idleTimer,
};
(bridge as any).sessions.set(key, entry);
return { entry };
}

it('shuts down the executor and removes the session on reset for a hot chat', async () => {
const { executor, calls, makeSession } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
const session = makeSession('test-agent', 'chat-1');
injectSession(bridge, 'test-agent:chat-1', 'inst-1', session);

await (bridge as any).handleSessionReset('inst-1', 'chat-1', 'kill');

expect(calls.shutdown.length).toBe(1);
expect(calls.shutdown[0].chatId).toBe('chat-1');
expect((bridge as any).sessions.size).toBe(0);
} finally {
await bridge.stop();
}
});

it('no-ops on reset for a cold chat (no live session)', async () => {
const { executor, calls } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
// Sessions map is empty β€” reset must not throw and must not call shutdown.
await (bridge as any).handleSessionReset('inst-1', 'chat-cold');

expect(calls.shutdown.length).toBe(0);
expect((bridge as any).sessions.size).toBe(0);
} finally {
await bridge.stop();
}
});

it('clears the idle timer when evicting a reset session', async () => {
const { executor, makeSession } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
const session = makeSession('test-agent', 'chat-1');
const { entry } = injectSession(bridge, 'test-agent:chat-1', 'inst-1', session);
const timerBefore = entry.idleTimer;
expect(timerBefore).not.toBeNull();

await (bridge as any).handleSessionReset('inst-1', 'chat-1');

// Session removed β†’ idle timer no longer reachable from sessions map
expect((bridge as any).sessions.has('test-agent:chat-1')).toBe(false);
} finally {
await bridge.stop();
}
});

it('parses subject with dotted chatId (e.g. WhatsApp +5511...@s.whatsapp.net)', async () => {
const { executor, calls, makeSession } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
const dottedChat = '+5511999999999@s.whatsapp.net';
const session = makeSession('test-agent', dottedChat);
injectSession(bridge, `test-agent:${dottedChat}`, 'inst-x', session);

// Build a fake NATS message and feed it through the dispatch path the way
// processSessionResetEvents would: subject + JSON-encoded payload bytes.
const fakeMsg = {
subject: `omni.session.reset.inst-x.${dottedChat}`,
data: new TextEncoder().encode(JSON.stringify({ action: 'kill' })),
};
// Drive a single iteration through processSessionResetEvents using a one-shot iterator.
const oneShot = {
[Symbol.asyncIterator]: async function* () {
yield fakeMsg;
},
} as any;
await (bridge as any).processSessionResetEvents(oneShot);

expect(calls.shutdown.length).toBe(1);
expect(calls.shutdown[0].chatId).toBe(dottedChat);
} finally {
await bridge.stop();
}
});

it('tolerates malformed JSON payload by routing on subject alone', async () => {
const { executor, calls, makeSession } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
const session = makeSession('test-agent', 'chat-1');
injectSession(bridge, 'test-agent:chat-1', 'inst-1', session);

const fakeMsg = {
subject: 'omni.session.reset.inst-1.chat-1',
data: new TextEncoder().encode('not-json-at-all'),
};
const oneShot = {
[Symbol.asyncIterator]: async function* () {
yield fakeMsg;
},
} as any;
await (bridge as any).processSessionResetEvents(oneShot);

// Subject-only routing still kills the session.
expect(calls.shutdown.length).toBe(1);
} finally {
await bridge.stop();
}
});

it('warns and skips on malformed subject with too few segments', async () => {
const { executor, calls } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
const fakeMsg = {
subject: 'omni.session.reset',
data: new TextEncoder().encode('{}'),
};
const oneShot = {
[Symbol.asyncIterator]: async function* () {
yield fakeMsg;
},
} as any;
await (bridge as any).processSessionResetEvents(oneShot);

expect(calls.shutdown.length).toBe(0);
} finally {
await bridge.stop();
}
});

it('cancels a spawning session on reset and tears down the freshly-spawned executor', async () => {
// Hold the spawn promise open so we can fire reset mid-spawn.
let releaseSpawn!: (s: OmniSession) => void;
const spawnGate = new Promise<OmniSession>((resolve) => {
releaseSpawn = resolve;
});

const { executor, calls, makeSession } = makeMockExecutor({
spawnFn: async (agentName, chatId) => {
// Block until the test releases the spawn.
const session = await spawnGate;
return session ?? makeSession(agentName, chatId);
},
});
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
await bridge.start();

try {
// Kick off the spawn β€” routeMessage will block on spawnGate.
const routePromise = (bridge as any).routeMessage(makeMsg({ content: 'first' }));

// Yield so spawnSession installs the placeholder before we reset.
await new Promise((r) => setTimeout(r, 5));
expect((bridge as any).sessions.has('test-agent:chat-1')).toBe(true);
expect((bridge as any).sessions.get('test-agent:chat-1').spawning).toBe(true);

// Fire the reset while spawn is still in flight.
await (bridge as any).handleSessionReset('inst-1', 'chat-1', 'kill');

// Placeholder is gone β€” the spawn-in-flight has been logically cancelled.
expect((bridge as any).sessions.has('test-agent:chat-1')).toBe(false);

// Now release the spawn β€” spawnSession should detect cancelled and tear it down.
releaseSpawn(makeSession('test-agent', 'chat-1'));
await routePromise;

// executor.shutdown was called for the freshly-spawned (cancelled) session.
expect(calls.shutdown.length).toBe(1);
expect(calls.shutdown[0].chatId).toBe('chat-1');

// Buffered triggering message must NOT be delivered to the killed session.
expect(calls.deliver.length).toBe(0);
} finally {
await bridge.stop();
}
});

it('drains the message queue after evicting a reset session', async () => {
const { executor, calls, makeSession } = makeMockExecutor();
const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => makeFakeNats()) as any,
});
(bridge as any).executor = executor;
// Force the bridge to look fully saturated so a reset opens a slot.
(bridge as any).maxConcurrent = 1;
await bridge.start();

try {
const session = makeSession('test-agent', 'chat-1');
injectSession(bridge, 'test-agent:chat-1', 'inst-1', session);

// Park a queued message that's waiting for a free slot.
(bridge as any).messageQueue.push(makeMsg({ chatId: 'chat-2', agent: 'agent-b' }));

await (bridge as any).handleSessionReset('inst-1', 'chat-1', 'kill');

// Original session is gone, queued message picked up its slot, drainQueue spawned it.
expect((bridge as any).sessions.has('test-agent:chat-1')).toBe(false);
expect((bridge as any).messageQueue.length).toBe(0);
expect(calls.spawn.length).toBe(1);
expect(calls.spawn[0].chatId).toBe('chat-2');
} finally {
await bridge.stop();
}
});

it('subscribes to omni.session.reset.> on start()', async () => {
const subscribeCalls: string[] = [];
const fakeSub: Partial<Subscription> & AsyncIterable<never> = {
unsubscribe: () => {},
[Symbol.asyncIterator]: async function* () {},
};
const nc: Partial<NatsConnection> = {
info: undefined,
closed: async () => undefined,
close: async () => undefined,
drain: async () => undefined,
publish: () => {},
subscribe: (subject: string) => {
subscribeCalls.push(subject);
return fakeSub as Subscription;
},
};

const bridge = new OmniBridge({
natsUrl: 'test://fake',
pgProvider: degradedPgProvider,
natsConnectFn: (async () => nc as NatsConnection) as any,
});

try {
await bridge.start();
expect(subscribeCalls).toContain('omni.session.reset.>');
} finally {
await bridge.stop();
}
});
});
Loading
Loading