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
35 changes: 16 additions & 19 deletions packages/client/__tests__/protocol-parser.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import type { ProtocolCallbacks, ProtocolParserState } from '../src/protocol-par
function makeState(overrides?: Partial<ProtocolParserState>): ProtocolParserState {
return {
currentSessionId: undefined,
pendingSend: [],
...overrides,
};
}
Expand All @@ -17,7 +16,6 @@ function makeCallbacks(overrides?: Partial<ProtocolCallbacks>): ProtocolCallback
onMessagesRestored: vi.fn(),
onSessionRenamed: vi.fn(),
setWsRunning: vi.fn(),
sendQueued: vi.fn(),
...overrides,
};
}
Expand Down Expand Up @@ -157,21 +155,15 @@ describe('session lifecycle', () => {
]);
});

it('session_end dequeues first pending send and queues it', () => {
const state = makeState({
pendingSend: [
{ type: 'send', prompt: 'follow-up' },
{ type: 'send', prompt: 'second' },
],
});
it('session_end dispatches SESSION_END action', () => {
const cb = makeCallbacks();
const r = parseServerMessage({ type: 'session_end', sessionId: 'sid' }, state, cb, POOL_KEY);
const r = parseServerMessage(
{ type: 'session_end', sessionId: 'sid' },
makeState(),
cb,
POOL_KEY,
);
expect(r.messagesActions).toContainEqual({ type: 'SESSION_END', sessionId: 'sid' });
// Optimistic running=true when draining pending send
expect(r.messagesActions).toContainEqual({ type: 'SESSION_STATE_CHANGED', state: 'running' });
expect(cb.sendQueued).toHaveBeenCalledWith(POOL_KEY, { type: 'send', prompt: 'follow-up' });
// Second message stays queued
expect(state.pendingSend).toEqual([{ type: 'send', prompt: 'second' }]);
});
});

Expand Down Expand Up @@ -378,10 +370,15 @@ describe('error handling', () => {
expect(r.messagesActions).toEqual([{ type: 'ERROR', error: 'Something broke' }]);
});

it('error clears pendingSend queue', () => {
const state = makeState({ pendingSend: [{ type: 'send', prompt: 'test' }] });
parseServerMessage({ type: 'error', error: 'fail' }, state, makeCallbacks(), POOL_KEY);
expect(state.pendingSend).toEqual([]);
it('error does not require pendingSend cleanup (removed in P2)', () => {
const state = makeState();
const r = parseServerMessage(
{ type: 'error', error: 'fail' },
state,
makeCallbacks(),
POOL_KEY,
);
expect(r.messagesActions).toContainEqual({ type: 'ERROR', error: 'fail' });
});
});

Expand Down
72 changes: 3 additions & 69 deletions packages/client/__tests__/store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -402,89 +402,23 @@ describe('sendMessage', () => {
);
});

it('queues second message while first turn is running', async () => {
it('sends second message immediately while first turn is running (server dedup)', async () => {
const store = createReadyStore();

store.getState().sendMessage('first');
lastWs.simulateMessage({ type: 'session_id', sessionId: 'sess-q' });

// Turn is still running — second send should queue
// Turn is still running — P2: second send goes immediately (server deduplicates)
const sentBefore = lastWs.sent.length;
store.getState().sendMessage('second');

// Should NOT have sent yet
const newSendsImmediate = lastWs.sent.slice(sentBefore).map((s) => JSON.parse(s));
expect(newSendsImmediate.filter((m) => m.type === 'send')).toHaveLength(0);

// session_end triggers flush of queued message
lastWs.simulateMessage({ type: 'session_end', sessionId: 'sess-q' });

const newSends = lastWs.sent.slice(sentBefore).map((s) => JSON.parse(s));
expect(newSends).toContainEqual(expect.objectContaining({ type: 'send', prompt: 'second' }));
});

it('queues second message as pendingSend while first turn is active', () => {
const store = createReadyStore();

store.getState().sendMessage('first');
lastWs.simulateMessage({ type: 'session_id', sessionId: 'sess-pend' });

const sentBefore = lastWs.sent.length;
store.getState().sendMessage('second');

// Second message should be queued, not immediately sent
const immediateSends = lastWs.sent
.slice(sentBefore)
.filter((s) => JSON.parse(s).type === 'send');
expect(immediateSends).toHaveLength(0);

// But the optimistic user message should appear in the store
// Optimistic user message should appear in the store
const userMsgs = store.getState().messages.messages.filter((m) => m.role === 'user');
expect(userMsgs).toHaveLength(2);
});

it('cancels pending timeout when session_end arrives in time', () => {
vi.useFakeTimers();
try {
const store = createMitzoStore(makeOptions());
lastWs.completeHandshake();

store.getState().sendMessage('first');
lastWs.simulateMessage({ type: 'session_id', sessionId: 'sess-ok' });

store.getState().sendMessage('second');
lastWs.simulateMessage({ type: 'session_end', sessionId: 'sess-ok' });

const sentAfterEnd = lastWs.sent.length;
vi.advanceTimersByTime(6_000);

expect(lastWs.sent.length).toBe(sentAfterEnd);
} finally {
vi.useRealTimers();
}
});

it('cancels pending timeout on newSession', () => {
vi.useFakeTimers();
try {
const store = createMitzoStore(makeOptions());
lastWs.completeHandshake();

store.getState().sendMessage('first');
lastWs.simulateMessage({ type: 'session_id', sessionId: 'sess-new' });

store.getState().sendMessage('second');
const sentBefore = lastWs.sent.length;

store.getState().newSession();
vi.advanceTimersByTime(6_000);

const flushed = lastWs.sent.slice(sentBefore).filter((s) => JSON.parse(s).type === 'send');
expect(flushed).toHaveLength(0);
} finally {
vi.useRealTimers();
}
});
});

describe('WS → store wiring', () => {
Expand Down
Loading
Loading