diff --git a/docs/BUILT.md b/docs/BUILT.md index fb308bf0..8d646bd9 100644 --- a/docs/BUILT.md +++ b/docs/BUILT.md @@ -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 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: -`Wire.recv` blocked, so an attached, quiet editor let the pipe fill until its next request. A push goes out only when -the socket is writable; an editor that has stopped reading holds pushes back rather than stalling the loop, and its -output waits in `t.out`. +`Wire.recv` blocked, so an attached, quiet editor let the pipe fill until its next request. **The program never waits on +the editor.** A push is written without blocking from a per-connection buffer, and nothing new is framed while that +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 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 diff --git a/emacs/MANUAL.md b/emacs/MANUAL.md index 922ec2c0..471d5dc3 100644 --- a/emacs/MANUAL.md +++ b/emacs/MANUAL.md @@ -1206,6 +1206,7 @@ in the buffer). | `flan-names-shown` | `4` | how many names to list before summarising | | `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-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-start-timeout` | `60` | seconds to wait for a program to come up | | `flan-lower-buffer` | `"*flan-lowering*"` | where `C-c C-l` writes | diff --git a/emacs/flan-repl.el b/emacs/flan-repl.el index f3b3c6e3..0f6e8721 100644 --- a/emacs/flan-repl.el +++ b/emacs/flan-repl.el @@ -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 ;; reply carries would then land above the output it caused. (when (>= at (marker-position mark)) - (set-marker mark (point))))))))) + (set-marker mark (point))))) + (flan--trim-lines))))) (defun flan-repl--buffer () "The REPL buffer, for the two clear commands, from wherever they are run." diff --git a/emacs/flan-watch.el b/emacs/flan-watch.el index 66db28eb..d0701a0a 100644 --- a/emacs/flan-watch.el +++ b/emacs/flan-watch.el @@ -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")))))) ;; 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 ;; a reply in flight for another request to read. ;; diff --git a/emacs/flan.el b/emacs/flan.el index fe109e33..0f60efd1 100644 --- a/emacs/flan.el +++ b/emacs/flan.el @@ -397,6 +397,21 @@ Pushes that arrive first are handled as they come." (bound-and-true-p 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) "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 @@ -409,6 +424,7 @@ open, above its prompt, which is where whoever is typing there is looking." (save-excursion (goto-char (point-max)) (insert text)) + (flan--trim-lines) ;; Follow the tail only for someone who was already at it; a reader ;; scrolled back is reading something. (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 ;; 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. - (define-key map (kbd "n") #'compilation-next-error) - (define-key map (kbd "p") #'compilation-previous-error) + (define-key map (kbd "n") #'next-error-no-select) + (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 ;; that they are this mode's own keys and reach Evil's states. (define-key map (kbd "RET") #'compile-goto-error) diff --git a/emacs/test-flan-watch.el b/emacs/test-flan-watch.el index 8c43156f..ef4db654 100644 --- a/emacs/test-flan-watch.el +++ b/emacs/test-flan-watch.el @@ -272,5 +272,19 @@ what bounds its cost; a `with-temp-buffer' would be scanned by nothing." (delete-process proc) (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) ;;; test-flan-watch.el ends here diff --git a/emacs/test-flan.el b/emacs/test-flan.el index ab876b60..762551ed 100644 --- a/emacs/test-flan.el +++ b/emacs/test-flan.el @@ -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) (setq seen (with-current-buffer flan-daemon-buffer (string-match-p "HELLO" (buffer-string))))) - (let ((took (- (float-time) started))) - (test-flan--check "the program's output reaches the daemon's buffer" seen) - (test-flan--check - (format "without a request, in well under a second (%.3fs)" took) - (and seen (< took 0.3)))) - ;; The program prints every 5ms. Coalesced, a second of that is at - ;; most one push per 50ms, not two hundred. + ;; Arrival is the claim; the time is printed, not asserted, so a + ;; loaded machine running the suite in parallel cannot fail it. + (test-flan--check + (format "the program's output reaches the daemon's buffer without a request (%.3fs)" + (- (float-time) started)) + seen) + ;; 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) (let ((until (+ (float-time) 1.0))) (while (< (float-time) until) (accept-process-output flan--connection 0.02))) - (test-flan--check - (format "a burst of output is coalesced (%d pushes in 1s)" pushes) - (<= 5 pushes 25))) + (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 + (format "a burst of output is coalesced (%d pushes in 1s)" in-window) + (<= in-window 25)) + (test-flan--check "and output keeps being pushed" + (> pushes in-window)))) (advice-remove 'flan--dispatch-push count) (flan--start-polling))) ;; 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" (with-current-buffer flan-watch-buffer (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) - (let ((until (+ (float-time) 1.0))) - (while (< (float-time) until) - (accept-process-output flan--connection 0.02))) - (test-flan--check - (format "and keeps arriving at the repaint interval (%d in 1s)" pushes) - (<= 2 pushes 8))) + (let ((until (+ (float-time) 10))) + (while (and (< pushes 2) (< (float-time) until)) + (accept-process-output flan--connection 0.05))) + (test-flan--check "and keeps arriving" (>= pushes 2))) (advice-remove 'flan-watch--absorb count)))) ;; Replies are not disturbed by the pushes around them. (test-flan--check "a request while the table is being pushed gets its own reply" diff --git a/lib/dev.ml b/lib/dev.ml index 387a01cf..dd6c6fbe 100644 --- a/lib/dev.ml +++ b/lib/dev.ml @@ -68,6 +68,10 @@ type t = { A signal here is a crash, and a crash keeps [dir] on disk: see [remove_session_dirs]. *) 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 @@ -97,6 +101,7 @@ let drain t = (* Bounded: a program that prints every frame must not grow this process without limit. The newest text is the useful end. *) 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 Buffer.clear t.out; Buffer.add_string t.out keep @@ -112,7 +117,14 @@ let take t = drain t; let s = Buffer.contents t.out in 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 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 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 - 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_gap = 0.05 let watch_interval_default = 0.2 @@ -4299,6 +4321,8 @@ type pushing = { mutable last_out : float; mutable watch_every : float option; (* armed by this connection *) 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 = @@ -4322,52 +4346,64 @@ let note_output p t now = 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)) -(* Whether the client is taking what it is sent. An editor that has stopped - 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 +let pending p = Buffer.length p.pend > p.sent -(* Something is due to be pushed. *) -let push_is_due p now = - p.on - && ((match p.out_due with Some d -> now >= d | None -> false) - || (p.watch_every <> None && now >= p.watch_due)) +let frame p payload = + Buffer.add_string p.pend (string_of_int (String.length payload)); + Buffer.add_char p.pend '\n'; + Buffer.add_string p.pend payload -(* Send whatever is due, and say whether something due could not be sent - because the client is not reading. Raises [Unix_error] when the client has - gone. *) -let rec push_due p t fd now = - if not (push_is_due p now) then false - else if not (writable fd) then true - else begin push_now p t fd now; false end +(* Write as much of [pend] as the socket takes now. With [~block] the rest is + waited for, which only a reply about to follow it asks. Raises [Unix_error] + when the client has gone. *) +let flush_pend ?(block = false) p fd = + if pending p then begin + let data = Buffer.contents p.pend in + 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 = - (match p.out_due with - | Some due when p.on && now >= due -> - p.out_due <- None; - p.last_out <- now; - (match take t with - | "" -> () - | text -> - Wire.send fd (as_push "output" (ok [ ":output " ^ Wire.quote text ]))) - | _ -> ()); - match p.watch_every with - | Some every when p.on && now >= p.watch_due -> - p.watch_due <- now +. every; - (* Each push closes the numeric slots' accumulation window while the - program runs. A stopped program takes no samples, so its window is left - open and the numbers from the moment it stopped stay on the screen. *) - let stopped = match state t with Stopped _ -> true | _ -> false in - (* Sent whether or not it changed: ghost text is repainted from each push, - and a buffer scrolled into view while the program is stopped gets its - values from the next one. *) - Wire.send fd (as_push "watch" (watch_read t ~reset:(not stopped))) - | _ -> () +(* 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 + | Some due when now >= due -> + p.out_due <- None; + p.last_out <- now; + (match take t with + | "" -> () + | text -> + frame p (as_push "output" (ok [ ":output " ^ Wire.quote text ]))) + | _ -> ()); + match p.watch_every with + | Some every when now >= p.watch_due -> + p.watch_due <- now +. every; + (* Each push closes the numeric slots' accumulation window while the + program runs. A stopped program takes no samples, so its window is + left open and the numbers from the moment it stopped stay on the + screen. *) + let stopped = match state t with Stopped _ -> true | _ -> false in + (* Sent whether or not it changed: ghost text is repainted from each + push, and a buffer scrolled into view while the program is stopped + gets its values from the next one. *) + 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 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. *) let serve t fd = 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 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 @@ -4420,8 +4456,9 @@ let serve t fd = | exception Unix.Unix_error _ -> false | blocked -> 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 - rather than spinning on a deadline that has already passed. *) + (* Bytes the client has not taken wait for it to become writable; no + deadline applies meanwhile, since nothing new is framed until they + have gone. *) let wfds, wait = if blocked then [ fd ], -1. else [], push_wait p (Unix.gettimeofday ()) in @@ -4485,7 +4522,7 @@ let serve t fd = 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 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 () | exception Unix.Unix_error _ -> 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; out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32; host_ll; host_exe = exe; finished = false; agent_watch = None; - park_noted = false; died = None } + park_noted = false; died = None; dropped = 0 } in ignore_sigpipe (); (try Unix.unlink sock with Unix.Unix_error _ -> ()); @@ -5657,7 +5694,7 @@ let merged_setup () = { session; child = None; agent; dir; stdout = rd; out = Buffer.create 4096; n = 0; gen = 0; owners = Hashtbl.create 32; host_ll; host_exe = exe; finished = false; agent_watch = None; - park_noted = false; died = None } + park_noted = false; died = None; dropped = 0 } in ignore_sigpipe (); (try Unix.unlink sock with Unix.Unix_error _ -> ()); diff --git a/test/test_dev.ml b/test/test_dev.ml index eadfe160..e2457655 100644 --- a/test/test_dev.ml +++ b/test/test_dev.ml @@ -3925,6 +3925,65 @@ let () = fail "the evaluation drained the program's pipe and kept none of it, so \ 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\")"); (try Unix.close c with Unix.Unix_error _ -> ()); if not