From 4c619312a7f89522b520ee4f66bd047ee5f8e79c Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Sun, 27 Sep 2026 01:02:18 -0700 Subject: [PATCH] fix(connectors): the external feed sync reads membership before it writes (TASK-164) `syncExternalFeeds` read no membership at all. The flag-on path wrote posts AS `integration.createdBy`, and the default path handed the owner's feed to the pod's curator agents; `createMessage` would 403 that same owner. Found by Wren (74669) while closing TASK-161, and kept out of #1940 because it needs a read plus a pause rule rather than a predicate swap. The read goes at the top of the per-integration body, before `syncRecent`, so a departed owner makes no provider call and advances no cursor either. The predicate is `isListedPodMember` - the same rule the pod's own write paths implement - so the feed is never more permissive than the pod it writes into. A pod that is gone pauses too: there is no surface left to write into, and syncing into nothing looks like a healthy run. Pause shape (wren, TASK-164): `status: 'error'` with a Commonly-written `errorMessage` and `errorMessageUserFacing: true`, written as its own update rather than thrown into the catch, which stamps the flag false. `isActive` stays true, so the row stays on the owner's Connectors page - that page has no action for x/instagram, which is why the copy carries the next step itself. `'error'` also takes the row out of the sync query (`status: 'connected'`), so the pause is written once per connection rather than every tick; a re-saved row pauses again on the next sync, and no resume logic is added. Also, admin rows would have paused at birth: only the FIRST requester of the Global Social Feed pod became a Mongo member (at creation), and a second admin configuring the other feed type was mirrored into PG alone. `ensureGlobalSocialFeedPod` now `$addToSet`s the requester into Mongo `members` and hands on the refetched pod. Not in scope, noted by wren: those global routes find their pod by name, and pod names are not unique. --- .../routes/admin.globalIntegrations.test.js | 84 +++++++++- .../unit/services/externalFeedService.test.js | 155 +++++++++++++++++- backend/routes/admin/globalIntegrations.ts | 13 ++ backend/services/externalFeedService.ts | 53 ++++++ 4 files changed, 296 insertions(+), 9 deletions(-) diff --git a/backend/__tests__/unit/routes/admin.globalIntegrations.test.js b/backend/__tests__/unit/routes/admin.globalIntegrations.test.js index b6e7676f3..f333dc3ac 100644 --- a/backend/__tests__/unit/routes/admin.globalIntegrations.test.js +++ b/backend/__tests__/unit/routes/admin.globalIntegrations.test.js @@ -13,6 +13,8 @@ jest.mock('../../../models/OAuthState', () => ({ jest.mock('../../../models/Pod', () => ({ findOne: jest.fn(), create: jest.fn(), + updateOne: jest.fn(), + findById: jest.fn(), })); jest.mock('../../../integrations', () => ({ @@ -61,7 +63,17 @@ function createRes() { describe('admin global integrations route', () => { beforeEach(() => { jest.clearAllMocks(); - Pod.findOne.mockResolvedValue({ _id: 'pod-global', name: 'Global Social Feed' }); + Pod.findOne.mockResolvedValue({ + _id: 'pod-global', + name: 'Global Social Feed', + members: ['admin-1'], + }); + Pod.updateOne.mockResolvedValue({}); + Pod.findById.mockResolvedValue({ + _id: 'pod-global', + name: 'Global Social Feed', + members: ['admin-1', 'admin-2'], + }); GlobalModelConfigService.getConfig.mockResolvedValue({ llmService: { provider: 'auto', @@ -263,6 +275,76 @@ describe('admin global integrations route', () => { expect(Integration.create).not.toHaveBeenCalled(); }); + it('lists a second admin in the global pod before their first sync (TASK-164)', async () => { + const handler = getRouteHandler('/x', 'post'); + const req = { + userId: 'admin-2', + body: { + enabled: true, + accessToken: 'x-token', + username: 'commonly', + userId: 'x-user-id', + }, + }; + const res = createRes(); + Integration.findOne.mockResolvedValueOnce(null); + Integration.create.mockResolvedValueOnce({ + _id: 'x-int-2', + type: 'x', + status: 'connected', + config: {}, + }); + + await handler(req, res); + + // Only the first requester became a Mongo member (at pod creation), so a + // second admin's feed would pause on its first sync without this write. + expect(Pod.updateOne).toHaveBeenCalledWith( + { _id: 'pod-global' }, + { $addToSet: { members: 'admin-2' } }, + ); + // The pod handed on is the refetched one, so its members match the write. + expect(Pod.findById).toHaveBeenCalledWith('pod-global'); + expect(Integration.create).toHaveBeenCalledWith(expect.objectContaining({ + createdBy: 'admin-2', + podId: 'pod-global', + })); + expect(res.json).toHaveBeenCalledWith({ + success: true, + integration: expect.objectContaining({ _id: 'x-int-2' }), + }); + }); + + it('does not re-add an admin the global pod already lists (control)', async () => { + const handler = getRouteHandler('/x', 'post'); + const req = { + userId: 'admin-1', + body: { + enabled: true, + accessToken: 'x-token', + username: 'commonly', + userId: 'x-user-id', + }, + }; + const res = createRes(); + Integration.findOne.mockResolvedValueOnce(null); + Integration.create.mockResolvedValueOnce({ + _id: 'x-int-3', + type: 'x', + status: 'connected', + config: {}, + }); + + await handler(req, res); + + expect(Pod.updateOne).not.toHaveBeenCalled(); + expect(Pod.findById).not.toHaveBeenCalled(); + expect(res.json).toHaveBeenCalledWith({ + success: true, + integration: expect.objectContaining({ _id: 'x-int-3' }), + }); + }); + it('the admin Instagram save keeps the stored token when the body omits it', async () => { const handler = getRouteHandler('/instagram', 'post'); const req = { diff --git a/backend/__tests__/unit/services/externalFeedService.test.js b/backend/__tests__/unit/services/externalFeedService.test.js index c8696521c..5a4737c2e 100644 --- a/backend/__tests__/unit/services/externalFeedService.test.js +++ b/backend/__tests__/unit/services/externalFeedService.test.js @@ -3,6 +3,10 @@ jest.mock('../../../models/Integration', () => ({ findByIdAndUpdate: jest.fn(), })); +jest.mock('../../../models/Pod', () => ({ + findById: jest.fn(), +})); + jest.mock('../../../models/Post', () => { const Post = jest.fn(function Post(doc) { Object.assign(this, doc); @@ -26,7 +30,9 @@ jest.mock('../../../services/agentEventService', () => ({ enqueue: jest.fn(), })); +const mongoose = require('mongoose'); const Integration = require('../../../models/Integration'); +const Pod = require('../../../models/Pod'); const Post = require('../../../models/Post'); const { AgentInstallation } = require('../../../models/AgentRegistry'); const registry = require('../../../integrations'); @@ -34,10 +40,12 @@ const AgentEventService = require('../../../services/agentEventService'); const externalFeedService = require('../../../services/externalFeedService'); describe('externalFeedService', () => { - beforeEach(() => { - jest.clearAllMocks(); - delete process.env.EXTERNAL_FEED_PERSIST_POSTS; - }); + // The real shape: the service reads the pod with `.lean()`, so members and a + // lean integration's `createdBy` are ObjectIds, not strings. String members + // would let a raw `members.includes(owner)` pass for the shared predicate + // (TASK-164 ledger M10 - the fixture, not the code, decides that). + const OWNER_ID = new mongoose.Types.ObjectId('507f1f77bcf86cd799439011'); + const OTHER_ID = new mongoose.Types.ObjectId('507f191e810c19729de860ea'); const mockFindChain = (value) => ({ select: () => ({ @@ -46,6 +54,45 @@ describe('externalFeedService', () => { lean: jest.fn().mockResolvedValue(value), }); + const mockPodMembers = (members) => { + // A fresh instance carrying the same hex, which is what a lean read of the + // pod produces: the pod's ObjectId and the integration's are never the same + // reference. Admitting the owner has to be a value comparison, so a raw + // `members.includes(owner)` - reference equality - is refused here + // (TASK-164 ledger M10). + const stored = members.map((member) => ( + member instanceof mongoose.Types.ObjectId + ? new mongoose.Types.ObjectId(member.toHexString()) + : member + )); + Pod.findById.mockReturnValue({ + select: () => ({ lean: jest.fn().mockResolvedValue({ _id: 'pod-1', members: stored }) }), + }); + }; + + const feedRow = (over = {}) => ({ + _id: 'int-1', + type: 'x', + podId: 'pod-1', + status: 'connected', + isActive: true, + createdBy: OWNER_ID, + config: { messageBuffer: [], maxBufferSize: 1000 }, + ...over, + }); + + const mockOneFeed = (over = {}) => { + Integration.find.mockReturnValue({ lean: jest.fn().mockResolvedValue([feedRow(over)]) }); + }; + + beforeEach(() => { + jest.clearAllMocks(); + delete process.env.EXTERNAL_FEED_PERSIST_POSTS; + // The default fixture is a pod that lists the integration's owner, so every + // pre-existing arm below stays a listed-owner arm (TASK-164). + mockPodMembers([OWNER_ID]); + }); + test('does not persist external feed posts by default and enqueues curator events', async () => { Integration.find.mockReturnValue({ lean: jest.fn().mockResolvedValue([ @@ -55,7 +102,7 @@ describe('externalFeedService', () => { podId: 'pod-1', status: 'connected', isActive: true, - createdBy: 'user-1', + createdBy: OWNER_ID, config: { messageBuffer: [], maxBufferSize: 1000 }, }, ]), @@ -107,6 +154,14 @@ describe('externalFeedService', () => { createdPosts: 0, curatorEventsEnqueued: 1, })); + // The agreed control for the membership arms below (wren, TASK-164): a + // listed owner still reaches the provider, still buffers, still enqueues. + expect(registry.get).toHaveBeenCalledTimes(1); + expect(Integration.findByIdAndUpdate).toHaveBeenCalledWith('int-1', { + $push: { + 'config.messageBuffer': expect.objectContaining({ $each: expect.any(Array) }), + }, + }); }); test('can persist external posts when EXTERNAL_FEED_PERSIST_POSTS=1', async () => { @@ -119,7 +174,7 @@ describe('externalFeedService', () => { podId: 'pod-1', status: 'connected', isActive: true, - createdBy: 'user-1', + createdBy: OWNER_ID, config: { messageBuffer: [], maxBufferSize: 1000 }, }, ]), @@ -160,7 +215,7 @@ describe('externalFeedService', () => { podId: 'pod-1', status: 'connected', isActive: true, - createdBy: 'user-1', + createdBy: OWNER_ID, config: { messageBuffer: [], maxBufferSize: 1000 }, }, ]), @@ -211,7 +266,7 @@ describe('externalFeedService', () => { podId: 'pod-1', status: 'connected', isActive: true, - createdBy: 'user-1', + createdBy: OWNER_ID, config: { messageBuffer: [], maxBufferSize: 1000 }, }, ]), @@ -241,4 +296,88 @@ describe('externalFeedService', () => { }), ); }); + + describe('owner membership (TASK-164)', () => { + test('a departed owner calls no provider, buffers nothing, and pauses the row with the reason', async () => { + mockPodMembers([OTHER_ID]); + mockOneFeed(); + + const results = await externalFeedService.syncExternalFeeds(); + + // No provider call, so no cursor advance either: the whole sync is skipped. + expect(registry.get).not.toHaveBeenCalled(); + expect(Post.insertMany).not.toHaveBeenCalled(); + expect(AgentEventService.enqueue).not.toHaveBeenCalled(); + + const writes = Integration.findByIdAndUpdate.mock.calls; + expect(writes).toHaveLength(1); + const [, update] = writes[0]; + expect(writes[0][0]).toBe('int-1'); + expect(update.$push).toBeUndefined(); + expect(update.$set).toEqual(expect.objectContaining({ + status: 'error', + errorMessageUserFacing: true, + })); + expect(update.$set.errorMessage).toContain('no longer a member of the pod'); + // isActive stays true, so the row stays on the owner's Connectors page; + // the pause is expressed as status, and nothing here writes the flag. + expect(update.$set.isActive).toBeUndefined(); + expect(results[0]).toEqual(expect.objectContaining({ + integrationId: 'int-1', + success: false, + paused: true, + messageCount: 0, + })); + }); + + test('a departed owner writes no posts on the flag-on path either', async () => { + process.env.EXTERNAL_FEED_PERSIST_POSTS = '1'; + mockPodMembers([]); + mockOneFeed(); + + const results = await externalFeedService.syncExternalFeeds(); + + expect(Post.find).not.toHaveBeenCalled(); + expect(Post.insertMany).not.toHaveBeenCalled(); + expect(registry.get).not.toHaveBeenCalled(); + expect(results[0]).toEqual(expect.objectContaining({ paused: true, success: false })); + }); + + test('a pod that is gone pauses rather than syncing into nothing', async () => { + Pod.findById.mockReturnValue({ select: () => ({ lean: jest.fn().mockResolvedValue(null) }) }); + mockOneFeed(); + + const results = await externalFeedService.syncExternalFeeds(); + + expect(registry.get).not.toHaveBeenCalled(); + expect(results[0]).toEqual(expect.objectContaining({ paused: true, success: false })); + }); + + test('a listed owner is unaffected: provider, buffer and curator events all run', async () => { + mockPodMembers([OWNER_ID]); + mockOneFeed(); + registry.get.mockReturnValue({ + syncRecent: jest.fn().mockResolvedValue({ + messages: [{ + externalId: 'x-9', + content: 'post nine', + timestamp: new Date().toISOString(), + authorName: 'author', + }], + content: 'Synced external feed', + }), + }); + Post.find.mockImplementation(() => mockFindChain([])); + AgentInstallation.find.mockReturnValue(mockFindChain([])); + + const results = await externalFeedService.syncExternalFeeds(); + + expect(registry.get).toHaveBeenCalledWith('x', expect.objectContaining({ _id: 'int-1' })); + expect(Integration.findByIdAndUpdate).toHaveBeenCalledWith('int-1', expect.objectContaining({ + $push: expect.anything(), + })); + expect(results[0]).toEqual(expect.objectContaining({ success: true, messageCount: 1 })); + expect(results[0].paused).toBeUndefined(); + }); + }); }); diff --git a/backend/routes/admin/globalIntegrations.ts b/backend/routes/admin/globalIntegrations.ts index da86d963d..81e1a959e 100644 --- a/backend/routes/admin/globalIntegrations.ts +++ b/backend/routes/admin/globalIntegrations.ts @@ -12,6 +12,8 @@ const Pod = require('../../models/Pod'); const registry = require('../../integrations'); const SocialPolicyService = require('../../services/socialPolicyService'); const GlobalModelConfigService = require('../../services/globalModelConfigService'); +// eslint-disable-next-line global-require +const { isListedPodMember } = require('../../utils/isPodMember'); const externalFeedService = require('../../services/externalFeedService'); let PGPod = null; @@ -115,6 +117,17 @@ const ensureGlobalSocialFeedPod = async (userId: any) => { createdBy: userId, tags: ['social', 'global', 'feeds'], }); + } else if (!isListedPodMember(globalPod, userId)) { + // The requester is about to own a feed integration in this pod, and the + // sync refuses to write for an owner the pod does not list (TASK-164). + // Only the FIRST requester became a Mongo member (at creation); a second + // admin configuring the other feed type was mirrored into PG alone, so + // their first sync would pause a supported setup. Mongo `members` is what + // the predicate reads and what the pod's own write paths enforce; the PG + // mirror follows below. `createdBy` is deliberately untouched — it is the + // row's owner, not a membership record. + await Pod.updateOne({ _id: globalPod._id }, { $addToSet: { members: userId } }); + globalPod = await Pod.findById(globalPod._id); } await ensureGlobalPodPostgresSync({ pod: globalPod, userId }); diff --git a/backend/services/externalFeedService.ts b/backend/services/externalFeedService.ts index c25c42b5d..e109ed801 100644 --- a/backend/services/externalFeedService.ts +++ b/backend/services/externalFeedService.ts @@ -3,6 +3,10 @@ const Integration = require('../models/Integration'); // eslint-disable-next-line global-require const Post = require('../models/Post'); // eslint-disable-next-line global-require +const Pod = require('../models/Pod'); +// eslint-disable-next-line global-require +const { isListedPodMember } = require('../utils/isPodMember'); +// eslint-disable-next-line global-require const { AgentInstallation } = require('../models/AgentRegistry'); // eslint-disable-next-line global-require const registry = require('../integrations'); @@ -73,6 +77,8 @@ interface CuratorDispatch { interface FeedSyncResult { integrationId: unknown; success: boolean; + /** The owner is not a listed member of the target pod, so nothing was written. */ + paused?: boolean; messageCount: number; createdPosts?: number; curatorEventsEnqueued?: number; @@ -84,6 +90,39 @@ function shouldPersistExternalFeedPosts(): boolean { return process.env.EXTERNAL_FEED_PERSIST_POSTS === '1'; } +const OWNER_LEFT_POD_MESSAGE = 'The account that authorized this connection is no longer a member of the pod it ' + + 'syncs into, so nothing is being synced. Add that account back to the pod, then reconnect this feed.'; + +/** + * A feed's owner can stop being a pod member without anyone touching the + * integration: `leavePod` filters them out of `members`, agent cleanup `$pull`s + * them, and the row kept syncing as them. This service read no membership at + * all before TASK-164, so a write could be attributed to an account the pod + * refused — the flag-on Post path wrote as `integration.createdBy`, and the + * default path handed their feed to the pod's curator agents. + * + * `isListedPodMember` is the predicate `createMessage` and the socket write path + * implement, so the feed is never more permissive than the pod it writes into. + * A missing pod pauses too: there is no surface left to write into, and silence + * would look like a healthy sync. + * + * `status: 'error'` takes the row out of the sync query (`status: 'connected'`), + * so this is written once per connection rather than on every tick, and the row + * keeps `isActive: true` so it stays on the owner's Connectors page. That page + * has no action for x/instagram, which is why the copy carries the next step + * itself, and why it is marked user-facing rather than diagnostic. + */ +async function pauseFeedForDepartedOwner(integrationId: unknown): Promise { + await Integration.findByIdAndUpdate(integrationId, { + $set: { + status: 'error', + errorMessage: OWNER_LEFT_POD_MESSAGE, + errorMessageUserFacing: true, + lastSync: new Date(), + }, + }); +} + function getAttachmentUrl(attachments: NormalizedMessage['attachments'] = []): string { if (!Array.isArray(attachments) || attachments.length === 0) return ''; const first = attachments[0]; @@ -312,6 +351,20 @@ async function syncExternalFeeds(): Promise { const results = await Promise.all( integrations.map(async (integration): Promise => { try { + // Before the provider call: a departed owner must make no provider + // request, advance no cursor, and write nothing on either path. + const pod = await Pod.findById(integration.podId).select('members').lean(); + if (!pod || !isListedPodMember(pod, integration.createdBy)) { + await pauseFeedForDepartedOwner(integration._id); + return { + integrationId: integration._id, + success: false, + paused: true, + messageCount: 0, + content: OWNER_LEFT_POD_MESSAGE, + }; + } + const provider = registry.get(integration.type, integration); const sinceId = integration.config?.lastExternalId; const sinceTimestamp = integration.config?.lastExternalTimestamp;