1- import { describe , expect , it } from "vitest" ;
1+ import { describe , expect , it , vi } from "vitest" ;
22import type { UIMessage } from "ai" ;
33import { TriggerChatTransport , type TriggerChatTransportOptions } from "../src/v3/chat.js" ;
44
@@ -17,13 +17,18 @@ type BatchRecord = {
1717 headers ?: Array < [ string , string ] > ;
1818} ;
1919
20- function batchResponse ( records : BatchRecord [ ] ) : Response {
20+ function batchResponse ( records : BatchRecord [ ] , settled = false ) : Response {
2121 const frames = records
2222 . map ( ( r ) => `event: batch\ndata: ${ JSON . stringify ( { records : [ r ] } ) } \n\n` )
2323 . join ( "" ) ;
24+ const headers : Record < string , string > = {
25+ "Content-Type" : "text/event-stream" ,
26+ "X-Stream-Version" : "v2" ,
27+ } ;
28+ if ( settled ) headers [ "X-Session-Settled" ] = "true" ;
2429 return new Response ( frames , {
2530 status : 200 ,
26- headers : { "Content-Type" : "text/event-stream" , "X-Stream-Version" : "v2" } ,
31+ headers,
2732 } ) ;
2833}
2934
@@ -42,7 +47,10 @@ function turnComplete(seqNum: number, inCursor: number): BatchRecord {
4247
4348function textDelta ( seqNum : number , text : string ) : BatchRecord {
4449 return {
45- body : JSON . stringify ( { data : { type : "text-delta" , id : "t1" , delta : text } , id : "m1" } ) ,
50+ body : JSON . stringify ( {
51+ data : { type : "text-delta" , id : "t1" , delta : text } ,
52+ id : `m${ seqNum } ` ,
53+ } ) ,
4654 seq_num : seqNum ,
4755 timestamp : seqNum ,
4856 headers : [ ] ,
@@ -88,6 +96,36 @@ async function submit(transport: TriggerChatTransport): Promise<string[]> {
8896}
8997
9098describe ( "transport turn correlation" , ( ) => {
99+ it ( "persists the owned send's input sequence before subscribing" , async ( ) => {
100+ const onSessionChange = vi . fn ( ) ;
101+ const transport = new TriggerChatTransport ( {
102+ task : "test-task" ,
103+ accessToken : async ( ) => "tok_test" ,
104+ sessions : { c1 : { publicAccessToken : "tok_test" , isStreaming : false } } ,
105+ onSessionChange,
106+ fetch : async ( _url , _init , ctx ) =>
107+ ctx . endpoint === "in" ? inResponse ( 5 ) : batchResponse ( [ turnComplete ( 10 , 5 ) ] ) ,
108+ } ) ;
109+
110+ const stream = await transport . sendMessages ( {
111+ trigger : "submit-message" ,
112+ chatId : "c1" ,
113+ messageId : undefined ,
114+ messages : [ user ( "hi" , "u-1" ) ] ,
115+ abortSignal : undefined ,
116+ } ) ;
117+
118+ expect ( onSessionChange ) . toHaveBeenCalledWith ( "c1" , {
119+ publicAccessToken : "tok_test" ,
120+ lastEventId : undefined ,
121+ activeInputSeq : 5 ,
122+ isStreaming : true ,
123+ } ) ;
124+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBe ( 5 ) ;
125+ await readDeltas ( stream ) ;
126+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBeUndefined ( ) ;
127+ } ) ;
128+
91129 it ( "skips an earlier turn's turn-complete and closes on its own" , async ( ) => {
92130 // Append seq 5; the undo turn's complete (cursor 4) must be skipped.
93131 const out = batchResponse ( [ turnComplete ( 10 , 4 ) , textDelta ( 11 , "56" ) , turnComplete ( 12 , 5 ) ] ) ;
@@ -107,4 +145,117 @@ describe("transport turn correlation", () => {
107145 const deltas = await submit ( makeTransport ( out , undefined ) ) ;
108146 expect ( deltas ) . toEqual ( [ ] ) ;
109147 } ) ;
148+
149+ it ( "reuses a hydrated input sequence to skip stale turn-completes after reconnecting" , async ( ) => {
150+ const transport = new TriggerChatTransport ( {
151+ task : "test-task" ,
152+ accessToken : async ( ) => "tok_test" ,
153+ sessions : {
154+ c1 : { publicAccessToken : "tok_test" , isStreaming : true , activeInputSeq : 5 } ,
155+ } ,
156+ fetch : async ( ) =>
157+ batchResponse ( [ turnComplete ( 10 , 4 ) , textDelta ( 11 , "current" ) , turnComplete ( 12 , 5 ) ] ) ,
158+ } ) ;
159+
160+ const stream = await transport . reconnectToStream ( { chatId : "c1" } ) ;
161+
162+ expect ( stream ) . not . toBeNull ( ) ;
163+ await expect ( readDeltas ( stream ! ) ) . resolves . toEqual ( [ "current" ] ) ;
164+ expect ( transport . getSession ( "c1" ) ?. isStreaming ) . toBe ( false ) ;
165+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBeUndefined ( ) ;
166+ } ) ;
167+
168+ it ( "does not request a settled peek while reconnecting a known active input" , async ( ) => {
169+ vi . useFakeTimers ( ) ;
170+ try {
171+ const subscribeHeaders : Headers [ ] = [ ] ;
172+ const transport = new TriggerChatTransport ( {
173+ task : "test-task" ,
174+ accessToken : async ( ) => "tok_test" ,
175+ sessions : {
176+ c1 : { publicAccessToken : "tok_test" , isStreaming : true , activeInputSeq : 5 } ,
177+ } ,
178+ fetch : async ( _url , init ) => {
179+ const headers = new Headers ( init ?. headers ) ;
180+ subscribeHeaders . push ( headers ) ;
181+
182+ if ( subscribeHeaders . length === 1 ) {
183+ // Match the server shortcut: a peek sees the previous turn's
184+ // boundary at the tail and marks this otherwise-normal EOF settled.
185+ return batchResponse ( [ turnComplete ( 10 , 4 ) ] , headers . has ( "X-Peek-Settled" ) ) ;
186+ }
187+
188+ return batchResponse ( [ textDelta ( 11 , "current" ) , turnComplete ( 12 , 5 ) ] ) ;
189+ } ,
190+ } ) ;
191+
192+ const stream = await transport . reconnectToStream ( { chatId : "c1" } ) ;
193+
194+ expect ( stream ) . not . toBeNull ( ) ;
195+ const deltas = readDeltas ( stream ! ) ;
196+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
197+ await expect ( deltas ) . resolves . toEqual ( [ "current" ] ) ;
198+ expect ( subscribeHeaders ) . toHaveLength ( 2 ) ;
199+ expect ( subscribeHeaders [ 0 ] ?. get ( "X-Peek-Settled" ) ) . toBeNull ( ) ;
200+ expect ( transport . getSession ( "c1" ) ?. isStreaming ) . toBe ( false ) ;
201+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBeUndefined ( ) ;
202+ } finally {
203+ vi . useRealTimers ( ) ;
204+ }
205+ } ) ;
206+ it . each ( [ 5 , 6 ] ) (
207+ "accepts a reconnected turn-complete at or after the active input sequence (%i)" ,
208+ async ( inCursor ) => {
209+ const transport = new TriggerChatTransport ( {
210+ task : "test-task" ,
211+ accessToken : async ( ) => "tok_test" ,
212+ sessions : {
213+ c1 : { publicAccessToken : "tok_test" , isStreaming : true , activeInputSeq : 5 } ,
214+ } ,
215+ fetch : async ( ) => batchResponse ( [ turnComplete ( 10 , inCursor ) , textDelta ( 11 , "late" ) ] ) ,
216+ } ) ;
217+
218+ const stream = await transport . reconnectToStream ( { chatId : "c1" } ) ;
219+
220+ expect ( stream ) . not . toBeNull ( ) ;
221+ await expect ( readDeltas ( stream ! ) ) . resolves . toEqual ( [ ] ) ;
222+ expect ( transport . getSession ( "c1" ) ?. isStreaming ) . toBe ( false ) ;
223+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBeUndefined ( ) ;
224+ }
225+ ) ;
226+
227+ it ( "uses the input sequence for one accepted watch turn only" , async ( ) => {
228+ let outCalls = 0 ;
229+ const turnCompleted : number [ ] = [ ] ;
230+ const transport = new TriggerChatTransport ( {
231+ task : "test-task" ,
232+ accessToken : async ( ) => "tok_test" ,
233+ watch : true ,
234+ sessions : {
235+ c1 : { publicAccessToken : "tok_test" , isStreaming : true , activeInputSeq : 5 } ,
236+ } ,
237+ onEvent : ( event ) => {
238+ if ( event . type === "turn-completed" ) turnCompleted . push ( Number ( event . sessionInEventId ) ) ;
239+ } ,
240+ fetch : async ( ) => {
241+ outCalls ++ ;
242+ return outCalls === 1
243+ ? batchResponse ( [
244+ turnComplete ( 10 , 4 ) ,
245+ textDelta ( 11 , "first" ) ,
246+ turnComplete ( 12 , 5 ) ,
247+ textDelta ( 13 , "second" ) ,
248+ turnComplete ( 14 , 4 ) ,
249+ ] )
250+ : batchResponse ( [ ] , true ) ;
251+ } ,
252+ } ) ;
253+
254+ const stream = await transport . reconnectToStream ( { chatId : "c1" } ) ;
255+
256+ expect ( stream ) . not . toBeNull ( ) ;
257+ await expect ( readDeltas ( stream ! ) ) . resolves . toEqual ( [ "first" , "second" ] ) ;
258+ expect ( turnCompleted ) . toEqual ( [ 5 , 4 ] ) ;
259+ expect ( transport . getSession ( "c1" ) ?. activeInputSeq ) . toBeUndefined ( ) ;
260+ } ) ;
110261} ) ;
0 commit comments