mirror of
https://github.com/wu736139669/hapi.git
synced 2026-08-05 06:24:37 +00:00
feat(web): show queued status for messages pending inference (#492)
This commit is contained in:
@@ -426,6 +426,69 @@ describe('MessageQueue2', () => {
|
||||
expect(batch3?.mode.type).toBe('A');
|
||||
});
|
||||
|
||||
it('should call onBatchConsumed with collected localIds', async () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
const received: string[][] = [];
|
||||
queue.onBatchConsumed = (localIds) => { received.push(localIds); };
|
||||
|
||||
queue.push('message1', 'local', 'id1');
|
||||
queue.push('message2', 'local', 'id2');
|
||||
|
||||
await queue.waitForMessagesAndGetAsString();
|
||||
expect(received).toEqual([['id1', 'id2']]);
|
||||
|
||||
// Push more with a different mode and consume again
|
||||
queue.push('message3', 'remote', 'id3');
|
||||
await queue.waitForMessagesAndGetAsString();
|
||||
expect(received).toEqual([['id1', 'id2'], ['id3']]);
|
||||
});
|
||||
|
||||
it('should report localIds batch-by-batch when modes differ', async () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
const received: string[][] = [];
|
||||
queue.onBatchConsumed = (localIds) => { received.push(localIds); };
|
||||
|
||||
// Two messages land in different batches because their mode hashes differ.
|
||||
queue.push('first', 'A', 'id1');
|
||||
queue.push('second', 'B', 'id2');
|
||||
|
||||
const batch1 = await queue.waitForMessagesAndGetAsString();
|
||||
expect(batch1?.message).toBe('first');
|
||||
expect(received).toEqual([['id1']]);
|
||||
// Second message still waiting in the queue.
|
||||
expect(queue.size()).toBe(1);
|
||||
|
||||
const batch2 = await queue.waitForMessagesAndGetAsString();
|
||||
expect(batch2?.message).toBe('second');
|
||||
expect(received).toEqual([['id1'], ['id2']]);
|
||||
expect(queue.size()).toBe(0);
|
||||
});
|
||||
|
||||
it('should skip onBatchConsumed when batch has no localIds', async () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
let called = false;
|
||||
queue.onBatchConsumed = () => { called = true; };
|
||||
|
||||
// Push without localIds (e.g., internal commands that do not need UI ack)
|
||||
queue.push('internal', 'local');
|
||||
await queue.waitForMessagesAndGetAsString();
|
||||
expect(called).toBe(false);
|
||||
});
|
||||
|
||||
it('should not call onBatchConsumed when collectBatch returns null', async () => {
|
||||
const queue = new MessageQueue2<string>(mode => mode);
|
||||
let consumedCount = 0;
|
||||
queue.onBatchConsumed = () => { consumedCount++; };
|
||||
|
||||
// Close queue while waiting — should return null
|
||||
const waitPromise = queue.waitForMessagesAndGetAsString();
|
||||
queue.close();
|
||||
const result = await waitPromise;
|
||||
|
||||
expect(result).toBeNull();
|
||||
expect(consumedCount).toBe(0);
|
||||
});
|
||||
|
||||
it('should differentiate between pushImmediate and pushIsolateAndClear behavior', async () => {
|
||||
const queue = new MessageQueue2<{ type: string }>((mode) => mode.type);
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ interface QueueItem<T> {
|
||||
message: string;
|
||||
mode: T;
|
||||
modeHash: string;
|
||||
localId?: string;
|
||||
isolate?: boolean; // If true, this message must be processed alone
|
||||
}
|
||||
|
||||
@@ -16,6 +17,7 @@ export class MessageQueue2<T> {
|
||||
private waiter: ((hasMessages: boolean) => void) | null = null;
|
||||
private closed = false;
|
||||
private onMessageHandler: ((message: string, mode: T) => void) | null = null;
|
||||
onBatchConsumed: ((localIds: string[]) => void) | null = null;
|
||||
modeHasher: (mode: T) => string;
|
||||
|
||||
constructor(
|
||||
@@ -37,7 +39,7 @@ export class MessageQueue2<T> {
|
||||
/**
|
||||
* Push a message to the queue with a mode.
|
||||
*/
|
||||
push(message: string, mode: T): void {
|
||||
push(message: string, mode: T, localId?: string): void {
|
||||
if (this.closed) {
|
||||
throw new Error('Cannot push to closed queue');
|
||||
}
|
||||
@@ -49,6 +51,7 @@ export class MessageQueue2<T> {
|
||||
message,
|
||||
mode,
|
||||
modeHash,
|
||||
localId,
|
||||
isolate: false
|
||||
});
|
||||
|
||||
@@ -72,7 +75,7 @@ export class MessageQueue2<T> {
|
||||
* Push a message immediately without batching delay.
|
||||
* Does not clear the queue or enforce isolation.
|
||||
*/
|
||||
pushImmediate(message: string, mode: T): void {
|
||||
pushImmediate(message: string, mode: T, localId?: string): void {
|
||||
if (this.closed) {
|
||||
throw new Error('Cannot push to closed queue');
|
||||
}
|
||||
@@ -84,6 +87,7 @@ export class MessageQueue2<T> {
|
||||
message,
|
||||
mode,
|
||||
modeHash,
|
||||
localId,
|
||||
isolate: false
|
||||
});
|
||||
|
||||
@@ -108,7 +112,7 @@ export class MessageQueue2<T> {
|
||||
* Clears any pending messages and ensures this message is never batched with others.
|
||||
* Used for special commands that require dedicated processing.
|
||||
*/
|
||||
pushIsolateAndClear(message: string, mode: T): void {
|
||||
pushIsolateAndClear(message: string, mode: T, localId?: string): void {
|
||||
if (this.closed) {
|
||||
throw new Error('Cannot push to closed queue');
|
||||
}
|
||||
@@ -123,6 +127,7 @@ export class MessageQueue2<T> {
|
||||
message,
|
||||
mode,
|
||||
modeHash,
|
||||
localId,
|
||||
isolate: true
|
||||
});
|
||||
|
||||
@@ -145,7 +150,7 @@ export class MessageQueue2<T> {
|
||||
/**
|
||||
* Push a message to the beginning of the queue with a mode.
|
||||
*/
|
||||
unshift(message: string, mode: T): void {
|
||||
unshift(message: string, mode: T, localId?: string): void {
|
||||
if (this.closed) {
|
||||
throw new Error('Cannot unshift to closed queue');
|
||||
}
|
||||
@@ -157,6 +162,7 @@ export class MessageQueue2<T> {
|
||||
message,
|
||||
mode,
|
||||
modeHash,
|
||||
localId,
|
||||
isolate: false
|
||||
});
|
||||
|
||||
@@ -252,6 +258,7 @@ export class MessageQueue2<T> {
|
||||
|
||||
const firstItem = this.queue[0];
|
||||
const sameModeMessages: string[] = [];
|
||||
const consumedLocalIds: string[] = [];
|
||||
let mode = firstItem.mode;
|
||||
let isolate = firstItem.isolate ?? false;
|
||||
const targetModeHash = firstItem.modeHash;
|
||||
@@ -260,6 +267,7 @@ export class MessageQueue2<T> {
|
||||
if (firstItem.isolate) {
|
||||
const item = this.queue.shift()!;
|
||||
sameModeMessages.push(item.message);
|
||||
if (item.localId) consumedLocalIds.push(item.localId);
|
||||
logger.debug(`[MessageQueue2] Collected isolated message with mode hash: ${targetModeHash}`);
|
||||
} else {
|
||||
// Collect all messages with the same mode until we hit an isolated message
|
||||
@@ -268,6 +276,7 @@ export class MessageQueue2<T> {
|
||||
!this.queue[0].isolate) {
|
||||
const item = this.queue.shift()!;
|
||||
sameModeMessages.push(item.message);
|
||||
if (item.localId) consumedLocalIds.push(item.localId);
|
||||
}
|
||||
logger.debug(`[MessageQueue2] Collected batch of ${sameModeMessages.length} messages with mode hash: ${targetModeHash}`);
|
||||
}
|
||||
@@ -275,6 +284,10 @@ export class MessageQueue2<T> {
|
||||
// Join all messages with newlines
|
||||
const combinedMessage = sameModeMessages.join('\n');
|
||||
|
||||
if (consumedLocalIds.length > 0) {
|
||||
this.onBatchConsumed?.(consumedLocalIds);
|
||||
}
|
||||
|
||||
return {
|
||||
message: combinedMessage,
|
||||
mode,
|
||||
|
||||
Reference in New Issue
Block a user