@@ -1334,6 +1334,94 @@ describe("TriggerChatTransport", () => {
13341334 } ) ;
13351335 } ) ;
13361336
1337+ describe ( "reconnectToStream stop-on-abort ownership (TRI-13070)" , ( ) => {
1338+ // A quiet stream: EOF, no records, never settled — the subscription
1339+ // stays alive (watch mode) so an abort mid-flight exercises the stop path.
1340+ function quietWatchTransport ( ) : {
1341+ transport : TriggerChatTransport ;
1342+ appends : ( ) => number ;
1343+ } {
1344+ let appendCount = 0 ;
1345+ global . fetch = vi . fn ( ) . mockImplementation ( async ( url : string | URL ) => {
1346+ const urlStr = typeof url === "string" ? url : url . toString ( ) ;
1347+ if ( isSessionStreamAppendUrl ( urlStr ) ) {
1348+ appendCount ++ ;
1349+ return defaultAppendResponse ( ) ;
1350+ }
1351+ if ( isSessionOutSubscribeUrl ( urlStr ) ) return defaultSseResponse ( [ ] ) ;
1352+ throw new Error ( `Unexpected URL: ${ urlStr } ` ) ;
1353+ } ) ;
1354+ const transport = new TriggerChatTransport ( {
1355+ task : "my-chat-task" ,
1356+ accessToken : ( ) => "pat" ,
1357+ watch : true ,
1358+ sessions : { "chat-own" : { publicAccessToken : "p" , isStreaming : true } } ,
1359+ } ) ;
1360+ return { transport, appends : ( ) => appendCount } ;
1361+ }
1362+
1363+ it ( "passive subscriber aborting writes no stop chunk to .in" , async ( ) => {
1364+ vi . useFakeTimers ( ) ;
1365+ try {
1366+ const { transport, appends } = quietWatchTransport ( ) ;
1367+ const abort = new AbortController ( ) ;
1368+ const stream = await transport . reconnectToStream ( {
1369+ chatId : "chat-own" ,
1370+ abortSignal : abort . signal ,
1371+ } ) ;
1372+ const drained = drainChunks ( stream ! ) ;
1373+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1374+ abort . abort ( ) ;
1375+ await drained ;
1376+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1377+ expect ( appends ( ) ) . toBe ( 0 ) ;
1378+ } finally {
1379+ vi . useRealTimers ( ) ;
1380+ }
1381+ } ) ;
1382+
1383+ it ( "owning subscriber with stopOnAbort:true sends a stop chunk on abort" , async ( ) => {
1384+ vi . useFakeTimers ( ) ;
1385+ try {
1386+ const { transport, appends } = quietWatchTransport ( ) ;
1387+ const abort = new AbortController ( ) ;
1388+ const stream = await transport . reconnectToStream ( {
1389+ chatId : "chat-own" ,
1390+ abortSignal : abort . signal ,
1391+ stopOnAbort : true ,
1392+ } ) ;
1393+ const drained = drainChunks ( stream ! ) ;
1394+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1395+ abort . abort ( ) ;
1396+ await drained ;
1397+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1398+ expect ( appends ( ) ) . toBe ( 1 ) ;
1399+ } finally {
1400+ vi . useRealTimers ( ) ;
1401+ }
1402+ } ) ;
1403+
1404+ it ( "abortSignal presence alone (stopOnAbort unset) sends no stop" , async ( ) => {
1405+ vi . useFakeTimers ( ) ;
1406+ try {
1407+ const { transport, appends } = quietWatchTransport ( ) ;
1408+ const abort = new AbortController ( ) ;
1409+ const stream = await transport . reconnectToStream ( {
1410+ chatId : "chat-own" ,
1411+ abortSignal : abort . signal ,
1412+ } ) ;
1413+ const drained = drainChunks ( stream ! ) ;
1414+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1415+ abort . abort ( ) ;
1416+ await drained ;
1417+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1418+ expect ( appends ( ) ) . toBe ( 0 ) ;
1419+ } finally {
1420+ vi . useRealTimers ( ) ;
1421+ }
1422+ } ) ;
1423+ } ) ;
1424+
13371425 describe ( "multi-tab coordination" , ( ) => {
13381426 it ( "isReadOnly defaults to false when multiTab is disabled" , ( ) => {
13391427 const transport = new TriggerChatTransport ( {
0 commit comments