Skip to content

Commit 1a68eb9

Browse files
authored
Merge pull request #532 from itkonen/retry-openai-responses-failures
Retry transient OpenAI Responses failures
2 parents cb47895 + ff006bd commit 1a68eb9

7 files changed

Lines changed: 334 additions & 21 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
## Unreleased
44

5+
- Retry transient OpenAI Responses `response.failed` server errors before output, preserving structured error and request IDs.
56
- Fix prompt cache invalidation warning after clearing the chat and changing model. #530
67
- Fix missing line break after "Prompt stopped" message when followed by another system message.
78

src/eca/llm_api.clj

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -465,8 +465,7 @@
465465
(catch Exception e
466466
(logger/warn logger-tag "on-history-sanitized callback failed" {:exception (ex-message e)})))))
467467
emit-first-message-fn (fn [& args]
468-
(when-not @first-response-received*
469-
(reset! first-response-received* true)
468+
(when (compare-and-set! first-response-received* false true)
470469
(apply on-first-response-received args)))
471470
on-message-received-wrapper (fn [& args]
472471
(apply emit-first-message-fn args)
@@ -516,6 +515,7 @@
516515
(if (and (contains? #{:rate-limited :overloaded :retryable-custom :premature-stop} error-type)
517516
(< attempt max-retries)
518517
(not wait-too-long?)
518+
(not @first-response-received*)
519519
(not (cancelled?)))
520520
(let [delay-ms (if rl-wait
521521
(+ (long (:delay-ms rl-wait)) rate-limit-wait-buffer-ms)

src/eca/llm_providers/errors.clj

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,9 @@
3636
#"(?i)connection refused"
3737
#"(?i)UnresolvedAddressException"])
3838

39+
(def ^:private openai-transient-message-pattern
40+
#"(?i)an error occurred while processing your request\.\s+you can retry your request\b")
41+
3942
(defn ^:private matches-any-pattern? [^String text patterns]
4043
(when text
4144
(some #(re-find % text) patterns)))
@@ -65,6 +68,25 @@
6568

6669
:else nil))
6770

71+
(defn ^:private classify-openai-responses-error
72+
[{:keys [code type message] source :error/source}]
73+
(when (= :openai-responses source)
74+
(cond
75+
(some #{"server_error"} [code type])
76+
{:error/type :overloaded}
77+
78+
(some #{"rate_limit_exceeded"} [code type])
79+
{:error/type :rate-limited}
80+
81+
(or (some? code) (some? type))
82+
{:error/type :unknown}
83+
84+
(and (string? message)
85+
(re-find openai-transient-message-pattern message))
86+
{:error/type :overloaded}
87+
88+
:else nil)))
89+
6890
(defn ^:private classify-by-message
6991
"Fallback classification from unstructured error message strings
7092
(e.g. SSE stream errors where HTTP status is not available)."
@@ -112,7 +134,8 @@
112134
(defn classify-error
113135
"Classifies an error map into a semantic error type.
114136
115-
Accepts the standard on-error map shape: {:message :status :body :exception}.
137+
Accepts the standard on-error map shape: {:message :status :body :exception},
138+
plus optional structured provider fields such as :code, :type, and :error/source.
116139
Optional `retry-rules` seq of user-configured rules checked before built-in classification.
117140
Returns a map with :error/type — one of:
118141
:retryable-custom — matched a user-configured retry rule (with optional :error/label)
@@ -126,6 +149,7 @@
126149
(or (when-let [pre-type (:error/type error-data)]
127150
{:error/type pre-type})
128151
(classify-by-custom-rules error-data retry-rules)
152+
(classify-openai-responses-error error-data)
129153
(when status
130154
(classify-by-status-and-body error-data))
131155
(classify-by-message error-data)

src/eca/llm_providers/openai.clj

Lines changed: 50 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -142,8 +142,38 @@
142142
""
143143
(:content (last (:output body))))})
144144

145+
(defn ^:private non-blank-str [value]
146+
(let [value (if (coll? value) (first value) value)
147+
value (some-> value str string/trim)]
148+
(when-not (string/blank? value)
149+
value)))
150+
151+
(defn ^:private response-header [headers header-name]
152+
(some (fn [[k value]]
153+
(let [k (if (keyword? k) (name k) (str k))]
154+
(when (= header-name (string/lower-case k))
155+
(non-blank-str value))))
156+
headers))
157+
158+
(defn ^:private response-failed->error-data [data response-headers]
159+
(let [response (:response data)
160+
error (:error response)
161+
request-id (some non-blank-str
162+
[(:request-id error)
163+
(:request_id error)
164+
(:request-id response)
165+
(:request_id response)
166+
(response-header response-headers "x-request-id")])]
167+
(assoc-some (or error {})
168+
:message (or (:message error) "OpenAI response failed")
169+
:error/source :openai-responses
170+
:response-id (:id response)
171+
:request-id request-id
172+
:headers response-headers)))
173+
145174
(defn ^:private base-responses-request! [{:keys [rid body api-url auth-type url-relative-path api-key account-id on-error on-stream http-client extra-headers]}]
146175
(let [oauth? (= :auth/oauth auth-type)
176+
stream? (and on-stream (not= false (:stream body)))
147177
url (if oauth?
148178
codex-url
149179
(join-api-url api-url (or url-relative-path responses-path)))
@@ -174,22 +204,32 @@
174204
:body (json/generate-string body)
175205
:throw-exceptions? false
176206
:http-client (client/merge-with-global-http-client http-client)
177-
:as (if on-stream :stream :json)})]
207+
:as (if stream? :stream :json)})]
178208
(if (not= 200 status)
179-
(let [body-str (if on-stream (slurp body) body)]
209+
(let [body-str (if stream? (slurp body) body)]
180210
(logger/warn logger-tag "Unexpected response status: %s body: %s" status body-str)
181211
(on-error {:message (format "OpenAI response status: %s body: %s" status body-str)
182212
:status status
183213
:body body-str
184214
:headers resp-headers}))
185-
(if on-stream
186-
(with-open [rdr (io/reader body)]
187-
(doseq [[event data] (llm-util/event-data-seq rdr)]
188-
(llm-util/log-response logger-tag rid event data)
189-
(on-stream event data)))
215+
(if stream?
216+
(let [stream-error
217+
(with-open [rdr (io/reader body)]
218+
(loop [events (seq (llm-util/event-data-seq rdr))]
219+
(when-let [[event data] (first events)]
220+
(llm-util/log-response logger-tag rid event data)
221+
(if (= "response.failed" event)
222+
(response-failed->error-data data resp-headers)
223+
(do
224+
(on-stream event data resp-headers)
225+
(recur (next events)))))))]
226+
(when stream-error
227+
(on-error stream-error)))
190228
(do
191229
(llm-util/log-response logger-tag rid "response" body)
192-
(response-body->result body)))))
230+
(if (= "failed" (:status body))
231+
(on-error (response-failed->error-data {:response body} resp-headers))
232+
(response-body->result body))))))
193233
(catch Exception e
194234
(on-error {:exception e
195235
:message (if (ex-data e)
@@ -327,7 +367,7 @@
327367
sync-result* (when-not callbacks (atom nil))
328368
on-stream-fn
329369
(if callbacks
330-
(fn handle-stream [event data]
370+
(fn handle-stream [event data & _]
331371
(case event
332372
;; text
333373
"response.output_text.delta"
@@ -465,23 +505,15 @@
465505
:on-stream handle-stream})))
466506
(on-message-received {:type :finish
467507
:finish-reason (-> data :response :status)})))
468-
469-
"response.failed" (do
470-
(when-let [error (-> data :response :error)]
471-
(on-error {:message (:message error)}))
472-
(on-message-received {:type :finish
473-
:finish-reason (-> data :response :status)}))
474508
nil))
475509
;; Sync mode: collect text deltas into result atom
476510
(let [sb (StringBuilder.)]
477-
(fn handle-sync-stream [event data]
511+
(fn handle-sync-stream [event data & _]
478512
(case event
479513
"response.output_text.delta"
480514
(.append sb ^String (:delta data))
481515
"response.completed"
482516
(reset! sync-result* {:output-text (.toString sb)})
483-
"response.failed"
484-
(reset! sync-result* {:error {:message (-> data :response :error :message)}})
485517
nil))))
486518
result (base-responses-request!
487519
{:rid (llm-util/gen-rid)

test/eca/llm_api_test.clj

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -897,6 +897,129 @@
897897
(is (false? @on-error-called*))
898898
(is (= "hello" @received-text*)))))
899899

900+
(deftest first-response-callback-is-atomic-test
901+
(testing "concurrent output callbacks emit first response once"
902+
(let [first-response-count* (atom 0)]
903+
(with-redefs [eca.llm-api/prompt! (fn [{:keys [on-message-received]}]
904+
(let [ready (java.util.concurrent.CountDownLatch. 2)
905+
go (java.util.concurrent.CountDownLatch. 1)
906+
workers (mapv (fn [text]
907+
(future
908+
(.countDown ready)
909+
(.await go)
910+
(on-message-received {:type :text :text text})))
911+
["one" "two"])]
912+
(.await ready)
913+
(.countDown go)
914+
(doseq [worker workers]
915+
@worker)))]
916+
(llm-api/sync-or-async-prompt!
917+
(make-prompt-opts
918+
{:on-first-response-received (fn [& _] (swap! first-response-count* inc))
919+
:on-message-received identity
920+
:on-error identity})))
921+
(is (= 1 @first-response-count*)))))
922+
923+
(deftest async-retry-on-structured-server-error-test
924+
(testing "retries an OpenAI Responses server_error before output"
925+
(let [attempt* (atom 0)
926+
retry-events* (atom [])
927+
received-text* (atom "")
928+
on-error-called* (atom false)]
929+
(with-redefs [eca.llm-api/prompt! (fn [{:keys [on-message-received on-error]}]
930+
(let [attempt (swap! attempt* inc)]
931+
(if (= 1 attempt)
932+
(on-error {:code "server_error"
933+
:message "Request failed"
934+
:error/source :openai-responses
935+
:response-id "resp_123"
936+
:request-id "req_123"})
937+
(do
938+
(on-message-received {:type :text :text "hello"})
939+
(on-message-received {:type :finish :finish-reason "stop"})))))
940+
eca.llm-api/sleep-with-cancel (fn [_ cancelled?] (not (cancelled?)))]
941+
(llm-api/sync-or-async-prompt!
942+
(make-prompt-opts
943+
{:on-retry (fn [event] (swap! retry-events* conj event))
944+
:on-error (fn [_] (reset! on-error-called* true))
945+
:on-message-received (fn [{:keys [type text]}]
946+
(when (= :text type)
947+
(swap! received-text* str text)))})))
948+
(is (= 2 @attempt*))
949+
(is (= :overloaded (get-in (first @retry-events*) [:classified :error/type])))
950+
(is (= "resp_123" (get-in (first @retry-events*) [:error-data :response-id])))
951+
(is (= "req_123" (get-in (first @retry-events*) [:error-data :request-id])))
952+
(is (false? @on-error-called*))
953+
(is (= "hello" @received-text*)))))
954+
955+
(deftest sync-retry-on-structured-server-error-test
956+
(testing "retries an OpenAI Responses server_error in sync mode"
957+
(let [attempt* (atom 0)
958+
retry-events* (atom [])
959+
on-error-called* (atom false)]
960+
(with-redefs [eca.llm-api/prompt! (fn [_opts]
961+
(if (= 1 (swap! attempt* inc))
962+
{:error {:code "server_error"
963+
:message "Request failed"
964+
:error/source :openai-responses
965+
:response-id "resp_sync"
966+
:request-id "req_sync"}}
967+
{:output-text "success"
968+
:usage {:input-tokens 1 :output-tokens 1}}))
969+
eca.llm-api/sleep-with-cancel (fn [_ cancelled?] (not (cancelled?)))]
970+
(llm-api/sync-or-async-prompt!
971+
(make-prompt-opts
972+
{:stream false
973+
:on-retry (fn [event] (swap! retry-events* conj event))
974+
:on-error (fn [_] (reset! on-error-called* true))
975+
:on-message-received identity})))
976+
(is (= 2 @attempt*))
977+
(is (= :overloaded (get-in (first @retry-events*) [:classified :error/type])))
978+
(is (= "resp_sync" (get-in (first @retry-events*) [:error-data :response-id])))
979+
(is (= "req_sync" (get-in (first @retry-events*) [:error-data :request-id])))
980+
(is (false? @on-error-called*)))))
981+
982+
(deftest async-no-retry-after-visible-output-test
983+
(doseq [[label emit-output]
984+
[["text output"
985+
(fn [{:keys [on-message-received]}]
986+
(on-message-received {:type :text :text "partial"}))]
987+
["reasoning output"
988+
(fn [{:keys [on-reason]}]
989+
(on-reason {:status :thinking :id "reason_1" :text "partial"}))]
990+
["tool-call output"
991+
(fn [{:keys [on-prepare-tool-call]}]
992+
(on-prepare-tool-call {:id "call_1" :full-name "tool" :arguments-text "{"}))]
993+
["web-search output"
994+
(fn [{:keys [on-server-web-search]}]
995+
(on-server-web-search {:status :started :id "search_1"}))]
996+
["image-generation output"
997+
(fn [{:keys [on-server-image-generation]}]
998+
(on-server-image-generation {:status :started :id "image_1"}))]]]
999+
(testing (str "does not retry after " label)
1000+
(let [attempt* (atom 0)
1001+
sleep-calls* (atom 0)
1002+
errors* (atom [])]
1003+
(with-redefs [eca.llm-api/prompt! (fn [{:keys [on-error] :as callbacks}]
1004+
(swap! attempt* inc)
1005+
(emit-output callbacks)
1006+
(on-error {:code "server_error"
1007+
:message "Request failed"
1008+
:error/source :openai-responses}))
1009+
eca.llm-api/sleep-with-cancel (fn [_ _]
1010+
(swap! sleep-calls* inc)
1011+
true)]
1012+
(llm-api/sync-or-async-prompt!
1013+
(make-prompt-opts
1014+
{:on-error (fn [error] (swap! errors* conj error))
1015+
:on-message-received identity})))
1016+
(is (= 1 @attempt*))
1017+
(is (zero? @sleep-calls*))
1018+
(is (= [{:code "server_error"
1019+
:message "Request failed"
1020+
:error/source :openai-responses}]
1021+
@errors*))))))
1022+
9001023
(deftest async-no-retry-on-context-overflow-test
9011024
(testing "does not retry on context overflow"
9021025
(let [attempt* (atom 0)

test/eca/llm_providers/errors_test.clj

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,62 @@
9494
:body "Bad Gateway"
9595
:message "LLM response status: 502 body: Bad Gateway"})))))
9696

97+
(deftest classify-openai-response-failed-test
98+
(let [generic-message (str "An error occurred while processing your request. "
99+
"You can retry your request, or contact support if the error persists.")
100+
responses-error {:error/source :openai-responses}]
101+
(testing "structured server_error code is overloaded"
102+
(is (= {:error/type :overloaded}
103+
(llm-providers.errors/classify-error
104+
(assoc responses-error
105+
:code "server_error"
106+
:message "Request failed")))))
107+
108+
(testing "structured server_error type is overloaded"
109+
(is (= {:error/type :overloaded}
110+
(llm-providers.errors/classify-error
111+
(assoc responses-error
112+
:type "server_error"
113+
:message "Request failed")))))
114+
115+
(testing "structured rate_limit_exceeded code is rate limited"
116+
(is (= {:error/type :rate-limited}
117+
(llm-providers.errors/classify-error
118+
(assoc responses-error
119+
:code "rate_limit_exceeded"
120+
:message "Request failed")))))
121+
122+
(testing "generic transient message is a fallback without structured fields"
123+
(is (= {:error/type :overloaded}
124+
(llm-providers.errors/classify-error
125+
(assoc responses-error :message generic-message)))))
126+
127+
(testing "generic transient message is not applied to unmarked errors"
128+
(is (= {:error/type :unknown}
129+
(llm-providers.errors/classify-error
130+
{:message generic-message}))))
131+
132+
(testing "unknown structured fields are authoritative over retry-looking messages"
133+
(doseq [error-data [{:code "invalid_prompt" :message generic-message}
134+
{:code "invalid_prompt" :message "Internal server error"}
135+
{:type "invalid_request_error" :message "Rate limit exceeded"}
136+
{:code "invalid_prompt"
137+
:message "Request failed"
138+
:exception (Exception. "Connection error")}]]
139+
(is (= {:error/type :unknown}
140+
(llm-providers.errors/classify-error
141+
(merge responses-error error-data))))))
142+
143+
(testing "custom retry rules still take priority over structured fields"
144+
(is (= {:error/type :retryable-custom
145+
:error/label "Temporary proxy failure"}
146+
(llm-providers.errors/classify-error
147+
(assoc responses-error
148+
:code "invalid_prompt"
149+
:message generic-message)
150+
[{:errorPattern "processing your request"
151+
:label "Temporary proxy failure"}]))))))
152+
97153
(deftest classify-error-unknown-test
98154
(testing "500 is classified as overloaded"
99155
(is (= {:error/type :overloaded}

0 commit comments

Comments
 (0)