Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions apps/web/src/chat/threads-api.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ function meTurn(subject: string, body: string) {
return {
id: "Sent:1",
messageId: "m1",
parentId: undefined,
address: "run_alice@example.com",
author: "me" as const,
subject,
Expand Down
88 changes: 61 additions & 27 deletions apps/web/src/chat/threads-api.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
// Chats are mail threads (see docs/chat-mail-threading.md). Keyed by the
// agent's definition asset id, not a run id, since a redeploy retires the
// run id and would 409 a keyed-on-run chat.
// Chats are mail threads (see docs/chat-mail-threading.md). A chat is keyed
// by its thread root turn's id (the person's opening send), so two chats to
// the same agent stay separate instead of merging into one conversation.

import { type } from "arktype";
import { WorkflowDeploymentResponse } from "@intx/types";
Expand Down Expand Up @@ -64,7 +64,8 @@ export type ChatMessage = {
};

export type ChatSummary = {
/** Route id: the agent's run id. */
/** Route id: the thread's root turn id (`folder:uid` of its opening
* send) — unique per conversation, unlike the agent it's addressed to. */
readonly id: string;
readonly title: string;
readonly agentName: string;
Expand Down Expand Up @@ -332,6 +333,9 @@ function participantAddress(
type MailTurn = {
readonly id: string;
readonly messageId: string;
/** In-Reply-To, or the newest References entry when that header is
* missing — what ties a reply back to the turn before it. */
readonly parentId: string | undefined;
readonly address: string;
readonly author: "me" | "agent";
readonly subject: string;
Expand All @@ -352,6 +356,7 @@ async function readFolder(tenantId: string, folder: "INBOX" | "Sent"): Promise<M
{
id: `${folder}:${String(message.uid)}`,
messageId: message.envelope.messageId,
parentId: message.envelope.inReplyTo ?? message.envelope.references.at(-1),
address,
author: folder === "Sent" ? ("me" as const) : ("agent" as const),
subject: message.envelope.subject,
Expand All @@ -371,6 +376,26 @@ async function readTurns(tenantId: string): Promise<MailTurn[]> {
return [...inbox, ...sent].sort((a, b) => Date.parse(a.at) - Date.parse(b.at));
}

/** Each turn's own id, mapped to the id of its thread's opening turn —
* found by walking In-Reply-To/References back to a turn with no known
* parent. A cycle (shouldn't happen) just stops at the point it's seen. */
function threadRootIds(turns: readonly MailTurn[]): Map<string, string> {
const byMessageId = new Map(turns.map((turn) => [turn.messageId, turn]));
const roots = new Map<string, string>();
for (const turn of turns) {
let current = turn;
const seen = new Set<string>([turn.messageId]);
for (;;) {
const parent = current.parentId === undefined ? undefined : byMessageId.get(current.parentId);
if (parent === undefined || seen.has(parent.messageId)) break;
seen.add(parent.messageId);
current = parent;
}
roots.set(turn.id, current.id);
}
return roots;
}

// ---------------------------------------------------------------------
// Sends
// ---------------------------------------------------------------------
Expand All @@ -382,7 +407,7 @@ async function sendToAgent(
address: string,
body: string,
inReplyTo: string | undefined,
): Promise<void> {
): Promise<{ readonly messageId: string; readonly uid: number }> {
let response: Response;
try {
response = await fetch(`${mailboxPath(tenantId)}/send`, {
Expand All @@ -408,11 +433,14 @@ async function sendToAgent(
if (parsed instanceof type.errors) {
throw new ChatApiError(`Unexpected send response: ${parsed.summary}`);
}
return parsed;
}

/** Starts a chat with one agent. Returns the chat id to route to — the
* agent's definition asset id. Callers should keep the composer disabled
* until `liveAddress` is set; this still guards against a stale click. */
/** Starts a brand-new chat with one agent: no `inReplyTo`, so the send opens
* its own thread rather than landing on whatever this agent last answered.
* Returns the chat id to route to — the opening send's own `Sent:<uid>`, the
* thread's root turn id. Callers should keep the composer disabled until
* `liveAddress` is set; this still guards against a stale click. */
export async function startChat(
tenantId: string,
agent: ChatAgent,
Expand All @@ -421,8 +449,8 @@ export async function startChat(
if (agent.liveAddress === null) {
throw new ChatApiError(`${agent.name} is starting…`);
}
await sendToAgent(tenantId, agent.liveAddress, content, undefined);
return agent.id;
const sent = await sendToAgent(tenantId, agent.liveAddress, content, undefined);
return `Sent:${String(sent.uid)}`;
}

/** A reply is the same send, threaded onto the chat's newest message, sent
Expand Down Expand Up @@ -469,25 +497,27 @@ function agentByAddress(agents: readonly ChatAgent[]): Map<string, ChatAgent> {
return index;
}

/** Every chat the person has: one per agent they have exchanged mail with,
* grouped by the agent's asset id so a redeploy's new run still lands in
* the same chat. */
/** Every chat the person has: one per mail thread with an agent, grouped by
* the thread's root turn id, not the agent — an agent's address set spans
* every run it has ever had, so two separate threads to the same agent stay
* two separate chats instead of merging into one. */
export async function listChats(tenantId: string): Promise<readonly ChatSummary[]> {
const [agents, turns] = await Promise.all([listChatAgents(tenantId), readTurns(tenantId)]);
const addressToAgent = agentByAddress(agents);
const byAgent = new Map<string, MailTurn[]>();
const rootIds = threadRootIds(turns);
const byThread = new Map<string, MailTurn[]>();
for (const turn of turns) {
const agent = addressToAgent.get(turn.address);
if (agent === undefined) continue;
byAgent.set(agent.id, [...(byAgent.get(agent.id) ?? []), turn]);
if (!addressToAgent.has(turn.address)) continue;
const rootId = rootIds.get(turn.id)!;
byThread.set(rootId, [...(byThread.get(rootId) ?? []), turn]);
}
return [...byAgent.entries()]
.map(([agentId, rows]) => {
const resolvedName = agents.find((agent) => agent.id === agentId)?.name;
const agentName = resolvedName ?? agentId;
return [...byThread.entries()]
.map(([threadId, rows]) => {
const newest = rows[rows.length - 1]!;
const resolvedName = addressToAgent.get(newest.address)?.name;
const agentName = resolvedName ?? "Unknown agent";
return {
id: agentId,
id: threadId,
title: chatTitle(rows, resolvedName),
agentName,
preview: newest.body.slice(0, 80),
Expand Down Expand Up @@ -528,16 +558,20 @@ export function isChatReplyReady(
}
}

/** One chat's full transcript: every mail turn addressed to any run this
* agent has ever had, oldest first. */
/** One chat's full transcript: every turn in the thread rooted at `chatId`,
* oldest first. */
export async function readChat(tenantId: string, chatId: string): Promise<ChatThread> {
const [agents, turns] = await Promise.all([listChatAgents(tenantId), readTurns(tenantId)]);
const agent = agents.find((candidate) => candidate.id === chatId);
const rootIds = threadRootIds(turns);
const rows = turns.filter((turn) => rootIds.get(turn.id) === chatId);
if (rows.length === 0) {
throw new ChatApiError("That chat could not be found.", 404);
}
const addressToAgent = agentByAddress(agents);
const agent = addressToAgent.get(rows[0]!.address);
if (agent === undefined) {
throw new ChatApiError("That agent could not be found.", 404);
}
const addresses = new Set(agent.addresses);
const rows = turns.filter((turn) => addresses.has(turn.address));
return {
id: chatId,
title: chatTitle(rows, agent.name),
Expand Down
4 changes: 3 additions & 1 deletion apps/web/src/pages/chat-thread-page.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -248,7 +248,9 @@ function ChatTranscript({
<IdentityAvatar
kind={message.author === "me" ? "person" : "agent"}
name={message.authorName}
principalId={message.author === "me" ? (selectedPrincipalId ?? "me") : chat.id}
principalId={
message.author === "me" ? (selectedPrincipalId ?? "me") : chat.agent.id
}
/>
</span>
<div className="chat-thread-body">
Expand Down
13 changes: 8 additions & 5 deletions docs/chat-mail-threading.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,14 @@ mail: a chat has exactly one agent, and both sides are durable mail — the
person sends from their own mailbox (keeping a Sent copy), and an agent's
reply lands in the same mailbox's INBOX. No chat-specific hub route exists.

A chat is keyed by its agent's definition asset id, not a run id: every hub
restart releases the old run and redeploys under a new one, so keying on a
run id would 409 the moment it turns terminal. An agent's address set spans
every run it has ever had, which keeps history intact across a redeploy;
sends resolve the current live run's address at send time.
A chat is keyed by its thread's root turn id (`Sent:<uid>` of the person's
opening send), found by walking In-Reply-To/References back from every mail
turn — not the agent's definition asset id, so two separate chats with the
same agent stay separate instead of merging into one conversation. Sends
resolve the agent's current live run's address at send time; an agent's
address set spans every run it has ever had (a hub restart retires the old
run and redeploys under a new one), which keeps a thread's history readable
across a redeploy even though the address on later turns changes.

A workbench is a child tenant, and its conversation is that tenant's own
mailbox — reads are the same stock mailbox routes as a chat, scoped to the
Expand Down
Loading