diff --git a/backend/__tests__/unit/routes/integrations.linkedUserId.test.js b/backend/__tests__/unit/routes/integrations.linkedUserId.test.js index 2f1ac1d97..16215dac9 100644 --- a/backend/__tests__/unit/routes/integrations.linkedUserId.test.js +++ b/backend/__tests__/unit/routes/integrations.linkedUserId.test.js @@ -82,6 +82,41 @@ describe('PATCH /api/integrations/:id — linkedUserId guard', () => { expect(update.config.linkedUserId).toBe('user-1'); }); + it("derives linkedUserId when liveRelay arrives as the string 'true' (#1293)", async () => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { liveRelay: 'true' } }); + + expect(res.status).toBe(200); + const [, update] = Integration.findByIdAndUpdate.mock.calls[0]; + expect(update.config.liveRelay).toBe(true); + expect(update.config.linkedUserId).toBe('user-1'); + }); + + // Mongoose's Boolean cast is wider than the two literals above: 1, '1' and + // 'yes' are all stored as true. Any of them would be written as a live relay + // while skipping the === true guards (sprint-review's must-fix on #1297), so + // the edge refuses everything that is not a boolean or 'true'/'false'. + it.each([1, '1', 'yes', 'TRUE', 0, 'no'])( + 'refuses liveRelay %p with 400 instead of letting Mongoose cast it', + async (value) => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { liveRelay: value } }); + expect(res.status).toBe(400); + expect(res.body.message).toMatch(/liveRelay must be true or false/); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }, + ); + + it.each([1, '1', 'yes'])('refuses relayAllAgentMessages %p the same way', async (value) => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { relayAllAgentMessages: value } }); + expect(res.status).toBe(400); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }); + it('does not stamp linkedUserId when liveRelay is switched off', async () => { const res = await request(app) .patch('/api/integrations/integration-1') @@ -94,18 +129,118 @@ describe('PATCH /api/integrations/:id — linkedUserId guard', () => { }); }); -describe('POST /api/integrations — telegram first-run defaults', () => { +// POST /api/integrations — the create path had none of the PATCH guards: +// config was spread verbatim (linkedUserId, connectCode, chatId all client- +// settable) and there was no pod-membership check (ADR-025 review, 2026-08-26). +describe('POST /api/integrations — create-path guards', () => { + beforeEach(() => { + jest.clearAllMocks(); + Integration.prototype.save = jest.fn().mockResolvedValue(undefined); + Pod.findById.mockResolvedValue({ _id: 'pod-1', type: 'private', members: ['user-1'] }); + }); + + it('rejects a client-supplied linkedUserId naming someone else', async () => { + const res = await request(app) + .post('/api/integrations') + .send({ podId: 'pod-1', type: 'telegram', config: { liveRelay: true, linkedUserId: 'VICTIM-USER-ID' } }); + expect(res.status).toBe(400); + }); + + it.each([1, '1', 'yes'])('refuses liveRelay %p on create before anything is saved', async (value) => { + const res = await request(app) + .post('/api/integrations') + .send({ podId: 'pod-1', type: 'telegram', config: { liveRelay: value } }); + expect(res.status).toBe(400); + expect(res.body.message).toMatch(/liveRelay must be true or false/); + expect(Integration.prototype.save).not.toHaveBeenCalled(); + }); + + it('refuses non-members of the target pod', async () => { + Pod.findById.mockResolvedValue({ _id: 'pod-1', type: 'private', members: ['someone-else'] }); + User.findById.mockReturnValue({ select: () => ({ lean: () => Promise.resolve({ role: 'member' }) }) }); + const res = await request(app) + .post('/api/integrations') + .send({ podId: 'pod-1', type: 'telegram', config: {} }); + expect(res.status).toBe(403); + }); + + it('mints a server-side 128-bit expiring code and strips client binding fields', async () => { + const res = await request(app) + .post('/api/integrations') + .send({ + podId: 'pod-1', + type: 'telegram', + config: { + connectCode: 'chosen', chatId: '999', chatType: 'private', liveRelay: true, + }, + }); + expect(res.status).toBe(201); + const { config } = res.body.integration; + expect(config.connectCode).toMatch(/^[0-9a-f]{32}$/); + expect(new Date(config.connectCodeExpiresAt).getTime()).toBeGreaterThan(Date.now()); + expect(config.chatId).toBeUndefined(); + expect(config.chatType).toBeUndefined(); + expect(config.linkedUserId).toBe('user-1'); + }); +}); + +describe('PATCH /api/integrations/:id — live relay on a group chat', () => { beforeEach(() => { jest.clearAllMocks(); - const User = require('../../../models/User'); User.findById.mockResolvedValue({ _id: 'user-1', role: 'member' }); + Pod.findById.mockResolvedValue(null); + Integration.findById.mockResolvedValue({ + ...telegramIntegration(), + config: { chatId: '42', chatType: 'group', toObject() { return { chatId: '42', chatType: 'group' }; } }, + }); + Integration.findByIdAndUpdate.mockResolvedValue({ _id: 'integration-1' }); + }); + + it('refuses to flip liveRelay on when the bound chat is not private', async () => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { liveRelay: true } }); + expect(res.status).toBe(400); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }); + + it("refuses the same flip when liveRelay arrives as the string 'true' (#1293)", async () => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { liveRelay: 'true' } }); + expect(res.status).toBe(400); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }); + + it.each([1, '1', 'yes'])('refuses the flip when liveRelay is %p on a group', async (value) => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { liveRelay: value } }); + expect(res.status).toBe(400); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }); + + it('ignores client-supplied chatId on PATCH', async () => { + const res = await request(app) + .patch('/api/integrations/integration-1') + .send({ config: { chatId: '777', leadAgentUsername: 'theo' } }); + expect(res.status).toBe(200); + const [, update] = Integration.findByIdAndUpdate.mock.calls[0]; + expect(update.config.chatId).toBe('42'); + expect(update.config.leadAgentUsername).toBe('theo'); + }); +}); + +describe('POST /api/integrations — telegram first-run defaults', () => { + beforeEach(() => { + jest.clearAllMocks(); + Pod.findById.mockResolvedValue({ _id: 'pod-1', type: 'private', members: ['user-1'] }); }); // A fresh connector must never default into the silence trap: attention // mode with no lead configured relays nothing, so mirror is the first-run // default and the Connected message teaches /mode attention. it('defaults a new telegram connector to liveRelay + mirror', async () => { - const Integration = require('../../../models/Integration'); const res = await request(app) .post('/api/integrations') .send({ podId: 'pod-1', type: 'telegram', config: {} }); @@ -114,6 +249,20 @@ describe('POST /api/integrations — telegram first-run defaults', () => { expect(created.config.relayAllAgentMessages).toBe(true); expect(created.config.liveRelay).toBe(true); expect(created.config.connectCode).toBeTruthy(); + // The default switched liveRelay on, so the stamp must follow it: a live + // relay with no linkedUserId authors nothing inbound and still streams + // outbound (ordering ruled at the #1297 gate). + expect(created.config.linkedUserId).toBe('user-1'); + }); + + it('respects an explicit liveRelay:false at create and does not stamp', async () => { + const res = await request(app) + .post('/api/integrations') + .send({ podId: 'pod-1', type: 'telegram', config: { liveRelay: false } }); + expect(res.status).toBe(201); + expect(res.body.integration.config.liveRelay).toBe(false); + expect(res.body.integration.config.relayAllAgentMessages).toBe(true); + expect(res.body.integration.config.linkedUserId).toBeUndefined(); }); it('respects an explicit attention-mode choice at create', async () => { diff --git a/backend/__tests__/unit/routes/telegram.webhook.connectCode.test.js b/backend/__tests__/unit/routes/telegram.webhook.connectCode.test.js new file mode 100644 index 000000000..cd890e66c --- /dev/null +++ b/backend/__tests__/unit/routes/telegram.webhook.connectCode.test.js @@ -0,0 +1,107 @@ +// /commonly-enable hardening: expired/legacy codes are refused, attempts are +// rate-limited per chat, and a live-relay integration cannot be bound from a +// group (the relay authors inbound as the linked user and streams outbound). +const request = require('supertest'); +const express = require('express'); + +jest.mock('../../../models/Integration'); +jest.mock('../../../models/Pod'); +jest.mock('../../../models/Summary', () => ({ findOne: jest.fn() })); +jest.mock('../../../services/integrationSummaryService', () => ({ createSummary: jest.fn() })); +jest.mock('../../../services/agentEventService', () => ({ enqueue: jest.fn() })); +jest.mock('../../../services/telegramService', () => ({ sendMessage: jest.fn() })); +jest.mock('../../../integrations', () => ({ get: jest.fn() })); +jest.mock('../../../services/telegramBridgeService', () => ({ relayTelegramMessageToPod: jest.fn() })); +jest.mock('../../../models/WebhookDelivery', () => ({ + create: jest.fn(), + deleteOne: jest.fn(), +})); + +const Integration = require('../../../models/Integration'); +const Pod = require('../../../models/Pod'); +const WebhookDelivery = require('../../../models/WebhookDelivery'); +const telegramService = require('../../../services/telegramService'); +const { resetEnableAttempts, ENABLE_ATTEMPT_LIMIT } = require('../../../services/telegramConnectCode'); +const telegramRoutes = require('../../../routes/webhooks/telegram'); + +const app = express(); +app.use(express.json()); +app.use('/api/webhooks/telegram', telegramRoutes); + +const enable = (code, chat = { id: 42, type: 'private', first_name: 'Sam' }) => request(app) + .post('/api/webhooks/telegram') + .send({ message: { text: `/commonly-enable ${code}`, chat, from: { id: 7 } } }); + +const freshCode = () => ({ + connectCode: 'c'.repeat(32), + connectCodeExpiresAt: new Date(Date.now() + 60000), +}); + +describe('/commonly-enable hardening', () => { + beforeEach(() => { + jest.clearAllMocks(); + resetEnableAttempts(); + process.env.TELEGRAM_BOT_TOKEN = 'bot-token'; + delete process.env.TELEGRAM_SECRET_TOKEN; + // Verification is fail-closed on main (hardening H1); these tests exercise + // the enable handler, not auth, so they run with the explicit dev override + // and a stubbed delivery-claim store, same as telegram.webhook.test.js. + process.env.TELEGRAM_WEBHOOK_ALLOW_UNVERIFIED = 'true'; + WebhookDelivery.create.mockResolvedValue({}); + WebhookDelivery.deleteOne.mockResolvedValue({}); + Integration.findByIdAndUpdate = jest.fn().mockResolvedValue({}); + Pod.findById.mockReturnValue({ lean: jest.fn().mockResolvedValue({ name: 'Test Pod' }) }); + }); + + it('refuses a legacy code with no expiry', async () => { + Integration.findOne = jest.fn() + .mockResolvedValueOnce({ _id: 'i1', podId: 'p1', config: { connectCode: 'abc123' } }); + await enable('abc123'); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + expect(telegramService.sendMessage.mock.calls[0][2]).toMatch(/expired/i); + }); + + it('refuses an expired code', async () => { + Integration.findOne = jest.fn().mockResolvedValueOnce({ + _id: 'i1', podId: 'p1', config: { connectCode: 'x', connectCodeExpiresAt: new Date(Date.now() - 1) }, + }); + await enable('x'); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + }); + + it('binds a fresh code and clears both code fields', async () => { + Integration.findOne = jest.fn() + .mockResolvedValueOnce({ _id: 'i1', podId: 'p1', config: freshCode() }) + .mockResolvedValueOnce(null); + await enable('c'.repeat(32)); + const [, update] = Integration.findByIdAndUpdate.mock.calls[0]; + expect(update.$unset).toEqual({ 'config.connectCode': '', 'config.connectCodeExpiresAt': '' }); + expect(update.$set['config.chatType']).toBe('private'); + }); + + it('rate-limits attempts per chat and stops looking codes up', async () => { + Integration.findOne = jest.fn().mockResolvedValue(null); + for (let i = 0; i < ENABLE_ATTEMPT_LIMIT; i += 1) await enable(`guess${i}`); // eslint-disable-line no-await-in-loop + expect(Integration.findOne).toHaveBeenCalledTimes(ENABLE_ATTEMPT_LIMIT); + await enable('one-more'); + expect(Integration.findOne).toHaveBeenCalledTimes(ENABLE_ATTEMPT_LIMIT); + expect(telegramService.sendMessage.mock.calls.at(-1)[2]).toMatch(/too many attempts/i); + }); + + it('refuses to bind a live-relay integration from a group chat', async () => { + Integration.findOne = jest.fn() + .mockResolvedValueOnce({ _id: 'i1', podId: 'p1', config: { ...freshCode(), liveRelay: true } }) + .mockResolvedValueOnce(null); + await enable('c'.repeat(32), { id: -100, type: 'supergroup', title: 'Crew' }); + expect(Integration.findByIdAndUpdate).not.toHaveBeenCalled(); + expect(telegramService.sendMessage.mock.calls[0][2]).toMatch(/private chat/i); + }); + + it('still binds a legacy (buffer) integration from a group', async () => { + Integration.findOne = jest.fn() + .mockResolvedValueOnce({ _id: 'i1', podId: 'p1', config: freshCode() }) + .mockResolvedValueOnce(null); + await enable('c'.repeat(32), { id: -100, type: 'group', title: 'Crew' }); + expect(Integration.findByIdAndUpdate).toHaveBeenCalled(); + }); +}); diff --git a/backend/__tests__/unit/routes/telegram.webhook.test.js b/backend/__tests__/unit/routes/telegram.webhook.test.js index 1d0bb0aa7..fac989ccc 100644 --- a/backend/__tests__/unit/routes/telegram.webhook.test.js +++ b/backend/__tests__/unit/routes/telegram.webhook.test.js @@ -48,7 +48,7 @@ describe('Telegram webhook routes', () => { const integration = { _id: 'integration-1', podId: 'pod-1', - config: { connectCode: 'abc123' }, + config: { connectCode: 'abc123', connectCodeExpiresAt: new Date(Date.now() + 60000) }, }; Integration.findOne = jest.fn() @@ -170,7 +170,9 @@ describe('Telegram webhook routes', () => { _id: 'integration-1', type: 'telegram', podId: 'pod-1', - config: { chatId: '42', chatType: 'private', liveRelay: true, linkedUserId: 'user-1' }, + config: { + chatId: '42', chatType: 'private', liveRelay: true, linkedUserId: 'user-1', + }, }; const post = () => request(app) diff --git a/backend/__tests__/unit/services/telegramBridgeService.attribution.test.js b/backend/__tests__/unit/services/telegramBridgeService.attribution.test.js index 01ed69867..8083c192d 100644 --- a/backend/__tests__/unit/services/telegramBridgeService.attribution.test.js +++ b/backend/__tests__/unit/services/telegramBridgeService.attribution.test.js @@ -100,3 +100,19 @@ describe('relayTelegramMessageToPod — who the pod row is attributed to', () => expect(User.findById).not.toHaveBeenCalled(); }); }); + +// Outbound mirror of the inbound gate (connector-verify F2, 2026-08-26): a +// code redeemed into a group must not receive the pod's escalation stream. +describe('outbound relay — chatType gate', () => { + it('only resolves live integrations bound to a private chat', async () => { + // eslint-disable-next-line global-require + const Integration = require('../../../models/Integration'); + // eslint-disable-next-line global-require + const { relayAgentMessageToTelegram } = require('../../../services/telegramBridgeService'); + Integration.findOne.mockReturnValue({ lean: () => Promise.resolve(null) }); + await relayAgentMessageToTelegram({ + podId: 'p1', agentUsername: 'theo', displayName: 'Theo', content: '[BLOCKED] x', + }); + expect(Integration.findOne).toHaveBeenCalledWith(expect.objectContaining({ 'config.chatType': 'private' })); + }); +}); diff --git a/backend/__tests__/unit/services/telegramConnectCode.test.js b/backend/__tests__/unit/services/telegramConnectCode.test.js new file mode 100644 index 000000000..cf7b2d9ab --- /dev/null +++ b/backend/__tests__/unit/services/telegramConnectCode.test.js @@ -0,0 +1,53 @@ +// Connect-code lifecycle: 128-bit, 10-minute TTL, legacy codes (no expiry) +// are dead, and /commonly-enable attempts are rate-limited per chat. +const { + mintConnectCode, isConnectCodeExpired, registerEnableAttempt, resetEnableAttempts, + CONNECT_CODE_TTL_MS, ENABLE_ATTEMPT_LIMIT, ENABLE_ATTEMPT_WINDOW_MS, ENABLE_ATTEMPT_MAX_CHATS, +} = require('../../../services/telegramConnectCode'); + +describe('telegramConnectCode', () => { + beforeEach(() => resetEnableAttempts()); + + it('mints a 128-bit hex code with a 10-minute expiry', () => { + const now = 1000000; + const { connectCode, connectCodeExpiresAt } = mintConnectCode(now); + expect(connectCode).toMatch(/^[0-9a-f]{32}$/); + expect(connectCodeExpiresAt.getTime()).toBe(now + CONNECT_CODE_TTL_MS); + expect(mintConnectCode().connectCode).not.toBe(connectCode); + }); + + it('treats a code with no expiry (legacy 24-bit) as expired', () => { + expect(isConnectCodeExpired({ connectCode: 'abc123' })).toBe(true); + expect(isConnectCodeExpired(undefined)).toBe(true); + }); + + it('expires exactly at the deadline', () => { + const cfg = { connectCodeExpiresAt: new Date(2000) }; + expect(isConnectCodeExpired(cfg, 1999)).toBe(false); + expect(isConnectCodeExpired(cfg, 2000)).toBe(true); + }); + + it('allows N attempts per chat per window, then refuses until the window slides', () => { + for (let i = 0; i < ENABLE_ATTEMPT_LIMIT; i += 1) expect(registerEnableAttempt('42', 0)).toBe(true); + expect(registerEnableAttempt('42', 1)).toBe(false); + expect(registerEnableAttempt('43', 1)).toBe(true); // other chats unaffected + expect(registerEnableAttempt('42', ENABLE_ATTEMPT_WINDOW_MS + 1)).toBe(true); + }); + + // The key is any chat id an attacker chooses, so the map cannot grow without + // bound: once the cap is reached, chats whose window has slid out are evicted + // before a new key is admitted — and a chat still inside its window keeps + // its count, so the sweep never resets a live limiter. + it('evicts idle chats at the cap and keeps a live window intact', () => { + for (let i = 0; i < ENABLE_ATTEMPT_LIMIT; i += 1) registerEnableAttempt('hot', 0); + for (let i = 0; i < ENABLE_ATTEMPT_MAX_CHATS - 1; i += 1) registerEnableAttempt(`idle-${i}`, 0); + // Cap reached; a new key one window later triggers the sweep. + const later = ENABLE_ATTEMPT_WINDOW_MS - 1; + expect(registerEnableAttempt('hot', later)).toBe(false); // still limited within its window + expect(registerEnableAttempt('new', ENABLE_ATTEMPT_WINDOW_MS + 1)).toBe(true); + // Everything from t=0 has slid out and was evicted; 'new' is the only key. + expect(registerEnableAttempt('idle-0', ENABLE_ATTEMPT_WINDOW_MS + 1)).toBe(true); + for (let i = 0; i < ENABLE_ATTEMPT_LIMIT - 1; i += 1) registerEnableAttempt('idle-0', ENABLE_ATTEMPT_WINDOW_MS + 1); + expect(registerEnableAttempt('idle-0', ENABLE_ATTEMPT_WINDOW_MS + 1)).toBe(false); + }); +}); diff --git a/backend/models/Integration.ts b/backend/models/Integration.ts index ef8acf7df..a958fa330 100644 --- a/backend/models/Integration.ts +++ b/backend/models/Integration.ts @@ -72,6 +72,7 @@ export interface IIntegration extends Document { lastExternalId?: string; lastExternalTimestamp?: Date; connectCode?: string; + connectCodeExpiresAt?: Date | null; permissions?: string[]; webhookListenerEnabled?: boolean; lastSummaryAt?: Date; @@ -156,6 +157,7 @@ const IntegrationSchema = new Schema( lastExternalId: String, lastExternalTimestamp: Date, connectCode: String, + connectCodeExpiresAt: Date, permissions: [String], webhookListenerEnabled: { type: Boolean, default: false }, lastSummaryAt: Date, diff --git a/backend/routes/integrations.ts b/backend/routes/integrations.ts index 39fe0e6b0..eb45b24c4 100644 --- a/backend/routes/integrations.ts +++ b/backend/routes/integrations.ts @@ -6,7 +6,6 @@ import rateLimit, { ipKeyGenerator } from 'express-rate-limit'; import { createHash } from 'crypto'; // eslint-disable-next-line global-require const axios = require('axios'); -import crypto from 'crypto'; // eslint-disable-next-line global-require const auth = require('../middleware/auth'); // eslint-disable-next-line global-require @@ -31,6 +30,44 @@ const registry = require('../integrations'); const { normalizeBufferMessage } = require('../integrations/normalizeBufferMessage'); // eslint-disable-next-line global-require const { hash, randomSecret } = require('../utils/secret'); +// eslint-disable-next-line global-require +const { mintConnectCode } = require('../services/telegramConnectCode'); +// eslint-disable-next-line global-require +const isPodMember = require('../utils/isPodMember'); + +// Bridge attribution + binding fields are server-owned. linkedUserId is the +// identity every inbound live-relay message is AUTHORED as; chatId/chatType +// are written only by the /commonly-enable webhook (the code is the proof); +// connectCode is minted here. Accepting any of them from a client body lets a +// caller name someone else as the author or bind a chat without a code. +const SERVER_OWNED_CONFIG_KEYS = ['linkedUserId', 'connectCode', 'connectCodeExpiresAt', 'chatId', 'chatType', 'chatTitle']; +const stripServerOwnedConfig = (config: Record): Record => { + const next = { ...config }; + SERVER_OWNED_CONFIG_KEYS.forEach((k) => { delete next[k]; }); + return next; +}; + +// liveRelay / relayAllAgentMessages are declared Boolean paths, so Mongoose +// casts on the way in while the guards below compare strictly (=== true). +// Mongoose's truth table is wider than 'true'/'false': 1, '1' and 'yes' are +// all stored as true. Anything that would be stored as a boolean but is not +// one the guards recognise must be refused at the edge, or the value skips +// the linkedUserId stamp and the group refusal (#1293 through a different +// literal). Accepted: true/false and the legacy string forms 'true'/'false'; +// everything else is a 400, never a silent cast. +const RELAY_FLAG_KEYS = ['liveRelay', 'relayAllAgentMessages']; +const readRelayFlags = (config: Record): { next: Record; invalid?: string } => { + const next = { ...config }; + for (const k of RELAY_FLAG_KEYS) { + const v = next[k]; + if (v === undefined) continue; + if (v === true || v === 'true') next[k] = true; + else if (v === false || v === 'false') next[k] = false; + else return { next, invalid: k }; + } + return { next }; +}; +const relayFlagError = (key: string) => ({ message: `${key} must be true or false` }); interface AuthReq { user?: { id: string; role?: string }; @@ -232,14 +269,50 @@ router.get('/:podId', auth, async (req: AuthReq, res: Res) => { } }); -router.post('/', auth, async (req: AuthReq, res: Res) => { +// Token/IP keying shared by every limiter in this file — same shape as +// routes/messages.ts so NAT'd users don't share a bucket. +const integrationsRateLimitKey = (req: { get?: (h: string) => string | undefined; ip?: string }): string => { + const authHeader = req.get?.('authorization'); + if (authHeader) { + return `tok:${createHash('sha256').update(authHeader).digest('hex').slice(0, 16)}`; + } + return req.ip ? ipKeyGenerator(req.ip) : 'anon'; +}; + +// Write limiter for the create + re-mint paths: each one mints a connect code +// (a bearer secret) and writes a row, so a burst is either a bug or a probe. +const writeIntegrationsRateLimit = rateLimit({ + windowMs: 60_000, + max: 30, + standardHeaders: true, + legacyHeaders: false, + keyGenerator: integrationsRateLimitKey, + handler: (_req: unknown, res: { status: (n: number) => { json: (b: unknown) => void } }) => { + res.status(429).json({ msg: 'rate limit exceeded: 30 writes per 60s' }); + }, +}); + +router.post('/', writeIntegrationsRateLimit, auth, async (req: AuthReq, res: Res) => { try { const { podId, type, config } = (req.body || {}) as { podId?: string; type?: string; config?: Record }; if (!podId || !type || !config) return res.status(400).json({ message: 'Missing required fields' }); const manifest = (manifests as Record)[type]; if (!manifest) return res.status(400).json({ message: 'Unsupported integration type' }); - const nextConfig = { ...config }; - if (type === 'telegram' && !nextConfig.connectCode) nextConfig.connectCode = crypto.randomBytes(3).toString('hex'); + if ('linkedUserId' in config && String(config.linkedUserId) !== String(req.user?.id)) { + return res.status(400).json({ message: 'linkedUserId is derived from the authenticated caller and cannot be set' }); + } + // Membership gate: an integration relays a pod's content outward and + // authors content into it — a WRITE, so it takes the strict predicate + // (members + creator; no admin read-bypass — #1302's isPodMember, not + // DMService.canViewPod). Plain findById: unit mocks resolve a bare doc. + const targetPod = await Pod.findById(String(podId)); + if (!targetPod || !isPodMember(targetPod, req.user?.id)) { + return res.status(403).json({ message: 'Access denied' }); + } + const relay = readRelayFlags(stripServerOwnedConfig(config)); + if (relay.invalid) return res.status(400).json(relayFlagError(relay.invalid)); + const nextConfig: Record = relay.next; + if (type === 'telegram') Object.assign(nextConfig, mintConnectCode()); // First-run default: mirror. A fresh connector has no leadAgentUsername // and its agents use no escalation markers, so attention mode relays // NOTHING — a new user's first experience of the bridge would be silence. @@ -248,6 +321,12 @@ router.post('/', auth, async (req: AuthReq, res: Res) => { nextConfig.relayAllAgentMessages = true; if (nextConfig.liveRelay === undefined) nextConfig.liveRelay = true; } + // Ordering is load-bearing: the default above may have just switched + // liveRelay on, and a live relay with no linkedUserId authors nothing + // inbound (the bridge fails closed on a missing linked user) while still + // streaming outbound. The creator IS the authenticated caller, so the + // defaulted case is stamped exactly like an explicit liveRelay:true. + if (nextConfig.liveRelay === true) nextConfig.linkedUserId = req.user?.id; const missingRequired = getMissingRequiredFields(type, nextConfig); if (type === 'discord' && missingRequired.length) return res.status(400).json({ message: `Missing required fields: ${missingRequired.join(', ')}`, missing: missingRequired }); if (missingRequired.length && (req.body as { status?: string })?.status === 'connected') return res.status(400).json({ message: `Missing required fields: ${missingRequired.join(', ')}`, missing: missingRequired }); @@ -394,13 +473,7 @@ const listIntegrationsRateLimit = rateLimit({ max: 120, standardHeaders: true, legacyHeaders: false, - keyGenerator: (req: { get?: (h: string) => string | undefined; ip?: string }) => { - const authHeader = req.get?.('authorization'); - if (authHeader) { - return `tok:${createHash('sha256').update(authHeader).digest('hex').slice(0, 16)}`; - } - return req.ip ? ipKeyGenerator(req.ip) : 'anon'; - }, + keyGenerator: integrationsRateLimitKey, handler: (_req: unknown, res: { status: (n: number) => { json: (b: unknown) => void } }) => { res.status(429).json({ msg: 'rate limit exceeded: 120 reads per 60s' }); }, @@ -416,6 +489,27 @@ router.get('/user/all', listIntegrationsRateLimit, auth, async (req: AuthReq, re } }); +// Re-mint the one-time enable code (codes expire after 10 minutes). Only +// while the chat is still unbound — a connected integration has no code. +router.post('/:id/connect-code', writeIntegrationsRateLimit, auth, async (req: AuthReq, res: Res) => { + try { + const { id } = req.params || {}; + const integration = await Integration.findById(id) as { type?: string; createdBy?: { toString: () => string }; podId?: unknown; config?: { chatId?: string } } | null; + if (!integration) return res.status(404).json({ message: 'Integration not found' }); + if (integration.type !== 'telegram') return res.status(400).json({ message: 'Connect codes are telegram-only' }); + if (!(await canDeleteIntegration(integration, req.user?.id || ''))) return res.status(403).json({ message: 'Access denied' }); + if (integration.config?.chatId) return res.status(409).json({ message: 'Already linked to a chat' }); + const minted = mintConnectCode(); + const updated = await Integration.findByIdAndUpdate(id, { + $set: { 'config.connectCode': minted.connectCode, 'config.connectCodeExpiresAt': minted.connectCodeExpiresAt }, + }, { new: true }); + return res.json(updated); + } catch (error) { + console.error('Error minting connect code:', error); + return res.status(500).json({ message: 'Server error' }); + } +}); + router.patch('/:id', auth, async (req: AuthReq, res: Res) => { try { const { id } = req.params || {}; @@ -433,8 +527,18 @@ router.patch('/:id', auth, async (req: AuthReq, res: Res) => { if (config && 'linkedUserId' in config && String(config.linkedUserId) !== String(req.user?.id)) { return res.status(400).json({ message: 'linkedUserId is derived from the authenticated caller and cannot be set' }); } - const nextConfig = config ? { ...currentConfig, ...config } : currentConfig; - if (config && config.liveRelay === true) nextConfig.linkedUserId = req.user?.id; + const relay = config ? readRelayFlags(stripServerOwnedConfig(config)) : null; + if (relay?.invalid) return res.status(400).json(relayFlagError(relay.invalid)); + const incoming = relay ? relay.next : null; + const nextConfig = incoming ? { ...currentConfig, ...incoming } : currentConfig; + if (incoming && incoming.liveRelay === true) { + // Relay authors inbound as linkedUserId and streams outbound to chatId; + // both are only honest in a 1:1 private chat (#1289 inbound, F2 outbound). + if (nextConfig.chatId && nextConfig.chatType !== 'private') { + return res.status(400).json({ message: 'Live relay requires a private chat with the bot; this connector is bound to a group' }); + } + nextConfig.linkedUserId = req.user?.id; + } const missingRequired = getMissingRequiredFields(integration.type || '', nextConfig); if (missingRequired.length && status === 'connected') return res.status(400).json({ message: `Missing required fields: ${missingRequired.join(', ')}`, missing: missingRequired }); validateManifestIfComplete(integration.type || '', nextConfig); diff --git a/backend/routes/webhooks/telegram.ts b/backend/routes/webhooks/telegram.ts index 509c12c34..997afe473 100644 --- a/backend/routes/webhooks/telegram.ts +++ b/backend/routes/webhooks/telegram.ts @@ -7,6 +7,7 @@ const registry = require('../../integrations'); const IntegrationSummaryService = require('../../services/integrationSummaryService'); const AgentEventService = require('../../services/agentEventService'); const telegramService = require('../../services/telegramService'); +const { isConnectCodeExpired, registerEnableAttempt } = require('../../services/telegramConnectCode'); const router = express.Router({ mergeParams: true }); @@ -76,17 +77,28 @@ const handleEnableCommand = async (chat: any, code: any) => { return; } + // The code is the only proof of ownership this path has (the redeemer is + // unauthenticated), so guessing is rate-limited per chat and codes expire. + if (!registerEnableAttempt(chatId)) { + await telegramService.sendMessage( + botToken, + chatId, + '⚠️ Too many attempts. Wait a few minutes and request a fresh code from Commonly.', + ); + return; + } + const integration = await Integration.findOne({ type: 'telegram', isActive: true, 'config.connectCode': code, }); - if (!integration) { + if (!integration || isConnectCodeExpired(integration.config)) { await telegramService.sendMessage( botToken, chatId, - '❌ Invalid code. Please request a fresh code from Commonly.', + '❌ Invalid or expired code. Please request a fresh code from Commonly.', ); return; } @@ -119,6 +131,17 @@ const handleEnableCommand = async (chat: any, code: any) => { const chatTitle = getChatTitle(chat); const chatType = chat?.type || null; + // Live relay authors inbound as the linked user and streams the pod's + // escalations outbound; both are only honest in a private 1:1 chat. + if (integration.config?.liveRelay && chatType !== 'private') { + await telegramService.sendMessage( + botToken, + chatId, + '⚠️ Live relay only works from a private chat with this bot. Open a direct chat and send the code there.', + ); + return; + } + await Integration.findByIdAndUpdate(integration._id, { status: 'connected', $set: { @@ -129,6 +152,7 @@ const handleEnableCommand = async (chat: any, code: any) => { }, $unset: { 'config.connectCode': '', + 'config.connectCodeExpiresAt': '', }, }); diff --git a/backend/services/telegramBridgeService.ts b/backend/services/telegramBridgeService.ts index 06113cdb7..603e03f7c 100644 --- a/backend/services/telegramBridgeService.ts +++ b/backend/services/telegramBridgeService.ts @@ -108,6 +108,9 @@ const findLiveIntegration = async (podId: unknown): Promise ({ + connectCode: crypto.randomBytes(16).toString('hex'), + connectCodeExpiresAt: new Date(now + CONNECT_CODE_TTL_MS), +}); + +// A code without an expiry predates the TTL — treat it as expired so legacy +// 24-bit codes can never be redeemed; the owner re-mints from the UI. +export const isConnectCodeExpired = ( + config: { connectCode?: string; connectCodeExpiresAt?: Date | string | null } | undefined, + now: number = Date.now(), +): boolean => { + if (!config?.connectCodeExpiresAt) return true; + return new Date(config.connectCodeExpiresAt).getTime() <= now; +}; + +// Per-chat sliding window for /commonly-enable attempts. In-memory is enough: +// a code lives 10 minutes and the backend runs one replica; a restart resets +// the window, which is the failure mode we accept over a Redis dependency. +// COUPLING: the "one replica" premise is `replicaCount: 1` in the Helm values +// AND `autoscaling.backend.enabled: false` (values.yaml). Enabling backend +// autoscaling silently makes the effective limit ENABLE_ATTEMPT_LIMIT × replicas +// — move this window to Redis in the same change. +// The key is attacker-supplied (any chat id), so the map is bounded: past +// ENABLE_ATTEMPT_MAX_CHATS keys, every chat whose window has fully slid out is +// evicted before a new key is added. A chat that attempts once and never +// returns costs one slot for one window, not forever. +export const ENABLE_ATTEMPT_MAX_CHATS = 10_000; +const attempts = new Map(); + +const sweepIdleChats = (now: number): void => { + attempts.forEach((stamps, key) => { + if (!stamps.some((t) => now - t < ENABLE_ATTEMPT_WINDOW_MS)) attempts.delete(key); + }); +}; + +export const registerEnableAttempt = (chatId: string, now: number = Date.now()): boolean => { + if (!attempts.has(chatId) && attempts.size >= ENABLE_ATTEMPT_MAX_CHATS) sweepIdleChats(now); + const recent = (attempts.get(chatId) || []).filter((t) => now - t < ENABLE_ATTEMPT_WINDOW_MS); + if (recent.length >= ENABLE_ATTEMPT_LIMIT) { + attempts.set(chatId, recent); + return false; + } + recent.push(now); + attempts.set(chatId, recent); + return true; +}; + +export const resetEnableAttempts = (): void => { attempts.clear(); }; + +module.exports = { + CONNECT_CODE_TTL_MS, + ENABLE_ATTEMPT_WINDOW_MS, + ENABLE_ATTEMPT_LIMIT, + ENABLE_ATTEMPT_MAX_CHATS, + mintConnectCode, + isConnectCodeExpired, + registerEnableAttempt, + resetEnableAttempts, +}; + +export {}; diff --git a/frontend/src/v2/__tests__/V2ConnectorsPage.test.tsx b/frontend/src/v2/__tests__/V2ConnectorsPage.test.tsx index 7ad91ff79..09060afc4 100644 --- a/frontend/src/v2/__tests__/V2ConnectorsPage.test.tsx +++ b/frontend/src/v2/__tests__/V2ConnectorsPage.test.tsx @@ -44,7 +44,7 @@ const connectors = [ _id: 'i-pending', type: 'telegram', status: 'pending', - config: { connectCode: 'abc123' }, + config: { connectCode: 'abc123', connectCodeExpiresAt: new Date(Date.now() + 60_000).toISOString() }, podId: { _id: 'p1', name: 'Rewire Live Demo' }, }, { @@ -93,6 +93,21 @@ describe('V2ConnectorsPage', () => { expect(screen.getByText(/Rewire crew/)).toBeInTheDocument(); }); + it('offers a new code when the enable code has expired (or never had an expiry)', async () => { + mockGets([{ ...connectors[0], config: { connectCode: 'abc123' } }]); + axios.post.mockResolvedValue({ data: {} }); + renderPage(); + expect(await screen.findByText(/code expired/i)).toBeInTheDocument(); + expect(screen.queryByText('abc1 23')).not.toBeInTheDocument(); + expect(screen.queryByText('Copy command')).not.toBeInTheDocument(); + fireEvent.click(screen.getByRole('button', { name: /new code/i })); + await waitFor(() => expect(axios.post).toHaveBeenCalledWith( + '/api/integrations/i-pending/connect-code', + {}, + expect.anything(), + )); + }); + it('toggling relay PATCHes liveRelay only — linkedUserId is server-derived', async () => { mockGets(); axios.patch.mockResolvedValue({ data: {} }); diff --git a/frontend/src/v2/components/V2ConnectorsPage.tsx b/frontend/src/v2/components/V2ConnectorsPage.tsx index a8e65a91a..6b181622d 100644 --- a/frontend/src/v2/components/V2ConnectorsPage.tsx +++ b/frontend/src/v2/components/V2ConnectorsPage.tsx @@ -17,6 +17,7 @@ import { PlatformGlyph } from '../icons/platforms'; interface ConnectorConfig { chatTitle?: string; connectCode?: string; + connectCodeExpiresAt?: string; liveRelay?: boolean; relayAllAgentMessages?: boolean; } @@ -62,6 +63,15 @@ const groupCode = (code: string): string => (code.match(/.{1,4}/g) || [code]).jo const RECENT_MS = 10 * 60_000; +// Codes expire after 10 minutes server-side (#1297); a pending card past its +// expiry — or carrying a legacy code that never had one — offers a re-mint +// instead of a command that the webhook will refuse. +const codeIsLive = (c: Connector): boolean => Boolean( + c.config?.connectCode + && c.config?.connectCodeExpiresAt + && new Date(c.config.connectCodeExpiresAt).getTime() > Date.now(), +); + const V2ConnectorsPage: React.FC = () => { const { t } = useTranslation(); const api = useV2Api(); @@ -93,7 +103,7 @@ const V2ConnectorsPage: React.FC = () => { useEffect(() => { load(); }, [load]); // Spec §2.3 step 3: poll while any code is pending; stop when it connects. - const hasPending = connectors.some((c) => c.status !== 'connected' && c.config?.connectCode); + const hasPending = connectors.some((c) => c.status !== 'connected' && codeIsLive(c)); useEffect(() => { if (hasPending && !pollRef.current) { pollRef.current = setInterval(load, 3000); @@ -152,6 +162,20 @@ const V2ConnectorsPage: React.FC = () => { } }; + const regenerateCode = async (c: Connector) => { + if (busyId) return; + setBusyId(c._id); + setError(null); + try { + await api.post(`/api/integrations/${c._id}/connect-code`, {}); + await load(); + } catch { + setError(t('connectors.codeError', { defaultValue: 'Could not create a new code.' })); + } finally { + setBusyId(null); + } + }; + const disconnect = async (c: Connector) => { setBusyId(c._id); try { @@ -278,7 +302,18 @@ const V2ConnectorsPage: React.FC = () => { )} - {!connected && isTelegram && c.config?.connectCode && ( + {!connected && isTelegram && !codeIsLive(c) && ( +
+
+ {t('connectors.codeExpired', { defaultValue: 'The enable code expired.' })} +
+ +
+ )} + + {!connected && isTelegram && codeIsLive(c) && (
{t('connectors.enableHint', { defaultValue: 'Open a private chat with the Commonly bot' })} @@ -290,7 +325,7 @@ const V2ConnectorsPage: React.FC = () => { ? t('connectors.copied', { defaultValue: 'Copied' }) : t('connectors.copyCommand', { defaultValue: 'Copy command' })} - {groupCode(c.config.connectCode)} + {groupCode(c.config?.connectCode || '')}
)}