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
39 changes: 37 additions & 2 deletions apps/mobile/src/features/threads/ThreadFeed.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -1449,6 +1449,36 @@ function useMarkdownStyles(
]);
}

function AgentMessageAttribution(props: {
readonly environmentId: EnvironmentId;
readonly senderThreadId?: ThreadId;
}) {
const navigation = useNavigation();
const senderThreadId = props.senderThreadId;
const label = (
<Text className="mb-1 pr-1 font-t3-medium text-2xs text-foreground-muted opacity-60">
Sent by another agent
</Text>
);
return senderThreadId ? (
<Pressable
accessibilityRole="button"
accessibilityLabel="Open sending thread"
hitSlop={4}
onPress={() =>
navigation.navigate("Thread", {
environmentId: String(props.environmentId),
threadId: String(senderThreadId),
})
}
>
{label}
</Pressable>
) : (
label
);
}

function renderFeedEntry(
info: { item: PendingThreadFeedEntry; index: number },
props: Pick<
Expand Down Expand Up @@ -1619,10 +1649,15 @@ function renderFeedEntry(
className="mb-5 items-end"
{...(enterAnimated ? { entering: FadeInUp.duration(220) } : {})}
>
{presentation.isAutomation || message.createdBy === "agent" ? (
{presentation.isAutomation ? (
<Text className="mb-1 pr-1 font-t3-medium text-2xs text-foreground-muted opacity-60">
{presentation.isAutomation ? "Sent by automation" : "Sent by another agent"}
Sent by automation
</Text>
) : message.createdBy === "agent" ? (
<AgentMessageAttribution
environmentId={props.environmentId}
senderThreadId={message.senderThreadId}
/>
) : null}
<View
className="min-w-0 gap-2 rounded-[20px] px-3.5 py-2.5"
Expand Down
17 changes: 17 additions & 0 deletions apps/mobile/src/lib/threadActivity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,23 @@ describe("buildThreadFeed", () => {
});
});

it("keeps the sender of an agent message distinct from its timeline source", () => {
const feed = buildThreadFeed([
projected(
{
...userMessage(),
createdBy: "agent",
creationSource: "mcp",
senderThreadId: sourceThreadId,
},
0,
),
]);
const messageEntry = feed.find((entry) => entry.type === "message");
expect(messageEntry?.message.senderThreadId).toBe(sourceThreadId);
expect(messageEntry?.message.sourceThreadId).toBe(threadId);
});

it("adds local feedback messages to an otherwise server-authored feed", () => {
const feed = buildThreadFeed([], {
localMessages: [
Expand Down
2 changes: 2 additions & 0 deletions apps/mobile/src/lib/threadActivity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ export interface ThreadFeedMessage {
readonly createdBy?: OrchestrationV2Actor;
readonly creationSource?: OrchestrationV2CreationSource;
readonly scheduledTaskId?: ScheduledTaskId;
readonly senderThreadId?: ThreadId;
readonly visibility: OrchestrationV2ProjectedTurnItem["visibility"];
readonly sourceThreadId: ThreadId;
readonly createdAt: string;
Expand Down Expand Up @@ -1559,6 +1560,7 @@ export function buildThreadFeed(
createdBy: item.createdBy,
creationSource: item.creationSource,
...(item.scheduledTaskId ? { scheduledTaskId: item.scheduledTaskId } : {}),
...(item.senderThreadId ? { senderThreadId: item.senderThreadId } : {}),
}
: {}),
visibility: row.visibility,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1640,6 +1640,7 @@ const make = Effect.gen(function* () {
index,
}),
threadId,
senderThreadId: scope.threadId,
messageId: stableMessageId({
scope,
requestKey: key,
Expand Down Expand Up @@ -1840,6 +1841,7 @@ const make = Effect.gen(function* () {
operation: "thread-send",
}),
threadId: input.threadId,
senderThreadId: scope.threadId,
messageId,
text: input.message,
attachments: [],
Expand Down
46 changes: 46 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1453,6 +1453,17 @@ describe("orchestrator MCP toolkit", () => {
const delegated = yield* decodeDelegateTaskResult(delegatedCall.structuredContent).pipe(
Effect.orDie,
);
const delegatedSource = yield* orchestrator.getThreadProjection(
delegated.childThreadId,
);
expect(delegatedSource.messages[0]).toMatchObject({
senderThreadId: parentThreadId,
});
expect(
delegatedSource.turnItems.find((item) => item.type === "user_message"),
).toMatchObject({
senderThreadId: parentThreadId,
});
expect(delegated.status).toBe("completed");
expect(delegated.summary).toBe(delegatedResult);
expect(delegated.providerInstanceId).toBe(claudeInstanceId);
Expand Down Expand Up @@ -1847,6 +1858,15 @@ describe("orchestrator MCP toolkit", () => {
providerInstanceId: claudeInstanceId,
model: claudeModel,
});
const createdSource = yield* orchestrator.getThreadProjection(promptedThread.threadId);
expect(createdSource.messages[0]).toMatchObject({
senderThreadId: parentThreadId,
});
expect(
createdSource.turnItems.find((item) => item.type === "user_message"),
).toMatchObject({
senderThreadId: parentThreadId,
});
const emptyProjection = yield* orchestrator.getThreadProjection(emptyThread.threadId);
expect(emptyProjection.thread.lineage).toEqual({
parentThreadId: null,
Expand Down Expand Up @@ -2097,6 +2117,19 @@ describe("orchestrator MCP toolkit", () => {
const sent = yield* decodeThreadSendResult(sendCall.structuredContent).pipe(
Effect.orDie,
);
const sentSource = yield* orchestrator.getThreadProjection(emptyThread.threadId);
expect(
sentSource.messages.find((message) => message.id === sent.messageId),
).toMatchObject({
senderThreadId: parentThreadId,
});
expect(
sentSource.turnItems.find(
(item) => item.type === "user_message" && item.messageId === sent.messageId,
),
).toMatchObject({
senderThreadId: parentThreadId,
});
expect(sent.delivery).toBe("started");
const waitCall = yield* invoke("t3_thread_wait", {
threadId: emptyThread.threadId,
Expand Down Expand Up @@ -2181,6 +2214,19 @@ describe("orchestrator MCP toolkit", () => {
runId: activeRun.id,
delivery: "steered",
});
const steeredSource = yield* orchestrator.getThreadProjection(activeThread.threadId);
expect(
steeredSource.messages.find((message) => message.id === steered.messageId),
).toMatchObject({
senderThreadId: parentThreadId,
});
expect(
steeredSource.turnItems.find(
(item) => item.type === "user_message" && item.messageId === steered.messageId,
),
).toMatchObject({
senderThreadId: parentThreadId,
});
const interruptCall = yield* invoke("t3_thread_interrupt", {
threadId: activeThread.threadId,
reason: "The orchestration loop has enough evidence.",
Expand Down
4 changes: 4 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2531,6 +2531,7 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
});
const promptNativeItemId = `${nativeTaskId}:prompt`;
const promptArtifacts = makeSubagentConversationArtifacts({
senderThreadId: context.input.threadId,
messageId: providerMessageId(promptNativeItemId),
turnItemId: providerTurnItemId(promptNativeItemId),
threadId: childThreadId,
Expand Down Expand Up @@ -6565,6 +6566,9 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
...(turnInput.message.scheduledTaskId === undefined
? {}
: { scheduledTaskId: turnInput.message.scheduledTaskId }),
...(turnInput.message.senderThreadId === undefined
? {}
: { senderThreadId: turnInput.message.senderThreadId }),
id: turnInput.message.messageId,
threadId: turnInput.threadId,
runId: turnInput.runId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3533,6 +3533,7 @@ export function makeClaudeAdapterV2(
if (existingSubagent === undefined) {
const promptNativeItemId = `${nativeItemId}:prompt`;
const promptArtifacts = makeSubagentConversationArtifacts({
senderThreadId: input.context.input.threadId,
messageId: idAllocator.derive.messageFromProviderItem({
driver: CLAUDE_PROVIDER,
nativeItemId: promptNativeItemId,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2420,6 +2420,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
if (input.emitInitialPrompt && input.prompt.length > 0) {
const promptNativeItemId = `${input.nativeItemId}:prompt`;
const promptArtifacts = makeSubagentConversationArtifacts({
senderThreadId: input.context.projectionThreadId,
messageId: idAllocator.derive.messageFromProviderItem({
driver: CODEX_PROVIDER,
nativeItemId: promptNativeItemId,
Expand Down Expand Up @@ -2872,6 +2873,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
const now = yield* DateTime.now;
const ordinal = yield* resolveItemOrdinal(context, item.id);
const artifacts = makeSubagentConversationArtifacts({
senderThreadId: context.subagent.parentContext.projectionThreadId,
messageId: idAllocator.derive.messageFromProviderItem({
driver: CODEX_PROVIDER,
nativeItemId: item.id,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1604,6 +1604,7 @@ export function makeCursorAdapterV2(
});
const promptNativeId = `${nativeItemId}:prompt`;
const promptArtifacts = makeSubagentConversationArtifacts({
senderThreadId: input.context.input.threadId,
messageId: idAllocator.derive.messageFromProviderItem({
driver: CURSOR_PROVIDER,
nativeItemId: promptNativeId,
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,9 @@ export const executorLayer: Layer.Layer<
...(message.scheduledTaskId === undefined
? {}
: { scheduledTaskId: message.scheduledTaskId }),
...(message.senderThreadId === undefined
? {}
: { senderThreadId: message.senderThreadId }),
});
}),
),
Expand Down
24 changes: 24 additions & 0 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1410,6 +1410,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(queuedMessage.scheduledTaskId === undefined
? {}
: { scheduledTaskId: queuedMessage.scheduledTaskId }),
...(queuedMessage.senderThreadId === undefined
? {}
: { senderThreadId: queuedMessage.senderThreadId }),
}),
inputIntent: "queued_turn",
startedAt: now,
Expand Down Expand Up @@ -3285,6 +3288,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
readonly createdBy: OrchestrationV2ConversationMessage["createdBy"];
readonly creationSource: OrchestrationV2ConversationMessage["creationSource"];
readonly scheduledTaskId?: OrchestrationV2ConversationMessage["scheduledTaskId"];
readonly senderThreadId?: OrchestrationV2ConversationMessage["senderThreadId"];
readonly delegatedCompletion?: OrchestrationV2ConversationMessage["delegatedCompletion"];
readonly forceRestart: boolean;
}) =>
Expand Down Expand Up @@ -3430,6 +3434,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(input.scheduledTaskId === undefined
? {}
: { scheduledTaskId: input.scheduledTaskId }),
...(input.senderThreadId === undefined ? {} : { senderThreadId: input.senderThreadId }),
id: input.messageId,
threadId: input.command.threadId,
runId: messageInput.runId,
Expand All @@ -3448,6 +3453,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(input.scheduledTaskId === undefined
? {}
: { scheduledTaskId: input.scheduledTaskId }),
...(input.senderThreadId === undefined ? {} : { senderThreadId: input.senderThreadId }),
id: idAllocator.derive.userTurnItem({ messageId: input.messageId }),
threadId: input.command.threadId,
runId: messageInput.runId,
Expand Down Expand Up @@ -4288,6 +4294,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined
? {}
: { senderThreadId: command.senderThreadId }),
forceRestart: dispatchMode.type === "restart_active",
});
return;
Expand Down Expand Up @@ -4469,6 +4478,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined
? {}
: { senderThreadId: command.senderThreadId }),
id: command.messageId,
threadId: command.threadId,
runId,
Expand Down Expand Up @@ -4805,6 +4817,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined
? {}
: { senderThreadId: command.senderThreadId }),
id: command.messageId,
threadId: command.threadId,
runId,
Expand All @@ -4825,6 +4840,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined
? {}
: { senderThreadId: command.senderThreadId }),
id: idAllocator.derive.userTurnItem({ messageId: command.messageId }),
threadId: command.threadId,
runId,
Expand Down Expand Up @@ -5491,6 +5509,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined ? {} : { senderThreadId: command.senderThreadId }),
id: command.messageId,
threadId: command.threadId,
runId,
Expand All @@ -5511,6 +5530,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(command.scheduledTaskId === undefined
? {}
: { scheduledTaskId: command.scheduledTaskId }),
...(command.senderThreadId === undefined ? {} : { senderThreadId: command.senderThreadId }),
id: idAllocator.derive.userTurnItem({ messageId: command.messageId }),
threadId: command.threadId,
runId,
Expand Down Expand Up @@ -6110,6 +6130,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
creationSource: command.creationSource,
commandId: command.commandId,
threadId: childThreadId,
senderThreadId: command.parentThreadId,
messageId: childMessageId,
text: command.task,
attachments: [],
Expand Down Expand Up @@ -6795,6 +6816,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
...(queuedMessage.scheduledTaskId === undefined
? {}
: { scheduledTaskId: queuedMessage.scheduledTaskId }),
...(queuedMessage.senderThreadId === undefined
? {}
: { senderThreadId: queuedMessage.senderThreadId }),
forceRestart: false,
});
});
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/orchestration-v2/ProviderAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ export const ProviderAdapterV2TurnMessage = Schema.Struct({
createdBy: OrchestrationV2ConversationMessage.fields.createdBy,
creationSource: OrchestrationV2ConversationMessage.fields.creationSource,
scheduledTaskId: OrchestrationV2ConversationMessage.fields.scheduledTaskId,
senderThreadId: OrchestrationV2ConversationMessage.fields.senderThreadId,
});
export type ProviderAdapterV2TurnMessage = typeof ProviderAdapterV2TurnMessage.Type;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -313,6 +313,9 @@ export const layer: Layer.Layer<
...(message.scheduledTaskId === undefined
? {}
: { scheduledTaskId: message.scheduledTaskId }),
...(message.senderThreadId === undefined
? {}
: { senderThreadId: message.senderThreadId }),
},
})
.pipe(
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1151,6 +1151,9 @@ export const layer: Layer.Layer<
...(message.scheduledTaskId === undefined
? {}
: { scheduledTaskId: message.scheduledTaskId }),
...(message.senderThreadId === undefined
? {}
: { senderThreadId: message.senderThreadId }),
},
modelSelection: run.modelSelection,
runtimePolicy: resolvedRuntimePolicy,
Expand Down
Loading
Loading