@@ -33,6 +33,7 @@ import type {
3333 StreamOptions ,
3434 Usage ,
3535} from "../types.js" ;
36+ import { abortableDelay } from "../utils/abortable-delay.js" ;
3637import {
3738 appendAssistantMessageDiagnostic ,
3839 createAssistantMessageDiagnostic ,
@@ -128,20 +129,6 @@ function codexTransportError(error: unknown): Error {
128129 ) ;
129130}
130131
131- function sleep ( ms : number , signal ?: AbortSignal ) : Promise < void > {
132- return new Promise ( ( resolve , reject ) => {
133- if ( signal ?. aborted ) {
134- reject ( new Error ( "Request was aborted" ) ) ;
135- return ;
136- }
137- const timeout = setTimeout ( resolve , ms ) ;
138- signal ?. addEventListener ( "abort" , ( ) => {
139- clearTimeout ( timeout ) ;
140- reject ( new Error ( "Request was aborted" ) ) ;
141- } ) ;
142- } ) ;
143- }
144-
145132export const streamOpenAICodexResponses : StreamFunction < "openai-codex-responses" , OpenAICodexResponsesOptions > = (
146133 model : Model < "openai-codex-responses" > ,
147134 context : Context ,
@@ -383,8 +370,14 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses"
383370 }
384371 if ( ! ( attempt < MAX_RETRIES && isRetryableError ( lastError ) ) ) throw lastError ;
385372 await attempts . settle ( "failed" ) ;
386- const delayMs = BASE_DELAY_MS * 2 ** attempt ;
387- await sleep ( delayMs , options ?. signal ) ;
373+ const retryAfter = response ?. headers . get ( "retry-after" ) ;
374+ const retryAfterMs = retryAfter
375+ ? / ^ \d + ( \. \d + ) ? $ / . test ( retryAfter )
376+ ? Number ( retryAfter ) * 1000
377+ : Date . parse ( retryAfter ) - Date . now ( )
378+ : NaN ;
379+ const delayMs = Number . isFinite ( retryAfterMs ) ? Math . max ( 0 , retryAfterMs ) : BASE_DELAY_MS * 2 ** attempt ;
380+ await abortableDelay ( delayMs , options ?. signal ) ;
388381 }
389382
390383 if ( ! response ?. ok ) {
@@ -453,6 +446,7 @@ function buildRequestBody(
453446 includeSystemPrompt : false ,
454447 onProjection,
455448 pendingPublicMessageGroups : options ?. attempts ?. pendingPublicMessageGroups ,
449+ responsesMessageIds : options ?. attempts ?. responsesMessageIds ,
456450 } ) ;
457451
458452 const body : RequestBody = {
@@ -463,7 +457,7 @@ function buildRequestBody(
463457 input : messages ,
464458 text : { verbosity : options ?. textVerbosity || "low" } ,
465459 include : [ "reasoning.encrypted_content" ] ,
466- prompt_cache_key : options ?. sessionId ,
460+ prompt_cache_key : options ?. cacheRetention === "none" ? undefined : options ?. sessionId ,
467461 tool_choice : "auto" ,
468462 parallel_tool_calls : true ,
469463 } ;
@@ -702,6 +696,13 @@ async function* parseSSE(response: Response): AsyncGenerator<Record<string, unkn
702696
703697const OPENAI_BETA_RESPONSES_WEBSOCKETS = "responses_websockets=2026-02-06" ;
704698const SESSION_WEBSOCKET_CACHE_TTL_MS = 5 * 60 * 1000 ;
699+ const WEBSOCKET_AUTH_HEADERS = [
700+ "authorization" ,
701+ "chatgpt-account-id" ,
702+ "openai-organization" ,
703+ "openai-project" ,
704+ "x-api-key" ,
705+ ] as const ;
705706
706707type WebSocketEventType = "open" | "message" | "error" | "close" ;
707708type WebSocketListener = ( event : unknown ) => void ;
@@ -723,6 +724,7 @@ interface CachedWebSocketContinuationState {
723724interface CachedWebSocketConnection {
724725 socket : WebSocketLike ;
725726 url : string ;
727+ headers : Headers ;
726728 busy : boolean ;
727729 idleTimer ?: ReturnType < typeof setTimeout > ;
728730 continuation ?: CachedWebSocketContinuationState ;
@@ -886,7 +888,9 @@ function scheduleSessionWebSocketExpiry(sessionId: string, entry: CachedWebSocke
886888 entry . idleTimer = setTimeout ( ( ) => {
887889 if ( entry . busy ) return ;
888890 closeWebSocketSilently ( entry . socket , 1000 , "idle_timeout" ) ;
889- websocketSessionCache . delete ( sessionId ) ;
891+ if ( websocketSessionCache . get ( sessionId ) === entry ) {
892+ websocketSessionCache . delete ( sessionId ) ;
893+ }
890894 } , SESSION_WEBSOCKET_CACHE_TTL_MS ) ;
891895}
892896
@@ -1005,16 +1009,20 @@ async function acquireWebSocket(
10051009 clearTimeout ( cached . idleTimer ) ;
10061010 cached . idleTimer = undefined ;
10071011 }
1008- if ( ! cached . busy && isWebSocketReusable ( cached . socket ) ) {
1012+ const matchesConnection =
1013+ cached . url === url && WEBSOCKET_AUTH_HEADERS . every ( ( name ) => cached . headers . get ( name ) === headers . get ( name ) ) ;
1014+ if ( ! cached . busy && matchesConnection && isWebSocketReusable ( cached . socket ) ) {
10091015 cached . busy = true ;
10101016 return {
10111017 socket : cached . socket ,
10121018 entry : cached ,
10131019 reused : true ,
10141020 release : ( { keep } = { } ) => {
1015- if ( ! keep || ! isWebSocketReusable ( cached . socket ) ) {
1021+ if ( ! keep || ! isWebSocketReusable ( cached . socket ) || websocketSessionCache . get ( sessionId ) !== cached ) {
10161022 closeWebSocketSilently ( cached . socket ) ;
1017- websocketSessionCache . delete ( sessionId ) ;
1023+ if ( websocketSessionCache . get ( sessionId ) === cached ) {
1024+ websocketSessionCache . delete ( sessionId ) ;
1025+ }
10181026 return ;
10191027 }
10201028 cached . busy = false ;
@@ -1032,21 +1040,22 @@ async function acquireWebSocket(
10321040 } ,
10331041 } ;
10341042 }
1035- if ( ! isWebSocketReusable ( cached . socket ) ) {
1036- closeWebSocketSilently ( cached . socket ) ;
1043+ cached . continuation = undefined ;
1044+ closeWebSocketSilently ( cached . socket ) ;
1045+ if ( websocketSessionCache . get ( sessionId ) === cached ) {
10371046 websocketSessionCache . delete ( sessionId ) ;
10381047 }
10391048 }
10401049
10411050 const socket = await connectWebSocket ( url , headers , signal , timeoutMs ) ;
1042- const entry : CachedWebSocketConnection = { socket, url, busy : true } ;
1051+ const entry : CachedWebSocketConnection = { socket, url, headers : new Headers ( headers ) , busy : true } ;
10431052 websocketSessionCache . set ( sessionId , entry ) ;
10441053 return {
10451054 socket,
10461055 entry,
10471056 reused : false ,
10481057 release : ( { keep } = { } ) => {
1049- if ( ! keep || ! isWebSocketReusable ( entry . socket ) ) {
1058+ if ( ! keep || ! isWebSocketReusable ( entry . socket ) || websocketSessionCache . get ( sessionId ) !== entry ) {
10501059 closeWebSocketSilently ( entry . socket ) ;
10511060 if ( entry . idleTimer ) clearTimeout ( entry . idleTimer ) ;
10521061 if ( websocketSessionCache . get ( sessionId ) === entry ) {
0 commit comments