Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ describe('ChatHttpService', () => {
// attaching a Bearer manually.
{ provide: BffSessionService, useValue: { csrfHeaders: vi.fn().mockReturnValue({}), handleUnauthorized: vi.fn() } },
{ provide: SessionService, useValue: { currentSession: signal({ sessionId: 's1' }), updateSessionTitleInCache: vi.fn(), getSessionMetadata: vi.fn().mockResolvedValue({}), isNewSession: vi.fn().mockReturnValue(false) } },
{ provide: StreamParserService, useValue: { getCurrentStreamId: vi.fn().mockReturnValue('stream-1'), parseEventSourceMessage: vi.fn() } },
{ provide: StreamParserService, useValue: { getCurrentStreamId: vi.fn().mockReturnValue('stream-1'), parseEventSourceMessage: vi.fn(), hasReceivedDone: vi.fn().mockReturnValue(false) } },
{ provide: ChatStateService, useValue: { abortRequest: vi.fn(), setChatLoading: vi.fn(), setLastTurnInterrupted: vi.fn(), seedSessionAggregates: vi.fn(), createAbortController: vi.fn().mockReturnValue(new AbortController()), releaseAbortController: vi.fn(), streamingSessionIds: vi.fn().mockReturnValue([]) } },
{ provide: MessageMapService, useValue: { endStreaming: vi.fn() } },
{ provide: ErrorService, useValue: { handleHttpError: vi.fn(), addError: vi.fn() } },
Expand Down Expand Up @@ -212,6 +212,128 @@ describe('ChatHttpService', () => {
vi.restoreAllMocks();
});

describe('Stop racing the end of the turn', () => {
// `done` means the server's turn is over, but loading (and so the Stop
// button) only clears on the transport's close, which can trail it — on a
// first turn `session_title` may arrive after `done`. A Stop in that
// window must not stamp a false "interrupted" marker on a finished turn.

const encoder = new TextEncoder();
let controller: AbortController;
let parser: any;
let fetchSpy: any;

/**
* Serve `/chat/stream` from a body that emits `frames` and then stays
* open — the transport has not closed, so `onclose` has not run.
*/
function openStream(frames: string): void {
fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation(async (input) => {
if (!String(input).endsWith('/chat/stream')) {
return new Response(null, { status: 204 });
}
const body = new ReadableStream<Uint8Array>({
start(c) {
c.enqueue(encoder.encode(frames));
},
});
return new Response(body, {
status: 200,
headers: { 'content-type': 'text/event-stream' },
});
});
}

function interruptPosts(): [string, RequestInit][] {
return fetchSpy.mock.calls.filter(([url]: [string]) => String(url).includes('/interrupt'));
}

beforeEach(() => {
controller = new AbortController();
chatStateService.createAbortController.mockReturnValue(controller);
chatStateService.abortRequest.mockImplementation(() => controller.abort());

// Stand-in for the parser's own bookkeeping: record `done` per session.
parser = TestBed.inject(StreamParserService);
const doneFor = new Set<string>();
parser.parseEventSourceMessage.mockImplementation((sessionId: string, event: string) => {
if (event === 'done') doneFor.add(sessionId);
});
parser.hasReceivedDone.mockImplementation((sessionId: string) => doneFor.has(sessionId));
});

afterEach(() => {
vi.restoreAllMocks();
});

it('a Stop after done but before close tears down the transport without marking the turn interrupted', async () => {
const messageMap = TestBed.inject(MessageMapService) as any;
openStream('event: message_start\ndata: {"role":"assistant"}\n\nevent: done\ndata: {}\n\n');

const streaming = service.sendChatRequest({ session_id: 's1', message: 'hi' });
await vi.waitFor(() => expect(parser.hasReceivedDone('s1')).toBe(true));
// Still open: the stream's own teardown has not run.
expect(chatStateService.releaseAbortController).not.toHaveBeenCalled();
expect(chatStateService.setChatLoading).not.toHaveBeenCalled();

service.cancelChatRequest('s1');
await streaming;

expect(interruptPosts()).toHaveLength(0);
expect(chatStateService.setLastTurnInterrupted).not.toHaveBeenCalled();
// The transport is still torn down and the UI leaves its streaming state.
expect(controller.signal.aborted).toBe(true);
expect(messageMap.endStreaming).toHaveBeenCalledWith('s1');
expect(chatStateService.setChatLoading).toHaveBeenCalledWith('s1', false);
});

it('a Stop mid-stream still marks the turn interrupted', async () => {
openStream('event: message_start\ndata: {"role":"assistant"}\n\n');

const streaming = service.sendChatRequest({ session_id: 's1', message: 'hi' });
await vi.waitFor(() =>
expect(parser.parseEventSourceMessage).toHaveBeenCalledWith(
's1',
'message_start',
expect.anything(),
'stream-1',
),
);

service.cancelChatRequest('s1');
await streaming;

const posts = interruptPosts();
expect(posts).toHaveLength(1);
expect(JSON.parse(String(posts[0][1].body))).toEqual({ reason: 'user_stopped' });
expect(chatStateService.setLastTurnInterrupted).toHaveBeenCalledWith('s1', true, 'user_stopped');
expect(controller.signal.aborted).toBe(true);
});

it('runs the title fallback the skipped onclose would have, and no interrupted-turn cost refresh', async () => {
vi.useFakeTimers();
try {
vi.spyOn(globalThis, 'fetch').mockResolvedValue(new Response(null, { status: 204 }));
parser.hasReceivedDone.mockReturnValue(true);
const sessionSvc = TestBed.inject(SessionService) as any;
sessionSvc.isNewSession.mockReturnValue(true);
sessionSvc.applyServerTitle = vi.fn();
sessionSvc.getSessionMetadata.mockResolvedValue({ title: 'Finished Turn Title' });

service.cancelChatRequest('s1');
await vi.runAllTimersAsync();

// Aborting skips onclose, whose job on a first turn is to pull a
// title that `session_title` may not have delivered yet.
expect(sessionSvc.applyServerTitle).toHaveBeenCalledWith('s1', 'Finished Turn Title');
expect(sessionSvc.getSessionMetadata).toHaveBeenCalledTimes(1);
expect(chatStateService.seedSessionAggregates).not.toHaveBeenCalled();
} finally {
vi.useRealTimers();
}
});
});

describe('page-departure attribution', () => {
// A refresh / tab close / navigation is the one interruption cause only
// the browser witnesses. Unattested it lands server-side as
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -350,9 +350,25 @@ export class ChatHttpService {
// that will never be loaded again. Abort the transport and tear down
// locally; that is the whole of "stop" for a session nobody persists.
if (isPreviewSession(sessionId)) {
this.chatStateService.abortRequest(sessionId);
this.messageMapService.endStreaming(sessionId);
this.chatStateService.setChatLoading(sessionId, false);
this.stopLocally(sessionId);
return;
}

// The server already sent `done`: the turn finished, and the socket is
// only open because the transport's close hasn't landed yet (widest on a
// first turn, where `session_title` can trail `done`). The Stop button is
// still showing because loading clears on close, so a click here is
// real — but there is nothing left to interrupt. Signalling
// `user_stopped` would stamp a false "interrupted" marker on a complete
// answer (and a false interruption note on the next prompt), so just
// close the transport. Aborting skips `onclose`, so run its title
// fallback here; the aggregates re-fetch below is unneeded, because the
// turn's `metadata` event arrived before `done`.
if (this.streamParserService.hasReceivedDone(sessionId)) {
this.stopLocally(sessionId);
if (this.sessionService.isNewSession(sessionId)) {
void this.refreshTitleFromServer(sessionId);
}
return;
}

Expand All @@ -370,9 +386,7 @@ export class ChatHttpService {
// refresh-survival source of truth.
this.chatStateService.setLastTurnInterrupted(sessionId, true, 'user_stopped');

this.chatStateService.abortRequest(sessionId);
this.messageMapService.endStreaming(sessionId);
this.chatStateService.setChatLoading(sessionId, false);
this.stopLocally(sessionId);

// Aborting the fetch cut the socket before the stream's terminal
// `metadata` SSE (usage / cost / context) could arrive, so the session
Expand All @@ -385,6 +399,17 @@ export class ChatHttpService {
setTimeout(() => void this.refreshAggregatesAfterStop(sessionId), 900);
}

/**
* Abort a session's transport and tear its streaming state down here.
* fetch-event-source calls neither `onclose` nor `onerror` on abort, so
* the stream's own `finalizeStream` never runs for a stopped stream.
*/
private stopLocally(sessionId: string): void {
this.chatStateService.abortRequest(sessionId);
this.messageMapService.endStreaming(sessionId);
this.chatStateService.setChatLoading(sessionId, false);
}

/**
* Re-seed a session's cost/context badge after a Stop, once the backend's
* interruption teardown has persisted the partial turn's metadata (which
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -660,3 +660,78 @@ describe('StreamParserService - session_title events', () => {
expect(applyServerTitle).not.toHaveBeenCalled();
});
});

describe('StreamParserService - hasReceivedDone (Stop guard)', () => {
// The Stop path reads this to decide whether a click can still interrupt
// anything. It must mean exactly "the server sent this stream's `done`":
// per session, per stream, and not tripped by a client-side parse error.
let service: StreamParserService;

beforeEach(() => {
TestBed.resetTestingModule();
TestBed.configureTestingModule({
providers: [
StreamParserService,
ChatStateService,
ErrorService,
QuotaWarningService,
{ provide: SessionService, useValue: { applyServerTitle: vi.fn() } },
],
});
service = TestBed.inject(StreamParserService);
service.reset('s1');
});

afterEach(() => {
TestBed.resetTestingModule();
});

it('is false mid-stream and true once done arrives', () => {
service.parseEventSourceMessage('s1', 'message_start', { role: 'assistant' });
expect(service.hasReceivedDone('s1')).toBe(false);

service.parseEventSourceMessage('s1', 'done', null);
expect(service.hasReceivedDone('s1')).toBe(true);
});

it('resets to false when the session starts a new stream', () => {
service.parseEventSourceMessage('s1', 'done', null);
service.reset('s1');
expect(service.hasReceivedDone('s1')).toBe(false);
});

it('is scoped to its session', () => {
service.reset('s2');
service.parseEventSourceMessage('s1', 'done', null);
expect(service.hasReceivedDone('s2')).toBe(false);
expect(service.hasReceivedDone('unknown')).toBe(false);
});

it('ignores a done from a superseded stream', () => {
const oldStreamId = service.getCurrentStreamId('s1');
service.reset('s1');
service.parseEventSourceMessage('s1', 'done', null, oldStreamId);
expect(service.hasReceivedDone('s1')).toBe(false);
});

it('is not set by a client parse error — the server turn may still be running', () => {
// A delta with no active message puts the parser in its Error state.
service.parseEventSourceMessage('s1', 'content_block_delta', {
contentBlockIndex: 0,
type: 'text',
text: 'x',
});
expect(service.errorFor('s1')()).not.toBeNull();
expect(service.hasReceivedDone('s1')).toBe(false);
});

it('still records done after a parse error, even though the state gate drops the event', () => {
service.parseEventSourceMessage('s1', 'content_block_delta', {
contentBlockIndex: 0,
type: 'text',
text: 'x',
});
service.parseEventSourceMessage('s1', 'done', null);
expect(service.hasReceivedDone('s1')).toBe(true);
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -136,9 +136,23 @@ interface ParserSessionState {
/** Error state */
error: WritableSignal<string | null>;

/** Stream completion state */
/**
* Parser-side completion: true once `done` was handled OR the parser gave
* up on the stream (`setError`). Drives rendering only — a client parse
* error says nothing about whether the server's turn is still running.
*/
isStreamComplete: WritableSignal<boolean>;

/**
* True once the server sent this stream's `done` frame — the turn is over
* server-side, whatever the parser made of it. Recorded before the state
* gate, so a parser already in `Error` still learns the turn ended.
* Read by the Stop path: a Stop after this point must not mark the
* finished turn interrupted. Reset with the rest of the state on every
* new stream.
*/
doneReceived: boolean;

/** Metadata (usage, metrics) from the stream */
metadata: WritableSignal<MetadataEvent | null>;

Expand Down Expand Up @@ -198,7 +212,6 @@ export class StreamParserService {
private readonly lastEventAtCache = new Map<string, Signal<number>>();
private readonly citationsCache = new Map<string, Signal<Citation[]>>();
private readonly errorCache = new Map<string, Signal<string | null>>();
private readonly isStreamCompleteCache = new Map<string, Signal<boolean>>();

// =========================================================================
// Public API
Expand Down Expand Up @@ -248,9 +261,16 @@ export class StreamParserService {
return this.cachedAccessor(this.errorCache, sessionId, (state) => state.error(), null);
}

/** Stream completion state for a session. */
isStreamCompleteFor(sessionId: string): Signal<boolean> {
return this.cachedAccessor(this.isStreamCompleteCache, sessionId, (state) => state.isStreamComplete(), false);
/**
* Whether a session's CURRENT stream has received the server's `done`.
*
* Scoped to the stream the parser was last reset for, so a new turn reads
* false until its own `done`. Deliberately not derived from
* `isStreamComplete`, which a client-side parse error also sets while the
* server turn keeps running — a Stop then is a real interruption.
*/
hasReceivedDone(sessionId: string): boolean {
return this.states().get(sessionId)?.doneReceived ?? false;
}

/**
Expand Down Expand Up @@ -305,6 +325,12 @@ export class StreamParserService {
// its replacement looking alive.
state.lastEventAt.set(Date.now());

// Like liveness, the end of the server's turn is a transport fact the
// state gate below must not swallow (a parser in `Error` drops `done`).
if (event === 'done') {
state.doneReceived = true;
}

// Validate inputs
if (!event || typeof event !== 'string') {
this.setError(state, 'parseEventSourceMessage: event must be a non-empty string');
Expand Down Expand Up @@ -429,6 +455,7 @@ export class StreamParserService {
lastEventAt: signal<number>(Date.now()),
error: signal<string | null>(null),
isStreamComplete,
doneReceived: false,
metadata: signal<MetadataEvent | null>(null),
pendingCitations: signal<Citation[]>([]),
currentMessage: computed<Message | null>(() => {
Expand Down Expand Up @@ -1019,6 +1046,7 @@ export class StreamParserService {
private handleDone(state: ParserSessionState): void {
this.finalizeCurrentMessage(state);
state.isStreamComplete.set(true);
state.doneReceived = true;
state.modelRetry.set(null);
// "Using list_courses" on a finished turn is a lie, not a stale nicety.
// Durations and summaries already recorded are untouched.
Expand Down