A program never waits on an editor that stops reading its pushes, and the editor keeps a bounded tail of the program's output

This commit is contained in:
Joseph Ferano 2026-09-25 11:42:54 +07:00
parent c671f23b2c
commit 20fa04d9df
9 changed files with 216 additions and 75 deletions

View File

@ -2462,9 +2462,15 @@ hundred, and a single line still arrives within 50ms. `capacity` still bounds wh
**A push never lands inside a reply.** `serve` is one thread, so nothing is pushed while a request is handled; what **A push never lands inside a reply.** `serve` is one thread, so nothing is pushed while a request is handled; what
the program prints meanwhile stays in `t.out`, `with_output` puts it on that request's reply and empties the buffer, the program prints meanwhile stays in `t.out`, `with_output` puts it on that request's reply and empties the buffer,
and it is not pushed a second time. The stdout pipe is also drained now while an editor sits idle, which it was not: and it is not pushed a second time. The stdout pipe is also drained now while an editor sits idle, which it was not:
`Wire.recv` blocked, so an attached, quiet editor let the pipe fill until its next request. A push goes out only when `Wire.recv` blocked, so an attached, quiet editor let the pipe fill until its next request. **The program never waits on
the socket is writable; an editor that has stopped reading holds pushes back rather than stalling the loop, and its the editor.** A push is written without blocking from a per-connection buffer, and nothing new is framed while that
output waits in `t.out`. buffer holds bytes the socket would not take — so an editor that stops reading, even partway through a frame, holds
back one frame and no more. What the program prints meanwhile goes on being drained into `t.out`, which keeps its last
256K and drops the oldest bytes; `take` then puts a line in the output saying how many were dropped. Only a reply is
written blocking, after the buffer has been flushed ahead of it, because the editor asked for it and is reading.
**The editor's side is bounded too.** `*flan*` and the REPL keep `flan-output-maximum-lines` (5000) and delete from the
top, as `comint-buffer-maximum-size` does, so a program that prints every frame cannot grow Emacs without limit.
**The watch table moved onto the same channel.** `watch-enable` with `:interval` arms a push of the table from the **The watch table moved onto the same channel.** `watch-enable` with `:interval` arms a push of the table from the
same loop, sent every interval whether or not it changed, because ghost text is repainted from each one. The editor's watch timer, its in-flight flag and same loop, sent every interval whether or not it changed, because ghost text is repainted from each one. The editor's watch timer, its in-flight flag and

View File

@ -1206,6 +1206,7 @@ in the buffer).
| `flan-names-shown` | `4` | how many names to list before summarising | | `flan-names-shown` | `4` | how many names to list before summarising |
| `flan-poll-interval` | `1.0` | seconds between checks for whether it stopped | | `flan-poll-interval` | `1.0` | seconds between checks for whether it stopped |
| `flan-daemon-buffer` | `"*flan*"` | the daemon's own log; mirrors program output | | `flan-daemon-buffer` | `"*flan*"` | the daemon's own log; mirrors program output |
| `flan-output-maximum-lines` | `5000` | lines `*flan*` and the REPL keep; older ones are deleted |
| `flan-diagnostics-buffer` | `"*flan-diagnostics*"` | everything the compiler reports, kept | | `flan-diagnostics-buffer` | `"*flan-diagnostics*"` | everything the compiler reports, kept |
| `flan-start-timeout` | `60` | seconds to wait for a program to come up | | `flan-start-timeout` | `60` | seconds to wait for a program to come up |
| `flan-lower-buffer` | `"*flan-lowering*"` | where `C-c C-l` writes | | `flan-lower-buffer` | `"*flan-lowering*"` | where `C-c C-l` writes |

View File

@ -273,7 +273,8 @@ prompt's line, so the prompt and whatever is being typed at it do not move."
;; text — a plain marker does not advance — and the value the ;; text — a plain marker does not advance — and the value the
;; reply carries would then land above the output it caused. ;; reply carries would then land above the output it caused.
(when (>= at (marker-position mark)) (when (>= at (marker-position mark))
(set-marker mark (point))))))))) (set-marker mark (point)))))
(flan--trim-lines)))))
(defun flan-repl--buffer () (defun flan-repl--buffer ()
"The REPL buffer, for the two clear commands, from wherever they are run." "The REPL buffer, for the two clear commands, from wherever they are run."

View File

@ -355,7 +355,7 @@ Guarded on the row's length so a short name cannot be claimed by accident."
(format "error:\n%s\n" (or (plist-get reply :message) "refused")))))) (format "error:\n%s\n" (or (plist-get reply :message) "refused"))))))
;; The daemon sends the table on the connection every `flan-watch-interval' ;; The daemon sends the table on the connection every `flan-watch-interval'
;; while it is armed and has changed, the same way it sends the program's ;; while it is armed, changed or not, the same way it sends the program's
;; output. Nothing here sends a request on a timer, so nothing here can leave ;; output. Nothing here sends a request on a timer, so nothing here can leave
;; a reply in flight for another request to read. ;; a reply in flight for another request to read.
;; ;;

View File

@ -397,6 +397,21 @@ Pushes that arrive first are handled as they come."
(bound-and-true-p flan-repl-buffer) (bound-and-true-p flan-repl-buffer)
(get-buffer flan-repl-buffer))) (get-buffer flan-repl-buffer)))
(defcustom flan-output-maximum-lines 5000
"Lines the daemon's buffer and the REPL keep; older lines are deleted.
A program that prints every frame would otherwise grow them without limit.
nil keeps everything."
:type '(choice integer (const :tag "No limit" nil)))
(defun flan--trim-lines ()
"Delete lines from the top of this buffer past `flan-output-maximum-lines'."
(when flan-output-maximum-lines
(save-excursion
(goto-char (point-max))
(when (zerop (forward-line (- flan-output-maximum-lines)))
(let ((inhibit-read-only t))
(delete-region (point-min) (line-beginning-position)))))))
(defun flan--append-output (text) (defun flan--append-output (text)
"Append TEXT, the running program's own output, where it can be read. "Append TEXT, the running program's own output, where it can be read.
Two places. The daemon's log always gets it, so output lands somewhere Two places. The daemon's log always gets it, so output lands somewhere
@ -409,6 +424,7 @@ open, above its prompt, which is where whoever is typing there is looking."
(save-excursion (save-excursion
(goto-char (point-max)) (goto-char (point-max))
(insert text)) (insert text))
(flan--trim-lines)
;; Follow the tail only for someone who was already at it; a reader ;; Follow the tail only for someone who was already at it; a reader
;; scrolled back is reading something. ;; scrolled back is reading something.
(when at-end (goto-char (point-max))))) (when at-end (goto-char (point-max)))))
@ -1535,8 +1551,8 @@ no longer wrong."
;; The bare keys grep-mode and compilation-mode readers reach for. The ;; The bare keys grep-mode and compilation-mode readers reach for. The
;; minor mode below provides RET, M-g M-n and M-g M-p; these two are the ;; minor mode below provides RET, M-g M-n and M-g M-p; these two are the
;; major-mode half, possible because the buffer is read-only. ;; major-mode half, possible because the buffer is read-only.
(define-key map (kbd "n") #'compilation-next-error) (define-key map (kbd "n") #'next-error-no-select)
(define-key map (kbd "p") #'compilation-previous-error) (define-key map (kbd "p") #'previous-error-no-select)
;; The minor mode's RET and `special-mode-map''s q, bound here as well so ;; The minor mode's RET and `special-mode-map''s q, bound here as well so
;; that they are this mode's own keys and reach Evil's states. ;; that they are this mode's own keys and reach Evil's states.
(define-key map (kbd "RET") #'compile-goto-error) (define-key map (kbd "RET") #'compile-goto-error)

View File

@ -272,5 +272,19 @@ what bounds its cost; a `with-temp-buffer' would be scanned by nothing."
(delete-process proc) (delete-process proc)
(kill-buffer buf))) (kill-buffer buf)))
;; The daemon's buffer keeps `flan-output-maximum-lines' and loses the oldest.
(let ((flan-daemon-buffer " *flan-trim-test*")
(flan-output-maximum-lines 10))
(unwind-protect
(progn
(dotimes (i 30) (flan--append-output (format "line %d\n" i)))
(with-current-buffer flan-daemon-buffer
(test-flan--check "the daemon's buffer is kept to its line limit"
(<= (count-lines (point-min) (point-max)) 10))
(test-flan--check "and keeps the newest lines"
(and (string-match-p "line 29" (buffer-string))
(not (string-match-p "line 0\n" (buffer-string)))))))
(kill-buffer flan-daemon-buffer)))
(provide 'test-flan-watch) (provide 'test-flan-watch)
;;; test-flan-watch.el ends here ;;; test-flan-watch.el ends here

View File

@ -557,20 +557,29 @@ already rely on it — so nothing here is a stand-in for the real thing."
(accept-process-output flan--connection 0.01) (accept-process-output flan--connection 0.01)
(setq seen (with-current-buffer flan-daemon-buffer (setq seen (with-current-buffer flan-daemon-buffer
(string-match-p "HELLO" (buffer-string))))) (string-match-p "HELLO" (buffer-string)))))
(let ((took (- (float-time) started))) ;; Arrival is the claim; the time is printed, not asserted, so a
(test-flan--check "the program's output reaches the daemon's buffer" seen) ;; loaded machine running the suite in parallel cannot fail it.
(test-flan--check (test-flan--check
(format "without a request, in well under a second (%.3fs)" took) (format "the program's output reaches the daemon's buffer without a request (%.3fs)"
(and seen (< took 0.3)))) (- (float-time) started))
;; The program prints every 5ms. Coalesced, a second of that is at seen)
;; most one push per 50ms, not two hundred. ;; The program prints every 5ms. Coalesced, the daemon sends at
;; most one push per 50ms however fast it prints; a slow machine
;; can only make that number smaller, so only the upper bound is
;; asserted, and that pushes keep coming at all.
(setq pushes 0) (setq pushes 0)
(let ((until (+ (float-time) 1.0))) (let ((until (+ (float-time) 1.0)))
(while (< (float-time) until) (while (< (float-time) until)
(accept-process-output flan--connection 0.02))) (accept-process-output flan--connection 0.02)))
(let ((in-window pushes)
(until (+ (float-time) 10)))
(while (and (< pushes (1+ in-window)) (< (float-time) until))
(accept-process-output flan--connection 0.05))
(test-flan--check (test-flan--check
(format "a burst of output is coalesced (%d pushes in 1s)" pushes) (format "a burst of output is coalesced (%d pushes in 1s)" in-window)
(<= 5 pushes 25))) (<= in-window 25))
(test-flan--check "and output keeps being pushed"
(> pushes in-window))))
(advice-remove 'flan--dispatch-push count) (advice-remove 'flan--dispatch-push count)
(flan--start-polling))) (flan--start-polling)))
;; A request made while the program prints still gets its own reply. ;; A request made while the program prints still gets its own reply.
@ -1393,14 +1402,12 @@ already rely on it — so nothing here is a stand-in for the real thing."
(test-flan--check "the table arrives without being asked for" (test-flan--check "the table arrives without being asked for"
(with-current-buffer flan-watch-buffer (with-current-buffer flan-watch-buffer
(string-match-p "^ticks [0-9]+" (buffer-string)))) (string-match-p "^ticks [0-9]+" (buffer-string))))
;; `ticks' moves every frame, so every push is a new table. ;; And again, unasked, rather than once.
(setq pushes 0) (setq pushes 0)
(let ((until (+ (float-time) 1.0))) (let ((until (+ (float-time) 10)))
(while (< (float-time) until) (while (and (< pushes 2) (< (float-time) until))
(accept-process-output flan--connection 0.02))) (accept-process-output flan--connection 0.05)))
(test-flan--check (test-flan--check "and keeps arriving" (>= pushes 2)))
(format "and keeps arriving at the repaint interval (%d in 1s)" pushes)
(<= 2 pushes 8)))
(advice-remove 'flan-watch--absorb count)))) (advice-remove 'flan-watch--absorb count))))
;; Replies are not disturbed by the pushes around them. ;; Replies are not disturbed by the pushes around them.
(test-flan--check "a request while the table is being pushed gets its own reply" (test-flan--check "a request while the table is being pushed gets its own reply"

View File

@ -68,6 +68,10 @@ type t = {
A signal here is a crash, and a crash keeps [dir] on disk: see A signal here is a crash, and a crash keeps [dir] on disk: see
[remove_session_dirs]. *) [remove_session_dirs]. *)
mutable died : Unix.process_status option; mutable died : Unix.process_status option;
(* Bytes of program output dropped from the front of [out] because it grew
past [capacity] before anything read it. Said once, in the output, by
[take]. *)
mutable dropped : int;
} }
(* The program's stdout is a pipe into this process, so that an editor can see (* The program's stdout is a pipe into this process, so that an editor can see
@ -97,6 +101,7 @@ let drain t =
(* Bounded: a program that prints every frame must not grow this (* Bounded: a program that prints every frame must not grow this
process without limit. The newest text is the useful end. *) process without limit. The newest text is the useful end. *)
if Buffer.length t.out > capacity then begin if Buffer.length t.out > capacity then begin
t.dropped <- t.dropped + (Buffer.length t.out - capacity);
let keep = Buffer.sub t.out (Buffer.length t.out - capacity) capacity in let keep = Buffer.sub t.out (Buffer.length t.out - capacity) capacity in
Buffer.clear t.out; Buffer.clear t.out;
Buffer.add_string t.out keep Buffer.add_string t.out keep
@ -112,7 +117,14 @@ let take t =
drain t; drain t;
let s = Buffer.contents t.out in let s = Buffer.contents t.out in
Buffer.clear t.out; Buffer.clear t.out;
s if t.dropped = 0 then s
else begin
let n = t.dropped in
t.dropped <- 0;
Printf.sprintf
"[flan dev: %d bytes of the program's output were dropped; the editor \
was not reading fast enough]\n%s" n s
end
let await ?(ms = 5000) f = let await ?(ms = 5000) f =
let rec go ms = let rec go ms =
@ -4288,7 +4300,17 @@ let reply_of_exn e =
A push never interleaves with a reply. The loop is one thread: while a A push never interleaves with a reply. The loop is one thread: while a
request is handled nothing is pushed, what the program prints meanwhile request is handled nothing is pushed, what the program prints meanwhile
stays in [t.out], and [with_output] puts it on that request's reply and stays in [t.out], and [with_output] puts it on that request's reply and
empties the buffer, so it is not pushed again afterwards. *) empties the buffer, so it is not pushed again afterwards.
The program never waits on the editor. A push is written without blocking
from [pend], and whatever the socket will not take stays there until it is
writable again; nothing new is framed while [pend] holds anything. So an
editor that stops reading, even partway through a frame, costs at most one
frame here, and what the program prints in the meantime piles up in
[t.out], which [drain] keeps to [capacity] by dropping the oldest bytes.
[take] then says in the output how much was dropped. Only a reply is
written blocking, after [pend] has been flushed ahead of it: the editor
asked for it and is reading. *)
let push_settle = 0.01 let push_settle = 0.01
let push_gap = 0.05 let push_gap = 0.05
let watch_interval_default = 0.2 let watch_interval_default = 0.2
@ -4299,6 +4321,8 @@ type pushing = {
mutable last_out : float; mutable last_out : float;
mutable watch_every : float option; (* armed by this connection *) mutable watch_every : float option; (* armed by this connection *)
mutable watch_due : float; mutable watch_due : float;
pend : Buffer.t; (* framed bytes not yet written *)
mutable sent : int; (* how many of [pend] have been *)
} }
let bool_field req key = let bool_field req key =
@ -4322,52 +4346,64 @@ let note_output p t now =
if p.on && p.out_due = None && Buffer.length t.out > 0 then if p.on && p.out_due = None && Buffer.length t.out > 0 then
p.out_due <- Some (Float.max (now +. push_settle) (p.last_out +. push_gap)) p.out_due <- Some (Float.max (now +. push_settle) (p.last_out +. push_gap))
(* Whether the client is taking what it is sent. An editor that has stopped let pending p = Buffer.length p.pend > p.sent
reading — busy, or suspended — must not stop this loop, which is also what
drains the program's pipe: while the socket will not take more, output
stays in [t.out] (bounded by [capacity]) and the table waits a turn. *)
let writable fd =
match Unix.select [] [ fd ] [] 0. with
| _, [], _ -> false
| _ -> true
| exception Unix.Unix_error (Unix.EINTR, _, _) -> false
(* Something is due to be pushed. *) let frame p payload =
let push_is_due p now = Buffer.add_string p.pend (string_of_int (String.length payload));
p.on Buffer.add_char p.pend '\n';
&& ((match p.out_due with Some d -> now >= d | None -> false) Buffer.add_string p.pend payload
|| (p.watch_every <> None && now >= p.watch_due))
(* Send whatever is due, and say whether something due could not be sent (* Write as much of [pend] as the socket takes now. With [~block] the rest is
because the client is not reading. Raises [Unix_error] when the client has waited for, which only a reply about to follow it asks. Raises [Unix_error]
gone. *) when the client has gone. *)
let rec push_due p t fd now = let flush_pend ?(block = false) p fd =
if not (push_is_due p now) then false if pending p then begin
else if not (writable fd) then true let data = Buffer.contents p.pend in
else begin push_now p t fd now; false end let n = String.length data in
if not block then Unix.set_nonblock fd;
Fun.protect ~finally:(fun () -> if not block then Unix.clear_nonblock fd)
(fun () ->
let rec go () =
if p.sent < n then
match Unix.single_write_substring fd data p.sent (n - p.sent) with
| k -> p.sent <- p.sent + k; go ()
| exception Unix.Unix_error ((Unix.EAGAIN | Unix.EWOULDBLOCK), _, _)
-> ()
in
go ());
if p.sent >= n then begin Buffer.clear p.pend; p.sent <- 0 end
end
and push_now p t fd now = (* Frame whatever is due, if nothing is still waiting to be written, and write
what the socket will take. Answers whether bytes are left waiting for the
socket to become writable. *)
let push_due p t fd now =
if p.on && not (pending p) then begin
(match p.out_due with (match p.out_due with
| Some due when p.on && now >= due -> | Some due when now >= due ->
p.out_due <- None; p.out_due <- None;
p.last_out <- now; p.last_out <- now;
(match take t with (match take t with
| "" -> () | "" -> ()
| text -> | text ->
Wire.send fd (as_push "output" (ok [ ":output " ^ Wire.quote text ]))) frame p (as_push "output" (ok [ ":output " ^ Wire.quote text ])))
| _ -> ()); | _ -> ());
match p.watch_every with match p.watch_every with
| Some every when p.on && now >= p.watch_due -> | Some every when now >= p.watch_due ->
p.watch_due <- now +. every; p.watch_due <- now +. every;
(* Each push closes the numeric slots' accumulation window while the (* Each push closes the numeric slots' accumulation window while the
program runs. A stopped program takes no samples, so its window is left program runs. A stopped program takes no samples, so its window is
open and the numbers from the moment it stopped stay on the screen. *) left open and the numbers from the moment it stopped stay on the
screen. *)
let stopped = match state t with Stopped _ -> true | _ -> false in let stopped = match state t with Stopped _ -> true | _ -> false in
(* Sent whether or not it changed: ghost text is repainted from each push, (* Sent whether or not it changed: ghost text is repainted from each
and a buffer scrolled into view while the program is stopped gets its push, and a buffer scrolled into view while the program is stopped
values from the next one. *) gets its values from the next one. *)
Wire.send fd (as_push "watch" (watch_read t ~reset:(not stopped))) frame p (as_push "watch" (watch_read t ~reset:(not stopped)))
| _ -> () | _ -> ()
end;
flush_pend p fd;
pending p
(* How long the loop may block before a push is due; [-1.] is no limit, which (* How long the loop may block before a push is due; [-1.] is no limit, which
is what [Unix.select] takes a negative timeout to mean. *) is what [Unix.select] takes a negative timeout to mean. *)
@ -4409,7 +4445,7 @@ let push_request p req op reply =
connection nobody is going to make. *) connection nobody is going to make. *)
let serve t fd = let serve t fd =
let p = { on = false; out_due = None; last_out = 0.; watch_every = None; let p = { on = false; out_due = None; last_out = 0.; watch_every = None;
watch_due = 0. } in watch_due = 0.; pend = Buffer.create 4096; sent = 0 } in
(* Between requests: push what is due, then wait for a request, for the (* Between requests: push what is due, then wait for a request, for the
program's output, or for the next push, whichever comes first. The pipe program's output, or for the next push, whichever comes first. The pipe
is drained here whether or not this client takes pushes, which is the is drained here whether or not this client takes pushes, which is the
@ -4420,8 +4456,9 @@ let serve t fd =
| exception Unix.Unix_error _ -> false | exception Unix.Unix_error _ -> false
| blocked -> | blocked ->
let fds = if t.finished then [ fd ] else [ fd; t.stdout ] in let fds = if t.finished then [ fd ] else [ fd; t.stdout ] in
(* A push the client is not taking waits for it to become writable (* Bytes the client has not taken wait for it to become writable; no
rather than spinning on a deadline that has already passed. *) deadline applies meanwhile, since nothing new is framed until they
have gone. *)
let wfds, wait = let wfds, wait =
if blocked then [ fd ], -1. else [], push_wait p (Unix.gettimeofday ()) if blocked then [ fd ], -1. else [], push_wait p (Unix.gettimeofday ())
in in
@ -4485,7 +4522,7 @@ let serve t fd =
here went straight past the two below and out of the accept loop, here went straight past the two below and out of the accept loop,
ending the session. An editor that left before its reply arrived is a ending the session. An editor that left before its reply arrived is a
closed connection and nothing more, which is what [false] says. *) closed connection and nothing more, which is what [false] says. *)
(match Wire.send fd annotated with (match flush_pend ~block:true p fd; Wire.send fd annotated with
| () -> if op = Some "close" then true else go () | () -> if op = Some "close" then true else go ()
| exception Unix.Unix_error _ -> false) | exception Unix.Unix_error _ -> false)
| exception Wire.Closed -> false | exception Wire.Closed -> false
@ -4815,7 +4852,7 @@ let two_process ?(debug = false) ?(x86 = true) ~file ~sock () =
{ session; child = Some child; agent; dir; stdout = rd; { session; child = Some child; agent; dir; stdout = rd;
out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32; out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32;
host_ll; host_exe = exe; finished = false; agent_watch = None; host_ll; host_exe = exe; finished = false; agent_watch = None;
park_noted = false; died = None } park_noted = false; died = None; dropped = 0 }
in in
ignore_sigpipe (); ignore_sigpipe ();
(try Unix.unlink sock with Unix.Unix_error _ -> ()); (try Unix.unlink sock with Unix.Unix_error _ -> ());
@ -5657,7 +5694,7 @@ let merged_setup () =
{ session; child = None; agent; dir; stdout = rd; { session; child = None; agent; dir; stdout = rd;
out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32; out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32;
host_ll; host_exe = exe; finished = false; agent_watch = None; host_ll; host_exe = exe; finished = false; agent_watch = None;
park_noted = false; died = None } park_noted = false; died = None; dropped = 0 }
in in
ignore_sigpipe (); ignore_sigpipe ();
(try Unix.unlink sock with Unix.Unix_error _ -> ()); (try Unix.unlink sock with Unix.Unix_error _ -> ());

View File

@ -3925,6 +3925,65 @@ let () =
fail fail
"the evaluation drained the program's pipe and kept none of it, so \ "the evaluation drained the program's pipe and kept none of it, so \
the output it caused is gone"); the output it caused is gone");
(* An editor that takes pushes and then stops reading partway through
one must not stop the program. The daemon writes pushes without
blocking, so the program keeps printing into a bounded buffer whose
oldest bytes are dropped, and a line in the output says so. Progress
is the [frames] counter, read before and after a stall long enough
that a blocked program would have filled every buffer between it and
this socket many times over. *)
let rec reply () =
let f = Wire.parse (Wire.recv c) in
if Wire.field f "push" = None then f else reply ()
in
let frames_now () =
Wire.send c
"(:op \"eval-expr\" :code \"frames\" :file \"/tmp/chatty.flan\")";
match Wire.string_field (reply ()) "value" with
| Some v -> Option.value ~default:(-1) (int_of_string_opt v)
| None -> -1
in
Wire.send c "(:op \"push\" :on t)";
ignore (reply ());
let f0 = frames_now () in
(* Into the next push's body and no further. *)
let one = Bytes.create 1 in
let rec header acc =
match Unix.read c one 0 1 with
| 1 when Bytes.get one 0 = '\n' -> int_of_string acc
| 1 -> header (acc ^ Bytes.to_string one)
| _ -> raise Wire.Closed
in
let n = header "" in
let part = Wire.read_exactly c (min 16 n) in
ignore part;
ignore (Unix.select [] [] [] 3.0);
ignore (Wire.read_exactly c (n - min 16 n));
(* Everything after the stall, until the reply to this count, is read
for the note about what was dropped. *)
Wire.send c
"(:op \"eval-expr\" :code \"frames\" :file \"/tmp/chatty.flan\")";
let noted = ref false in
let rec count () =
let f = Wire.parse (Wire.recv c) in
(match Wire.string_field f "output" with
| Some o when contains_sub o "were dropped" -> noted := true
| _ -> ());
if Wire.field f "push" = None then f else count ()
in
let f1 =
match Wire.string_field (count ()) "value" with
| Some v -> Option.value ~default:(-1) (int_of_string_opt v)
| None -> -1
in
if f0 < 0 || f1 < 0 then fail "the frame counter could not be read"
else if f1 - f0 < 400 then
fail "an editor that stopped reading mid-push held the program to %d \
frames in 3s" (f1 - f0);
if not !noted then
fail "output dropped while the editor was not reading was not reported";
Wire.send c "(:op \"push\" :on nil)";
ignore (reply ());
ignore (request c "(:op \"close\")"); ignore (request c "(:op \"close\")");
(try Unix.close c with Unix.Unix_error _ -> ()); (try Unix.close c with Unix.Unix_error _ -> ());
if not if not