fix(omni-bridge): subscribe to omni.session.reset.> and evict live sessions - #1092
Conversation
…ssions Closes #1089 Adds the missing NATS subscription so the bridge actually reacts to session-reset events published by Omni (e.g. user sends 🗑️). Without this wiring, reset events fell on the floor and stale sessions kept serving turns until idle timeout. - subscribe('omni.session.reset.>') in start() (recursive `>` because WhatsApp chat ids contain dots) - processSessionResetEvents() parses subject `omni.session.reset.{instance}.{chat}`, tolerates malformed JSON payloads, and routes to handleSessionReset - handleSessionReset() mirrors handleTurnTimeout: clears idle timer, closes turn, calls executor.shutdown, removes from sessions Map. Cold-chat resets are a logged no-op for traceability. - 7 new tests covering hot/cold reset, idle timer cleanup, dotted WhatsApp chat ids, malformed JSON tolerance, malformed subject rejection, and subscription wiring proof. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
|
Important Review skippedAuto reviews are disabled on base/target branches other than the default branch. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Code Review
This pull request implements a session reset feature in the OmniBridge service, enabling the termination of active sessions and their associated executors via NATS messages. The implementation includes support for chat IDs containing dots and handles malformed payloads gracefully. Review feedback recommends calling drainQueue() upon session eviction to optimize concurrency and utilizing the removeSession() helper for consistent state management.
| if (entry.idleTimer) clearTimeout(entry.idleTimer); | ||
| this.turnTracker.close(sessionKey, 'reset'); | ||
| try { | ||
| await this.executor.shutdown(entry.session); | ||
| } catch (err) { | ||
| console.warn(`[omni-bridge] Error shutting down reset session ${sessionKey}:`, err); | ||
| } | ||
| this.sessions.delete(sessionKey); |
There was a problem hiding this comment.
When a session is evicted via a reset event, a slot becomes available in the concurrency limit. Calling drainQueue() ensures that any messages waiting in the messageQueue are processed immediately, improving the bridge's throughput when operating at maxConcurrent capacity.
Additionally, you can leverage the existing removeSession() helper to handle the timer cleanup and map deletion consistently.
console.log("[omni-bridge] Session reset for " + sessionKey + (action ? " (action=" + action + ")" : "") + ", evicting");
this.turnTracker.close(sessionKey, "reset");
try {
await this.executor.shutdown(entry.session);
} catch (err) {
console.warn("[omni-bridge] Error shutting down reset session " + sessionKey + ":", err);
}
this.removeSession(sessionKey);
await this.drainQueue();There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: bba1e1b03a
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| const sessionKey = this.findSessionKey(instanceId, chatId); | ||
| if (!sessionKey) { |
There was a problem hiding this comment.
Reset spawning sessions instead of treating them as cold chats
handleSessionReset resolves the target via findSessionKey(instanceId, chatId), but that lookup only matches entries with entry.session?.chatId; during spawnSession, the map contains a placeholder with spawning=true and session=null, so a reset arriving in that window is misclassified as a cold chat and ignored. In practice, if a user sends reset right after first message (while spawn is still in flight), the session still comes up and continues processing, which defeats the reset semantics this change is adding.
Useful? React with 👍 / 👎.
Address review feedback on #1092: - Codex P1: handleSessionReset misclassified spawning sessions as cold chats because findSessionKey only matched entries with entry.session?.chatId. A reset arriving in the narrow window between placeholder insertion and spawn completion was ignored, so the session finished spawning and kept serving turns — defeating the reset semantics. - Gemini medium: handleSessionReset duplicated the timer-clear and map-delete that removeSession() already encapsulates, and never called drainQueue() so queued messages couldn't take the freed concurrency slot. Changes: - SessionEntry gains an optional `cancelled` flag. - findSessionKey falls back to the map-key suffix (`:${chatId}`) when the entry is in the spawning state and entry.session is null. - handleSessionReset now has two branches: * spawning entries: mark cancelled, clear buffer, close turn, removeSession(), drainQueue(). The in-flight executor.spawn cannot be aborted, but the placeholder is removed immediately so subsequent messages spawn a fresh session. * live entries: close turn, executor.shutdown, removeSession(), drainQueue(). - spawnSession checks placeholder.cancelled after executor.spawn resolves and tears down the freshly-created session if a reset hit mid-spawn (no deliver, no idle timer). - Two new tests: 1. Cancel-mid-spawn: spawn promise is held open, reset fires while placeholder is in spawning state, then spawn is released — asserts shutdown is called and the buffered message is NOT delivered. 2. Drain-after-reset: maxConcurrent=1, queued message waits, reset frees the slot, drainQueue spawns the queued chat. All 27 omni-bridge tests pass. Full gate green (2252/2252). Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Summary
Closes #1089
Adds the missing NATS subscription so the omni bridge actually reacts to session-reset events published by Omni (e.g. user sends 🗑️). Without this wiring, reset events fell on the floor and stale sessions kept serving turns until idle timeout.
Changes
omni.session.reset.>instart()(recursive>because WhatsApp chat ids contain dots)processSessionResetEvents()parses subjectomni.session.reset.{instance}.{chat}, tolerates malformed JSON payloads, routes tohandleSessionResethandleSessionReset()mirrorshandleTurnTimeout: clears idle timer, closes turn, callsexecutor.shutdown, removes from sessions Map. Cold-chat resets are a logged no-op for traceability.Test plan
bun test src/services/__tests__/omni-bridge.test.ts— 25/25 passbun run check— 2250/2250 tests, typecheck + lint + dead-code cleanomni.session.reset.{instance}.{chat}→ bridge logs eviction → next message spawns a fresh session🤖 Generated with Claude Code