Add interrupt/abort and incremental output streaming
Wire the 0.17.9 regvm preemptive-yield and custom output-port primitives into (sigil nrepl) so live-coding is stream-safe and responsive.
Sliced eval tasks: expression-form evals run cooperatively, one preemptive-yield slice per nrepl-process-pending, so a CPU-bound eval (e.g. (let loop () (loop))) no longer freezes the host loop. An (:op abort :target-id <eval-id>) from a second connection drops the suspended continuation; the interrupted eval receives :status error :code "interrupted" and the session survives.
Output streaming: around each slice, current-output-port/current-error-port are redirected to custom callback ports that buffer writes; buffered chunks are flushed back as incremental (response :status out :out <chunk>) / :err frames carrying the original eval request id, so a long eval's prints arrive as they happen rather than batched at completion.
Definition forms (define/define-syntax/import/begin-with-defs/...) eval immediately at top level (preserving module-level define semantics) and are not sliced; yield is armed only around the sliced expression path, so there is zero preemption overhead off that path. This hybrid works around a VM limitation: preemptive yield does not propagate across the eval native boundary (nested sigil_vmexecute), so sliceable expressions are run as ((eval (list 'lambda '() form))) — a bytecode call in the prompt-bearing invocation where yields propagate.
Tests: abort interrupts a runaway loop + session survives; output streams incrementally (frames precede the final response and arrive while the eval is still suspended); malformed non-string :target-id does not crash the server. Existing suite stays green (17 total).
src/sigil/nrepl.sgl | 323 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
test/test-server.sgl | 177 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
2 files changed, 486 insertions(+), 14 deletions(-)src/sigil/nrepl.sglmodified
(sigil inspect) (sigil error) (sigil diagnostic) (sigil async) ; Provides: async-prompt-tag (for preemptive-yield slicing) (sigil repl) ; Provides: make-prompt, make-debug-prompt, format-srcloc, ; format-proc-name, format-backtrace-str, format-frame-str, ; debug-help-text, parse-debug-command, string->int, ;; DATA STRUCTURES ;; ============================================================ ;; Server state: (vector 'nrepl-server socket port sessions) ;; Server state: (vector 'nrepl-server socket port sessions clients tasks) ;; sessions is an alist of (session-id . session-state) ;; tasks is a list of pending sliced eval tasks (see below) (define (make-nrepl-server socket port) (vector 'nrepl-server socket port '() '())) (vector 'nrepl-server socket port '() '() '())) (define (nrepl-server? obj) (and (vector? obj) (define (nrepl-clients-set! server clients) (vector-set! server 4 clients)) (define (nrepl-tasks server) (vector-ref server 5)) (define (nrepl-tasks-set! server tasks) (vector-set! server 5 tasks)) ;; ============================================================ ;; SLICED EVAL TASKS ;; ============================================================ ;; ;; A CPU-bound eval (e.g. a runaway loop) must not freeze the host's ;; cooperative loop. Expression-form evals run as *sliced tasks*: each ;; call to nrepl-process-pending runs one preemptive-yield slice of the ;; eval, so an `abort` op from a second connection can interrupt it and ;; the host loop keeps turning. Output produced during the eval is ;; buffered per slice and flushed back as incremental :out/:err frames. ;; ;; Task: (vector 'nrepl-eval-task client id thunk k response ;; done? aborted? outbuf out-port err-port) ;; thunk - runs the (sliced) eval, returns the final response ;; k - captured continuation while the eval is suspended, else #f ;; response - the final response to send when done ;; outbuf - reverse-order list of (kind . string) chunks awaiting flush ;; out-port / err-port - custom output ports that buffer into outbuf (define (make-eval-task client id thunk) (let ((task (vector 'nrepl-eval-task client id thunk #f #f #f #f '() #f #f))) (vector-set! task 9 (make-custom-output-port (lambda (chunk) (task-buffer-add! task 'out chunk)))) (vector-set! task 10 (make-custom-output-port (lambda (chunk) (task-buffer-add! task 'err chunk)))) task)) (define (eval-task-client task) (vector-ref task 1)) (define (eval-task-id task) (vector-ref task 2)) (define (eval-task-thunk task) (vector-ref task 3)) (define (eval-task-k task) (vector-ref task 4)) (define (eval-task-k-set! task k) (vector-set! task 4 k)) (define (eval-task-response task) (vector-ref task 5)) (define (eval-task-response-set! task response) (vector-set! task 5 response)) (define (eval-task-done? task) (vector-ref task 6)) (define (eval-task-done-set! task done?) (vector-set! task 6 done?)) (define (eval-task-aborted? task) (vector-ref task 7)) (define (eval-task-aborted-set! task aborted?) (vector-set! task 7 aborted?)) (define (eval-task-outbuf task) (vector-ref task 8)) (define (eval-task-outbuf-set! task buf) (vector-set! task 8 buf)) (define (eval-task-out-port task) (vector-ref task 9)) (define (eval-task-err-port task) (vector-ref task 10)) (define (task-buffer-add! task kind chunk) (eval-task-outbuf-set! task (cons (cons kind chunk) (eval-task-outbuf task)))) ;; Client connection: (vector 'nrepl-client socket buffer session-id module debug-mode? debug-error debug-trace debug-frame debug-on-error?) ;; debug-mode?: #t when client is in debug mode after an error ;; debug-error: the error message string when in debug mode (let ((mod (find-module (read-expr module-name)))) (when mod (client-module-set! client mod)))) ;; Dispatch on operation ;; Dispatch on operation. ;; NOTE: the server path (nrepl-handle-request*) intercepts eval ;; and abort so they run through the sliced-task machinery; the ;; cases here are for standalone/direct callers (no server state). (case op ((eval) (handle-eval client id args)) ((abort interrupt) (handle-abort client id args)) ((eval) (handle-eval client id args #f)) ((abort interrupt) (handle-abort-standalone client id args)) ((complete) (handle-complete client id args)) ((doc) (handle-doc client id args)) ((describe) (handle-describe client id args)) ;; Catches both Scheme exceptions (via guard) and VM errors (via vm-error-prompt-tag) ;; Default wire sessions return structured errors and remain in normal eval mode. ;; Debug mode is entered only when the session's debug-on-error policy is enabled. (define (handle-eval client id args) ;; ;; sliced?: when #t the user form is run as a preemptive-yield-abortable ;; call (see evaluate-user-code). The server passes #t for sliceable ;; expression forms; definition forms and standalone callers pass #f. (define (handle-eval client id args sliced?) (let ((code (get-arg args ':code "")) (fmt (get-arg args ':format #f))) (if (string=? code "") (make-response* client id 'exit ':value "Goodbye!") (make-response* client id 'ok ':value "ok"))) ;; Regular evaluation (let ((value (eval (read-expr code)))) (let ((value (evaluate-user-code code sliced?))) (set-current-module! saved-module) (make-response* client id 'ok ':value (if structured? (format "~s" value)))))))))) result)))))) ;; Evaluate the user's code form and return its value. ;; ;; sliced? #f: evaluate the form directly (top-level semantics preserved — ;; `define` lands at module level). Used for definition forms and ;; standalone callers. ;; sliced? #t: build (lambda () <form>) and CALL it, so the yield-bearing ;; work runs as a bytecode OP_CALL in the prompt-bearing invocation and ;; is preemptible/abortable. Only valid for expression forms — a top-level ;; `define` here would become a local binding. The form is classified as ;; sliceable before we get here (see sliceable-form?). ;; ;; Why the lambda dance: the `eval` native runs code in a nested VM ;; invocation (sigil__vm_execute), and a preemptive yield tripped inside it ;; cannot abort to the async prompt installed in the OUTER invocation. By ;; eval'ing only `(lambda () form)` (which completes immediately, producing ;; a closure) and then calling that closure ourselves, the user code runs ;; in the current, prompt-bearing invocation where yields propagate. (define (evaluate-user-code code sliced?) (let ((form (read-expr code))) (if sliced? ((eval (list 'lambda '() form))) (eval form)))) (define (handle-debug-policy client id args) (let ((enabled? (get-arg args ':debug-on-error #f))) (client-debug-on-error-set! client (and enabled? #t)) (clear-debug-mode! client) (make-response* client id 'ok ':in-debugger #f ':value "Returning to REPL")) (define (handle-abort client id args) ;; abort/interrupt for a standalone (no-server) caller: there is no task ;; registry, so validate the request shape and report not-found. (define (handle-abort-standalone client id args) (let ((target-id (get-arg args ':target-id #f))) (if (not target-id) (make-error-response* client id "missing-target-id" "No :target-id provided") (make-error-response* client id "not-found" (format "No running eval for request id ~a" target-id))))) ;; ============================================================ ;; SLICED EVAL: classification, dispatch, execution ;; ============================================================ ;; Heads of top-level definition forms. These must be evaluated with ;; top-level semantics (so `define` lands at module level) and complete ;; immediately, so they are NOT sliced. (define definition-form-heads '(define define-values define-syntax define-record-type define-library define-constant define-structure define-parameter import include include-ci load-module use require)) (define (definition-form? form) (and (pair? form) (or (memq (car form) definition-form-heads) ;; A begin containing any definition must run at top level. (and (eq? (car form) 'begin) (any-definition? (cdr form)))))) (define (any-definition? forms) (cond ((null? forms) #f) ((definition-form? (car forms)) #t) (else (any-definition? (cdr forms))))) ;; A form is sliceable (run preemptibly as a lambda call) when it is an ;; expression, not a top-level definition. Non-pairs (literals, symbols) ;; are trivially sliceable but complete in one slice. (define (sliceable-form? form) (not (definition-form? form))) ;; Server-aware eval dispatch. Sliceable expression forms become pending ;; tasks (response delivered later, after slicing); definition forms and ;; anything we cannot classify run immediately through handle-eval. (define (handle-eval-request server client id args) (let ((code (get-arg args ':code ""))) (if (or (string=? code "") (client-debug-mode? client)) ;; Empty code / debug-command handling stays on the immediate path. (handle-eval client id args #f) (let ((form (guard (exn (else 'unreadable)) (read-expr code)))) (if (and (not (eq? form 'unreadable)) (sliceable-form? form)) (begin (enqueue-eval-task! server client id args) #f) ;; Definition form, or unreadable (let handle-eval surface ;; the read/eval error through its structured-error path). (handle-eval client id args #f)))))) (define (enqueue-eval-task! server client id args) (let ((task (make-eval-task client id (lambda () (handle-eval client id args #t))))) (nrepl-tasks-set! server (append (nrepl-tasks server) (list task))))) ;; equal? (not string=?) so a client-supplied :target-id of an unexpected ;; type (number, symbol) simply fails to match instead of raising a ;; type-error that would unwind through the whole server loop. (define (find-task-by-id server target-id) (let loop ((tasks (nrepl-tasks server))) (cond ((null? tasks) #f) ((equal? (eval-task-id (car tasks)) target-id) (car tasks)) (else (loop (cdr tasks)))))) ;; The interrupted-eval response delivered to the ORIGINAL eval client ;; when its running eval is aborted. Shape mirrors a structured error so ;; existing error-handling clients cope: :status error :code "interrupted". (define (make-interrupted-response client id) (make-response* client id 'error ':code "interrupted" ':message "Evaluation interrupted" ':error (list ':type "interrupted" ':message "Evaluation interrupted" ':repr "interrupted") ':stack "" ':in-debugger #f)) (define (handle-abort server client id args) (let ((target-id (get-arg args ':target-id #f))) (if (not target-id) (make-error-response* client id "missing-target-id" "No :target-id provided") (let ((task (find-task-by-id server target-id))) (if task (begin ;; Deliver any output produced before the interruption, ;; then drop the suspended continuation and mark the task ;; done with the interrupted response. (flush-eval-task-output! task) (eval-task-aborted-set! task #t) (eval-task-k-set! task #f) (eval-task-response-set! task (make-interrupted-response (eval-task-client task) (eval-task-id task))) (eval-task-done-set! task #t) (make-response* client id 'ok ':aborted target-id)) (make-error-response* client id "not-found" (format "No running eval for request id ~a" target-id))))))) ;; Run one preemptive-yield slice of a pending eval task. On first entry ;; the task thunk starts; on later entries the captured continuation ;; resumes. When the eval yields, the continuation is captured and the ;; task stays pending; when it returns, its value becomes the response. (define (run-eval-task-slice! task) (when (and (not (eval-task-done? task)) (not (eval-task-aborted? task))) (let ((saved-module (current-module)) (saved-out (current-output-port)) (saved-err (current-error-port))) ;; Redirect output to the task's buffering ports and enter the ;; client's module for the duration of this slice. (set-current-module! (client-module (eval-task-client task))) (current-output-port (eval-task-out-port task)) (current-error-port (eval-task-err-port task)) (set-vm-async-prompt-tag! async-prompt-tag) (set-vm-yield-requested! #t) (let ((result (call-with-prompt async-prompt-tag ;; Yield handler: capture the continuation, stay pending. (lambda (k op-type arg1 arg2) (set-vm-yield-requested! #f) (eval-task-k-set! task k) 'yielded) ;; Body: resume the suspended eval, or start it. (lambda () (let ((k (eval-task-k task))) (if k (begin (eval-task-k-set! task #f) (k #f)) ((eval-task-thunk task)))))))) ;; Restore VM/dynamic state before returning to the server loop. (set-vm-yield-requested! #f) (set-current-module! saved-module) (current-output-port saved-out) (current-error-port saved-err) (unless (eq? result 'yielded) (eval-task-response-set! task result) (eval-task-done-set! task #t)))))) ;; Flush buffered eval output as incremental :out/:err frames to the eval ;; client, carrying the original eval request id. Runs in the server loop ;; (yield disarmed, ports restored), so socket writes are safe here. (define (flush-eval-task-output! task) (let ((buf (reverse (eval-task-outbuf task))) (client (eval-task-client task)) (id (eval-task-id task))) (eval-task-outbuf-set! task '()) (for-each (lambda (entry) (let ((kind (car entry)) (chunk (cdr entry))) (socket-write (client-socket client) (encode-message (if (eq? kind 'err) (list 'response ':id id ':status 'err ':err chunk) (list 'response ':id id ':status 'out ':out chunk)))))) buf))) (define (send-eval-task-response! task) (let ((response (eval-task-response task)) (client (eval-task-client task))) (when response (socket-write (client-socket client) (encode-message response))))) ;; Advance every pending eval task by one slice, flushing streamed output ;; and sending final responses for completed tasks. A task whose eval ;; client has disconnected is dropped (not sliced, not written to) so an ;; orphaned runaway eval can't slice forever or write to a closed socket. (define (process-eval-tasks server) (let loop ((tasks (nrepl-tasks server)) (remaining '())) (if (null? tasks) (nrepl-tasks-set! server (reverse remaining)) (let ((task (car tasks))) (cond ((socket-closed? (client-socket (eval-task-client task))) ;; Client gone: abandon the task. (loop (cdr tasks) remaining)) (else (run-eval-task-slice! task) (flush-eval-task-output! task) (if (eval-task-done? task) (begin (send-eval-task-response! task) (loop (cdr tasks) remaining)) (loop (cdr tasks) (cons task remaining))))))))) ;; Server-aware request dispatch: eval and abort route through the sliced ;; task machinery; everything else delegates to nrepl-handle-request. (define (nrepl-handle-request* server client request) (if (not (and (pair? request) (eq? (car request) 'request))) (make-error-response "unknown" "invalid-request" "Expected (request ...)") (let* ((args (cdr request)) (id (get-arg args ':id "unknown")) (op (get-arg args ':op #f)) (module-name (get-arg args ':module #f))) (when (and module-name (string? module-name) (not (string=? module-name ""))) (let ((mod (find-module (read-expr module-name)))) (when mod (client-module-set! client mod)))) (case op ((eval) (handle-eval-request server client id args)) ((abort interrupt) (handle-abort server client id args)) (else (nrepl-handle-request client request)))))) ;; complete - return completions for a prefix ;; Includes both value bindings and syntax/macro bindings (define (handle-complete client id args) ;; Accept new connections (accept-pending-connections server) ;; Process data from existing clients (process-client-data server))) (process-client-data server) ;; Advance each pending sliced eval task by one slice (process-eval-tasks server))) ;; Accept any pending connections (define (accept-pending-connections server) (let ((message (car result)) (remaining (cdr result))) (client-buffer-set! client remaining) ;; Handle the request (let ((response (nrepl-handle-request client message))) ;; Send response (socket-write (client-socket client) (encode-message response))) ;; Handle the request via the server-aware dispatcher. Sliced ;; eval requests return #f here (their response is delivered later ;; by process-eval-tasks); everything else responds immediately. (let ((response (nrepl-handle-request* server client message))) (when response (socket-write (client-socket client) (encode-message response)))) ;; Check for more messages (process-client-messages server client)))))test/test-server.sglmodified
(nrepl-process-pending server) (pump server (- n 1))));; Buffered stream reader. Streamed eval produces MULTIPLE frames that can;; arrive coalesced in a single TCP segment, so a reader must preserve bytes;; beyond the first frame. A conn wraps a socket plus a leftover byte buffer.(define (make-conn sock) (vector sock (make-bytevector 0)))(define (conn-sock c) (vector-ref c 0))(define (conn-buf c) (vector-ref c 1))(define (conn-buf-set! c b) (vector-set! c 1 b));; Decode one frame from a byte buffer; returns (frame . remaining) or #f.(define (conn-try-decode buf) (if (< (bytevector-length buf) 4) #f (let* ((len (+ (* (bytevector-u8-ref buf 0) 16777216) (* (bytevector-u8-ref buf 1) 65536) (* (bytevector-u8-ref buf 2) 256) (bytevector-u8-ref buf 3))) (total (+ 4 len))) (if (< (bytevector-length buf) total) #f (cons (read-expr (utf8->string (bytevector-copy buf 4 total))) (bytevector-copy buf total))))));; Read one frame from a conn, polling the socket up to `attempts` times.(define (conn-read-frame c attempts) (let loop ((att attempts)) (let ((dec (conn-try-decode (conn-buf c)))) (if dec (begin (conn-buf-set! c (cdr dec)) (car dec)) (if (<= att 0) #f (begin (when (socket-ready? (conn-sock c) 20) (let ((d (socket-read-bytevector (conn-sock c)))) (when (and d (not (eof-object? d)) (> (bytevector-length d) 0)) (conn-buf-set! c (bytevector-append (conn-buf c) d))))) (loop (- att 1))))))));; Read framed responses from a conn until one carries a terminal status;; (ok/error/exit) or max frames are read. Returns the frames in wire order.;; Used to observe streamed :out/:err frames that precede the final response.(define (conn-read-until-final c max) (let loop ((i 0) (acc '())) (if (>= i max) (reverse acc) (let ((resp (conn-read-frame c 60))) (if resp (let ((acc2 (cons resp acc))) (if (memq (response-ref resp ':status #f) '(ok error exit)) (reverse acc2) (loop (+ i 1) acc2))) (reverse acc))))))(define (frames-with-status frames status) (cond ((null? frames) '()) ((eq? (response-ref (car frames) ':status #f) status) (cons (car frames) (frames-with-status (cdr frames) status))) (else (frames-with-status (cdr frames) status))))(define (any-chunk-contains? frames key needle) (cond ((null? frames) #f) ((let ((v (response-ref (car frames) key #f))) (and (string? v) (string-contains? v needle))) #t) (else (any-chunk-contains? (cdr frames) key needle))))(define (request-response server sock request) (send-request sock request) (pump server 20) "nrepl-getcell-match"))) (socket-close sock))))))(test-group "nrepl interrupt/abort" (test "a runaway eval is interrupted from a second connection and the session survives" (with-server (lambda (server) (let ((eval-sock (tcp-connect "127.0.0.1" *test-port*))) (pump server 5) ;; Start a CPU-bound eval that never returns on its own. (send-request eval-sock '(request :id "runaway" :op eval :code "(let loop () (loop))")) ;; Drive several slices — the eval keeps yielding, never completes, ;; and (critically) the host loop keeps turning. (pump server 10) (assert-false (socket-ready? eval-sock 20)) ; no response yet ;; Interrupt it from a SECOND connection. (let ((abort-sock (tcp-connect "127.0.0.1" *test-port*))) (pump server 3) (send-request abort-sock '(request :id "do-abort" :op abort :target-id "runaway")) (pump server 10) (let ((abort-resp (read-response abort-sock)) (eval-resp (read-response eval-sock))) (assert-eq (response-ref abort-resp ':status #f) 'ok) (assert-equal (response-ref abort-resp ':aborted #f) "runaway") ;; The interrupted eval receives a terminal interrupted error. (assert-eq (response-ref eval-resp ':status #f) 'error) (assert-equal (response-ref eval-resp ':code #f) "interrupted")) ;; Session survives: the same connection evaluates normally after. (send-request eval-sock '(request :id "after-abort" :op eval :code "(+ 20 22)")) (pump server 20) (let ((after-resp (read-response eval-sock))) (assert-eq (response-ref after-resp ':status #f) 'ok) (assert-equal (response-ref after-resp ':value #f) "42")) (socket-close eval-sock) (socket-close abort-sock)))))) (test "aborting an unknown target reports not-found and leaves running evals alone" (with-server (lambda (server) (let ((sock (tcp-connect "127.0.0.1" *test-port*))) (pump server 5) (let ((resp (request-response server sock '(request :id "abort-nobody" :op abort :target-id "ghost")))) (assert-eq (response-ref resp ':status #f) 'error) (assert-equal (response-ref resp ':code #f) "not-found")) ;; A non-string :target-id must not crash the server loop; it simply ;; matches nothing and reports not-found. (let ((resp (request-response server sock '(request :id "abort-badtype" :op abort :target-id 42)))) (assert-eq (response-ref resp ':status #f) 'error) (assert-equal (response-ref resp ':code #f) "not-found")) ;; Server still healthy afterwards. (let ((resp (request-response server sock '(request :id "ping-after" :op ping)))) (assert-eq (response-ref resp ':status #f) 'ok)) (socket-close sock))))))(test-group "nrepl output streaming" (test "eval output streams as incremental :out frames before the final response" (with-server (lambda (server) (let* ((sock (tcp-connect "127.0.0.1" *test-port*)) (conn (make-conn sock))) (pump server 5) ;; Two displays separated by CPU-bound spins so they land in ;; different slices — proving output is flushed incrementally, not ;; batched at completion. (send-request sock '(request :id "stream" :op eval :code "(begin (display \"chunk-a\") (let loop ((n 0)) (if (< n 80000) (loop (+ n 1)) #t)) (display \"chunk-b\") (let loop ((n 0)) (if (< n 80000) (loop (+ n 1)) #t)) 42)")) (pump server 400) ; drive the sliced eval to completion (let* ((frames (conn-read-until-final conn 32)) (out-frames (frames-with-status frames 'out)) (final (car (reverse frames)))) ;; At least one streamed :out frame arrived... (assert-true (> (length out-frames) 0)) ;; ...carrying the original eval request id... (assert-equal (response-ref (car out-frames) ':id #f) "stream") ;; ...and both displayed chunks were streamed. (assert-true (any-chunk-contains? out-frames ':out "chunk-a")) (assert-true (any-chunk-contains? out-frames ':out "chunk-b")) ;; The terminal frame is the eval result, and it comes LAST — ;; every :out frame precedes it on the wire (not batched). (assert-eq (response-ref final ':status #f) 'ok) (assert-equal (response-ref final ':value #f) "42")) (socket-close sock))))) (test "streamed output frames arrive while the eval is still running" (with-server (lambda (server) (let* ((sock (tcp-connect "127.0.0.1" *test-port*)) (conn (make-conn sock))) (pump server 5) (send-request sock '(request :id "early" :op eval :code "(begin (display \"early-out\") (let loop ((n 0)) (if (< n 300000) (loop (+ n 1)) #t)) 99)")) ;; Only a few slices: enough to emit the first display and yield, ;; but NOT enough to finish the long spin. (pump server 6) (let ((first (conn-read-frame conn 60))) ;; The first frame is a streamed :out, delivered before any final ;; response exists — the eval is demonstrably still suspended. (assert-eq (response-ref first ':status #f) 'out) (assert-equal (response-ref first ':out #f) "early-out")) ;; Now let it finish and collect the terminal result. (pump server 400) (let ((rest (conn-read-until-final conn 32))) (assert-eq (response-ref (car (reverse rest)) ':status #f) 'ok) (assert-equal (response-ref (car (reverse rest)) ':value #f) "99")) (socket-close sock))))))(run-tests)