Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
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
22 changes: 22 additions & 0 deletions apps/mobile/src/lib/threadActivity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1090,4 +1090,26 @@ describe("quiet timeline: nested agents", () => {
{ type: "activity-group", id: "nested-done" },
]);
});

it("filters reasoning messages at thread-feed derivation entry", () => {
const thread = makeThread({
id: ThreadId.make("thread-reasoning"),
projectId: ProjectId.make("project-1"),
title: "Reasoning",
messages: [
{
id: MessageId.make("reasoning-1"),
role: "assistant",
channel: "reasoning",
text: "Private reasoning",
turnId: TurnId.make("turn-1"),
streaming: false,
createdAt: "2026-04-01T00:00:01.000Z",
updatedAt: "2026-04-01T00:00:01.000Z",
},
],
});

expect(buildThreadFeed(thread)).toEqual([]);
});
});
6 changes: 4 additions & 2 deletions apps/mobile/src/lib/threadActivity.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import {
ApprovalRequestId,
isReasoningMessage,
isToolLifecycleItemType,
ProviderApprovalOption,
ProviderRequestKind,
Expand Down Expand Up @@ -1730,12 +1731,13 @@ export function buildThreadFeed(
readonly localMessages?: ReadonlyArray<OrchestrationThread["messages"][number]>;
},
): ThreadFeedEntry[] {
const loadedMessages = options?.loadedMessages ?? thread.messages;
const sourceMessages = options?.loadedMessages ?? thread.messages;
const loadedMessages = sourceMessages.filter((message) => !isReasoningMessage(message));
const messages = options?.localMessages
? [...loadedMessages, ...options.localMessages]
: loadedMessages;
const oldestLoadedMessageCreatedAt =
options?.loadedMessages !== undefined ? (loadedMessages[0]?.createdAt ?? null) : null;
options?.loadedMessages !== undefined ? (sourceMessages[0]?.createdAt ?? null) : null;
const workLogEntries = deriveWorkLogEntries(thread.activities);
const entries = Arr.sortWith(
[
Expand Down
125 changes: 74 additions & 51 deletions apps/server/src/orchestration/Layers/CheckpointReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ import * as ManagedRuntime from "effect/ManagedRuntime";
import * as PubSub from "effect/PubSub";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import { afterEach, describe, expect, it, vi } from "vite-plus/test";
import { afterEach, describe, expect, it, vi } from "@effect/vitest";

import * as CheckpointStore from "../../checkpointing/CheckpointStore.ts";
import * as VcsDriverRegistry from "../../vcs/VcsDriverRegistry.ts";
Expand Down Expand Up @@ -490,6 +490,28 @@ describe("CheckpointReactor", () => {
checkpointRefForThreadTurn(ThreadId.make("thread-1"), 0),
);

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.message.assistant.complete",
commandId: CommandId.make("cmd-checkpoint-answer-complete"),
threadId: ThreadId.make("thread-1"),
messageId: MessageId.make("answer-1"),
turnId: asTurnId("turn-1"),
createdAt: "2026-01-01T00:00:00.200Z",
}),
);
await Effect.runPromise(
harness.engine.dispatch({
type: "thread.message.assistant.complete",
commandId: CommandId.make("cmd-checkpoint-reasoning-complete"),
threadId: ThreadId.make("thread-1"),
messageId: MessageId.make("reasoning-1"),
channel: "reasoning",
turnId: asTurnId("turn-1"),
createdAt: "2026-01-01T00:00:00.400Z",
}),
);

NodeFS.writeFileSync(NodePath.join(harness.cwd, "README.md"), "v2\n", "utf8");
harness.provider.emit({
type: "turn.completed",
Expand All @@ -502,7 +524,14 @@ describe("CheckpointReactor", () => {
payload: { state: "completed" },
});

await waitForEvent(harness.engine, (event) => event.type === "thread.turn-diff-completed");
const events = await waitForEvent(
harness.engine,
(event) => event.type === "thread.turn-diff-completed",
);
const checkpointEvent = events.find((event) => event.type === "thread.turn-diff-completed");
if (checkpointEvent?.type === "thread.turn-diff-completed") {
expect(checkpointEvent.payload.assistantMessageId).toBe(MessageId.make("answer-1"));
}
const thread = await waitForThread(
harness.readModel,
(entry) => entry.latestTurn?.turnId === "turn-1" && entry.checkpoints.length === 1,
Expand Down Expand Up @@ -1152,12 +1181,12 @@ describe("CheckpointReactor", () => {
});
});

it("processes consecutive revert requests with deterministic rollback sequencing", async () => {
const harness = await createHarness();
const createdAt = "2026-01-01T00:00:00.000Z";
it.effect("processes consecutive revert requests with deterministic rollback sequencing", () =>
Effect.gen(function* () {
const harness = yield* Effect.promise(() => createHarness());
const createdAt = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
harness.engine.dispatch({
yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-set-inline-revert"),
threadId: ThreadId.make("thread-1"),
Expand All @@ -1171,11 +1200,9 @@ describe("CheckpointReactor", () => {
updatedAt: createdAt,
},
createdAt,
}),
);
});

await Effect.runPromise(
harness.engine.dispatch({
yield* harness.engine.dispatch({
type: "thread.turn.diff.complete",
commandId: CommandId.make("cmd-inline-revert-diff-1"),
threadId: ThreadId.make("thread-1"),
Expand All @@ -1186,10 +1213,8 @@ describe("CheckpointReactor", () => {
files: [],
checkpointTurnCount: 1,
createdAt,
}),
);
await Effect.runPromise(
harness.engine.dispatch({
});
yield* harness.engine.dispatch({
type: "thread.turn.diff.complete",
commandId: CommandId.make("cmd-inline-revert-diff-2"),
threadId: ThreadId.make("thread-1"),
Expand All @@ -1200,62 +1225,60 @@ describe("CheckpointReactor", () => {
files: [],
checkpointTurnCount: 2,
createdAt,
}),
);
});

await Effect.runPromise(
harness.engine.dispatch({
yield* harness.engine.dispatch({
type: "thread.checkpoint.revert",
commandId: CommandId.make("cmd-sequenced-revert-request-1"),
threadId: ThreadId.make("thread-1"),
turnCount: 1,
createdAt,
}),
);
await Effect.runPromise(
harness.engine.dispatch({
});
yield* harness.engine.dispatch({
type: "thread.checkpoint.revert",
commandId: CommandId.make("cmd-sequenced-revert-request-0"),
threadId: ThreadId.make("thread-1"),
turnCount: 0,
createdAt,
}),
);
});

await harness.drain();
yield* Effect.promise(() => harness.drain());

expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(2);
expect(harness.provider.rollbackConversation.mock.calls[0]?.[0]).toEqual({
threadId: ThreadId.make("thread-1"),
numTurns: 1,
});
expect(harness.provider.rollbackConversation.mock.calls[1]?.[0]).toEqual({
threadId: ThreadId.make("thread-1"),
numTurns: 1,
});
});
expect(harness.provider.rollbackConversation).toHaveBeenCalledTimes(2);
expect(harness.provider.rollbackConversation.mock.calls[0]?.[0]).toEqual({
threadId: ThreadId.make("thread-1"),
numTurns: 1,
});
expect(harness.provider.rollbackConversation.mock.calls[1]?.[0]).toEqual({
threadId: ThreadId.make("thread-1"),
numTurns: 1,
});
}),
);

it("appends an error activity when revert is requested without an active session", async () => {
const harness = await createHarness({ hasSession: false });
const createdAt = "2026-01-01T00:00:00.000Z";
it.effect("appends an error activity when revert is requested without an active session", () =>
Effect.gen(function* () {
const harness = yield* Effect.promise(() => createHarness({ hasSession: false }));
const createdAt = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
harness.engine.dispatch({
yield* harness.engine.dispatch({
type: "thread.checkpoint.revert",
commandId: CommandId.make("cmd-revert-no-session"),
threadId: ThreadId.make("thread-1"),
turnCount: 1,
createdAt,
}),
);
});

const thread = await waitForThread(harness.readModel, (entry) =>
entry.activities.some((activity) => activity.kind === "checkpoint.revert.failed"),
);
const thread = yield* Effect.promise(() =>
waitForThread(harness.readModel, (entry) =>
entry.activities.some((activity) => activity.kind === "checkpoint.revert.failed"),
),
);

expect(thread.activities.some((activity) => activity.kind === "checkpoint.revert.failed")).toBe(
true,
);
expect(harness.provider.rollbackConversation).not.toHaveBeenCalled();
});
expect(
thread.activities.some((activity) => activity.kind === "checkpoint.revert.failed"),
).toBe(true);
expect(harness.provider.rollbackConversation).not.toHaveBeenCalled();
}),
);
});
10 changes: 9 additions & 1 deletion apps/server/src/orchestration/Layers/CheckpointReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,12 @@ import {
type CheckpointRef,
EventId,
MessageId,
isReasoningMessage,
type ProjectId,
ThreadId,
TurnId,
type OrchestrationEvent,
type OrchestrationMessageChannel,
type ProviderRuntimeEvent,
type VcsStatusLocalResult,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -225,6 +227,7 @@ const make = Effect.gen(function* () {
readonly messages: ReadonlyArray<{
readonly id: MessageId;
readonly role: string;
readonly channel?: OrchestrationMessageChannel | undefined;
readonly turnId: TurnId | null;
}>;
};
Expand Down Expand Up @@ -298,7 +301,12 @@ const make = Effect.gen(function* () {
input.assistantMessageId ??
input.thread.messages
.toReversed()
.find((entry) => entry.role === "assistant" && entry.turnId === input.turnId)?.id ??
.find(
(entry) =>
entry.role === "assistant" &&
!isReasoningMessage(entry) &&
entry.turnId === input.turnId,
)?.id ??
MessageId.make(`assistant:${input.turnId}`);

yield* orchestrationEngine.dispatch({
Expand Down
Loading
Loading