From 337dffad8ca9897b19a619e3cf5b6f07aa2cf268 Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Thu, 3 Sep 2026 04:30:53 -0700 Subject: [PATCH 1/2] feat(activity): persist recipient attention items --- .../unit/models/AttentionItem.test.js | 38 ++++ .../activityService.decisionQueue.test.js | 4 +- .../services/activityService.recap.test.js | 4 +- .../services/attentionItemService.test.js | 60 ++++++ backend/models/Activity.ts | 8 +- backend/models/AttentionItem.ts | 58 ++++++ backend/models/Message.ts | 15 ++ backend/models/User.ts | 8 - backend/models/pg/Message.ts | 21 +- backend/services/activityService.ts | 59 ++++-- backend/services/attentionItemService.ts | 183 ++++++++++++++++++ backend/services/decisionRequestService.ts | 6 + .../src/v2/__tests__/V2ActivityPage.test.tsx | 15 +- frontend/src/v2/components/V2ActivityPage.tsx | 21 +- 14 files changed, 438 insertions(+), 62 deletions(-) create mode 100644 backend/__tests__/unit/models/AttentionItem.test.js create mode 100644 backend/__tests__/unit/services/attentionItemService.test.js create mode 100644 backend/models/AttentionItem.ts create mode 100644 backend/services/attentionItemService.ts diff --git a/backend/__tests__/unit/models/AttentionItem.test.js b/backend/__tests__/unit/models/AttentionItem.test.js new file mode 100644 index 000000000..ac8b397e9 --- /dev/null +++ b/backend/__tests__/unit/models/AttentionItem.test.js @@ -0,0 +1,38 @@ +const mongoose = require('mongoose'); +const AttentionItem = require('../../../models/AttentionItem'); +const { setupMongoDb, closeMongoDb, clearMongoDb } = require('../../utils/testUtils'); + +describe('AttentionItem', () => { + const recipient = new mongoose.Types.ObjectId(); + const pod = new mongoose.Types.ObjectId(); + + it('keeps a recipient/source fact unique and persists source snapshots', async () => { + await AttentionItem.create({ + recipientUserId: recipient, + podId: pod, + kind: 'mention', + source: { type: 'message', id: '42' }, + title: 'Ada mentioned you', + detail: '@sam please review', + messageId: '42', + threadRootId: '41', + }); + + await expect(AttentionItem.create({ + recipientUserId: recipient, + podId: pod, + kind: 'mention', + source: { type: 'message', id: '42' }, + title: 'duplicate', + })).rejects.toMatchObject({ code: 11000 }); + + const row = await AttentionItem.findOne({ recipientUserId: recipient }).lean(); + expect(row).toMatchObject({ + status: 'open', kind: 'mention', messageId: '42', threadRootId: '41', + source: { type: 'message', id: '42' }, + }); + }); +}); + beforeAll(async () => { await setupMongoDb(); }); + afterAll(async () => { await closeMongoDb(); }); + afterEach(async () => { await clearMongoDb(); }); diff --git a/backend/__tests__/unit/services/activityService.decisionQueue.test.js b/backend/__tests__/unit/services/activityService.decisionQueue.test.js index 659b2b217..67a13285b 100644 --- a/backend/__tests__/unit/services/activityService.decisionQueue.test.js +++ b/backend/__tests__/unit/services/activityService.decisionQueue.test.js @@ -39,7 +39,9 @@ const decision = (overrides = {}) => ({ createdAt: new Date('2026-09-01T12:00:00Z'), ...overrides, }); -describe('ActivityService.getDecisionQueue', () => { +// TASK-112 replaces this reconstructed reader with AttentionItem source-write +// tests. Its task-prose / history assertions are intentionally retired. +describe.skip('retired ActivityService.getDecisionQueue reconstruction', () => { beforeEach(() => { jest.restoreAllMocks(); Pod.find.mockReturnValue({ select: () => ({ lean: async () => [POD] }) }); diff --git a/backend/__tests__/unit/services/activityService.recap.test.js b/backend/__tests__/unit/services/activityService.recap.test.js index 9d9db0d46..670c9b08d 100644 --- a/backend/__tests__/unit/services/activityService.recap.test.js +++ b/backend/__tests__/unit/services/activityService.recap.test.js @@ -22,7 +22,9 @@ const taskQuery = (tasks) => ({ }), }); -describe('ActivityService.getRecap', () => { +// TASK-112 removes the acknowledgement array and read-time queue projection. +// Recap itself remains covered by browser/UI contract; retire obsolete queue fixtures. +describe.skip('retired ActivityService.getRecap queue projection', () => { let spy; let findByIdSpy; let pendingApprovalsSpy; diff --git a/backend/__tests__/unit/services/attentionItemService.test.js b/backend/__tests__/unit/services/attentionItemService.test.js new file mode 100644 index 000000000..3bd103a47 --- /dev/null +++ b/backend/__tests__/unit/services/attentionItemService.test.js @@ -0,0 +1,60 @@ +const mockUpdateOne = jest.fn(); +const mockUpdateMany = jest.fn(); +const mockFind = jest.fn(); +const mockPodFindById = jest.fn(); +const mockPodFind = jest.fn(); +const mockUserFind = jest.fn(); + +jest.mock('../../../models/AttentionItem', () => ({ updateOne: mockUpdateOne, updateMany: mockUpdateMany, find: mockFind })); +jest.mock('../../../models/Pod', () => ({ findById: mockPodFindById, find: mockPodFind })); +jest.mock('../../../models/User', () => ({ find: mockUserFind })); + +const chain = (value) => ({ select: () => ({ lean: async () => value }) }); +const AttentionItemService = require('../../../services/attentionItemService'); + +describe('attentionItemService', () => { + beforeEach(() => { + jest.clearAllMocks(); + mockUpdateOne.mockResolvedValue({ modifiedCount: 1 }); + mockUpdateMany.mockResolvedValue({ modifiedCount: 1 }); + }); + + it('materializes a mention only for mentioned human members other than the author', async () => { + mockPodFindById.mockReturnValue(chain({ _id: 'pod-1', name: 'Ship room', createdBy: 'owner', members: [{ userId: 'sam' }, { userId: 'bot' }] })); + mockUserFind.mockReturnValue(chain([ + { _id: 'owner', username: 'owner', isBot: false }, + { _id: 'sam', username: 'Sam', isBot: false }, + { _id: 'bot', username: 'Scout', isBot: true }, + ])); + + await AttentionItemService.recordMentionedUsers({ id: 42, podId: 'pod-1', userId: 'owner', username: 'Ada', content: '@sam please review; @samantha is different' }); + + expect(mockUpdateOne).toHaveBeenCalledTimes(1); + expect(mockUpdateOne.mock.calls[0][0]).toEqual({ recipientUserId: 'sam', 'source.type': 'message', 'source.id': '42' }); + expect(mockUpdateOne.mock.calls[0][1].$setOnInsert).toMatchObject({ kind: 'mention', title: 'Ada mentioned you', messageId: '42' }); + }); + + it('returns only rows whose recipient is still a member and resolves by recipient-owned id', async () => { + mockFind.mockReturnValue({ sort: () => ({ limit: () => ({ lean: async () => [ + { _id: 'attention-1', recipientUserId: '507f191e810c19729de860ea', podId: 'pod-1', kind: 'mention', source: { type: 'message', id: '41' }, title: 'Mention', createdAt: new Date() }, + { _id: 'attention-2', recipientUserId: '507f191e810c19729de860ea', podId: 'pod-2', kind: 'approval', source: { type: 'approval', id: 'a-1' }, title: 'Old access', createdAt: new Date() }, + ] }) }) }); + mockPodFind.mockReturnValue(chain([ + { _id: 'pod-1', name: 'Current', createdBy: '507f191e810c19729de860ea', members: [] }, + { _id: 'pod-2', name: 'Removed', createdBy: 'someone-else', members: [] }, + ])); + + const queue = await AttentionItemService.getOpenQueue('507f191e810c19729de860ea'); + expect(queue.items).toEqual([expect.objectContaining({ id: '41', attentionItemId: 'attention-1', podName: 'Current' })]); + await AttentionItemService.acknowledgeMention('507f191e810c19729de860ea', '507f191e810c19729de860eb'); + expect(mockUpdateOne).toHaveBeenLastCalledWith( + expect.objectContaining({ recipientUserId: '507f191e810c19729de860ea', kind: 'mention' }), + expect.any(Object), + ); + }); + + it('does not let projection-resolution storage turn a completed source into a failure', async () => { + mockUpdateMany.mockRejectedValueOnce(new Error('mongo unavailable')); + await expect(AttentionItemService.resolve('approval', 'a-1')).resolves.toBeUndefined(); + }); +}); diff --git a/backend/models/Activity.ts b/backend/models/Activity.ts index af575a2f6..8d7c3a344 100644 --- a/backend/models/Activity.ts +++ b/backend/models/Activity.ts @@ -174,7 +174,7 @@ activitySchema.statics.createSkillActivity = async function (summary, pod) { activitySchema.statics.createApprovalRequest = async function (options) { const { podId, requestedBy, agentName, scopes, content } = options; - return this.create({ + const approval = await this.create({ type: 'approval_needed', actor: { id: null, name: 'commonly-bot', type: 'system', verified: true }, action: 'approval_needed', @@ -183,6 +183,12 @@ activitySchema.statics.createApprovalRequest = async function (options) { approval: { status: 'pending', requestedBy, requestedScopes: scopes }, agentMetadata: { agentName }, }); + // The source row remains authoritative. Attention is a per-recipient + // projection written at the same boundary, never reconstructed by a read. + // eslint-disable-next-line global-require + const { recordApproval } = require('../services/attentionItemService'); + await recordApproval(approval); + return approval; }; activitySchema.statics.getFeedForUser = async function (userId, pods, options = {}) { diff --git a/backend/models/AttentionItem.ts b/backend/models/AttentionItem.ts new file mode 100644 index 000000000..25aadfd42 --- /dev/null +++ b/backend/models/AttentionItem.ts @@ -0,0 +1,58 @@ +import mongoose, { Document, Model, Schema, Types } from 'mongoose'; + +export type AttentionKind = 'mention' | 'approval' | 'decision'; +export type AttentionSourceType = 'message' | 'approval' | 'decision_request'; + +export interface IAttentionItem extends Document { + recipientUserId: Types.ObjectId; + podId: Types.ObjectId; + kind: AttentionKind; + source: { type: AttentionSourceType; id: string }; + title: string; + detail?: string; + podName?: string; + messageId?: string; + threadRootId?: string; + options?: Array<{ label: string; description?: string; recommended?: boolean }>; + status: 'open' | 'resolved'; + resolvedAt?: Date; + createdAt: Date; + updatedAt: Date; +} + +const optionSchema = new Schema({ + label: { type: String, required: true }, + description: { type: String }, + recommended: { type: Boolean }, +}, { _id: false }); + +const attentionItemSchema = new Schema({ + recipientUserId: { type: Schema.Types.ObjectId, ref: 'User', required: true }, + podId: { type: Schema.Types.ObjectId, ref: 'Pod', required: true }, + kind: { type: String, enum: ['mention', 'approval', 'decision'], required: true }, + source: { + type: { type: String, enum: ['message', 'approval', 'decision_request'], required: true }, + id: { type: String, required: true }, + }, + title: { type: String, required: true }, + detail: { type: String }, + podName: { type: String }, + messageId: { type: String }, + threadRootId: { type: String }, + options: [optionSchema], + status: { type: String, enum: ['open', 'resolved'], default: 'open', required: true }, + resolvedAt: { type: Date }, +}, { timestamps: true }); + +// A source can notify each recipient once. Retried source writes must not +// duplicate cards or resurrect a recipient's acknowledgement. +attentionItemSchema.index({ recipientUserId: 1, 'source.type': 1, 'source.id': 1 }, { unique: true }); +attentionItemSchema.index({ recipientUserId: 1, status: 1, createdAt: -1 }); +attentionItemSchema.index({ podId: 1, status: 1, createdAt: -1 }); + +const AttentionItem: Model = mongoose.models.AttentionItem + || mongoose.model('AttentionItem', attentionItemSchema); + +export default AttentionItem; +// eslint-disable-next-line @typescript-eslint/no-require-imports +module.exports = AttentionItem; diff --git a/backend/models/Message.ts b/backend/models/Message.ts index 6f62d038a..b187faa9f 100644 --- a/backend/models/Message.ts +++ b/backend/models/Message.ts @@ -37,6 +37,21 @@ const MessageSchema = new Schema({ updatedAt: { type: Date, default: Date.now }, }); +MessageSchema.post('save', async (doc) => { + // Mongo is the availability fallback for chat writes, so it must materialize + // the same recipient-owned fact as the normal PostgreSQL writer. + // eslint-disable-next-line global-require + const { recordMentionedUsers } = require('../services/attentionItemService'); + await recordMentionedUsers(doc); +}); + +MessageSchema.post('findOneAndDelete', async (doc) => { + if (!doc) return; + // eslint-disable-next-line global-require + const { resolve } = require('../services/attentionItemService'); + await resolve('message', doc._id); +}); + export const Message: Model = mongoose.model('Message', MessageSchema); export default Message; diff --git a/backend/models/User.ts b/backend/models/User.ts index a93c22abf..409db4444 100644 --- a/backend/models/User.ts +++ b/backend/models/User.ts @@ -164,11 +164,6 @@ export interface IUser extends Document { lastViewedAt: Date; readItemIds: string[]; }; - // Queue acknowledgement is deliberately separate from activityFeed read - // state. Opening a feed is not the same as resolving a direct mention. - activityQueue: { - acknowledgedMentionIds: string[]; - }; digestPreferences: { enabled: boolean; frequency: DigestFrequency; @@ -371,9 +366,6 @@ const userSchema = new Schema({ lastViewedAt: { type: Date, default: new Date(0) }, readItemIds: { type: [String], default: [] }, }, - activityQueue: { - acknowledgedMentionIds: { type: [String], default: [] }, - }, digestPreferences: { enabled: { type: Boolean, default: true }, frequency: { type: String, enum: ['daily', 'weekly', 'never'], default: 'daily' }, diff --git a/backend/models/pg/Message.ts b/backend/models/pg/Message.ts index 281429e15..0778026ed 100644 --- a/backend/models/pg/Message.ts +++ b/backend/models/pg/Message.ts @@ -152,7 +152,11 @@ class Message { 'UPDATE pods SET updated_at = CURRENT_TIMESTAMP WHERE id = $1', [podId], ); - return result.rows[0] as unknown as MessageRow; + const row = result.rows[0] as unknown as MessageRow; + // eslint-disable-next-line global-require + const { recordMentionedUsers } = require('../../services/attentionItemService'); + await recordMentionedUsers({ ...row, podId, userId, content, threadRootId: resolvedThreadRootId }); + return row; } catch (error) { const e = error as { message?: string }; console.error('SQL Error in Message.create:', e.message); @@ -280,13 +284,21 @@ class Message { static async delete(id: string): Promise { const query = `DELETE FROM messages WHERE id = $1 RETURNING *`; const result = await (pool as PgPool).query(query, [id]); - return result.rows[0] as unknown as MessageRow | undefined; + const deleted = result.rows[0] as unknown as MessageRow | undefined; + // eslint-disable-next-line global-require + const { resolve } = require('../../services/attentionItemService'); + await resolve('message', id); + return deleted; } static async deleteByPodId(podId: string): Promise { const query = `DELETE FROM messages WHERE pod_id = $1 RETURNING *`; const result = await (pool as PgPool).query(query, [podId]); - return result.rows as unknown as MessageRow[]; + const rows = result.rows as unknown as MessageRow[]; + // eslint-disable-next-line global-require + const { resolveMany } = require('../../services/attentionItemService'); + await resolveMany('message', rows.map((row) => row.id)); + return rows; } /** @@ -411,6 +423,9 @@ class Message { const deleted = typeof result.rowCount === 'number' ? result.rowCount : (Array.isArray(result.rows) ? result.rows.length : 0); + // eslint-disable-next-line global-require + const { resolveMany } = require('../../services/attentionItemService'); + await resolveMany('message', (result.rows || []).map((row) => row.id)); // Repair before returning. Deleting a thread root orphans everything below // it (see reRootOrphanedChains), and the caller has no way to know it diff --git a/backend/services/activityService.ts b/backend/services/activityService.ts index 7a9875ce8..ddbaa61bd 100644 --- a/backend/services/activityService.ts +++ b/backend/services/activityService.ts @@ -10,9 +10,12 @@ const Summary = require('../models/Summary'); const Post = require('../models/Post'); // eslint-disable-next-line global-require const Task = require('../models/Task'); +// Kept while the dead former implementation below is mechanically removed; +// getDecisionQueue returns before this is ever read. // eslint-disable-next-line global-require const DecisionRequest = require('../models/DecisionRequest'); // eslint-disable-next-line global-require +// eslint-disable-next-line global-require const isPodMember = require('../utils/isPodMember'); let PGMessage: unknown = null; @@ -73,7 +76,6 @@ interface UserDoc { followers?: unknown[]; followedThreads?: Array<{ postId: unknown; followedAt?: Date }>; activityFeed?: { lastViewedAt?: Date | string; readItemIds?: unknown[] }; - activityQueue?: { acknowledgedMentionIds?: unknown[] }; } interface PodDoc { @@ -154,9 +156,8 @@ class ActivityService { return timestamp >= since.getTime() && (!requestedPodId || (activity.pod && scopedPodIds.has(activity.pod.id))); }); - const acknowledgedMentionIds = new Set( - ((feed.acknowledgedMentionIds as unknown[] | undefined) || []).map((id) => String(id)), - ); + // Needs-you is not derived from this display feed. It is materialized by + // its source writers and rechecked against current membership below. type AgentRecap = { id: string; @@ -285,10 +286,10 @@ class ActivityService { .filter(isPendingApproval) .sort(newestFirst); const mentionQueue = Array.from(queueCandidates.values()) - .filter((activity) => activity.flags?.isMention && !acknowledgedMentionIds.has(String(activity.id))) + .filter((activity) => activity.flags?.isMention) .sort(newestFirst) .slice(0, Math.max(0, 12 - approvalQueue.length)); - const needsYou = [...approvalQueue, ...mentionQueue] + let needsYou = [...approvalQueue, ...mentionQueue] .sort(newestFirst) .map((activity) => { const isApproval = isPendingApproval(activity); @@ -303,6 +304,17 @@ class ActivityService { }; }); + // Recipient-owned AttentionItem rows are the only queue source. The + // temporary calculation above remains only for the normal activity feed; + // it cannot become a queue fallback or revive historical heuristics. + // eslint-disable-next-line global-require + const AttentionItemService = require('./attentionItemService'); + const attention = await AttentionItemService.getOpenQueue(userId); + needsYou = attention.items.map((item: any) => ({ + ...item, + timestamp: item.createdAt || null, + })); + let board: Array> = []; if (scopedPods.length > 0) { const taskRows: Array> = await Task.find({ @@ -360,7 +372,7 @@ class ActivityService { try { const user: UserDoc | null = await User.findById(userId) - .select('_id username following followers followedThreads activityQueue') + .select('_id username following followers followedThreads') .lean(); if (!user) { return { activities: [], hasMore: false, quick: null }; @@ -394,7 +406,6 @@ class ActivityService { hasMore: withReadState.length === limit, quick, unreadCount: withReadState.filter((item) => !item.read).length, - acknowledgedMentionIds: user.activityQueue?.acknowledgedMentionIds || [], }; } catch (error) { console.error('Error in getUserFeed:', error); @@ -489,6 +500,9 @@ class ActivityService { } static async getDecisionQueue(userId: unknown): Promise> { + // This guarded former reader is retained for one patch only while this + // large method is mechanically removed. It is never a production fallback. + if (process.env.ATTENTION_ITEM_LEGACY_READER_FOR_TESTS === '1') { const pods: PodDoc[] = await Pod.find({ $or: [ { createdBy: userId }, @@ -528,7 +542,7 @@ class ActivityService { kind: 'approval', // RAW Activity id — the existing Activity actions key on it. id: String(a._id), - title: a.agentMetadata?.agentName ? `${a.agentMetadata.agentName} requests approval` : 'Approval requested', + title: a.agentMetadata?.agentName ? `${a.agentMetadata?.agentName} requests approval` : 'Approval requested', detail: String(a.content || '').slice(0, 160), podId: a.podId ? String(a.podId) : null, podName: a.podId ? podName.get(String(a.podId)) : undefined, @@ -701,6 +715,12 @@ class ActivityService { // No mention yet is not an error; the page falls back to its first pod. composePodId, }; + } + // The recipient-scoped read is indexed and authorization-safe. Do not + // reconstruct this from messages, Task prose, or historical acknowledgements. + // eslint-disable-next-line global-require + const AttentionItemService = require('./attentionItemService'); + return AttentionItemService.getOpenQueue(userId); } static async getPodFeed( @@ -1145,18 +1165,9 @@ class ActivityService { } static async acknowledgeMention(userId: unknown, activityId: string): Promise> { - const user = await User.findById(userId).select('_id activityQueue'); - if (!user) return { success: false, error: 'User not found' }; - - if (!user.activityQueue) user.activityQueue = { acknowledgedMentionIds: [] }; - const next = new Set((user.activityQueue.acknowledgedMentionIds || []).map((id: unknown) => String(id))); - next.add(String(activityId)); - // This is per-(user, message) state, rather than a recent-feed cache: an - // acknowledged mention must not resurface merely because more messages - // arrive later. - user.activityQueue.acknowledgedMentionIds = Array.from(next); - await user.save(); - return { success: true, acknowledgedMentionIds: user.activityQueue.acknowledgedMentionIds }; + // eslint-disable-next-line global-require + const AttentionItemService = require('./attentionItemService'); + return AttentionItemService.acknowledgeMention(userId, activityId); } static async getUnreadCount(userId: unknown, options: GetFeedOptions = {}): Promise<{ unreadCount: number }> { @@ -1475,6 +1486,9 @@ class ActivityService { if (membershipError) return membershipError; await activity.approve(userId, notes); + // eslint-disable-next-line global-require + const { resolve } = require('./attentionItemService'); + await resolve('approval', activity._id); return { success: true, status: 'approved' }; } catch (error) { const err = error as { message?: string }; @@ -1498,6 +1512,9 @@ class ActivityService { if (membershipError) return membershipError; await activity.reject(userId, notes); + // eslint-disable-next-line global-require + const { resolve } = require('./attentionItemService'); + await resolve('approval', activity._id); return { success: true, status: 'rejected' }; } catch (error) { const err = error as { message?: string }; diff --git a/backend/services/attentionItemService.ts b/backend/services/attentionItemService.ts new file mode 100644 index 000000000..4b7f3e14d --- /dev/null +++ b/backend/services/attentionItemService.ts @@ -0,0 +1,183 @@ +// Recipient-owned attention is written beside its authoritative source. It is +// intentionally not rebuilt from message history or task prose at read time. +// There is no historical backfill: facts start materializing when this ships. +// eslint-disable-next-line global-require +const AttentionItem = require('../models/AttentionItem'); +// eslint-disable-next-line global-require +const Pod = require('../models/Pod'); +// eslint-disable-next-line global-require +const User = require('../models/User'); + +type SourceType = 'message' | 'approval' | 'decision_request'; +type Kind = 'mention' | 'approval' | 'decision'; + +const compact = (value: unknown, max = 220): string => String(value || '').replace(/\s+/g, ' ').trim().slice(0, max); +const sourceKey = (type: SourceType, id: unknown): string => String(id || '').trim(); +const isCurrentMember = (pod: any, userId: unknown): boolean => ( + String(pod?.createdBy || '') === String(userId) + || (pod?.members || []).some((member: any) => String(member?.userId || member?._id || member) === String(userId)) +); + +const currentHumanMembers = async (podId: unknown): Promise> => { + const pod = await Pod.findById(podId).select('_id name createdBy members').lean(); + if (!pod) return []; + const ids = new Set([String(pod.createdBy || '')]); + for (const member of (pod.members || [])) { + const id = (member as any)?.userId || (member as any)?._id || member; + if (id) ids.add(String(id)); + } + const users = await User.find({ _id: { $in: [...ids].filter(Boolean) }, isBot: { $ne: true } }) + .select('_id username isBot').lean(); + return users.filter((user: any) => isCurrentMember(pod, user._id)); +}; + +const recordForRecipients = async ( + recipients: Array<{ _id: unknown }>, + payload: Record, +): Promise => { + await Promise.all(recipients.map((recipient) => AttentionItem.updateOne( + { recipientUserId: recipient._id, 'source.type': payload.sourceType, 'source.id': payload.sourceId }, + { + $setOnInsert: { + recipientUserId: recipient._id, + podId: payload.podId, + kind: payload.kind, + source: { type: payload.sourceType, id: payload.sourceId }, + title: payload.title, + detail: payload.detail, + podName: payload.podName, + messageId: payload.messageId, + threadRootId: payload.threadRootId, + options: payload.options, + status: 'open', + }, + }, + { upsert: true }, + ))); +}; + +export const recordMentionedUsers = async (message: any): Promise => { + try { + const podId = message?.podId || message?.pod_id; + const messageId = message?._id || message?.id; + if (!podId || messageId === undefined || messageId === null) return; + const authorId = message?.userId?._id || message?.userId || message?.user_id; + const content = String(message?.content || message?.text || ''); + const members = await currentHumanMembers(podId); + const recipients = members.filter((member) => { + const handle = String(member.username || '').trim(); + if (!handle || String(member._id) === String(authorId)) return false; + return new RegExp(`(^|[^A-Za-z0-9_-])@${handle.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')}(?![A-Za-z0-9_-])`, 'i').test(content); + }); + if (!recipients.length) return; + const pod = await Pod.findById(podId).select('name').lean(); + const authorName = message?.username || message?.userId?.username || 'Someone'; + await recordForRecipients(recipients, { + podId, kind: 'mention' as Kind, sourceType: 'message' as SourceType, sourceId: sourceKey('message', messageId), + title: `${authorName} mentioned you`, detail: compact(content), podName: pod?.name || 'Pod', + messageId: String(messageId), threadRootId: String(message?.threadRootId || message?.thread_root_id || messageId), + }); + } catch (error) { + console.warn('[attention] mention materialization failed:', (error as Error).message); + } +}; + +export const recordApproval = async (approval: any): Promise => { + try { + const podId = approval?.podId; + const id = approval?._id || approval?.id; + if (!podId || !id) return; + const recipients = await currentHumanMembers(podId); + const pod = await Pod.findById(podId).select('name').lean(); + const agentName = approval?.agentMetadata?.agentName; + await recordForRecipients(recipients, { + podId, kind: 'approval' as Kind, sourceType: 'approval' as SourceType, sourceId: sourceKey('approval', id), + title: agentName ? `${agentName} requests approval` : 'Approval requested', detail: compact(approval?.content, 180), podName: pod?.name || 'Pod', + }); + } catch (error) { + console.warn('[attention] approval materialization failed:', (error as Error).message); + } +}; + +export const recordDecision = async (decision: any): Promise => { + try { + const podId = decision?.podId; + const id = decision?._id || decision?.id; + if (!podId || !id) return; + const recipients = await currentHumanMembers(podId); + const pod = await Pod.findById(podId).select('name').lean(); + const options = (decision.options || []).filter((option: any) => option?.label).map((option: any) => ({ + label: String(option.label), ...(option.description ? { description: String(option.description) } : {}), + ...(option.recommended ? { recommended: true } : {}), + })); + await recordForRecipients(recipients, { + podId, kind: 'decision' as Kind, sourceType: 'decision_request' as SourceType, sourceId: sourceKey('decision_request', id), + title: String(decision.title || 'Decision requested'), detail: compact(decision.question || decision.context, 1000), + podName: pod?.name || 'Pod', messageId: decision.messageId ? String(decision.messageId) : undefined, + threadRootId: String(decision.threadRootId || decision.messageId || ''), options, + }); + } catch (error) { + console.warn('[attention] decision materialization failed:', (error as Error).message); + } +}; + +export const resolve = async (sourceType: SourceType, sourceId: unknown): Promise => { + const id = sourceKey(sourceType, sourceId); + if (!id) return; + try { + await AttentionItem.updateMany({ 'source.type': sourceType, 'source.id': id, status: 'open' }, { $set: { status: 'resolved', resolvedAt: new Date() } }); + } catch (error) { + // The owning source already completed. Leave stale attention visible over + // returning a false failure from a completed source action. + console.warn('[attention] resolution storage failed; leaving item visible:', (error as Error).message); + } +}; + +export const resolveMany = async (sourceType: SourceType, sourceIds: unknown[]): Promise => { + const ids = sourceIds.map((id) => sourceKey(sourceType, id)).filter(Boolean); + if (!ids.length) return; + try { + await AttentionItem.updateMany({ 'source.type': sourceType, 'source.id': { $in: ids }, status: 'open' }, { $set: { status: 'resolved', resolvedAt: new Date() } }); + } catch (error) { + console.warn('[attention] bulk resolution storage failed; leaving items visible:', (error as Error).message); + } +}; + +export const getOpenQueue = async (recipientUserId: unknown): Promise<{ items: any[]; count: number; composePodId: string | null }> => { + // Route callers carry a real Mongo id. Returning an empty queue for a bad + // value keeps malformed/read-only callers from turning a cast error into a + // 500 and makes the authorization boundary explicit. + if (!/^[a-f\d]{24}$/i.test(String(recipientUserId))) return { items: [], count: 0, composePodId: null }; + const rows = await AttentionItem.find({ recipientUserId, status: 'open' }).sort({ createdAt: -1 }).limit(80).lean(); + const podIds = [...new Set(rows.map((row: any) => String(row.podId)))]; + const pods = await Pod.find({ _id: { $in: podIds } }).select('_id name createdBy members').lean(); + const allowed = new Map(pods.filter((pod: any) => isCurrentMember(pod, recipientUserId)).map((pod: any) => [String(pod._id), pod])); + const priority: Record = { approval: 0, decision: 1, mention: 2 }; + const valid = rows.filter((row: any) => allowed.has(String(row.podId))).sort((a: any, b: any) => ( + (priority[a.kind] ?? 9) - (priority[b.kind] ?? 9) + || new Date(b.createdAt).getTime() - new Date(a.createdAt).getTime() + )); + const picked: any[] = []; + let mentionCount = 0; + for (const row of valid) { + if (picked.length >= 12) break; + if (row.kind === 'mention' && mentionCount >= 8) continue; + if (row.kind === 'mention') mentionCount += 1; + picked.push({ + id: String(row.source.id), attentionItemId: String(row._id), kind: row.kind, title: row.title, detail: row.detail || '', + podId: String(row.podId), podName: (allowed.get(String(row.podId)) as any)?.name || row.podName || 'Pod', + messageId: row.messageId, threadRootId: row.threadRootId, options: row.options || [], createdAt: row.createdAt, + }); + } + return { items: picked, count: valid.length, composePodId: picked.find((row) => row.kind === 'mention')?.podId || null }; +}; + +export const acknowledgeMention = async (recipientUserId: unknown, attentionItemId: string): Promise<{ success: boolean; error?: string }> => { + if (!/^[a-f\d]{24}$/i.test(String(attentionItemId))) return { success: false, error: 'Invalid attention item' }; + const result = await AttentionItem.updateOne({ _id: attentionItemId, recipientUserId, kind: 'mention', status: 'open' }, { $set: { status: 'resolved', resolvedAt: new Date() } }); + return result.modifiedCount === 1 ? { success: true } : { success: false, error: 'Attention item not found' }; +}; + +export default { recordMentionedUsers, recordApproval, recordDecision, resolve, resolveMany, getOpenQueue, acknowledgeMention }; +// eslint-disable-next-line @typescript-eslint/no-require-imports +module.exports = { recordMentionedUsers, recordApproval, recordDecision, resolve, resolveMany, getOpenQueue, acknowledgeMention }; diff --git a/backend/services/decisionRequestService.ts b/backend/services/decisionRequestService.ts index edcd7ad60..9c9429aee 100644 --- a/backend/services/decisionRequestService.ts +++ b/backend/services/decisionRequestService.ts @@ -214,6 +214,9 @@ export const requestDecision = async (input: RequestDecisionOptions): Promise { const decisionQueue = { items: [ { - id: 'mention-1', kind: 'mention', title: 'Review requested', detail: 'A direct mention.', + id: 'mention-1', attentionItemId: 'attention-1', kind: 'mention', title: 'Review requested', detail: 'A direct mention.', podId: 'pod-1', podName: 'Launch pod', createdAt: '2026-08-26T11:00:00.000Z', }, - { - id: 'task_TASK-059', kind: 'press', title: 'Retention ledger', detail: '#1208 held for the human merge press.', - podId: 'pod-1', podName: 'Launch pod', taskId: 'TASK-059', createdAt: '2026-08-26T10:00:00.000Z', - }, { id: 'decision-024', kind: 'decision', title: 'Choose the eslint scope', detail: 'What should the agent do?', podId: 'pod-1', podName: 'Launch pod', messageId: '700', threadRootId: '695', options: [ @@ -41,7 +37,7 @@ const decisionQueue = { ], createdAt: '2026-08-26T09:00:00.000Z', }, ], - count: 3, + count: 2, composePodId: 'pod-1', }; @@ -102,14 +98,11 @@ describe('V2ActivityPage', () => { // adds a microtask hop the old single-request race happened to win. expect(await screen.findByRole('heading', { name: 'Needs you' })).toBeInTheDocument(); expect(screen.getByText('Review requested')).toBeInTheDocument(); - // A board press remains an Open-board action; DecisionRequest cards use - // their declared options rather than deriving actions from task prose. - expect(screen.getByText('Retention ledger')).toBeInTheDocument(); - expect(screen.getByText('Ready for your press')).toBeInTheDocument(); + // Queue rows are only durable source facts; task handoff prose never + // creates a card. DecisionRequest cards use declared alternatives. expect(screen.getByText('Choose the eslint scope')).toBeInTheDocument(); expect(screen.getByRole('button', { name: 'Rule: Ship now' })).toBeInTheDocument(); expect(screen.getByText('Release the bounded change.')).toBeInTheDocument(); - expect(screen.getAllByRole('button', { name: 'Open board' })).toHaveLength(1); expect(screen.getByRole('button', { name: 'Other…' })).toBeInTheDocument(); expect(mockGet).toHaveBeenCalledWith('/api/activity/decision-queue', expect.anything()); expect(screen.getByRole('heading', { name: 'What your agents did' })).toBeInTheDocument(); diff --git a/frontend/src/v2/components/V2ActivityPage.tsx b/frontend/src/v2/components/V2ActivityPage.tsx index 57bd59a38..564caedad 100644 --- a/frontend/src/v2/components/V2ActivityPage.tsx +++ b/frontend/src/v2/components/V2ActivityPage.tsx @@ -27,14 +27,12 @@ interface AgentRecap { interface NeedsYouItem { id: string; - // Decision cards are agent-authored DecisionRequest rows. `press` remains - // an ordinary board handoff and never parses a task title into options. - kind: 'mention' | 'approval' | 'decision' | 'press'; + attentionItemId?: string; + kind: 'mention' | 'approval' | 'decision'; title: string; detail: string; podId: string | null; podName: string; - taskId?: string; options?: Array<{ label: string; description?: string; recommended?: boolean }>; timestamp: string | null; // Mention rows carry where they live so a reply can land IN the thread. @@ -266,7 +264,7 @@ const V2ActivityPage: React.FC = () => { ); setRepliedIds((prev) => new Set(prev).add(item.id)); setReplyDrafts((prev) => ({ ...prev, [item.id]: '' })); - await axios.post(`/api/activity/${item.id}/acknowledge`, {}, { headers: { 'x-auth-token': token ?? '' } }).catch(() => null); + if (item.attentionItemId) await axios.post(`/api/activity/${item.attentionItemId}/acknowledge`, {}, { headers: { 'x-auth-token': token ?? '' } }).catch(() => null); setReloadKey((value) => value + 1); } catch { setActionError(t('activity.mention.actionFailed')); @@ -282,7 +280,7 @@ const V2ActivityPage: React.FC = () => { try { const token = localStorage.getItem('token'); const response = await axios.post<{ success?: boolean }>( - `/api/activity/${item.id}/acknowledge`, + `/api/activity/${item.attentionItemId || item.id}/acknowledge`, {}, { headers: { 'x-auth-token': token ?? '' } }, ); @@ -295,10 +293,6 @@ const V2ActivityPage: React.FC = () => { } }; - const openBoard = (item: NeedsYouItem) => { - if (item.podId) navigate(`/v2/pods/${item.podId}/board`); - }; - const isDayZero = podId === 'all' && queue.length === 0 && recap?.agents.length === 0 @@ -428,7 +422,7 @@ const V2ActivityPage: React.FC = () => { {queue.map((item) => (
{t(`activity.needsYou.kinds.${item.kind}`)}
@@ -528,11 +522,6 @@ const V2ActivityPage: React.FC = () => { )} )} - {item.kind === 'press' && ( - - )} From 5387b8d1d16dc39136a9500a1d41c025043df804 Mon Sep 17 00:00:00 2001 From: Lily Shen <115414357+lilyshen0722@users.noreply.github.com> Date: Thu, 3 Sep 2026 04:56:10 -0700 Subject: [PATCH 2/2] fix(activity): materialize task attention and backfill --- .../unit/models/AttentionItem.test.js | 12 + .../routes/tasksApi.status-vocabulary.test.js | 35 ++ .../routes/tasksApi.updateRenewsLease.test.js | 22 +- .../activityService.attentionQueue.test.js | 52 +++ .../activityService.decisionQueue.test.js | 146 ------- .../services/activityService.recap.test.js | 257 ++---------- .../services/attentionItemService.test.js | 90 +++++ backend/models/AttentionItem.ts | 4 +- backend/package.json | 1 + backend/routes/tasksApi.ts | 5 + backend/scripts/backfill-attention-items.ts | 108 ++++++ backend/services/activityService.ts | 365 +----------------- backend/services/attentionItemService.ts | 59 ++- docs/adr/ADR-017-attention-routing.md | 2 + .../src/v2/__tests__/V2ActivityPage.test.tsx | 21 + frontend/src/v2/components/V2ActivityPage.tsx | 2 +- 16 files changed, 452 insertions(+), 729 deletions(-) create mode 100644 backend/__tests__/unit/services/activityService.attentionQueue.test.js delete mode 100644 backend/__tests__/unit/services/activityService.decisionQueue.test.js create mode 100644 backend/scripts/backfill-attention-items.ts diff --git a/backend/__tests__/unit/models/AttentionItem.test.js b/backend/__tests__/unit/models/AttentionItem.test.js index ac8b397e9..a70f1ec6b 100644 --- a/backend/__tests__/unit/models/AttentionItem.test.js +++ b/backend/__tests__/unit/models/AttentionItem.test.js @@ -32,6 +32,18 @@ describe('AttentionItem', () => { source: { type: 'message', id: '42' }, }); }); + + it('accepts a task source as a recipient-owned board attention fact', async () => { + const row = new AttentionItem({ + recipientUserId: recipient, + podId: pod, + kind: 'decision', + source: { type: 'task', id: 'task-1:update-1' }, + title: 'Choose a deploy shape', + }); + + await expect(row.validate()).resolves.toBeUndefined(); + }); }); beforeAll(async () => { await setupMongoDb(); }); afterAll(async () => { await closeMongoDb(); }); diff --git a/backend/__tests__/unit/routes/tasksApi.status-vocabulary.test.js b/backend/__tests__/unit/routes/tasksApi.status-vocabulary.test.js index 1006ab09b..e006cfe95 100644 --- a/backend/__tests__/unit/routes/tasksApi.status-vocabulary.test.js +++ b/backend/__tests__/unit/routes/tasksApi.status-vocabulary.test.js @@ -49,6 +49,13 @@ jest.mock('../../../services/taskEventService', () => ({ emitTaskUpdated: (...args) => mockEmitTaskUpdated(...args), })); +const mockRecordTaskAttention = jest.fn(); +const mockResolveTaskAttention = jest.fn(); +jest.mock('../../../services/attentionItemService', () => ({ + recordTaskAttention: (...args) => mockRecordTaskAttention(...args), + resolveTaskAttention: (...args) => mockResolveTaskAttention(...args), +})); + const tasksApi = require('../../../routes/tasksApi'); const app = express(); @@ -61,6 +68,10 @@ beforeEach(() => { mockFindOneAndUpdate.mockReset(); mockFindOneAndUpdate.mockResolvedValue({ taskId: 'TASK-001', status: 'claimed' }); mockEmitTaskUpdated.mockReset(); + mockRecordTaskAttention.mockReset(); + mockRecordTaskAttention.mockResolvedValue(undefined); + mockResolveTaskAttention.mockReset(); + mockResolveTaskAttention.mockResolvedValue(undefined); }); describe('PATCH /:podId/:taskId status vocabulary', () => { @@ -91,6 +102,30 @@ describe('PATCH /:podId/:taskId status vocabulary', () => { expect(mockFindOneAndUpdate.mock.calls[0][1].$set.status).toBe('blocked'); }); + test('materializes a blocked row at the status-write boundary', async () => { + const task = { taskId: 'TASK-001', status: 'blocked', podId: POD_ID, updates: [] }; + mockFindOneAndUpdate.mockResolvedValueOnce(task); + + const res = await request(app) + .patch(`/api/v1/tasks/${POD_ID}/TASK-001`) + .send({ status: 'blocked' }); + + expect(res.status).toBe(200); + expect(mockRecordTaskAttention).toHaveBeenCalledWith(task, { includeBlocked: true }); + }); + + test('resolves task attention when a task leaves blocked state', async () => { + const task = { taskId: 'TASK-001', status: 'done', podId: POD_ID, updates: [] }; + mockFindOneAndUpdate.mockResolvedValueOnce(task); + + const res = await request(app) + .patch(`/api/v1/tasks/${POD_ID}/TASK-001`) + .send({ status: 'done' }); + + expect(res.status).toBe(200); + expect(mockResolveTaskAttention).toHaveBeenCalledWith(task); + }); + test('rejects unknown statuses with the vocabulary in the error', async () => { const res = await request(app) .patch(`/api/v1/tasks/${POD_ID}/TASK-001`) diff --git a/backend/__tests__/unit/routes/tasksApi.updateRenewsLease.test.js b/backend/__tests__/unit/routes/tasksApi.updateRenewsLease.test.js index c1e23cd7a..af9ca4777 100644 --- a/backend/__tests__/unit/routes/tasksApi.updateRenewsLease.test.js +++ b/backend/__tests__/unit/routes/tasksApi.updateRenewsLease.test.js @@ -84,6 +84,10 @@ jest.mock('../../../models/User', () => ({ jest.mock('../../../services/githubAppService', () => ({ isPatConfigured: jest.fn(() => false) })); jest.mock('../../../services/taskEventService', () => ({ emitTaskUpdated: jest.fn() })); +const mockRecordTaskAttention = jest.fn(); +jest.mock('../../../services/attentionItemService', () => ({ + recordTaskAttention: (...args) => mockRecordTaskAttention(...args), +})); const Task = require('../../../models/Task'); const tasksApi = require('../../../routes/tasksApi'); @@ -125,7 +129,11 @@ afterAll(async () => { if (mongod) await mongod.stop(); }); -beforeEach(async () => { await Task.deleteMany({}); }); +beforeEach(async () => { + await Task.deleteMany({}); + mockRecordTaskAttention.mockReset(); + mockRecordTaskAttention.mockResolvedValue(undefined); +}); describe('POST /updates renews a holder-authored lease', () => { it('the holder\'s note pushes the lease out by a full period', async () => { @@ -159,6 +167,18 @@ describe('POST /updates renews a holder-authored lease', () => { expect(row.updates.map((u) => u.text)).toContain('reviewed at 9f91b9b0'); }); + it('hands every progress note to the task-attention source writer', async () => { + await seed({ claimedBy: 'holder', claimedAt: new Date(), claimExpiresAt: new Date(Date.now() + LEASE_MS) }); + + const res = await postUpdate('holder', 'TASK-001', 'Ready for Sam\'s ruling.'); + + expect(res.status).toBe(200); + expect(mockRecordTaskAttention).toHaveBeenCalledWith(expect.objectContaining({ + taskId: 'TASK-001', + updates: expect.arrayContaining([expect.objectContaining({ text: 'Ready for Sam\'s ruling.' })]), + })); + }); + it('an ALREADY-LAPSED holder still renews — the row was theirs and they are working', async () => { const claimedAt = new Date(Date.now() - 90 * MIN); const before = new Date(claimedAt.getTime() + LEASE_MS); diff --git a/backend/__tests__/unit/services/activityService.attentionQueue.test.js b/backend/__tests__/unit/services/activityService.attentionQueue.test.js new file mode 100644 index 000000000..61fa28451 --- /dev/null +++ b/backend/__tests__/unit/services/activityService.attentionQueue.test.js @@ -0,0 +1,52 @@ +const mockGetOpenQueue = jest.fn(); +const mockPodFind = jest.fn(); +const mockTaskFind = jest.fn(); + +jest.mock('../../../models/Pod', () => ({ find: (...args) => mockPodFind(...args) })); +jest.mock('../../../models/User', () => ({})); +jest.mock('../../../models/Activity', () => ({})); +jest.mock('../../../models/Summary', () => ({})); +jest.mock('../../../models/Post', () => ({})); +jest.mock('../../../models/Task', () => ({ find: (...args) => mockTaskFind(...args) })); +jest.mock('../../../services/attentionItemService', () => ({ + getOpenQueue: (...args) => mockGetOpenQueue(...args), +})); + +const ActivityService = require('../../../services/activityService'); +const chain = (value) => ({ select: () => ({ lean: async () => value }) }); +const taskChain = (value) => ({ + select: () => ({ sort: () => ({ limit: () => ({ lean: async () => value }) }) }), +}); + +describe('ActivityService.getDecisionQueue', () => { + beforeEach(() => { + jest.clearAllMocks(); + }); + + it('reads the recipient-owned AttentionItem projection without a source-store fallback', async () => { + const queue = { items: [{ id: '42', kind: 'mention' }], count: 1, composePodId: 'pod-1' }; + mockGetOpenQueue.mockResolvedValue(queue); + + await expect(ActivityService.getDecisionQueue('507f191e810c19729de860ea')).resolves.toBe(queue); + expect(mockGetOpenQueue).toHaveBeenCalledWith('507f191e810c19729de860ea'); + }); + + it('keeps the recap fallback scoped to the requested pod', async () => { + mockPodFind.mockReturnValue(chain([{ _id: 'pod-1', name: 'Current pod' }])); + mockTaskFind.mockReturnValue(taskChain([])); + mockGetOpenQueue.mockResolvedValue({ + items: [ + { id: 'first', podId: 'pod-1', createdAt: new Date() }, + { id: 'other', podId: 'pod-2', createdAt: new Date() }, + ], + count: 2, + composePodId: null, + }); + const feed = jest.spyOn(ActivityService, 'getUserFeed').mockResolvedValue({ activities: [] }); + + const recap = await ActivityService.getRecap('507f191e810c19729de860ea', { podId: 'pod-1' }); + + expect(recap.needsYou).toEqual([expect.objectContaining({ id: 'first' })]); + feed.mockRestore(); + }); +}); diff --git a/backend/__tests__/unit/services/activityService.decisionQueue.test.js b/backend/__tests__/unit/services/activityService.decisionQueue.test.js deleted file mode 100644 index 67a13285b..000000000 --- a/backend/__tests__/unit/services/activityService.decisionQueue.test.js +++ /dev/null @@ -1,146 +0,0 @@ -/** - * DecisionRequest is the only source of interactive decision cards. Board - * prose remains board prose; an OPTIONS line must never turn into an action - * surface merely because it looked structured to the reader. - */ - -jest.mock('../../../models/Pod', () => ({ - find: jest.fn(), - findById: jest.fn(), -})); -jest.mock('../../../models/Task', () => ({ find: jest.fn() })); -jest.mock('../../../models/User', () => ({ findById: jest.fn() })); -jest.mock('../../../models/DecisionRequest', () => ({ find: jest.fn() })); -jest.mock('../../../config/db-pg', () => ({ pool: { query: jest.fn().mockResolvedValue({ rows: [] }) } })); -jest.mock('../../../models/Activity', () => ({})); -jest.mock('../../../models/Summary', () => ({})); -jest.mock('../../../models/Post', () => ({})); - -const ActivityService = require('../../../services/activityService'); -const Pod = require('../../../models/Pod'); -const Task = require('../../../models/Task'); -const User = require('../../../models/User'); -const DecisionRequest = require('../../../models/DecisionRequest'); -const { pool } = require('../../../config/db-pg'); - -const POD = { _id: 'pod-1', name: 'Sprint HQ', type: 'team' }; -const userChain = (doc) => ({ select: () => ({ lean: async () => doc }) }); -const taskChain = (rows) => ({ sort: () => ({ limit: () => ({ lean: async () => rows }) }) }); -const decisionChain = (rows) => ({ sort: () => ({ limit: () => ({ lean: async () => rows }) }) }); - -const decision = (overrides = {}) => ({ - _id: 'decision-1', podId: 'pod-1', title: 'Choose deployment shape', - question: 'Which release train should this agent follow?', - options: [ - { label: 'Fast lane', description: 'Ship after the green suite.' }, - { label: 'Canary', description: 'Roll out gradually.', recommended: true }, - ], - messageId: '700', threadRootId: '695', status: 'pending', - createdAt: new Date('2026-09-01T12:00:00Z'), ...overrides, -}); - -// TASK-112 replaces this reconstructed reader with AttentionItem source-write -// tests. Its task-prose / history assertions are intentionally retired. -describe.skip('retired ActivityService.getDecisionQueue reconstruction', () => { - beforeEach(() => { - jest.restoreAllMocks(); - Pod.find.mockReturnValue({ select: () => ({ lean: async () => [POD] }) }); - Task.find.mockReturnValue(taskChain([])); - User.findById.mockReturnValue(userChain({ username: 'Sam', activityQueue: { acknowledgedMentionIds: [] } })); - DecisionRequest.find.mockReturnValue(decisionChain([])); - pool.query.mockResolvedValue({ rows: [] }); - jest.spyOn(ActivityService, 'getPendingApprovals').mockResolvedValue([]); - }); - - test('renders agent-authored pending rows with their stored alternatives, not task prose', async () => { - DecisionRequest.find.mockReturnValue(decisionChain([decision()])); - Task.find.mockReturnValue(taskChain([ - { - taskId: 'TASK-024', podId: 'pod-1', status: 'pending', - title: 'DECIDE: legacy board wording', - updates: [{ text: 'OPTIONS: Parse me | Do not parse me', createdAt: new Date() }], - }, - ])); - - const result = await ActivityService.getDecisionQueue('u1'); - - expect(DecisionRequest.find).toHaveBeenCalledWith({ - podId: { $in: ['pod-1'] }, status: 'pending', messageId: { $exists: true, $ne: null }, - }); - expect(result.items).toEqual(expect.arrayContaining([ - expect.objectContaining({ - kind: 'decision', id: 'decision-1', messageId: '700', threadRootId: '695', - options: [ - { label: 'Canary', description: 'Roll out gradually.', recommended: true }, - { label: 'Fast lane', description: 'Ship after the green suite.' }, - ], - }), - ])); - expect(result.items.find((item) => item.taskId === 'TASK-024')).toBeUndefined(); - }); - - test('a blocked row is a standing decision, a human handoff is a press — both open-board items', async () => { - Task.find.mockReturnValue(taskChain([ - { - taskId: 'TASK-016', podId: 'pod-1', status: 'blocked', title: 'Fix release gate', - updates: [{ text: 'blocked on upstream', createdAt: new Date() }], - }, - { - taskId: 'TASK-059', podId: 'pod-1', status: 'claimed', title: 'Retention ledger', - updates: [{ text: 'held for the human merge press', createdAt: new Date() }], - }, - ])); - const result = await ActivityService.getDecisionQueue('u1'); - expect(result.items.map((item) => [item.taskId, item.kind])).toEqual([ - ['TASK-016', 'decision'], ['TASK-059', 'press'], - ]); - }); - - test('leaves existing Activity approvals on their established read path', async () => { - const getPendingApprovals = jest.spyOn(ActivityService, 'getPendingApprovals').mockResolvedValue([ - { _id: 'activity-approval-1', content: 'Deploy?', podId: 'pod-1', createdAt: new Date(), agentMetadata: { agentName: 'scout' } }, - ]); - const result = await ActivityService.getDecisionQueue('u1'); - expect(getPendingApprovals).toHaveBeenCalledWith('u1'); - expect(result.items).toEqual(expect.arrayContaining([ - expect.objectContaining({ kind: 'approval', id: 'activity-approval-1', title: 'scout requests approval' }), - ])); - }); - - test('presents a member with an approval they can decide', async () => { - Pod.find.mockReturnValue({ select: () => ({ lean: async () => [POD] }) }); - jest.spyOn(ActivityService, 'getPendingApprovals').mockResolvedValue([ - { _id: 'activity-approval-1', content: 'Deploy?', podId: 'pod-1', createdAt: new Date() }, - ]); - - const result = await ActivityService.getDecisionQueue('member-1'); - - expect(result.items).toEqual(expect.arrayContaining([ - expect.objectContaining({ kind: 'approval', id: 'activity-approval-1' }), - ])); - }); - - test('mentions are still thread-aware and cap at eight without hiding decision cards', async () => { - DecisionRequest.find.mockReturnValue(decisionChain(Array.from({ length: 4 }, (_, index) => decision({ - _id: `decision-${index}`, messageId: `${700 + index}`, createdAt: new Date(Date.now() - index * 1000), - })))); - pool.query.mockResolvedValue({ rows: Array.from({ length: 10 }, (_, index) => ({ - id: 100 + index, pod_id: 'pod-1', user_id: 'bot-1', content: '@Sam ping', - created_at: new Date(Date.now() - index * 1000), thread_root_id: null, author: 'vale', - })) }); - const result = await ActivityService.getDecisionQueue('u1'); - expect(result.items.filter((item) => item.kind === 'mention')).toHaveLength(8); - expect(result.items.filter((item) => item.kind === 'decision')).toHaveLength(4); - expect(result.count).toBe(14); - }); - - test('a failed DecisionRequest read degrades to the still-available board facts', async () => { - DecisionRequest.find.mockImplementation(() => { throw new Error('store down'); }); - Task.find.mockReturnValue(taskChain([{ - taskId: 'TASK-059', podId: 'pod-1', status: 'claimed', title: 'Release', - updates: [{ text: 'ready for the human press', createdAt: new Date() }], - }])); - const result = await ActivityService.getDecisionQueue('u1'); - expect(result.items).toEqual([expect.objectContaining({ taskId: 'TASK-059', kind: 'press' })]); - }); -}); diff --git a/backend/__tests__/unit/services/activityService.recap.test.js b/backend/__tests__/unit/services/activityService.recap.test.js index 670c9b08d..c3612a386 100644 --- a/backend/__tests__/unit/services/activityService.recap.test.js +++ b/backend/__tests__/unit/services/activityService.recap.test.js @@ -1,93 +1,50 @@ jest.mock('../../../models/Pod', () => ({ find: jest.fn(), findById: jest.fn() })); jest.mock('../../../models/Task', () => ({ find: jest.fn() })); +const mockGetOpenQueue = jest.fn(); +const mockAcknowledgeMention = jest.fn(); +const mockResolve = jest.fn(); +jest.mock('../../../services/attentionItemService', () => ({ + getOpenQueue: (...args) => mockGetOpenQueue(...args), + acknowledgeMention: (...args) => mockAcknowledgeMention(...args), + resolve: (...args) => mockResolve(...args), +})); + const Pod = require('../../../models/Pod'); const Task = require('../../../models/Task'); const Activity = require('../../../models/Activity'); -const User = require('../../../models/User'); const ActivityService = require('../../../services/activityService'); const ownerId = 'owner-1'; -const pod = { _id: 'pod-1', name: 'Activity source pod', type: 'team', createdBy: ownerId, members: [ownerId, 'member-1'] }; - -const podQuery = (pods) => ({ - select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue(pods) }), -}); +const pod = { + _id: 'pod-1', name: 'Activity source pod', type: 'team', createdBy: ownerId, members: [ownerId, 'member-1'], +}; +const podQuery = (pods) => ({ select: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue(pods) }) }); const taskQuery = (tasks) => ({ select: jest.fn().mockReturnValue({ - sort: jest.fn().mockReturnValue({ - limit: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue(tasks) }), - }), + sort: jest.fn().mockReturnValue({ limit: jest.fn().mockReturnValue({ lean: jest.fn().mockResolvedValue(tasks) }) }), }), }); -// TASK-112 removes the acknowledgement array and read-time queue projection. -// Recap itself remains covered by browser/UI contract; retire obsolete queue fixtures. -describe.skip('retired ActivityService.getRecap queue projection', () => { - let spy; +describe('ActivityService recap and legacy approval authorization', () => { + let feedSpy; let findByIdSpy; - let pendingApprovalsSpy; - let userFindByIdSpy; beforeEach(() => { jest.clearAllMocks(); Pod.find.mockReturnValue(podQuery([pod])); Pod.findById.mockReturnValue({ select: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(pod) })) }); - Task.find.mockReturnValue(taskQuery([{ - _id: 'board-1', - podId: 'pod-1', - taskId: 'TASK-068', - title: 'Activity tab', - status: 'claimed', - updatedAt: new Date(), - updates: [{ text: 'Implementation began.', author: 'sprint-impl', createdAt: new Date() }], - }])); - spy = jest.spyOn(ActivityService, 'getUserFeed').mockResolvedValue({ - activities: [{ - id: 'message-1', - type: 'message', - actor: { id: 'agent-1', name: 'sprint-impl', type: 'agent' }, - action: 'posted a message', - preview: 'Checks passed.', - timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, - flags: { isAgentAction: true, isMention: true }, - }], - }); - pendingApprovalsSpy = jest.spyOn(ActivityService, 'getPendingApprovals').mockResolvedValue([]); + Task.find.mockReturnValue(taskQuery([])); + mockGetOpenQueue.mockResolvedValue({ items: [], count: 0, composePodId: null }); + mockAcknowledgeMention.mockResolvedValue({ success: true }); + mockResolve.mockResolvedValue(undefined); + feedSpy = jest.spyOn(ActivityService, 'getUserFeed').mockResolvedValue({ activities: [] }); }); afterEach(() => { - spy?.mockRestore(); + feedSpy?.mockRestore(); findByIdSpy?.mockRestore(); - pendingApprovalsSpy?.mockRestore(); - userFindByIdSpy?.mockRestore(); - }); - - test('projects existing agent activity, direct mentions, and board updates without writing new events', async () => { - const result = await ActivityService.getRecap(ownerId, { window: 'today' }); - - expect(result.pods).toEqual([expect.objectContaining({ id: 'pod-1', name: pod.name })]); - expect(result.needsYou).toEqual([expect.objectContaining({ - kind: 'mention', podId: 'pod-1', title: 'sprint-impl mentioned you', - })]); - expect(result.agents).toEqual([expect.objectContaining({ - id: 'agent-1', name: 'sprint-impl', messageCount: 1, - updates: [expect.objectContaining({ content: 'Checks passed.' })], - })]); - expect(result.board).toEqual([expect.objectContaining({ - taskId: 'TASK-068', title: 'Activity tab', status: 'claimed', - lastUpdate: expect.objectContaining({ text: 'Implementation began.' }), - })]); - // filter: 'agents' is load-bearing — without it the 100-slot budget was - // consumed entirely by summary activities and no agent message ever - // reached the recap grouping (measured live: 30/30 summaries, #1307). - expect(spy).toHaveBeenCalledWith(ownerId, { limit: 100, filter: 'agents' }); - expect(Pod.find).toHaveBeenCalledWith(expect.objectContaining({ $or: expect.any(Array) })); - expect(Task.find).toHaveBeenCalledWith(expect.objectContaining({ - podId: { $in: ['pod-1'] }, updatedAt: { $gte: expect.any(Date) }, - })); }); test('rejects a requested pod that is outside the viewer membership', async () => { @@ -96,157 +53,26 @@ describe.skip('retired ActivityService.getRecap queue projection', () => { }); test('does not mistake the approval.status default on an ordinary message for a request', async () => { - // Build the stored document with the real schema. Its nested default is - // the production condition that caused every ordinary activity to be - // projected as an approval; a hand-written missing approval field would - // not reproduce it. - const storedMessage = new Activity({ - type: 'message', - action: 'posted a message', - content: 'An ordinary update.', - }); + const storedMessage = new Activity({ type: 'message', action: 'posted a message', content: 'An ordinary update.' }); expect(storedMessage.approval.status).toBe('pending'); - spy.mockResolvedValue({ + feedSpy.mockResolvedValue({ activities: [{ - id: 'message-with-defaulted-approval', - type: 'message', - actor: { id: 'human-1', name: 'A human', type: 'human' }, - action: 'posted a message', - preview: 'An ordinary update.', - timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, - approval: storedMessage.approval.toObject(), - flags: { isAgentAction: false, isMention: false }, + id: 'message-with-defaulted-approval', type: 'message', + actor: { id: 'human-1', name: 'A human', type: 'human' }, action: 'posted a message', + preview: 'An ordinary update.', timestamp: new Date(), pod: { id: 'pod-1', name: pod.name }, + approval: storedMessage.approval.toObject(), flags: { isAgentAction: false, isMention: false }, }], }); const result = await ActivityService.getRecap(ownerId, { window: 'today' }); expect(result.needsYou).toEqual([]); + expect(mockGetOpenQueue).toHaveBeenCalledWith(ownerId); }); - test('keeps an actual pending approval in the decision queue', async () => { - spy.mockResolvedValue({ - activities: [{ - id: 'approval-1', - type: 'approval_needed', - actor: { id: 'agent-1', name: 'release-agent', type: 'agent' }, - action: 'approval_needed', - preview: 'Approve access to Production.', - timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, - approval: { status: 'pending' }, - flags: { isAgentAction: true, isMention: false }, - }], - }); - - const result = await ActivityService.getRecap(ownerId, { window: 'today' }); - - expect(result.needsYou).toEqual([expect.objectContaining({ - id: 'approval-1', kind: 'approval', title: 'Approval requested', - })]); - }); - - test('puts a member’s approval in the needs-you queue', async () => { - Pod.find.mockReturnValue(podQuery([{ ...pod, createdBy: 'other-owner' }])); - spy.mockResolvedValue({ - activities: [{ - id: 'approval-1', type: 'approval_needed', actor: { id: 'agent-1', name: 'release-agent', type: 'agent' }, - action: 'approval_needed', preview: 'Approve access to Production.', timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, approval: { status: 'pending' }, - flags: { isAgentAction: true, isMention: false }, - }], - }); - - const result = await ActivityService.getRecap('member-1', { window: 'today' }); - - expect(result.needsYou).toEqual([expect.objectContaining({ id: 'approval-1', kind: 'approval' })]); - }); - - test('keeps a pending approval older than the recap window and absent from the sampled feed', async () => { - spy.mockResolvedValue({ activities: [] }); - pendingApprovalsSpy.mockResolvedValue([{ - _id: 'approval-before-feed-page', - type: 'approval_needed', - actor: { id: 'agent-1', name: 'release-agent', type: 'agent' }, - action: 'approval_needed', - content: 'Approve a decision that has waited longer than seven days.', - podId: 'pod-1', - approval: { status: 'pending' }, - createdAt: new Date('2026-08-01T00:00:00.000Z'), - }]); - - const result = await ActivityService.getRecap(ownerId, { window: '7d' }); - - expect(result.needsYou).toEqual([expect.objectContaining({ - id: 'approval-before-feed-page', kind: 'approval', podId: 'pod-1', - })]); - }); - - test('removes a mention only after its dedicated acknowledgement is recorded', async () => { - spy.mockResolvedValue({ - acknowledgedMentionIds: ['message-1'], - activities: [{ - id: 'message-1', - type: 'message', - actor: { id: 'agent-1', name: 'sprint-impl', type: 'agent' }, - action: 'posted a message', - preview: 'Please review this.', - timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, - flags: { isAgentAction: true, isMention: true }, - }], - }); - - const result = await ActivityService.getRecap(ownerId, { window: 'today' }); - - expect(result.needsYou).toEqual([]); - }); - - test('stores an acknowledgement separately from activity feed read-state', async () => { - const user = { - activityQueue: { acknowledgedMentionIds: [] }, - save: jest.fn().mockResolvedValue(), - }; - userFindByIdSpy = jest.spyOn(User, 'findById').mockReturnValue({ - select: jest.fn().mockResolvedValue(user), - }); - - const result = await ActivityService.acknowledgeMention(ownerId, 'message-1'); - - expect(result).toEqual({ success: true, acknowledgedMentionIds: ['message-1'] }); - expect(user.save).toHaveBeenCalledTimes(1); - expect(user).not.toHaveProperty('activityFeed'); - }); - - test('projects only an approval that the existing approve writer accepts', async () => { - const storedApproval = new Activity({ - type: 'approval_needed', - action: 'approval_needed', - content: 'Approve the release.', - }); - const approve = jest.fn().mockResolvedValue(); - storedApproval.approve = approve; - spy.mockResolvedValue({ - activities: [{ - id: String(storedApproval._id), - type: storedApproval.type, - actor: { id: 'agent-1', name: 'release-agent', type: 'agent' }, - action: storedApproval.action, - preview: storedApproval.content, - timestamp: new Date(), - pod: { id: 'pod-1', name: pod.name }, - approval: storedApproval.approval.toObject(), - flags: { isAgentAction: true, isMention: false }, - }], - }); - const recap = await ActivityService.getRecap(ownerId, { window: 'today' }); - findByIdSpy = jest.spyOn(Activity, 'findById').mockResolvedValue(storedApproval); - - const result = await ActivityService.approveActivity(recap.needsYou[0].id, ownerId, 'Approved in Activity'); - - expect(result).toEqual({ success: true, status: 'approved' }); - expect(approve).toHaveBeenCalledWith(ownerId, 'Approved in Activity'); + test('acknowledges a mention only through the recipient-owned attention record', async () => { + await expect(ActivityService.acknowledgeMention(ownerId, 'attention-1')).resolves.toEqual({ success: true }); + expect(mockAcknowledgeMention).toHaveBeenCalledWith(ownerId, 'attention-1'); }); test('allows a pod member to approve a legacy Activity approval', async () => { @@ -254,12 +80,11 @@ describe.skip('retired ActivityService.getRecap queue projection', () => { const approve = jest.fn().mockResolvedValue(); storedApproval.approve = approve; findByIdSpy = jest.spyOn(Activity, 'findById').mockResolvedValue(storedApproval); - Pod.findById.mockReturnValue({ select: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(pod) })) }); - const result = await ActivityService.approveActivity(String(storedApproval._id), 'member-1', 'Approved'); - - expect(result).toEqual({ success: true, status: 'approved' }); + await expect(ActivityService.approveActivity(String(storedApproval._id), 'member-1', 'Approved')) + .resolves.toEqual({ success: true, status: 'approved' }); expect(approve).toHaveBeenCalledWith('member-1', 'Approved'); + expect(mockResolve).toHaveBeenCalledWith('approval', storedApproval._id); }); test('fails closed when a non-member attempts a legacy Activity approval', async () => { @@ -267,11 +92,9 @@ describe.skip('retired ActivityService.getRecap queue projection', () => { const approve = jest.fn().mockResolvedValue(); storedApproval.approve = approve; findByIdSpy = jest.spyOn(Activity, 'findById').mockResolvedValue(storedApproval); - Pod.findById.mockReturnValue({ select: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(pod) })) }); - const result = await ActivityService.approveActivity(String(storedApproval._id), 'non-member', 'Nope'); - - expect(result).toEqual({ success: false, status: 403, error: 'Only pod members can decide this' }); + await expect(ActivityService.approveActivity(String(storedApproval._id), 'non-member', 'Nope')) + .resolves.toEqual({ success: false, status: 403, error: 'Only pod members can decide this' }); expect(approve).not.toHaveBeenCalled(); }); @@ -280,11 +103,9 @@ describe.skip('retired ActivityService.getRecap queue projection', () => { const reject = jest.fn().mockResolvedValue(); storedApproval.reject = reject; findByIdSpy = jest.spyOn(Activity, 'findById').mockResolvedValue(storedApproval); - Pod.findById.mockReturnValue({ select: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(pod) })) }); - - const result = await ActivityService.rejectActivity(String(storedApproval._id), 'non-member', 'Nope'); - expect(result).toEqual({ success: false, status: 403, error: 'Only pod members can decide this' }); + await expect(ActivityService.rejectActivity(String(storedApproval._id), 'non-member', 'Nope')) + .resolves.toEqual({ success: false, status: 403, error: 'Only pod members can decide this' }); expect(reject).not.toHaveBeenCalled(); }); }); diff --git a/backend/__tests__/unit/services/attentionItemService.test.js b/backend/__tests__/unit/services/attentionItemService.test.js index 3bd103a47..df3cb6cd8 100644 --- a/backend/__tests__/unit/services/attentionItemService.test.js +++ b/backend/__tests__/unit/services/attentionItemService.test.js @@ -34,6 +34,61 @@ describe('attentionItemService', () => { expect(mockUpdateOne.mock.calls[0][1].$setOnInsert).toMatchObject({ kind: 'mention', title: 'Ada mentioned you', messageId: '42' }); }); + it('does not read pod membership for a message with no mention marker', async () => { + await AttentionItemService.recordMentionedUsers({ + id: 42, podId: 'pod-1', userId: 'owner', username: 'Ada', content: 'ordinary status update', + }); + + expect(mockPodFindById).not.toHaveBeenCalled(); + expect(mockUserFind).not.toHaveBeenCalled(); + expect(mockUpdateOne).not.toHaveBeenCalled(); + }); + + it('does not re-materialize a legacy-acknowledged mention during backfill', async () => { + mockPodFindById.mockReturnValue(chain({ _id: 'pod-1', name: 'Ship room', createdBy: 'owner', members: [{ userId: 'sam' }] })); + mockUserFind.mockReturnValue(chain([ + { _id: 'owner', username: 'owner', isBot: false }, + { _id: 'sam', username: 'Sam', isBot: false }, + ])); + + await AttentionItemService.recordMentionedUsers( + { id: 42, podId: 'pod-1', userId: 'owner', username: 'Ada', content: '@sam please review' }, + { isAlreadyAcknowledged: (recipient, id) => recipient === 'sam' && id === 'msg_42' }, + ); + + expect(mockUpdateOne).not.toHaveBeenCalled(); + }); + + it('does not resurrect a resolved attention row when its source write is retried', async () => { + mockPodFindById.mockReturnValue(chain({ _id: 'pod-1', name: 'Ship room', createdBy: 'owner', members: [{ userId: 'sam' }] })); + mockUserFind.mockReturnValue(chain([ + { _id: 'owner', username: 'owner', isBot: false }, + { _id: 'sam', username: 'Sam', isBot: false }, + ])); + const rows = []; + mockUpdateOne.mockImplementation(async (filter, update) => { + const row = rows.find((candidate) => ( + candidate.recipientUserId === filter.recipientUserId + && candidate.source.type === filter['source.type'] + && candidate.source.id === filter['source.id'] + )); + if (row) { + if (update.$set) Object.assign(row, update.$set); + return { matchedCount: 1, modifiedCount: 1 }; + } + rows.push({ ...update.$setOnInsert }); + return { upsertedCount: 1 }; + }); + const message = { id: 42, podId: 'pod-1', userId: 'owner', username: 'Ada', content: '@sam please review' }; + + await AttentionItemService.recordMentionedUsers(message); + rows[0].status = 'resolved'; + await AttentionItemService.recordMentionedUsers(message); + + expect(rows).toHaveLength(1); + expect(rows[0]).toMatchObject({ status: 'resolved', source: { type: 'message', id: '42' } }); + }); + it('returns only rows whose recipient is still a member and resolves by recipient-owned id', async () => { mockFind.mockReturnValue({ sort: () => ({ limit: () => ({ lean: async () => [ { _id: 'attention-1', recipientUserId: '507f191e810c19729de860ea', podId: 'pod-1', kind: 'mention', source: { type: 'message', id: '41' }, title: 'Mention', createdAt: new Date() }, @@ -57,4 +112,39 @@ describe('attentionItemService', () => { mockUpdateMany.mockRejectedValueOnce(new Error('mongo unavailable')); await expect(AttentionItemService.resolve('approval', 'a-1')).resolves.toBeUndefined(); }); + + it('materializes a blocked board row once for each current human recipient', async () => { + mockPodFindById.mockReturnValue(chain({ _id: 'pod-1', name: 'Ship room', createdBy: 'owner', members: [{ userId: 'sam' }] })); + mockUserFind.mockReturnValue(chain([ + { _id: 'owner', username: 'owner', isBot: false }, + { _id: 'sam', username: 'Sam', isBot: false }, + ])); + + await AttentionItemService.recordTaskAttention({ + _id: 'task-1', podId: 'pod-1', taskId: 'TASK-1', status: 'blocked', + title: 'Choose a deploy shape', updates: [{ _id: 'update-1', text: 'Blocked on an upstream choice.' }], + }, { includeBlocked: true }); + + expect(mockUpdateOne).toHaveBeenCalledTimes(2); + expect(mockUpdateOne).toHaveBeenCalledWith( + expect.objectContaining({ 'source.type': 'task', 'source.id': 'task-1:update-1' }), + expect.objectContaining({ $setOnInsert: expect.objectContaining({ kind: 'decision', title: 'Choose a deploy shape' }) }), + { upsert: true }, + ); + }); + + it('resolves every outstanding fact for a task once the task no longer needs a human', async () => { + await AttentionItemService.resolveTaskAttention({ _id: 'task.1' }); + + expect(mockUpdateMany).toHaveBeenCalledWith( + expect.objectContaining({ + 'source.type': 'task', + 'source.id': expect.objectContaining({ $regex: expect.any(RegExp) }), + }), + expect.any(Object), + ); + const sourcePattern = mockUpdateMany.mock.calls[0][0]['source.id'].$regex; + expect(sourcePattern.test('task.1:update-1')).toBe(true); + expect(sourcePattern.test('taskx1:update-1')).toBe(false); + }); }); diff --git a/backend/models/AttentionItem.ts b/backend/models/AttentionItem.ts index 25aadfd42..db23941a3 100644 --- a/backend/models/AttentionItem.ts +++ b/backend/models/AttentionItem.ts @@ -1,7 +1,7 @@ import mongoose, { Document, Model, Schema, Types } from 'mongoose'; export type AttentionKind = 'mention' | 'approval' | 'decision'; -export type AttentionSourceType = 'message' | 'approval' | 'decision_request'; +export type AttentionSourceType = 'message' | 'approval' | 'decision_request' | 'task'; export interface IAttentionItem extends Document { recipientUserId: Types.ObjectId; @@ -31,7 +31,7 @@ const attentionItemSchema = new Schema({ podId: { type: Schema.Types.ObjectId, ref: 'Pod', required: true }, kind: { type: String, enum: ['mention', 'approval', 'decision'], required: true }, source: { - type: { type: String, enum: ['message', 'approval', 'decision_request'], required: true }, + type: { type: String, enum: ['message', 'approval', 'decision_request', 'task'], required: true }, id: { type: String, required: true }, }, title: { type: String, required: true }, diff --git a/backend/package.json b/backend/package.json index 723232b22..caa0e3494 100644 --- a/backend/package.json +++ b/backend/package.json @@ -21,6 +21,7 @@ "discord:register": "node scripts/register-discord-commands.js", "discord:list": "node scripts/register-discord-commands.js list", "bootstrap:clawd-bot": "node scripts/bootstrap-clawd-bot.js", + "backfill:attention-items": "ts-node scripts/backfill-attention-items.ts", "tsc:check": "tsc --noEmit -p tsconfig.typescheck.json", "tsc:check:all": "tsc --noEmit" }, diff --git a/backend/routes/tasksApi.ts b/backend/routes/tasksApi.ts index 276f6833d..e61a343b4 100644 --- a/backend/routes/tasksApi.ts +++ b/backend/routes/tasksApi.ts @@ -17,6 +17,8 @@ const Task = require('../models/Task'); const User = require('../models/User'); // eslint-disable-next-line global-require const { emitTaskUpdated, notifyPodAgents } = require('../services/taskEventService'); +// eslint-disable-next-line global-require +const { recordTaskAttention, resolveTaskAttention } = require('../services/attentionItemService'); /** * Who made this board change, for the agent fan-out. @@ -686,6 +688,7 @@ router.post('/:podId/:taskId/updates', taskWriteRateLimit(60), auth, async (req: task = await Task.findOneAndUpdate({ podId: podFilter, taskId }, { $push: note }, { new: true }); } if (!task) return res.status(404).json({ error: 'Task not found' }); + await recordTaskAttention(task); emitTaskUpdated(podId, task, 'updated'); notifyAgents(req, podId, task, 'updated'); // Reported explicitly. A note that did not renew is a legitimate outcome @@ -763,6 +766,8 @@ router.patch('/:podId/:taskId', taskWriteRateLimit(60), auth, async (req: AuthRe if (changeParts.length > 0) update.$push = { updates: { text: `${author} updated: ${changeParts.join(', ')}`, author, authorId: userId?.toString() || null, createdAt: new Date() } }; const task = await Task.findOneAndUpdate({ podId: mongoose.Types.ObjectId.createFromHexString(podId || ''), taskId }, update, { new: true }); if (!task) return res.status(404).json({ error: 'Task not found' }); + if (fieldUpdates.status === 'blocked') await recordTaskAttention(task, { includeBlocked: true }); + if (fieldUpdates.status !== undefined && fieldUpdates.status !== 'blocked') await resolveTaskAttention(task); emitTaskUpdated(podId, task, 'updated'); notifyAgents(req, podId, task, 'updated'); return res.json({ task }); diff --git a/backend/scripts/backfill-attention-items.ts b/backend/scripts/backfill-attention-items.ts new file mode 100644 index 000000000..9e3a8af27 --- /dev/null +++ b/backend/scripts/backfill-attention-items.ts @@ -0,0 +1,108 @@ +/** + * Materialize AttentionItem rows for facts written before TASK-112. + * + * Run once per deployed environment after the AttentionItem model and source + * writers are live: + * + * npm run backfill:attention-items -- --apply + * + * The source/recipient unique index makes a retry safe. The script directly + * queries source stores; after it succeeds, Activity's old reconstructed + * readers must stay deleted. + */ +/* eslint-disable no-console */ +const mongoose = require('mongoose'); +const Activity = require('../models/Activity'); +const DecisionRequest = require('../models/DecisionRequest'); +const Message = require('../models/Message'); +const Task = require('../models/Task'); +const User = require('../models/User'); +const { + recordApproval, recordDecision, recordMentionedUsers, recordTaskAttention, TASK_HANDOFF_RE, +} = require('../services/attentionItemService'); + +const APPLY = process.argv.includes('--apply'); +const MENTION_CUTOFF_DAYS = 14; + +type PgPool = { + query: (sql: string, values?: unknown[]) => Promise<{ rows: Array> }>; + end: () => Promise; +}; + +const legacyAcknowledgements = async (): Promise>> => { + // This retired field is deliberately absent from the current User schema. + // Use the raw collection so Mongoose cannot omit it from the projection. + const users = await User.collection.find( + { 'activityQueue.acknowledgedMentionIds.0': { $exists: true } }, + { projection: { _id: 1, activityQueue: 1 } }, + ).toArray(); + return new Map(users.map((user: any) => [ + String(user._id), + new Set((user.activityQueue?.acknowledgedMentionIds || []).map(String)), + ])); +}; + +export const main = async (): Promise => { + if (!process.env.MONGO_URI) throw new Error('MONGO_URI is required'); + // eslint-disable-next-line global-require + const { pool } = require('../config/db-pg') as { pool: PgPool | null }; + + await mongoose.connect(process.env.MONGO_URI); + try { + const cutoff = new Date(Date.now() - MENTION_CUTOFF_DAYS * 24 * 60 * 60 * 1000); + const [approvals, decisions, tasks, acknowledgements, mongoMessages] = await Promise.all([ + Activity.find({ type: 'approval_needed', 'approval.status': 'pending', deleted: { $ne: true } }).lean(), + DecisionRequest.find({ status: 'pending' }).lean(), + Task.find({ + $or: [{ status: 'blocked' }, { 'updates.text': { $regex: TASK_HANDOFF_RE } }], + }).select('podId taskId title status notes updates updatedAt').lean(), + legacyAcknowledgements(), + Message.find({ createdAt: { $gte: cutoff } }).lean(), + ]); + + const pgMessages = pool + ? (await pool.query( + 'SELECT m.id, m.pod_id, m.user_id, m.content, m.created_at, m.thread_root_id, u.username ' + + 'FROM messages m LEFT JOIN users u ON u._id = m.user_id ' + + "WHERE m.created_at >= now() - ($1::int * interval '1 day') ORDER BY m.created_at ASC", + [MENTION_CUTOFF_DAYS], + )).rows + : []; + const messages = [...pgMessages, ...mongoMessages]; + + console.log(JSON.stringify({ + approvals: approvals.length, + decisions: decisions.length, + tasks: tasks.length, + recentPgMessages: pgMessages.length, + recentMongoFallbackMessages: mongoMessages.length, + apply: APPLY, + })); + if (!APPLY) { + console.log('DRY RUN — no AttentionItems written. Re-run with --apply after deploy.'); + return; + } + + for (const approval of approvals) await recordApproval(approval); + for (const decision of decisions) await recordDecision(decision); + for (const task of tasks) await recordTaskAttention(task, { includeBlocked: true }); + for (const message of messages) { + await recordMentionedUsers(message, { + isAlreadyAcknowledged: (recipientUserId: unknown, legacyMentionId: string) => ( + acknowledgements.get(String(recipientUserId))?.has(legacyMentionId) || false + ), + }); + } + console.log('APPLIED — source facts materialized. A retry is safe because recipient/source is unique.'); + } finally { + await mongoose.disconnect(); + if (pool) await pool.end(); + } +}; + +if (require.main === module) { + main().catch((error) => { + console.error('attention-item backfill failed:', error); + process.exit(1); + }); +} diff --git a/backend/services/activityService.ts b/backend/services/activityService.ts index ddbaa61bd..cf84b2336 100644 --- a/backend/services/activityService.ts +++ b/backend/services/activityService.ts @@ -10,11 +10,6 @@ const Summary = require('../models/Summary'); const Post = require('../models/Post'); // eslint-disable-next-line global-require const Task = require('../models/Task'); -// Kept while the dead former implementation below is mechanically removed; -// getDecisionQueue returns before this is ever read. -// eslint-disable-next-line global-require -const DecisionRequest = require('../models/DecisionRequest'); -// eslint-disable-next-line global-require // eslint-disable-next-line global-require const isPodMember = require('../utils/isPodMember'); @@ -227,93 +222,16 @@ class ActivityService { }) .sort((a, b) => new Date(b.lastActiveAt || 0).getTime() - new Date(a.lastActiveAt || 0).getTime()); - // Approvals are a decision queue, not an activity sample: query the - // existing authoritative pending-approval reader separately so a busy - // pod cannot push an older decision behind getUserFeed's display page. - // Mentions remain a bounded recent interrupt list and are removed only by - // their explicit acknowledgement, never by feed read-state. - const pendingApprovals = await ActivityService.getPendingApprovals(userId) as Array<{ - _id?: unknown; - id?: unknown; - type?: string; - actor?: ActorInfo; - action?: string; - content?: string; - podId?: unknown; - approval?: unknown; - agentMetadata?: { agentName?: string }; - createdAt?: Date | string; - updatedAt?: Date | string; - }>; - const approvalItems: ActivityItem[] = pendingApprovals - .filter((approval) => !requestedPodId || scopedPodIds.has(String(approval.podId))) - .map((approval) => { - const podId = approval.podId ? String(approval.podId) : ''; - const pod = scopedPods.find((candidate) => String(candidate._id) === podId); - return { - id: String(approval._id || approval.id || ''), - type: approval.type || 'approval_needed', - actor: approval.actor || { - name: approval.agentMetadata?.agentName || 'An agent', type: 'agent', - }, - action: approval.action || 'approval_needed', - content: approval.content, - timestamp: approval.createdAt || approval.updatedAt || null, - pod: pod ? { id: String(pod._id), name: pod.name } : null, - approval: approval.approval, - reactions: { likes: 0, liked: false }, - replyCount: 0, - replies: [], - }; - }) - .filter((approval) => Boolean(approval.id)); - - const isPendingApproval = (activity: ActivityItem): boolean => { - const approval = activity.approval as { status?: string } | undefined; - // Mongoose materializes approval.status = 'pending' for every Activity - // document. It is meaningful only on the one activity type that carries - // an approval request; otherwise every message becomes a human action. - return activity.type === 'approval_needed' && approval?.status === 'pending'; - }; - const queueCandidates = new Map(); - [...activities, ...approvalItems].forEach((activity) => { - if (!queueCandidates.has(activity.id)) queueCandidates.set(activity.id, activity); - }); - const newestFirst = (left: ActivityItem, right: ActivityItem) => ( - new Date(right.timestamp || 0).getTime() - new Date(left.timestamp || 0).getTime() - ); - const approvalQueue = Array.from(queueCandidates.values()) - .filter(isPendingApproval) - .sort(newestFirst); - const mentionQueue = Array.from(queueCandidates.values()) - .filter((activity) => activity.flags?.isMention) - .sort(newestFirst) - .slice(0, Math.max(0, 12 - approvalQueue.length)); - let needsYou = [...approvalQueue, ...mentionQueue] - .sort(newestFirst) - .map((activity) => { - const isApproval = isPendingApproval(activity); - return { - id: activity.id, - kind: isApproval ? 'approval' : 'mention', - title: isApproval ? 'Approval requested' : `${activity.actor?.name || 'Someone'} mentioned you`, - detail: String(activity.preview || activity.content || '').replace(/\s+/g, ' ').trim().slice(0, 180), - podId: activity.pod?.id || null, - podName: activity.pod?.name || 'Direct activity', - timestamp: activity.timestamp, - }; - }); - - // Recipient-owned AttentionItem rows are the only queue source. The - // temporary calculation above remains only for the normal activity feed; - // it cannot become a queue fallback or revive historical heuristics. + // Recipient-owned AttentionItem rows are the only needs-you source. // eslint-disable-next-line global-require const AttentionItemService = require('./attentionItemService'); const attention = await AttentionItemService.getOpenQueue(userId); - needsYou = attention.items.map((item: any) => ({ - ...item, - timestamp: item.createdAt || null, - })); + const needsYou = attention.items + .filter((item: any) => !requestedPodId || item.podId === requestedPodId) + .map((item: any) => ({ + ...item, + timestamp: item.createdAt || null, + })); let board: Array> = []; if (scopedPods.length > 0) { @@ -450,274 +368,9 @@ class ActivityService { } } - /** - * The decision queue: everything concretely waiting on THIS human, from - * stored facts: existing approval requests, unacknowledged direct mentions, - * agent-authored DecisionRequests, and explicit board press handoffs. - * A board title is never parsed into a decision surface. - */ - /** - * Every message in the user's pods that @mentions them, threads included, - * minus their own and minus acknowledged ones. Ids are `msg_` — - * the same shape the feed emits, so acknowledgedMentionIds keeps working. - */ - static async getMentionsForUser(userId: unknown, podIds: unknown[]): Promise> { - const user = await (User as { - findById(id: unknown): { select(f: string): { lean(): Promise<{ username?: string; activityQueue?: { acknowledgedMentionIds?: unknown[] } } | null> } }; - }).findById(userId).select('username activityQueue.acknowledgedMentionIds').lean(); - const username = String(user?.username || '').trim(); - if (!username || !podIds.length) return []; - const acked = new Set(((user?.activityQueue?.acknowledgedMentionIds as unknown[]) || []).map(String)); - const { pool } = require('../config/db-pg') as { pool: { query(q: string, p: unknown[]): Promise<{ rows: Array> }> } }; - // ILIKE on the handle; the trailing boundary keeps @sam from matching - // @samantha. ::text[] is load-bearing (the rank query lesson, #1307). - const { rows } = await pool.query( - `SELECT m.id, m.pod_id, m.user_id, m.content, m.created_at, m.thread_root_id, u.username AS author - FROM messages m - LEFT JOIN users u ON u._id = m.user_id - WHERE m.pod_id = ANY($1::text[]) - AND m.user_id <> $2 - AND m.created_at > now() - interval '7 days' - AND m.content ~* $3 - ORDER BY m.created_at DESC - LIMIT 40`, - [podIds.map(String), String(userId), `@${username.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')}(?![A-Za-z0-9_])`], - ); - return rows - .map((r) => ({ - id: `msg_${r.id}`, - messageId: Number(r.id), - threadRootId: Number(r.thread_root_id || r.id), - podId: String(r.pod_id), - authorName: String(r.author || 'Someone'), - content: String(r.content || ''), - createdAt: r.created_at as Date, - })) - .filter((m) => !acked.has(m.id)); - } - static async getDecisionQueue(userId: unknown): Promise> { - // This guarded former reader is retained for one patch only while this - // large method is mechanically removed. It is never a production fallback. - if (process.env.ATTENTION_ITEM_LEGACY_READER_FOR_TESTS === '1') { - const pods: PodDoc[] = await Pod.find({ - $or: [ - { createdBy: userId }, - { 'members.userId': userId }, - { members: userId }, - ], - }).select('_id name type').lean(); - const podIds = pods.map((p) => p._id); - const podName = new Map(pods.map((p) => [String(p._id), p.name as string])); - - type QueueItem = { - kind: 'approval' | 'mention' | 'decision' | 'press'; - id: string; - title: string; - detail?: string; - podId: string | null; - podName?: string; - taskId?: string; - options?: Array<{ label: string; description?: string; recommended?: boolean }>; - messageId?: string; - threadRootId?: string; - createdAt: Date | string | null; - }; - const items: QueueItem[] = []; - let composePodId: string | null = null; - - // 1. Existing binary approvals stay on their established Activity path. - // DecisionRequest is additive; it does not widen or replace the - // owner-scoped ApprovalAction flow. - try { - const approvals = await ActivityService.getPendingApprovals(userId) as Array<{ - _id: { toString(): string }; content?: string; podId?: unknown; createdAt?: Date; - agentMetadata?: { agentName?: string }; - }>; - for (const a of approvals) { - items.push({ - kind: 'approval', - // RAW Activity id — the existing Activity actions key on it. - id: String(a._id), - title: a.agentMetadata?.agentName ? `${a.agentMetadata?.agentName} requests approval` : 'Approval requested', - detail: String(a.content || '').slice(0, 160), - podId: a.podId ? String(a.podId) : null, - podName: a.podId ? podName.get(String(a.podId)) : undefined, - createdAt: a.createdAt || null, - }); - } - } catch (err) { - console.warn('[decision-queue] approvals read failed:', (err as Error).message); - } - - // 2. Unacknowledged direct mentions — a DEDICATED query, not the feed. - // Sam, 2026-09-01: "pods mentioning me… but activity tab is not showing - // properly." Measured: 15 @Sam mentions in 36h, 14 inside THREADS, and - // the feed path saw zero — it samples ~20 recent rows per pod through - // the generic aggregator, so thread traffic (where agents now do most - // of their talking) never reaches the mention flag. The mention query - // asks the store the actual question, across every pod, threads - // included, for the last 7 days. - try { - const mentions = await ActivityService.getMentionsForUser(userId, podIds); - for (const m of mentions) { - // The query is newest-first, so its first result is the closest - // factual default for "tell them what is on my mind". - if (!composePodId) composePodId = m.podId; - items.push({ - kind: 'mention', - id: m.id, - title: `${m.authorName} mentioned you`, - detail: m.content.slice(0, 220), - podId: m.podId, - podName: podName.get(m.podId), - createdAt: m.createdAt, - // Reply-in-place needs the thread root (or the message itself, as - // the root of a new thread) and the message to address. - ...({ threadRootId: m.threadRootId, messageId: m.messageId } as Record), - }); - } - } catch (err) { - console.warn('[decision-queue] mentions read failed:', (err as Error).message); - } - - // 3. Agent-authored DecisionRequest rows. The asking runtime supplied - // the title, question, and alternatives at the fork; this reader neither - // parses task prose nor reconstructs context from a board title. - try { - const decisions = await DecisionRequest.find({ - podId: { $in: podIds }, - status: 'pending', - messageId: { $exists: true, $ne: null }, - }).sort({ createdAt: -1 }).limit(50).lean() as Array<{ - _id: { toString(): string }; podId?: unknown; title?: string; question?: string; - context?: string; options?: Array<{ label?: string; description?: string; recommended?: boolean }>; - messageId?: string; threadRootId?: string; createdAt?: Date; - }>; - for (const decision of decisions) { - const options = (decision.options || []) - .filter((option) => option && typeof option.label === 'string') - .sort((a, b) => Number(Boolean(b.recommended)) - Number(Boolean(a.recommended))) - .map((option) => ({ - label: String(option.label), - ...(option.description ? { description: String(option.description) } : {}), - ...(option.recommended ? { recommended: true } : {}), - })); - items.push({ - kind: 'decision', - id: String(decision._id), - title: String(decision.title || 'Decision requested'), - detail: String(decision.question || decision.context || '').slice(0, 1000), - podId: decision.podId ? String(decision.podId) : null, - podName: decision.podId ? podName.get(String(decision.podId)) : undefined, - options, - messageId: decision.messageId, - threadRootId: decision.threadRootId || decision.messageId, - createdAt: decision.createdAt || null, - }); - } - } catch (err) { - console.warn('[decision-queue] decision request read failed:', (err as Error).message); - } - - // 4. Board handoffs still need an open-board affordance, but task prose - // is never an option source. DECIDE rows remain ordinary board rows. - try { - const HANDOFF_RE = /human\s+(merge\s+)?press|ready for (the\s+)?(human|sam)|sam'?s?\s+(ruling|call|decision)|awaiting\s+(sam|human)/i; - // Freshness cutoff (Sam, 2026-09-01: "some of those activities seem - // stale and non new attention routes"). Measured that day: every queue - // item was 10–47 days old, half of them zombie rows from dead pods — - // "Fix CI failures on PR #N" from July does not need attention, it - // needs an archive. A board row qualifies as WAITING ON YOU only if - // someone touched it recently; the row itself stays on the board - // either way, so nothing is lost — only the attention claim expires. - const STALE_CUTOFF = new Date(Date.now() - 14 * 24 * 60 * 60 * 1000); - const tasks = await (Task as { - find(q: unknown): { sort(s: unknown): { limit(n: number): { lean(): Promise>> } } }; - }).find({ - podId: { $in: podIds }, - status: { $in: ['pending', 'claimed', 'blocked'] }, - updatedAt: { $gte: STALE_CUTOFF }, - }).sort({ updatedAt: -1 }).limit(120).lean(); - for (const t of tasks) { - const title = String(t.title || ''); - const updates = (t.updates as Array<{ text?: string; createdAt?: Date }> | undefined) || []; - const last = updates[updates.length - 1]; - const isBlocked = t.status === 'blocked'; - const isHandoff = !!(last && HANDOFF_RE.test(String(last.text || ''))); - if (!isBlocked && !isHandoff) continue; - items.push({ - // A blocked row is a standing DECISION (no options — the fork owner - // should convert it to a DecisionRequest); only an explicit human - // handoff in the latest update is a PRESS. #1470 collapsed both to - // 'press', which put "ADR-024 D3: batch the poller tick" under - // "Ready for your press" — a mislabel Sam noticed the same morning. - kind: isHandoff ? 'press' : 'decision', - id: `task_${t.taskId}`, - title, - detail: last ? String(last.text || '').slice(0, 160) : undefined, - podId: String(t.podId), - podName: podName.get(String(t.podId)), - taskId: String(t.taskId || ''), - createdAt: (last?.createdAt as Date) || (t.updatedAt as Date) || null, - }); - } - } catch (err) { - console.warn('[decision-queue] board read failed:', (err as Error).message); - } - - // Attention order, then recency, then a hard cap. The option card is the - // direct decision surface; existing approvals and board presses retain - // their prior precedence and mentions stay bounded below action rows. - const KIND_PRIORITY: Record = { - approval: 0, decision: 1, press: 2, mention: 3, - }; - items.sort((a, b) => { - const kindDelta = (KIND_PRIORITY[a.kind] ?? 9) - (KIND_PRIORITY[b.kind] ?? 9); - if (kindDelta !== 0) return kindDelta; - return new Date(b.createdAt || 0).getTime() - new Date(a.createdAt || 0).getTime(); - }); - // Per-kind soft cap. Live after #1464: 40+ thread mentions filled all - // 12 slots and direct decision cards vanished below the cap — the - // opposite of "decision pending on me" being visible. Mentions take at - // most 8 of the 12; whatever is left goes to the other kinds - // in attention order. `count` still reports the full total. - const MENTION_SLOTS = 8; - const picked: QueueItem[] = []; - let mentionsPicked = 0; - const deferredMentions: QueueItem[] = []; - for (const it of items) { - if (picked.length >= 12) break; - if (it.kind === 'mention') { - if (mentionsPicked < MENTION_SLOTS) { picked.push(it); mentionsPicked += 1; } else deferredMentions.push(it); - } else { - picked.push(it); - } - } - for (const it of deferredMentions) { - if (picked.length >= 12) break; - picked.push(it); - } - // Restore attention order within the final page, recency inside each - // kind — the loop above kept the sorted order for everything it picked, - // so a stable re-sort is enough. - picked.sort((a, b) => { - const kindDelta = (KIND_PRIORITY[a.kind] ?? 9) - (KIND_PRIORITY[b.kind] ?? 9); - if (kindDelta !== 0) return kindDelta; - return new Date(b.createdAt || 0).getTime() - new Date(a.createdAt || 0).getTime(); - }); - return { - items: picked, - count: items.length, - // No mention yet is not an error; the page falls back to its first pod. - composePodId, - }; - } - // The recipient-scoped read is indexed and authorization-safe. Do not - // reconstruct this from messages, Task prose, or historical acknowledgements. + // Recipient-scoped rows are materialized at each source write, then + // membership-checked by this indexed reader. // eslint-disable-next-line global-require const AttentionItemService = require('./attentionItemService'); return AttentionItemService.getOpenQueue(userId); diff --git a/backend/services/attentionItemService.ts b/backend/services/attentionItemService.ts index 4b7f3e14d..dbca1b273 100644 --- a/backend/services/attentionItemService.ts +++ b/backend/services/attentionItemService.ts @@ -1,6 +1,6 @@ // Recipient-owned attention is written beside its authoritative source. It is // intentionally not rebuilt from message history or task prose at read time. -// There is no historical backfill: facts start materializing when this ships. +// A one-time direct-source backfill materializes facts that predate adoption. // eslint-disable-next-line global-require const AttentionItem = require('../models/AttentionItem'); // eslint-disable-next-line global-require @@ -8,8 +8,12 @@ const Pod = require('../models/Pod'); // eslint-disable-next-line global-require const User = require('../models/User'); -type SourceType = 'message' | 'approval' | 'decision_request'; +type SourceType = 'message' | 'approval' | 'decision_request' | 'task'; type Kind = 'mention' | 'approval' | 'decision'; +type MentionOptions = { + isAlreadyAcknowledged?: (recipientUserId: unknown, legacyMentionId: string) => boolean; +}; +type TaskAttentionOptions = { includeBlocked?: boolean }; const compact = (value: unknown, max = 220): string => String(value || '').replace(/\s+/g, ' ').trim().slice(0, max); const sourceKey = (type: SourceType, id: unknown): string => String(id || '').trim(); @@ -56,17 +60,22 @@ const recordForRecipients = async ( ))); }; -export const recordMentionedUsers = async (message: any): Promise => { +export const recordMentionedUsers = async (message: any, options: MentionOptions = {}): Promise => { try { const podId = message?.podId || message?.pod_id; const messageId = message?._id || message?.id; if (!podId || messageId === undefined || messageId === null) return; const authorId = message?.userId?._id || message?.userId || message?.user_id; const content = String(message?.content || message?.text || ''); + // Message delivery invokes this writer for every post. Avoid two Mongo + // reads on the common no-mention path; every valid @mention contains this + // sentinel, so the guard cannot hide a recipient. + if (!content.includes('@')) return; const members = await currentHumanMembers(podId); const recipients = members.filter((member) => { const handle = String(member.username || '').trim(); if (!handle || String(member._id) === String(authorId)) return false; + if (options.isAlreadyAcknowledged?.(member._id, 'msg_' + String(messageId))) return false; return new RegExp(`(^|[^A-Za-z0-9_-])@${handle.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')}(?![A-Za-z0-9_-])`, 'i').test(content); }); if (!recipients.length) return; @@ -121,6 +130,46 @@ export const recordDecision = async (decision: any): Promise => { } }; +// A task may be ordinary board history, or it may state a concrete blocked / +// handoff fact. Capture only the latter at the write boundary; no reader +// should parse arbitrary historic task prose into a queue card. +export const TASK_HANDOFF_RE = /human\s+(merge\s+)?press|ready for (the\s+)?(human|sam)|sam'?s?\s+(ruling|call|decision)|awaiting\s+(sam|human)/i; + +export const recordTaskAttention = async (task: any, options: TaskAttentionOptions = {}): Promise => { + try { + const last = Array.isArray(task?.updates) ? task.updates[task.updates.length - 1] : null; + const blocked = task?.status === 'blocked'; + const handoff = Boolean(last && TASK_HANDOFF_RE.test(String(last.text || ''))); + if (!task?.podId || (!handoff && !(options.includeBlocked && blocked))) return; + const recipients = await currentHumanMembers(task.podId); + const pod = await Pod.findById(task.podId).select('name').lean(); + const taskKey = String(task._id || task.taskId); + const sequence = String(last?._id || last?.createdAt?.getTime?.() || task.updatedAt?.getTime?.() || taskKey); + await recordForRecipients(recipients, { + podId: task.podId, kind: 'decision' as Kind, sourceType: 'task' as SourceType, + sourceId: `${taskKey}:${sequence}`, + title: String(task.title || 'Task needs attention'), detail: compact(last?.text || task.notes, 220), + podName: pod?.name || 'Pod', + }); + } catch (error) { + console.warn('[attention] task materialization failed:', (error as Error).message); + } +}; + +export const resolveTaskAttention = async (task: any): Promise => { + const taskKey = String(task?._id || task?.taskId || '').trim(); + if (!taskKey) return; + try { + const escaped = taskKey.replace(/[|\\{}()[\]^$+*?.]/g, '\\$&'); + await AttentionItem.updateMany( + { 'source.type': 'task', 'source.id': { $regex: new RegExp('^' + escaped + ':') }, status: 'open' }, + { $set: { status: 'resolved', resolvedAt: new Date() } }, + ); + } catch (error) { + console.warn('[attention] task resolution storage failed; leaving item visible:', (error as Error).message); + } +}; + export const resolve = async (sourceType: SourceType, sourceId: unknown): Promise => { const id = sourceKey(sourceType, sourceId); if (!id) return; @@ -178,6 +227,6 @@ export const acknowledgeMention = async (recipientUserId: unknown, attentionItem return result.modifiedCount === 1 ? { success: true } : { success: false, error: 'Attention item not found' }; }; -export default { recordMentionedUsers, recordApproval, recordDecision, resolve, resolveMany, getOpenQueue, acknowledgeMention }; +export default { recordMentionedUsers, recordApproval, recordDecision, recordTaskAttention, resolveTaskAttention, resolve, resolveMany, getOpenQueue, acknowledgeMention }; // eslint-disable-next-line @typescript-eslint/no-require-imports -module.exports = { recordMentionedUsers, recordApproval, recordDecision, resolve, resolveMany, getOpenQueue, acknowledgeMention }; +module.exports = { recordMentionedUsers, recordApproval, recordDecision, recordTaskAttention, resolveTaskAttention, resolve, resolveMany, getOpenQueue, acknowledgeMention, TASK_HANDOFF_RE }; diff --git a/docs/adr/ADR-017-attention-routing.md b/docs/adr/ADR-017-attention-routing.md index 9eb5e16a1..2074e76c8 100644 --- a/docs/adr/ADR-017-attention-routing.md +++ b/docs/adr/ADR-017-attention-routing.md @@ -376,6 +376,8 @@ The table above returns three "needs a …" verdicts. They are three changes, no **Correction (@sprint-review, 2026-08-29, reproduced at `origin/main`): that store already exists, so this section is a migration and not a build.** `User.activityQueue.acknowledgedMentionIds` (`models/User.ts`) is live end to end — written by `ActivityService.acknowledgeMention`, reached by `POST /api/activity/:activityId/acknowledge` (`routes/activity.ts`), consumed by #1274's `V2ActivityPage.tsx`, and read in `ActivityService.getRecap`, where it already filters acked mentions out of the queue. Its own inline comment makes this section's argument: *"per-(user, message) state, rather than a recent-feed cache: an acknowledged mention must not resurface merely because more messages arrive later."* An earlier revision said v1 **must store** the acknowledgement, which reads as *nothing does* — an absence asserted without naming the instrument that failed to find it, which is the error this document spends §Layer-3.1 warning about, committed here against a store one grep away. +**Superseded by TASK-112 (2026-09-03):** the queue now materializes recipient-owned `AttentionItem` facts at each source write, and acknowledgement resolves that record rather than mutating `User.activityQueue`. A one-time direct-source backfill preserves open 14-day mentions, pending approvals and decisions, plus blocked or explicit-handoff board facts while the retired acknowledgement ids suppress already-handled mentions. The legacy field and recap reader are removed with that rollout. + **The invariant below is satisfied too, and by construction rather than by discipline.** The reader at `:285` is a `.filter` that *excludes* acked ids — no path lets an ack create or retain a row. The shape below is therefore the **migration target**, not a greenfield design: diff --git a/frontend/src/v2/__tests__/V2ActivityPage.test.tsx b/frontend/src/v2/__tests__/V2ActivityPage.test.tsx index 752a5bdde..f3f7dd61c 100644 --- a/frontend/src/v2/__tests__/V2ActivityPage.test.tsx +++ b/frontend/src/v2/__tests__/V2ActivityPage.test.tsx @@ -262,6 +262,27 @@ describe('V2ActivityPage', () => { )); }); + test('keeps a task attention row as an open-thread fact when it has no declared options', async () => { + const taskQueue = { + items: [{ + id: 'task-1:blocked', attentionItemId: 'attention-task-1', kind: 'decision', + title: 'Choose a deploy shape', detail: 'Blocked on an upstream choice.', + podId: 'pod-1', podName: 'Launch pod', options: [], createdAt: '2026-08-26T11:00:00.000Z', + }], + count: 1, + composePodId: null, + }; + mockGet.mockImplementation((url: string) => Promise.resolve({ + data: url === '/api/activity/decision-queue' ? taskQueue : { ...recap, needsYou: [] }, + })); + renderPage(); + + expect(await screen.findByText('Choose a deploy shape')).toBeInTheDocument(); + expect(screen.queryByRole('button', { name: 'Other…' })).not.toBeInTheDocument(); + expect(screen.queryByRole('button', { name: /Rule:/ })).not.toBeInTheDocument(); + expect(screen.getAllByRole('button', { name: 'Open thread' })).not.toHaveLength(0); + }); + test('composes an ordinary pod message into the most recently addressed pod', async () => { mockPost.mockResolvedValue({ data: { id: 123 } }); renderPage(); diff --git a/frontend/src/v2/components/V2ActivityPage.tsx b/frontend/src/v2/components/V2ActivityPage.tsx index 564caedad..255d0bf1d 100644 --- a/frontend/src/v2/components/V2ActivityPage.tsx +++ b/frontend/src/v2/components/V2ActivityPage.tsx @@ -466,7 +466,7 @@ const V2ActivityPage: React.FC = () => { )} - {item.kind === 'decision' && ( + {item.kind === 'decision' && (item.options || []).length > 0 && ( <> {ruledDecisions[item.id] ? (