From 7718fb46819a1f154d93ddefb3293e623d1affcd Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:45:17 -0700 Subject: [PATCH 1/4] fix(connectors): tag Telegram's outbound line and route a quote-reply to the pod it quotes (TASK-156, ADR-025 D10/D11) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A user-scoped connector has one active inbound destination and N gated pods. Two gaps appear once a second gate is on. D10, Telegram only — the two things Slack already does: - the outbound line carries the pod name, so a reader can tell two pods apart; - the relayMap entry records `podId`, so a reply has a pod to route by. D11, both bridges — a quote-reply answers the line it quotes: - `routeReplyContent` / `routeSlackReplyContent` return the quoted entry's pod; - when it differs from the active pod, the target is re-derived before posting: the linked user must still be a member and the pod must still be gated. Any failure refuses in the chat, naming the pod, and posts nothing anywhere. Falling back to the active pod is the defect this rule exists for — the user's answer to B would be authored into A and the agent it names would wake there without B's thread. - an entry with no `podId` (written before this shipped) routes as it always has. Witnesses: the quoted pod receives the reply and the active pod does not; an unquoted message still goes to the active pod; a reply to the active pod's own line is unchanged even with its gate off; a gate-off or membership-gone target is refused by name with nothing posted; a pre-D11 entry still routes to the active pod; Telegram's outbound line carries the pod name. --- .../unit/services/slackBridgeService.test.js | 129 +++++++++++++- .../services/telegramBridgeService.test.js | 165 ++++++++++++++++++ backend/services/slackBridgeService.ts | 75 +++++++- backend/services/telegramBridgeService.ts | 97 +++++++++- 4 files changed, 450 insertions(+), 16 deletions(-) diff --git a/backend/__tests__/unit/services/slackBridgeService.test.js b/backend/__tests__/unit/services/slackBridgeService.test.js index 5fd4e406f..880472c0a 100644 --- a/backend/__tests__/unit/services/slackBridgeService.test.js +++ b/backend/__tests__/unit/services/slackBridgeService.test.js @@ -1,6 +1,9 @@ jest.mock('../../../models/Integration', () => ({ findOne: jest.fn(), findByIdAndUpdate: jest.fn() })); jest.mock('../../../models/Pod', () => ({ findById: jest.fn() })); jest.mock('../../../models/User', () => ({ findById: jest.fn() })); +jest.mock('../../../models/pg/Message', () => ({ create: jest.fn(), findById: jest.fn() })); +jest.mock('../../../services/messageAgentDeliveryService', () => ({ deliverMessageToAgents: jest.fn() })); +jest.mock('../../../config/socket', () => ({ getIO: jest.fn(() => null) })); jest.mock('../../../services/connectorSecrets', () => ({ get: jest.fn() })); // The constructor is stubbed (the network); the escape is the real one, because // the escaping these tests assert is the behaviour we ship. @@ -20,6 +23,8 @@ jest.mock('../../../services/connectorDeliveryFailureService', () => ({ const Integration = require('../../../models/Integration'); const Pod = require('../../../models/Pod'); const User = require('../../../models/User'); +const PGMessage = require('../../../models/pg/Message'); +const { deliverMessageToAgents } = require('../../../services/messageAgentDeliveryService'); const connectorSecrets = require('../../../services/connectorSecrets'); const SlackApi = require('../../../services/slackApi'); const deliveryFailures = require('../../../services/connectorDeliveryFailureService'); @@ -132,7 +137,129 @@ describe('Slack installable bridge', () => { content: 'Can you clarify?', threadTs: '171234.0001', relayMap: [{ externalMessageId: '171234.0001', agentUsername: 'kai' }], - })).toEqual({ content: '@kai Can you clarify?', routedAgent: 'kai' }); + })).toEqual({ content: '@kai Can you clarify?', routedAgent: 'kai', podId: null }); + }); + + test('sends a thread reply into the quoted pod, not the connector\'s active one', async () => { + // ADR-025 D11. The map entry names the pod its line came from; before this, + // the reader kept only the agent and the reply landed in the active pod. + Pod.findById.mockImplementation((id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === 'pod-2' ? { name: 'Launch', type: 'team', members: ['user-1'] } : { name: 'Alpha', type: 'team', members: ['user-1'] }, + ), + }), + })); + PGMessage.create.mockResolvedValue({ id: 'pg-1' }); + PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); + User.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ username: 'sam' }) }), + }); + deliverMessageToAgents.mockResolvedValue(undefined); + Integration.findOne.mockResolvedValue(null); + + const result = await relaySlackMessageToPod({ + integration: { + ...integration, + scope: 'user', + config: { + ...integration.config, + linkedUserId: 'user-1', + slackUserId: 'U1', + gates: { 'pod-2': { enabled: true } }, + relayMap: [{ externalMessageId: '171234.0001', agentUsername: 'kai', podId: 'pod-2' }], + }, + }, + event: { text: 'yes, ship it', user: 'U1', thread_ts: '171234.0001' }, + }); + + expect(result).toEqual({ relayed: true, routedAgent: 'kai' }); + expect(PGMessage.create.mock.calls[0][0]).toBe('pod-2'); + }); + + test('still posts an unquoted Slack message to the active pod', async () => { + Pod.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ name: 'Alpha', type: 'team', members: ['user-1'] }) }), + }); + PGMessage.create.mockResolvedValue({ id: 'pg-1' }); + PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); + User.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ username: 'sam' }) }), + }); + deliverMessageToAgents.mockResolvedValue(undefined); + + await relaySlackMessageToPod({ + integration: { + ...integration, + scope: 'user', + config: { ...integration.config, linkedUserId: 'user-1', slackUserId: 'U1', gates: {} }, + }, + event: { text: 'hello', user: 'U1' }, + }); + + expect(PGMessage.create.mock.calls[0][0]).toBe('pod-1'); + }); + + test('refuses a thread reply into a pod whose gate is off, naming it, posting nothing', async () => { + Pod.findById.mockImplementation((id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === 'pod-2' ? { name: 'Launch', type: 'team', members: ['user-1'] } : { name: 'Alpha', type: 'team', members: ['user-1'] }, + ), + }), + })); + + const result = await relaySlackMessageToPod({ + integration: { + ...integration, + scope: 'user', + config: { + ...integration.config, + linkedUserId: 'user-1', + slackUserId: 'U1', + gates: {}, + relayMap: [{ externalMessageId: '171234.0001', agentUsername: 'kai', podId: 'pod-2' }], + }, + }, + event: { text: 'yes, ship it', user: 'U1', thread_ts: '171234.0001' }, + }); + + expect(result).toEqual({ relayed: false }); + expect(PGMessage.create).not.toHaveBeenCalled(); + expect(deliverMessageToAgents).not.toHaveBeenCalled(); + const api = SlackApi.mock.results[0].value; + expect(api.postMessage).toHaveBeenCalledWith( + 'D1', expect.stringContaining('Launch'), + ); + }); + + test('routes a pre-D11 map entry (no podId) to the active pod', async () => { + Pod.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ name: 'Alpha', type: 'team', members: ['user-1'] }) }), + }); + PGMessage.create.mockResolvedValue({ id: 'pg-1' }); + PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); + User.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ username: 'sam' }) }), + }); + deliverMessageToAgents.mockResolvedValue(undefined); + + await relaySlackMessageToPod({ + integration: { + ...integration, + scope: 'user', + config: { + ...integration.config, + linkedUserId: 'user-1', + slackUserId: 'U1', + gates: {}, + relayMap: [{ externalMessageId: '171234.0001', agentUsername: 'kai' }], + }, + }, + event: { text: 'yes, ship it', user: 'U1', thread_ts: '171234.0001' }, + }); + + expect(PGMessage.create.mock.calls[0][0]).toBe('pod-1'); }); test('does not relay through a visible recovery row whose secret is unavailable', async () => { diff --git a/backend/__tests__/unit/services/telegramBridgeService.test.js b/backend/__tests__/unit/services/telegramBridgeService.test.js index 09d7ac34c..e3851d545 100644 --- a/backend/__tests__/unit/services/telegramBridgeService.test.js +++ b/backend/__tests__/unit/services/telegramBridgeService.test.js @@ -2,6 +2,11 @@ jest.mock('../../../models/Integration', () => ({ findOne: jest.fn(), findByIdAndUpdate: jest.fn(), })); +jest.mock('../../../models/Pod', () => ({ findById: jest.fn() })); +jest.mock('../../../models/User', () => ({ findById: jest.fn() })); +jest.mock('../../../models/pg/Message', () => ({ create: jest.fn(), findById: jest.fn() })); +jest.mock('../../../services/messageAgentDeliveryService', () => ({ deliverMessageToAgents: jest.fn() })); +jest.mock('../../../config/socket', () => ({ getIO: jest.fn(() => null) })); // Stub the network, keep the behaviour: the bridge escapes through this module's // escapeHtml, so a bare stub leaves it undefined and the send is swallowed. jest.mock('../../../services/telegramService', () => ({ @@ -10,11 +15,17 @@ jest.mock('../../../services/telegramService', () => ({ })); const telegramSend = require('../../../services/telegramService'); +const IntegrationModel = require('../../../models/Integration'); +const Pod = require('../../../models/Pod'); +const User = require('../../../models/User'); +const PGMessage = require('../../../models/pg/Message'); +const { deliverMessageToAgents } = require('../../../services/messageAgentDeliveryService'); const { shouldEscalate, routeReplyContent, isRelayableIntegration, isInboundRelayableIntegration, + relayAgentMessageToTelegram, relayTelegramMessageToPod, } = require('../../../services/telegramBridgeService'); @@ -181,3 +192,157 @@ describe('telegramBridgeService — user-scope outbound gate', () => { } }); }); + +// ADR-025 D10/D11. A user-scoped connector has ONE active inbound destination +// (`podId`) and N gated pods. Before this, a quote-reply landed in the active pod +// whichever pod the quoted line came from, so the user's answer to B was authored +// into A and the agent it named woke there without B's thread; and Telegram's +// outbound line carried no pod tag, so a reader could not tell the two apart. +describe('telegramBridgeService — multi-pod routing', () => { + const ORIGINAL_POD = 'pod-a'; + const GATED_POD = 'pod-b'; + const userScoped = (overrides = {}) => ({ + _id: 'i1', + scope: 'user', + podId: ORIGINAL_POD, + type: 'telegram', + isActive: true, + config: { + liveRelay: true, + chatType: 'private', + chatId: 'chat-1', + linkedUserId: 'user-1', + // pod-a is the active inbound destination and its gate is OFF, which is + // legal: inbound does not consult gates (isInboundRelayableIntegration). + gates: { [GATED_POD]: { enabled: true } }, + relayMap: [ + { tgMessageId: '101', agentUsername: 'gene-fix-agent', podId: GATED_POD }, + { tgMessageId: '102', agentUsername: 'lead-agent', podId: ORIGINAL_POD }, + { tgMessageId: '103', agentUsername: 'legacy-agent' }, + ], + ...overrides, + }, + }); + const podDoc = (overrides = {}) => ({ name: 'Alpha', type: 'team', members: ['user-1'], ...overrides }); + + const originalToken = process.env.TELEGRAM_BOT_TOKEN; + + beforeEach(() => { + jest.clearAllMocks(); + process.env.TELEGRAM_BOT_TOKEN = 'bot-token'; + Pod.findById.mockImplementation((id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === GATED_POD ? podDoc({ name: 'Launch' }) : podDoc(), + ), + }), + })); + User.findById.mockReturnValue({ + select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ username: 'sam' }) }), + }); + PGMessage.create.mockResolvedValue({ id: 'pg-1' }); + PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); + deliverMessageToAgents.mockResolvedValue(undefined); + telegramSend.sendMessage.mockResolvedValue({ success: true, messageId: 555 }); + IntegrationModel.findByIdAndUpdate.mockResolvedValue(undefined); + }); + + afterAll(() => { + if (originalToken === undefined) delete process.env.TELEGRAM_BOT_TOKEN; + else process.env.TELEGRAM_BOT_TOKEN = originalToken; + }); + + const inbound = (integration, message) => relayTelegramMessageToPod({ + integration, + telegramMessage: { text: 'looks wrong', ...message }, + }); + + it('returns the quoted line\'s pod, and null for an entry that predates it', () => { + const integration = userScoped(); + const relayMap = integration.config.relayMap; + expect(routeReplyContent({ + content: 'use the v2 schema', replyToTgMessageId: '101', relayMap, + }).podId).toBe(GATED_POD); + expect(routeReplyContent({ + content: 'use the v2 schema', replyToTgMessageId: '103', relayMap, + }).podId).toBeNull(); + }); + + it('tags the outbound line with the pod name and records the pod in the map', async () => { + await relayAgentMessageToTelegram({ + podId: GATED_POD, + agentUsername: 'kai', + displayName: 'Kai', + content: 'deploy is green', + podMessageId: 'pm-1', + integration: userScoped({ relayAllAgentMessages: true }), + }); + + const [, , text] = telegramSend.sendMessage.mock.calls[0]; + expect(text).toContain('[Launch]'); + expect(IntegrationModel.findByIdAndUpdate).toHaveBeenCalledWith('i1', expect.objectContaining({ + $push: expect.objectContaining({ + 'config.relayMap': expect.objectContaining({ + $each: [expect.objectContaining({ podId: GATED_POD, tgMessageId: '555' })], + }), + }), + })); + }); + + it('posts a quote-reply into the quoted pod, not the active one', async () => { + const result = await inbound(userScoped(), { reply_to_message: { message_id: 101 } }); + + expect(result).toEqual({ relayed: true, routedAgent: 'gene-fix-agent' }); + expect(PGMessage.create).toHaveBeenCalledTimes(1); + expect(PGMessage.create.mock.calls[0][0]).toBe(GATED_POD); + }); + + it('still posts an unquoted message to the active pod', async () => { + await inbound(userScoped(), {}); + + expect(PGMessage.create.mock.calls[0][0]).toBe(ORIGINAL_POD); + }); + + it('keeps a quote-reply to the active pod\'s own line on the active pod, gate off', async () => { + // The `:184` rule: inbound never consults the active pod's gate. Routing a + // reply to the pod it already belongs in must not add that requirement. + await inbound(userScoped(), { reply_to_message: { message_id: 102 } }); + + expect(PGMessage.create.mock.calls[0][0]).toBe(ORIGINAL_POD); + expect(telegramSend.sendMessage).not.toHaveBeenCalled(); + }); + + it('refuses a quote-reply into a pod whose gate is off, naming it, posting nothing', async () => { + const integration = userScoped({ gates: {} }); + const result = await inbound(integration, { reply_to_message: { message_id: 101 } }); + + expect(result).toEqual({ relayed: false }); + expect(PGMessage.create).not.toHaveBeenCalled(); + expect(deliverMessageToAgents).not.toHaveBeenCalled(); + expect(telegramSend.sendMessage).toHaveBeenCalledTimes(1); + const [, , text] = telegramSend.sendMessage.mock.calls[0]; + expect(text).toContain('Launch'); + expect(text).toContain('Nothing was posted'); + }); + + it('refuses a quote-reply into a pod the linked user has left', async () => { + Pod.findById.mockImplementation((id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === GATED_POD ? podDoc({ name: 'Launch', members: [] }) : podDoc(), + ), + }), + })); + const result = await inbound(userScoped(), { reply_to_message: { message_id: 101 } }); + + expect(result).toEqual({ relayed: false }); + expect(PGMessage.create).not.toHaveBeenCalled(); + expect(telegramSend.sendMessage.mock.calls[0][2]).toContain('Launch'); + }); + + it('routes an entry with no podId as it always has — the active pod', async () => { + await inbound(userScoped(), { reply_to_message: { message_id: 103 } }); + + expect(PGMessage.create.mock.calls[0][0]).toBe(ORIGINAL_POD); + }); +}); diff --git a/backend/services/slackBridgeService.ts b/backend/services/slackBridgeService.ts index c69786af8..1d13c8df1 100644 --- a/backend/services/slackBridgeService.ts +++ b/backend/services/slackBridgeService.ts @@ -132,6 +132,15 @@ const isInboundRelayableIntegration = (integration: SlackIntegrationDoc, podId: const NO_ACTIVE_POD_REPLY = 'This connector has no active pod. Choose one in Commonly first.'; +// ADR-025 D11: a quote-reply whose pod is no longer reachable is refused in the +// chat and posted nowhere. Naming the pod is the whole point — "your reply did +// not go through" is unactionable, while "that line came from Launch" tells the +// user which pod the fix is in. +const routedPodRefusal = (podLabel: string): string => ( + `⚠️ That line came from “${podLabel}”, which this chat no longer reaches. ` + + 'Nothing was posted — open Commonly to reply there.' +); + const replyNoActivePod = async (integration: SlackIntegrationDoc): Promise => { const chatId = integration.config?.chatId; const botTokenRef = integration.config?.botTokenRef; @@ -146,6 +155,24 @@ const replyNoActivePod = async (integration: SlackIntegrationDoc): Promise } }; +// Same trust level as replyNoActivePod: a refusal that never reaches the chat +// leaves the user believing their message was relayed. +const replyRoutedPodRefused = async ( + integration: SlackIntegrationDoc, + podLabel: string, +): Promise => { + const chatId = integration.config?.chatId; + const botTokenRef = integration.config?.botTokenRef; + if (!chatId || !botTokenRef) return; + try { + const token = await connectorSecrets.get(String(botTokenRef)); + const sent = await new SlackApi(token).postMessage(String(chatId), routedPodRefusal(podLabel)); + await deliveryFailures.noteBoundChatDeliveryFailure(integration, chatId, sent); + } catch (error) { + console.warn('[slack-bridge] could not send routed-pod refusal:', (error as Error).message); + } +}; + const findLiveIntegration = async (podId: string): Promise => ( Integration.findOne({ type: 'slack', @@ -276,6 +303,10 @@ export const relayAgentMessageToSlack = async (opts: { // D11: a Slack thread attached to a relayed line is a direct answer to that // line's agent. Keep the map generic so Telegram can migrate from tgMessageId // without changing this reader. +// +// It also returns the pod the quoted line came from. The caller owns the +// decision — the map is data, and this function does no lookups — so a `podId` +// here means "the quoted entry names a pod", not "routing to it is allowed". export const routeSlackReplyContent = (opts: { content: string; threadTs?: string | null; @@ -283,16 +314,18 @@ export const routeSlackReplyContent = (opts: { externalMessageId?: string; tgMessageId?: string; agentUsername?: string; + podId?: string | null; }>; -}): { content: string; routedAgent: string | null } => { +}): { content: string; routedAgent: string | null; podId: string | null } => { const { content, threadTs, relayMap } = opts; - if (!threadTs || !Array.isArray(relayMap)) return { content, routedAgent: null }; + if (!threadTs || !Array.isArray(relayMap)) return { content, routedAgent: null, podId: null }; const hit = relayMap.find((entry) => String(entry.externalMessageId || entry.tgMessageId) === String(threadTs)); - if (!hit?.agentUsername) return { content, routedAgent: null }; + if (!hit?.agentUsername) return { content, routedAgent: null, podId: null }; + const routedPodId = hit.podId ? String(hit.podId) : null; const mention = `@${hit.agentUsername}`; return content.toLowerCase().includes(mention.toLowerCase()) - ? { content, routedAgent: hit.agentUsername } - : { content: `${mention} ${content}`, routedAgent: hit.agentUsername }; + ? { content, routedAgent: hit.agentUsername, podId: routedPodId } + : { content: `${mention} ${content}`, routedAgent: hit.agentUsername, podId: routedPodId }; }; // Inbound Slack DM → Commonly pod. The event route has already proven the @@ -360,12 +393,40 @@ export const relaySlackMessageToPod = async (opts: { await replyNoActivePod(integration); return { relayed: false }; } - const { content: routedText, routedAgent } = routeSlackReplyContent({ + const routed = routeSlackReplyContent({ content: rawText, threadTs: event.thread_ts, relayMap: cardReply.lateReply ? [] : config.relayMap, }); - const podId = cardReply.lateReply?.podId || String(integration.podId); + // ADR-025 D11: a thread reply answers the line it quotes, so it belongs in THAT + // pod — not in whichever pod is this connector's active destination. The quoted + // pod is re-derived here rather than trusted from the map: the map is written + // at send time and an entry can outlive its gate, its pod, or the owner's + // membership (100-entry cap, owner-editable gates). Any failure refuses in the + // chat and posts nothing. Falling back to the active pod is the defect this + // rule exists for — the user's answer to B would be authored into A and the + // agent it names would wake there without B's thread. + // + // An entry with no `podId` is not this case: it was written before multi-pod + // routing shipped, carries no pod to check, and routes as it always has. + let podId = cardReply.lateReply?.podId || String(integration.podId); + if (routed.podId && String(routed.podId) !== String(podId)) { + const routedPod = await Pod.findById(routed.podId).select('name type createdBy members').lean(); + if (!routedPod + || !isPodMember(routedPod, String(config.linkedUserId)) + || !isRelayableIntegration(integration, routed.podId)) { + console.warn( + `[slack-bridge] thread reply refused — quoted pod ${routed.podId} is no longer routed to this chat`, + ); + await replyRoutedPodRefused( + integration, + routedPod?.name ? String(routedPod.name) : `pod ${routed.podId}`, + ); + return { relayed: false }; + } + podId = String(routed.podId); + } + const { content: routedText, routedAgent } = routed; const replyToMessageId = cardReply.lateReply?.messageId || null; const linkedUserId = String(config.linkedUserId); const senderName = event.user_profile?.display_name || event.user_profile?.real_name; diff --git a/backend/services/telegramBridgeService.ts b/backend/services/telegramBridgeService.ts index 5af9b46f1..bcdc92dd4 100644 --- a/backend/services/telegramBridgeService.ts +++ b/backend/services/telegramBridgeService.ts @@ -41,6 +41,10 @@ export interface RelayMapEntry { tgMessageId: string; agentUsername: string; podMessageId?: string | null; + // ADR-025 D11: which pod's line this was. Absent on entries written before + // multi-pod routing shipped, and absent is not "the active pod" — it is + // "unknown", and unknown routes as it always has (see relayTelegramMessageToPod). + podId?: string | null; } interface TelegramIntegrationDoc { @@ -121,22 +125,28 @@ export const renderTelegramDecisionCard = (opts: { // Prefix an inbound Telegram quote-reply with the @mention of the agent whose // relayed line was quoted, so the normal mention pipeline routes it. Pure — // unit-tested without any I/O. +// +// It also returns the pod the quoted line came from (ADR-025 D11). The caller +// owns the decision — the map is data, and this function does no lookups — so a +// `podId` here means "the quoted entry names a pod", not "routing to it is +// allowed". export const routeReplyContent = (opts: { content: string; replyToTgMessageId?: string | null; relayMap?: RelayMapEntry[]; -}): { content: string; routedAgent: string | null } => { +}): { content: string; routedAgent: string | null; podId: string | null } => { const { content, replyToTgMessageId, relayMap } = opts; if (!replyToTgMessageId || !Array.isArray(relayMap)) { - return { content, routedAgent: null }; + return { content, routedAgent: null, podId: null }; } const hit = relayMap.find((e) => String(e.tgMessageId) === String(replyToTgMessageId)); - if (!hit || !hit.agentUsername) return { content, routedAgent: null }; + if (!hit || !hit.agentUsername) return { content, routedAgent: null, podId: null }; + const routedPodId = hit.podId ? String(hit.podId) : null; const mention = `@${hit.agentUsername}`; if (content.toLowerCase().includes(mention.toLowerCase())) { - return { content, routedAgent: hit.agentUsername }; + return { content, routedAgent: hit.agentUsername, podId: routedPodId }; } - return { content: `${mention} ${content}`, routedAgent: hit.agentUsername }; + return { content: `${mention} ${content}`, routedAgent: hit.agentUsername, podId: routedPodId }; }; const findLiveIntegration = async (podId: unknown): Promise => { @@ -197,6 +207,16 @@ const isInboundRelayableIntegration = ( const NO_ACTIVE_POD_REPLY = 'This connector has no active pod. Choose one in Commonly first.'; +// ADR-025 D11: a quote-reply whose pod is no longer reachable is refused in the +// chat and posted nowhere. Naming the pod is the whole point — "your reply did +// not go through" is unactionable, while "that line came from Launch" tells the +// user which pod the fix is in. The label is escaped because this send is +// `parse_mode: 'HTML'` and a pod name is agent-choosable. +const routedPodRefusal = (podLabel: string): string => ( + `⚠️ That line came from “${escapeHtml(podLabel)}”, which this chat no longer reaches. ` + + 'Nothing was posted — open Commonly to reply there.' +); + const replyNoActivePod = async (integration: TelegramIntegrationDoc): Promise => { const botToken = process.env.TELEGRAM_BOT_TOKEN; const chatId = integration.config?.chatId; @@ -210,6 +230,40 @@ const replyNoActivePod = async (integration: TelegramIntegrationDoc): Promise => { + const botToken = process.env.TELEGRAM_BOT_TOKEN; + const chatId = integration.config?.chatId; + if (!botToken || !chatId) return; + try { + const sent = await telegramSend.sendMessage(botToken, chatId, routedPodRefusal(podLabel)); + await deliveryFailures.noteBoundChatDeliveryFailure(integration, chatId, sent); + } catch (error) { + console.warn('[tg-bridge] could not send routed-pod refusal:', (error as Error).message); + } +}; + +// The pod tag is cosmetic (ADR-025 D10): it tells a reader which pod a line came +// from once more than one pod reaches this chat. Slack has carried it since its +// bridge shipped; Telegram did not, which is half of what D10 names. A lookup +// failure must not cost the relay, so it degrades to the id rather than to a +// name that was never read — never to silence. +const podNameForTag = async (podId: string): Promise => { + try { + // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require + const Pod = require('../models/Pod'); + const pod = await Pod.findById(podId).select('name').lean(); + return String(pod?.name || podId); + } catch (error) { + console.warn('[tg-bridge] pod name lookup failed, tagging with the id:', (error as Error).message); + return podId; + } +}; + // Outbound: agent message → Telegram, attributed, with a deep link back into // the pod. Fire-and-forget from AgentMessageService.postMessage — a bridge // failure must never fail the post itself. @@ -282,7 +336,7 @@ export const relayAgentMessageToTelegram = async (opts: { const body = escapeHtml(String(content).slice(0, OUTBOUND_TEXT_CAP)); const text = opts.card ? renderTelegramDecisionCard({ card: opts.card, displayName, agentUsername, link }) - : `${escapeHtml(displayName || agentUsername)}: ${body}` + : `[${escapeHtml(await podNameForTag(podId))}] ${escapeHtml(displayName || agentUsername)}: ${body}` + `\n\nopen in Commonly`; const result = await telegramSend.sendMessage(botToken, chatId, text); @@ -297,7 +351,7 @@ export const relayAgentMessageToTelegram = async (opts: { await IntegrationModel.findByIdAndUpdate(integration._id, { $push: { 'config.relayMap': { - $each: [{ tgMessageId, agentUsername, podMessageId: podMessageId || null }], + $each: [{ tgMessageId, agentUsername, podMessageId: podMessageId || null, podId }], $slice: -RELAY_MAP_CAP, }, ...(opts.card && podMessageId ? { @@ -430,11 +484,38 @@ export const relayTelegramMessageToPod = async (opts: { } if (cardReply.lateReply) podId = cardReply.lateReply.podId; const replyToMessageId = cardReply.lateReply?.messageId || null; - const { content: routedText, routedAgent } = routeReplyContent({ + const routed = routeReplyContent({ content: rawText, replyToTgMessageId: replyToTgMessageId != null ? String(replyToTgMessageId) : null, relayMap: cardReply.lateReply ? [] : integration.config?.relayMap, }); + // ADR-025 D11: a quote-reply answers the line it quotes, so it belongs in THAT + // pod — not in whichever pod is this connector's active destination. The quoted + // pod is re-derived here rather than trusted from the map: the map is written + // at send time and an entry can outlive its gate, its pod, or the owner's + // membership (100-entry cap, owner-editable gates). Any failure refuses in the + // chat and posts nothing. Falling back to the active pod is the defect this + // rule exists for — the user's answer to B would be authored into A and the + // agent it names would wake there without B's thread. + // + // An entry with no `podId` is not this case. It was written before multi-pod + // routing shipped, carries no pod to check, and routes as it always has. + if (routed.podId && String(routed.podId) !== String(podId)) { + // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require + const PodModel = require('../models/Pod'); + const routedPod = await PodModel.findById(routed.podId).select('name type createdBy members').lean(); + if (!routedPod + || !isPodMember(routedPod, linkedUserId) + || !isRelayableIntegration(integration, routed.podId)) { + console.warn( + `[tg-bridge] quote-reply refused — quoted pod ${routed.podId} is no longer routed to this chat`, + ); + await replyRoutedPodRefused(integration, routedPod?.name ? String(routedPod.name) : `pod ${routed.podId}`); + return { relayed: false }; + } + podId = String(routed.podId); + } + const { content: routedText, routedAgent } = routed; const senderName = [ telegramMessage.from?.first_name, From 562a147f43ee863d38ce044a39b9b99871ae10d5 Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:49:01 -0700 Subject: [PATCH 2/4] test(connectors): tidy the new fixtures so no diagnostic lands on a line this PR adds MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Formatting only, assertions unchanged: multiline object literals, a destructured fixture, and a one-line `toHaveBeenCalledWith`. What remains on added lines is `import/no-unresolved` + `import/extensions` on the new `require(...)` calls, the same pair every pre-existing require in these two files already reports — the backend `.js` corpus is not linted (2,279 errors, 277 of these files). --- .../unit/services/slackBridgeService.test.js | 49 ++++++++++--------- .../services/telegramBridgeService.test.js | 9 +++- 2 files changed, 32 insertions(+), 26 deletions(-) diff --git a/backend/__tests__/unit/services/slackBridgeService.test.js b/backend/__tests__/unit/services/slackBridgeService.test.js index 880472c0a..69ad9c2d2 100644 --- a/backend/__tests__/unit/services/slackBridgeService.test.js +++ b/backend/__tests__/unit/services/slackBridgeService.test.js @@ -140,16 +140,27 @@ describe('Slack installable bridge', () => { })).toEqual({ content: '@kai Can you clarify?', routedAgent: 'kai', podId: null }); }); + // Two fixture shapes, because the routed pod and the active pod are different + // documents: the first names the pod that may not be reachable any more. + const podWithMembers = (name) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue({ name, type: 'team', members: ['user-1'] }), + }), + }); + const podsById = (gatedPodName) => (id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === 'pod-2' + ? { name: gatedPodName, type: 'team', members: ['user-1'] } + : { name: 'Alpha', type: 'team', members: ['user-1'] }, + ), + }), + }); + test('sends a thread reply into the quoted pod, not the connector\'s active one', async () => { // ADR-025 D11. The map entry names the pod its line came from; before this, // the reader kept only the agent and the reply landed in the active pod. - Pod.findById.mockImplementation((id) => ({ - select: jest.fn().mockReturnValue({ - lean: jest.fn().mockResolvedValue( - String(id) === 'pod-2' ? { name: 'Launch', type: 'team', members: ['user-1'] } : { name: 'Alpha', type: 'team', members: ['user-1'] }, - ), - }), - })); + Pod.findById.mockImplementation(podsById('Launch')); PGMessage.create.mockResolvedValue({ id: 'pg-1' }); PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); User.findById.mockReturnValue({ @@ -178,9 +189,7 @@ describe('Slack installable bridge', () => { }); test('still posts an unquoted Slack message to the active pod', async () => { - Pod.findById.mockReturnValue({ - select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ name: 'Alpha', type: 'team', members: ['user-1'] }) }), - }); + Pod.findById.mockReturnValue(podWithMembers('Alpha')); PGMessage.create.mockResolvedValue({ id: 'pg-1' }); PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); User.findById.mockReturnValue({ @@ -192,7 +201,9 @@ describe('Slack installable bridge', () => { integration: { ...integration, scope: 'user', - config: { ...integration.config, linkedUserId: 'user-1', slackUserId: 'U1', gates: {} }, + config: { + ...integration.config, linkedUserId: 'user-1', slackUserId: 'U1', gates: {}, + }, }, event: { text: 'hello', user: 'U1' }, }); @@ -201,13 +212,7 @@ describe('Slack installable bridge', () => { }); test('refuses a thread reply into a pod whose gate is off, naming it, posting nothing', async () => { - Pod.findById.mockImplementation((id) => ({ - select: jest.fn().mockReturnValue({ - lean: jest.fn().mockResolvedValue( - String(id) === 'pod-2' ? { name: 'Launch', type: 'team', members: ['user-1'] } : { name: 'Alpha', type: 'team', members: ['user-1'] }, - ), - }), - })); + Pod.findById.mockImplementation(podsById('Launch')); const result = await relaySlackMessageToPod({ integration: { @@ -228,15 +233,11 @@ describe('Slack installable bridge', () => { expect(PGMessage.create).not.toHaveBeenCalled(); expect(deliverMessageToAgents).not.toHaveBeenCalled(); const api = SlackApi.mock.results[0].value; - expect(api.postMessage).toHaveBeenCalledWith( - 'D1', expect.stringContaining('Launch'), - ); + expect(api.postMessage).toHaveBeenCalledWith('D1', expect.stringContaining('Launch')); }); test('routes a pre-D11 map entry (no podId) to the active pod', async () => { - Pod.findById.mockReturnValue({ - select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue({ name: 'Alpha', type: 'team', members: ['user-1'] }) }), - }); + Pod.findById.mockReturnValue(podWithMembers('Alpha')); PGMessage.create.mockResolvedValue({ id: 'pg-1' }); PGMessage.findById.mockResolvedValue({ id: 'pg-1', content: 'relayed' }); User.findById.mockReturnValue({ diff --git a/backend/__tests__/unit/services/telegramBridgeService.test.js b/backend/__tests__/unit/services/telegramBridgeService.test.js index e3851d545..eb9d97e37 100644 --- a/backend/__tests__/unit/services/telegramBridgeService.test.js +++ b/backend/__tests__/unit/services/telegramBridgeService.test.js @@ -223,7 +223,12 @@ describe('telegramBridgeService — multi-pod routing', () => { ...overrides, }, }); - const podDoc = (overrides = {}) => ({ name: 'Alpha', type: 'team', members: ['user-1'], ...overrides }); + const podDoc = (overrides = {}) => ({ + name: 'Alpha', + type: 'team', + members: ['user-1'], + ...overrides, + }); const originalToken = process.env.TELEGRAM_BOT_TOKEN; @@ -259,7 +264,7 @@ describe('telegramBridgeService — multi-pod routing', () => { it('returns the quoted line\'s pod, and null for an entry that predates it', () => { const integration = userScoped(); - const relayMap = integration.config.relayMap; + const { relayMap } = integration.config; expect(routeReplyContent({ content: 'use the v2 schema', replyToTgMessageId: '101', relayMap, }).podId).toBe(GATED_POD); From db2774f6fef4103ae1f068032255de6dffab9949 Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:55:07 -0700 Subject: [PATCH 3/4] refactor(connectors): one home for the gate reading and for gate+membership (TASK-156, wren 74618) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wren's spec fix, applied before the gate: the routed-reply rule was specified for two implementations and would have been written twice, which is how the list and the call drifted in TASK-146. - `connectorRelayPolicy.isGatedPodTarget(integration, podId)` is now the one reading of `config.gates` — the user-scope gate / pod-scope own-pod ternary that both bridges' `isRelayableIntegration` and `decisionCardReconcileService` each carried a copy of. Four consumers, one reading. - `connectorRelayPolicy.isRoutedPodTarget({integration, pod, podId, userId})` is the gate + membership conjunction, and both bridges' routed-reply check calls it instead of restating it. The protocol-health half stays where it belongs — in each bridge's own predicate, since chatType/teamId/chatId differ by protocol. - The check-then-act window is now named and accepted in the predicate's docstring, with what closing it would take, rather than left unstated. - Both call sites state that the ACTIVE pod is exempt from the gate by design, so "re-check the gate" cannot be read as universal. - `decisionCardRelay.bridges.test.js` stubbed the whole relay-policy module with a factory returning only `shouldEscalate`; it now spreads `jest.requireActual`, or the new exports would be undefined in the bridges under test. Witnesses: a new `connectorRelayPolicy.test.js` pins both predicates (both scopes, absent/disabled/non-boolean gates, pod ids compared by value, the membership conjunction, the missing-user-id guard). Six suites, 135 tests green. --- .../services/connectorRelayPolicy.test.js | 87 +++++++++++++++++++ .../decisionCardRelay.bridges.test.js | 8 +- backend/services/connectorRelayPolicy.ts | 59 ++++++++++++- .../services/decisionCardReconcileService.ts | 10 ++- backend/services/slackBridgeService.ts | 24 +++-- backend/services/telegramBridgeService.ts | 24 +++-- 6 files changed, 191 insertions(+), 21 deletions(-) create mode 100644 backend/__tests__/unit/services/connectorRelayPolicy.test.js diff --git a/backend/__tests__/unit/services/connectorRelayPolicy.test.js b/backend/__tests__/unit/services/connectorRelayPolicy.test.js new file mode 100644 index 000000000..e04a9d50a --- /dev/null +++ b/backend/__tests__/unit/services/connectorRelayPolicy.test.js @@ -0,0 +1,87 @@ +// TASK-156 spec fix (wren 74618): the gate reading and the gate+membership +// conjunction have one home, because three copies of the same ternary is how a +// list and a call drift apart. These are that home's own arms. +const { + isGatedPodTarget, + isRoutedPodTarget, +} = require('../../../services/connectorRelayPolicy'); + +const userScoped = (gates) => ({ + scope: 'user', + type: 'telegram', + podId: 'active-pod', + config: { gates }, +}); +const podScoped = (podId, gates) => ({ + scope: 'pod', + type: 'telegram', + podId, + config: { gates }, +}); +const pod = ({ members = [], createdBy } = {}) => ({ members, createdBy }); + +describe('connectorRelayPolicy — the gate reading', () => { + it('keys a user-scoped connector on its gate for that pod', () => { + const integration = userScoped({ + 'pod-a': { enabled: true }, + 'pod-b': { enabled: false }, + }); + expect(isGatedPodTarget(integration, 'pod-a')).toBe(true); + expect(isGatedPodTarget(integration, 'pod-b')).toBe(false); + // Absent key, and a key that exists without `enabled: true`, are both closed. + expect(isGatedPodTarget(integration, 'pod-c')).toBe(false); + expect(isGatedPodTarget(userScoped({ 'pod-a': {} }), 'pod-a')).toBe(false); + expect(isGatedPodTarget(userScoped({ 'pod-a': { enabled: 'yes' } }), 'pod-a')).toBe(false); + expect(isGatedPodTarget(userScoped(undefined), 'pod-a')).toBe(false); + }); + + it('keys a pod-scoped connector on its own pod and never on gates', () => { + expect(isGatedPodTarget(podScoped('pod-a', undefined), 'pod-a')).toBe(true); + expect(isGatedPodTarget(podScoped('pod-a', undefined), 'pod-b')).toBe(false); + // Two identifier spaces for the same allow-list: a pod-scoped connector's + // gates object is not a second switch, and must not be read as one. + expect(isGatedPodTarget(podScoped('pod-a', { 'pod-a': { enabled: true } }), 'pod-b')).toBe(false); + expect(isGatedPodTarget(podScoped('pod-a', { 'pod-b': { enabled: true } }), 'pod-b')).toBe(false); + }); + + it('compares pod ids by value, not by identity', () => { + const integration = podScoped('pod-a', undefined); + integration.podId = { toString: () => 'pod-a' }; + expect(isGatedPodTarget(integration, 'pod-a')).toBe(true); + }); +}); + +describe('connectorRelayPolicy — the routed-target conjunction', () => { + const target = (overrides = {}) => ({ + integration: userScoped({ 'pod-a': { enabled: true } }), + pod: pod({ members: ['user-1'] }), + podId: 'pod-a', + userId: 'user-1', + ...overrides, + }); + + it('admits a gated pod the linked user is still in', () => { + expect(isRoutedPodTarget(target())).toBe(true); + }); + + it('needs BOTH halves — the gate off, or the membership gone, is a refusal', () => { + expect(isRoutedPodTarget(target({ + integration: userScoped({ 'pod-a': { enabled: false } }), + }))).toBe(false); + expect(isRoutedPodTarget(target({ pod: pod({ members: ['someone-else'] }) }))).toBe(false); + expect(isRoutedPodTarget(target({ pod: undefined }))).toBe(false); + }); + + it('admits a pod whose creator is not listed in members', () => { + expect(isRoutedPodTarget(target({ pod: pod({ createdBy: 'user-1' }) }))).toBe(true); + }); + + it('refuses a missing user id rather than stringifying it into a match', () => { + // `String(undefined)` is the truthy string 'undefined'; a caller that passed + // that through would reach the membership read with a user who cannot exist. + // Not exploitable, but the guard is what keeps this half honest. + expect(isRoutedPodTarget(target({ userId: undefined }))).toBe(false); + expect(isRoutedPodTarget(target({ userId: null }))).toBe(false); + expect(isRoutedPodTarget(target({ userId: '' }))).toBe(false); + }); +}); diff --git a/backend/__tests__/unit/services/decisionCardRelay.bridges.test.js b/backend/__tests__/unit/services/decisionCardRelay.bridges.test.js index bef94f5d0..484dc3a3c 100644 --- a/backend/__tests__/unit/services/decisionCardRelay.bridges.test.js +++ b/backend/__tests__/unit/services/decisionCardRelay.bridges.test.js @@ -18,7 +18,13 @@ jest.mock('../../../services/slackApi', () => { mock.escapeSlackMrkdwn = actual.escapeSlackMrkdwn; return mock; }); -jest.mock('../../../services/connectorRelayPolicy', () => ({ shouldEscalate: jest.fn(() => false) })); +// Delegate to the real module: the bridges also read isGatedPodTarget / +// isRoutedPodTarget from here, and a factory that stubs the whole module makes +// them undefined (TASK-156). +jest.mock('../../../services/connectorRelayPolicy', () => ({ + ...jest.requireActual('../../../services/connectorRelayPolicy'), + shouldEscalate: jest.fn(() => false), +})); jest.mock('../../../services/channelVerdictService', () => ({ record: jest.fn() })); const Integration = require('../../../models/Integration'); diff --git a/backend/services/connectorRelayPolicy.ts b/backend/services/connectorRelayPolicy.ts index 5e597b47e..b6b0d72dc 100644 --- a/backend/services/connectorRelayPolicy.ts +++ b/backend/services/connectorRelayPolicy.ts @@ -1,8 +1,13 @@ -// Shared, deterministic interruption policy for connector bridges. Keeping it -// outside a provider means Telegram and Slack cannot silently diverge about -// which pod messages are allowed to interrupt a person's attention surface. +// Shared, deterministic interruption and target policy for connector bridges. +// Keeping it outside a provider means Telegram and Slack cannot silently diverge +// about which pod messages are allowed to interrupt a person's attention surface, +// or about which pods a connector may address at all. + +const isPodMember = require('../utils/isPodMember'); export interface RelayPolicyIntegration { + scope?: string; + podId?: unknown; config?: { relayAllAgentMessages?: boolean; leadAgentUsername?: string; @@ -42,4 +47,50 @@ export const shouldEscalate = (opts: { return ESCALATION_MARKERS.test(content) || QUESTION_AT_HUMAN.test(content); }; -module.exports = { shouldEscalate }; +// One reading of `config.gates`, shared by every consumer that asks whether a pod +// is still a target of a connector: both bridges' outbound relay, decision-card +// delivery, and the routed-quote-reply check below. It was three copies of the +// same ternary before TASK-156, which is how the list and the call drift apart. +// +// Scope is the connector's, not the pod's: a `user`-scoped connector holds one +// private chat and subscribes to N pods through its gates, while a pod-scoped one +// speaks for exactly its own pod. +export const isGatedPodTarget = ( + integration: RelayPolicyIntegration, + podId: string, +): boolean => ( + integration.scope === 'user' + ? integration.config?.gates?.[String(podId)]?.enabled === true + : String(integration.podId) === String(podId) +); + +// May this connector address `pod` on behalf of `userId`? Gate and membership in +// one predicate, because a caller that holds one half and not the other is the +// failure this exists for: the gate says the connector is still subscribed, the +// membership says the person it speaks for is still in the room. +// +// KNOWN WINDOW, accepted: this is check-then-act. A gate switched off between +// this call and the write still lets that one message through. Closing it means +// making the write itself carry the condition (a conditional update or a +// transaction) rather than a preceding read; not worth it for a one-message +// window on a subscription toggle the owner is the only one who can flip. +// +// The ACTIVE pod is deliberately exempt from the gate and never passes through +// here — see isInboundRelayableIntegration in each bridge. A pruned outbound gate +// must not disconnect the owner's own private chat, so the gate bounds the N +// subscribed pods and never the active one. +export const isRoutedPodTarget = (opts: { + integration: RelayPolicyIntegration; + pod: any; + podId: string; + userId: unknown; +}): boolean => { + const { + integration, pod, podId, userId, + } = opts; + return Boolean(userId) + && isGatedPodTarget(integration, podId) + && isPodMember(pod, userId); +}; + +module.exports = { shouldEscalate, isGatedPodTarget, isRoutedPodTarget }; diff --git a/backend/services/decisionCardReconcileService.ts b/backend/services/decisionCardReconcileService.ts index b9f59608a..ad003a1f5 100644 --- a/backend/services/decisionCardReconcileService.ts +++ b/backend/services/decisionCardReconcileService.ts @@ -9,6 +9,8 @@ const Pod = require('../models/Pod'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require const isPodMember = require('../utils/isPodMember'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require +const { isGatedPodTarget } = require('./connectorRelayPolicy'); +// eslint-disable-next-line @typescript-eslint/no-require-imports, global-require const telegramSend = require('./telegramService'); const deliveryFailures = require('./connectorDeliveryFailureService'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require @@ -81,8 +83,12 @@ const canSendClosingLine = ( if (integration.isActive !== true || integration.status === 'error') return false; if (integration.config?.liveRelay !== true || integration.config.adminPause) return false; if (integration.type !== 'telegram' && integration.type !== 'slack') return false; - if (integration.scope === 'user' && integration.config.gates?.[podId]?.enabled !== true) return false; - if (integration.scope !== 'user' && String(integration.podId) !== podId) return false; + // The gate reading is shared with both bridges and with the routed-reply check + // (connectorRelayPolicy.isGatedPodTarget) so the four consumers cannot drift. + // Its companion conjunction lives in isRoutedPodTarget; this function keeps the + // membership read below where it is, because the mute and chatId checks sit + // between the two halves here. + if (!isGatedPodTarget(integration, podId)) return false; const mutedUntil = integration.config.relayMutedUntil; if (mutedUntil && new Date(mutedUntil).getTime() > now.getTime()) return false; const linkedUserId = memberIdFor(integration); diff --git a/backend/services/slackBridgeService.ts b/backend/services/slackBridgeService.ts index 1d13c8df1..7ca6dac48 100644 --- a/backend/services/slackBridgeService.ts +++ b/backend/services/slackBridgeService.ts @@ -10,7 +10,7 @@ const isPodMember = require('../utils/isPodMember'); const connectorSecrets = require('./connectorSecrets'); const deliveryFailures = require('./connectorDeliveryFailureService'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require -const { shouldEscalate } = require('./connectorRelayPolicy'); +const { shouldEscalate, isGatedPodTarget, isRoutedPodTarget } = require('./connectorRelayPolicy'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require const channelVerdictService = require('./channelVerdictService'); import type { DecisionRelayCard } from './decisionCardRelay'; @@ -59,9 +59,7 @@ interface SlackIntegrationDoc { } const isRelayableIntegration = (integration: SlackIntegrationDoc, podId: string): boolean => ( - (integration.scope === 'user' - ? integration.config?.gates?.[String(podId)]?.enabled === true - : String(integration.podId) === String(podId)) + isGatedPodTarget(integration, podId) && integration.type === 'slack' && integration.isActive === true && integration.status !== 'error' @@ -409,12 +407,24 @@ export const relaySlackMessageToPod = async (opts: { // // An entry with no `podId` is not this case: it was written before multi-pod // routing shipped, carries no pod to check, and routes as it always has. + // + // The predicate is `isRoutedPodTarget` (gate + membership, one home in + // connectorRelayPolicy): the same rule the outbound relay and decision-card + // delivery read, so a fix to the rule reaches all of them. Note what it does + // NOT bound — the ACTIVE pod is exempt from the gate by design (see + // isInboundRelayableIntegration above), so this check applies only to the + // quoted pod, and only when it differs from the active one. The shared + // predicate answers gate + membership; the bridge's own predicate adds the + // protocol-health conditions only Slack knows (liveRelay, chatType, teamId). let podId = cardReply.lateReply?.podId || String(integration.podId); if (routed.podId && String(routed.podId) !== String(podId)) { const routedPod = await Pod.findById(routed.podId).select('name type createdBy members').lean(); - if (!routedPod - || !isPodMember(routedPod, String(config.linkedUserId)) - || !isRelayableIntegration(integration, routed.podId)) { + if (!isRoutedPodTarget({ + integration, + pod: routedPod, + podId: routed.podId, + userId: config.linkedUserId, + }) || !isRelayableIntegration(integration, routed.podId)) { console.warn( `[slack-bridge] thread reply refused — quoted pod ${routed.podId} is no longer routed to this chat`, ); diff --git a/backend/services/telegramBridgeService.ts b/backend/services/telegramBridgeService.ts index bcdc92dd4..7e58f44fc 100644 --- a/backend/services/telegramBridgeService.ts +++ b/backend/services/telegramBridgeService.ts @@ -26,7 +26,7 @@ const isPodMember = require('../utils/isPodMember'); const telegramSend = require('./telegramService'); const deliveryFailures = require('./connectorDeliveryFailureService'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require -const { shouldEscalate } = require('./connectorRelayPolicy'); +const { shouldEscalate, isGatedPodTarget, isRoutedPodTarget } = require('./connectorRelayPolicy'); // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require const channelVerdictService = require('./channelVerdictService'); import type { DecisionRelayCard } from './decisionCardRelay'; @@ -178,9 +178,7 @@ const isRelayableIntegration = ( integration: TelegramIntegrationDoc, podId: string, ): boolean => ( - (integration.scope === 'user' - ? integration.config?.gates?.[String(podId)]?.enabled === true - : String(integration.podId) === String(podId)) + isGatedPodTarget(integration, podId) && (integration.type === undefined || integration.type === 'telegram') && integration.isActive !== false && integration.config?.liveRelay === true @@ -500,13 +498,25 @@ export const relayTelegramMessageToPod = async (opts: { // // An entry with no `podId` is not this case. It was written before multi-pod // routing shipped, carries no pod to check, and routes as it always has. + // + // The predicate is `isRoutedPodTarget` (gate + membership, one home in + // connectorRelayPolicy): the same rule the outbound relay and decision-card + // delivery read, so a fix to the rule reaches all of them. Note what it does + // NOT bound — the ACTIVE pod is exempt from the gate by design (see + // isInboundRelayableIntegration above), so this check applies only to the + // quoted pod, and only when it differs from the active one. The shared + // predicate answers gate + membership; the bridge's own predicate adds the + // protocol-health conditions only Telegram knows (liveRelay, chatType, chatId). if (routed.podId && String(routed.podId) !== String(podId)) { // eslint-disable-next-line @typescript-eslint/no-require-imports, global-require const PodModel = require('../models/Pod'); const routedPod = await PodModel.findById(routed.podId).select('name type createdBy members').lean(); - if (!routedPod - || !isPodMember(routedPod, linkedUserId) - || !isRelayableIntegration(integration, routed.podId)) { + if (!isRoutedPodTarget({ + integration, + pod: routedPod, + podId: routed.podId, + userId: linkedUserId, + }) || !isRelayableIntegration(integration, routed.podId)) { console.warn( `[tg-bridge] quote-reply refused — quoted pod ${routed.podId} is no longer routed to this chat`, ); From 25652999925176fe1e36765b25f4cc66413f1faa Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:56:57 -0700 Subject: [PATCH 4/4] test(connectors): witness Slack's membership re-check, and drop a guard the ledger showed was redundant (TASK-156) Two findings from the mutation run on the refactor, both acted on rather than disclosed as noise: - Dropping the shared predicate from Slack's routed check reddened NOTHING, because the gate half is also carried by `isRelayableIntegration` there and the suite had no membership-gone arm on the Slack path (Telegram had one). Added: gate on, linked user no longer a member, refusal by name, nothing posted. - Restoring the `Boolean(userId)` guard in `isRoutedPodTarget` reddened nothing: `isPodMember` already fails closed on a falsy id. The guard is removed rather than kept unwitnessed; the predicate's own arm still pins the behaviour. --- .../unit/services/slackBridgeService.test.js | 36 +++++++++++++++++++ backend/services/connectorRelayPolicy.ts | 6 ++-- 2 files changed, 40 insertions(+), 2 deletions(-) diff --git a/backend/__tests__/unit/services/slackBridgeService.test.js b/backend/__tests__/unit/services/slackBridgeService.test.js index 69ad9c2d2..5c288396c 100644 --- a/backend/__tests__/unit/services/slackBridgeService.test.js +++ b/backend/__tests__/unit/services/slackBridgeService.test.js @@ -236,6 +236,42 @@ describe('Slack installable bridge', () => { expect(api.postMessage).toHaveBeenCalledWith('D1', expect.stringContaining('Launch')); }); + test('refuses a thread reply into a pod the linked user has left, posting nothing', async () => { + // The gate is still on for pod-2; the membership half is what has to refuse. + // Without this arm the Slack path's membership re-check is unwitnessed — + // dropping the shared predicate left it green (ledger M8, first run). + Pod.findById.mockImplementation((id) => ({ + select: jest.fn().mockReturnValue({ + lean: jest.fn().mockResolvedValue( + String(id) === 'pod-2' + ? { name: 'Launch', type: 'team', members: ['someone-else'] } + : { name: 'Alpha', type: 'team', members: ['user-1'] }, + ), + }), + })); + + const result = await relaySlackMessageToPod({ + integration: { + ...integration, + scope: 'user', + config: { + ...integration.config, + linkedUserId: 'user-1', + slackUserId: 'U1', + gates: { 'pod-2': { enabled: true } }, + relayMap: [{ externalMessageId: '171234.0001', agentUsername: 'kai', podId: 'pod-2' }], + }, + }, + event: { text: 'yes, ship it', user: 'U1', thread_ts: '171234.0001' }, + }); + + expect(result).toEqual({ relayed: false }); + expect(PGMessage.create).not.toHaveBeenCalled(); + expect(deliverMessageToAgents).not.toHaveBeenCalled(); + const api = SlackApi.mock.results[0].value; + expect(api.postMessage).toHaveBeenCalledWith('D1', expect.stringContaining('Launch')); + }); + test('routes a pre-D11 map entry (no podId) to the active pod', async () => { Pod.findById.mockReturnValue(podWithMembers('Alpha')); PGMessage.create.mockResolvedValue({ id: 'pg-1' }); diff --git a/backend/services/connectorRelayPolicy.ts b/backend/services/connectorRelayPolicy.ts index b6b0d72dc..96521eb5a 100644 --- a/backend/services/connectorRelayPolicy.ts +++ b/backend/services/connectorRelayPolicy.ts @@ -88,8 +88,10 @@ export const isRoutedPodTarget = (opts: { const { integration, pod, podId, userId, } = opts; - return Boolean(userId) - && isGatedPodTarget(integration, podId) + // No separate user-id guard: `isPodMember` fails closed on a falsy id itself + // (measured — a guard here changed no arm, so it was removed rather than kept + // unwitnessed). + return isGatedPodTarget(integration, podId) && isPodMember(pod, userId); };