mirror of
https://github.com/wu736139669/hapi.git
synced 2026-08-05 06:24:37 +00:00
feat(web,hub): cancel queued messages (#568)
This commit is contained in:
@@ -17,6 +17,7 @@ vi.mock('@/agent/sessionFactory', () => ({
|
||||
onUserMessage: vi.fn((handler) => {
|
||||
harness.userMessageHandler = handler
|
||||
}),
|
||||
onCancelQueuedMessage: vi.fn(),
|
||||
keepAlive: vi.fn(),
|
||||
sendSessionEvent: vi.fn(),
|
||||
sendAgentMessage: vi.fn(),
|
||||
|
||||
@@ -56,6 +56,12 @@ export async function runAgentSession(opts: {
|
||||
messageQueue.push(formattedText, {}, localId);
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[agent] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
let currentPermissionMode: SessionPermissionMode = opts.permissionMode ?? sessionInfo.permissionMode ?? 'default';
|
||||
|
||||
const backend: AgentBackend = AgentRegistry.create(opts.agentType);
|
||||
|
||||
@@ -81,6 +81,7 @@ export class ApiSessionClient extends EventEmitter {
|
||||
private readonly socket: Socket<ServerToClientEvents, ClientToServerEvents>
|
||||
private pendingMessages: { message: UserMessage; localId?: string }[] = []
|
||||
private pendingMessageCallback: ((message: UserMessage, localId?: string) => void) | null = null
|
||||
private cancelQueuedMessageCallback: ((localId: string) => boolean) | null = null
|
||||
private lastSeenMessageSeq: number | null = null
|
||||
private backfillInFlight: Promise<void> | null = null
|
||||
private needsBackfill = false
|
||||
@@ -200,7 +201,7 @@ export class ApiSessionClient extends EventEmitter {
|
||||
this.terminalManager.close(payload.terminalId)
|
||||
}))
|
||||
|
||||
this.socket.on('update', (data: Update) => {
|
||||
this.socket.on('update', (data: Update, ack?: (response: { removed: boolean }) => void) => {
|
||||
try {
|
||||
if (!data.body) return
|
||||
|
||||
@@ -209,6 +210,14 @@ export class ApiSessionClient extends EventEmitter {
|
||||
return
|
||||
}
|
||||
|
||||
if (data.body.t === 'cancel-queued-message') {
|
||||
const removed = (data.body.localId && this.cancelQueuedMessageCallback)
|
||||
? this.cancelQueuedMessageCallback(data.body.localId)
|
||||
: false
|
||||
ack?.({ removed })
|
||||
return
|
||||
}
|
||||
|
||||
if (data.body.t === 'update-session') {
|
||||
if (data.body.metadata && data.body.metadata.version > this.metadataVersion) {
|
||||
const parsed = MetadataSchema.safeParse(data.body.metadata.value)
|
||||
@@ -253,6 +262,10 @@ export class ApiSessionClient extends EventEmitter {
|
||||
}
|
||||
}
|
||||
|
||||
onCancelQueuedMessage(callback: (localId: string) => boolean): void {
|
||||
this.cancelQueuedMessageCallback = callback
|
||||
}
|
||||
|
||||
private enqueueUserMessage(message: UserMessage, localId?: string): void {
|
||||
if (this.pendingMessageCallback) {
|
||||
this.pendingMessageCallback(message, localId)
|
||||
|
||||
@@ -332,6 +332,12 @@ export async function runClaude(options: StartOptions = {}): Promise<void> {
|
||||
logger.debugLargeJson('User message pushed to queue:', message)
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[claude] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
const resolvePermissionMode = (value: unknown): PermissionMode => {
|
||||
const parsed = PermissionModeSchema.safeParse(value);
|
||||
if (!parsed.success || !isPermissionModeAllowedForFlavor(parsed.data, 'claude')) {
|
||||
|
||||
@@ -205,6 +205,12 @@ export async function runCodex(opts: {
|
||||
});
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[codex] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
const formatFailureReason = (message: string): string => {
|
||||
const maxLength = 200;
|
||||
if (message.length <= maxLength) {
|
||||
|
||||
@@ -86,6 +86,12 @@ export async function runCursor(opts: {
|
||||
messageQueue.push(formattedText, enhancedMode, localId);
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[cursor] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
const resolvePermissionMode = (value: unknown): PermissionMode => {
|
||||
const parsed = PermissionModeSchema.safeParse(value);
|
||||
if (!parsed.success || !isPermissionModeAllowedForFlavor(parsed.data, 'cursor')) {
|
||||
|
||||
@@ -14,6 +14,7 @@ const harness = vi.hoisted(() => ({
|
||||
geminiLoopError: null as Error | null,
|
||||
session: {
|
||||
onUserMessage: vi.fn(),
|
||||
onCancelQueuedMessage: vi.fn(),
|
||||
rpcHandlerManager: {
|
||||
registerHandler: vi.fn()
|
||||
}
|
||||
|
||||
@@ -128,6 +128,12 @@ export async function runGemini(opts: {
|
||||
messageQueue.push(formattedText, mode, localId);
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[gemini] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
const resolvePermissionMode = (value: unknown): PermissionMode => {
|
||||
const parsed = PermissionModeSchema.safeParse(value);
|
||||
if (!parsed.success || !isPermissionModeAllowedForFlavor(parsed.data, 'gemini')) {
|
||||
|
||||
@@ -14,6 +14,7 @@ const harness = vi.hoisted(() => ({
|
||||
opencodeLoopError: null as Error | null,
|
||||
session: {
|
||||
onUserMessage: vi.fn(),
|
||||
onCancelQueuedMessage: vi.fn(),
|
||||
rpcHandlerManager: {
|
||||
registerHandler: vi.fn()
|
||||
}
|
||||
|
||||
@@ -108,6 +108,12 @@ export async function runOpencode(opts: {
|
||||
messageQueue.push(formattedText, mode, localId);
|
||||
});
|
||||
|
||||
session.onCancelQueuedMessage((localId) => {
|
||||
const removed = messageQueue.cancelByLocalId(localId);
|
||||
logger.debug(`[opencode] cancelByLocalId(${localId}): ${removed ? 'removed' : 'not found (best-effort)'}`);
|
||||
return removed;
|
||||
});
|
||||
|
||||
const resolvePermissionMode = (value: unknown): PermissionMode => {
|
||||
const parsed = PermissionModeSchema.safeParse(value);
|
||||
if (!parsed.success || !isPermissionModeAllowedForFlavor(parsed.data, 'opencode')) {
|
||||
|
||||
@@ -489,6 +489,68 @@ describe('MessageQueue2', () => {
|
||||
expect(consumedCount).toBe(0);
|
||||
});
|
||||
|
||||
describe('cancelByLocalId', () => {
|
||||
it('should remove the message with matching localId and return true', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
queue.push('msg1', 'local', 'id-abc');
|
||||
queue.push('msg2', 'local', 'id-def');
|
||||
|
||||
const removed = queue.cancelByLocalId('id-abc');
|
||||
expect(removed).toBe(true);
|
||||
expect(queue.size()).toBe(1);
|
||||
expect(queue.queue[0].localId).toBe('id-def');
|
||||
});
|
||||
|
||||
it('should return false when localId is not found', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
queue.push('msg1', 'local', 'id-abc');
|
||||
|
||||
const removed = queue.cancelByLocalId('id-nonexistent');
|
||||
expect(removed).toBe(false);
|
||||
expect(queue.size()).toBe(1);
|
||||
});
|
||||
|
||||
it('should return false when queue is empty', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
const removed = queue.cancelByLocalId('id-abc');
|
||||
expect(removed).toBe(false);
|
||||
});
|
||||
|
||||
it('should not remove a message without localId even if localId param matches empty string', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
queue.push('msg-no-localid', 'local'); // no localId
|
||||
|
||||
const removed = queue.cancelByLocalId('');
|
||||
expect(removed).toBe(false);
|
||||
expect(queue.size()).toBe(1);
|
||||
});
|
||||
|
||||
it('should only remove the first matching localId when duplicates exist', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
queue.push('msg1', 'local', 'id-dup');
|
||||
queue.push('msg2', 'local', 'id-dup');
|
||||
|
||||
const removed = queue.cancelByLocalId('id-dup');
|
||||
expect(removed).toBe(true);
|
||||
expect(queue.size()).toBe(1);
|
||||
// msg2 still remains
|
||||
expect(queue.queue[0].message).toBe('msg2');
|
||||
});
|
||||
|
||||
it('should not affect messages without localId when cancelling by id', () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
queue.push('msg-no-id', 'local');
|
||||
queue.push('msg-with-id', 'local', 'target-id');
|
||||
queue.push('msg-no-id-2', 'local');
|
||||
|
||||
const removed = queue.cancelByLocalId('target-id');
|
||||
expect(removed).toBe(true);
|
||||
expect(queue.size()).toBe(2);
|
||||
expect(queue.queue[0].message).toBe('msg-no-id');
|
||||
expect(queue.queue[1].message).toBe('msg-no-id-2');
|
||||
});
|
||||
});
|
||||
|
||||
it('should differentiate between pushImmediate and pushIsolateAndClear behavior', async () => {
|
||||
const queue = new MessageQueue2<{ type: string }>((mode) => mode.type);
|
||||
|
||||
|
||||
@@ -182,6 +182,20 @@ export class MessageQueue2<T> {
|
||||
logger.debug(`[MessageQueue2] unshift() completed. Queue size: ${this.queue.length}`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove the first queued message that matches the given localId.
|
||||
* Returns true if a message was removed, false if not found.
|
||||
* Best-effort: if the CLI is offline when cancel is issued, the message
|
||||
* may already have been collected for invocation and won't be found here.
|
||||
*/
|
||||
cancelByLocalId(localId: string): boolean {
|
||||
if (!localId) return false;
|
||||
const idx = this.queue.findIndex(item => item.localId === localId);
|
||||
if (idx === -1) return false;
|
||||
this.queue.splice(idx, 1);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reset the queue - clears all messages and resets to empty state
|
||||
*/
|
||||
|
||||
Reference in New Issue
Block a user