diff --git a/Autotests/run_mandatory b/Autotests/run_mandatory index 637140931..986796687 100644 --- a/Autotests/run_mandatory +++ b/Autotests/run_mandatory @@ -42,4 +42,6 @@ unit/test_fileio_verified_deletes.py unit/test_helper_parsing.py unit/test_llm_budget.py unit/test_openclaw_unit.py +unit/test_nginx_llm_timeouts.py +unit/test_llm_timeout_message.py import_knowledge/test_import_knowledge.py diff --git a/Autotests/unit/test_llm_timeout_message.py b/Autotests/unit/test_llm_timeout_message.py new file mode 100644 index 000000000..52aab0388 --- /dev/null +++ b/Autotests/unit/test_llm_timeout_message.py @@ -0,0 +1,455 @@ +"""Unit tests for the provider timeout status message. + +When a provider request runs out of time and the client's retries are gone, the +provider used to return an empty string. The loop then had nothing to run, so +the turn ended without a word to the user and the task looked abandoned (#321). +The timeout now comes back as a `send` command carrying a status message, the +same way a reply cut off by the token limit already does. + +No container, no network, no API key, and no provider SDK: `openai` and the +configuration module are stubbed before the module under test is loaded, the +same pattern as test_openclaw_unit.py. +""" +import importlib.util +import os +import re +import sys +import types + +import pytest + +_REPO_ROOT = os.path.normpath(os.path.join(os.path.dirname(__file__), "..", "..")) +_LIB_LLM_EXT_PATH = os.path.join(_REPO_ROOT, "providers", "lib_llm_ext.py") +_HELPER_PATH = os.path.join(_REPO_ROOT, "src", "helper.py") + +# lib_llm_ext.py does `from src.helper import quote_arg` (repo-root package). +if _REPO_ROOT not in sys.path: + sys.path.insert(0, _REPO_ROOT) + + +class _StubAPITimeoutError(Exception): + """Stands in for openai.APITimeoutError: the client's own timeout.""" + + +class _StubAPIConnectionError(Exception): + """Stands in for openai.APIConnectionError: the request never got an answer.""" + + +def _install_stubs(): + config_stub = types.ModuleType("config") + config_stub.config_get_by_key = lambda key, default=None: None + sys.modules["config"] = config_stub + + openai_stub = types.ModuleType("openai") + openai_stub.APITimeoutError = _StubAPITimeoutError + openai_stub.APIConnectionError = _StubAPIConnectionError + openai_stub.OpenAI = object # only referenced in a type annotation + sys.modules["openai"] = openai_stub + + +def _load(name, path): + spec = importlib.util.spec_from_file_location(name, path) + module = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = module + spec.loader.exec_module(module) + return module + + +@pytest.fixture(scope="module") +def llm(): + saved = {name: sys.modules.get(name) for name in ("config", "openai")} + _install_stubs() + try: + yield _load("lib_llm_ext_under_test", _LIB_LLM_EXT_PATH) + finally: + for name, module in saved.items(): + if module is None: + sys.modules.pop(name, None) + else: + sys.modules[name] = module + + +@pytest.fixture(scope="module") +def helper(): + return _load("helper_under_test", _HELPER_PATH) + + +def _gateway_error(status): + error = Exception(f"{status} Gateway Time-out") + error.status_code = status + return error + + +class _RaisingClient: + """Stand-in for the OpenAI client whose chat call always fails.""" + + def __init__(self, error): + def create(**kwargs): + raise error + + completions = types.SimpleNamespace(create=create) + self.chat = types.SimpleNamespace(completions=completions) + + +def _provider(llm, error): + provider = llm.AIProvider("OpenAIAPI", "OPENAIAPI_API_KEY", "test-model", "http://localhost/v1/") + provider._client = _RaisingClient(error) + return provider + + +# --- which failures count as a timeout --------------------------------------- + +def test_client_timeout_is_a_timeout(llm): + assert llm._is_timeout_error(_StubAPITimeoutError("timed out")) + + +@pytest.mark.parametrize("status", [408, 504, 524]) +def test_gateway_timeout_statuses_are_a_timeout(llm, status): + assert llm._is_timeout_error(_gateway_error(status)) + + +@pytest.mark.parametrize("status", [400, 429, 500, 502]) +def test_other_statuses_are_not_a_timeout(llm, status): + assert not llm._is_timeout_error(_gateway_error(status)) + + +def test_a_plain_error_is_not_a_timeout(llm): + assert not llm._is_timeout_error(ValueError("boom")) + + +# --- what chat() returns ------------------------------------------------------ + +def _is_timeout_notice(result): + return result.startswith('(send "LLM request timed out at ') and result.endswith('")') + + +def test_chat_tells_the_user_when_the_client_times_out(llm): + assert _is_timeout_notice(_provider(llm, _StubAPITimeoutError("timed out")).chat("prompt")) + + +def test_chat_tells_the_user_when_the_gateway_times_out(llm): + assert _is_timeout_notice(_provider(llm, _gateway_error(504)).chat("prompt")) + + +def test_chat_still_returns_nothing_for_other_failures(llm): + assert _provider(llm, ValueError("boom")).chat("prompt") == "" + + +# --- the message has to survive the parser the loop runs it through ---------- + +def test_the_timeout_command_parses_into_one_send(llm, helper): + parsed = helper.balance_parentheses(llm._llm_timeout_command()) + assert parsed.startswith('((send "') + assert parsed.endswith('"))') + assert parsed.count("(send ") == 1 + + +# --- one attempt, so the timeout is reported when the first request gives up -- + +def _client_kwargs(llm, monkeypatch, gateway): + captured = {} + + def recorder(**kwargs): + captured.update(kwargs) + return object() + + monkeypatch.setattr(llm.openai, "OpenAI", recorder) + monkeypatch.setattr( + llm, "config_get_by_key", + lambda key, default=None: gateway if key == "GATEWAY_URL" else default, + ) + if gateway is None: + monkeypatch.setenv("OPENAIAPI_API_KEY", "dummy") + provider = llm.AIProvider("OpenAIAPI", "OPENAIAPI_API_KEY", "test-model", "http://localhost/v1/") + assert provider._create_client() is not None + return captured + + +def test_the_proxy_client_makes_one_attempt(llm, monkeypatch): + assert _client_kwargs(llm, monkeypatch, "http://localhost:8080")["max_retries"] == 0 + + +def test_the_direct_client_makes_one_attempt(llm, monkeypatch): + assert _client_kwargs(llm, monkeypatch, None)["max_retries"] == 0 + + +# --- which failures are worth another attempt -------------------------------- + +@pytest.mark.parametrize("status", [409, 429, 500, 502, 503, 522, 529]) +def test_transient_statuses_are_retried(llm, status): + assert llm._is_transient_error(_gateway_error(status)) + + +@pytest.mark.parametrize("status", [408, 504, 524]) +def test_timeout_statuses_are_not_retried(llm, status): + assert not llm._is_transient_error(_gateway_error(status)) + + +def test_a_client_timeout_is_not_retried(llm): + assert not llm._is_transient_error(_StubAPITimeoutError("timed out")) + + +def test_a_connection_error_is_retried(llm): + assert llm._is_transient_error(_StubAPIConnectionError("no answer")) + + +def test_a_plain_error_is_not_retried(llm): + assert not llm._is_transient_error(ValueError("boom")) + + +# --- how _retrying behaves ---------------------------------------------------- + +def _counting_call(errors): + """Raise each error in turn, then return a sentinel. Records the attempts.""" + calls = [] + + def call(): + calls.append(len(calls) + 1) + if calls[-1] <= len(errors): + raise errors[calls[-1] - 1] + return "answer" + + return call, calls + + +def test_a_transient_failure_is_retried_until_it_succeeds(llm, monkeypatch): + monkeypatch.setattr(llm.time, "sleep", lambda seconds: None) + call, calls = _counting_call([_gateway_error(503)]) + assert llm._retrying(call, "OpenAIAPI") == "answer" + assert len(calls) == 2 + + +def test_a_timeout_is_attempted_once(llm, monkeypatch): + monkeypatch.setattr(llm.time, "sleep", lambda seconds: None) + call, calls = _counting_call([_StubAPITimeoutError("timed out")] * 3) + with pytest.raises(_StubAPITimeoutError): + llm._retrying(call, "OpenAIAPI") + assert len(calls) == 1 + + +def test_transient_failures_stop_at_the_attempt_limit(llm, monkeypatch): + monkeypatch.setattr(llm.time, "sleep", lambda seconds: None) + call, calls = _counting_call([_gateway_error(503)] * 5) + with pytest.raises(Exception): + llm._retrying(call, "OpenAIAPI") + assert len(calls) == llm.CHAT_ATTEMPTS + + +def test_a_slow_transient_failure_is_not_retried(llm, monkeypatch): + """A failure that already took longer than the budget is not transient in + any useful sense, so the caller hears about it instead of waiting again.""" + monkeypatch.setattr(llm.time, "sleep", lambda seconds: None) + clock = iter([0, llm.CHAT_RETRY_BUDGET_SECONDS + 1]) + monkeypatch.setattr(llm.time, "monotonic", lambda: next(clock)) + call, calls = _counting_call([_gateway_error(503)] * 3) + with pytest.raises(Exception): + llm._retrying(call, "OpenAIAPI") + assert len(calls) == 1 + + +# --- the wait between attempts follows Retry-After ---------------------------- + +def _retry_after_error(status, value): + error = _gateway_error(status) + error.response = types.SimpleNamespace(headers={"retry-after": value}) + return error + + +def test_retry_after_in_seconds_is_honoured(llm): + assert llm._retry_delay(_retry_after_error(429, "5"), 1) == 5.0 + + +def test_retry_after_as_a_date_falls_back_to_the_backoff(llm): + delay = llm._retry_delay(_retry_after_error(429, "Wed, 21 Oct 2026 07:28:00 GMT"), 1) + assert delay == llm.CHAT_RETRY_BACKOFF_SECONDS + + +def test_without_retry_after_the_backoff_grows(llm): + assert llm._retry_delay(_gateway_error(503), 1) == llm.CHAT_RETRY_BACKOFF_SECONDS + assert llm._retry_delay(_gateway_error(503), 2) == llm.CHAT_RETRY_BACKOFF_SECONDS * 2 + + +def test_the_retry_waits_as_long_as_the_provider_asked(llm, monkeypatch): + slept = [] + monkeypatch.setattr(llm.time, "sleep", slept.append) + call, calls = _counting_call([_retry_after_error(429, "5")]) + assert llm._retrying(call, "OpenAIAPI") == "answer" + assert slept == [5.0] + assert len(calls) == 2 + + +def test_a_retry_after_beyond_the_budget_is_not_waited_out(llm, monkeypatch): + monkeypatch.setattr(llm.time, "sleep", lambda seconds: None) + call, calls = _counting_call([_retry_after_error(429, str(llm.CHAT_RETRY_BUDGET_SECONDS + 10))] * 3) + with pytest.raises(Exception): + llm._retrying(call, "OpenAIAPI") + assert len(calls) == 1 + + +# --- two timeouts in a row must both reach the user -------------------------- + +def test_the_notice_carries_the_time_to_the_millisecond(llm): + assert re.search(r"timed out at \d{2}:\d{2}:\d{2}\.\d{3}", llm._llm_timeout_command()) + + +def test_two_notices_differ_even_inside_the_same_millisecond(llm, monkeypatch): + """`send` drops a message equal to the last one it sent, so two notices in a + row have to differ or the second turn goes unanswered. The clock alone cannot + guarantee that, so the count has to carry it.""" + frozen = types.SimpleNamespace(strftime=lambda fmt: "05:14:17.000000") + monkeypatch.setattr(llm, "datetime", types.SimpleNamespace(now=lambda: frozen)) + first = llm._llm_timeout_command() + second = llm._llm_timeout_command() + assert "05:14:17.000" in first and "05:14:17.000" in second + assert first != second + + +# --- one notice per run of timeouts ------------------------------------------ + +SYSTEM_PART = "PROMPT: you are an agent" + + +def _prompt(human=""): + return f"{SYSTEM_PART} :-:-:-: {human}" + + +def _reply(text='(send "hello")'): + message = types.SimpleNamespace(content=text) + return types.SimpleNamespace(choices=[types.SimpleNamespace(message=message, finish_reason="stop")], + usage=None) + + +class _ScriptedClient: + """Chat client that walks a list of steps: an exception is raised, anything + else is returned. The last step repeats.""" + + def __init__(self, steps): + self.calls = 0 + + def create(**kwargs): + self.calls += 1 + step = steps[min(self.calls, len(steps)) - 1] + if isinstance(step, BaseException): + raise step + return step + + self.chat = types.SimpleNamespace(completions=types.SimpleNamespace(create=create)) + + +def _scripted_provider(llm, steps): + provider = llm.AIProvider("OpenAIAPI", "OPENAIAPI_API_KEY", "test-model", "http://localhost/v1/") + provider._client = _ScriptedClient(steps) + return provider + + +def test_only_the_first_timeout_of_a_run_is_reported(llm): + provider = _scripted_provider(llm, [_gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt())) + assert provider.chat(_prompt()) == "" + assert provider.chat(_prompt()) == "" + + +def test_a_successful_answer_starts_a_new_run(llm): + provider = _scripted_provider(llm, [_gateway_error(504), _reply(), _gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt())) + assert provider.chat(_prompt()) == '(send "hello")' + assert _is_timeout_notice(provider.chat(_prompt())) + + +def test_a_new_human_message_starts_a_new_run(llm): + """The user who just wrote deserves an answer, even if the previous cycle + already reported a timeout.""" + provider = _scripted_provider(llm, [_gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt("HUMAN-MSG: first"))) + assert provider.chat(_prompt()) == "" + assert _is_timeout_notice(provider.chat(_prompt("HUMAN-MSG: second"))) + assert provider.chat(_prompt()) == "" + + +def test_every_tagged_message_gets_its_own_notice(llm): + provider = _scripted_provider(llm, [_gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt("HUMAN-MSG: same"))) + assert provider.chat(_prompt()) == "" + assert _is_timeout_notice(provider.chat(_prompt("HUMAN-MSG: same"))) + + +def test_the_tagged_message_is_recognised_however_the_loop_wraps_it(llm): + provider = _scripted_provider(llm, [_gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt("(HUMAN-MSG: hello)"))) + assert provider.chat(_prompt()) == "" + assert _is_timeout_notice(provider.chat(_prompt("['HUMAN-MSG:', 'hello']"))) + + +def test_the_spamshield_reminder_is_not_a_new_turn(llm): + """With spamShield on, the loop puts this in the tail on every follow-up + cycle. Taking it for a message would restore one notice per cycle.""" + provider = _scripted_provider(llm, [_gateway_error(504)]) + assert _is_timeout_notice(provider.chat(_prompt("HUMAN-MSG: hello"))) + for _ in range(3): + assert provider.chat(_prompt(" DO NOT RE-SEND OR SPAM!")) == "" + + +# --- the client enforces the timeout the notice names ------------------------ + +def test_the_client_enforces_the_request_timeout(llm, monkeypatch): + assert _client_kwargs(llm, monkeypatch, "http://localhost:8080")["timeout"] == llm.CHAT_REQUEST_TIMEOUT_SECONDS + assert llm.CHAT_REQUEST_TIMEOUT_SECONDS == 600 + + +def test_the_notice_names_the_limit_and_both_places_that_hold_it(llm): + notice = llm._llm_timeout_command() + assert "600 s request timeout" in notice + assert "raising both" in notice + + +# --- every chat client carries the same limits ------------------------------- + +@pytest.fixture(scope="module") +def openrouter(llm): + """OpenRouter overrides _create_client, so it needs checking on its own.""" + providers_stub = types.ModuleType("providers") + providers_stub.LLMProvider = object + providers_stub.registerLLMProvider = lambda name, provider: None + saved = {name: sys.modules.get(name) for name in ("providers", "lib_llm_ext", "openai", "config")} + sys.modules["providers"] = providers_stub + sys.modules["lib_llm_ext"] = llm + sys.modules["openai"] = llm.openai + config_stub = types.ModuleType("config") + config_stub.config_get_by_key = lambda key, default=None: default + sys.modules["config"] = config_stub + try: + yield _load("openrouter_under_test", os.path.join(_REPO_ROOT, "providers", "openrouter.py")) + finally: + for name, module in saved.items(): + if module is None: + sys.modules.pop(name, None) + else: + sys.modules[name] = module + + +def test_the_openrouter_client_carries_the_same_limits(openrouter, llm, monkeypatch): + captured = {} + + def recorder(**kwargs): + captured.update(kwargs) + return object() + + monkeypatch.setattr(llm.openai, "OpenAI", recorder) + monkeypatch.setattr(openrouter, "config_get_by_key", + lambda key, default=None: "http://localhost:8080" if key == "GATEWAY_URL" else default) + provider = openrouter.OpenRouterProviderImpl("OpenRouter", "OPENROUTER_API_KEY", "z-ai/glm-5.2", + "https://openrouter.ai/api/v1") + assert provider._create_client() is not None + assert captured["max_retries"] == llm.CHAT_MAX_RETRIES + assert captured["timeout"] == llm.CHAT_REQUEST_TIMEOUT_SECONDS + + + +# --- the tag the reset depends on must stay the tag the loop writes ---------- + +def test_the_marker_matches_the_tag_the_loop_writes(llm): + """_start_of_turn keys off this tag. If the loop ever renames it, the reset + would quietly stop working, so the two are checked against each other.""" + with open(os.path.join(_REPO_ROOT, "src", "loop.metta"), encoding="utf-8") as f: + loop = f.read() + assert llm.HUMAN_MESSAGE_MARKER in loop diff --git a/Autotests/unit/test_nginx_llm_timeouts.py b/Autotests/unit/test_nginx_llm_timeouts.py new file mode 100644 index 000000000..efc9a0611 --- /dev/null +++ b/Autotests/unit/test_nginx_llm_timeouts.py @@ -0,0 +1,77 @@ +"""Unit tests for the LLM provider routes in proxy/nginx.conf.template. + +In the Docker image every LLM request goes through this proxy. When a route +gives up before the provider answers, nginx returns 504, the provider call +fails and the agent sends nothing back to the user (#321). The OpenAI client +used by the providers waits up to 600 seconds, so each LLM route has to wait +at least as long. + +No container, no network: the template is read as text. +""" +import os +import re + +import pytest + +_REPO_ROOT = os.path.normpath(os.path.join(os.path.dirname(__file__), "..", "..")) +_TEMPLATE = os.path.join(_REPO_ROOT, "proxy", "nginx.conf.template") + +LLM_ROUTES = ["anthropic", "asicloud", "openai", "asione", "openaiapi", "openrouter"] +_PROVIDER_SOURCE = os.path.join(_REPO_ROOT, "providers", "lib_llm_ext.py") + + +def _client_timeout_seconds(): + """The timeout the providers give their client, read from the source rather + than repeated here: the route has to wait at least as long, and the two must + not drift apart. The module itself is not imported, since the provider SDK is + not installed where this suite runs. + """ + with open(_PROVIDER_SOURCE, encoding="utf-8") as f: + match = re.search(r"^CHAT_REQUEST_TIMEOUT_SECONDS\s*=\s*(\d+)", f.read(), re.MULTILINE) + assert match, "providers/lib_llm_ext.py no longer defines CHAT_REQUEST_TIMEOUT_SECONDS" + return int(match.group(1)) + + +CLIENT_TIMEOUT_SECONDS = _client_timeout_seconds() + +_UNITS = {"s": 1, "m": 60, "h": 3600} + + +def _seconds(value): + match = re.fullmatch(r"(\d+)([smh]?)", value) + assert match, f"unsupported nginx time value: {value}" + return int(match.group(1)) * _UNITS.get(match.group(2) or "s") + + +def _directive(text, name): + match = re.search(rf"^\s*{name}\s+(\S+);", text, re.MULTILINE) + return match.group(1) if match else None + + +@pytest.fixture(scope="module") +def template(): + with open(_TEMPLATE, encoding="utf-8") as f: + return f.read() + + +def _route_block(template, route): + match = re.search(rf"location /{route}/ \{{\n(.*?)\n\s*\}}\n", template, re.DOTALL) + assert match, f"no location block for /{route}/" + return match.group(1) + + +def _effective(template, route, name): + http_defaults = template.split("server {", 1)[0] + value = _directive(_route_block(template, route), name) or _directive(http_defaults, name) + assert value, f"{name} is not set for /{route}/" + return _seconds(value) + + +@pytest.mark.parametrize("route", LLM_ROUTES) +def test_llm_route_waits_as_long_as_the_client_for_a_response(template, route): + assert _effective(template, route, "proxy_read_timeout") >= CLIENT_TIMEOUT_SECONDS + + +@pytest.mark.parametrize("route", LLM_ROUTES) +def test_llm_route_waits_as_long_as_the_client_to_send_the_request(template, route): + assert _effective(template, route, "proxy_send_timeout") >= CLIENT_TIMEOUT_SECONDS diff --git a/providers/asione.py b/providers/asione.py index 371ecaca3..4f91ae109 100644 --- a/providers/asione.py +++ b/providers/asione.py @@ -49,6 +49,7 @@ def __init__(self, name: str, var_name: str, model_name: str, base_url: str): def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", **kwargs) -> str: """Send chat request, initializing client if needed.""" + self._start_of_turn(content) self._ensure_client() if self._client is None: @@ -57,18 +58,22 @@ def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", sysmsg, usermsg = content.split(":-:-:-:") thinking_budget = _reasoning_budget(max_tokens, reasoning) try: - response = self._client.chat.completions.create( - model=self._model_name, - messages=[{"role": "system", "content": sysmsg}, - {"role": "user", "content": usermsg}], - max_tokens=max_tokens, - extra_body={ - "enable_thinking": thinking_budget > 0, - "thinking_budget": thinking_budget - }, - **kwargs + response = llm._retrying( + lambda: self._client.chat.completions.create( + model=self._model_name, + messages=[{"role": "system", "content": sysmsg}, + {"role": "user", "content": usermsg}], + max_tokens=max_tokens, + extra_body={ + "enable_thinking": thinking_budget > 0, + "thinking_budget": thinking_budget + }, + **kwargs + ), + self._name, ) + self._answered() raw = response.choices[0].message.content or "" finish_reason = getattr(response.choices[0], "finish_reason", None) llm._log_raw(self._name, self._model_name, raw) @@ -81,4 +86,6 @@ def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", return resp except Exception as e: logger.exception(f"[ASIOneProviderImpl.chat]: Exception while communicating with LLM: {e}") + if llm._is_timeout_error(e): + return self._timeout_reply() return "" diff --git a/providers/lib_llm_ext.py b/providers/lib_llm_ext.py index 28590cb7a..6e88c17b9 100644 --- a/providers/lib_llm_ext.py +++ b/providers/lib_llm_ext.py @@ -1,4 +1,5 @@ -import os, hashlib +import os, hashlib, time +from datetime import datetime import openai from typing import Optional, Tuple, Dict, Any from config import config_get_by_key @@ -6,6 +7,10 @@ from src.logger import get_logger PROMPT_DELIMITER = ":-:-:-:" +# The loop tags the human message it appends after the delimiter. Everything else +# that can land there, the spamShield reminder or nothing at all, is not a turn of +# its own. +HUMAN_MESSAGE_MARKER = "HUMAN-MSG:" LLM_EMPTY_RESPONSE_MESSAGE = ( "The agent didn\'t return an answer: reasoning exceeded the token limit for " "this response before it could produce one." @@ -21,6 +26,43 @@ "reasoning levels need a higher token limit." ) +LLM_TIMEOUT_MESSAGE = ( + "LLM request timed out at {time}. Please try again later." + "\n\n" + "If you are the Omega administrator: the provider did not answer within the " + "{timeout} s request timeout, and a timed-out request is not retried. This is " + "timeout notice {count} since the agent started. The failed request is in the " + "agent log; check the provider status. The limit is enforced by the client and " + "by the provider's route in the proxy, so allowing longer answers means " + "raising both." +) + +# Statuses a gateway returns when the upstream did not answer in time. +GATEWAY_TIMEOUT_STATUSES = (408, 504, 524) + +# The request timeout the client enforces. It matches the proxy route timeout, so +# a slow answer is cut once, by whichever limit is reached first, and the notice +# can name a single number. +CHAT_REQUEST_TIMEOUT_SECONDS = 600 + +# One attempt per chat request inside the SDK. Its retry loop repeats a timed-out +# request unconditionally, which would multiply the request timeout before the +# user hears anything. Transient failures are retried by _retrying() below +# instead, where a timeout can be excluded. +CHAT_MAX_RETRIES = 0 + +# Failures worth trying again right away. The SDK retries 409, 429 and any 5xx, +# so keep that rule rather than a list that misses one (529 and 522 both reach +# here). The timeout statuses are excluded by _is_timeout_error, so they are +# reported instead of retried. +TRANSIENT_STATUSES = (409, 429) +# First attempt plus two retries, and only while the whole call stays inside the +# budget: a failure that already cost minutes is not "transient", and the user is +# waiting for an answer. +CHAT_ATTEMPTS = 3 +CHAT_RETRY_BUDGET_SECONDS = 60 +CHAT_RETRY_BACKOFF_SECONDS = 0.5 + logger = get_logger(__name__) @@ -67,6 +109,87 @@ def _llm_empty_response_command() -> str: """ return f"(send {quote_arg(LLM_EMPTY_RESPONSE_MESSAGE)})" +def _is_timeout_error(error: BaseException) -> bool: + """True when the request ran out of time rather than failing outright: the + client's own timeout, or a timeout status from the gateway in front of the + provider (the proxy answers 504 when the upstream is still thinking). + """ + # The classes are looked up rather than referenced: classifying a failure must + # never raise one of its own, whatever the installed client exposes. + if isinstance(error, getattr(openai, "APITimeoutError", ())): + return True + return getattr(error, "status_code", None) in GATEWAY_TIMEOUT_STATUSES + +_timeout_notices = 0 + +def _llm_timeout_command() -> str: + """Return a status message as a MeTTa `send` command when the request times + out, so the turn ends with the user told instead of in silence. + + `send` drops a message equal to the last one it sent, so two notices must + never render the same: a second timeout would leave that turn silent, which is + the symptom this whole change is about. The time alone does not guarantee it, + since two cycles can fail inside the same second, so the notice carries the + time to the millisecond and a count that rises with every notice. + """ + global _timeout_notices + _timeout_notices += 1 + stamp = datetime.now().strftime("%H:%M:%S.%f")[:-3] + message = LLM_TIMEOUT_MESSAGE.format( + time=stamp, count=_timeout_notices, timeout=CHAT_REQUEST_TIMEOUT_SECONDS) + return f"(send {quote_arg(message)})" + +def _is_transient_error(error: BaseException) -> bool: + """True for a failure that another attempt may get past. A timeout is not + one of them: it already spent the request timeout, so retrying it only keeps + the user waiting. + """ + if _is_timeout_error(error): + return False + if isinstance(error, getattr(openai, "APIConnectionError", ())): + return True + status = getattr(error, "status_code", None) + if status is None: + return False + return status in TRANSIENT_STATUSES or status >= 500 + +def _retry_delay(error: BaseException, attempt: int) -> float: + """How long to wait before the next attempt: the provider's Retry-After when + it sends one in seconds, otherwise a short exponential backoff. A Retry-After + given as an HTTP date falls back to the backoff. + """ + headers = getattr(getattr(error, "response", None), "headers", None) + value = headers.get("retry-after") if hasattr(headers, "get") else None + if value is not None: + try: + return max(0.0, float(value)) + except (TypeError, ValueError): + pass + return CHAT_RETRY_BACKOFF_SECONDS * (2 ** (attempt - 1)) + +def _retrying(call, provider: str): + """Run call(), retrying only transient failures and only briefly. + + The SDK's own retries are off (CHAT_MAX_RETRIES), so this is the single place + that decides what gets another attempt: transient failures do, a timeout does + not, and nothing is retried once the budget is spent. + """ + started = time.monotonic() + for attempt in range(1, CHAT_ATTEMPTS + 1): + try: + return call() + except Exception as error: + if attempt == CHAT_ATTEMPTS or not _is_transient_error(error): + raise + delay = _retry_delay(error, attempt) + if time.monotonic() - started + delay >= CHAT_RETRY_BUDGET_SECONDS: + logger.warning( + f"[{provider}.chat]: retry budget spent, giving up: {error}") + raise + logger.warning( + f"[{provider}.chat]: transient failure, retrying in {delay:.1f}s: {error}") + time.sleep(delay) + def _split_system_user(content: str) -> Tuple[str, str]: """ MeTTa sends: @@ -129,6 +252,41 @@ def __init__(self, name: str, var_name: str, model_name: str, base_url: str): self._model_name = model_name self._base_url = base_url self._client = None # lazy initialization + self._timeout_notice_sent = False + + def _start_of_turn(self, content: str) -> None: + """A prompt carrying a human message starts a fresh turn, which deserves + its own answer even if the previous cycle already timed out. + + Only a tagged message counts. The tail also carries the spamShield + reminder on every follow-up cycle when that option is on, and treating + that as a turn would bring back one notice per cycle. The tail is read + straight from the prompt rather than through _split_system_user, which + substitutes a placeholder when it is empty. + """ + _, delimiter, tail = content.partition(PROMPT_DELIMITER) + tail = (tail if delimiter else content).lstrip(" ([\"'") + if tail.startswith(HUMAN_MESSAGE_MARKER): + self._timeout_notice_sent = False + + def _answered(self) -> None: + """The provider answered, so the next timeout is a new run.""" + self._timeout_notice_sent = False + + def _timeout_reply(self) -> str: + """The notice for a timed-out request, once per run of timeouts. + + The loop keeps calling for up to maxNewInputLoops cycles, so a notice per + cycle would fill the chat with the same text. The run ends when a call + succeeds or a new human message arrives; until then the timeout is logged + and nothing is sent. + """ + if self._timeout_notice_sent: + logger.warning( + f"[{self.name}.chat]: timed out again, the user was already told") + return "" + self._timeout_notice_sent = True + return _llm_timeout_command() def _ensure_client(self): """Initialize client on first use.""" @@ -145,9 +303,13 @@ def _create_client(self) -> Optional[openai.OpenAI]: return openai.OpenAI( api_key="proxy", base_url=base_url, + max_retries=CHAT_MAX_RETRIES, + timeout=CHAT_REQUEST_TIMEOUT_SECONDS, ) if self._var_name in os.environ: - return openai.OpenAI(api_key=os.environ.get(self._var_name), base_url=self._base_url) + return openai.OpenAI(api_key=os.environ.get(self._var_name), base_url=self._base_url, + max_retries=CHAT_MAX_RETRIES, + timeout=CHAT_REQUEST_TIMEOUT_SECONDS) return None @@ -169,19 +331,24 @@ def _build_messages(self, content: str): def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", **kwargs) -> str: """Send chat request, initializing client if needed.""" + self._start_of_turn(content) self._ensure_client() if self._client is None: raise RuntimeError(f"{self.name} not configured (set {self._var_name})") try: - response = self._client.chat.completions.create( - model=self._model_name, - messages=self._build_messages(content), - max_tokens=max_tokens, - **kwargs + response = _retrying( + lambda: self._client.chat.completions.create( + model=self._model_name, + messages=self._build_messages(content), + max_tokens=max_tokens, + **kwargs + ), + self._name, ) + self._answered() raw = response.choices[0].message.content or "" finish_reason = getattr(response.choices[0], "finish_reason", None) _log_raw(self._name, self._model_name, raw) @@ -194,6 +361,8 @@ def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", return resp except Exception as e: logger.exception(f"[AIProvider.chat]: Exception while communicating with LLM: {e}") + if _is_timeout_error(e): + return self._timeout_reply() return "" def _clean_text(self, text: str) -> str: diff --git a/providers/openai.py b/providers/openai.py index 8d0513539..831c17940 100644 --- a/providers/openai.py +++ b/providers/openai.py @@ -31,6 +31,7 @@ class OpenAIProviderImpl(llm.AIProvider): def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", **kwargs) -> str: """Send chat request via the Responses API, initializing client if needed.""" + self._start_of_turn(content) self._ensure_client() if self._client is None: @@ -53,8 +54,9 @@ def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", create_kwargs.update(kwargs) - response = self._client.responses.create(**create_kwargs) + response = llm._retrying(lambda: self._client.responses.create(**create_kwargs), self._name) + self._answered() raw = response.output_text or "" incomplete_details = getattr(response, "incomplete_details", None) incomplete_reason = getattr(incomplete_details, "reason", None) @@ -67,4 +69,6 @@ def chat(self, content: str, max_tokens: int = 6000, reasoning: str = "medium", return self._clean_text(raw) except Exception as e: logger.exception(f"[OpenAIProviderImpl.chat]: Exception while communicating with LLM: {e}") + if llm._is_timeout_error(e): + return self._timeout_reply() return "" diff --git a/providers/openrouter.py b/providers/openrouter.py index 332ce88a5..2b7bb40e6 100644 --- a/providers/openrouter.py +++ b/providers/openrouter.py @@ -40,9 +40,13 @@ def _create_client(self) -> Optional[openai.OpenAI]: return openai.OpenAI( api_key="proxy", base_url=base_url, + max_retries=llm.CHAT_MAX_RETRIES, + timeout=llm.CHAT_REQUEST_TIMEOUT_SECONDS, ) if self._var_name in os.environ: - return openai.OpenAI(api_key=os.environ.get(self._var_name), base_url=self._base_url) + return openai.OpenAI(api_key=os.environ.get(self._var_name), base_url=self._base_url, + max_retries=llm.CHAT_MAX_RETRIES, + timeout=llm.CHAT_REQUEST_TIMEOUT_SECONDS) return None diff --git a/proxy/nginx.conf.template b/proxy/nginx.conf.template index f569ea97b..88dfcb7ef 100644 --- a/proxy/nginx.conf.template +++ b/proxy/nginx.conf.template @@ -35,6 +35,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } location /asicloud/ { @@ -45,6 +47,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } location /openai/ { @@ -55,6 +59,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } location /asione/ { @@ -65,6 +71,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } location /openaiapi/ { @@ -74,6 +82,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } # OpenClaw Gateway: inject Bearer token, upstream is set at container @@ -97,6 +107,8 @@ http { proxy_ssl_server_name on; proxy_ssl_protocols TLSv1.2 TLSv1.3; proxy_http_version 1.1; + proxy_read_timeout 600s; + proxy_send_timeout 600s; } # Telegram: rewrite URL to inject bot token in path.