diff --git a/src/core/bytecode.cc b/src/core/bytecode.cc index 2b4cc0649d..4a50e487dd 100644 --- a/src/core/bytecode.cc +++ b/src/core/bytecode.cc @@ -217,6 +217,22 @@ static unsigned char* long_dispatch(VirtualMachine&, unsigned char*, MultipleVal core::T_O**, core::T_O**, size_t, core::T_O**, uint8_t); SYMBOL_EXPORT_SC_(KeywordPkg, name); + +// Poll interrupts on backward branches so loops of pure opcodes stay cancellable. +static inline unsigned char* vm_branch(VirtualMachine& vm, ThreadLocalState* thread, unsigned char* pc, core::T_O** sp, + int32_t rel) { + pc += rel; + if (rel < 0) { + core::Cons_sp head = thread->_PendingInterruptsHead.load(std::memory_order_acquire); + if (thread->pending_signals_p() || (static_cast(head) && head->cdr().notnilp())) { + // The handler can cons or unwind, so the GC must see the current frame. + vm._pc = pc; + vm._stackPointer = sp; + gctools::handle_all_queued_interrupts(); + } + } + return pc; +} #ifdef DEBUG_VIRTUAL_MACHINE __attribute__((optnone)) #endif @@ -564,19 +580,19 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure case vm_code::jump_8: { int8_t rel = *(pc + 1); DBG_VM1("jump %" PRId8 "\n", rel); - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); break; } case vm_code::jump_16: { int16_t rel = read_s16(pc + 1); DBG_VM("jump %" PRId16 "\n", rel); - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); break; } case vm_code::jump_24: { int32_t rel = read_label(pc, 3); DBG_VM("jump %" PRId32 "\n", rel); - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); break; } case vm_code::jump_if_8: { @@ -585,7 +601,7 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure T_sp tval((gctools::Tagged)vm.pop(sp)); VM_RECORD_PLAYBACK(tval.raw_(), "vm_jump_if_8"); if (tval.notnilp()) - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); else pc += 2; break; @@ -595,7 +611,7 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure DBG_VM("jump-if %" PRId16 "\n", rel); T_sp tval((gctools::Tagged)vm.pop(sp)); if (tval.notnilp()) - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); else pc += 3; break; @@ -605,7 +621,7 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure DBG_VM("jump-if %" PRId32 "\n", rel); T_sp tval((gctools::Tagged)vm.pop(sp)); if (tval.notnilp()) - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); else pc += 4; break; @@ -618,7 +634,7 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure pc += 2; } else { vm.push(sp, tval.raw_()); - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); } break; } @@ -630,7 +646,7 @@ bytecode_vm(VirtualMachine& vm, T_O** literals, T_O** closed, Closure_O* closure pc += 3; } else { vm.push(sp, tval.raw_()); - pc += rel; + pc = vm_branch(vm, thread, pc, sp, rel); } break; } diff --git a/src/core/lispStream.cc b/src/core/lispStream.cc index b35cb3a7c6..9113554630 100644 --- a/src/core/lispStream.cc +++ b/src/core/lispStream.cc @@ -1353,7 +1353,9 @@ CL_DEFUN T_mv core__read_fd(int filedes, SimpleBaseString_sp buffer) { size_t buffer_length = cl__length(buffer); unsigned char* buffer_data = &(*buffer)[0]; while (1) { - int num = read(filedes, buffer_data, buffer_length); + // Park: read() blocks until data arrives. Safe to hold BUFFER_DATA across it + // because no GC variant relocates (NON_MOVING_GC). + int num = BEGIN_PARK { return (int)read(filedes, buffer_data, buffer_length); } END_PARK; if (!(num < 0 && errno == EINTR)) { if (num < 0) { return Values(make_fixnum(num), make_fixnum(errno)); diff --git a/src/core/unixfsys.cc b/src/core/unixfsys.cc index 3b05d316a9..ec0bc00cb6 100644 --- a/src/core/unixfsys.cc +++ b/src/core/unixfsys.cc @@ -37,6 +37,7 @@ THE SOFTWARE. */ #include +#include #include #include @@ -294,7 +295,8 @@ The status can be passed to //core:wifexited// and //core:wifsignaled//. )") DOCGROUP(clasp); CL_DEFUN T_mv core__wait() { int status; - pid_t p = wait(&status); + // Park: wait() blocks until a child exits, which may be never. + pid_t p = BEGIN_PARK { return wait(&status); } END_PARK; return Values(make_fixnum(p), make_fixnum(status)); }; @@ -1707,7 +1709,9 @@ DOCGROUP(clasp); CL_DEFUN T_mv ext__system(String_sp cmd) { ASSERT(cl__stringp(cmd)); string command = cmd->get_std_string(); - int ret = system(command.c_str()); + // Park: system() blocks for an unbounded time, so the thread must be GC-safe + // and marked blocking, which is what lets an interrupt wake it with SIGCONT. + int ret = BEGIN_PARK { return system(command.c_str()); } END_PARK; if (ret == 0) { return Values(core::make_fixnum(0)); } else { @@ -1828,7 +1832,7 @@ CL_DEFUN T_mv ext__vfork_execvp(List_sp call_and_arguments, T_sp return_stream) while (b_done == false) { errno = 0; - wait_ret = wait(&status); + wait_ret = BEGIN_PARK { return wait(&status); } END_PARK; if (WIFEXITED(status)) { child_exit_status = WEXITSTATUS(status); @@ -1946,7 +1950,7 @@ CL_DEFUN T_mv ext__fork_execvp(List_sp call_and_arguments, T_sp return_stream) { } else { // Parent int status; - pid_t wait_ret = wait(&status); + pid_t wait_ret = BEGIN_PARK { return wait(&status); } END_PARK; // Clean up args for (int i(0); i < execvp_args.size() - 1; ++i) free((void*)execvp_args[i]); @@ -2003,7 +2007,10 @@ CL_DEFUN T_mv core__select(int nfds, FdSet_sp readfds, FdSet_sp writefds, FdSet_ struct timeval timeout; timeout.tv_sec = seconds; timeout.tv_usec = microseconds; - int num = select(nfds, &readfds->_fd_set, &writefds->_fd_set, &errorfds->_fd_set, &timeout); + // Park: select() blocks for up to the caller's timeout. + int num = BEGIN_PARK { + return select(nfds, &readfds->_fd_set, &writefds->_fd_set, &errorfds->_fd_set, &timeout); + } END_PARK; if (num < 0) { return Values(make_fixnum(num), make_fixnum(errno)); } diff --git a/src/lisp/kernel/cleavir/translate.lisp b/src/lisp/kernel/cleavir/translate.lisp index bc0f7d658d..a7af3132ff 100644 --- a/src/lisp/kernel/cleavir/translate.lisp +++ b/src/lisp/kernel/cleavir/translate.lisp @@ -255,12 +255,44 @@ function-or-placeholder - the llvm function or a placeholder for (inst-source instruction) 999902)) (call-next-method))) +(defvar *laid-out-iblocks*) +(defvar *safepoint-counter* nil) + +;; Iblocks are laid out in forward flow order, so a branch to one already laid +;; out is a back edge - the native counterpart of the VM's negative jump offset. +(defun back-edge-p (instruction) + (and (boundp '*laid-out-iblocks*) + (loop for succ in (bir:next instruction) + thereis (gethash succ *laid-out-iblocks*)))) + +;; cc_safepoint is an out-of-line unwinding call, so polling every back edge +;; costs ~3x on a tight loop; sample one edge in +safepoint-sample+ instead. +(defparameter +safepoint-sample+ 255) + +(defun emit-back-edge-safepoint () + (let* ((n (cmp:irc-typed-load cmp:%size_t% *safepoint-counter*)) + (n1 (cmp:irc-add n (cmp:jit-constant-size_t 1))) + (poll (cmp:irc-basic-block-create "safepoint-poll")) + (cont (cmp:irc-basic-block-create "safepoint-cont"))) + (cmp:irc-store n1 *safepoint-counter*) + (cmp:irc-cond-br (cmp:irc-icmp-eq + (cmp:irc-and n1 (cmp:jit-constant-size_t +safepoint-sample+)) + (cmp:jit-constant-size_t 0)) + poll cont) + (cmp:irc-begin-block poll) + (%intrinsic-invoke-if-landing-pad-or-call "cc_safepoint" ()) + (cmp:irc-br cont) + (cmp:irc-begin-block cont))) + (defmethod translate-terminator :around ((instruction bir:instruction) abi next) (declare (ignore abi next)) (cmp:with-debug-info-source-position ((ensure-origin (inst-source instruction) 999903)) + ;; Poll interrupts on back edges so native loops stay cancellable. + (when (and *safepoint-counter* (back-edge-p instruction)) + (emit-back-edge-safepoint)) (call-next-method))) (defmethod translate-terminator ((instruction bir:unreachable) @@ -1876,6 +1908,8 @@ function-or-placeholder - the llvm function or a placeholder for do (setf (gethash phi *datum-values*) dat)))))) (defun layout-iblock (iblock abi) + (when (boundp '*laid-out-iblocks*) + (setf (gethash iblock *laid-out-iblocks*) t)) (cmp:irc-begin-block (iblock-tag iblock)) (cmp:with-landing-pad (maybe-entry-landing-pad (bir:dynamic-environment iblock) *tags*) @@ -1991,11 +2025,18 @@ function-or-placeholder - the llvm function or a placeholder for (arguments llvm-function-info)) when lexical ; skip unused fixed do (setf (gethash lexical *datum-values*) arg))) - ;; Branch to the start block. - (cmp:irc-br (iblock-tag (bir:start ir))) - ;; Lay out blocks. - (bir:do-iblocks (ib ir) - (layout-iblock ib abi)))))) + ;; Counter for sampled back-edge safepoints; must live in the alloca block. + (let ((*safepoint-counter* + (cmp:with-irbuilder (cmp:*irbuilder-function-alloca*) + (let ((c (cmp:alloca-size_t "safepoint-counter"))) + (cmp:irc-store (cmp:jit-constant-size_t 0) c) + c))) + (*laid-out-iblocks* (make-hash-table :test #'eq))) + ;; Branch to the start block. + (cmp:irc-br (iblock-tag (bir:start ir))) + ;; Lay out blocks. + (bir:do-iblocks (ib ir) + (layout-iblock ib abi))))))) ;; Finish up by jumping from the entry block to the body block (cmp:with-irbuilder (cmp:*irbuilder-function-alloca*) (cmp:irc-br body-block)) diff --git a/src/lisp/kernel/lsp/mp-package.lisp b/src/lisp/kernel/lsp/mp-package.lisp index 70e9c24f37..7f703642e6 100644 --- a/src/lisp/kernel/lsp/mp-package.lisp +++ b/src/lisp/kernel/lsp/mp-package.lisp @@ -21,5 +21,7 @@ signal-pending-interrupts raise without-interrupts with-interrupts with-local-interrupts with-restored-interrupts allow-with-interrupts interruptiblep + ;; deadlines + timeout timeout-seconds with-timeout call-with-timeout )) ) ; eval-when diff --git a/src/lisp/kernel/lsp/mp.lisp b/src/lisp/kernel/lsp/mp.lisp index ea45cc08eb..976fc3c41e 100644 --- a/src/lisp/kernel/lsp/mp.lisp +++ b/src/lisp/kernel/lsp/mp.lisp @@ -158,3 +158,59 @@ If DATUM is provided, it and ARGUMENTS designate a condition of default type SIM (if datum (core::coerce-to-condition datum arguments 'simple-error 'abort-process) nil))) + +#+threads +;; A subtype of ERROR, not just SERIOUS-CONDITION, so HANDLER-CASE on ERROR and +;; IGNORE-ERRORS catch it; SBCL and bordeaux-threads use SERIOUS-CONDITION. +(define-condition timeout (error) + ((%seconds :initarg :seconds :reader timeout-seconds)) + (:report (lambda (condition stream) + (format stream "Timed out after ~a second~:p." + (timeout-seconds condition))))) + +#+threads +(defun call-with-timeout (seconds function) + "Call FUNCTION, signalling TIMEOUT in this process if it runs longer than SECONDS. +The timeout is delivered as an interrupt, so it lands at the next safepoint. A body +blocked in a foreign call reaches no safepoint until the call returns, so there the +expiry is detected on return instead; either way TIMEOUT is signalled inside the +dynamic extent of the caller, where handlers are still established." + (let* ((lock (make-lock :name 'with-timeout)) + (cv (make-condition-variable :name 'with-timeout)) + (target *current-process*) + (deadline (+ (get-internal-real-time) + (round (* seconds internal-time-units-per-second)))) + (donep nil) + (firedp nil) + (watchdog + (process-run-function + 'with-timeout-watchdog + (lambda () + (with-lock (lock) + ;; Re-check DONEP because a timedwait may return early. + (loop until donep + for remaining = (/ (- deadline (get-internal-real-time)) + internal-time-units-per-second) + while (plusp remaining) + do (condition-variable-timedwait cv lock (float remaining 1d0))) + (unless donep + (setf firedp t) + (interrupt-process + target + ;; Re-check at delivery: an interrupt queued while the body was + ;; in a foreign call arrives after the body has already returned. + (lambda () (unless donep (error 'timeout :seconds seconds)))))))))) + (let ((values (unwind-protect (multiple-value-list (funcall function)) + (with-lock (lock) + (setf donep t) + (condition-variable-signal cv)) + (process-join watchdog)))) + ;; The deadline passed but the interrupt could not land in time. + (when firedp (error 'timeout :seconds seconds)) + (values-list values)))) + +#+threads +(defmacro with-timeout ((seconds) &body body) + "Execute BODY, signalling MP:TIMEOUT in this process if it has not finished +within SECONDS." + `(call-with-timeout ,seconds (lambda () ,@body))) diff --git a/src/lisp/regression-tests/mp.lisp b/src/lisp/regression-tests/mp.lisp index baf8b26c03..daa07f4ac9 100644 --- a/src/lisp/regression-tests/mp.lisp +++ b/src/lisp/regression-tests/mp.lisp @@ -253,3 +253,105 @@ (spam-processes nthreads (lambda () (mp:atomic-push nil (car place)))) (car place)) ((nil nil nil nil nil nil nil))) + + +;;; Returns true if THUNK's process is gone within SECONDS of being killed. The +;;; process must still be running when it is killed: a thunk that dies on its own +;;; would otherwise look exactly like a cancelled one and pass vacuously. +(defun cancelled-within-p (thunk seconds) + (let ((p (mp:process-run-function nil thunk))) + (loop repeat 200 until (mp:process-active-p p) do (sleep 0.01)) + (and (mp:process-active-p p) + (progn + (mp:process-kill p) + (loop repeat (ceiling seconds 0.01) + while (mp:process-active-p p) + do (sleep 0.01)) + (not (mp:process-active-p p)))))) + +;;; A loop body of pure VM opcodes reaches no function-call safepoint, so it is +;;; cancellable only if the interpreter polls interrupts on backward branches. +(test-true cancel-opcode-only-loop + (cancelled-within-p (lambda () (loop)) 3)) + +(test-true cancel-arithmetic-loop + (cancelled-within-p (lambda () (let ((x 0)) (loop (setq x (1+ x))))) 3)) + +;;; Control: a loop that calls a function was always cancellable. +(test-true cancel-loop-with-call + (cancelled-within-p (lambda () (loop (funcall #'identity 1))) 3)) + +;;; Native code reaches its own safepoints, so the VM's back-edge poll does not +;;; cover it. Asserts SIMPLE-CORE-FUN so it cannot pass by testing bytecode; +;;; vacuous where no native compiler exists. +(test-true cancel-native-opcode-only-loop + (let ((f (ignore-errors + (let ((cmp:*compile-native* t)) + (compile nil '(lambda () (loop))))))) + (if (typep f 'core:simple-core-fun) + (cancelled-within-p f 3) + t))) + +;;; A body that finishes in time returns normally and signals nothing. +(test with-timeout-completes + (mp:with-timeout (30) (+ 1 2)) + (3)) + +;;; A spinning body is interrupted; this only works because loops now poll. +(test-expect-error with-timeout-fires + (mp:with-timeout (0.2) (loop)) + :type mp:timeout) + +;;; The timeout must not fire after the body has already returned. An instant body +;;; never enqueues an interrupt, so this alone does not cover the race below. +(test-true with-timeout-no-late-fire + (progn (mp:with-timeout (0.2) t) + (sleep 0.5) + t)) + +;;; A blocking foreign call must park, so an interrupt can wake it with SIGCONT +;;; rather than sitting queued until the call returns on its own. +(test-expect-error with-timeout-blocking-foreign + (mp:with-timeout (0.5) (ext:system "sleep 3")) + :type mp:timeout) + +;;; ...and it must be woken PROMPTLY, not merely reported late on return. Three +;;; seconds of sleep must not elapse; without parking this takes the full 3s. +(test-true with-timeout-foreign-is-prompt + (let ((start (get-internal-real-time))) + (ignore-errors (mp:with-timeout (0.5) (ext:system "sleep 3"))) + (< (/ (- (get-internal-real-time) start) + internal-time-units-per-second) + 2.0))) + +;;; ...and the interrupt queued during that call must not fire afterwards. +(test-true with-timeout-foreign-no-late-fire + (progn (ignore-errors (mp:with-timeout (0.5) (ext:system "sleep 2"))) + (sleep 1) + t)) + +;;; A thread blocked in read(2) on an empty pipe must be cancellable. Without +;;; parking the thread is never marked blocking, so no SIGCONT is sent and it +;;; blocks forever. Bounded deliberately: an unbounded form would hang the whole +;;; suite rather than fail, which is how a 6-hour CI timeout happens. +(test-true cancel-blocking-read + (multiple-value-bind (r w) (core:pipe) + (declare (ignore w)) + (let ((buf (make-string 16 :element-type 'base-char))) + (cancelled-within-p (lambda () (core:read-fd r buf)) 3)))) + +;;; A thread blocked in accept(2) must be cancellable. A local socket is used so +;;; the test needs no port and no network. Verified to answer NO before the +;;; sockets were parked and YES after, so it is a real discriminator. +(test-true cancel-blocking-accept + (let* ((path (format nil "/tmp/clasp-accept-test-~a.sock" (get-universal-time))) + (sock (make-instance 'sb-bsd-sockets:local-socket :type :stream))) + (unwind-protect + (progn (sb-bsd-sockets:socket-bind sock path) + (sb-bsd-sockets:socket-listen sock 1) + (cancelled-within-p + (lambda () (sb-bsd-sockets:socket-accept sock)) 3)) + (ignore-errors (sb-bsd-sockets:socket-close sock)) + (ignore-errors (delete-file path))))) + + diff --git a/src/serveEvent/serveEvent.cc b/src/serveEvent/serveEvent.cc index d7cab9e7da..12f8f2bcf5 100644 --- a/src/serveEvent/serveEvent.cc +++ b/src/serveEvent/serveEvent.cc @@ -26,6 +26,7 @@ THE SOFTWARE. /* -^- */ #include +#include #include #include #include @@ -55,7 +56,10 @@ CL_DEFUN int serve_event_internal__ll_fdset_size() { return sizeof(fd_set); } DOCGROUP(clasp); CL_DEFUN core::Integer_mv serve_event_internal__ll_serveEventNoTimeout(clasp_ffi::ForeignData_sp rfd, clasp_ffi::ForeignData_sp wfd, int maxfdp1) { - gc::Fixnum selectRet = select(maxfdp1, rfd->data(), wfd->data(), NULL, NULL); + // Park: a NULL timeout blocks indefinitely. The fd_sets are foreign memory. + gc::Fixnum selectRet = BEGIN_PARK { + return (gc::Fixnum)select(maxfdp1, rfd->data(), wfd->data(), NULL, NULL); + } END_PARK; return Values(Integer_O::create(selectRet), Integer_O::create((gc::Fixnum)errno)); } @@ -69,7 +73,10 @@ CL_DEFUN core::Integer_mv serve_event_internal__ll_serveEventWithTimeout(clasp_f struct timeval tv; tv.tv_sec = seconds; tv.tv_usec = ((seconds - floor(seconds)) * 1e6); - gc::Fixnum selectRet = select(maxfdp1, rfd->data(), wfd->data(), NULL, &tv); + // Park: blocks for up to the caller's timeout. + gc::Fixnum selectRet = BEGIN_PARK { + return (gc::Fixnum)select(maxfdp1, rfd->data(), wfd->data(), NULL, &tv); + } END_PARK; return Values(Integer_O::create(selectRet), Integer_O::create((gc::Fixnum)errno)); } diff --git a/src/sockets/sockets.cc b/src/sockets/sockets.cc index 43bfdcf1b2..6c31b8b2d1 100644 --- a/src/sockets/sockets.cc +++ b/src/sockets/sockets.cc @@ -315,7 +315,8 @@ CL_DEFUN core::T_mv sockets_internal__ll_socketAccept_inetSocket(int sfd) { socklen_t addr_len = (socklen_t)sizeof(struct sockaddr_in); int new_fd; - new_fd = accept(sfd, (struct sockaddr*)&sockaddr, &addr_len); + // Park: accept() blocks until a connection arrives. + new_fd = BEGIN_PARK { return accept(sfd, (struct sockaddr*)&sockaddr, &addr_len); } END_PARK; int return0 = new_fd; core::T_sp return1 = nil(); @@ -343,7 +344,10 @@ CL_DEFUN int sockets_internal__ll_socketConnect_inetSocket(int port, int ip0, in struct sockaddr_in sockaddr; int output; fill_inet_sockaddr(&sockaddr, port, ip0, ip1, ip2, ip3); - output = connect(socket_file_descriptor, (struct sockaddr*)&sockaddr, sizeof(struct sockaddr_in)); + // Park: connect() blocks for the handshake. + output = BEGIN_PARK { + return connect(socket_file_descriptor, (struct sockaddr*)&sockaddr, sizeof(struct sockaddr_in)); + } END_PARK; return output; } @@ -512,7 +516,10 @@ DOCGROUP(clasp); CL_DEFUN core::T_mv sockets_internal__ll_socketAccept_localSocket(int socketFileDescriptor) { struct sockaddr_un sockaddr; socklen_t addr_len = (socklen_t)sizeof(struct sockaddr_un); - int new_fd = accept(socketFileDescriptor, (struct sockaddr*)&sockaddr, &addr_len); + // Park: accept() blocks until a connection arrives. + int new_fd = BEGIN_PARK { + return accept(socketFileDescriptor, (struct sockaddr*)&sockaddr, &addr_len); + } END_PARK; core::T_sp second_ret = nil(); if (new_fd != -1) { second_ret = core::SimpleBaseString_O::make(sockaddr.sun_path); @@ -534,7 +541,8 @@ CL_DEFUN int sockets_internal__ll_socketConnect_localSocket(int fd, int family, strncpy(sockaddr.sun_path, path.c_str(), sizeof(sockaddr.sun_path)); sockaddr.sun_path[sizeof(sockaddr.sun_path) - 1] = '\0'; - output = connect(fd, (struct sockaddr*)&sockaddr, sizeof(struct sockaddr_un)); + // Park: connect() blocks for the handshake. + output = BEGIN_PARK { return connect(fd, (struct sockaddr*)&sockaddr, sizeof(struct sockaddr_un)); } END_PARK; return output; } @@ -752,7 +760,10 @@ CL_DEFUN int sockets_internal__do_select(core::T_sp to_secs, unsigned int to_mus tv.tv_sec = to_secs.unsafe_fixnum(); tv.tv_usec = to_musecs; } - return select(max_fd + 1, (fd_set*)rfds->ptr(), NULL, NULL, (to_secs.fixnump()) ? &tv : NULL); + // Park: a non-fixnum timeout means NULL, i.e. block indefinitely. + return BEGIN_PARK { + return select(max_fd + 1, (fd_set*)rfds->ptr(), NULL, NULL, (to_secs.fixnump()) ? &tv : NULL); + } END_PARK; } void initialize_sockets_globals() {