@@ -38,9 +38,55 @@ namespace Beam.Broker
3838abbrev brokerStdio : IO.Process.StdioConfig where
3939 stdin := .piped
4040 stdout := .piped
41- -- Keep backend stderr away from MCP stdio. Inheriting it can corrupt
42- -- client framing; piping and draining it caused macOS save_olean hangs.
43- stderr := .null
41+ -- Keep backend stderr away from MCP stdio while retaining a bounded tail for
42+ -- startup and worker-exit diagnostics. The blocking drain runs on a dedicated
43+ -- task so it cannot starve regular Lean tasks.
44+ stderr := .piped
45+
46+ private def backendStderrTailLimit : Nat :=
47+ 16 * 1024
48+
49+ private def backendStderrReadSize : USize :=
50+ 4096
51+
52+ private def isUtf8ContinuationByte (byte : UInt8) : Bool :=
53+ decide (128 ≤ byte.toNat ∧ byte.toNat < 192 )
54+
55+ private def utf8BoundaryAtOrAfter (bytes : ByteArray) (offset : Nat) : Nat :=
56+ let rec loop (offset : Nat) : Nat → Nat
57+ | 0 => offset
58+ | fuel + 1 =>
59+ if h : offset < bytes.size then
60+ if isUtf8ContinuationByte bytes[offset] then
61+ loop (offset + 1 ) fuel
62+ else
63+ offset
64+ else
65+ offset
66+ loop offset 3
67+
68+ structure BackendStderrCapture where
69+ tail : Std.Mutex ByteArray
70+ drainTask : Task (Except IO.Error Unit)
71+
72+ private partial def drainBackendStderr
73+ (stderr : IO.FS.Handle)
74+ (tail : Std.Mutex ByteArray) : IO Unit := do
75+ let chunk ← stderr.read backendStderrReadSize
76+ unless chunk.isEmpty do
77+ tail.atomically do
78+ let combined := (← get) ++ chunk
79+ if combined.size > backendStderrTailLimit then
80+ let start := utf8BoundaryAtOrAfter combined (combined.size - backendStderrTailLimit)
81+ set <| combined.extract start combined.size
82+ else
83+ set combined
84+ drainBackendStderr stderr tail
85+
86+ def startBackendStderrCapture (stderr : IO.FS.Handle) : IO BackendStderrCapture := do
87+ let tail ← Std.Mutex.new ByteArray.empty
88+ let drainTask ← IO.asTask (prio := Task.Priority.dedicated) <| drainBackendStderr stderr tail
89+ pure { tail, drainTask }
4490
4591structure Session where
4692 workspaceId : WorkspaceId
@@ -51,6 +97,7 @@ structure Session where
5197 proc : IO.Process.Child brokerStdio
5298 stdin : IO.FS.Stream
5399 stdout : IO.FS.Stream
100+ stderrCapture : BackendStderrCapture
54101 pending : PendingRequestStore
55102 nextId : Nat := 1
56103 nextEventSeq : Nat := 1
@@ -155,6 +202,31 @@ private partial def waitForTaskWithTimeout
155202 loop (remainingMs - min pollMs remainingMs)
156203 loop timeoutMs
157204
205+ private def backendName : Backend → String
206+ | .lean => "Lean"
207+ | .rocq => "Rocq"
208+
209+ private def BackendStderrCapture.snapshot (capture : BackendStderrCapture) : IO String := do
210+ let bytes ← capture.tail.atomically get
211+ pure <| (String.fromUTF8? bytes).getD "<backend stderr tail is not valid UTF-8>"
212+
213+ private def BackendStderrCapture.awaitDrain
214+ (capture : BackendStderrCapture)
215+ (timeoutMs : Nat := 500 ) : IO Unit := do
216+ discard <| waitForTaskWithTimeout capture.drainTask timeoutMs
217+
218+ private def backendFailureMessage
219+ (backend : Backend)
220+ (phase cause : String)
221+ (capture : BackendStderrCapture) : IO String := do
222+ let stderr := (← capture.snapshot).trimAscii.toString
223+ let stderr := if stderr.isEmpty then "<empty>" else stderr
224+ pure <| String.intercalate "\n " [
225+ s! "{ backendName backend} backend failed { phase} : { cause} " ,
226+ s! "backend stderr tail (last { backendStderrTailLimit} bytes):" ,
227+ stderr
228+ ]
229+
158230private def sessionShutdownReplyTimeoutMs : Nat :=
159231 1000
160232
@@ -209,6 +281,25 @@ private def terminateBackendProcess (proc : IO.Process.Child brokerStdio) : IO U
209281 catch _ =>
210282 pure ()
211283
284+ private def startBackendStderrCaptureOrTerminate
285+ (backend : Backend)
286+ (proc : IO.Process.Child brokerStdio) : IO BackendStderrCapture := do
287+ try
288+ startBackendStderrCapture proc.stderr
289+ catch err =>
290+ terminateBackendProcess proc
291+ throw <| IO.userError <|
292+ s! "{ backendName backend} backend failed during startup before stderr capture: { err} "
293+
294+ private def terminateBackendFailure
295+ (backend : Backend)
296+ (phase cause : String)
297+ (proc : IO.Process.Child brokerStdio)
298+ (capture : BackendStderrCapture) : IO String := do
299+ terminateBackendProcess proc
300+ capture.awaitDrain
301+ backendFailureMessage backend phase cause capture
302+
212303private def sessionExited (session : Session) : IO Bool := do
213304 try
214305 pure (← session.proc.tryWait).isSome
@@ -453,11 +544,13 @@ partial def sessionReaderLoop (session : Session) : IO Unit := do
453544 pure ()
454545 sessionReaderLoop session
455546 catch e =>
547+ let message ←
548+ terminateBackendFailure session.backend "after startup" e.toString
549+ session.proc session.stderrCapture
456550 PendingRequestStore.failAll session.pending <| BrokerFailure.toResponseFailure {
457551 code := .workerExited
458- message := e.toString
552+ message
459553 }
460- terminateBackendProcess session.proc
461554
462555private def startRequestJsonTrackedDetailed
463556 (session : Session)
@@ -589,6 +682,7 @@ private def acquireBackendSession
589682 env := env
590683 cwd := root.toString
591684 }
685+ let stderrCapture ← startBackendStderrCaptureOrTerminate backend proc
592686 let (session, initializeTask) ←
593687 try
594688 let stdin := IO.FS.Stream.ofHandle proc.stdin
@@ -604,6 +698,7 @@ private def acquireBackendSession
604698 proc
605699 stdin
606700 stdout
701+ stderrCapture
607702 pending
608703 }
609704 writeLspRequest stdin
@@ -613,8 +708,8 @@ private def acquireBackendSession
613708 awaitInitializeResponse stdout
614709 pure (session, initializeTask)
615710 catch err =>
616- terminateBackendProcess proc
617- throw err
711+ throw <| IO.userError <| ←
712+ terminateBackendFailure backend "during startup" err.toString proc stderrCapture
618713 try
619714 match ← waitForTaskWithTimeout initializeTask backendInitializeTimeoutMs with
620715 | some (.ok ()) => pure ()
@@ -632,9 +727,10 @@ private def acquireBackendSession
632727 pure session
633728 catch err =>
634729 IO.cancel initializeTask
635- terminateBackendProcess proc
730+ let message ←
731+ terminateBackendFailure backend "during startup" err.toString proc stderrCapture
636732 discard <| waitForTaskWithTimeout initializeTask sessionShutdownReplyTimeoutMs
637- throw err
733+ throw <| IO.userError message
638734
639735private def requireWorkspace (workspaceId : WorkspaceId) : M WorkspaceState := do
640736 let state ← get
0 commit comments