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
5 changes: 5 additions & 0 deletions .changeset/batch-queued-prompts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@moonshot-ai/kimi-code": minor
---

Prompts sent while the agent is working are now delivered together: when a turn ends with several prompts queued, the model receives all of them at once in a single turn instead of answering them one turn each in FIFO order.
142 changes: 107 additions & 35 deletions packages/agent-core-v2/src/agent/loop/loopService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -140,14 +140,13 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
private readonly errorHandlers: LoopErrorHandler[] = [];
private readonly promptWaiters = new Map<string, PromptWaiter>();
private readonly steered = new Map<string, SteeredPrompt>();
private readonly batched = new Map<string, BatchedPrompt>();
private readonly terminalStates = new Map<string, PromptState>();
private readonly pendingSubmissions: UserEntry[] = [];
private readonly nudges: Nudge[] = [];
private nudgeCursor = 0;
private active: ActiveTurn | undefined;
private pendingMachineTurn:
| { readonly id: number; readonly queueItemId?: string; readonly entry?: UserEntry }
| undefined;
private pendingMachineTurn: PendingMachineTurn | undefined;
private machineTurnSuppressed = false;
private readonly settleWaiters: Array<() => void> = [];
private quiescenceDepth = 0;
Expand Down Expand Up @@ -286,6 +285,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
}
this.promptWaiters.clear();
this.steered.clear();
this.batched.clear();
this.pendingSubmissions.length = 0;
const active = this.active;
active?.turn.cancel(reason);
Expand Down Expand Up @@ -410,6 +410,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {

private promptStateOf(id: string): PromptState {
if (this.active?.prompt.id === id) return 'running';
if (this.batched.has(id)) return 'running';
if (this.steered.has(id)) return 'steered';
return this.terminalStates.get(id) ?? 'pending';
}
Expand All @@ -419,12 +420,14 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
if (failedEntry !== undefined) return projectionFromEntry(failedEntry);
const active = this.active;
if (active !== undefined && active.prompt.id === id) return active.prompt;
const batched = this.batched.get(id);
if (batched !== undefined) return batched;
const steered = this.steered.get(id);
if (steered !== undefined) return steered;
const pending = this.pendingMachineTurn;
if (pending?.queueItemId === id && pending.entry !== undefined) {
return projectionFromEntry(pending.entry);
}
const pendingEntry = pending?.entries?.find((entry) => entry.meta?.promptId === id)
?? (pending?.queueItemId === id ? pending.entry : undefined);
if (pendingEntry !== undefined) return projectionFromEntry(pendingEntry);
const queued = this.engine
?.snapshot()
.queue.find((item) => item.meta?.promptId === id);
Expand Down Expand Up @@ -557,6 +560,25 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
this.promptWaiters.delete(waiter.id);
}

private settleBatchedCompletion(waiter: PromptWaiter, result: TurnResult): void {
const prompt = this.batched.get(waiter.id);
if (prompt === undefined) return;
const state =
result.type === 'cancelled' ? 'cancelled' : result.type === 'failed' ? 'failed' : 'completed';
waiter.completion.resolve({
promptId: waiter.id,
result,
state,
});
if (prompt.tracked) {
if (state === 'cancelled') this.publishPromptAborted(waiter.id);
else this.publishPromptCompleted(waiter.id, state);
}
this.terminalStates.set(waiter.id, state);
this.promptWaiters.delete(waiter.id);
this.batched.delete(waiter.id);
}

private async materializeDaemonRefs(message: {
readonly content: readonly ContentPart[];
}): Promise<void> {
Expand Down Expand Up @@ -726,6 +748,9 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
if (active !== undefined && active.prompt.tracked && active.prompt.id === target.promptId) {
return this.cancelActiveTurn(undefined, cancellation);
}
if (active !== undefined && this.batched.get(target.promptId)?.parentTurnId === active.id) {
return this.cancelActiveTurn(undefined, cancellation);
}
const waiter = this.promptWaiters.get(target.promptId);
if (waiter === undefined) {
throw new Error2(ErrorCodes.PROMPT_NOT_FOUND, `prompt ${target.promptId} not found`);
Expand All @@ -749,6 +774,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
this.terminalStates.set(waiter.id, 'cancelled');
this.promptWaiters.delete(waiter.id);
this.steered.delete(waiter.id);
this.batched.delete(waiter.id);
return true;
}

Expand Down Expand Up @@ -839,7 +865,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
}

private settleUnboundRecord(
pending: { readonly id: number; readonly queueItemId?: string; readonly entry?: UserEntry },
pending: PendingMachineTurn,
outcome: { readonly outcome: MachineTurnOutcome; readonly error?: unknown },
): void {
const active = this.active;
Expand All @@ -849,6 +875,20 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
});
return;
}
const pendingEntries = pending.entries ?? (pending.entry === undefined ? [] : [pending.entry]);
const [primaryEntry, ...batchedEntries] = this.livePendingEntries(pending);
if (primaryEntry !== undefined) {
const promptId = primaryEntry.meta?.promptId as string;
const waiter = this.promptWaiters.get(promptId) as PromptWaiter;
const boundTurn = this.beginActiveTurn(waiter, primaryEntry, pending.id);
waiter.onMaterialize?.();
this.materializeMessage(this.gatedProjectionMessage(boundTurn.prompt));
this.settlePromptLaunched(waiter, boundTurn);
this.bindBatchedEntries(batchedEntries, boundTurn);
this.endPreGateTurn(boundTurn, outcome);
return;
}
if (pendingEntries.length > 0) return;
if (pending.queueItemId === undefined) {
const seeded = this.nudges.slice(this.nudgeCursor).find(
(nudge) => !nudge.dropped && nudge.contextMessage !== undefined && nudge.contextMessage.content.length > 0,
Expand All @@ -870,13 +910,6 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
this.endPreGateTurn(seededTurn, outcome);
return;
}
const waiter = this.promptWaiters.get(pending.queueItemId);
if (waiter === undefined || pending.entry === undefined) return;
const boundTurn = this.beginActiveTurn(waiter, pending.entry, pending.id);
waiter.onMaterialize?.();
this.materializeMessage(this.gatedProjectionMessage(boundTurn.prompt));
this.settlePromptLaunched(waiter, boundTurn);
this.endPreGateTurn(boundTurn, outcome);
}

private endPreGateTurn(
Expand Down Expand Up @@ -1050,35 +1083,35 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return { type: 'fail' };
}

private bindMachineTurn(pending: {
readonly id: number;
readonly queueItemId?: string;
readonly entry?: UserEntry;
}): boolean {
private bindMachineTurn(pending: PendingMachineTurn): boolean {
if (this.active !== undefined) {
if (pending.queueItemId !== undefined && pending.entry !== undefined) {
const waiter = this.promptWaiters.get(pending.queueItemId);
if (waiter !== undefined) {
this.machineEngine().submit({
message: this.gatedEntryMessage(pending.entry),
meta: pending.entry.meta,
});
}
const entries = pending.entries ?? (pending.entry === undefined ? [] : [pending.entry]);
for (const entry of entries) {
const promptId = entry.meta?.promptId;
if (promptId === undefined || !this.promptWaiters.has(promptId)) continue;
this.machineEngine().submit({
message: this.gatedEntryMessage(entry),
meta: entry.meta,
});
}
return true;
}
if (pending.queueItemId !== undefined) {
const waiter = this.promptWaiters.get(pending.queueItemId);
if (waiter === undefined || pending.entry === undefined) {
this.machineTurnSuppressed = true;
return false;
}
const boundTurn = this.beginActiveTurn(waiter, pending.entry, pending.id);
const pendingEntries = pending.entries ?? (pending.entry === undefined ? [] : [pending.entry]);
const [primaryEntry, ...batchedEntries] = this.livePendingEntries(pending);
if (primaryEntry !== undefined) {
const promptId = primaryEntry.meta?.promptId as string;
const waiter = this.promptWaiters.get(promptId) as PromptWaiter;
const boundTurn = this.beginActiveTurn(waiter, primaryEntry, pending.id);
waiter.onMaterialize?.();
this.materializeMessage(this.gatedProjectionMessage(boundTurn.prompt));
this.settlePromptLaunched(waiter, boundTurn);
this.bindBatchedEntries(batchedEntries, boundTurn);
return true;
}
if (pendingEntries.length > 0 || pending.queueItemId !== undefined) {
this.machineTurnSuppressed = true;
return false;
}
const seeded = this.nudges.slice(this.nudgeCursor).find(
(nudge) => !nudge.dropped && nudge.contextMessage !== undefined && nudge.contextMessage.content.length > 0,
);
Expand All @@ -1097,7 +1130,34 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return true;
}

private gatedProjectionMessage(prompt: ActivePrompt): ContextMessage {
private livePendingEntries(pending: PendingMachineTurn): UserEntry[] {
const entries = pending.entries ?? (pending.entry === undefined ? [] : [pending.entry]);
return entries.filter((entry) => {
const promptId = entry.meta?.promptId;
return promptId !== undefined && this.promptWaiters.has(promptId);
});
}

private bindBatchedEntries(entries: readonly UserEntry[], turn: ActiveTurn): void {
for (const entry of entries) {
const promptId = entry.meta?.promptId;
if (promptId === undefined) continue;
const waiter = this.promptWaiters.get(promptId);
if (waiter === undefined) continue;
const prompt: BatchedPrompt = {
...projectionFromEntry(entry),
parentTurnId: turn.id,
};
this.batched.set(promptId, prompt);
waiter.onMaterialize?.();
this.materializeMessage(this.gatedProjectionMessage(prompt));
waiter.launched.resolve(turn.turn);
void turn.turn.result.then((result) => this.settleBatchedCompletion(waiter, result));
if (prompt.tracked) this.publishPromptStarted(promptId, prompt.origin);
}
}

private gatedProjectionMessage(prompt: PromptProjection): ContextMessage {
if (!prompt.tracked) return prompt.message;
return {
...prompt.message,
Expand Down Expand Up @@ -1279,6 +1339,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
id: event.machineTurnId,
queueItemId: event.queueItemId,
entry: event.entry,
entries: event.entries,
};
this.machineTurnSuppressed = false;
return;
Expand Down Expand Up @@ -2150,6 +2211,13 @@ interface PromptWaiter {
failedEntry?: UserEntry;
}

interface PendingMachineTurn {
readonly id: number;
readonly queueItemId?: string;
readonly entry?: UserEntry;
readonly entries?: readonly UserEntry[];
}

interface PromptProjection {
readonly tracked: boolean;
readonly origin: PromptOrigin;
Expand All @@ -2167,6 +2235,10 @@ interface SteeredPrompt extends PromptProjection {
readonly parentId: string;
}

interface BatchedPrompt extends PromptProjection {
readonly parentTurnId: number;
}

const EMPTY_HANDLE_MESSAGE: ContextMessage = {
role: 'user',
content: [],
Expand Down
16 changes: 14 additions & 2 deletions packages/agent-core-v2/src/agent/loop/machine/engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,13 @@ export type MachineEngineDelta =
export type MachineTurnOutcome = 'done' | 'failed' | 'aborted';

export type MachineEngineEvent =
| { readonly type: 'turnStarted'; readonly machineTurnId: number; readonly queueItemId?: string; readonly entry?: UserEntry }
| {
readonly type: 'turnStarted';
readonly machineTurnId: number;
readonly queueItemId?: string;
readonly entry?: UserEntry;
readonly entries: readonly UserEntry[];
}
| {
readonly type: 'turnSettled';
readonly outcome: MachineTurnOutcome;
Expand Down Expand Up @@ -366,7 +372,13 @@ export function attachMachineEngine(
split = createDeltaSplitter();
pendingFailure = undefined;
lastRetry = undefined;
publish({ type: 'turnStarted', machineTurnId: event.turnId, queueItemId: event.queueItemId, entry: event.entry });
publish({
type: 'turnStarted',
machineTurnId: event.turnId,
queueItemId: event.queueItemId,
entry: event.entry,
entries: event.entries,
});
}),
ref.on('step.started', (event) => {
currentStep = event.step;
Expand Down
Loading