Repository navigation
Add DSPy RLM example with an E2B interpreter - #150
max-sudolabs wants to merge 1 commit into
Conversation
| message = json.loads(line) | ||
| if not isinstance(message, dict) or not isinstance(message.get("type"), str): | ||
| raise CodeInterpreterError("Invalid E2B response frame") | ||
| self._responses.put(message, timeout=0.1) |
There was a problem hiding this comment.
🔴 Generated code that fires a burst of more than about 32 concurrent host-tool calls kills the whole session. The model should get per-call tool errors instead. The reader thread puts frames into _responses (maxsize=32) with timeout=0.1. If the main loop is busy, queue.Full is caught as a reader failure, so _receive raises "session is terminal" and the sandbox is killed. Fix: never let a full queue be fatal. Either make the queue unbounded, since frames are already capped by MAX_INPUT_BYTES, or have the reader block until _stop is set rather than time out.
Why this was flagged
Suppose model-generated code runs something like ThreadPoolExecutor(max_workers=50) calling llm_query. The worker's Protocol.call_tool (_worker.py:72-87) then sends 50 tool_request frames almost at once. On the host, execute (interpreter.py:292-320) starts threads for the first 8. For each request after that it makes a blocking network call, _send -> send_stdin, to return "Eight host tools are already active". Meanwhile the reader thread keeps calling self._responses.put(message, timeout=0.1) (interpreter.py:175) on a queue created with maxsize=32 (interpreter.py:68). Once it fills, queue.Full is raised and caught by except BaseException at interpreter.py:180-182, which sets _reader_error. The next _receive then raises "E2B worker/transport failed; session is terminal". execute shuts down and kills the sandbox, and all persistent CPython/SQLite state is lost. The 8-slot tool limit does not prevent this, because rejected requests still pass through the bounded queue and each rejection costs a network round trip. This is a new file; nothing like it exists on the base branch.
Verification: Severity: normal. Trigger: generated code issues roughly 40 or more host-tool calls at once, and a single host round trip takes longer than 100 ms. The reader's self._responses.put(message, timeout=0.1) (interpreter.py:175) on a queue with maxsize=32 (interpreter.py:68) raises Full, which is caught at 180-182, and execute (341-343) then calls shutdown(), which kills the sandbox.
| if len(payload.encode()) > MAX_FRAME_BYTES and isinstance(message.get("stdout"), str): | ||
| # JSON escaping and response metadata also consume the frame budget. | ||
| stdout = message["stdout"].removesuffix(TRUNCATION_MARKER) | ||
| low, high = 0, len(stdout) | ||
| while low < high: | ||
| middle = (low + high + 1) // 2 | ||
| message["stdout"] = stdout[:middle] + TRUNCATION_MARKER | ||
| candidate = json.dumps(message, separators=(",", ":"), allow_nan=False) | ||
| if len(candidate.encode()) <= MAX_FRAME_BYTES: | ||
| low = middle | ||
| else: | ||
| high = middle - 1 | ||
| message["stdout"] = stdout[:low] + TRUNCATION_MARKER | ||
| payload = json.dumps(message, separators=(",", ":"), allow_nan=False) | ||
| if len(payload.encode()) > MAX_FRAME_BYTES: | ||
| payload = json.dumps({"type": "terminal_error", "error": "response frame exceeds 8 MiB"}) |
There was a problem hiding this comment.
🔴 A step whose last expression returns a large value (for example records on the 50,008-record input) kills the whole RLM session instead of being truncated. send only shrinks stdout. When value makes the frame larger than MAX_FRAME_BYTES, it sends terminal_error and the host deletes the sandbox. Fix: give the result/final value an output budget the same way stdout has one: past the limit, swap it for a truncated repr plus TRUNCATION_MARKER. Keep terminal_error only for frames that cannot be shrunk.
Why this was flagged
Trigger: main.py --background-records 50000 loads ~50k records (~10 MB JSON) as the records variable. The model then writes a step ending in a bare expression such as records or cur.execute(...).fetchall(). Session.execute (_worker.py:207) stores the full jsonable(value). send (_worker.py:241) only truncates when stdout is a str and only shortens stdout, so the frame stays over 8 MiB. Line 256 then replaces it with {"type":"terminal_error","error":"response frame exceeds 8 MiB"}. On the host, interpreter.py:336-337 raises CodeInterpreterError and line 343 calls shutdown(), which kills the sandbox. Every later execute then fails in _active() with "session ended", so the RLM loses all SQLite/variable state and wastes its remaining iterations. Printing the same data would just be truncated to 1 MiB by MAX_CAPTURE_BYTES, so the size safeguard covers stdout but not the expression value.
Verification: On the --background-records 50000 input, a step ending in bare records produces JSON over 8 MiB. In send, the truncation branch (241-254) only shrinks stdout, so an oversized value from _worker.py:207 stays oversized. Line 255-256 then swaps the whole frame for a terminal_error. On the host, interpreter.py:336-337 raises and the handler at 341-343 calls shutdown(), which kills the sandbox (line 373).
| _protocol_output = os.fdopen(os.dup(sys.__stdout__.fileno()), "w", encoding="utf-8") | ||
| os.set_inheritable(_protocol_output.fileno(), False) | ||
| _worker_stdout, _worker_stderr = sys.stdout, sys.stderr | ||
| _sink = os.open(os.devnull, os.O_WRONLY) | ||
| os.dup2(_sink, 1) | ||
| os.dup2(_sink, 2) | ||
| os.close(_sink) |
There was a problem hiding this comment.
🟡 (optional) Generated code that reads stdin, or starts a subprocess that does, consumes the host's protocol frames. The tool call then hangs until the deadline kills the session. The worker moves its protocol output to a private, non-inheritable fd and points fd 1/2 at /dev/null. fd 0 and sys.stdin are left as the live protocol input that receive() reads. Fix: dup stdin to a private non-inheritable fd for receive(), and point fd 0 and sys.stdin at /dev/null, the same way stdout/stderr are handled. This covers both input() and inherited child stdin.
Why this was flagged
receive() (_worker.py:261) reads frames from sys.__stdin__.buffer, which is fd 0. Lines 20-26 make only the protocol output private and send fds 1 and 2 to devnull. fd 0 stays inheritable, and sys.stdin is never replaced, including in CapturedOutput.__enter__ (_worker.py:93). The interpreter advertises subprocess. If generated code runs a child that reads stdin (for example subprocess.run(['sqlite3', 'db.sqlite']), cat, or python3 with no script) or calls input(), it competes with Protocol._read for incoming bytes. It can swallow the host's tool_result frames sent by _tool (interpreter.py:246), or block forever. The worker-side call_tool (_worker.py:79) waits with no timeout, so the step hangs until execution_timeout. The host then raises "E2B execution deadline exceeded; terminating owned sandbox" and the whole RLM invocation fails, where a recoverable error was possible. Nothing in _validate or the leaked-thread check prevents this.
Verification: Trigger: a generated step calls input()/sys.stdin.read(), or runs a subprocess that reads stdin. _worker.py:20-26 makes only the protocol OUTPUT private. fd 0 is never touched. receive() at :261 reads frames from sys.__stdin__.buffer.readline(...), which is fd 0. In both cases the host's _receive(deadline) raises "E2B execution deadline exceeded" (interpreter.py:198-199).
| ) | ||
| result = rlm(records=records, query=QUERY) | ||
| artifact = { | ||
| "answer": {name: getattr(result, name) for name in IncidentAnalysis.output_fields}, |
There was a problem hiding this comment.
🟡 (optional) When the model's IDs or counts are wrong, users get only a traceback: no result.json and no printed trajectory. assess_fixture raises ValueError while the artifact dict is being built at main.py:117. The write and print never run, so the answer and trajectory are lost exactly when they are needed. Fix: catch the assessment failure and record it as a field such as "fixture_assessment": {"exact_ids_and_counts": "failed", ...}. Then always write and print the artifact, and exit non-zero afterwards if you want a failing status.
Why this was flagged
The trigger is any run where the model submits candidate_ids that differ from the SQL filter, unsorted or duplicate queue_incident_ids, or counts_by_region that do not match. That is normal for a model-driven run, and the README says model output is not promised. In that case assess_fixture raises ValueError("Output IDs do not match the deterministic candidate filter") or ValueError("Region counts do not match the submitted IDs") (main.py:75, main.py:79). The raise happens inside the artifact = {...} literal at main.py:115-119, so Path("result.json").write_text and print at main.py:120-121 never run. The README tells users to review the trajectory in result.json, but after a mismatch they cannot. The paid model run's answer and trajectory are thrown away and only a stack trace is left. No safeguard writes a partial artifact first.
Verification: In main.py:117-121 the artifact dict literal calls assess_fixture(records, result) while it is being built (line 119). It raises ValueError at line 75 or 79 when IDs or counts differ. Nothing catches either error, so Path("result.json").write_text(...) (line 122) and print(...) (line 123) never run. The answer and result.trajectory are thrown away.
Adds a locally installable
E2BInterpreterfor DSPy 3.4.0's nativeinterpreter_factoryinterface and an incident-analysis RLM example. Generated Python runs in one owned, network-off E2B sandbox per invocation, with persistent CPython/SQLite state, host-side model tools, execution deadlines and explicit cleanup.The example supports locally generated background records, uses
openai/gpt-5.6-lunaby default, and checks exact IDs/counts separately from semantic model output. The README covers installation, limits and cleanup; the cookbook catalog and a scoped offline CI workflow make the example discoverable and maintainable. This is a cookbook project, with no separate package release.Validation:
Related to INT-161. The companion docs guide will remain a draft until this example lands.