diff --git a/agent/channels/eve.ts b/agent/channels/eve.ts index 4bcc3f23..a8948bfb 100644 --- a/agent/channels/eve.ts +++ b/agent/channels/eve.ts @@ -12,7 +12,6 @@ import { type AccessScope, } from "@shared/identity/access-scope"; import { getAuthSession } from "@db/services/auth/session"; -import { sendMessageToolResultSchema } from "@shared/chat/message-delivery"; import { finalizeScheduledReportDelivery, releaseScheduledReportDelivery, @@ -70,12 +69,6 @@ const authenticate: Parameters[1] = [ const channel = eveChannel({ auth: authenticate, events: { - async "action.result"(event, _channel, session) { - const message = sendMessageToolResultSchema.safeParse(event.result); - if (event.status === "completed" && message.success) { - await finalizeScheduledReportDelivery(session); - } - }, async "message.completed"(event, _channel, session) { if (event.finishReason === "tool-calls") return; if (scheduledReportFromSession(session)) { diff --git a/agent/channels/linq.ts b/agent/channels/linq.ts index fb639489..fd72c026 100644 --- a/agent/channels/linq.ts +++ b/agent/channels/linq.ts @@ -1,29 +1,14 @@ import type { LinqAPIV3 } from "@linqapp/sdk"; -import type { AdapterPostableMessage } from "chat"; -import { - defaultLinqAuth, - linqChannel, - type LinqChannelCredentials, -} from "eve/channels/linq"; -import { vercelOidc } from "eve/channels/auth"; +import type { Message, Thread } from "chat"; +import { createMemoryState } from "@chat-adapter/state-memory"; +import { defaultLinqAuth } from "eve/channels/linq"; +import { chatSdkChannel } from "eve/channels/chat-sdk"; import { z } from "zod"; -import { resolveLinqReplyTarget } from "@agent/lib/reply-targets"; -import { scopeFromPrincipal } from "@agent/lib/principal-scope"; import { getAuth } from "@db/services/auth"; -import { sendMessageToolResultSchema } from "@shared/chat/message-delivery"; -import { reactToMessageToolResultSchema } from "@shared/chat/reaction"; import { accessScopeForUser } from "@shared/identity/access-scope"; import { normalizeAuthPhoneNumber } from "@shared/identity/phone-number"; -import { prepareLinqImageArtifactDelivery } from "../lib/linq-image-artifact/delivery"; -import { - extractImageArtifactMarkdownReferences, - stripImageArtifactMarkdownReferences, -} from "../lib/linq-image-artifact/markdown"; -import { env } from "@shared/environment"; -import { - linqCredentials, - sendNativeLinqMessage, -} from "@agent/lib/linq/transport"; +import { linqMessageContent } from "@agent/lib/linq/content"; +import { linqAdapter, sendNativeLinqMessage } from "@agent/lib/linq/transport"; import { finalizeScheduledReportDelivery, releaseScheduledReportDelivery, @@ -34,30 +19,16 @@ const verifiedPhoneUserSchema = z.object({ id: z.string().min(1), phoneNumberVerified: z.literal(true), }); -const unavailableReplyTargetSchema = z.object({ - status: z.union([z.literal(400), z.literal(404)]), -}); - type LinqMessageContent = Parameters< LinqAPIV3["chats"]["messages"]["send"] >[1]["message"]; -const trustedForwarder = vercelOidc(); - -// The Linq adapter only rejects a webhook when the verifier returns `false`, -// while eve's OIDC verifier reports failure as `null`. Translate explicitly so -// an unverified forwarder can never reach message dispatch. -export const linqWebhookVerifier: NonNullable< - LinqChannelCredentials["webhookVerifier"] -> = async (request) => (await trustedForwarder(request)) ?? false; - -const credentials = { - ...linqCredentials, - webhookVerifier: env.LINQ_CONNECTOR ? linqWebhookVerifier : () => false, -} satisfies LinqChannelCredentials; - -export default linqChannel({ - credentials, +const linq = chatSdkChannel({ + adapters: { linq: linqAdapter }, + state: createMemoryState(), + userName: "eve", + concurrency: "concurrent", + streaming: false, events: { async "authorization.required"(event, context, session) { const { thread } = context; @@ -105,210 +76,6 @@ export default linqChannel({ idempotency_key: `${idempotencyKey}:link`, }); }, - async "action.result"(event, context, session) { - const reaction = reactToMessageToolResultSchema.safeParse(event.result); - if (event.status === "completed" && reaction.success) { - if (!context.thread) { - throw new Error( - "react_to_message requires an active Linq conversation thread." - ); - } - // The thread is a persisted snapshot; auth identifies this incoming message. - const target = resolveLinqReplyTarget( - { kind: "current" }, - session.session.auth - ); - const messageId = - target?.conversationId === context.thread.id - ? target.messageId - : undefined; - if (!messageId) { - throw new Error("react_to_message requires a current Linq message."); - } - const adapter = context.bot.getAdapter("linq"); - if (reaction.data.output.operation === "remove") { - await adapter.removeReaction( - context.thread.id, - messageId, - reaction.data.output.type - ); - } else { - await adapter.addReaction( - context.thread.id, - messageId, - reaction.data.output.type - ); - } - await finalizeScheduledReportDelivery(session); - return; - } - - const message = sendMessageToolResultSchema.safeParse(event.result); - if ( - event.status === "completed" && - message.success && - message.data.toolName === "send_message" - ) { - const { thread } = context; - if (!thread) { - throw new Error( - "send_message requires an active Linq conversation thread." - ); - } - const report = scheduledReportFromSession(session); - const replyTarget = resolveLinqReplyTarget( - message.data.output.replyTo, - session.session.auth - ); - const requestedReplyMessageId = - replyTarget?.conversationId === thread.id - ? replyTarget.messageId - : undefined; - const idempotencyKey = report - ? `scheduled-report:${report.runId}:${String(report.sequence)}` - : undefined; - const adapter = context.bot.getAdapter("linq"); - const post = idempotencyKey - ? (content: AdapterPostableMessage) => - adapter.postMessage(thread.id, content, { idempotencyKey }) - : (content: AdapterPostableMessage) => thread.post(content); - const postReply = ( - content: AdapterPostableMessage, - replyToMessageId: string - ) => { - if (idempotencyKey) { - return adapter.postMessage(thread.id, content, { - idempotencyKey, - replyToMessageId, - }); - } - return adapter.postMessage(thread.id, content, { - replyToMessageId, - }); - }; - const resolveExistingChatId = () => { - const { chatId, pendingHandle } = adapter.decodeThreadId(thread.id); - if (pendingHandle || !chatId) { - throw new Error("A Linq reply requires an existing conversation."); - } - return chatId; - }; - - if (message.data.output.kind === "link") { - const { url } = message.data.output; - const chatId = resolveExistingChatId(); - const sendLink = (replyToMessageId?: string) => { - const nativeMessage: LinqMessageContent = { - parts: [{ type: "link", value: url }], - }; - if (idempotencyKey) { - nativeMessage.idempotency_key = idempotencyKey; - } - if (replyToMessageId) { - nativeMessage.reply_to = { message_id: replyToMessageId }; - } - return sendNativeLinqMessage(chatId, nativeMessage); - }; - try { - await sendLink(requestedReplyMessageId); - } catch (error) { - if ( - !requestedReplyMessageId || - !unavailableReplyTargetSchema.safeParse(error).success - ) { - throw error; - } - console.warn("[linq] reply target is unavailable", { - sessionId: session.session.id, - }); - await sendLink(); - } - await finalizeScheduledReportDelivery(session); - return; - } - - const attachments = message.data.output.attachments?.map( - ({ kind, ...attachment }) => ({ ...attachment, type: kind }) - ); - const { text: requestedText } = message.data.output; - if (!requestedText) { - if (attachments?.length) { - await sendLinqMessage({ - outgoing: { attachments, raw: "" }, - post, - postReply, - replyToMessageId: requestedReplyMessageId, - }); - await finalizeScheduledReportDelivery(session); - return; - } - await finalizeScheduledReportDelivery(session); - return; - } - - const caller = - session.session.auth.current ?? session.session.auth.initiator; - if (!caller) { - const references = - extractImageArtifactMarkdownReferences(requestedText); - const text = - references.length === 0 - ? requestedText - : [ - stripImageArtifactMarkdownReferences(requestedText), - "I couldn't attach the image.", - ] - .filter(Boolean) - .join("\n\n"); - const outgoing: Extract< - Parameters[0], - { raw: string } - > = { raw: text }; - if (attachments?.length) outgoing.attachments = attachments; - await sendLinqMessage({ - outgoing, - post, - postReply, - replyToMessageId: requestedReplyMessageId, - }); - await finalizeScheduledReportDelivery(session); - return; - } - - const delivery = await prepareLinqImageArtifactDelivery(requestedText, { - rootSessionId: report?.workerSessionId ?? session.session.id, - scope: scopeFromPrincipal(caller), - }); - if (delivery.failedArtifactIds.length > 0) { - console.warn("[linq] browser image delivery failed", { - artifactIds: delivery.failedArtifactIds, - sessionId: session.session.id, - }); - } - const failureMessage = - delivery.failedArtifactIds.length === 0 - ? "" - : delivery.failedArtifactIds.length === 1 - ? "I couldn't attach one image." - : `I couldn't attach ${String(delivery.failedArtifactIds.length)} images.`; - const text = [delivery.text, failureMessage] - .filter(Boolean) - .join("\n\n"); - const outgoing: Extract< - Parameters[0], - { raw: string } - > = { raw: text }; - if (attachments?.length) outgoing.attachments = attachments; - if (delivery.files.length > 0) outgoing.files = delivery.files; - await sendLinqMessage({ - outgoing, - post, - postReply, - replyToMessageId: requestedReplyMessageId, - }); - await finalizeScheduledReportDelivery(session); - } - }, async "message.completed"(event, _context, session) { if (event.finishReason === "tool-calls") return; const report = scheduledReportFromSession(session); @@ -332,78 +99,60 @@ export default linqChannel({ await releaseScheduledReportDelivery(session, event.message); }, }, - async onMessage(context, message) { - if (message.author.isBot) return null; - - const auth = defaultLinqAuth(message); - const authorUserName = z.string().safeParse(message.author.userName); - const phoneNumber = authorUserName.success - ? normalizeAuthPhoneNumber(authorUserName.data) - : undefined; - const verifiedUserId = phoneNumber - ? await findVerifiedAuthUserIdByPhoneNumber(phoneNumber) - : undefined; - if (!verifiedUserId || !phoneNumber) { - // Phone possession is the only sign-in factor, so a handle that is not - // linked to a verified user is unauthenticated: never mint a principal - // or a workspace for it. - console.warn("[linq] ignoring message from an unlinked handle", { - threadId: context.thread.id, - }); - return null; - } - const principalId = `better-auth:${verifiedUserId}`; - const scope = accessScopeForUser(principalId); - return { - auth: { - ...auth, - attributes: { - ...auth.attributes, - conversationChannel: "linq", - conversationId: context.thread.id, - linqThreadId: context.thread.id, - linqMessageId: message.id, - phoneNumber, - workspaceId: scope.workspaceId, - }, - principalId, - }, - }; - }, }); -async function sendLinqMessage({ - outgoing, - post, - postReply, - replyToMessageId, -}: { - readonly outgoing: Extract; - readonly post: ( - content: AdapterPostableMessage - ) => Promise<{ readonly id: string }>; - readonly postReply: ( - content: AdapterPostableMessage, - replyToMessageId: string - ) => Promise<{ readonly id: string }>; - readonly replyToMessageId?: string; -}) { - if (!replyToMessageId) { - await post(outgoing); +async function onMessage(thread: Thread, message: Message) { + if (message.author.isBot || message.author.isMe) return; + + const auth = defaultLinqAuth(message); + const authorUserName = z.string().safeParse(message.author.userName); + const phoneNumber = authorUserName.success + ? normalizeAuthPhoneNumber(authorUserName.data) + : undefined; + const verifiedUserId = phoneNumber + ? await findVerifiedAuthUserIdByPhoneNumber(phoneNumber) + : undefined; + if (!verifiedUserId || !phoneNumber) { + // Phone possession is the only sign-in factor, so a handle that is not + // linked to a verified user is unauthenticated: never mint a principal + // or a workspace for it. + console.warn("[linq] ignoring message from an unlinked handle", { + threadId: thread.id, + }); return; } + const principalId = `better-auth:${verifiedUserId}`; + const scope = accessScopeForUser(principalId); + const content = linqMessageContent(message); + if (!content.length) return; try { - await postReply(outgoing, replyToMessageId); - return; - } catch (error) { - if (!unavailableReplyTargetSchema.safeParse(error).success) throw error; - console.warn("[linq] reply target is unavailable", { - replyToMessageId, - }); - await post(outgoing); + await linqAdapter.markRead(thread.id, message.id); + } catch { + // A read receipt must not prevent dispatch. } + await linq.send(content, { + thread, + auth: { + ...auth, + attributes: { + ...auth.attributes, + conversationChannel: "linq", + conversationId: thread.id, + linqThreadId: thread.id, + linqMessageId: message.id, + phoneNumber, + workspaceId: scope.workspaceId, + }, + principalId, + }, + }); } +linq.bot.onDirectMessage(onMessage); +linq.bot.onNewMessage(/[\s\S]*/, onMessage); + +export default linq.channel; + async function findVerifiedAuthUserIdByPhoneNumber(phoneNumber: string) { const auth = await getAuth(); const context = await auth.$context; diff --git a/agent/instructions/content/role/interactive.md b/agent/instructions/content/role/interactive.md index 41486d95..2e816a9e 100644 --- a/agent/instructions/content/role/interactive.md +++ b/agent/instructions/content/role/interactive.md @@ -52,7 +52,7 @@ The main conversation is the control plane. Coordinate the user's work there and - On every user-initiated conversational turn, use `send_message`, or use `react_to_message` alone when a lightweight reaction fully answers the user and words would add nothing. Use `share_contact` for an introduction with OpenInstinct's saveable contact or when asked for it; that call delivers both the introduction and attachment. Linq renders reactions as native Tapbacks; other supported conversations render a compact reaction. Ordinary assistant text is internal and is never user-visible. - Use `schedules-create` to create a one-time reminder, recurring job, monitor, or scheduled follow-up. Use `kind: "calendar"` with the user's IANA timezone for wall-clock recurrence, `kind: "interval"` for elapsed intervals, and `kind: "once"` for one future instant. Summarize the exact task in `prompt`. Use `schedules-list` before changing an ambiguous schedule and `schedules-update` to edit, pause, resume, or delete it. - When the user answers a question previously sent for a scheduled task, call `schedules-answer` with the internal run ID retained in conversation context and their answer. The parked background run continues from the exact point where it asked. -- Use the full native Linq/iMessage surface when it helps. Choose `kind: "message"` for plain text, exact worker artifact references, and HTTPS attachments; text and attachments may be combined. Choose `kind: "link"` with `url` for a standalone native rich link-preview card, or put a URL in message text for a plain tappable URL. `react_to_message` can add or remove any supported Tapback on the current user message: `thumbs_up`, `thumbs_down`, `heart`, `laugh`, `exclamation` (emphasis), or `question`. +- Use the full native Linq/iMessage surface when it helps. Choose `kind: "message"` for plain text, exact worker artifact references, and HTTPS attachments; text and attachments may be combined. Choose `kind: "link"` with `url` for a standalone native rich link-preview card, or put a URL in message text for a plain tappable URL. Use supplied `[Message]`, `[Parts]`, and `[Reply to]` references. `react_to_message` adds or removes one Unicode emoji on its actual target's `messageId`, current or earlier; never guess or expose IDs. - Eve and Linq own read receipts, typing indicators, delivery state, authorization prompts, and input cards. Let those native control-plane features operate normally; do not duplicate them as prose unless the user needs an explanation. - After a `send_message`, `share_contact`, or `react_to_message` call, never repeat or summarize it in assistant text. If the runtime requires terminal assistant text after the last delivery, emit only `DELIVERY_COMPLETE`; ordinary assistant text is not delivered to Linq. - The worker's structured result is coordinator-facing only. Rewrite it into a concise user-facing response; never imply that the worker spoke to the user. diff --git a/agent/lib/linq/content.test.ts b/agent/lib/linq/content.test.ts new file mode 100644 index 00000000..ab9e8945 --- /dev/null +++ b/agent/lib/linq/content.test.ts @@ -0,0 +1,90 @@ +import { Message } from "chat"; +import { describe, expect, it } from "vitest"; +import { linqMessageContent } from "./content"; + +describe("Linq message content", () => { + it("embeds indexed parts and the quoted target before the original text and media", () => { + const message = received({ + parts: [ + { type: "text", value: "this one" }, + { type: "media", url: "https://media.example/photo.png" }, + ], + reply_to: { message_id: "older-message", part_index: 1 }, + }); + message.attachments.push({ + type: "image", + mimeType: "image/png", + url: "https://media.example/photo.png", + }); + expect(linqMessageContent(message)).toEqual([ + { + type: "text", + text: '[Message: {"messageId":"message-2","sender":"user"}]\n[Parts: [{"partIndex":0,"type":"text","value":"this one"},{"partIndex":1,"type":"media","url":"https://media.example/photo.png"}]]\n[Reply to: {"messageId":"older-message","partIndex":1}]', + }, + { type: "text", text: "this one" }, + expect.objectContaining({ type: "file", mediaType: "image/png" }), + ]); + }); + + it("preserves opaque app-card state when the adapter supplies no text or attachments", () => { + const raw = { + parts: [ + { + type: "imessage_app", + url: "data:application/json;base64,e30=", + app: { bundle_id: "example.game", name: "Game" }, + layout: { caption: "Your move" }, + fallback_text: "Game invite", + interactive: true, + }, + ], + }; + const message = received(raw); + message.text = ""; + const content = linqMessageContent(message); + expect(content).toEqual([ + { + type: "text", + text: '[Message: {"messageId":"message-2","sender":"user"}]\n[Parts: [{"partIndex":0,"type":"imessage_app","url":"data:application/json;base64,e30=","app":{"bundle_id":"example.game","name":"Game"},"layout":{"caption":"Your move"},"fallback_text":"Game invite","interactive":true}]]', + }, + { type: "text", text: "[iMessage app card]" }, + ]); + expect(message.raw).toBe(raw); + expect(message.text).toBe(""); + }); + + it("keeps original content when the raw metadata is malformed", () => { + expect(linqMessageContent(received({ parts: "invalid" }))).toEqual([ + { + type: "text", + text: '[Message: {"messageId":"message-2","sender":"user"}]', + }, + { type: "text", text: "this one" }, + ]); + }); + + it("does not turn an empty incoming message into a metadata-only turn", () => { + const message = received({ parts: [] }); + message.text = ""; + expect(linqMessageContent(message)).toEqual([]); + }); +}); + +function received(raw: Message["raw"]) { + return new Message({ + id: "message-2", + threadId: "linq:chat-1", + raw, + text: "this one", + attachments: [], + formatted: { type: "root", children: [] }, + author: { + userId: "sender-1", + userName: "+15550100011", + fullName: "Sender", + isBot: false, + isMe: false, + }, + metadata: { dateSent: new Date("2026-10-05T12:00:00Z"), edited: false }, + }); +} diff --git a/agent/lib/linq/content.ts b/agent/lib/linq/content.ts new file mode 100644 index 00000000..74195124 --- /dev/null +++ b/agent/lib/linq/content.ts @@ -0,0 +1,51 @@ +import { messageToUserContent } from "eve/channels/chat-sdk"; +import { z } from "zod"; +import type { Message } from "chat"; +import { + formatMessageContext, + messagePartReferenceSchema, +} from "@shared/chat/message-context"; + +const messageReferencesSchema = z.object({ + parts: z + .array(messagePartReferenceSchema.omit({ partIndex: true })) + .optional(), + reply_to: z + .object({ + message_id: z.string().optional(), + part_index: z.number().int().nonnegative().optional(), + }) + .nullish(), +}); + +export function linqMessageContent(message: Message) { + const parsed = messageReferencesSchema.safeParse(message.raw); + const references = parsed.success ? parsed.data : undefined; + const hasApp = references?.parts?.some( + (part) => part.type === "imessage_app" + ); + if (!message.text.trim() && !message.attachments.length && !hasApp) return []; + const content = messageToUserContent(message); + const parts = ( + Array.isArray(content) + ? content + : [{ type: "text" as const, text: content }] + ).filter((part) => part.type !== "text" || part.text.length > 0); + const label = formatMessageContext({ + messageId: message.id, + sender: message.author.isMe ? "openinstinct" : "user", + parts: references?.parts?.map((part, partIndex) => ({ + partIndex, + ...part, + })), + replyTo: references?.reply_to?.message_id + ? { + messageId: references.reply_to.message_id, + partIndex: references.reply_to.part_index, + } + : undefined, + }); + if (hasApp && !message.text.trim()) + parts.push({ type: "text", text: "[iMessage app card]" }); + return [{ type: "text" as const, text: label }, ...parts]; +} diff --git a/agent/lib/linq/transport.ts b/agent/lib/linq/transport.ts index 2081a1d8..abe0a66e 100644 --- a/agent/lib/linq/transport.ts +++ b/agent/lib/linq/transport.ts @@ -1,8 +1,13 @@ import { connectLinqCredentials } from "@vercel/connect/eve"; +import { createLinqAdapter } from "@linqapp/chat-sdk-adapter"; +import { vercelOidc } from "eve/channels/auth"; +import type { LinqChannelCredentials } from "eve/channels/linq"; import { LinqAPIV3 } from "@linqapp/sdk"; +import { z } from "zod"; import { env } from "@shared/environment"; +import type { reactToMessageInputSchema } from "@shared/chat/reaction"; -export const linqCredentials = env.LINQ_CONNECTOR +const linqCredentials = env.LINQ_CONNECTOR ? connectLinqCredentials(env.LINQ_CONNECTOR) : { apiKey() { @@ -12,6 +17,16 @@ export const linqCredentials = env.LINQ_CONNECTOR }, }; +const authenticateWebhook = vercelOidc(); +export const linqWebhookVerifier: NonNullable< + LinqChannelCredentials["webhookVerifier"] +> = async (request) => (await authenticateWebhook(request)) ?? false; + +export const linqAdapter = createLinqAdapter({ + credentials: async () => ({ apiKey: await linqCredentials.apiKey() }), + webhookVerifier: env.LINQ_CONNECTOR ? linqWebhookVerifier : () => false, +}); + export async function sendNativeLinqMessage( chatId: string, message: Parameters[1]["message"], @@ -21,3 +36,32 @@ export async function sendNativeLinqMessage( const client = new LinqAPIV3({ apiKey }); return client.chats.messages.send(chatId, { message }, options); } + +export async function sendNativeLinqReaction( + threadId: string, + reaction: z.infer, + signal: AbortSignal +) { + const messageId = z + .uuid({ + error: + "The target messageId must be a valid Linq UUID without punctuation.", + }) + .parse(reaction.messageId); + const apiKey = await linqCredentials.apiKey(); + const { chatId, pendingHandle } = linqAdapter.decodeThreadId(threadId); + if (!chatId || pendingHandle) { + throw new Error("Reactions require an existing Linq conversation."); + } + const client = new LinqAPIV3({ apiKey, maxRetries: 0 }); + const target = await client.messages.retrieve(messageId, { signal }); + if (target.chat_id !== chatId) { + throw new Error("Message target is not in the current conversation."); + } + signal.throwIfAborted(); + if (reaction.operation === "remove") { + await linqAdapter.removeReaction(threadId, messageId, reaction.emoji); + } else { + await linqAdapter.addReaction(threadId, messageId, reaction.emoji); + } +} diff --git a/agent/lib/messaging/delivery.ts b/agent/lib/messaging/delivery.ts new file mode 100644 index 00000000..01970125 --- /dev/null +++ b/agent/lib/messaging/delivery.ts @@ -0,0 +1,142 @@ +import type { AdapterPostableMessage } from "chat"; +import type { ToolContext } from "eve/tools"; +import { z } from "zod"; +import { + linqAdapter, + sendNativeLinqMessage, + sendNativeLinqReaction, +} from "@agent/lib/linq/transport"; +import { prepareLinqImageArtifactDelivery } from "@agent/lib/linq-image-artifact/delivery"; +import { resolveLinqReplyTarget } from "@agent/lib/reply-targets"; +import { scopeFromPrincipal } from "@agent/lib/principal-scope"; +import { + finalizeScheduledReportDelivery, + scheduledReportFromSession, +} from "@agent/lib/schedules/report-lifecycle"; +import type { reactToMessageInputSchema } from "@shared/chat/reaction"; +import type { sendMessageOutputSchema } from "@shared/chat/message-delivery"; + +type OutboundMessage = + | z.infer + | ({ kind: "reaction" } & z.infer); + +const unavailableReplyTargetSchema = z.object({ + status: z.union([z.literal(400), z.literal(404)]), +}); + +export async function deliverMessage( + message: OutboundMessage, + context: ToolContext, + deliveryId?: string +) { + const caller = context.session.auth.current; + if (caller?.principalType !== "user") + throw new Error( + "Delivery requires the current authenticated Linq conversation or browser conversation." + ); + context.abortSignal.throwIfAborted(); + if (caller.attributes.conversationChannel === "eve") { + // In browser chat the successful tool result is the delivered message. + await finalizeScheduledReportDelivery(context); + return; + } + const thread = z + .string() + .startsWith("linq:") + .safeParse(caller.attributes.conversationId); + if ( + caller.attributes.conversationChannel !== "linq" || + !thread.success || + (caller.attributes.linqThreadId !== undefined && + caller.attributes.linqThreadId !== thread.data) + ) + throw new Error( + "Delivery requires the current authenticated Linq conversation. Nothing was sent." + ); + const { chatId, pendingHandle } = linqAdapter.decodeThreadId(thread.data); + if (!chatId || pendingHandle) + throw new Error("Delivery requires an existing Linq conversation."); + if (message.kind === "reaction") { + await sendNativeLinqReaction(thread.data, message, context.abortSignal); + } else { + const report = scheduledReportFromSession(context); + const idempotencyKey = + deliveryId ?? + (report + ? `scheduled-report:${report.runId}:${String(report.sequence)}` + : `message:${context.session.id}:${context.callId}`); + const replyTarget = resolveLinqReplyTarget( + message.replyTo, + context.session.auth + ); + const replyMessageId = + replyTarget?.conversationId === thread.data + ? replyTarget.messageId + : undefined; + if (message.kind === "link") { + await sendWithReplyFallback(replyMessageId, async (target) => { + context.abortSignal.throwIfAborted(); + const nativeMessage: Parameters[1] = { + parts: [{ type: "link", value: message.url }], + idempotency_key: idempotencyKey, + }; + if (target) nativeMessage.reply_to = { message_id: target }; + await sendNativeLinqMessage(chatId, nativeMessage, { + signal: context.abortSignal, + }); + }); + } else { + const images = await prepareLinqImageArtifactDelivery( + message.text ?? "", + { + rootSessionId: report?.workerSessionId ?? context.session.id, + scope: scopeFromPrincipal(caller), + signal: context.abortSignal, + } + ); + if (images.failedArtifactIds.length) + console.warn("[linq] browser image delivery failed", { + artifactIds: images.failedArtifactIds, + sessionId: context.session.id, + }); + const failure = + images.failedArtifactIds.length === 0 + ? "" + : images.failedArtifactIds.length === 1 + ? "I couldn't attach one image." + : `I couldn't attach ${String(images.failedArtifactIds.length)} images.`; + const outgoing: Extract = { + raw: [images.text, failure].filter(Boolean).join("\n\n"), + }; + if (message.attachments?.length) + outgoing.attachments = message.attachments.map( + ({ kind, ...attachment }) => ({ ...attachment, type: kind }) + ); + if (images.files.length) outgoing.files = images.files; + await sendWithReplyFallback(replyMessageId, async (replyToMessageId) => { + context.abortSignal.throwIfAborted(); + await linqAdapter.postMessage(thread.data, outgoing, { + idempotencyKey, + replyToMessageId, + }); + }); + } + } + await finalizeScheduledReportDelivery(context); +} + +async function sendWithReplyFallback( + target: string | undefined, + send: (replyMessageId?: string) => Promise +) { + try { + await send(target); + } catch (error) { + if (!target || !unavailableReplyTargetSchema.safeParse(error).success) + throw error; + console.warn("[linq] reply target is unavailable", { + replyToMessageId: target, + }); + await send(); + } +} diff --git a/agent/tools/messaging.ts b/agent/tools/messaging.ts index 7dede9d6..881f6b0d 100644 --- a/agent/tools/messaging.ts +++ b/agent/tools/messaging.ts @@ -2,16 +2,13 @@ import { defineDynamic, defineTool, toolOutput } from "eve/tools"; import { defineState } from "eve/context"; import { z } from "zod"; import { resolveModeValue } from "../lib/mode"; -import { sendNativeLinqMessage } from "@agent/lib/linq/transport"; +import { deliverMessage } from "@agent/lib/messaging/delivery"; import { scopeFromPrincipal } from "@agent/lib/principal-scope"; import { readLinqOnboardingPhoneNumber } from "@db/services/auth/linq"; import { getInstallationSecrets } from "@db/services/installation-secrets"; import { openInstinctContactUrl } from "@shared/chat/contact-card"; import { env } from "@shared/environment"; -import { - addReactionToMessageOutputSchema, - reactToMessageOutputSchema, -} from "@shared/chat/reaction"; +import { reactToMessageInputSchema } from "@shared/chat/reaction"; import { sendMessageOutputSchema } from "@shared/chat/message-delivery"; const contactDelivery = defineState<{ @@ -28,12 +25,15 @@ function defineSendMessage() { description: "Send exactly one user-visible message to the current conversation. This is the delivery path for questions, progress updates, blockers, and final answers that need words. Choose kind message for plain text, private image artifacts, and HTTPS attachments; text and attachments may be combined, including in replies. Text is delivered exactly as written, so write it like a brief natural text message and do not use Markdown. Put nearly every response in a native quoted thread by setting replyTo: use current for an ordinary answer, clarification, status update, or follow-up prompted by the current user message, including when the user changes topics; use task with a task ID from Eve's Task state for delayed background work; and use automation with the automation ID supplied by a scheduled report. Omit replyTo only when the message is genuinely standalone and does not answer any particular user message, such as an unsolicited announcement or proactive notice, or when no applicable handle is available. Use only handles present in the current context. Choose kind link with a URL to send a standalone native preview. Put an ordinary URL in message text when a preview is not wanted. Call send_message multiple times only when you intentionally want separate messages. Call it directly without an assistant-text preamble, and do not repeat delivered content afterward.", inputSchema: sendMessageOutputSchema, - execute(message) { + availableInSubagents: false, + async execute(input, context) { + const message = sendMessageOutputSchema.parse(input); + await deliverMessage(message, context); return message; }, toModelOutput() { return toolOutput.text( - "The message was submitted to the active channel. Do not repeat it in assistant text." + "The message was accepted by the active channel. Do not repeat it in assistant text." ); }, }); @@ -57,25 +57,6 @@ function defineShareContact() { "Contact sharing requires an active Linq or browser conversation." ); } - let chatId: string | undefined; - if (channel === "linq") { - // Linq uses linq:, with optional :dm/:group on older threads. - const threadId = z - .string() - .regex(/^linq:([^:]+)(?::(?:dm|group))?$/) - .safeParse(caller.attributes.linqThreadId); - chatId = threadId.success ? threadId.data.split(":")[1] : undefined; - if ( - !threadId.success || - !chatId || - chatId === "pending" || - caller.attributes.conversationId !== threadId.data - ) { - throw new Error( - "Contact sharing requires the current authenticated Linq conversation." - ); - } - } let delivery = contactDelivery.get(); if (delivery?.userId !== userId) { const phone = @@ -107,25 +88,11 @@ function defineShareContact() { contactDelivery.update(() => delivery); } if (delivery.sent) return null; - if (chatId) { - const { text: introduction, attachments } = delivery.message; - await sendNativeLinqMessage( - chatId, - { - idempotency_key: `openinstinct-contact:${context.session.id}`, - parts: [ - ...(introduction - ? [{ type: "text" as const, value: introduction }] - : []), - ...(attachments ?? []).map(({ url }) => ({ - type: "media" as const, - url, - })), - ], - }, - { signal: context.abortSignal } - ); - } + await deliverMessage( + delivery.message, + context, + `openinstinct-contact:${context.session.id}` + ); contactDelivery.update(() => ({ ...delivery, sent: true })); return delivery.message; }, @@ -141,23 +108,32 @@ function defineShareContact() { export default defineDynamic({ events: { - "turn.started": (_event, context) => { + "turn.started": (event, context) => { const isLinq = context.channel.kind === "channel:linq"; const send_message = defineSendMessage(); + const turn = z + .object({ data: z.object({ turnId: z.string() }) }) + .safeParse(event); + const browserTarget = + !isLinq && turn.success + ? ` The current browser messageId is ${turn.data.data.turnId}:user.` + : ""; const react_to_message = defineTool({ - description: isLinq - ? "Add or remove a native iMessage Tapback on the user's current message. Use this instead of send_message when a reaction fully communicates a lightweight acknowledgement and words would add nothing. Supports thumbs_up, thumbs_down, heart, laugh, exclamation (emphasis), and question." - : "Acknowledge the user's current message with one compact reaction displayed in the conversation. Use this instead of send_message when the reaction fully communicates the response and words would add nothing. Supports thumbs_up, thumbs_down, heart, laugh, exclamation (emphasis), and question.", - inputSchema: isLinq - ? reactToMessageOutputSchema - : addReactionToMessageOutputSchema, - execute(reaction) { + availableInSubagents: false, + description: + "Add or remove a native emoji reaction to a specific message in the current conversation. Set messageId to the exact supplied ID of that message. When the user replies to an older message and asks to react to it, use the supplied Reply to messageId. Never invent an ID or default to the session's first message. Supply exactly one real Unicode emoji, not a name or shortcode. Use a reaction for a lightweight acknowledgement; deliver answers that need words through send_message." + + browserTarget, + inputSchema: reactToMessageInputSchema, + async execute(input, toolContext) { + // Runtime validation must survive serialization of the model-facing schema. + const reaction = reactToMessageInputSchema.parse(input); + await deliverMessage({ kind: "reaction", ...reaction }, toolContext); return reaction; }, toModelOutput() { return toolOutput.text( - "The reaction was submitted to the active conversation. Do not repeat it in assistant text." + "The reaction was accepted by the active channel. Do not repeat it or duplicate it in assistant text." ); }, }); diff --git a/app/(authenticated)/chat/[sessionId]/_components/activity/trace.tsx b/app/(authenticated)/chat/[sessionId]/_components/activity/trace.tsx index e4ddccbd..62269c45 100644 --- a/app/(authenticated)/chat/[sessionId]/_components/activity/trace.tsx +++ b/app/(authenticated)/chat/[sessionId]/_components/activity/trace.tsx @@ -17,7 +17,7 @@ import { import { Badge } from "@web/components/ui/badge"; import { Button } from "@web/components/ui/button"; import { getLatestTurnFailure } from "../../_lib/turn-failure"; -import { messageTimestamps } from "../../_lib/message-events"; +import { conversationRows } from "../../_lib/conversation-rows"; import type { SubagentStatus } from "@app/_lib/subagent-sessions"; import { AgentMessage } from "../conversation/message"; @@ -50,7 +50,10 @@ export function SubagentTrace({ ), [events] ); - const timestamps = useMemo(() => messageTimestamps(events), [events]); + const rows = useMemo( + () => conversationRows(data.messages, events, "trace"), + [data.messages, events] + ); const isRunning = status === "starting" || status === "working"; const turnFailure = useMemo(() => getLatestTurnFailure(events), [events]); const error = streamError ?? turnFailure; @@ -91,14 +94,13 @@ export function SubagentTrace({ {isLoadingOlder ? "Loading…" : "Load older messages"} ) : null} - {data.messages.map((message, index) => ( + {rows.map((message, index) => ( undefined} - timestamp={timestamps.get(message.id)} /> ))} {(isLoading || isRunning) && data.messages.length === 0 ? ( diff --git a/app/(authenticated)/chat/[sessionId]/_components/conversation/index.test.tsx b/app/(authenticated)/chat/[sessionId]/_components/conversation/index.test.tsx index 46388d23..d19972cb 100644 --- a/app/(authenticated)/chat/[sessionId]/_components/conversation/index.test.tsx +++ b/app/(authenticated)/chat/[sessionId]/_components/conversation/index.test.tsx @@ -4,8 +4,45 @@ import { renderToStaticMarkup } from "react-dom/server"; import { describe, expect, it } from "vitest"; import { ChatConversation } from "."; import type { ChatAgent } from "../chat-agent"; +import { defaultMessageReducer } from "eve/client"; +import { messageHistoryEvents } from "@tests/fixtures/message-history"; describe("chat conversation", () => { + it("renders new provider history, attachments, quoted parts, and targeted reactions", () => { + const reducer = defaultMessageReducer(); + const data = messageHistoryEvents.reduce( + (state, event) => reducer.reduce(state, event), + reducer.initial() + ); + const agent = { + data, + events: messageHistoryEvents, + error: undefined, + respond: async () => undefined, + status: "ready" as const, + }; + const markup = renderToStaticMarkup( + + ); + expect(markup).toContain("Here are two photos."); + expect(markup).toContain('aria-label="Open orange.svg"'); + expect(markup).toContain('aria-label="Open blue.svg"'); + expect(markup).toContain('aria-label="Reply to blue.svg"'); + expect(markup).toContain('href="#photos%3Areceived%3Auser"'); + expect(markup).toContain('aria-label="Reaction πŸ‘"'); + expect(markup).toContain('aria-label="Reaction πŸ‘©πŸ½β€πŸ’»"'); + expect(markup).toContain('aria-label="Reaction ❀️"'); + expect(markup).not.toContain('aria-label="Reaction πŸ‘€"'); + expect(markup).not.toContain("11111111-1111-4111-8111-111111111111"); + expect(markup).toContain("literal-parts-example"); + expect(markup).not.toContain("opaque-app-payload"); + expect(markup).toContain("example"); + const trace = renderToStaticMarkup( + + ); + expect(trace).toContain("[Parts:"); + expect(trace).toContain("opaque-app-payload"); + }); it("shows send_message output instead of assistant stream text", () => { const agent = { data: { @@ -33,7 +70,10 @@ describe("chat conversation", () => { ], }, error: undefined, - events: [sendMessageResult("The visible iMessage response.")], + events: [ + delivery("turn-1", "What happened?"), + sendMessageResult("The visible iMessage response."), + ], respond: async () => undefined, status: "ready", } satisfies Pick< @@ -63,6 +103,7 @@ describe("chat conversation", () => { role: "assistant", } satisfies EveMessage; const events = [ + delivery("visible-turn", "Keep this visible"), workerReceipt("task_worker"), workerCancellation("task_worker"), delivery("task-delivery", cancellationText), @@ -91,7 +132,7 @@ describe("chat conversation", () => { const agent = { data: { messages: [message("turn-1:user", "Try this")] }, error: new Error("Internal runtime failure"), - events: [], + events: [delivery("turn-1", "Try this")], respond: async () => undefined, status: "error", } satisfies Pick< @@ -166,7 +207,7 @@ function workerCancellation(taskId: string): MessageStreamEvent { function delivery(turnId: string, messageText: string): MessageStreamEvent { return { data: { message: messageText, sequence: 0, turnId }, - meta: { at: "2026-08-27T20:00:01.000Z", id: "delivery" }, + meta: { at: "2026-08-27T20:00:01.000Z", id: `receipt-${turnId}` }, type: "message.received", }; } diff --git a/app/(authenticated)/chat/[sessionId]/_components/conversation/index.tsx b/app/(authenticated)/chat/[sessionId]/_components/conversation/index.tsx index 0e95fc79..991703cf 100644 --- a/app/(authenticated)/chat/[sessionId]/_components/conversation/index.tsx +++ b/app/(authenticated)/chat/[sessionId]/_components/conversation/index.tsx @@ -1,12 +1,8 @@ import { AlertCircleIcon, BrainIcon, LoaderCircleIcon } from "lucide-react"; -import { Fragment, useMemo } from "react"; -import { - imessageTimestamps, - messageTimestamps, - sentMessages, -} from "../../_lib/message-events"; -import { messagesForTraceView, type TraceView } from "../../_lib/trace-view"; +import { useMemo } from "react"; +import type { TraceView } from "../../_lib/trace-view"; import { getLatestTurnFailure } from "../../_lib/turn-failure"; +import { conversationRows } from "../../_lib/conversation-rows"; import { Conversation, ConversationContent, @@ -59,20 +55,9 @@ export function ChatConversation({ const errorMessage = (agent.error ? toErrorMessage(agent.error) : undefined) ?? turnFailure; const messages = useMemo( - () => messagesForTraceView(agent.data.messages, agent.events, traceView), + () => conversationRows(agent.data.messages, agent.events, traceView), [agent.data.messages, agent.events, traceView] ); - const timestamps = useMemo( - () => - traceView === "imessage" - ? imessageTimestamps(agent.events) - : messageTimestamps(agent.events), - [agent.events, traceView] - ); - const deliveredMessages = useMemo( - () => sentMessages(agent.events), - [agent.events] - ); return ( - {deliveries.map((delivery) => ( - agent.respond(responses)} - sentMessageParts={delivery.parts} - timestamp={delivery.timestamp} - userVisibleOnly - /> - ))} - - ); - } - return ( agent.respond(responses)} - timestamp={timestamps.get(message.id)} userVisibleOnly={traceView === "imessage"} /> ); diff --git a/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.test.tsx b/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.test.tsx index 434065fc..b15eedf3 100644 --- a/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.test.tsx +++ b/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.test.tsx @@ -2,9 +2,10 @@ import type { EveMessage } from "eve/react"; import { renderToStaticMarkup } from "react-dom/server"; import { describe, expect, it } from "vitest"; import { AgentMessage } from "."; +import type { ConversationRow } from "../../../_lib/conversation-rows"; -describe("agent messages", () => { - it("renders ordinary assistant text without a delivery tool result", () => { +describe("conversation row rendering", () => { + it("renders ordinary trace text supplied by Eve", () => { const message = { id: "assistant-message", metadata: { status: "complete" }, @@ -17,7 +18,6 @@ describe("agent messages", () => { ], role: "assistant", } satisfies EveMessage; - const markup = renderToStaticMarkup( { onInputResponses={() => undefined} /> ); - expect(markup).toContain("Hello from ordinary assistant output."); }); - it("renders only Linq-delivered content in the iMessage view", () => { - const message = { - id: "turn-1:assistant", - metadata: { status: "complete", turnId: "turn-1" }, - parts: [ - { - state: "done", - stepIndex: 1, - text: "I’ll check that now.", - type: "text", - }, - { - state: "done", - stepIndex: 0, - text: "Private reasoning", - type: "reasoning", - }, - { - input: { query: "example" }, - output: { result: "internal" }, - state: "output-available", - stepIndex: 0, - toolCallId: "call-1", - toolName: "web_search", - type: "dynamic-tool", - }, - { - state: "done", - stepIndex: 1, - text: "Here’s what I found.", - type: "text", - }, - ], - role: "assistant", - } satisfies EveMessage; - - const markup = renderToStaticMarkup( - undefined} - sentMessageParts={[ - { - state: "done", - stepIndex: 1, - text: "Here’s what I found.", - type: "text", - }, - ]} - userVisibleOnly - /> - ); - - expect(markup).toContain("Here’s what I found."); - expect(markup).not.toContain("I’ll check that now."); - expect(markup).not.toContain("Private reasoning"); - expect(markup).not.toContain("web_search"); - }); - - it("hides non-send_message controls in the iMessage projection", () => { + it("renders the projected reply, reaction, and timestamp without a second content override", () => { const message = { - id: "turn-2:assistant", - metadata: { status: "streaming", turnId: "turn-2" }, - parts: [ - { - approval: { id: "approval-1" }, - input: { amount: 50, recipient: "Hidden recipient" }, - state: "approval-requested", - stepIndex: 0, - toolCallId: "call-2", - toolMetadata: { - eve: { - inputRequest: { - kind: "tool-approval", - options: [ - { id: "approve", label: "Approve", style: "primary" }, - { id: "cancel", label: "Cancel", style: "danger" }, - ], - prompt: "Approve this action?", - requestId: "approval-1", - }, - kind: "tool-call", - name: "send_payment", - }, - }, - toolName: "send_payment", - type: "dynamic-tool", - }, - ], - role: "assistant", - } satisfies EveMessage; - + id: "receipt:user", + metadata: { status: "complete" }, + parts: [{ type: "text", text: "Here is the reply.", state: "done" }], + role: "user", + timestamp: "2026-10-05T20:00:00.000Z", + reactions: ["πŸ‘"], + reply: { targetId: "older:user", text: "Original request" }, + } satisfies ConversationRow; const markup = renderToStaticMarkup( { userVisibleOnly /> ); - - expect(markup).not.toContain("Approve this action?"); - expect(markup).not.toContain("Approve"); - expect(markup).not.toContain("Cancel"); - expect(markup).not.toContain("send_payment"); - expect(markup).not.toContain("Hidden recipient"); + expect(markup).toContain("Here is the reply."); + expect(markup).toContain('href="#older%3Auser"'); + expect(markup).toContain('aria-label="Reaction πŸ‘"'); + expect(markup).toContain('dateTime="2026-10-05T20:00:00.000Z"'); }); }); diff --git a/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.tsx b/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.tsx index 325dd9ea..1c1c41a0 100644 --- a/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.tsx +++ b/app/(authenticated)/chat/[sessionId]/_components/conversation/message/index.tsx @@ -1,35 +1,30 @@ "use client"; -import type { EveMessage } from "eve/react"; import { useState } from "react"; import { Message, MessageContent } from "@web/components/ai-elements/message"; import { cn } from "@web/components/class-names"; import { AgentMessagePart, partKey } from "./parts"; import type { RespondToAgentInput } from "./types"; +import type { ConversationRow } from "../../../_lib/conversation-rows"; export function AgentMessage({ canRespond, isStreaming, message, onInputResponses, - sentMessageParts, - timestamp, userVisibleOnly = false, }: { readonly canRespond: boolean; readonly isStreaming: boolean; - readonly message: EveMessage; + readonly message: ConversationRow; readonly onInputResponses: RespondToAgentInput; - readonly sentMessageParts?: readonly EveMessage["parts"][number][]; - readonly timestamp?: string; readonly userVisibleOnly?: boolean; }) { const [optimisticTimestamp] = useState(() => new Date().toISOString()); const displayedTimestamp = - timestamp ?? (message.role === "user" ? optimisticTimestamp : undefined); - const visibleParts = userVisibleOnly - ? userVisibleParts(message, sentMessageParts) - : message.parts; + message.timestamp ?? + (message.role === "user" ? optimisticTimestamp : undefined); + const visibleParts = message.parts; const lastTextIndex = visibleParts.reduce( (last, part, index) => (part.type === "text" ? index : last), -1 @@ -42,10 +37,12 @@ export function AgentMessage({ return ( + {message.reply ? : null} {visibleParts.map((part, index) => hasAssistantText && part.type === "reasoning" ? null : ( + {message.reactions?.length ? ( +
    + {message.reactions.map((emoji) => ( +
  • + {emoji} +
  • + ))} +
+ ) : null} {displayedTimestamp ? (