|
897 | 897 | (is (false? @on-error-called*)) |
898 | 898 | (is (= "hello" @received-text*))))) |
899 | 899 |
|
| 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 | + |
900 | 1023 | (deftest async-no-retry-on-context-overflow-test |
901 | 1024 | (testing "does not retry on context overflow" |
902 | 1025 | (let [attempt* (atom 0) |
|
0 commit comments