Skip to content
Open
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
40 changes: 40 additions & 0 deletions src/component/messages.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1328,3 +1328,43 @@ describe("late saves racing a failed pending message (issue #320)", () => {
expect(stream?.state.kind).toBe("aborted");
});
});

describe("deleting a message aborts generation writing to it (issue #300)", () => {
test("a stream at the deleted order is aborted", async () => {
const t = initConvexTest();
const thread = await t.mutation(api.threads.createThread, { userId: "u" });
const threadId = thread._id as Id<"threads">;

const { messages } = await t.mutation(api.messages.addMessages, {
threadId,
messages: [{ message: { role: "user", content: "hello" } }],
});
const prompt = messages[0];

await t.mutation(api.streams.create, {
threadId,
order: prompt.order,
stepOrder: prompt.stepOrder + 1,
userId: "u",
agentName: "a",
model: "m",
provider: "p",
format: "UIMessageChunk",
});

await t.mutation(api.messages.deleteByIds, {
messageIds: [prompt._id as Id<"messages">],
});

const streaming = await t.query(api.streams.list, {
threadId,
statuses: ["streaming"],
});
const aborted = await t.query(api.streams.list, {
threadId,
statuses: ["aborted"],
});
expect(streaming).toHaveLength(0);
expect(aborted).toHaveLength(1);
});
});
32 changes: 29 additions & 3 deletions src/component/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ import {
getStreamingMessagesWithMetadata,
finishHandler,
releaseStreamFileOwnershipByIds,
abortStreamsAtOrder,
} from "./streams.js";
import { partial } from "convex-helpers/validators";

Expand All @@ -65,21 +66,45 @@ export async function deleteMessage(
}
}

/**
* Deleting a message strands any generation still writing to its order, which
* would otherwise only surface as a missing-parent failure when that generation
* finalizes. Aborting the stream lets the in-flight run stop on its own.
*/
async function abortStreamsForDeleted(
ctx: MutationCtx,
deleted: (Doc<"messages"> | null)[],
) {
const seen = new Set<string>();
for (const message of deleted) {
if (!message) continue;
const key = `${message.threadId}:${message.order}`;
if (seen.has(key)) continue;
seen.add(key);
await abortStreamsAtOrder(ctx, {
threadId: message.threadId,
order: message.order,
reason: "Message deleted",
});
}
}

export const deleteByIds = mutation({
args: { messageIds: v.array(v.id("messages")) },
returns: v.array(v.id("messages")),
handler: async (ctx, args) => {
const deletedMessageIds = await Promise.all(
const deleted = await Promise.all(
args.messageIds.map(async (id) => {
const message = await ctx.db.get("messages", id);
if (message) {
await deleteMessage(ctx, message);
return id;
return message;
}
return null;
}),
);
return deletedMessageIds.filter((id) => id !== null);
await abortStreamsForDeleted(ctx, deleted);
return deleted.filter((m) => m !== null).map((m) => m._id);
},
});

Expand Down Expand Up @@ -145,6 +170,7 @@ export const deleteByOrder = mutation({
})
.take(64);
await Promise.all(messages.map((m) => deleteMessage(ctx, m)));
await abortStreamsForDeleted(ctx, messages);
return {
isDone: messages.length < 64,
lastOrder: messages.at(-1)?.order,
Expand Down
35 changes: 20 additions & 15 deletions src/component/streams.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,24 +197,29 @@ function publicStreamMessage(m: Doc<"streamingMessages">): StreamMessage {
};
}

export async function abortStreamsAtOrder(
ctx: MutationCtx,
args: { threadId: Id<"threads">; order: number; reason: string },
) {
const streams = await ctx.db
.query("streamingMessages")
.withIndex("threadId_state_order_stepOrder", (q) =>
q
.eq("threadId", args.threadId)
.eq("state.kind", "streaming")
.eq("order", args.order),
)
.take(100);
for (const stream of streams) {
await abortById(ctx, { streamId: stream._id, reason: args.reason });
}
return streams.length > 0;
}

export const abortByOrder = mutation({
args: { threadId: v.id("threads"), order: v.number(), reason: v.string() },
returns: v.boolean(),
handler: async (ctx, args) => {
const streams = await ctx.db
.query("streamingMessages")
.withIndex("threadId_state_order_stepOrder", (q) =>
q
.eq("threadId", args.threadId)
.eq("state.kind", "streaming")
.eq("order", args.order),
)
.take(100);
for (const stream of streams) {
await abortById(ctx, { streamId: stream._id, reason: args.reason });
}
return streams.length > 0;
},
handler: abortStreamsAtOrder,
});

export const abort = mutation({
Expand Down
Loading