From cfb3cd26572ab9a881d00f2750e5f6d08fc87885 Mon Sep 17 00:00:00 2001 From: Maneek21 <208369276+Maneek21@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:34:52 +0530 Subject: [PATCH 01/21] feat: add opt-in native mentions with passive agent attention --- .env.example | 4 + .github/workflows/native-mentions.yml | 79 ++++ apps/api/src/index.ts | 2 + .../src/lib/agent-mention-normalization.ts | 24 +- apps/api/src/lib/attention.ts | 22 +- apps/api/src/lib/mcp-tools/human.ts | 87 +++- apps/api/src/lib/mcp-tools/index.ts | 8 + .../src/lib/mcp-tools/mention-attention.ts | 50 +++ apps/api/src/lib/mcp-tools/messages.ts | 2 + apps/api/src/lib/mentions.ts | 6 + apps/api/src/lib/native-mention-visibility.ts | 60 +++ apps/api/src/lib/native-mentions.ts | 396 ++++++++++++++++++ apps/api/src/lib/notification-policy.ts | 5 +- apps/api/src/routes/inbox.ts | 5 + apps/api/src/routes/messages.ts | 32 +- apps/api/src/routes/native-mentions.ts | 49 +++ apps/api/src/routes/notifications.ts | 6 + apps/api/src/routes/tasks.ts | 18 +- .../src/workers/handlers/native-mentions.ts | 18 + apps/api/src/workers/index.ts | 12 + .../test/agent-mention-normalization.test.ts | 6 + apps/api/test/fixtures/native-mentions.ts | 45 ++ .../fixtures/seed-native-mention-evidence.ts | 29 ++ apps/api/test/native-mentions-db.test.ts | 217 ++++++++++ apps/web/src/app/(app)/knowledge/page.tsx | 38 +- apps/web/src/app/(app)/notes/page.tsx | 9 + .../editor/native-mention-popup.tsx | 44 ++ .../components/native-mention-backlinks.tsx | 31 ++ .../src/components/native-mention-content.tsx | 25 ++ .../components/native-mention-markdown.tsx | 39 ++ .../src/components/native-mention-publish.tsx | 26 ++ .../components/native-mention-textarea.tsx | 54 +++ .../src/components/native-reference-chip.tsx | 37 ++ apps/web/src/components/rich-composer.tsx | 5 +- apps/web/src/components/space-chat.tsx | 4 + apps/web/src/components/task-detail.tsx | 76 +++- apps/web/src/components/thread-panel.tsx | 4 + .../lib/editor/native-mention-extension.tsx | 90 ++++ apps/web/src/lib/editor/shared-config.ts | 2 + apps/web/src/lib/native-mentions.ts | 90 ++++ apps/web/src/lib/sanitize.ts | 1 + docs/decisions/2026-10-02-native-mentions.md | 17 + .../2026-10-02-universal-native-mentions.md | 226 ++++++++++ packages/db/scripts/apply-extras.ts | 14 + packages/db/src/schema.ts | 39 ++ .../0.3.0-preview.61-native-mentions.sql | 104 +++++ packages/db/upgrades/manifest.ts | 5 + packages/shared/package.json | 4 +- packages/shared/src/index.ts | 2 + packages/shared/src/native-mentions.ts | 118 ++++++ packages/shared/src/resources-v2.ts | 119 ++++++ packages/shared/test/native-mentions.test.ts | 60 +++ packages/shared/test/resources-v2.test.ts | 110 +++++ scripts/native-mention-browser-evidence.mjs | 205 +++++++++ 54 files changed, 2693 insertions(+), 87 deletions(-) create mode 100644 .github/workflows/native-mentions.yml create mode 100644 apps/api/src/lib/mcp-tools/mention-attention.ts create mode 100644 apps/api/src/lib/native-mention-visibility.ts create mode 100644 apps/api/src/lib/native-mentions.ts create mode 100644 apps/api/src/routes/native-mentions.ts create mode 100644 apps/api/src/workers/handlers/native-mentions.ts create mode 100644 apps/api/test/fixtures/native-mentions.ts create mode 100644 apps/api/test/fixtures/seed-native-mention-evidence.ts create mode 100644 apps/api/test/native-mentions-db.test.ts create mode 100644 apps/web/src/components/editor/native-mention-popup.tsx create mode 100644 apps/web/src/components/native-mention-backlinks.tsx create mode 100644 apps/web/src/components/native-mention-content.tsx create mode 100644 apps/web/src/components/native-mention-markdown.tsx create mode 100644 apps/web/src/components/native-mention-publish.tsx create mode 100644 apps/web/src/components/native-mention-textarea.tsx create mode 100644 apps/web/src/components/native-reference-chip.tsx create mode 100644 apps/web/src/lib/editor/native-mention-extension.tsx create mode 100644 apps/web/src/lib/native-mentions.ts create mode 100644 docs/decisions/2026-10-02-native-mentions.md create mode 100644 docs/superpowers/plans/2026-10-02-universal-native-mentions.md create mode 100644 packages/db/upgrades/0.3.0-preview.61-native-mentions.sql create mode 100644 packages/shared/src/native-mentions.ts create mode 100644 packages/shared/src/resources-v2.ts create mode 100644 packages/shared/test/native-mentions.test.ts create mode 100644 packages/shared/test/resources-v2.test.ts create mode 100644 scripts/native-mention-browser-evidence.mjs diff --git a/.env.example b/.env.example index de9f06a1..69693d0a 100644 --- a/.env.example +++ b/.env.example @@ -148,3 +148,7 @@ METRICS_SCRAPE_TOKEN= VAPID_PUBLIC_KEY= VAPID_PRIVATE_KEY= VAPID_SUBJECT=mailto:admin@example.com + +# Native mentions: enable after running the supported schema upgrade. +# References remain readable when publication is disabled. +DEFT_NATIVE_MENTIONS_ENABLED=false diff --git a/.github/workflows/native-mentions.yml b/.github/workflows/native-mentions.yml new file mode 100644 index 00000000..fb56fec3 --- /dev/null +++ b/.github/workflows/native-mentions.yml @@ -0,0 +1,79 @@ +name: Native Mentions +on: + pull_request: + branches: [master, main] + workflow_dispatch: +permissions: + contents: read +jobs: + native-mentions: + runs-on: ubuntu-latest + timeout-minutes: 25 + services: + postgres: + image: pgvector/pgvector:pg16 + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: deft_mentions_test + ports: ['5432:5432'] + options: >- + --health-cmd "pg_isready -U postgres" + --health-interval 5s --health-timeout 5s --health-retries 10 + env: + DATABASE_URL: postgres://postgres:postgres@localhost:5432/deft_mentions_test + DEFT_TEST_DATABASE_URL: postgres://postgres:postgres@localhost:5432/deft_mentions_test + JWT_SECRET: native-mention-ci-only-secret + JWT_REFRESH_SECRET: native-mention-ci-only-refresh + DEFT_NATIVE_MENTIONS_ENABLED: 'true' + DEFT_NATIVE_MENTION_CONCURRENCY_CERTIFY: 'true' + API_PORT: '4011' + NEXT_PUBLIC_APP_URL: http://localhost:4010 + NEXT_PUBLIC_API_URL: http://localhost:4011 + DEFT_WEB_URL: http://localhost:4010 + steps: + - uses: actions/checkout@v7 + - uses: pnpm/action-setup@v6.1.0 + with: + version: 11.10.0 + - uses: actions/setup-node@v7 + with: + node-version: 22 + cache: pnpm + - run: pnpm install --frozen-lockfile + - run: pnpm --filter @deft/app-kit build + - run: pnpm exec playwright install --with-deps chromium + - name: Fresh schema and concurrent publication/delivery + run: | + pnpm db:push-full + pnpm --filter @deft/api exec tsx --test test/native-mentions-db.test.ts + - name: Separate browser fixture database + run: | + psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -c "CREATE DATABASE deft_mentions_browser_test" + echo "DATABASE_URL=postgres://postgres:postgres@localhost:5432/deft_mentions_browser_test" >> "$GITHUB_ENV" + echo "DEFT_TEST_DATABASE_URL=postgres://postgres:postgres@localhost:5432/deft_mentions_browser_test" >> "$GITHUB_ENV" + echo "DEFT_MENTION_FIXTURE_PATH=$RUNNER_TEMP/native-mention-fixture.json" >> "$GITHUB_ENV" + echo "DEFT_MENTION_EVIDENCE_DIR=$RUNNER_TEMP/native-mention-evidence" >> "$GITHUB_ENV" + - name: Run real desktop, recipient and mobile journeys + run: | + set -euo pipefail + pnpm db:push-full + pnpm --filter @deft/api exec tsx test/fixtures/seed-native-mention-evidence.ts + mkdir -p "$DEFT_MENTION_EVIDENCE_DIR" + pnpm --filter @deft/api exec tsx src/server.ts > "$DEFT_MENTION_EVIDENCE_DIR/api.log" 2>&1 & + api_pid=$! + pnpm --filter @deft/web exec next dev --port 4010 > "$DEFT_MENTION_EVIDENCE_DIR/web.log" 2>&1 & + web_pid=$! + trap 'kill "$api_pid" "$web_pid" 2>/dev/null || true' EXIT + for attempt in $(seq 1 90); do + if curl -fsS http://localhost:4011/health >/dev/null && curl -fsS http://localhost:4010/login >/dev/null; then break; fi + sleep 1 + done + node scripts/native-mention-browser-evidence.mjs + - name: Retain screenshots, recordings and diagnostics + if: always() + uses: actions/upload-artifact@v7 + with: + name: native-mention-evidence + path: ${{ runner.temp }}/native-mention-evidence + if-no-files-found: warn diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index 17e69a3a..6f8787a1 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -43,6 +43,7 @@ import { teamRoutes } from './routes/teams.js'; import { emojiRoutes } from './routes/emoji.js'; import { workflowRoutes } from './routes/workflows.js'; import { crossReferenceRoutes } from './routes/cross-references.js'; +import { nativeMentionRoutes } from './routes/native-mentions.js'; import { auditRoutes } from './routes/audit.js'; import { decisionRoutes } from './routes/decisions.js'; import { managerRoutes } from './routes/manager.js'; @@ -222,6 +223,7 @@ app.route('/api/teams', teamRoutes); app.route('/api/emoji', emojiRoutes); app.route('/api/workflows', workflowRoutes); app.route('/api', crossReferenceRoutes); +app.route('/api/native-mentions', nativeMentionRoutes); app.route('/api', moduleTaskLinkRoutes); app.route('/api/audit', auditRoutes); app.route('/api/decisions', decisionRoutes); diff --git a/apps/api/src/lib/agent-mention-normalization.ts b/apps/api/src/lib/agent-mention-normalization.ts index ac9055af..fc5490b3 100644 --- a/apps/api/src/lib/agent-mention-normalization.ts +++ b/apps/api/src/lib/agent-mention-normalization.ts @@ -3,6 +3,7 @@ export type AgentMentionIdentity = { name: string; slug: string; }; +import { stripNativeMentionAtoms } from '@deft/shared'; export type PlainAgentMentionResolution = { content: string; @@ -42,7 +43,8 @@ export function normalizePlainAgentMentions( } } - let normalized = content; + const segments = content.split(/(]*data-deft-ref-kind[^>]*>[\s\S]*?<\/span>|\[\[deft:(?:person|task|wiki_page):[^\]]+\]\])/gi); + let normalized = segments; const resolvedUserIds = new Set(); const ambiguousAliases = new Set(); const aliases = Array.from(ownersByAlias.keys()).sort((a, b) => b.length - a.length); @@ -50,7 +52,10 @@ export function normalizePlainAgentMentions( for (const alias of aliases) { const aliasPattern = escapeRegex(alias).replace(/\\ /g, '\\s+'); const pattern = new RegExp(`(^|[^a-z0-9_])@(${aliasPattern})(?=$|[^a-z0-9_-])`, 'gi'); - if (!pattern.test(normalized)) continue; + if (!normalized.some(segment => { + pattern.lastIndex = 0; + return stripNativeMentionAtoms(segment) !== '' && pattern.test(segment); + })) continue; pattern.lastIndex = 0; const owners = ownersByAlias.get(alias) ?? []; @@ -60,15 +65,18 @@ export function normalizePlainAgentMentions( } const agent = owners[0]!; - normalized = normalized.replace( - pattern, - (_match, prefix) => `${prefix}<@${agent.userId}|${agent.name}>`, - ); - resolvedUserIds.add(agent.userId); + normalized = normalized.map(segment => { + if (stripNativeMentionAtoms(segment) === '') return segment; + pattern.lastIndex = 0; + return segment.replace(pattern, (_match, prefix) => { + resolvedUserIds.add(agent.userId); + return `${prefix}<@${agent.userId}|${agent.name}>`; + }); + }); } return { - content: normalized, + content: normalized.join(''), resolvedUserIds: Array.from(resolvedUserIds), ambiguousAliases: Array.from(ambiguousAliases), }; diff --git a/apps/api/src/lib/attention.ts b/apps/api/src/lib/attention.ts index 715183cc..36784690 100644 --- a/apps/api/src/lib/attention.ts +++ b/apps/api/src/lib/attention.ts @@ -22,6 +22,7 @@ import { isModuleWriteActionName, } from './module-action-visibility.js'; import { scheduleAttentionDeliveries, scheduleAttentionDelivery } from './web-push.js'; +import { nativeDeliveryAccessSql } from './native-mention-visibility.js'; export type AttentionLane = 'needs_you' | 'updates'; export type AttentionPriority = 'critical' | 'high' | 'normal' | 'low'; @@ -55,7 +56,9 @@ export type AttentionDraft = { export function visibleAttentionCondition(userId: string) { return sql`( - ${attentionItems.source_type} NOT IN ('message', 'space', 'agent_action') + ${attentionItems.source_type} NOT IN ('message', 'space', 'agent_action', 'native_mention') + OR (${attentionItems.source_type} = 'native_mention' + AND ${nativeDeliveryAccessSql(userId, sql`${attentionItems.source_id}`, sql`${attentionItems.org_id}`)}) OR ( ${attentionItems.source_type} = 'message' AND EXISTS ( @@ -151,6 +154,15 @@ function sourceFromLink(link: string | null): { messageId: string | null; spaceI export function notificationToAttentionDraft(notification: LegacyNotification): AttentionDraft { const metadata = objectMetadata(notification.metadata); + const nativeDeliveryId = metadataString(metadata, 'native_mention_delivery_id'); + if (nativeDeliveryId) return { + orgId: notification.org_id, userId: notification.user_id, kind: 'mention', + lane: 'needs_you', priority: 'normal', dedupeKey: `native-mention:${nativeDeliveryId}`, + sourceType: 'native_mention', sourceId: nativeDeliveryId, + sourceEventId: `native-mention:${nativeDeliveryId}`, title: notification.title, + body: notification.body, link: notification.link, metadata, + occurredAt: notification.created_at, + }; const linkedSource = sourceFromLink(notification.link); const taskId = metadataString(metadata, 'task_id', 'taskId'); const messageId = metadataString(metadata, 'message_id', 'messageId', 'source_message_id') ?? linkedSource.messageId; @@ -567,6 +579,12 @@ export async function filterVisibleAttentionItems { + const nativeIds = rows.filter(item => item.source_type === 'native_mention').map(item => item.id); + const visibleNative = nativeIds.length ? await db.select({ id: attentionItems.id }).from(attentionItems).where(and( + inArray(attentionItems.id, nativeIds), + nativeDeliveryAccessSql(userId, sql`${attentionItems.source_id}`, sql`${attentionItems.org_id}`), + )) : []; + const allowedNative = new Set(visibleNative.map(item => item.id)); const messageIds = rows.filter((item) => item.source_type === 'message').map((item) => item.source_id); const spaceIds = rows.filter((item) => item.source_type === 'space').map((item) => item.source_id); const actionIds = rows.filter((item) => item.source_type === 'agent_action').map((item) => item.source_id); @@ -607,6 +625,8 @@ export async function filterVisibleAttentionItems row.type === 'public' || row.member_id).map((row) => row.id)); const allowedActions = new Set(visibleActions.map((row) => row.id)); const inaccessible = rows.filter((item) => + (item.source_type === 'native_mention' && !allowedNative.has(item.id)) + || (item.source_type === 'message' && !allowedMessages.has(item.source_id)) || (item.source_type === 'space' && !allowedSpaces.has(item.source_id)) || (item.source_type === 'agent_action' && !allowedActions.has(item.source_id))); diff --git a/apps/api/src/lib/mcp-tools/human.ts b/apps/api/src/lib/mcp-tools/human.ts index a9800a0d..ab701fa5 100644 --- a/apps/api/src/lib/mcp-tools/human.ts +++ b/apps/api/src/lib/mcp-tools/human.ts @@ -1,6 +1,10 @@ import { createHash, randomUUID } from 'node:crypto'; import { and, desc, eq, inArray, or, sql } from 'drizzle-orm'; import { db } from '../db.js'; +import { z } from 'zod'; +import { NativeMentionSourceSchema } from '@deft/shared'; +import { searchNativeMentions, publishNativeMentions, enqueueNativeMentionPublication } from '../native-mentions.js'; +import { nativeNotificationAccessSql } from '../native-mention-visibility.js'; import { connectedAccounts, agentActions, @@ -202,6 +206,7 @@ function retrievalResultToSearchResult(row: ContextResult): Record MODULE_OPERATION_DEFINITIONS[name].mode === 'read'), 'search', 'fetch', @@ -255,6 +260,7 @@ export const HUMAN_READ_TOOLS = new Set([ ]); export const HUMAN_WRITE_TOOLS = new Set([ + 'native_mentions_publish', ...MODULE_OPERATION_NAMES.filter((name) => MODULE_OPERATION_DEFINITIONS[name].mode === 'write'), 'memory_write', 'wiki_upsert', @@ -312,6 +318,23 @@ async function humanAppActionOperation( } export const HUMAN_TOOLS: Record = { + native_mentions_search: async (args, ctx) => { + const parsed = z.object({ query: z.string().max(120).default('') }).strict().safeParse(args); + if (!parsed.success) return errorResult('Invalid mention search'); + const scopeError = requireScope(ctx, 'read:workspace'); if (scopeError) return scopeError; + const items = await searchNativeMentions({ orgId: ctx.org_id, userId: ctx.user_id }, parsed.data.query); + return textResult({ items: items.filter(item => item.ref.resource_type === 'person' + || ctx.scopes.includes(item.ref.resource_type === 'task' ? 'read:tasks' : 'read:wiki')) }); + }, + native_mentions_publish: async (args, ctx) => { + const parsed = z.object({ source: NativeMentionSourceSchema, content_hash: z.string().regex(/^[a-f0-9]{64}$/) }).strict().safeParse(args); + if (!parsed.success) return errorResult('Invalid mention publication'); + const source = parsed.data.source; + const scope = source.kind === 'message' ? 'write:messages' : + source.kind === 'task' || source.kind === 'task_comment' ? 'write:tasks' : source.kind === 'wiki_page' ? 'write:wiki' : 'write:workspace'; + const scopeError = requireScope(ctx, scope); if (scopeError) return scopeError; + return textResult(await publishNativeMentions({ orgId: ctx.org_id, userId: ctx.user_id }, source, parsed.data.content_hash)); + }, ...Object.fromEntries(MODULE_OPERATION_NAMES.map((name) => [ name, (args: Record, ctx: HumanToolContext) => humanModuleOperation(name, args, ctx), @@ -2335,12 +2358,14 @@ export async function humanCommentOnTask(args: { task_id?: string; content?: str .where(and(eq(tasks.id, taskId), eq(tasks.org_id, ctx.org_id), eq(tasks.is_deleted, false))) .limit(1); if (!task) return errorResult('comment_on_task: task not found'); - const [comment] = await db.insert(taskComments).values({ - org_id: ctx.org_id, - task_id: taskId, - user_id: ctx.user_id, - content, - }).returning(); + const comment = await db.transaction(async tx => { + const [inserted] = await tx.insert(taskComments).values({ + org_id: ctx.org_id, task_id: taskId, user_id: ctx.user_id, content, + }).returning(); + await enqueueNativeMentionPublication(tx, { orgId: ctx.org_id, userId: ctx.user_id }, + { kind: 'task_comment', id: inserted!.id }, content); + return inserted; + }); await db.insert(taskActivity).values({ org_id: ctx.org_id, task_id: taskId, user_id: ctx.user_id, action: 'commented' }); await publishTaskChannelEventForAssignee({ orgId: ctx.org_id, @@ -2376,13 +2401,14 @@ export async function humanMessagePost(args: { space_id?: string; content?: stri .limit(1); if (!parent || parent.space_id !== spaceId) return errorResult('message_post: parent message not found in target space'); } - const [row] = await db.insert(messages).values({ - org_id: ctx.org_id, - space_id: spaceId, - user_id: ctx.user_id, - content, - parent_id: args.parent_id ?? null, - }).returning(); + const row = await db.transaction(async tx => { + const [inserted] = await tx.insert(messages).values({ + org_id: ctx.org_id, space_id: spaceId, user_id: ctx.user_id, content, parent_id: args.parent_id ?? null, + }).returning(); + await enqueueNativeMentionPublication(tx, { orgId: ctx.org_id, userId: ctx.user_id }, + { kind: 'message', id: inserted!.id }, content); + return inserted; + }); await dispatchAgentEmployeeMessage({ messageId: row!.id, spaceId, @@ -2466,13 +2492,14 @@ export async function humanSendMessage(args: HumanSendMessageArgs, ctx: HumanToo targetKind = 'dm'; } - const [row] = await db.insert(messages).values({ - org_id: ctx.org_id, - space_id: spaceId!, - user_id: ctx.user_id, - content, - parent_id: parentId, - }).returning(); + const row = await db.transaction(async tx => { + const [inserted] = await tx.insert(messages).values({ + org_id: ctx.org_id, space_id: spaceId!, user_id: ctx.user_id, content, parent_id: parentId, + }).returning(); + await enqueueNativeMentionPublication(tx, { orgId: ctx.org_id, userId: ctx.user_id }, + { kind: 'message', id: inserted!.id }, content); + return inserted; + }); await dispatchAgentEmployeeMessage({ messageId: row!.id, @@ -3473,6 +3500,7 @@ export async function humanInboxList(args: { unread_only?: boolean; type?: strin if (scopeError) return scopeError; const rows = await db.select().from(notifications).where(and( eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id), + nativeNotificationAccessSql(ctx.user_id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), args.unread_only === false ? sql`true` : eq(notifications.is_read, false), args.type ? eq(notifications.type, args.type as any) : sql`true`, )).orderBy(desc(notifications.created_at)).limit(Math.min(Math.max(1, args.limit ?? 50), 100)); @@ -3482,7 +3510,8 @@ export async function humanInboxList(args: { unread_only?: boolean; type?: strin export async function humanInboxGet(args: { notification_id?: string }, ctx: HumanToolContext): Promise { const scopeError = requireScope(ctx, 'read:workspace'); if (scopeError) return scopeError; if (!args.notification_id) return errorResult('notification_id is required'); - const [row] = await db.select().from(notifications).where(and(eq(notifications.id, args.notification_id), eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id))).limit(1); + const [row] = await db.select().from(notifications).where(and(eq(notifications.id, args.notification_id), eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id), + nativeNotificationAccessSql(ctx.user_id, sql`${notifications.metadata}`, sql`${notifications.org_id}`))).limit(1); return row ? textResult(row) : errorResult('Notification not found'); } @@ -3490,7 +3519,8 @@ export async function humanInboxMarkRead(args: Record, ctx: Hum const scopeError = requireScope(ctx, 'write:workspace'); if (scopeError) return scopeError; if (typeof args.notification_id !== 'string') return errorResult('notification_id is required'); return withIdempotency('inbox_mark_read', args, ctx, async () => { - const [row] = await db.update(notifications).set({ is_read: true }).where(and(eq(notifications.id, args.notification_id as string), eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id))).returning(); + const [row] = await db.update(notifications).set({ is_read: true }).where(and(eq(notifications.id, args.notification_id as string), eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id), + nativeNotificationAccessSql(ctx.user_id, sql`${notifications.metadata}`, sql`${notifications.org_id}`))).returning(); return row ? textResult({ marked_read: true, notification: row }) : errorResult('Notification not found'); }); } @@ -3498,7 +3528,8 @@ export async function humanInboxMarkRead(args: Record, ctx: Hum export async function humanInboxMarkAllRead(args: Record, ctx: HumanToolContext): Promise { const scopeError = requireScope(ctx, 'write:workspace'); if (scopeError) return scopeError; return withIdempotency('inbox_mark_all_read', args, ctx, async () => { - const rows = await db.update(notifications).set({ is_read: true }).where(and(eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id), eq(notifications.is_read, false))).returning({ id: notifications.id }); + const rows = await db.update(notifications).set({ is_read: true }).where(and(eq(notifications.org_id, ctx.org_id), eq(notifications.user_id, ctx.user_id), eq(notifications.is_read, false), + nativeNotificationAccessSql(ctx.user_id, sql`${notifications.metadata}`, sql`${notifications.org_id}`))).returning({ id: notifications.id }); return textResult({ marked_read: rows.length }); }); } @@ -3655,6 +3686,8 @@ export async function humanAgentEmployeeUpdateState(args: Record = { + native_mentions_search: 'read:workspace', + native_mentions_publish: ['write:workspace', 'write:messages', 'write:tasks', 'write:wiki'], search: ['read:workspace', 'read:wiki', 'read:tasks', 'read:messages', 'read:calendar', 'read:modules'], fetch: ['read:workspace', 'read:wiki', 'read:tasks', 'read:messages', 'read:calendar', 'read:modules'], platform_context: 'read:workspace', attention_digest: 'read:workspace', @@ -3754,6 +3787,14 @@ function operationalHumanSchemas(): Array> { const iso = (description: string) => ({ type: 'string', description }); return [ read('workspace_capabilities', 'Inspect Deft MCP Capabilities', 'Return granted scopes, available operational tools, usage guidance, and the intentionally UI-only boundary.'), + read('native_mentions_search', 'Search Native Mentions', 'Find authorized people, agents, tasks and wikis. Store stable identity tokens, never copy a display label as authority.', { query: { type: 'string', maxLength: 120 } }), + { name: 'native_mentions_publish', title: 'Publish Native Mentions', + description: 'Notify newly added people and agents from saved native content. Requires the source write scope and exact SHA-256 content hash. Repeat publication is idempotent.', + annotations: { readOnlyHint: false, destructiveHint: false }, + inputSchema: { type: 'object', properties: { + source: { type: 'object', properties: { kind: { type: 'string', enum: ['message', 'task', 'task_comment', 'wiki_page', 'note'] }, id: { type: 'string' } }, required: ['kind', 'id'], additionalProperties: false }, + content_hash: { type: 'string', pattern: '^[a-f0-9]{64}$' }, + }, required: ['source', 'content_hash'], additionalProperties: false } }, read('note_list', 'List Deft Notes', 'List private notes owned by the connected user plus org/space notes they can see.', { query: { type: 'string' }, limit }), read('note_get', 'Get Deft Note', 'Read one visible note including its full TipTap HTML content.', { note_id: id('Note id') }, ['note_id']), write('note_create', 'Create Deft Note', 'Create a private, org, or space-visible note as the connected user.', { title: { type: 'string' }, content: { type: 'string' }, icon: { type: 'string' }, is_pinned: { type: 'boolean' }, visibility: { type: 'string', enum: ['private', 'org', 'space'] }, visibility_space_id: id('Required for space visibility') }), diff --git a/apps/api/src/lib/mcp-tools/index.ts b/apps/api/src/lib/mcp-tools/index.ts index 40809a77..f9fc2ef1 100644 --- a/apps/api/src/lib/mcp-tools/index.ts +++ b/apps/api/src/lib/mcp-tools/index.ts @@ -20,6 +20,7 @@ import { memoryUpdate } from './memory-update.js'; import { taskQuery } from './tasks.js'; import { memberList } from './members.js'; import { threadFetch, fetchUnread } from './messages.js'; +import { mentionAttentionList, mentionAttentionAcknowledge } from './mention-attention.js'; import { attachmentList, attachmentRead } from './attachments.js'; import { workspacePlanImport } from './workspace-plan-import.js'; import { documentSend } from './document-send.js'; @@ -135,6 +136,9 @@ export const READ_ONLY_TOOLS: Record = { poll_pending_work: pollPendingWork as ToolHandler, ping_alive: pingAlive as ToolHandler, fetch_unread: fetchUnread as ToolHandler, + mention_attention_list: mentionAttentionList, + // Own read receipt only; this never authorizes work or writes to a source. + mention_attention_acknowledge: mentionAttentionAcknowledge, }; export const TOOL_ALIASES: Record = { @@ -190,6 +194,10 @@ const CALLER_SLUG_PROP = { }; export const toolSchemas: ToolSchema[] = [ + { name: 'mention_attention_list', description: 'Read your passive native mention attention. Reading is not a request to execute work.', + inputSchema: { type: 'object', properties: { caller_employee_slug: { type: 'string' } } } }, + { name: 'mention_attention_acknowledge', description: 'Acknowledge your own passive mention without performing work or changing its source.', + inputSchema: { type: 'object', properties: { caller_employee_slug: { type: 'string' }, attention_id: { type: 'string' } }, required: ['attention_id'] } }, ...(MODULE_MCP_TOOL_SCHEMAS as ToolSchema[]), { name: 'platform_context', diff --git a/apps/api/src/lib/mcp-tools/mention-attention.ts b/apps/api/src/lib/mcp-tools/mention-attention.ts new file mode 100644 index 00000000..79020163 --- /dev/null +++ b/apps/api/src/lib/mcp-tools/mention-attention.ts @@ -0,0 +1,50 @@ +import { and, eq } from 'drizzle-orm'; +import { z } from 'zod'; +import { agentEmployees } from '@deft/db/schema'; +import { NativeMentionSourceSchema } from '@deft/shared'; +import { db } from '../db.js'; +import { listNativeMentionAttention, loadNativeSource, nativeMentionsEnabled } from '../native-mentions.js'; +import { transitionAttentionItem } from '../attention.js'; +import { textResult, errorResult, type ToolContext } from './types.js'; + +export async function boundMentionAttention(ctx: ToolContext) { + if (!nativeMentionsEnabled()) return []; + if (ctx.token_id && !ctx.scopes?.includes('read:workspace')) return []; + const [employee] = await db.select({ user_id: agentEmployees.user_id }).from(agentEmployees).where(and( + eq(agentEmployees.id, ctx.employee_id), eq(agentEmployees.org_id, ctx.org_id), + eq(agentEmployees.is_active, true), eq(agentEmployees.is_deleted, false), + )); + if (!employee) return []; + const context = { orgId: ctx.org_id, userId: employee.user_id, employeeId: ctx.employee_id }; + const items = await listNativeMentionAttention(context); + const result = []; + for (const item of items) { + if (item.state === 'acknowledged') continue; + const source = NativeMentionSourceSchema.safeParse((item.metadata as Record).native_mention_source); + if (!source.success) continue; + const scope = source.data.kind === 'message' ? 'read:messages' : + source.data.kind === 'task' || source.data.kind === 'task_comment' ? 'read:tasks' : + source.data.kind === 'wiki_page' ? 'read:wiki' : 'read:workspace'; + if (ctx.token_id && !ctx.scopes?.includes(scope)) continue; + const current = await loadNativeSource(context, source.data); + if (current) result.push({ ...item, source: source.data, current_source: { + label: current.label, content: current.content.slice(0, 5000), truncated: current.content.length > 5000, + } }); + } + return result; +} +export async function mentionAttentionList(args: unknown, ctx: ToolContext) { + const parsed = z.object({ caller_employee_slug: z.string().optional() }).strict().safeParse(args); + if (!parsed.success) return errorResult('Invalid mention attention arguments'); + return textResult({ mention_attention: await boundMentionAttention(ctx), passive: true }); +} +export async function mentionAttentionAcknowledge(args: unknown, ctx: ToolContext) { + const parsed = z.object({ attention_id: z.string().min(1), caller_employee_slug: z.string().optional() }).strict().safeParse(args); + if (!parsed.success) return errorResult('attention_id is required'); + if (ctx.token_id && !ctx.scopes?.includes('write:workspace')) return errorResult('Missing MCP scope: write:workspace'); + const item = (await boundMentionAttention(ctx)).find(item => item.id === parsed.data.attention_id); + if (!item) return errorResult('Mention attention item not found'); + const updated = await transitionAttentionItem({ orgId: ctx.org_id, userId: item.user_id, + itemId: item.id, state: 'acknowledged', actorUserId: item.user_id }); + return textResult({ acknowledged: Boolean(updated), attention_id: item.id, passive: true }); +} diff --git a/apps/api/src/lib/mcp-tools/messages.ts b/apps/api/src/lib/mcp-tools/messages.ts index 341970a1..377e51a8 100644 --- a/apps/api/src/lib/mcp-tools/messages.ts +++ b/apps/api/src/lib/mcp-tools/messages.ts @@ -12,6 +12,7 @@ import { messages, users, spaces, spaceMembers, agentEmployees, agentActions } f import type { ToolContext, ToolResult } from './types.js'; import { errorResult, textResult } from './types.js'; import { manifestsByMessageId } from '../attachment-manifests.js'; +import { boundMentionAttention } from './mention-attention.js'; /** * Phase 12 review fix: before returning any thread content, verify the @@ -281,6 +282,7 @@ export async function fetchUnread( return textResult({ unread_messages: unreadMessages, pending_actions: actionRows, + mention_attention: await boundMentionAttention(ctx), }); } catch (err) { const msg = err instanceof Error ? err.message : String(err); diff --git a/apps/api/src/lib/mentions.ts b/apps/api/src/lib/mentions.ts index dab8474a..93accb3d 100644 --- a/apps/api/src/lib/mentions.ts +++ b/apps/api/src/lib/mentions.ts @@ -2,9 +2,15 @@ // Content format: "Hey <@userId|userName> check this out" // Returns array of mentioned user IDs // Also handle @here and @all +import { extractNativeMentions, NativeMentionLimitError } from '@deft/shared'; export function parseMentions(content: string): { userIds: string[]; here: boolean; all: boolean } { const userIds = new Set(); + try { + for (const ref of extractNativeMentions(content)) { + if (ref.resource_type === 'person') userIds.add(ref.resource_id); + } + } catch (error) { if (!(error instanceof NativeMentionLimitError)) throw error; } let here = false; let all = false; diff --git a/apps/api/src/lib/native-mention-visibility.ts b/apps/api/src/lib/native-mention-visibility.ts new file mode 100644 index 00000000..924b3057 --- /dev/null +++ b/apps/api/src/lib/native-mention-visibility.ts @@ -0,0 +1,60 @@ +import { sql, type SQL } from 'drizzle-orm'; + +// Watchers are subscriptions, not sufficient authority for new mention surfaces. +export function nativeTaskAccessSql(userId: string): SQL { + return sql`(coalesce(t.metadata->>'visibility', 'org') <> 'restricted' + OR t.created_by = ${userId} OR t.assignee_id = ${userId} OR p.lead_id = ${userId} + OR coalesce(t.metadata->'visible_user_ids', '[]'::jsonb) ? ${userId} + OR EXISTS (SELECT 1 FROM task_assignees ta WHERE ta.task_id = t.id AND ta.user_id = ${userId}))`; +} + +export function nativeSourceAccessSql( + userId: string, kind: SQL, sourceId: SQL, orgId: SQL, +): SQL { + return sql`( + EXISTS (SELECT 1 FROM org_members nm WHERE nm.org_id = ${orgId} + AND nm.user_id = ${userId} AND nm.is_active = true) + AND ( + (${kind} = 'message' AND EXISTS ( + SELECT 1 FROM messages m JOIN spaces s ON s.id = m.space_id AND s.org_id = m.org_id + WHERE m.id = ${sourceId} AND m.org_id = ${orgId} AND m.is_deleted = false + AND (s.type = 'public' OR EXISTS (SELECT 1 FROM space_members sm WHERE sm.space_id = s.id AND sm.user_id = ${userId})))) + OR (${kind} = 'task' AND EXISTS ( + SELECT 1 FROM tasks t JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE t.id = ${sourceId} AND t.org_id = ${orgId} AND t.is_deleted = false AND p.is_deleted = false + AND ${nativeTaskAccessSql(userId)})) + OR (${kind} = 'task_comment' AND EXISTS ( + SELECT 1 FROM task_comments tc JOIN tasks t ON t.id = tc.task_id AND t.org_id = tc.org_id + JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE tc.id = ${sourceId} AND tc.org_id = ${orgId} AND tc.is_deleted = false + AND t.is_deleted = false AND p.is_deleted = false AND ${nativeTaskAccessSql(userId)})) + OR (${kind} = 'wiki_page' AND EXISTS ( + SELECT 1 FROM wiki_pages w WHERE w.id = ${sourceId} AND w.org_id = ${orgId} AND w.is_deleted = false + AND (w.scope = 'org' OR w.user_id = ${userId} + OR (w.scope = 'space' AND EXISTS (SELECT 1 FROM space_members sm WHERE sm.space_id = w.space_id AND sm.user_id = ${userId}))))) + OR (${kind} = 'note' AND EXISTS ( + SELECT 1 FROM notes n WHERE n.id = ${sourceId} AND n.org_id = ${orgId} AND n.is_deleted = false + AND (n.user_id = ${userId} OR n.visibility = 'org' + OR EXISTS (SELECT 1 FROM note_shares ns WHERE ns.note_id = n.id AND ns.shared_with_user_id = ${userId}) + OR (n.visibility = 'space' AND EXISTS (SELECT 1 FROM space_members sm WHERE sm.space_id = n.visibility_space_id AND sm.user_id = ${userId}))))) + ) + )`; +} + +export function nativeDeliveryAccessSql(userId: string, deliveryId: SQL, orgId: SQL): SQL { + return sql`EXISTS ( + SELECT 1 FROM native_mention_deliveries nd + JOIN native_reference_states nr ON nr.id = nd.source_state_id AND nr.org_id = nd.org_id + WHERE nd.id = ${deliveryId} AND nd.org_id = ${orgId} + AND nd.recipient_user_id = ${userId} AND nr.is_deleted = false + AND ${nativeSourceAccessSql(userId, sql`nr.source_kind`, sql`nr.source_id`, sql`nr.org_id`)} + )`; +} + +/** Legacy notifications retain their existing semantics. New ones fail closed. */ +export function nativeNotificationAccessSql( + userId: string, metadata: SQL, orgId: SQL, +): SQL { + return sql`(${metadata}->>'native_mention_delivery_id' IS NULL + OR ${nativeDeliveryAccessSql(userId, sql`${metadata}->>'native_mention_delivery_id'`, orgId)})`; +} diff --git a/apps/api/src/lib/native-mentions.ts b/apps/api/src/lib/native-mentions.ts new file mode 100644 index 00000000..0d3c3ec4 --- /dev/null +++ b/apps/api/src/lib/native-mentions.ts @@ -0,0 +1,396 @@ +import { createHash } from 'node:crypto'; +import { and, desc, eq, sql } from 'drizzle-orm'; +import { + NativeMentionRefsSchema, NativeMentionSourceSchema, extractNativeMentions, + nativeMentionKey, nativeMentionRef, + NativeMentionLimitError, + type NativeMentionRef, type NativeMentionSource, +} from '@deft/shared'; +import { nativeReferenceStates, nativeMentionDeliveries, notifications, users, orgMembers, agentEmployees, attentionItems } from '@deft/db/schema'; +import { db, withDbAdvisoryLock } from './db.js'; +import { enqueue, QUEUE_NAMES, RetryLaterJobError } from './queues.js'; +import { explainNotificationPolicy } from './notification-policy.js'; +import { upsertAttentionItem } from './attention.js'; +import { emitToUser } from '../socket.js'; +import { nativeSourceAccessSql, nativeTaskAccessSql, nativeDeliveryAccessSql } from './native-mention-visibility.js'; + +export const nativeMentionsEnabled = () => process.env.DEFT_NATIVE_MENTIONS_ENABLED === 'true'; +export type NativeMentionContext = { orgId: string; userId: string; employeeId?: string }; +type Executor = Pick; +type NativeSourceRow = Record & { + id: string; content: string; label: string; owner_user_id: string | null; + space_id: string | null; project_id: string | null; href: string; +}; +export type NativeMentionProjection = { + ref: NativeMentionRef; state: 'available' | 'unavailable'; label?: string; + href?: string; group?: 'People' | 'Agents' | 'Tasks' | 'Wikis'; + avatar_url?: string | null; description?: string | null; +}; +export class NativeMentionError extends Error { + constructor(public code: string, message: string, public status: 400 | 403 | 404 | 409 = 400) { super(message); } +} +export const nativeContentHash = (content: string) => createHash('sha256').update(content).digest('hex'); +const rows = >(result: unknown): T[] => + ((result as { rows?: T[] }).rows ?? []) as T[]; + +async function activeActor(ctx: NativeMentionContext, executor: Executor = db) { + return rows<{ role: string; kind: string; employee_id: string | null; allowed_space_ids: string[] | null; project_ids: string[] | null; runtime_kind: string | null }>( + await executor.execute(sql`SELECT om.role, u.kind, ae.id AS employee_id, ae.space_ids AS allowed_space_ids, ae.project_ids, ae.runtime_kind + FROM org_members om JOIN users u ON u.id = om.user_id + LEFT JOIN agent_employees ae ON ae.user_id = u.id AND ae.org_id = om.org_id + AND ae.is_active = true AND (ae.is_deleted = false OR ae.runtime_kind = 'defty_system') + WHERE om.org_id = ${ctx.orgId} AND om.user_id = ${ctx.userId} AND om.is_active = true + LIMIT 1`), + )[0] ?? null; +} + +async function employeeBoundary(ctx: NativeMentionContext, source: { project_id: string | null; space_id: string | null }, executor: Executor = db): Promise { + const actor = await activeActor(ctx, executor); + if (!actor) return false; + if (ctx.employeeId && actor.employee_id !== ctx.employeeId) return false; + if (actor.kind !== 'agent') return !ctx.employeeId; + if (!actor.employee_id) return false; + if (source.project_id && actor.runtime_kind !== 'defty_system' && actor.project_ids?.length + && !actor.project_ids.includes(source.project_id)) return false; + return !source.space_id || !actor.allowed_space_ids?.length || actor.allowed_space_ids.includes(source.space_id); +} + +export async function loadNativeSource( + ctx: NativeMentionContext | { orgId: string }, source: NativeMentionSource, + executor: Executor = db, lock = false, +): Promise { + const src = NativeMentionSourceSchema.parse(source); + let query; + switch (src.kind) { + case 'message': + query = sql`SELECT m.id, m.content, 'Chat message' AS label, m.user_id AS owner_user_id, + m.space_id, NULL::text AS project_id, + '/chat?space=' || m.space_id || '&message=' || m.id + || CASE WHEN m.parent_id IS NULL THEN '' ELSE '&thread=' || m.parent_id END AS href + FROM messages m WHERE m.org_id = ${ctx.orgId} AND m.id = ${src.id} AND m.is_deleted = false + ${lock ? sql`FOR UPDATE OF m` : sql``}`; break; + case 'task': + query = sql`SELECT t.id, coalesce(t.description, '') AS content, t.title AS label, + t.created_by AS owner_user_id, NULL::text AS space_id, t.project_id, + '/tasks?task=' || t.id || '&field=description' AS href + FROM tasks t JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE t.org_id = ${ctx.orgId} AND t.id = ${src.id} AND t.is_deleted = false AND p.is_deleted = false + ${lock ? sql`FOR UPDATE OF t` : sql``}`; break; + case 'task_comment': + query = sql`SELECT tc.id, tc.content, t.title || ' · Comment' AS label, + tc.user_id AS owner_user_id, NULL::text AS space_id, t.project_id, + '/tasks?task=' || t.id || '&comment=' || tc.id AS href + FROM task_comments tc JOIN tasks t ON t.id = tc.task_id AND t.org_id = tc.org_id + JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE tc.org_id = ${ctx.orgId} AND tc.id = ${src.id} AND tc.is_deleted = false + AND t.is_deleted = false AND p.is_deleted = false + ${lock ? sql`FOR UPDATE OF tc` : sql``}`; break; + case 'wiki_page': + query = sql`SELECT w.id, w.content, w.title AS label, w.user_id AS owner_user_id, + w.space_id, NULL::text AS project_id, '/knowledge?slug=' || w.slug AS href + FROM wiki_pages w WHERE w.org_id = ${ctx.orgId} AND w.id = ${src.id} AND w.is_deleted = false + ${lock ? sql`FOR UPDATE OF w` : sql``}`; break; + case 'note': + query = sql`SELECT n.id, coalesce(n.content, '') AS content, n.title AS label, + n.user_id AS owner_user_id, n.visibility_space_id AS space_id, NULL::text AS project_id, + '/notes?id=' || n.id AS href FROM notes n + WHERE n.org_id = ${ctx.orgId} AND n.id = ${src.id} AND n.is_deleted = false + ${lock ? sql`FOR UPDATE OF n` : sql``}`; break; + } + const row = rows(await executor.execute(query))[0] ?? null; + if (!row || !('userId' in ctx)) return row; + const visible = rows(await executor.execute(sql`SELECT 1 WHERE + ${nativeSourceAccessSql(ctx.userId, sql`${src.kind}`, sql`${src.id}`, sql`${ctx.orgId}`)}`)); + return visible.length && await employeeBoundary(ctx, row, executor) ? row : null; +} + +async function reconcileWithin(executor: Parameters[0]>[0], orgId: string, source: NativeMentionSource) { + const current = await loadNativeSource({ orgId }, source, executor, true); + if (!current) { + await executor.update(nativeReferenceStates).set({ current_refs: [], is_deleted: true, updated_at: new Date() }) + .where(and(eq(nativeReferenceStates.org_id, orgId), eq(nativeReferenceStates.source_kind, source.kind), eq(nativeReferenceStates.source_id, source.id))); + return null; + } + let refs: NativeMentionRef[]; + let quarantined = false; + try { refs = extractNativeMentions(current.content); } + catch (error) { if (!(error instanceof NativeMentionLimitError)) throw error; refs = []; quarantined = true; } + const contentHash = nativeContentHash(current.content); + const [state] = await executor.insert(nativeReferenceStates).values({ + org_id: orgId, source_kind: source.kind, source_id: source.id, + content_hash: contentHash, current_refs: refs, + }).onConflictDoUpdate({ + target: [nativeReferenceStates.org_id, nativeReferenceStates.source_kind, nativeReferenceStates.source_id], + set: { + current_refs: refs, content_hash: contentHash, is_deleted: false, updated_at: new Date(), + revision: sql`CASE WHEN ${nativeReferenceStates.content_hash} = ${contentHash} + THEN ${nativeReferenceStates.revision} ELSE ${nativeReferenceStates.revision} + 1 END`, + }, + }).returning(); + return { state: state!, current, refs, quarantined }; +} + +export async function reconcileNativeMentions(orgId: string, source: NativeMentionSource) { + NativeMentionSourceSchema.parse(source); + return db.transaction(tx => reconcileWithin(tx, orgId, source)); +} + +export async function publishNativeMentions(ctx: NativeMentionContext, source: NativeMentionSource, expectedHash: string) { + if (!nativeMentionsEnabled()) throw new NativeMentionError('NATIVE_MENTIONS_DISABLED', 'Native mention publication is disabled', 403); + const actor = await activeActor(ctx); + if (!actor) throw new NativeMentionError('NOT_FOUND', 'Source not found', 404); + if (actor.kind !== 'human') throw new NativeMentionError('FORBIDDEN', 'Human publication intent is required', 403); + return db.transaction(async tx => { + const current = await loadNativeSource(ctx, source, tx, true); + if (!current) throw new NativeMentionError('NOT_FOUND', 'Source not found', 404); + if ((source.kind === 'message' || source.kind === 'task_comment' || source.kind === 'note') + && current.owner_user_id !== ctx.userId) { + throw new NativeMentionError('FORBIDDEN', 'Only the content author can publish mentions', 403); + } + if (actor.role === 'guest' && source.kind !== 'message' && source.kind !== 'note') { + throw new NativeMentionError('FORBIDDEN', 'Mention publication is not available for this source', 403); + } + if (nativeContentHash(current.content) !== expectedHash) { + throw new NativeMentionError('NATIVE_MENTION_REVISION_CONFLICT', 'Content changed. Save or reload before notifying mentions.', 409); + } + try { extractNativeMentions(current.content); } + catch (error) { + if (error instanceof NativeMentionLimitError) throw new NativeMentionError('NATIVE_MENTION_LIMIT', error.message, 400); + throw error; + } + const reconciled = await reconcileWithin(tx, ctx.orgId, source); + if (!reconciled) throw new NativeMentionError('NOT_FOUND', 'Source not found', 404); + const { state, refs } = reconciled; + const currentIds = refs.filter(ref => ref.resource_type === 'person').map(ref => ref.resource_id).filter(id => id !== ctx.userId); + const oldIds = state.published_person_ids; + const additions = currentIds.filter(id => !oldIds.includes(id)); + const eligible: string[] = []; + let blocked = 0; + for (const id of additions) { + if (await loadNativeSource({ orgId: ctx.orgId, userId: id }, source, tx)) eligible.push(id); + else blocked++; + } + const publishedIds = currentIds.filter(id => oldIds.includes(id) || eligible.includes(id)); + const changed = eligible.length > 0 || oldIds.some(id => !currentIds.includes(id)); + if (!changed) return { source, content_hash: expectedHash, status: 'unchanged' as const, queued_count: 0, blocked_count: blocked }; + const publicationRevision = state.publication_revision + 1; + await tx.update(nativeReferenceStates).set({ + published_person_ids: publishedIds, publication_revision: publicationRevision, updated_at: new Date(), + }).where(and(eq(nativeReferenceStates.id, state.id), eq(nativeReferenceStates.org_id, ctx.orgId))); + for (const recipient of eligible) { + const [delivery] = await tx.insert(nativeMentionDeliveries).values({ + org_id: ctx.orgId, source_state_id: state.id, publication_revision: publicationRevision, + recipient_user_id: recipient, actor_user_id: ctx.userId, + }).onConflictDoNothing().returning(); + if (delivery) await enqueue(QUEUE_NAMES.AGENT_JOBS, 'native-mention-deliver', + { orgId: ctx.orgId, deliveryId: delivery.id }, + { orgId: ctx.orgId, dedupeKey: `native-mention:${delivery.id}`, executor: tx, maxAttempts: 5 }); + } + return { source, content_hash: expectedHash, status: 'queued' as const, queued_count: eligible.length, blocked_count: blocked }; + }); +} + +export async function handleNativeMentionReconciliation(data: { orgId: string; source: NativeMentionSource; publishOnCreate?: boolean }) { + const parsed = NativeMentionSourceSchema.parse(data.source); + await reconcileNativeMentions(data.orgId, parsed); +} + +/** Called only at authenticated human Send/Post boundaries, in the content transaction. */ +export async function enqueueNativeMentionPublication( + executor: Parameters[0]>[0], + ctx: NativeMentionContext, source: NativeMentionSource, content: string, +) { + if (!nativeMentionsEnabled()) return; + const actor = await activeActor(ctx, executor); + if (actor?.kind !== 'human') return; + const refs = extractNativeMentions(content); + if (!refs.some(ref => ref.resource_type === 'person')) { + const [state] = await executor.select({ published: nativeReferenceStates.published_person_ids }).from(nativeReferenceStates).where(and( + eq(nativeReferenceStates.org_id, ctx.orgId), eq(nativeReferenceStates.source_kind, source.kind), eq(nativeReferenceStates.source_id, source.id), + )); + if (!state?.published.length) return; + } + const contentHash = nativeContentHash(content); + await enqueue(QUEUE_NAMES.AGENT_JOBS, 'native-mention-publish', + { orgId: ctx.orgId, actorUserId: ctx.userId, source, contentHash }, + { orgId: ctx.orgId, dedupeKey: `native-publish:${source.kind}:${source.id}:${contentHash}`, executor, maxAttempts: 5 }); +} + +export async function handleNativeMentionPublication(data: { + orgId: string; actorUserId: string; source: NativeMentionSource; contentHash: string; +}) { + if (!nativeMentionsEnabled()) throw new RetryLaterJobError('Native mention publication paused', 60_000); + try { await publishNativeMentions({ orgId: data.orgId, userId: data.actorUserId }, data.source, data.contentHash); } + catch (error) { + // Removed/revoked sources and superseded explicit edits are terminal, never replayed against newer content. + if (error instanceof NativeMentionError && ['NOT_FOUND', 'FORBIDDEN', 'NATIVE_MENTION_REVISION_CONFLICT'].includes(error.code)) return; + throw error; + } +} + +export async function resolveNativeMentions(ctx: NativeMentionContext, refs: NativeMentionRef[]): Promise { + NativeMentionRefsSchema.parse(refs); + if (!(await activeActor(ctx))) return refs.map(ref => ({ ref, state: 'unavailable' })); + return Promise.all(refs.map(async ref => { + let row: Record | undefined; + if (ref.resource_type === 'person') { + row = rows(await db.execute(sql`SELECT u.name AS label, u.avatar_url, u.kind, u.title AS description + FROM users u JOIN org_members om ON om.user_id = u.id + LEFT JOIN agent_employees ae ON ae.user_id = u.id AND ae.org_id = om.org_id + WHERE om.org_id = ${ctx.orgId} AND om.is_active = true AND u.id = ${ref.resource_id} + AND (u.kind <> 'agent' OR (ae.is_active = true AND (ae.is_deleted = false OR ae.runtime_kind = 'defty_system'))) LIMIT 1`))[0]; + } else if (ref.resource_type === 'task') { + row = rows(await db.execute(sql`SELECT p.prefix || '-' || t.number || ' · ' || t.title AS label, + '/tasks?task=' || t.id AS href, t.project_id + FROM tasks t JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE t.org_id = ${ctx.orgId} AND t.id = ${ref.resource_id} AND t.is_deleted = false + AND p.is_deleted = false AND ${nativeTaskAccessSql(ctx.userId)} LIMIT 1`))[0]; + if (row && !(await employeeBoundary(ctx, { project_id: row.project_id as string, space_id: null }))) row = undefined; + } else { + const source = await loadNativeSource(ctx, { kind: 'wiki_page', id: ref.resource_id }); + if (source) row = { label: source.label, href: source.href }; + } + if (!row) return { ref, state: 'unavailable' as const }; + return { + ref, state: 'available' as const, label: String(row.label), + href: typeof row.href === 'string' ? row.href : undefined, + group: ref.resource_type === 'task' ? 'Tasks' as const : ref.resource_type === 'wiki_page' ? 'Wikis' as const : + row.kind === 'agent' ? 'Agents' as const : 'People' as const, + avatar_url: typeof row.avatar_url === 'string' ? row.avatar_url : null, + description: typeof row.description === 'string' ? row.description : null, + }; + })); +} + +export async function searchNativeMentions(ctx: NativeMentionContext, query: string) { + if (!(await activeActor(ctx))) return []; + const pattern = `%${query.replace(/[\\%_]/g, '\\$&')}%`; + const people = rows<{ id: string }>(await db.execute(sql`SELECT id FROM ( + SELECT u.id, u.kind, u.name, row_number() OVER (PARTITION BY u.kind ORDER BY u.name, u.id) AS ordinal + FROM users u JOIN org_members om ON om.user_id = u.id + WHERE om.org_id = ${ctx.orgId} AND om.is_active = true AND u.name ILIKE ${pattern} + AND u.kind IN ('human', 'agent') + AND (u.kind = 'human' OR EXISTS (SELECT 1 FROM agent_employees ae + WHERE ae.org_id = om.org_id AND ae.user_id = u.id AND ae.is_active = true + AND (ae.is_deleted = false OR ae.runtime_kind = 'defty_system'))) + ) candidates WHERE ordinal <= 8 ORDER BY kind, name LIMIT 16`)); + const tasks = rows<{ id: string }>(await db.execute(sql`SELECT t.id FROM tasks t + JOIN projects p ON p.id = t.project_id AND p.org_id = t.org_id + WHERE t.org_id = ${ctx.orgId} AND t.is_deleted = false AND p.is_deleted = false + AND ${nativeTaskAccessSql(ctx.userId)} + AND (t.title ILIKE ${pattern} OR (p.prefix || '-' || t.number) ILIKE ${pattern}) + ORDER BY t.updated_at DESC LIMIT 8`)); + const wiki = rows<{ id: string }>(await db.execute(sql`SELECT w.id FROM wiki_pages w + WHERE w.org_id = ${ctx.orgId} AND w.is_deleted = false AND w.title ILIKE ${pattern} + AND ${nativeSourceAccessSql(ctx.userId, sql`'wiki_page'`, sql`w.id`, sql`w.org_id`)} + ORDER BY w.updated_at DESC LIMIT 8`)); + return (await resolveNativeMentions(ctx, [ + ...people.map(row => nativeMentionRef('person', row.id)), + ...tasks.map(row => nativeMentionRef('task', row.id)), + ...wiki.map(row => nativeMentionRef('wiki_page', row.id)), + ])).filter(item => item.state === 'available'); +} + +export async function nativeMentionBacklinks(ctx: NativeMentionContext, target: NativeMentionRef) { + const [resolved] = await resolveNativeMentions(ctx, [target]); + if (resolved?.state !== 'available') throw new NativeMentionError('NOT_FOUND', 'Reference not found', 404); + const states = await db.select().from(nativeReferenceStates).where(and( + eq(nativeReferenceStates.org_id, ctx.orgId), eq(nativeReferenceStates.is_deleted, false), + nativeSourceAccessSql(ctx.userId, sql`${nativeReferenceStates.source_kind}`, sql`${nativeReferenceStates.source_id}`, sql`${nativeReferenceStates.org_id}`), + sql`${nativeReferenceStates.current_refs} @> ${JSON.stringify([{ resource_type: target.resource_type, resource_id: target.resource_id }])}::jsonb`, + )).orderBy(desc(nativeReferenceStates.updated_at)).limit(200); + const backlinks: Array<{ source: NativeMentionSource; label: string; href: string }> = []; + for (const state of states) { + const source = NativeMentionSourceSchema.parse({ kind: state.source_kind, id: state.source_id }); + const current = await loadNativeSource(ctx, source); + let references: NativeMentionRef[] = []; + try { if (current) references = extractNativeMentions(current.content); } catch (error) { if (!(error instanceof NativeMentionLimitError)) throw error; } + if (current && references.some(ref => nativeMentionKey(ref) === nativeMentionKey(target))) { + backlinks.push({ source, label: current.label || 'Untitled', href: current.href }); + } + } + return { backlinks, count: backlinks.length, limited: backlinks.length === 200 }; +} + +export async function deliverNativeMention(orgId: string, deliveryId: string) { + return withDbAdvisoryLock(`native-mention-deliver:${orgId}:${deliveryId}`, () => deliverNativeMentionLocked(orgId, deliveryId)); +} +async function deliverNativeMentionLocked(orgId: string, deliveryId: string) { + if (!nativeMentionsEnabled()) throw new RetryLaterJobError('Native mention publication paused', 60_000); + const delivery = await db.transaction(async tx => { + const [pending] = await tx.select().from(nativeMentionDeliveries).where(and( + eq(nativeMentionDeliveries.id, deliveryId), eq(nativeMentionDeliveries.org_id, orgId), + )).for('update'); + if (!pending || pending.status === 'suppressed') return null; + const [state] = await tx.select().from(nativeReferenceStates).where(and( + eq(nativeReferenceStates.id, pending.source_state_id), eq(nativeReferenceStates.org_id, orgId), + )); + const source = state && NativeMentionSourceSchema.parse({ kind: state.source_kind, id: state.source_id }); + const current = source && await loadNativeSource({ orgId, userId: pending.recipient_user_id }, source, tx, true); + const author = await activeActor({ orgId, userId: pending.actor_user_id }, tx); + const recipient = await activeActor({ orgId, userId: pending.recipient_user_id }, tx); + let remainsMentioned = false; + try { remainsMentioned = Boolean(current && extractNativeMentions(current.content).some(ref => + ref.resource_type === 'person' && ref.resource_id === pending.recipient_user_id)); } + catch (error) { if (!(error instanceof NativeMentionLimitError)) throw error; } + if (!state || !source || !current || !author || !recipient || !remainsMentioned) { + await tx.update(nativeMentionDeliveries).set({ status: 'suppressed', reason: 'source_or_recipient_unavailable', updated_at: new Date() }) + .where(eq(nativeMentionDeliveries.id, deliveryId)); + return null; + } + const [actor] = await tx.select({ name: users.name }).from(users).where(eq(users.id, pending.actor_user_id)); + const title = `${actor?.name ?? 'Someone'} mentioned you in ${source.kind.replaceAll('_', ' ')}`; + if (recipient.kind !== 'agent') { + const decision = await explainNotificationPolicy({ user_id: pending.recipient_user_id, type: 'mention' }, { + channel: source.kind === 'task' || source.kind === 'task_comment' ? 'tasks' : 'chat', + spaceId: source.kind === 'message' ? current.space_id : null, isMention: true, respectDnd: true, + }, tx); + if (!decision.allowed) { + await tx.update(nativeMentionDeliveries).set({ status: 'suppressed', reason: decision.reason, updated_at: new Date() }).where(eq(nativeMentionDeliveries.id, deliveryId)); + return null; + } + } + return { pending, source, current, title, agent: recipient.kind === 'agent' }; + }); + if (!delivery) return; + // Stable sourceEventId makes attention repair idempotent after an uncertain response. + const attention = await upsertAttentionItem({ + orgId, userId: delivery.pending.recipient_user_id, kind: 'mention', lane: 'needs_you', priority: 'normal', + dedupeKey: `native-mention:${deliveryId}`, sourceType: 'native_mention', sourceId: deliveryId, + sourceEventId: `native-mention:${deliveryId}`, title: delivery.title, link: delivery.current.href, + metadata: { native_mention_delivery_id: deliveryId, native_mention_source: delivery.source, passive: true }, + }, { deliver: !delivery.agent }); + // Publish the legacy notification after attention has its durable source event. + // Its normal backfill can then observe the same event without racing a second projection. + await db.transaction(async tx => { + if (!delivery.agent) await tx.insert(notifications).values({ + id: deliveryId, org_id: orgId, user_id: delivery.pending.recipient_user_id, + type: 'mention', title: delivery.title, body: null, link: delivery.current.href, + metadata: { native_mention_delivery_id: deliveryId, native_mention_source: delivery.source }, + }).onConflictDoNothing(); + await tx.update(nativeMentionDeliveries).set({ + status: 'delivered', attention_id: attention?.id ?? null, updated_at: new Date(), + }).where(and(eq(nativeMentionDeliveries.id, deliveryId), eq(nativeMentionDeliveries.org_id, orgId))); + }); + if (!delivery.agent && delivery.pending.status !== 'delivered') { + const [notification] = await db.select().from(notifications).where(eq(notifications.id, deliveryId)); + if (notification) emitToUser(delivery.pending.recipient_user_id, 'notification:new', notification); + } +} + +export async function listNativeMentionAttention(ctx: NativeMentionContext) { + if (!(await activeActor(ctx))) return []; + const items = await db.select().from(attentionItems).where(and( + eq(attentionItems.org_id, ctx.orgId), eq(attentionItems.user_id, ctx.userId), + eq(attentionItems.source_type, 'native_mention'), + sql`${attentionItems.state} IN ('open_unseen', 'open_seen', 'acknowledged')`, + nativeDeliveryAccessSql(ctx.userId, sql`${attentionItems.source_id}`, sql`${attentionItems.org_id}`), + )).orderBy(desc(attentionItems.last_event_at)).limit(50); + const visible = []; + for (const item of items) { + const metadata = item.metadata as Record; + const source = NativeMentionSourceSchema.safeParse(metadata.native_mention_source); + if (source.success && await loadNativeSource(ctx, source.data)) visible.push(item); + } + return visible; +} diff --git a/apps/api/src/lib/notification-policy.ts b/apps/api/src/lib/notification-policy.ts index 03877aae..be321c13 100644 --- a/apps/api/src/lib/notification-policy.ts +++ b/apps/api/src/lib/notification-policy.ts @@ -103,13 +103,14 @@ function spaceLevelAllowsNotification( export async function explainNotificationPolicy( values: Pick, options: NotificationPolicyOptions = {}, + executor: Pick = db, ): Promise { const inferredChannel = options.channel === undefined ? notificationChannelForType(String(values.type)) : options.channel; - const [recipient] = await db + const [recipient] = await executor .select({ notification_preferences: users.notification_preferences, status_text: users.status_text, @@ -134,7 +135,7 @@ export async function explainNotificationPolicy( } if (options.spaceId) { - const [membership] = await db + const [membership] = await executor .select({ is_muted: spaceMembers.is_muted, notification_level: spaceMembers.notification_level, diff --git a/apps/api/src/routes/inbox.ts b/apps/api/src/routes/inbox.ts index 30c8ebf8..d6b2aa8f 100644 --- a/apps/api/src/routes/inbox.ts +++ b/apps/api/src/routes/inbox.ts @@ -1,3 +1,4 @@ +import { nativeNotificationAccessSql } from '../lib/native-mention-visibility.js'; // apps/api/src/routes/inbox.ts import { Hono } from 'hono'; import { eq, and, desc, sql, lt, gt, inArray } from 'drizzle-orm'; @@ -204,6 +205,7 @@ inboxRoutes.get('/', async (c) => { const notificationWhere = and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), includeRead ? sql`TRUE` : eq(notifications.is_read, false), notificationTypes.length > 0 ? inArray(notifications.type, notificationTypes) : sql`FALSE`, cursor ? lt(notifications.created_at, new Date(cursor)) : sql`TRUE`, @@ -211,6 +213,7 @@ inboxRoutes.get('/', async (c) => { const unreadNotificationWhere = and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), eq(notifications.is_read, false), notificationTypes.length > 0 ? inArray(notifications.type, notificationTypes) : sql`FALSE`, ); @@ -389,6 +392,7 @@ inboxRoutes.post('/read', async (c) => { .where(and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), eq(notifications.is_read, false), inArray(notifications.type, notificationTypes), )) @@ -413,6 +417,7 @@ inboxRoutes.post('/read', async (c) => { inArray(notifications.id, notifIds), eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), )) .returning({ id: notifications.id }); diff --git a/apps/api/src/routes/messages.ts b/apps/api/src/routes/messages.ts index 0ccd9ab7..cef428fa 100644 --- a/apps/api/src/routes/messages.ts +++ b/apps/api/src/routes/messages.ts @@ -6,6 +6,8 @@ import { nativeCreate, nativeCreateKey, NativeCreateError } from '../lib/native- import { messages, users, reactions, spaces, spaceMembers, orgs, threadReads, messageVersions, agentEmployees, userGroups, userGroupMembers, orgMembers, files, messageAttachments as messageAttachmentLinks } from '@deft/db/schema'; import { getIO, emitToUser } from '../socket.js'; import { parseMentions } from '../lib/mentions.js'; +import { extractNativeMentions, NativeMentionLimitError } from '@deft/shared'; +import { nativeMentionsEnabled, enqueueNativeMentionPublication } from '../lib/native-mentions.js'; import { fetchLinkPreview, extractUrls, type LinkPreview } from '../lib/link-preview.js'; import { enqueue, QUEUE_NAMES } from '../lib/queues.js'; import { resolveReasonProvider } from '../lib/org-ai-config.js'; @@ -475,6 +477,8 @@ messageRoutes.post('/:spaceId', async (c) => { parent_id: parsed.data.parent_id, }).returning(); if (!insertedMessage) throw new Error('Message insert returned no row'); + await enqueueNativeMentionPublication(tx, { orgId: user.org_id, userId: user.id }, + { kind: 'message', id: insertedMessage.id }, normalizedContent); if (attachmentIds.length === 0) { return insertedMessage; @@ -573,7 +577,10 @@ messageRoutes.post('/:spaceId', async (c) => { ...groupMentionedUserIds, ])); + const nativeRecipients = new Set(nativeMentionsEnabled() ? extractNativeMentions(normalizedContent) + .filter(ref => ref.resource_type === 'person').map(ref => ref.resource_id) : []); for (const mentionedUserId of mentionedUserIds) { + if (nativeRecipients.has(mentionedUserId)) continue; // The durable native publication owns this alert. // Don't notify the sender if (mentionedUserId === user.id) continue; @@ -846,6 +853,7 @@ messageRoutes.post('/:spaceId', async (c) => { return c.json({ error: err.message, code: 'ATTACHMENT_NOT_FOUND' }, 404); } if (err instanceof NativeCreateError) return c.json({ error: err.message, code: err.code }, err.status); + if (err instanceof NativeMentionLimitError) return c.json({ error: err.message, code: 'NATIVE_MENTION_LIMIT' }, 400); console.error('Failed to send message:', err); return c.json({ error: 'Failed to send message', code: 'INTERNAL_ERROR' }, 500); } @@ -856,10 +864,15 @@ messageRoutes.patch('/:id', async (c) => { const user = c.get('user'); const messageId = c.req.param('id'); const body = await c.req.json(); - const { content } = body; - - if (!content) { - return c.json({ error: 'Content required', code: 'VALIDATION_ERROR' }, 400); + const edit = z.object({ content: z.string().min(1).max(100_000) }).safeParse(body); + if (!edit.success) return c.json({ error: 'Valid message content required', code: 'VALIDATION_ERROR' }, 400); + const { content } = edit.data; + if (nativeMentionsEnabled()) { + try { extractNativeMentions(content); } + catch (error) { + if (error instanceof NativeMentionLimitError) return c.json({ error: error.message, code: 'NATIVE_MENTION_LIMIT' }, 400); + throw error; + } } const existing = await getVisibleMessage(messageId, user.org_id, user.id); @@ -877,10 +890,13 @@ messageRoutes.patch('/:id', async (c) => { edited_at: existing.edited_at || existing.created_at, }); - const [updated] = await db.update(messages) - .set({ content, edited_at: new Date() }) - .where(eq(messages.id, messageId)) - .returning(); + const updated = await db.transaction(async tx => { + const [row] = await tx.update(messages).set({ content, edited_at: new Date() }) + .where(and(eq(messages.id, messageId), eq(messages.org_id, user.org_id))).returning(); + await enqueueNativeMentionPublication(tx, { orgId: user.org_id, userId: user.id }, + { kind: 'message', id: messageId }, content); + return row; + }); const io = getIO(); if (io) { diff --git a/apps/api/src/routes/native-mentions.ts b/apps/api/src/routes/native-mentions.ts new file mode 100644 index 00000000..4fd465e9 --- /dev/null +++ b/apps/api/src/routes/native-mentions.ts @@ -0,0 +1,49 @@ +import { Hono } from 'hono'; +import { z } from 'zod'; +import { NativeMentionRefSchema, NativeMentionRefsSchema, NativeMentionSourceSchema } from '@deft/shared'; +import { + nativeMentionsEnabled, searchNativeMentions, resolveNativeMentions, nativeMentionBacklinks, + publishNativeMentions, loadNativeSource, nativeContentHash, NativeMentionError, +} from '../lib/native-mentions.js'; + +export const nativeMentionRoutes = new Hono(); +nativeMentionRoutes.get('/capabilities', c => c.json({ enabled: nativeMentionsEnabled() })); +nativeMentionRoutes.get('/search', async c => { + const query = z.string().max(120).safeParse(c.req.query('q') ?? ''); + if (!query.success) return c.json({ error: 'Invalid search query', code: 'VALIDATION_ERROR' }, 400); + if (!nativeMentionsEnabled()) return c.json({ items: [] }); + const user = c.get('user'); + return c.json({ items: await searchNativeMentions({ orgId: user.org_id, userId: user.id }, query.data) }); +}); +nativeMentionRoutes.post('/resolve', async c => { + const parsed = z.strictObject({ refs: NativeMentionRefsSchema }).safeParse(await c.req.json().catch(() => null)); + if (!parsed.success) return c.json({ error: 'Invalid references', code: 'VALIDATION_ERROR' }, 400); + const user = c.get('user'); + return c.json({ items: await resolveNativeMentions({ orgId: user.org_id, userId: user.id }, parsed.data.refs) }); +}); +nativeMentionRoutes.post('/backlinks', async c => { + const parsed = z.strictObject({ ref: NativeMentionRefSchema }).safeParse(await c.req.json().catch(() => null)); + if (!parsed.success) return c.json({ error: 'Invalid reference', code: 'VALIDATION_ERROR' }, 400); + const user = c.get('user'); + try { return c.json(await nativeMentionBacklinks({ orgId: user.org_id, userId: user.id }, parsed.data.ref)); } + catch (error) { if (error instanceof NativeMentionError) return c.json({ error: error.message, code: error.code }, error.status); throw error; } +}); +nativeMentionRoutes.post('/source', async c => { + const parsed = NativeMentionSourceSchema.safeParse(await c.req.json().catch(() => null)); + if (!parsed.success) return c.json({ error: 'Invalid source', code: 'VALIDATION_ERROR' }, 400); + const user = c.get('user'); + const source = await loadNativeSource({ orgId: user.org_id, userId: user.id }, parsed.data); + if (!source) return c.json({ error: 'Source not found', code: 'NOT_FOUND' }, 404); + return c.json({ source: parsed.data, content_hash: nativeContentHash(source.content) }); +}); +nativeMentionRoutes.post('/publish', async c => { + const parsed = z.strictObject({ + source: NativeMentionSourceSchema, content_hash: z.string().regex(/^[a-f0-9]{64}$/), + }).safeParse(await c.req.json().catch(() => null)); + if (!parsed.success) return c.json({ error: 'Invalid publication', code: 'VALIDATION_ERROR' }, 400); + const user = c.get('user'); + try { return c.json(await publishNativeMentions( + { orgId: user.org_id, userId: user.id }, parsed.data.source, parsed.data.content_hash, + )); } + catch (error) { if (error instanceof NativeMentionError) return c.json({ error: error.message, code: error.code }, error.status); throw error; } +}); diff --git a/apps/api/src/routes/notifications.ts b/apps/api/src/routes/notifications.ts index c857cb6b..6d2f87e8 100644 --- a/apps/api/src/routes/notifications.ts +++ b/apps/api/src/routes/notifications.ts @@ -2,6 +2,7 @@ import { Hono } from 'hono'; import { eq, and, desc, sql } from 'drizzle-orm'; import { db } from '../lib/db.js'; import { notifications } from '@deft/db/schema'; +import { nativeNotificationAccessSql } from '../lib/native-mention-visibility.js'; export const notificationRoutes = new Hono(); @@ -16,6 +17,7 @@ notificationRoutes.get('/', async (c) => { and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), ) ) .orderBy(desc(notifications.created_at)) @@ -29,6 +31,7 @@ notificationRoutes.get('/', async (c) => { and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), eq(notifications.is_read, false), ) ); @@ -55,6 +58,8 @@ notificationRoutes.patch('/:id/read', async (c) => { and( eq(notifications.id, notificationId), eq(notifications.user_id, user.id), + eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), ) ) .limit(1); @@ -86,6 +91,7 @@ notificationRoutes.post('/read-all', async (c) => { and( eq(notifications.user_id, user.id), eq(notifications.org_id, user.org_id), + nativeNotificationAccessSql(user.id, sql`${notifications.metadata}`, sql`${notifications.org_id}`), eq(notifications.is_read, false), ) ); diff --git a/apps/api/src/routes/tasks.ts b/apps/api/src/routes/tasks.ts index 91ca9b44..2bf0f001 100644 --- a/apps/api/src/routes/tasks.ts +++ b/apps/api/src/routes/tasks.ts @@ -1,4 +1,6 @@ import { Hono } from 'hono'; +import { stripNativeMentionAtoms, NativeMentionLimitError } from '@deft/shared'; +import { enqueueNativeMentionPublication } from '../lib/native-mentions.js'; import { z } from 'zod'; import { eq, and, desc, asc, sql, inArray, ilike, or, isNull, type SQL } from 'drizzle-orm'; import { db } from '../lib/db.js'; @@ -181,6 +183,7 @@ async function getVisibleTaskForOrg(taskId: string, orgId: string, userId: strin */ async function resolveMentions(content: string | null | undefined, orgId: string, authorId: string): Promise { if (!content) return []; + content = stripNativeMentionAtoms(content); // Strip HTML tags so `@name` mentions inside TipTap paragraph markup still // match the raw regex. TipTap wraps content in

/ / etc. const plain = toPlainText(content); @@ -1397,12 +1400,14 @@ taskRoutes.post('/:id/comments', async (c) => { return c.json({ error: 'Task not found', code: 'NOT_FOUND' }, 404); } - const [comment] = await db.insert(taskComments).values({ - org_id: user.org_id, - task_id: taskId, - user_id: user.id, - content: parsed.data.content, - }).returning(); + const comment = await db.transaction(async tx => { + const [inserted] = await tx.insert(taskComments).values({ + org_id: user.org_id, task_id: taskId, user_id: user.id, content: parsed.data.content, + }).returning(); + await enqueueNativeMentionPublication(tx, { orgId: user.org_id, userId: user.id }, + { kind: 'task_comment', id: inserted!.id }, parsed.data.content); + return inserted; + }); // Create activity log entry await db.insert(taskActivity).values({ @@ -1482,6 +1487,7 @@ taskRoutes.post('/:id/comments', async (c) => { user_avatar: userData?.avatar_url ?? null, }, 201); } catch (err) { + if (err instanceof NativeMentionLimitError) return c.json({ error: err.message, code: 'NATIVE_MENTION_LIMIT' }, 400); console.error('Failed to create task comment:', err); return c.json({ error: 'Failed to create comment', code: 'INTERNAL_ERROR' }, 500); } diff --git a/apps/api/src/workers/handlers/native-mentions.ts b/apps/api/src/workers/handlers/native-mentions.ts new file mode 100644 index 00000000..b6f83b81 --- /dev/null +++ b/apps/api/src/workers/handlers/native-mentions.ts @@ -0,0 +1,18 @@ +import { z } from 'zod'; +import { NativeMentionSourceSchema } from '@deft/shared'; +import { handleNativeMentionReconciliation, handleNativeMentionPublication, deliverNativeMention } from '../../lib/native-mentions.js'; +import type { JobHandler } from '../types.js'; + +export const handleNativeMentionReconcile: JobHandler = async job => { + const data = z.object({ orgId: z.string().min(1), source: NativeMentionSourceSchema, publishOnCreate: z.boolean().optional() }).parse(job.data); + await handleNativeMentionReconciliation(data); +}; +export const handleNativeMentionDeliver: JobHandler = async job => { + const data = z.object({ orgId: z.string().min(1), deliveryId: z.string().min(1) }).parse(job.data); + await deliverNativeMention(data.orgId, data.deliveryId); +}; +export const handleNativeMentionPublish: JobHandler = async job => { + const data = z.object({ orgId: z.string().min(1), actorUserId: z.string().min(1), + source: NativeMentionSourceSchema, contentHash: z.string().regex(/^[a-f0-9]{64}$/) }).parse(job.data); + await handleNativeMentionPublication(data); +}; diff --git a/apps/api/src/workers/index.ts b/apps/api/src/workers/index.ts index 3452234a..d012f901 100644 --- a/apps/api/src/workers/index.ts +++ b/apps/api/src/workers/index.ts @@ -201,6 +201,18 @@ async function getAgentJobHandler(jobName: string): Promise { const mod = await import('./handlers/cross-reference.js'); return mod.handleCrossReference; } + case 'native-mention-reconcile': { + const mod = await import('./handlers/native-mentions.js'); + return mod.handleNativeMentionReconcile; + } + case 'native-mention-publish': { + const mod = await import('./handlers/native-mentions.js'); + return mod.handleNativeMentionPublish; + } + case 'native-mention-deliver': { + const mod = await import('./handlers/native-mentions.js'); + return mod.handleNativeMentionDeliver; + } case 'embed-content': { const mod = await import('./handlers/embed-content.js'); return mod.handleEmbedContent; diff --git a/apps/api/test/agent-mention-normalization.test.ts b/apps/api/test/agent-mention-normalization.test.ts index 41a719cd..734776fb 100644 --- a/apps/api/test/agent-mention-normalization.test.ts +++ b/apps/api/test/agent-mention-normalization.test.ts @@ -2,6 +2,12 @@ import assert from 'node:assert/strict'; import { test } from 'node:test'; import { normalizePlainAgentMentions } from '../src/lib/agent-mention-normalization.js'; +test('native identity placeholders never become fuzzy agent handles', () => { + const content = '

@Person

'; + const result = normalizePlainAgentMentions(content, [{ userId: 'wrong-agent', name: 'Person', slug: 'person' }]); + assert.equal(result.content, content); + assert.deepEqual(result.resolvedUserIds, []); +}); test('resolves exact typed agent names and slugs into structured mentions', () => { const agents = [{ userId: 'rita-user', name: 'Rita', slug: 'research-agent' }]; diff --git a/apps/api/test/fixtures/native-mentions.ts b/apps/api/test/fixtures/native-mentions.ts new file mode 100644 index 00000000..9cdce9f2 --- /dev/null +++ b/apps/api/test/fixtures/native-mentions.ts @@ -0,0 +1,45 @@ +import { randomUUID } from 'node:crypto'; +import { db } from '../../src/lib/db.js'; +import { orgs, users, orgMembers, spaces, spaceMembers, projects, tasks, wikiPages, notes, agentEmployees } from '@deft/db/schema'; +export async function createNativeMentionFixture() { + const orgId = randomUUID(), otherOrgId = randomUUID(); + const ownerId = randomUUID(), samId = randomUUID(), agentId = randomUUID(), agent2Id = randomUUID(), outsiderId = randomUUID(); + const employeeId = randomUUID(), employee2Id = randomUUID(), publicSpaceId = randomUUID(), privateSpaceId = randomUUID(); + const projectId = randomUUID(), taskId = randomUUID(), restrictedId = randomUUID(), wikiId = randomUUID(), privateWikiId = randomUUID(), noteId = randomUUID(); + await db.insert(orgs).values([{ id: orgId, name: 'Native mention lab', slug: 'mention-lab-' + orgId }, + { id: otherOrgId, name: 'Other workspace', slug: 'other-' + otherOrgId }]); + await db.insert(users).values([ + { id: ownerId, name: 'Jordan', email: 'jordan-' + ownerId + '@example.test', kind: 'human' }, + { id: samId, name: 'Sam', email: 'sam-' + samId + '@example.test', kind: 'human' }, + { id: agentId, name: 'Rita Research', email: 'rita-' + agentId + '@example.test', kind: 'agent', is_agent: true }, + { id: agent2Id, name: 'Avery Review', email: 'avery-' + agent2Id + '@example.test', kind: 'agent', is_agent: true }, + { id: outsiderId, name: 'Other workspace user', email: 'other-' + outsiderId + '@example.test', kind: 'human' }, + ]); + await db.insert(orgMembers).values([ + { org_id: orgId, user_id: ownerId, role: 'owner' }, { org_id: orgId, user_id: samId, role: 'member' }, + { org_id: orgId, user_id: agentId, role: 'member' }, { org_id: orgId, user_id: agent2Id, role: 'member' }, + { org_id: otherOrgId, user_id: outsiderId, role: 'owner' }, + ]); + await db.insert(agentEmployees).values([ + { id: employeeId, org_id: orgId, user_id: agentId, name: 'Rita Research', slug: 'rita-' + employeeId, + role: 'custom', system_prompt: 'Synthetic offline agent for mention validation.', created_by: ownerId, is_byoa: true, runtime_kind: 'custom_mcp' }, + { id: employee2Id, org_id: orgId, user_id: agent2Id, name: 'Avery Review', slug: 'avery-' + employee2Id, + role: 'custom', system_prompt: 'Synthetic offline agent for mention validation.', created_by: ownerId, is_byoa: true, runtime_kind: 'custom_mcp' }, + ]); + await db.insert(spaces).values([{ id: publicSpaceId, org_id: orgId, name: 'Launch room', type: 'public', created_by: ownerId }, + { id: privateSpaceId, org_id: orgId, name: 'Private planning', type: 'private', created_by: ownerId }]); + await db.insert(spaceMembers).values([ownerId, samId, agentId, agent2Id].map(user_id => ({ space_id: publicSpaceId, user_id })) + .concat([{ space_id: privateSpaceId, user_id: ownerId }])); + await db.insert(projects).values({ id: projectId, org_id: orgId, name: 'Launch', prefix: 'DEFT', lead_id: ownerId, task_counter: 43 }); + await db.insert(tasks).values([ + { id: taskId, org_id: orgId, project_id: projectId, number: 42, title: 'Review release', created_by: ownerId }, + { id: restrictedId, org_id: orgId, project_id: projectId, number: 43, title: 'Private budget', created_by: ownerId, metadata: { visibility: 'restricted' } }, + ]); + await db.insert(wikiPages).values([ + { id: wikiId, org_id: orgId, title: 'Launch checklist', slug: 'launch-checklist-' + wikiId, content: 'Launch checklist', type: 'procedure', scope: 'org', user_id: ownerId }, + { id: privateWikiId, org_id: orgId, title: 'Private notes', slug: 'private-' + privateWikiId, content: 'Private', type: 'procedure', scope: 'user', user_id: ownerId }, + ]); + await db.insert(notes).values({ id: noteId, org_id: orgId, user_id: ownerId, title: 'Daily launch notes', content: '

Daily plan

', visibility: 'private' }); + return { orgId, otherOrgId, ownerId, samId, agentId, agent2Id, outsiderId, employeeId, employee2Id, + publicSpaceId, privateSpaceId, projectId, taskId, restrictedId, wikiId, privateWikiId, noteId }; +} diff --git a/apps/api/test/fixtures/seed-native-mention-evidence.ts b/apps/api/test/fixtures/seed-native-mention-evidence.ts new file mode 100644 index 00000000..bf02ac68 --- /dev/null +++ b/apps/api/test/fixtures/seed-native-mention-evidence.ts @@ -0,0 +1,29 @@ +import { writeFile } from 'node:fs/promises'; +import { and, eq } from 'drizzle-orm'; +import bcrypt from 'bcryptjs'; +import { users, orgs, orgMembers, onboardingState } from '@deft/db/schema'; +import { db, closeDb } from '../../src/lib/db.js'; +import { createWebSession } from '../../src/lib/web-sessions.js'; +import { createNativeMentionFixture } from './native-mentions.js'; +import { safeTestDatabaseUrl } from './safe-test-database.js'; +if (!safeTestDatabaseUrl()) throw new Error('Evidence fixtures require an explicitly disposable test database'); +const output = process.env.DEFT_MENTION_FIXTURE_PATH; +if (!output) throw new Error('DEFT_MENTION_FIXTURE_PATH is required'); +try { + const fixture = await createNativeMentionFixture(); + // The browser lab models the supported one-workspace deployment. + await db.delete(orgMembers).where(eq(orgMembers.org_id, fixture.otherOrgId)); + await db.delete(orgs).where(eq(orgs.id, fixture.otherOrgId)); + const password = 'mentions-demo-only'; + await db.update(users).set({ password_hash: await bcrypt.hash(password, 10), email_verified: true }).where(eq(users.id, fixture.ownerId)); + await db.update(users).set({ password_hash: await bcrypt.hash(password, 10), email_verified: true }).where(eq(users.id, fixture.samId)); + await db.update(orgs).set({ settings: { onboarding_completed: true } }).where(eq(orgs.id, fixture.orgId)); + await db.insert(onboardingState).values([{ user_id: fixture.ownerId, completed: true }, { user_id: fixture.samId, completed: true }]); + const [owner] = await db.select().from(users).where(eq(users.id, fixture.ownerId)); + const [sam] = await db.select().from(users).where(eq(users.id, fixture.samId)); + const ownerSession = await createWebSession({ id: owner!.id, org_id: fixture.orgId, email: owner!.email! }); + const samSession = await createWebSession({ id: sam!.id, org_id: fixture.orgId, email: sam!.email! }); + await writeFile(output, JSON.stringify({ ...fixture, ownerEmail: owner!.email, samEmail: sam!.email, password, + ownerSession, samSession, wikiSlug: 'launch-checklist-' + fixture.wikiId }, null, 2)); + console.log('Synthetic native mention evidence fixture saved; no real workspace used.'); +} finally { await closeDb(); } diff --git a/apps/api/test/native-mentions-db.test.ts b/apps/api/test/native-mentions-db.test.ts new file mode 100644 index 00000000..90d9adc9 --- /dev/null +++ b/apps/api/test/native-mentions-db.test.ts @@ -0,0 +1,217 @@ +import assert from 'node:assert/strict'; +import { after, before, test } from 'node:test'; +import { and, eq, sql } from 'drizzle-orm'; +import { Hono } from 'hono'; +import { nativeMentionRef, nativeMentionToken, type NativeMentionSource } from '@deft/shared'; +import { messages, taskComments, tasks, notes, wikiPages, users, spaceMembers, noteShares, + nativeReferenceStates, nativeMentionDeliveries, notifications, attentionItems, agentChannelEvents, taskWatchers, jobQueue } from '@deft/db/schema'; +import { db, closeDb } from '../src/lib/db.js'; +import { reconcileNativeMentions, handleNativeMentionReconciliation, publishNativeMentions, + nativeContentHash, deliverNativeMention, resolveNativeMentions, nativeMentionBacklinks, loadNativeSource } from '../src/lib/native-mentions.js'; +import { enqueueNativeMentionPublication, handleNativeMentionPublication } from '../src/lib/native-mentions.js'; +import { nativeMentionRoutes } from '../src/routes/native-mentions.js'; +import { notificationRoutes } from '../src/routes/notifications.js'; +import { boundMentionAttention, mentionAttentionAcknowledge } from '../src/lib/mcp-tools/mention-attention.js'; +import { createNativeMentionFixture } from './fixtures/native-mentions.js'; +import { safeTestDatabaseUrl } from './fixtures/safe-test-database.js'; +const enabled = Boolean(safeTestDatabaseUrl()); +let fixture: Awaited>; +before(async () => { if (enabled) { process.env.DEFT_NATIVE_MENTIONS_ENABLED = 'true'; fixture = await createNativeMentionFixture(); } }); +after(closeDb); +const token = (kind: 'person' | 'task' | 'wiki_page', id: string) => nativeMentionToken(nativeMentionRef(kind, id)); +const actor = () => ({ orgId: fixture.orgId, userId: fixture.ownerId }); +const appFor = (userId: string) => { + const app = new Hono(); + app.use('*', async (c, next) => { c.set('user', { id: userId, org_id: fixture.orgId, email: 'synthetic@example.test' }); await next(); }); + app.route('/mentions', nativeMentionRoutes); app.route('/notifications', notificationRoutes); + return app; +}; + +test('all native writers enqueue identity-only reconciliation; stale jobs use current content', { skip: !enabled }, async () => { + const body = token('task', fixture.taskId) + ' ' + token('wiki_page', fixture.wikiId); + const [message] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }).returning(); + const [reply] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, parent_id: message!.id, content: body }).returning(); + const [comment] = await db.insert(taskComments).values({ org_id: fixture.orgId, task_id: fixture.taskId, user_id: fixture.ownerId, content: body }).returning(); + await db.update(tasks).set({ description: body }).where(eq(tasks.id, fixture.taskId)); + await db.update(wikiPages).set({ content: body }).where(eq(wikiPages.id, fixture.wikiId)); + await db.update(notes).set({ content: body }).where(eq(notes.id, fixture.noteId)); + const sources: NativeMentionSource[] = [ + { kind: 'message', id: message!.id }, { kind: 'message', id: reply!.id }, { kind: 'task_comment', id: comment!.id }, + { kind: 'task', id: fixture.taskId }, { kind: 'wiki_page', id: fixture.wikiId }, { kind: 'note', id: fixture.noteId }, + ]; + const jobs = await db.select().from(jobQueue).where(and(eq(jobQueue.org_id, fixture.orgId), eq(jobQueue.name, 'native-mention-reconcile'))); + for (const source of sources) { + assert(jobs.some(job => (job.data.source as NativeMentionSource)?.id === source.id)); + await reconcileNativeMentions(fixture.orgId, source); + } + assert(!JSON.stringify(jobs.map(job => job.data)).includes('Launch checklist')); + const links = await nativeMentionBacklinks(actor(), nativeMentionRef('task', fixture.taskId)); + assert.equal(links.count, 6); + assert.equal((await nativeMentionBacklinks({ orgId: fixture.orgId, userId: fixture.samId }, nativeMentionRef('task', fixture.taskId))).count, 5); + await db.update(messages).set({ content: 'Reference removed' }).where(eq(messages.id, message!.id)); + await reconcileNativeMentions(fixture.orgId, { kind: 'message', id: message!.id }); + assert.equal((await nativeMentionBacklinks(actor(), nativeMentionRef('task', fixture.taskId))).count, 5); + await db.update(messages).set({ is_deleted: true }).where(eq(messages.id, reply!.id)); + await reconcileNativeMentions(fixture.orgId, { kind: 'message', id: reply!.id }); + assert.equal((await nativeMentionBacklinks(actor(), nativeMentionRef('task', fixture.taskId))).count, 4); +}); + +test('publication uses saved hashes, blocks private recipients and retries after sharing', { skip: !enabled }, async () => { + const body = token('person', fixture.samId); + const source = { kind: 'note' as const, id: fixture.noteId }; + await db.update(notes).set({ content: body }).where(eq(notes.id, source.id)); + await reconcileNativeMentions(fixture.orgId, source); + assert.equal((await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.org_id, fixture.orgId))).length, 0); + await assert.rejects(publishNativeMentions(actor(), source, nativeContentHash('stale')), /Content changed/); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash(body))).blocked_count, 1); + await db.insert(noteShares).values({ note_id: source.id, shared_with_user_id: fixture.samId }); + const result = await publishNativeMentions(actor(), source, nativeContentHash(body)); + assert.equal(result.queued_count, 1); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash(body))).queued_count, 0); + const [delivery] = await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.org_id, fixture.orgId)); + await deliverNativeMention(fixture.orgId, delivery!.id); + await deliverNativeMention(fixture.orgId, delivery!.id); + assert.equal((await db.select().from(notifications).where(eq(notifications.id, delivery!.id))).length, 1); + const [attention] = await db.select().from(attentionItems).where(eq(attentionItems.source_id, delivery!.id)); + assert.equal(attention!.event_count, 1); + await db.delete(noteShares).where(eq(noteShares.note_id, source.id)); + const response = await appFor(fixture.samId).request('/notifications'); + const data = await response.json(); + assert.equal(data.unread_count, 0); + assert.equal(data.notifications.length, 0); + assert.equal((await appFor(fixture.samId).request('/notifications/' + delivery!.id + '/read', { method: 'PATCH' })).status, 404); +}); + +test('resolver fails closed for restricted, cross-tenant, deleted and watcher-only targets; names stay live', { skip: !enabled }, async () => { + await db.insert(taskWatchers).values({ task_id: fixture.restrictedId, user_id: fixture.samId }); + const ctx = { orgId: fixture.orgId, userId: fixture.samId }; + const denied = await resolveNativeMentions(ctx, [nativeMentionRef('task', fixture.restrictedId), + nativeMentionRef('wiki_page', fixture.privateWikiId), nativeMentionRef('person', fixture.outsiderId)]); + assert(denied.every(item => item.state === 'unavailable' && item.label === undefined && item.href === undefined)); + await db.update(users).set({ name: 'Sam Updated' }).where(eq(users.id, fixture.samId)); + assert.equal((await resolveNativeMentions(actor(), [nativeMentionRef('person', fixture.samId)]))[0]!.label, 'Sam Updated'); + await db.update(wikiPages).set({ is_deleted: true }).where(eq(wikiPages.id, fixture.privateWikiId)); + assert.equal((await resolveNativeMentions(actor(), [nativeMentionRef('wiki_page', fixture.privateWikiId)]))[0]!.state, 'unavailable'); + const invalid = await appFor(fixture.ownerId).request('/mentions/publish', { method: 'POST', headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ source: { kind: 'note', id: fixture.noteId, org_id: fixture.otherOrgId }, content_hash: 'a'.repeat(64) }) }); + assert.equal(invalid.status, 400); +}); + +test('agent document mentions create two isolated passive feeds and never an execution channel event', { skip: !enabled }, async () => { + const body = token('person', fixture.agentId) + ' ' + token('person', fixture.agent2Id); + const source = { kind: 'task' as const, id: fixture.taskId }; + await db.update(tasks).set({ description: body }).where(eq(tasks.id, source.id)); + await handleNativeMentionReconciliation({ orgId: fixture.orgId, source }); // Autosave remains passive. + const result = await publishNativeMentions(actor(), source, nativeContentHash(body)); + assert.equal(result.queued_count, 2); + const deliveries = await db.select().from(nativeMentionDeliveries).where(and(eq(nativeMentionDeliveries.org_id, fixture.orgId), + sql`${nativeMentionDeliveries.recipient_user_id} IN (${fixture.agentId}, ${fixture.agent2Id})`)); + for (const delivery of deliveries) await deliverNativeMention(fixture.orgId, delivery.id); + const ctx = { org_id: fixture.orgId, employee_id: fixture.employeeId, employee_slug: 'bound', trust_level: 'conservative' as const, + token_id: 'synthetic', scopes: ['read:workspace', 'read:tasks', 'write:workspace'] }; + const own = await boundMentionAttention(ctx); + assert.equal(own.length, 1); + assert.equal(own[0]!.user_id, fixture.agentId); + assert.equal((await boundMentionAttention({ ...ctx, scopes: ['read:workspace'] })).length, 0); + const second = await boundMentionAttention({ ...ctx, employee_id: fixture.employee2Id }); + assert.equal(second.length, 1); + assert.equal((await mentionAttentionAcknowledge({ attention_id: second[0]!.id }, ctx)).isError, true); + assert.equal((await mentionAttentionAcknowledge({ attention_id: own[0]!.id }, ctx)).isError, false); + assert.equal((await boundMentionAttention(ctx)).length, 0); + assert.equal((await db.select().from(agentChannelEvents).where(eq(agentChannelEvents.org_id, fixture.orgId))).length, 0); + assert.equal(await loadNativeSource({ orgId: fixture.otherOrgId, userId: fixture.outsiderId }, source), null); +}); +test('human send intent commits with content; direct and governed writes cannot infer intent from the content author', { skip: !enabled }, async () => { + const body = token('person', fixture.samId); + const before = (await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.org_id, fixture.orgId))).length; + const [direct] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }).returning(); + await handleNativeMentionReconciliation({ orgId: fixture.orgId, source: { kind: 'message', id: direct!.id }, publishOnCreate: true }); + assert.equal((await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.org_id, fixture.orgId))).length, before); + const rolledBackId = crypto.randomUUID(); + await assert.rejects(db.transaction(async tx => { + await tx.insert(messages).values({ id: rolledBackId, org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }); + await enqueueNativeMentionPublication(tx, actor(), { kind: 'message', id: rolledBackId }, body); + throw new Error('Rollback proof'); + }), /Rollback proof/); + assert.equal((await db.select().from(messages).where(eq(messages.id, rolledBackId))).length, 0); + assert.equal((await db.select().from(jobQueue).where(sql`${jobQueue.data}->'source'->>'id' = ${rolledBackId}`)).length, 0); + const committed = await db.transaction(async tx => { + const [row] = await tx.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }).returning(); + await enqueueNativeMentionPublication(tx, actor(), { kind: 'message', id: row!.id }, body); + return row!; + }); + const [job] = await db.select().from(jobQueue).where(and(eq(jobQueue.name, 'native-mention-publish'), sql`${jobQueue.data}->'source'->>'id' = ${committed.id}`)); + assert(job); + assert.equal(job!.data.actorUserId, fixture.ownerId); + await handleNativeMentionPublication(job!.data as Parameters[0]); + assert.equal((await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.org_id, fixture.orgId))).length, before + 1); + await assert.rejects(publishNativeMentions({ orgId: fixture.orgId, userId: fixture.agentId }, { kind: 'task', id: fixture.taskId }, + nativeContentHash(token('person', fixture.agentId))), /Human publication intent/); +}); +test('unfinished removal and re-addition do not ping; separately published removal and re-addition do', { skip: !enabled }, async () => { + const source = { kind: 'task' as const, id: fixture.taskId }; + const body = token('person', fixture.samId); + await db.update(tasks).set({ description: body }).where(eq(tasks.id, source.id)); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash(body))).queued_count, 1); + await db.update(tasks).set({ description: 'draft removal' }).where(eq(tasks.id, source.id)); + await reconcileNativeMentions(fixture.orgId, source); + await db.update(tasks).set({ description: body }).where(eq(tasks.id, source.id)); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash(body))).queued_count, 0); + await db.update(tasks).set({ description: 'published removal' }).where(eq(tasks.id, source.id)); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash('published removal'))).queued_count, 0); + await db.update(tasks).set({ description: body }).where(eq(tasks.id, source.id)); + assert.equal((await publishNativeMentions(actor(), source, nativeContentHash(body))).queued_count, 1); +}); +test('deleted sources suppress pending delivery and malformed publication JSON returns a structured client error', { skip: !enabled }, async () => { + const [message] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: token('person', fixture.samId) }).returning(); + await publishNativeMentions(actor(), { kind: 'message', id: message!.id }, nativeContentHash(message!.content)); + const [state] = await db.select().from(nativeReferenceStates).where(eq(nativeReferenceStates.source_id, message!.id)); + const [delivery] = await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.source_state_id, state!.id)); + await db.update(messages).set({ is_deleted: true }).where(eq(messages.id, message!.id)); + await deliverNativeMention(fixture.orgId, delivery!.id); + const [suppressed] = await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.id, delivery!.id)); + assert.equal(suppressed!.status, 'suppressed'); + assert.equal((await db.select().from(notifications).where(eq(notifications.id, delivery!.id))).length, 0); + const invalid = await appFor(fixture.ownerId).request('/mentions/publish', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: '{' }); + assert.equal(invalid.status, 400); + assert.equal((await invalid.json()).code, 'VALIDATION_ERROR'); +}); + +test('notification preferences and rollout pause preserve durable effects without delivery', { skip: !enabled }, async () => { + const body = token('person', fixture.samId); + const [message] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }).returning(); + const source = { kind: 'message' as const, id: message!.id }; + await publishNativeMentions(actor(), source, nativeContentHash(body)); + const [state] = await db.select().from(nativeReferenceStates).where(eq(nativeReferenceStates.source_id, source.id)); + const [delivery] = await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.source_state_id, state!.id)); + process.env.DEFT_NATIVE_MENTIONS_ENABLED = 'false'; + try { + await assert.rejects(publishNativeMentions(actor(), source, nativeContentHash(body)), /disabled/); + await assert.rejects(deliverNativeMention(fixture.orgId, delivery!.id), /paused/); + assert.equal((await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.id, delivery!.id)))[0]!.status, 'pending'); + } finally { process.env.DEFT_NATIVE_MENTIONS_ENABLED = 'true'; } + await db.update(users).set({ status_text: 'Do Not Disturb' }).where(eq(users.id, fixture.samId)); + try { + await deliverNativeMention(fixture.orgId, delivery!.id); + assert.equal((await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.id, delivery!.id)))[0]!.reason, 'do_not_disturb'); + assert.equal((await db.select().from(notifications).where(eq(notifications.id, delivery!.id))).length, 0); + } finally { await db.update(users).set({ status_text: null }).where(eq(users.id, fixture.samId)); } +}); + +test('PostgreSQL concurrent publish and retry workers create exactly one durable attention effect', { + skip: !enabled || process.env.DEFT_NATIVE_MENTION_CONCURRENCY_CERTIFY !== 'true', +}, async () => { + const body = token('person', fixture.samId); + const [message] = await db.insert(messages).values({ org_id: fixture.orgId, space_id: fixture.publicSpaceId, user_id: fixture.ownerId, content: body }).returning(); + const source = { kind: 'message' as const, id: message!.id }; + const publications = await Promise.all(Array.from({ length: 8 }, () => publishNativeMentions(actor(), source, nativeContentHash(body)))); + assert.equal(publications.reduce((sum, result) => sum + result.queued_count, 0), 1); + const [state] = await db.select().from(nativeReferenceStates).where(eq(nativeReferenceStates.source_id, source.id)); + const deliveries = await db.select().from(nativeMentionDeliveries).where(eq(nativeMentionDeliveries.source_state_id, state!.id)); + assert.equal(deliveries.length, 1); + await Promise.all(Array.from({ length: 8 }, () => deliverNativeMention(fixture.orgId, deliveries[0]!.id))); + assert.equal((await db.select().from(notifications).where(eq(notifications.id, deliveries[0]!.id))).length, 1); + const attention = await db.select().from(attentionItems).where(eq(attentionItems.source_id, deliveries[0]!.id)); + assert.equal(attention.length, 1); + assert.equal(attention[0]!.event_count, 1); +}); diff --git a/apps/web/src/app/(app)/knowledge/page.tsx b/apps/web/src/app/(app)/knowledge/page.tsx index 8701a30e..68be387b 100644 --- a/apps/web/src/app/(app)/knowledge/page.tsx +++ b/apps/web/src/app/(app)/knowledge/page.tsx @@ -1,4 +1,8 @@ 'use client'; +import { NativeMentionTextarea } from '@/components/native-mention-textarea'; +import { NativeMentionMarkdown } from '@/components/native-mention-markdown'; +import { NativeMentionPublish } from '@/components/native-mention-publish'; +import { NativeMentionBacklinks } from '@/components/native-mention-backlinks'; import { useState, useEffect, useCallback, useRef, lazy, Suspense } from 'react'; import { useSearchParams } from 'next/navigation'; @@ -301,7 +305,7 @@ function CreatePageModal({ onClose, onCreated }: { onClose: () => void; onCreate {/* Content */}
-