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/server/src/auth/rpcForkScopes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ export const FORK_RPC_REQUIRED_SCOPES = {
[WS_FORK_METHODS.threadCommentsResolveAll]: AuthOrchestrationOperateScope,
[WS_FORK_METHODS.threadCommentsRemove]: AuthOrchestrationOperateScope,
[WS_FORK_METHODS.threadCommentsSetDeliveryPaused]: AuthOrchestrationOperateScope,
[WS_FORK_METHODS.threadCommentsResend]: AuthOrchestrationOperateScope,
// T3-CUSTOM(expbkt3): Claude account profiles per thread. Reading the host
// snapshot or a thread's account is a read; choosing the account is operate.
[WS_FORK_METHODS.claudeAccountsGetThread]: AuthOrchestrationReadScope,
Expand Down
6 changes: 3 additions & 3 deletions apps/server/src/mcp/toolkits/control/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -696,7 +696,7 @@ export const T3ShowUiTool = readonlyTool(
export const T3ListCommentsTool = readonlyTool(
Tool.make("t3_list_comments", {
description:
"List the review comments the user left on your earlier messages in this session. A comment quotes a passage of one of your messages and carries either free text or a reaction: good (keep this as it is), okay (acceptable, no change needed) or remove (drop this / do not do this). Status `open` means the comment is an active instruction for you and is re-sent with every new turn until you mark it addressed with t3_reply_comment; `addressed` means you reported it done and the user has not closed it yet; `resolved` means the user closed it. Defaults to open comments only.",
"List the review comments the user left on your earlier messages in this session. A comment quotes a passage of one of your messages and carries either free text or a reaction: good (keep this as it is), okay (acceptable, no change needed) or remove (drop this / do not do this). Status `open` means the comment is an active instruction for you until you mark it addressed with t3_reply_comment; each comment reaches you once, and again only when the user changes it, replies to it or sends the open comments again; `addressed` means you reported it done and the user has not closed it yet; `resolved` means the user closed it. Defaults to open comments only.",
parameters: Schema.Struct({
status: Schema.optional(
described(
Expand All @@ -714,7 +714,7 @@ export const T3ListCommentsTool = readonlyTool(
export const T3ReplyCommentTool = mutatingTool(
Tool.make("t3_reply_comment", {
description:
"Reply to one of the user's review comments on this session, and/or mark it addressed. Call it once per comment you have handled: pass `addressed: true` when the request in the comment is done, with a short `body` saying what you did (or, for a reaction, that you took note). Use a `body` without `addressed` to ask a question or explain why you are not doing it; the comment then stays open and is re-sent next turn. Only the user can resolve a comment; an addressed comment waits for them, and a user reply reopens it. Pass at least one of `body` or `addressed`.",
"Reply to one of the user's review comments on this session, and/or mark it addressed. Call it once per comment you have handled: pass `addressed: true` when the request in the comment is done, with a short `body` saying what you did (or, for a reaction, that you took note). Use a `body` without `addressed` to ask a question or explain why you are not doing it; the comment then stays open, and the user's reply sends it back to you. Only the user can resolve a comment; an addressed comment waits for them, and a user reply reopens it. Pass at least one of `body` or `addressed`.",
parameters: Schema.Struct({
commentId: described(
Schema.String,
Expand All @@ -729,7 +729,7 @@ export const T3ReplyCommentTool = mutatingTool(
addressed: Schema.optional(
described(
Schema.Boolean,
"Set to true once you have done what the comment asks; it stops the comment being re-sent to you. Omit or pass false to reply without closing it.",
"Set to true once you have done what the comment asks; the comment then waits for the user to resolve it. Omit or pass false to reply without closing it.",
),
),
}),
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/mcp/toolkits/webUi/catalog.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ const invocation = (

it("generates one unique virtual tool and complete schemas for every web RPC", () => {
// T3-CUSTOM(expbkt3): the catalog includes native V2 methods and every retained fork RPC.
expect(WEB_UI_VIRTUAL_TOOL_COUNT).toBe(230);
expect(WEB_UI_VIRTUAL_TOOL_COUNT).toBe(231);
expect(WEB_UI_STREAM_TOOL_COUNT).toBe(30);
expect(WEB_UI_VIRTUAL_TOOL_COUNT).toBe(WsRpcGroup.requests.size);
expect(new Set(WEB_UI_VIRTUAL_TOOLS.map((tool) => tool.name)).size).toBe(
Expand Down
4 changes: 2 additions & 2 deletions apps/server/src/mcp/toolkits/webUi/registration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,9 +74,9 @@ it.effect("registers four compact tools while listing the complete virtual surfa
expect(listed.structuredContent).toMatchObject({
ok: true,
// T3-CUSTOM(expbkt3): registration exposes native V2 methods and all retained fork RPCs.
rpcCount: 230,
rpcCount: 231,
streamCount: 30,
matchedCount: 230,
matchedCount: 231,
});

const schema = yield* withInvocation(
Expand Down
19 changes: 12 additions & 7 deletions apps/server/src/orchestration-v2/RunExecutionService.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// T3-CUSTOM(expbkt3): open review comments reach each new provider turn.
import { appendOpenThreadComments } from "../threadcomments/turnContext.ts";
import { prepareOpenThreadComments } from "../threadcomments/turnContext.ts";
// T3-CUSTOM(expbkt3): capture this optional fork service before native workers hide their construction context.
import { ThreadCommentsService } from "../threadcomments/ThreadCommentsService.ts";
// T3-CUSTOM(expbkt3): the server plan policy reaches every provider.
Expand Down Expand Up @@ -1371,17 +1371,20 @@ export const layer: Layer.Layer<
const compact =
input.message.attachments.length === 0 &&
input.message.text.trim().toLowerCase() === "/compact";
const commentText = Option.isSome(threadComments)
? appendOpenThreadComments(input.run.threadId, input.message.text).pipe(
Effect.provideService(ThreadCommentsService, threadComments.value),
)
: Effect.succeed(input.message.text);
// T3-CUSTOM(expbkt3): BEGIN — comments are marked sent only once the turn starts.
const comments =
Option.isSome(threadComments) && !compact
? yield* prepareOpenThreadComments(input.run.threadId, input.message.text).pipe(
Effect.provideService(ThreadCommentsService, threadComments.value),
)
: { text: input.message.text, markSent: Effect.void };
// T3-CUSTOM(expbkt3): END
const providerMessage = compact
? input.message
: // T3-CUSTOM(expbkt3): command history stays unchanged; the agent gets the plan policy.
{
...input.message,
text: appendAgentPlanInstructions(yield* commentText, agentPlanSubmissionEnabled),
text: appendAgentPlanInstructions(comments.text, agentPlanSubmissionEnabled),
};
const turnInput = {
appThread: input.appThread,
Expand Down Expand Up @@ -1417,6 +1420,8 @@ export const layer: Layer.Layer<
))
: input.session.startTurn(turnInput);
yield* Effect.andThen(shouldStart, startTurn).pipe(
// T3-CUSTOM(expbkt3): a turn that never starts leaves its comments unsent.
Effect.tap(() => comments.markSent),
Effect.catchCause((cause) =>
Effect.logError("orchestration V2 provider turn start failed", {
runId: input.run.id,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ it.layer(NodeSqliteClient.layer({ filename: ":memory:" }))("fork migration ledge
[1045, "SessionWebhooks"],
[1046, "ScheduledTaskWebhooks"],
[1047, "WebhookRelayDeliveries"],
[1048, "ThreadCommentLastSent"],
]);
const after = yield* sql<{
migration_id: number;
Expand Down
4 changes: 4 additions & 0 deletions apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,8 @@ import Migration1047 from "./Migrations/058_WebhookRelayDeliveries.ts";
// row already present in effect_sql_migrations.
// T3-CUSTOM(expbkt3): owner-bound session callback storage.
import Migration1045 from "./Migrations/1045_SessionWebhooks.ts";
// T3-CUSTOM(expbkt3): review comments are sent once; this records when.
import Migration1048 from "./Migrations/1048_ThreadCommentLastSent.ts";
const migrationEntries = [
[1, "OrchestrationEvents", Migration0001],
[2, "OrchestrationCommandReceipts", Migration0002],
Expand Down Expand Up @@ -317,6 +319,8 @@ const migrationEntries = [
// T3-CUSTOM(expbkt3): upstream 57-58 remapped above the shipped 1045 migration.
[1046, "ScheduledTaskWebhooks", Migration1046],
[1047, "WebhookRelayDeliveries", Migration1047],
// T3-CUSTOM(expbkt3): review comment send-once mark.
[1048, "ThreadCommentLastSent", Migration1048],
] as const;

export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ layer("055_OrchestrationV2", (it) => {
...Array.from({ length: 45 }, (_, index) => index + 1),
// T3-CUSTOM(expbkt3): append the session webhook ledger entry at 1045 and
// upstream's scheduled-task webhook migrations (57-58) at 1046-1047.
...Array.from({ length: 48 }, (_, index) => 1000 + index),
// T3-CUSTOM(expbkt3): 1048 adds the review comment send-once mark.
...Array.from({ length: 49 }, (_, index) => 1000 + index),
],
);
}),
Expand All @@ -35,6 +36,7 @@ layer("055_OrchestrationV2", (it) => {
[1045, "SessionWebhooks"], // T3-CUSTOM(expbkt3): durable callback schema.
[1046, "ScheduledTaskWebhooks"],
[1047, "WebhookRelayDeliveries"],
[1048, "ThreadCommentLastSent"], // T3-CUSTOM(expbkt3): comment send-once mark.
]);
assert.deepStrictEqual(yield* runMigrations(), []);

Expand Down Expand Up @@ -65,6 +67,7 @@ layer("055_OrchestrationV2", (it) => {
{ migration_id: 1045, name: "SessionWebhooks" }, // T3-CUSTOM(expbkt3): append, never rewrite.
{ migration_id: 1046, name: "ScheduledTaskWebhooks" },
{ migration_id: 1047, name: "WebhookRelayDeliveries" },
{ migration_id: 1048, name: "ThreadCommentLastSent" }, // T3-CUSTOM(expbkt3)
]);

// T3-CUSTOM(expbkt3): verify callback destinations and durable deliveries after the full upgrade.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
// T3-CUSTOM(expbkt3): when a review comment last went to the agent. A comment
// goes once; it is sent again only after the user edits, replies to or reopens
// it, or asks for a re-send. Nullable: existing comments count as not sent, so
// each open one goes once more and is then marked.
import * as Effect from "effect/Effect";
import * as SqlClient from "effect/sql/SqlClient";

export default Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const columns = yield* sql<{ readonly name: string }>`
PRAGMA table_info(thread_comments)
`;

if (!columns.some((column) => column.name === "last_sent_at")) {
yield* sql`
ALTER TABLE thread_comments
ADD COLUMN last_sent_at TEXT
`;
}
});
68 changes: 64 additions & 4 deletions apps/server/src/persistence/ThreadComments.ts
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,20 @@ export class ThreadCommentsRepository extends Context.Service<
readonly threadId: ThreadId;
readonly paused: boolean;
}) => Effect.Effect<void, ThreadCommentsRepositoryError>;
/**
* Records that these comments went to the agent at `sentAt`. Each row is
* marked only while it is unchanged since it was read (same `updatedAt`)
* and still unsent, so a reply that landed in between keeps it unsent.
*/
readonly markSent: (input: {
readonly threadId: ThreadId;
readonly comments: ReadonlyArray<Pick<ThreadComment, "commentId" | "updatedAt">>;
readonly sentAt: string;
}) => Effect.Effect<void, ThreadCommentsRepositoryError>;
/** Clears the sent mark on every open comment, so the next turn sends them again. */
readonly clearSentForOpen: (
threadId: ThreadId,
) => Effect.Effect<void, ThreadCommentsRepositoryError>;
}
>()("t3/persistence/ThreadComments/ThreadCommentsRepository") {}

Expand All @@ -118,7 +132,8 @@ export const make = Effect.gen(function* () {
replies_json AS "replies",
created_at AS "createdAt",
updated_at AS "updatedAt",
resolved_at AS "resolvedAt"
resolved_at AS "resolvedAt",
last_sent_at AS "lastSentAt"
`;

const listRows = SqlSchema.findAll({
Expand All @@ -145,12 +160,14 @@ export const make = Effect.gen(function* () {
execute: (comment) => sql`
INSERT INTO thread_comments (
comment_id, thread_id, number, kind, anchor_json, body, status,
author_user_id, author_label, replies_json, created_at, updated_at, resolved_at
author_user_id, author_label, replies_json, created_at, updated_at, resolved_at,
last_sent_at
) VALUES (
${comment.commentId}, ${comment.threadId}, ${comment.number}, ${comment.kind},
${JSON.stringify(comment.anchor)}, ${comment.body}, ${comment.status},
${comment.authorUserId}, ${comment.authorLabel}, ${JSON.stringify(comment.replies)},
${comment.createdAt}, ${comment.updatedAt}, ${comment.resolvedAt}
${comment.createdAt}, ${comment.updatedAt}, ${comment.resolvedAt},
${comment.lastSentAt ?? null}
)
ON CONFLICT(comment_id) DO NOTHING
`,
Expand All @@ -164,7 +181,8 @@ export const make = Effect.gen(function* () {
status = ${comment.status},
replies_json = ${JSON.stringify(comment.replies)},
updated_at = ${comment.updatedAt},
resolved_at = ${comment.resolvedAt}
resolved_at = ${comment.resolvedAt},
last_sent_at = ${comment.lastSentAt ?? null}
WHERE comment_id = ${comment.commentId} AND thread_id = ${comment.threadId}
`,
});
Expand Down Expand Up @@ -228,6 +246,30 @@ export const make = Effect.gen(function* () {
`,
});

const markSentRow = SqlSchema.void({
Request: Schema.Struct({
threadId: ThreadId,
commentId: ThreadCommentId,
updatedAt: Schema.String,
sentAt: Schema.String,
}),
execute: ({ threadId, commentId, updatedAt, sentAt }) => sql`
UPDATE thread_comments SET last_sent_at = ${sentAt}
WHERE thread_id = ${threadId}
AND comment_id = ${commentId}
AND updated_at = ${updatedAt}
AND last_sent_at IS NULL
`,
});

const clearSentForOpenRows = SqlSchema.void({
Request: Schema.Struct({ threadId: ThreadId }),
execute: ({ threadId }) => sql`
UPDATE thread_comments SET last_sent_at = NULL
WHERE thread_id = ${threadId} AND status = 'open'
`,
});

return ThreadCommentsRepository.of({
listForThread: (threadId) =>
listRows({ threadId }).pipe(Effect.mapError(mapError("ThreadComments.listForThread"))),
Expand Down Expand Up @@ -287,6 +329,24 @@ export const make = Effect.gen(function* () {
setDeliveryPausedRow({ threadId, paused: paused ? 1 : 0 }).pipe(
Effect.mapError(mapError("ThreadComments.setDeliveryPaused")),
),

markSent: ({ threadId, comments, sentAt }) =>
Effect.forEach(
comments,
(comment) =>
markSentRow({
threadId,
commentId: comment.commentId,
updatedAt: comment.updatedAt,
sentAt,
}),
{ discard: true },
).pipe(Effect.mapError(mapError("ThreadComments.markSent"))),

clearSentForOpen: (threadId) =>
clearSentForOpenRows({ threadId }).pipe(
Effect.mapError(mapError("ThreadComments.clearSentForOpen")),
),
});
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ const pendingV2Migrations = [
[1045, "SessionWebhooks"],
[1046, "ScheduledTaskWebhooks"],
[1047, "WebhookRelayDeliveries"],
[1048, "ThreadCommentLastSent"],
] as const;

describe("fork V2 ledger upgrade", () => {
Expand Down Expand Up @@ -55,6 +56,7 @@ describe("fork V2 ledger upgrade", () => {
pendingV2Migrations[2],
pendingV2Migrations[3],
pendingV2Migrations[4],
pendingV2Migrations[5],
]);
assert.deepStrictEqual(yield* runMigrations(), []);
assert.deepStrictEqual(yield* sql`SELECT * FROM orchestration_v2_legacy_imports`, imports);
Expand Down
60 changes: 60 additions & 0 deletions apps/server/src/threadcomments/ThreadCommentsService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,12 @@ import {
} from "@t3tools/contracts";
import { describe, expect, it } from "@effect/vitest";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";

import { MigrationsLive } from "../persistence/Migrations.ts";
import { layerMemory as SqlitePersistenceMemory } from "../persistence/Sqlite.ts";
Expand Down Expand Up @@ -253,6 +255,64 @@ describe("ThreadCommentsService", () => {
),
);

// T3-CUSTOM(expbkt3): each comment goes to the agent once.
it.effect("sends a comment once, and again after a user reply, a reopen or a resend", () =>
withService((service) =>
Effect.gen(function* () {
const added = yield* service.add({
threadId,
kind: "comment",
anchor: anchor(),
body: "one",
...actor,
});
const commentId = added.comments[0]!.commentId;
expect(added.comments[0]?.lastSentAt ?? null).toBeNull();
expect(yield* service.openForDelivery(threadId)).toHaveLength(1);
// What a turn does: read the unsent comments, then mark exactly those.
const sendTurn = Effect.flatMap(service.openForDelivery(threadId), (comments) =>
service.markSent({ threadId, comments }),
);

yield* sendTurn;
expect(yield* service.openForDelivery(threadId)).toEqual([]);
expect((yield* service.snapshot(threadId)).comments[0]?.lastSentAt).not.toBeNull();

// An agent reply alone does not send it again.
yield* service.agentReply({ threadId, commentId, body: "on it" });
expect(yield* service.openForDelivery(threadId)).toEqual([]);

yield* service.reply({ threadId, commentId, body: "not yet", ...actor });
expect(yield* service.openForDelivery(threadId)).toHaveLength(1);

yield* sendTurn;
yield* service.setStatus({ threadId, commentId, status: "resolved" });
yield* service.setStatus({ threadId, commentId, status: "open" });
expect(yield* service.openForDelivery(threadId)).toHaveLength(1);

yield* sendTurn;
const resent = yield* service.resend({ threadId });
expect(resent.comments[0]?.lastSentAt ?? null).toBeNull();
expect(yield* service.openForDelivery(threadId)).toHaveLength(1);
}),
),
);

// T3-CUSTOM(expbkt3): a reply between the turn's read and its mark wins.
it.effect("keeps a comment unsent when it changed after the turn read it", () =>
withService((service) =>
Effect.gen(function* () {
yield* service.add({ threadId, kind: "comment", anchor: anchor(), body: "one", ...actor });
const read = yield* service.openForDelivery(threadId);
// Tests run on the frozen test clock; move it so the reply's timestamp differs.
yield* TestClock.adjust(Duration.seconds(1));
yield* service.reply({ threadId, commentId: read[0]!.commentId, body: "also", ...actor });
yield* service.markSent({ threadId, comments: read });
expect(yield* service.openForDelivery(threadId)).toHaveLength(1);
}),
),
);

it.effect("openForDelivery returns open comments only, and nothing while paused", () =>
withService((service) =>
Effect.gen(function* () {
Expand Down
Loading
Loading