diff --git a/.github/workflows/studio-inference-smoke.yml b/.github/workflows/studio-inference-smoke.yml index f540c11da4..cf8c021e38 100644 --- a/.github/workflows/studio-inference-smoke.yml +++ b/.github/workflows/studio-inference-smoke.yml @@ -485,7 +485,7 @@ jobs: print(f"[retry] {path}: {exc!r}", flush = True) time.sleep(15) - def post_sse(path, body, *, timeout = 600): + def post_sse(path, body, *, timeout = 600, retries = 1, complete_on = None): """POST a streaming request and accumulate the assistant text deltas. The server-side agentic loop ALWAYS returns SSE regardless of the request's `stream` field, so any @@ -501,6 +501,22 @@ jobs: invocation markers / tool output, since `delta.content` alone is not evidence that the tool path executed. + + A shared CI runner can stall the stream transport (the + connection opening, or a mid-stream read) even when Studio + is healthy, so retry a stall once with a fresh request + capped at 300s. A stall means the stream did NOT complete, + so partial events are normally NOT returned (an early + tool_start with no tool_end is not proof the tool loop + finished). The one exception is `complete_on`: an optional + predicate over the events collected so far -- when a stall + happens after it is already satisfied (the tool ran and + produced its result before the trailing read timed out), + those events are returned rather than discarded, so the + stall-after-answer case still counts. HTTP status errors + surface immediately; a stall that yields no completed result + across all attempts re-raises so the caller can rotate to + the next seed. """ body = {**body, "stream": True} data = json.dumps(body).encode() @@ -513,26 +529,45 @@ jobs: "Content-Type": "application/json", }, ) - parts = [] - events = [] - with urllib.request.urlopen(req, timeout = timeout) as resp: - for raw in resp: - line = raw.decode().strip() - if not line.startswith("data: "): - continue - payload = line[6:] - if payload == "[DONE]": - break - events.append(payload) - try: - chunk = json.loads(payload) - except json.JSONDecodeError: - continue - for choice in chunk.get("choices", []): - delta = choice.get("delta", {}) or {} - if delta.get("content"): - parts.append(delta["content"]) - return "".join(parts), events + for attempt in range(retries + 1): + parts = [] + events = [] + t = timeout if attempt == 0 else min(timeout, 300) + try: + with urllib.request.urlopen(req, timeout = t) as resp: + for raw in resp: + line = raw.decode().strip() + if not line.startswith("data: "): + continue + payload = line[6:] + if payload == "[DONE]": + break + events.append(payload) + try: + chunk = json.loads(payload) + except json.JSONDecodeError: + continue + for choice in chunk.get("choices", []): + delta = choice.get("delta", {}) or {} + if delta.get("content"): + parts.append(delta["content"]) + return "".join(parts), events + except urllib.error.HTTPError: + raise + except (TimeoutError, ConnectionError, urllib.error.URLError) as exc: + # A stall after the tool already produced its result is + # the case this probe exists to tolerate: keep those + # events. But a stall with only an early tool_start (no + # completed output) is not proof the tool loop finished, + # so it must not pass -- retry once, then raise so + # _run_tool_probe rotates to the next seed. + if complete_on is not None and complete_on(events): + print(f"[retry-sse] {path}: {exc!r}; keeping {len(events)} completed events", flush = True) + return "".join(parts), events + if attempt == retries: + raise + print(f"[retry-sse] {path}: {exc!r}", flush = True) + time.sleep(15) _STUDIO_TOOL_TYPES = { "tool_start", "tool_end", "tool_use", "tool_result", @@ -669,17 +704,54 @@ jobs: """ attempts_log = [] best = None + # Cap the wall-clock spent rotating through stalled seeds so a + # persistent no-data wedge fails fast (clean assertion) instead + # of being killed by the job's timeout-minutes. A healthy or + # merely degenerate round answers in seconds, so all seeds still + # run in the normal case; only stalls consume the budget. + probe_deadline = time.monotonic() + 300 for attempt_i in range(max_attempts): + # Cap each read by the budget still remaining (not just a flat + # 180s) and skip an attempt too small to finish, so the whole + # rotation stays within ~300s -- two probes then fit the job's + # timeout-minutes even if every seed stalls. + remaining = int(probe_deadline - time.monotonic()) + if attempt_i and remaining < 30: + print(f"[tools] {label}: seed-rotation budget spent after {attempt_i} attempts", flush = True) + break attempt_seed = SEED + attempt_i - content, events = post_sse("/v1/chat/completions", { - "messages": [{"role": "user", "content": prompt}], - "enable_tools": True, - "enabled_tools": enabled, - "session_id": f"{session}-att{attempt_i}", - "temperature": TOOL_PROBE_TEMP, - "seed": attempt_seed, - "max_tokens": 600, - }) + try: + # Bounded per-attempt timeout, no inner retry -- the seed + # loop IS the retry, so a stall raises quickly and rotates + # rather than spending post_sse's full 600+300s. complete_on + # keeps a stall that already produced the tool result (only + # the trailing read timed out) instead of discarding it. + content, events = post_sse("/v1/chat/completions", { + "messages": [{"role": "user", "content": prompt}], + "enable_tools": True, + "enabled_tools": enabled, + "session_id": f"{session}-att{attempt_i}", + "temperature": TOOL_PROBE_TEMP, + "seed": attempt_seed, + "max_tokens": 600, + }, timeout = min(180, remaining), retries = 0, + complete_on = lambda ev: _tool_invoked(ev) and _tool_output_contains(ev, *needles)) + except urllib.error.HTTPError: + # HTTPError subclasses URLError, so re-raise a real 4xx/5xx + # here instead of letting the transport-stall handler below + # swallow it and rotate seeds -- an endpoint status failure + # must surface, not be masked as missing tool evidence. + raise + except (TimeoutError, ConnectionError, urllib.error.URLError) as exc: + # A transport stall that outlived post_sse's own retry: + # log it as a failed attempt and rotate to the next seed + # rather than sinking the whole probe on one bad stream. + attempts_log.append({ + "attempt": attempt_i, "seed": attempt_seed, + "transport_error": repr(exc), + }) + print(f"[tools] retry {label} attempt {attempt_i}: transport {exc!r}", flush = True) + continue invoked = _tool_invoked(events) produced = _tool_output_contains(events, *needles) attempts_log.append({ @@ -740,6 +812,9 @@ jobs: # enough that requiring a tool_call marker would create # red-herring failures from infra rather than from Studio. try: + # Best-effort and bounded: a single 180s attempt keeps a stall + # from eating the job's timeout-minutes (it already WARNs, so a + # retry buys nothing). content, events = post_sse("/v1/chat/completions", { "messages": [{"role": "user", "content": "Search the web for 'unsloth ai github' and summarise."}], "enable_tools": True, @@ -748,7 +823,7 @@ jobs: "temperature": 0.0, "seed": SEED, "max_tokens": 400, - }) + }, timeout = 180, retries = 0) print( f"[tools] PASS web_search stream ({len(content)} chars in content, " f"{len(events)} raw events)" diff --git a/.github/workflows/studio-mac-inference-smoke.yml b/.github/workflows/studio-mac-inference-smoke.yml index 03c0a8580d..d3d765aa84 100644 --- a/.github/workflows/studio-mac-inference-smoke.yml +++ b/.github/workflows/studio-mac-inference-smoke.yml @@ -471,11 +471,22 @@ jobs: print(f"[retry] {path}: {exc!r}", flush = True) time.sleep(15) - def post_sse(path, body, *, timeout = 600): + def post_sse(path, body, *, timeout = 600, retries = 1, soft = False): """POST a streaming request and accumulate the assistant text deltas. The server-side agentic loop ALWAYS returns SSE regardless of the request's `stream` field, so any - call with enable_tools=true must use this helper.""" + call with enable_tools=true must use this helper. + + A shared CI runner can stall the stream transport (the + connection opening, or a mid-stream read) even when Studio + is healthy, so harden the read three ways: retry a stall + once with a fresh request capped at 300s; return any text + already streamed before a stall (a stall on the trailing + tokens, after the answer arrived, still counts); and when + every attempt yields nothing, a hard call re-raises while a + soft call (the best-effort server-side tool probes) returns + None so the caller can WARN instead of sinking the whole + job. HTTP status errors always surface immediately.""" body = {**body, "stream": True} data = json.dumps(body).encode() req = urllib.request.Request( @@ -487,24 +498,43 @@ jobs: "Content-Type": "application/json", }, ) - parts = [] - with urllib.request.urlopen(req, timeout = timeout) as resp: - for raw in resp: - line = raw.decode().strip() - if not line.startswith("data: "): - continue - payload = line[6:] - if payload == "[DONE]": - break - try: - chunk = json.loads(payload) - except json.JSONDecodeError: - continue - for choice in chunk.get("choices", []): - delta = choice.get("delta", {}) or {} - if delta.get("content"): - parts.append(delta["content"]) - return "".join(parts) + for attempt in range(retries + 1): + parts = [] + t = timeout if attempt == 0 else min(timeout, 300) + try: + with urllib.request.urlopen(req, timeout = t) as resp: + for raw in resp: + line = raw.decode().strip() + if not line.startswith("data: "): + continue + payload = line[6:] + if payload == "[DONE]": + break + try: + chunk = json.loads(payload) + except json.JSONDecodeError: + continue + for choice in chunk.get("choices", []): + delta = choice.get("delta", {}) or {} + if delta.get("content"): + parts.append(delta["content"]) + return "".join(parts) + except urllib.error.HTTPError: + raise + except (TimeoutError, ConnectionError, urllib.error.URLError) as exc: + # Text already streamed is a valid signal -- keep it + # rather than re-running a heavy generation. + if parts: + joined = "".join(parts) + print(f"[retry-sse] {path}: {exc!r}; keeping {len(joined)} partial chars", flush = True) + return joined + if attempt == retries: + if soft: + print(f"[tools] WARN {path}: SSE transport stalled with no data ({exc!r}) -- non-blocking", flush = True) + return None + raise + print(f"[retry-sse] {path}: {exc!r}", flush = True) + time.sleep(15) # ── 1. Standard OpenAI function calling ────────────────────── weather_tool = { @@ -575,6 +605,10 @@ jobs: # macos-14 free runner is ~10 tok/s on Qwen3.5-2B Q4_K_XL; # cap max_tokens tightly so each SSE round stays under ~30s # even when the model stalls in a degenerate output state. + # retries=0 on the best-effort probes: this job's 25-minute cap + # allows a 10-minute model load, so a no-data stall must be a + # single 180s attempt (not 180+15+180s) to leave room for the + # thinking checks. A soft/best-effort probe only WARNs anyway. content = post_sse("/v1/chat/completions", { "messages": [{"role": "user", "content": "What is 123 * 456? Use the python tool to compute it and tell me the number."}], "enable_tools": True, @@ -583,8 +617,10 @@ jobs: "temperature": TEMP, "seed": SEED, "max_tokens": 128, - }, timeout = 180) - if "56088" in content or "56,088" in content: + }, timeout = 180, retries = 0, soft = True) + if content is None: + print("[tools] WARN python tool: SSE transport stalled after retries -- non-blocking") + elif "56088" in content or "56,088" in content: print(f"[tools] PASS python tool ({len(content)} chars, found 56088)") else: # Empty stream is a known Mac-quant degeneracy too; log @@ -616,7 +652,7 @@ jobs: "temperature": TEMP, "seed": SEED, "max_tokens": 96, - }, timeout = 180) + }, timeout = 180, retries = 0) print(f"[tools] PASS web_search stream ({len(content)} chars)") except Exception as exc: print(f"[tools] WARN web_search probe failed (non-blocking): {exc}") diff --git a/.github/workflows/studio-windows-inference-smoke.yml b/.github/workflows/studio-windows-inference-smoke.yml index 0453c9212a..233292f7a3 100644 --- a/.github/workflows/studio-windows-inference-smoke.yml +++ b/.github/workflows/studio-windows-inference-smoke.yml @@ -677,7 +677,22 @@ jobs: print(f"[retry] {path}: {exc!r}", flush = True) time.sleep(15) - def post_sse(path, body, *, timeout = 600): + def post_sse(path, body, *, timeout = 600, retries = 1, soft = False): + # The server-side agentic loop always answers over SSE. A + # shared CI runner can stall the stream transport (the + # connection opening, or a mid-stream read) even when Studio + # is healthy, so harden the read three ways: + # * retry a transport stall once with a fresh request, + # capped at 300s (a healthy server answers a retry + # quickly, a wedged one never does); + # * return any text already streamed before a stall, so a + # stall on the trailing tokens -- after the answer + # arrived -- still counts; + # * when every attempt yields nothing, a hard call + # re-raises while a soft call (the best-effort + # server-side tool probes) returns None so the caller + # can WARN instead of sinking the whole job. + # HTTP status errors always surface immediately. body = {**body, "stream": True} data = json.dumps(body).encode() req = urllib.request.Request( @@ -689,24 +704,43 @@ jobs: "Content-Type": "application/json", }, ) - parts = [] - with urllib.request.urlopen(req, timeout = timeout) as resp: - for raw in resp: - line = raw.decode().strip() - if not line.startswith("data: "): - continue - payload = line[6:] - if payload == "[DONE]": - break - try: - chunk = json.loads(payload) - except json.JSONDecodeError: - continue - for choice in chunk.get("choices", []): - delta = choice.get("delta", {}) or {} - if delta.get("content"): - parts.append(delta["content"]) - return "".join(parts) + for attempt in range(retries + 1): + parts = [] + t = timeout if attempt == 0 else min(timeout, 300) + try: + with urllib.request.urlopen(req, timeout = t) as resp: + for raw in resp: + line = raw.decode().strip() + if not line.startswith("data: "): + continue + payload = line[6:] + if payload == "[DONE]": + break + try: + chunk = json.loads(payload) + except json.JSONDecodeError: + continue + for choice in chunk.get("choices", []): + delta = choice.get("delta", {}) or {} + if delta.get("content"): + parts.append(delta["content"]) + return "".join(parts) + except urllib.error.HTTPError: + raise + except (TimeoutError, ConnectionError, urllib.error.URLError) as exc: + # Text already streamed is a valid signal -- keep it + # rather than re-running a heavy generation. + if parts: + joined = "".join(parts) + print(f"[retry-sse] {path}: {exc!r}; keeping {len(joined)} partial chars", flush = True) + return joined + if attempt == retries: + if soft: + print(f"[tools] WARN {path}: SSE transport stalled with no data ({exc!r}) -- non-blocking", flush = True) + return None + raise + print(f"[retry-sse] {path}: {exc!r}", flush = True) + time.sleep(15) # ── 1. Standard OpenAI function calling ────────────────────── weather_tool = { @@ -749,6 +783,11 @@ jobs: ) # ── 2. Server-side python tool ─────────────────────────────── + # Bound each soft probe to a single 180s attempt (timeout=180, + # retries=0): this job runs two of them back-to-back under a + # 30-minute cap, so the default 600+15+300s per stall could hit + # the workflow timeout before the thinking checks run. A soft + # probe only WARNs anyway, so a retry buys nothing. content = post_sse("/v1/chat/completions", { "messages": [{"role": "user", "content": "What is 123 * 456? Use the python tool to compute it and tell me the number."}], "enable_tools": True, @@ -757,8 +796,10 @@ jobs: "temperature": TEMP, "seed": SEED, "max_tokens": 600, - }) - if "56088" in content or "56,088" in content: + }, timeout = 180, retries = 0, soft = True) + if content is None: + print("[tools] WARN python tool: SSE transport stalled after retries -- non-blocking") + elif "56088" in content or "56,088" in content: print(f"[tools] PASS python tool ({len(content)} chars, found 56088)") else: assert content, "python tool: SSE stream empty" @@ -780,8 +821,10 @@ jobs: "temperature": TEMP, "seed": SEED, "max_tokens": 600, - }) - if "hello-bash-tool" in content: + }, timeout = 180, retries = 0, soft = True) + if content is None: + print("[tools] WARN terminal tool: SSE transport stalled after retries -- non-blocking") + elif "hello-bash-tool" in content: print(f"[tools] PASS terminal tool ({len(content)} chars)") else: assert content, "terminal tool: SSE stream empty" @@ -802,7 +845,7 @@ jobs: "temperature": TEMP, "seed": SEED, "max_tokens": 400, - }) + }, timeout = 180, retries = 0) print(f"[tools] PASS web_search stream ({len(content)} chars)") except Exception as exc: print(f"[tools] WARN web_search probe failed (non-blocking): {exc}")