From aeb1b0f7c726bffabd1db8ac7791e8e6be7dc05d Mon Sep 17 00:00:00 2001 From: Kinneyzhang Date: Tue, 1 Sep 2026 02:17:13 +0800 Subject: [PATCH] feat(reactive): isolate dispatcher contexts --- Makefile | 13 +- docs/architecture.en.md | 28 + docs/architecture.zh.md | 25 + etaf-actions.el | 24 +- etaf-data.el | 110 +-- etaf-events.el | 16 +- etaf-reactive.el | 189 +++-- etaf-runtime.el | 212 +++-- etaf-scheduler.el | 691 +++++++++++++++++ etaf.el | 1 + scripts/benchmark-scheduler-context.el | 142 ++++ tests/etaf-scheduler-tests.el | 733 ++++++++++++++++++ tests/etaf-tests.el | 36 +- .../etaf-m0a-condition-consumers.sexp | 10 +- 14 files changed, 2016 insertions(+), 214 deletions(-) create mode 100644 etaf-scheduler.el create mode 100644 scripts/benchmark-scheduler-context.el create mode 100644 tests/etaf-scheduler-tests.el diff --git a/Makefile b/Makefile index 09987cf..6173305 100644 --- a/Makefile +++ b/Makefile @@ -1,10 +1,10 @@ EMACS ?= emacs LOAD_PATH = -L . -L examples -L scripts -L ../ebox -L ../tp -L ../ecss -SOURCES = etaf-view.el etaf-compiler.el etaf-component.el etaf-reactive.el etaf-observer.el etaf-context.el etaf-theme-tp.el etaf-resource.el etaf-data.el etaf-generation.el etaf-host.el etaf-retirement.el etaf-render-port.el etaf-renderer.el etaf-runtime.el etaf-behavior.el etaf-actions.el etaf-events.el etaf-performance.el etaf.el scripts/emacs-gui-verifier.el +SOURCES = etaf-view.el etaf-compiler.el etaf-component.el etaf-scheduler.el etaf-reactive.el etaf-observer.el etaf-context.el etaf-theme-tp.el etaf-resource.el etaf-data.el etaf-generation.el etaf-host.el etaf-retirement.el etaf-render-port.el etaf-renderer.el etaf-runtime.el etaf-behavior.el etaf-actions.el etaf-events.el etaf-performance.el etaf.el scripts/emacs-gui-verifier.el scripts/benchmark-scheduler-context.el EXAMPLES = examples/etaf-counter-example.el examples/etaf-data-example.el examples/etaf-resource-example.el -TESTS = tests/etaf-tests.el tests/etaf-compiler-tests.el tests/etaf-component-frontends-tests.el tests/etaf-resource-tests.el tests/etaf-data-tests.el tests/etaf-theme-tp-tests.el tests/etaf-examples-tests.el tests/etaf-observer-tests.el tests/etaf-performance-tests.el tests/etaf-gui-verifier-tests.el tests/etaf-m0a-current-characterization-tests.el tests/etaf-interaction-contract-tests.el tests/etaf-m0b-component-manifest-tests.el tests/etaf-render-port-tests.el tests/etaf-generation-tests.el tests/etaf-host-tests.el tests/etaf-retirement-tests.el +TESTS = tests/etaf-tests.el tests/etaf-compiler-tests.el tests/etaf-component-frontends-tests.el tests/etaf-resource-tests.el tests/etaf-data-tests.el tests/etaf-theme-tp-tests.el tests/etaf-examples-tests.el tests/etaf-observer-tests.el tests/etaf-performance-tests.el tests/etaf-gui-verifier-tests.el tests/etaf-m0a-current-characterization-tests.el tests/etaf-interaction-contract-tests.el tests/etaf-m0b-component-manifest-tests.el tests/etaf-render-port-tests.el tests/etaf-generation-tests.el tests/etaf-host-tests.el tests/etaf-retirement-tests.el tests/etaf-scheduler-tests.el -.PHONY: test compile load checkdoc docs-check check clean +.PHONY: test compile load checkdoc docs-check scheduler-benchmark check clean test: compile $(EMACS) -Q --batch $(LOAD_PATH) --eval "(setq load-prefer-newer t)" \ @@ -24,10 +24,15 @@ docs-check: $(EMACS) -Q --batch $(LOAD_PATH) --eval "(setq load-prefer-newer t)" -l tests/etaf-docs-tests.el \ -f ert-run-tests-batch-and-exit +scheduler-benchmark: + $(EMACS) -Q --batch $(LOAD_PATH) --eval "(setq load-prefer-newer t)" \ + -l scripts/benchmark-scheduler-context.el \ + -f etaf-scheduler-benchmark-run + checkdoc: $(EMACS) -Q --batch --eval '(progn (require (quote checkdoc)) (dolist (directory (list "." "examples" "scripts")) (dolist (file (directory-files directory t)) (when (string-suffix-p ".el" file) (checkdoc-file file)))))' -check: checkdoc compile test docs-check +check: checkdoc compile test docs-check scheduler-benchmark clean: rm -f *.elc examples/*.elc scripts/*.elc tests/*.elc diff --git a/docs/architecture.en.md b/docs/architecture.en.md index bc40d49..3628897 100644 --- a/docs/architecture.en.md +++ b/docs/architecture.en.md @@ -493,6 +493,34 @@ can distinguish “committed, then callback failed” from a rollback failure. The buffer-kill path drains or contains retirement work but never throws a retirement condition from the kill hook. +Reactive publication is coordinated by an explicit `etaf-scheduler-context`. +The context owns the source and Runtime FIFOs, their dedupe sets, effect claims, +turn and projection epochs, nesting/busy state, fault diagnostics, and cost +counters; it does not own Component, resource, or generation state. Existing +callers use `etaf-scheduler-default-context`, while a mount may supply +`:scheduler-context` to isolate its dispatch authority. Scopes, Effects, and +opaque Runtime routes inherit and retain that context. + +One logical projection groups every changed source by its live subscriber +contexts in one subscriber-table scan before draining. Source propagation +settles across all touched contexts before any Runtime callback publishes. +Within a context, each source and Effect is delivered once per scheduler turn, +and multiple changed sources enqueue a Runtime once. A reentrant write to an +already delivered source is deferred to the next turn; a per-context turn +budget contains cross-context cycles with reusable fault diagnostics. Dedupe +in one context never suppresses another, and registry/token/Host validation +filters stale Runtime routes before fan-out. + +Runtime callbacks are detached as a turn, so a lifecycle write enters the +following turn. Data success/error multi-ref publication and event/action +callbacks use this same projection boundary, while the legacy facade continues +through the default context. Data source failures update the Controller error +state; a later projection/render failure propagates unchanged and cannot be +reclassified as a source failure. Runtime operation reports include both the +local context deltas and full cross-context projection summaries for source +delivery, subscriber visits, effect work, Runtime work, stale drops, turns, +and faults. + Each Runtime flush records a candidate-aware effect tuple containing the generation id, effect-to-source edges and source versions, plus an immutable semantic-node stamp for candidate input/context/output facts. A repeated tuple diff --git a/docs/architecture.zh.md b/docs/architecture.zh.md index 9f8cd56..b65955b 100644 --- a/docs/architecture.zh.md +++ b/docs/architecture.zh.md @@ -477,6 +477,31 @@ Ebox revision 与 diagnostic-journal ID。`etaf-condition-postcommit-info` 会 failure。buffer-kill 路径会 drain 或 contain retirement 工作,但绝不会从 kill hook 抛 retirement condition。 +Reactive publication 由显式 `etaf-scheduler-context` 协调。context 拥有 source +与 Runtime FIFO、对应 dedupe set、effect claim、turn/projection epoch、 +nesting/busy 状态、fault diagnostics 与成本计数器;它不拥有 Component、 +resource 或 generation state。既有调用者继续使用 +`etaf-scheduler-default-context`,mount 也可以传入 `:scheduler-context` 隔离 +dispatch authority。Scope、Effect 与 opaque Runtime route 会继承并保留该 +context。 + +一次 logical projection 会先按 live subscriber context 对所有 changed source +在一次 subscriber-table scan 中分组;所有已触及 context 的 source propagation +都稳定后,Runtime callback 才能 publish。在同一 context 中,每个 source 与 +Effect 每个 scheduler turn 只 delivery 一次,多个 changed source 也只 enqueue +同一 Runtime 一次。已经 delivery 的 source 若重入写入,会延后到下一个 turn; +per-context turn budget 会包含跨 context cycle,并留下可复用的 fault diagnostics。 +context A 的 dedupe 不会压制 context B;registry/token/Host 校验会在 fan-out 前 +过滤 stale Runtime route。 + +Runtime callback 先作为一个完整 turn detach,因此 lifecycle write 会进入下一个 +turn。Data success/error 的 multi-ref publication 与 event/action callback 都复用 +同一个 projection boundary,legacy façade 则继续走 default context。Data source +failure 会更新 Controller error state;之后发生的 projection/render failure 保持 +原 condition,不能被重新归类为 source failure。Runtime operation report 同时携带 +本 context delta 与完整 cross-context projection summary,覆盖 source delivery、 +subscriber visit、effect work、Runtime work、stale drop、turn 与 fault。 + 每次 Runtime flush 都记录 candidate-aware effect tuple:其中包含 generation id、 effect→source 边和 source version,以及 candidate input/context/output facts 的 immutable semantic-node stamp。重复 tuple 会报告有序的 effect/edge path;step diff --git a/etaf-actions.el b/etaf-actions.el index fd6c929..8e9ae7b 100644 --- a/etaf-actions.el +++ b/etaf-actions.el @@ -11,6 +11,8 @@ ;;; Code: (require 'cl-lib) +(require 'etaf-scheduler) +(require 'etaf-reactive) (require 'etaf-runtime) (defvar etaf--current-runtime) @@ -108,14 +110,20 @@ RUNTIME ACTION ...)'. Action functions receive Runtime first." (unless spec (signal 'etaf-action-error (list (format "Unknown ETAF Action: %S" action)))) - (if (or etaf--observer-context (etaf-runtime-observer runtime)) - (etaf-runtime-call-operation - runtime 'action (format "%S" action) - (lambda () - (let ((etaf--current-runtime runtime)) - (apply (etaf-action-spec-function spec) runtime arguments)))) - (let ((etaf--current-runtime runtime)) - (apply (etaf-action-spec-function spec) runtime arguments))))) + (let ((run + (lambda () + (etaf-scheduler-call-with-context + (etaf-runtime-scheduler-context runtime) + (lambda () + (etaf-reactive-call-with-batch + (lambda () + (let ((etaf--current-runtime runtime)) + (apply (etaf-action-spec-function spec) + runtime arguments))))))))) + (if (or etaf--observer-context (etaf-runtime-observer runtime)) + (etaf-runtime-call-operation + runtime 'action (format "%S" action) run) + (funcall run))))) ;;;###autoload (defun etaf-action-undefine (name) diff --git a/etaf-data.el b/etaf-data.el index 70bf9eb..e09e762 100644 --- a/etaf-data.el +++ b/etaf-data.el @@ -13,6 +13,7 @@ (require 'cl-lib) (require 'etaf-observer) +(require 'etaf-scheduler) (require 'etaf-reactive) (define-error 'etaf-data-error "Invalid ETAF data operation") @@ -285,36 +286,41 @@ use its result as `:initial-result' for `etaf-data-controller'." (when (= request-id (etaf-data--controller-request-id controller)) (let ((normalized (etaf-data--normalize-result result)) (auto-load-p (etaf-data--controller-auto-load-p controller))) - (unwind-protect - (progn - (setf (etaf-data--controller-auto-load-p controller) nil) - ;; ITEMS publication may synchronously schedule a DataGrid render. - ;; Materialize keyed selection dependencies first so render only - ;; reads retained reactive state. - (etaf-data--prepare-selected-refs - controller (plist-get normalized :items)) - (setf (etaf-value (etaf-data--controller-items controller)) - (plist-get normalized :items)) - (setf (etaf-value (etaf-data--controller-total controller)) - (or (plist-get normalized :total) - (length (plist-get normalized :items)))) - (when (plist-member normalized :page) - (setf (etaf-value (etaf-data--controller-page controller)) - (plist-get normalized :page))) - (when (plist-member normalized :page-size) - (setf (etaf-value (etaf-data--controller-page-size controller)) - (plist-get normalized :page-size))) - (setf (etaf-value (etaf-data--controller-error controller)) nil) - (setf (etaf-value (etaf-data--controller-status controller)) - 'success)) - (setf (etaf-data--controller-auto-load-p controller) auto-load-p)))) + ;; ITEMS publication may synchronously schedule a DataGrid render. + ;; Materialize keyed selection dependencies first so render only reads + ;; retained reactive state. + (etaf-data--prepare-selected-refs + controller (plist-get normalized :items)) + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-data--controller-auto-load-p controller) nil) + (etaf-scheduler-defer-finalizer + (lambda () + (setf (etaf-data--controller-auto-load-p controller) + auto-load-p))) + (setf (etaf-value (etaf-data--controller-items controller)) + (plist-get normalized :items)) + (setf (etaf-value (etaf-data--controller-total controller)) + (or (plist-get normalized :total) + (length (plist-get normalized :items)))) + (when (plist-member normalized :page) + (setf (etaf-value (etaf-data--controller-page controller)) + (plist-get normalized :page))) + (when (plist-member normalized :page-size) + (setf (etaf-value (etaf-data--controller-page-size controller)) + (plist-get normalized :page-size))) + (setf (etaf-value (etaf-data--controller-error controller)) nil) + (setf (etaf-value (etaf-data--controller-status controller)) + 'success))))) result) (defun etaf-data--apply-error (controller request-id error-data) "Publish ERROR-DATA for CONTROLLER when REQUEST-ID is current." (when (= request-id (etaf-data--controller-request-id controller)) - (setf (etaf-value (etaf-data--controller-error controller)) error-data) - (setf (etaf-value (etaf-data--controller-status controller)) 'error))) + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value (etaf-data--controller-error controller)) error-data) + (setf (etaf-value (etaf-data--controller-status controller)) 'error))))) ;;;###autoload (cl-defun etaf-data-controller @@ -447,20 +453,23 @@ load. INITIAL-RESULT and AUTO-LOAD are mutually exclusive." (etaf-data--require-controller controller) (let ((request-id (1+ (etaf-data--controller-request-id controller)))) (setf (etaf-data--controller-request-id controller) request-id) - (setf (etaf-value (etaf-data--controller-status controller)) 'loading) - (setf (etaf-value (etaf-data--controller-error controller)) nil) - (condition-case err - (etaf-data--apply-load-success - controller - request-id - (etaf-data--source-load - (etaf-data--controller-source controller) - (etaf-value (etaf-data--controller-query controller)) - (etaf-value (etaf-data--controller-page controller)) - (etaf-value (etaf-data--controller-page-size controller)))) - (error - (etaf-data--apply-error controller request-id err) - (signal (car err) (cdr err)))))) + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value (etaf-data--controller-status controller)) 'loading) + (setf (etaf-value (etaf-data--controller-error controller)) nil))) + (let ((result + (condition-case err + (etaf-data--source-load + (etaf-data--controller-source controller) + (etaf-value (etaf-data--controller-query controller)) + (etaf-value (etaf-data--controller-page controller)) + (etaf-value (etaf-data--controller-page-size controller))) + (error + (etaf-data--apply-error controller request-id err) + (signal (car err) (cdr err)))))) + ;; Projection/render failures after a successful source call are not + ;; source failures and must not rewrite the committed Controller refs. + (etaf-data--apply-load-success controller request-id result)))) ;;;###autoload (defun etaf-data-reload (controller) @@ -476,18 +485,21 @@ source mutation result after the reload succeeds." (etaf-data--require-controller controller) (let ((request-id (1+ (etaf-data--controller-request-id controller)))) (setf (etaf-data--controller-request-id controller) request-id) - (setf (etaf-value (etaf-data--controller-status controller)) 'loading) - (setf (etaf-value (etaf-data--controller-error controller)) nil) - (condition-case err - (let ((result + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value (etaf-data--controller-status controller)) 'loading) + (setf (etaf-value (etaf-data--controller-error controller)) nil))) + (let ((result + (condition-case err (etaf-data--source-mutate (etaf-data--controller-source controller) - operation payload))) - (etaf-data-load controller) - result) - (error - (etaf-data--apply-error controller request-id err) - (signal (car err) (cdr err)))))) + operation payload) + (error + (etaf-data--apply-error controller request-id err) + (signal (car err) (cdr err)))))) + ;; Reload owns its own source/projection classification. + (etaf-data-load controller) + result))) ;;;###autoload (defun etaf-data-stop (controller) diff --git a/etaf-events.el b/etaf-events.el index 43d9304..884875e 100644 --- a/etaf-events.el +++ b/etaf-events.el @@ -12,6 +12,8 @@ (require 'cl-lib) (require 'ebox) +(require 'etaf-scheduler) +(require 'etaf-reactive) (require 'etaf-runtime) (defvar etaf--current-runtime) @@ -67,7 +69,7 @@ otherwise call the local callback with no arguments." (etaf--assert-not-rendering 'dispatch-event) (setq runtime (etaf-runtime-require-mounted runtime)) - (let ((dispatch + (let* ((dispatch (lambda () (let ((callback (etaf--event-handler runtime host-ref kind))) (unless callback @@ -83,13 +85,19 @@ otherwise call the local callback with no arguments." (if payload-p (funcall callback payload) (funcall callback)) - (etaf-runtime-event-end runtime))))))) + (etaf-runtime-event-end runtime)))))) + (run + (lambda () + (etaf-scheduler-call-with-context + (etaf-runtime-scheduler-context runtime) + (lambda () + (etaf-reactive-call-with-batch dispatch)))))) (if (or etaf--observer-context (etaf-runtime-observer runtime)) (etaf-runtime-call-operation runtime 'event (format "%s %S" (etaf-event-kind kind) host-ref) - (lambda () (ebox-call-with-render-burst dispatch))) - (ebox-call-with-render-burst dispatch)))) + (lambda () (ebox-call-with-render-burst run))) + (ebox-call-with-render-burst run)))) ;;;###autoload (defun etaf-host-ref-bounds (runtime host-ref) diff --git a/etaf-reactive.el b/etaf-reactive.el index f25ffab..0c4a90c 100644 --- a/etaf-reactive.el +++ b/etaf-reactive.el @@ -12,6 +12,7 @@ (require 'cl-lib) (require 'gv) +(require 'etaf-scheduler) (define-error 'etaf-reactive-error "Invalid ETAF reactive operation") (define-error 'etaf-render-side-effect-error @@ -53,6 +54,7 @@ scheduler deps owner-scope + scheduler-context on-stop name (active-p t) @@ -65,13 +67,16 @@ effects children cleanups + scheduler-context name (active-p t)) (cl-defstruct (etaf-runtime-route (:constructor etaf-runtime-route-create)) "Opaque Runtime route stored by reactive sources." - runtime-id mount-epoch authority-token scheduler) + runtime-id mount-epoch authority-token scheduler scheduler-context + accepts-p + (active-p t)) (defvar etaf--active-effect nil "The effect currently collecting dependencies.") @@ -107,72 +112,50 @@ The function receives a zero-argument job and a phase symbol. Outside a Runtime, watchers run synchronously.") -(defvar etaf--dispatch-depth 0) -(defvar etaf--dispatch-source-queue nil) -(defvar etaf--dispatch-source-queue-tail nil) -(defvar etaf--dispatch-source-set (make-hash-table :test #'eq)) -(defvar etaf--dispatch-runtime-queue nil) -(defvar etaf--dispatch-runtime-queue-tail nil) -(defvar etaf--dispatch-runtime-set (make-hash-table :test #'eq)) -(defvar etaf--dispatch-effect-set (make-hash-table :test #'eq)) +(defun etaf--subscriber-scheduler-context (subscriber) + "Return SUBSCRIBER's scheduler context or the default context." + (or (cond + ((etaf-runtime-route-p subscriber) + (etaf-runtime-route-scheduler-context subscriber)) + ((etaf-effect-p subscriber) + (etaf-effect-scheduler-context subscriber))) + etaf-scheduler-default-context)) -(defun etaf--dispatch-append-source (source) - "Append SOURCE to the reactive FIFO in constant time." - (let ((cell (list source))) - (if etaf--dispatch-source-queue-tail - (setcdr etaf--dispatch-source-queue-tail cell) - (setq etaf--dispatch-source-queue cell)) - (setq etaf--dispatch-source-queue-tail cell))) +(defun etaf-runtime-route-live-p (route) + "Return non-nil when opaque Runtime ROUTE still owns dispatch authority." + (and (etaf-runtime-route-p route) + (etaf-runtime-route-active-p route) + (let ((predicate (etaf-runtime-route-accepts-p route))) + (or (null predicate) + (condition-case nil + (funcall predicate route) + ((error quit) nil)))))) -(defun etaf--dispatch-append-runtime (runtime) - "Append RUNTIME to the reactive publication FIFO in constant time." - (let ((cell (list runtime))) - (if etaf--dispatch-runtime-queue-tail - (setcdr etaf--dispatch-runtime-queue-tail cell) - (setq etaf--dispatch-runtime-queue cell)) - (setq etaf--dispatch-runtime-queue-tail cell))) +(defun etaf-reactive-enqueue-runtime-flush + (runtime function &optional scheduler-context) + "Queue FUNCTION once for RUNTIME in SCHEDULER-CONTEXT. +The default context preserves the legacy two-argument facade." + (etaf-scheduler-enqueue-runtime + (or scheduler-context (etaf-scheduler-current-context)) + runtime function)) -(defun etaf-reactive-enqueue-runtime-flush (runtime function) - "Queue FUNCTION once for RUNTIME after the outer reactive dispatch settles." - (unless (gethash runtime etaf--dispatch-runtime-set) - (puthash runtime function etaf--dispatch-runtime-set) - (etaf--dispatch-append-runtime runtime))) - -(defun etaf--dispatch-source-now (source) - "Notify SOURCE subscribers without opening another dispatch boundary." - (let ((subscribers (copy-hash-table (etaf--source-subscribers source)))) - (maphash - (lambda (subscriber _) - (cond - ((etaf-runtime-route-p subscriber) - (funcall (etaf-runtime-route-scheduler subscriber) - subscriber source)) - ((etaf-effect-active-p subscriber) - (unless (gethash subscriber etaf--dispatch-effect-set) - (puthash subscriber t etaf--dispatch-effect-set) - (if-let* ((scheduler (etaf-effect-scheduler subscriber))) - (funcall scheduler subscriber) - (etaf-reactive-effect-run subscriber)))))) - subscribers))) - -(defun etaf--drain-dispatch () - "Drain reactive sources and Runtime work to a stable outer fixed point." - (while (or etaf--dispatch-source-queue etaf--dispatch-runtime-queue) - (while etaf--dispatch-source-queue - (let ((source (pop etaf--dispatch-source-queue))) - (unless etaf--dispatch-source-queue - (setq etaf--dispatch-source-queue-tail nil)) - (remhash source etaf--dispatch-source-set) - (etaf--dispatch-source-now source))) - ;; Detach this turn. A lifecycle write may enqueue a source and the same - ;; Runtime again for the following turn without merging it into this one. - (let ((turn etaf--dispatch-runtime-queue)) - (setq etaf--dispatch-runtime-queue nil - etaf--dispatch-runtime-queue-tail nil) - (dolist (runtime turn) - (let ((function (gethash runtime etaf--dispatch-runtime-set))) - (remhash runtime etaf--dispatch-runtime-set) - (funcall function)))))) +(defun etaf--dispatch-subscriber-group + (context source _projection-epoch subscribers) + "Notify SOURCE SUBSCRIBERS already grouped for scheduler CONTEXT." + (dolist (subscriber subscribers) + (etaf-scheduler-record-subscriber-visit context) + (cond + ((etaf-runtime-route-p subscriber) + (if (etaf-runtime-route-live-p subscriber) + (funcall (etaf-runtime-route-scheduler subscriber) + subscriber source) + (etaf-scheduler-record-stale-route-drop context))) + ((and (etaf-effect-p subscriber) + (etaf-effect-active-p subscriber) + (etaf-scheduler-claim-effect context subscriber)) + (if-let* ((scheduler (etaf-effect-scheduler subscriber))) + (funcall scheduler subscriber) + (etaf-reactive-effect-run subscriber)))))) (defun etaf--reactive-same-p (left right) "Return whether LEFT and RIGHT are equal under ETAF's shallow rule." @@ -234,10 +217,14 @@ when the effect is disposed." (unless (functionp function) (signal 'wrong-type-argument (list 'functionp function))) (let* ((owner (or scope etaf--active-scope)) + (scheduler-context + (or (and owner (etaf-effect-scope-scheduler-context owner)) + (etaf-scheduler-current-context))) (effect (etaf--effect-create :function function :scheduler scheduler :owner-scope owner + :scheduler-context scheduler-context :on-stop on-stop :name name))) (when owner @@ -253,7 +240,12 @@ When RENDERING is non-nil, `etaf-value' writes are rejected for the duration of the run." (unless (etaf-effect-p effect) (signal 'wrong-type-argument (list 'etaf-effect-p effect))) - (cond + (let* ((scheduler-context + (or (etaf-effect-scheduler-context effect) + (etaf-scheduler-current-context))) + (etaf--scheduler-context scheduler-context)) + (etaf-scheduler-record-effect-evaluation scheduler-context) + (cond ((etaf-effect-running-p effect) (error "Recursive ETAF effect execution: %S" (or (etaf-effect-name effect) effect))) @@ -280,23 +272,45 @@ of the run." (etaf--clear-effect-deps effect) (dolist (source old-deps) (puthash effect t (etaf--source-subscribers source))) - (setf (etaf-effect-deps effect) old-deps)))))))) + (setf (etaf-effect-deps effect) old-deps))))))))) (defun etaf--dispatch-source (source) "Notify every current subscriber of SOURCE once." - (unless (gethash source etaf--dispatch-source-set) - (puthash source t etaf--dispatch-source-set) - (etaf--dispatch-append-source source)) - (when (zerop etaf--dispatch-depth) - (let ((etaf--dispatch-depth 1)) - (unwind-protect (etaf--drain-dispatch) - (setq etaf--dispatch-source-queue nil - etaf--dispatch-source-queue-tail nil - etaf--dispatch-runtime-queue nil - etaf--dispatch-runtime-queue-tail nil) - (clrhash etaf--dispatch-source-set) - (clrhash etaf--dispatch-runtime-set) - (clrhash etaf--dispatch-effect-set))))) + (etaf-scheduler-call-with-projection + (lambda () + (let ((groups (make-hash-table :test #'eq))) + (maphash + (lambda (subscriber _) + (let ((context (etaf--subscriber-scheduler-context subscriber))) + (cond + ((etaf-runtime-route-p subscriber) + (if (etaf-runtime-route-live-p subscriber) + (puthash context + (cons subscriber (gethash context groups)) groups) + (etaf-scheduler-record-stale-route-drop context))) + ((and (etaf-effect-p subscriber) + (etaf-effect-active-p subscriber)) + (puthash context + (cons subscriber (gethash context groups)) groups))))) + (etaf--source-subscribers source)) + (let (contexts) + (maphash (lambda (context _) (push context contexts)) groups) + (dolist (context + (sort contexts + (lambda (left right) + (< (etaf-scheduler-context-id left) + (etaf-scheduler-context-id right))))) + (let ((subscribers (nreverse (gethash context groups)))) + (etaf-scheduler-enqueue-source + context source + (lambda (delivery-context delivery-source projection-epoch) + (etaf--dispatch-subscriber-group + delivery-context delivery-source projection-epoch + subscribers)))))))))) + +(defun etaf-reactive-call-with-batch (function) + "Call FUNCTION inside one cross-context reactive projection." + (etaf-scheduler-call-with-projection function)) ;;;###autoload (cl-defun etaf-ref (initial-value &key test name) @@ -477,11 +491,20 @@ effect. If FUNCTION returns a function, it cleans up the previous run." (lambda () (etaf--stop-effect effect)))) ;;;###autoload -(cl-defun etaf-effect-scope (&key detached name) - "Create a Scope named NAME, owned by the current Scope unless DETACHED." +(cl-defun etaf-effect-scope (&key detached name scheduler-context) + "Create a Scope named NAME, owned by the current Scope unless DETACHED. +SCHEDULER-CONTEXT defaults to the parent or dynamically active context." (etaf--assert-not-rendering 'create-scope) (let* ((parent (and (not detached) etaf--active-scope)) - (scope (etaf--effect-scope-create :parent parent :name name))) + (scheduler-context + (or scheduler-context + (and parent (etaf-effect-scope-scheduler-context parent)) + (etaf-scheduler-current-context))) + (scope (etaf--effect-scope-create + :parent parent + :scheduler-context + (etaf-scheduler-context-resolve scheduler-context) + :name name))) (when parent (push scope (etaf-effect-scope-children parent))) scope)) @@ -499,6 +522,8 @@ effect. If FUNCTION returns a function, it cleans up the previous run." (etaf-effect-scope-active-p scope)) (error "Cannot enter an inactive ETAF Scope")) (let ((etaf--active-scope scope) + (etaf--scheduler-context + (etaf-effect-scope-scheduler-context scope)) (etaf--watch-scheduler (if watch-scheduler-p watch-scheduler etaf--watch-scheduler))) (funcall function))) diff --git a/etaf-runtime.el b/etaf-runtime.el index c46c33d..2a5f20d 100644 --- a/etaf-runtime.el +++ b/etaf-runtime.el @@ -14,6 +14,7 @@ (require 'ebox) (require 'etaf-view) (require 'etaf-component) +(require 'etaf-scheduler) (require 'etaf-reactive) (require 'etaf-generation) (require 'etaf-host) @@ -276,6 +277,7 @@ the sequential `etaf--pvec-put' contract." root-view root-node scope + scheduler-context root-effect-id root-range-id candidate-root-deps @@ -1239,15 +1241,25 @@ effect-to-source edges. It is intentionally immutable and suitable as an (puthash tuple t etaf--runtime-fixed-point-stamps) (push entry etaf--runtime-fixed-point-history)))) +(defun etaf--runtime-authorized-route-runtime (route) + "Return the live Runtime authorized by opaque ROUTE, or nil." + (when (and (etaf-runtime-route-p route) + (etaf-runtime-route-active-p route)) + (when-let* ((runtime + (gethash (etaf-runtime-route-mount-epoch route) + etaf--runtime-route-registry))) + (and (eq route (etaf-runtime-route-token runtime)) + (eq (etaf-runtime-route-scheduler-context route) + (etaf-runtime-scheduler-context runtime)) + (etaf-runtime-mounted-p runtime) + (etaf-host-authority-accepts-token-p + (etaf-runtime-host-authority runtime) + (etaf-runtime-route-authority-token route)) + runtime)))) + (defun etaf--runtime-route-scheduler (route source) "Route dirty SOURCE through opaque ROUTE to its current generation." - (when-let* ((runtime - (gethash (etaf-runtime-route-mount-epoch route) - etaf--runtime-route-registry))) - (when (and (etaf-runtime-mounted-p runtime) - (etaf-host-authority-accepts-token-p - (etaf-runtime-host-authority runtime) - (etaf-runtime-route-authority-token route))) + (when-let* ((runtime (etaf--runtime-authorized-route-runtime route))) (dolist (effect-id (etaf--generation-source-effects (etaf-runtime-current-generation runtime) source)) @@ -1258,7 +1270,8 @@ effect-to-source edges. It is intentionally immutable and suitable as an (etaf--runtime-mark-root-dirty runtime)))) (etaf-reactive-enqueue-runtime-flush (etaf-runtime-mount-epoch runtime) - (lambda () (etaf--runtime-request-flush runtime)))))) + (lambda () (etaf--runtime-request-flush runtime)) + (etaf-runtime-scheduler-context runtime)))) (defun etaf--runtime-evaluate-root-candidate (runtime) "Evaluate RUNTIME root while privately collecting dependencies." @@ -1362,6 +1375,11 @@ error or quit." (let* ((operation-id (1+ (or (etaf-runtime-next-operation-id runtime) 0))) (generation-before (etaf-runtime-generation runtime)) + (scheduler-context + (or (etaf-runtime-scheduler-context runtime) + etaf-scheduler-default-context)) + (scheduler-before + (etaf-scheduler-context-metrics scheduler-context)) (started (float-time)) (gc-count-before gcs-done) (gc-elapsed-before gc-elapsed) @@ -1376,34 +1394,101 @@ error or quit." (lambda (diagnostic) (etaf--runtime-record-observer-diagnostic runtime diagnostic)))) + projection-summaries result) (setf (etaf-runtime-next-operation-id runtime) operation-id) - (etaf--observer-call-with-context - context - (lambda () - (unwind-protect - (condition-case condition - (prog1 (setq result (funcall function)) - (setq status 'success)) - (quit - (setq status 'quit) - (signal (car condition) (cdr condition)))) - (etaf-observer-emit - (list - :provider 'etaf - :stage 'runtime-operation - :kind kind - :label (copy-sequence label) - :generation-before generation-before - :generation-after (etaf-runtime-generation runtime) - :status status - :duration-ms - (max 0.0 (* 1000.0 (- (float-time) started))) - :gc-count (- gcs-done gc-count-before) - :gc-duration-ms - (max 0.0 - (* 1000.0 (- gc-elapsed gc-elapsed-before)))))))) - result))))) + (cl-labels + ((run-body + () + (etaf--observer-call-with-context + context + (lambda () + (condition-case condition + (prog1 (setq result (funcall function)) + (setq status 'success)) + (quit + (setq status 'quit) + (signal (car condition) (cdr condition))))))) + (emit-final + (summaries) + (let ((scheduler-after + (etaf-scheduler-context-metrics scheduler-context))) + (etaf--observer-call-with-context + context + (lambda () + (etaf-observer-emit + (list + :provider 'etaf + :stage 'runtime-operation + :kind kind + :label (copy-sequence label) + :generation-before generation-before + :generation-after (etaf-runtime-generation runtime) + :status status + :scheduler-context-id + (plist-get scheduler-after :context-id) + :scheduler-projection-epoch-before + (plist-get scheduler-before :projection-epoch) + :scheduler-projection-epoch-after + (plist-get scheduler-after :projection-epoch) + :scheduler-turns + (- (plist-get scheduler-after :turn-count) + (plist-get scheduler-before :turn-count)) + :scheduler-source-enqueues + (- (plist-get scheduler-after :source-enqueues) + (plist-get scheduler-before :source-enqueues)) + :scheduler-source-dedupes + (- (plist-get scheduler-after :source-dedupes) + (plist-get scheduler-before :source-dedupes)) + :scheduler-source-deliveries + (- (plist-get scheduler-after :source-deliveries) + (plist-get scheduler-before :source-deliveries)) + :scheduler-subscriber-visits + (- (plist-get scheduler-after :subscriber-visits) + (plist-get scheduler-before :subscriber-visits)) + :scheduler-effect-claims + (- (plist-get scheduler-after :effect-claims) + (plist-get scheduler-before :effect-claims)) + :scheduler-effect-evaluations + (- (plist-get scheduler-after :effect-evaluations) + (plist-get scheduler-before :effect-evaluations)) + :scheduler-runtime-enqueues + (- (plist-get scheduler-after :runtime-enqueues) + (plist-get scheduler-before :runtime-enqueues)) + :scheduler-runtime-dedupes + (- (plist-get scheduler-after :runtime-dedupes) + (plist-get scheduler-before :runtime-dedupes)) + :scheduler-runtime-executions + (- (plist-get scheduler-after :runtime-executions) + (plist-get scheduler-before :runtime-executions)) + :scheduler-stale-route-drops + (- (plist-get scheduler-after :stale-route-drops) + (plist-get scheduler-before :stale-route-drops)) + :scheduler-faults + (- (plist-get scheduler-after :fault-count) + (plist-get scheduler-before :fault-count)) + :scheduler-projections (copy-tree summaries) + :scheduler-fault-state + (plist-get scheduler-after :fault-state) + :duration-ms + (max 0.0 (* 1000.0 (- (float-time) started))) + :gc-count (- gcs-done gc-count-before) + :gc-duration-ms + (max 0.0 + (* 1000.0 + (- gc-elapsed gc-elapsed-before)))))))))) + (if (etaf-scheduler-projection-active-p) + (unwind-protect + (run-body) + (etaf-scheduler-on-projection-complete + (lambda (summary) (emit-final (list summary))))) + (let ((etaf-scheduler-projection-observer + (lambda (summary) + (push summary projection-summaries)))) + (unwind-protect + (run-body) + (emit-final (nreverse projection-summaries))))) + result)))))) (cl-defmacro etaf--runtime-with-operation ((runtime kind label) &rest body) "Evaluate BODY in RUNTIME operation KIND and LABEL when observed." @@ -1425,15 +1510,19 @@ error or quit." "Enter one logical event batch for mounted RUNTIME." (setq runtime (etaf-runtime-require-mounted runtime)) (cl-incf (etaf-runtime-event-depth runtime)) + (etaf-scheduler-context-event-begin + (etaf-runtime-scheduler-context runtime)) runtime) (defun etaf-runtime-event-end (runtime) "Leave RUNTIME's logical event batch and publish pending state once." - (when (and (etaf-runtime-p runtime) - (etaf-runtime-mounted-p runtime)) + (when (etaf-runtime-p runtime) (setf (etaf-runtime-event-depth runtime) (max 0 (1- (etaf-runtime-event-depth runtime)))) - (when (and (zerop (etaf-runtime-event-depth runtime)) + (etaf-scheduler-context-event-end + (etaf-runtime-scheduler-context runtime)) + (when (and (etaf-runtime-mounted-p runtime) + (zerop (etaf-runtime-event-depth runtime)) (etaf-runtime-pending-p runtime)) (setf (etaf-runtime-pending-p runtime) nil) (etaf--runtime-request-flush runtime))) @@ -6109,7 +6198,10 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." (defun etaf--runtime-request-flush-now (runtime) "Flush RUNTIME now, or mark one follow-up flush while busy." - (when (etaf-runtime-mounted-p runtime) + (let ((etaf--scheduler-context + (or (etaf-runtime-scheduler-context runtime) + etaf-scheduler-default-context))) + (when (etaf-runtime-mounted-p runtime) (if (or (> (etaf-runtime-event-depth runtime) 0) (etaf-runtime-flushing-p runtime)) (setf (etaf-runtime-pending-p runtime) t) @@ -6137,7 +6229,7 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." (when (eq (etaf--runtime-component-overlay runtime) :root-fallback) (etaf--runtime-render-root-turn runtime t))))) - (setf (etaf-runtime-flushing-p runtime) nil)))))) + (setf (etaf-runtime-flushing-p runtime) nil))))))) (defun etaf--runtime-request-flush (runtime) "Flush RUNTIME within one optional Runtime operation boundary." @@ -6151,13 +6243,19 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." (etaf--runtime-request-flush runtime) (etaf-runtime-root-node runtime))) -(defun etaf--runtime-mount-now (buffer-or-name view observer) - "Mount VIEW with optional OBSERVER into BUFFER-OR-NAME." - (let* ((buffer (get-buffer-create buffer-or-name)) +(defun etaf--runtime-mount-now + (buffer-or-name view observer &optional scheduler-context) + "Mount VIEW into BUFFER-OR-NAME with OBSERVER and SCHEDULER-CONTEXT." + (setq scheduler-context + (or scheduler-context etaf-scheduler-default-context)) + (let ((etaf--scheduler-context scheduler-context)) + (let* ((buffer (get-buffer-create buffer-or-name)) (old (gethash buffer etaf--runtime-table))) (when old (etaf-runtime-unmount old)) - (let* ((scope (etaf-effect-scope :detached t :name buffer)) + (let* ((scope (etaf-effect-scope + :detached t :name buffer + :scheduler-context scheduler-context)) (mount-epoch (cl-incf etaf--mount-epoch-counter)) (host-authority (etaf-host-authority-create @@ -6166,6 +6264,7 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." :buffer buffer :root-view view :scope scope + :scheduler-context scheduler-context :instances (make-hash-table :test #'equal) :resource-registry (make-hash-table :test #'equal) :mount-epoch mount-epoch @@ -6199,7 +6298,9 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." :mount-epoch (etaf-runtime-mount-epoch runtime) :authority-token (etaf-host-authority-token host-authority) - :scheduler 'etaf--runtime-route-scheduler))) + :scheduler 'etaf--runtime-route-scheduler + :scheduler-context scheduler-context + :accepts-p 'etaf--runtime-authorized-route-runtime))) (setf (etaf-runtime-route-token runtime) route) (puthash (etaf-runtime-mount-epoch runtime) runtime etaf--runtime-route-registry)) @@ -6225,6 +6326,8 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." ;; Preaccept failures have no Host authority and discard registration. (unless (etaf-host-authority-attached-p host-authority) (setf (etaf-runtime-mounted-p runtime) nil) + (when-let* ((route (etaf-runtime-route-token runtime))) + (setf (etaf-runtime-route-active-p route) nil)) (etaf-host-authority-rollback-attach host-authority) (when (buffer-live-p buffer) (with-current-buffer buffer @@ -6235,7 +6338,7 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." etaf--runtime-route-registry) (etaf-scope-stop scope)) (signal (car err) (cdr err)))) - buffer))) + buffer)))) (defun etaf--runtime-validate-mount-options (options) "Validate and return mount OPTIONS." @@ -6245,7 +6348,8 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." (unless tail (error "ETAF mount option %S has no value" key)) (pop tail) - (unless (memq key '(:viewport-width :viewport-height :observer)) + (unless (memq key '(:viewport-width :viewport-height :observer + :scheduler-context)) (error "Unknown ETAF mount option: %S" key))))) (when-let* ((width (plist-get options :viewport-width))) (unless (and (numberp width) (> width 0)) @@ -6258,6 +6362,12 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." (functionp (plist-get options :observer))))) (signal 'wrong-type-argument (list 'functionp (plist-get options :observer)))) + (when (and (plist-member options :scheduler-context) + (not (etaf-scheduler-context-p + (plist-get options :scheduler-context)))) + (signal 'wrong-type-argument + (list 'etaf-scheduler-context-p + (plist-get options :scheduler-context)))) options) ;;;###autoload @@ -6266,15 +6376,17 @@ backend anchor proof failed; ordinary root turns keep their artifact reuse." Setup, reactive publication, Ebox rendering, and lifecycle work share one framework render-burst allocation budget when the installed Ebox supports it. OPTIONS may provide `:viewport-width' in pixels, `:viewport-height' in lines, -and `:observer' as a one-argument flat-report sink. The observer is installed -before the first Ebox publication." +`:observer' as a one-argument flat-report sink, and `:scheduler-context' for +explicit dispatch isolation. The observer is installed before publication." (etaf--assert-not-rendering 'mount-runtime) (setq options (etaf--runtime-validate-mount-options options)) (let ((ebox-viewport-width (plist-get options :viewport-width)) (ebox-viewport-height (plist-get options :viewport-height))) (ebox-call-with-render-burst #'etaf--runtime-mount-now buffer-or-name view - (plist-get options :observer)))) + (plist-get options :observer) + (or (plist-get options :scheduler-context) + etaf-scheduler-default-context)))) (defun etaf--runtime-unmount-now (runtime &optional cause) "Unmount RUNTIME idempotently and retire its Component scopes. @@ -6295,6 +6407,8 @@ unmount. Host authority is invalidated before any unbounded cleanup." (etaf-host-authority-begin-detach authority)) (unless (eq (etaf-host-authority-state authority) 'terminal) (etaf-host-authority-invalidate authority)) + (when-let* ((route (etaf-runtime-route-token runtime))) + (setf (etaf-runtime-route-active-p route) nil)) (setf (etaf-runtime-mounted-p runtime) nil) (unwind-protect (progn diff --git a/etaf-scheduler.el b/etaf-scheduler.el new file mode 100644 index 0000000..3e56e33 --- /dev/null +++ b/etaf-scheduler.el @@ -0,0 +1,691 @@ +;;; etaf-scheduler.el --- Explicit reactive dispatcher contexts -*- lexical-binding: t; -*- + +;; SPDX-License-Identifier: GPL-3.0-or-later + +;;; Commentary: + +;; Owns ETAF's cross-source dispatch queues. Contexts isolate FIFO and dedupe +;; authority, while a short-lived projection coordinates fan-out across every +;; context touched by one logical reactive publication. + +;;; Code: + +(require 'cl-lib) + +(define-error 'etaf-scheduler-error "Invalid ETAF scheduler operation") + +(defgroup etaf-scheduler nil + "Explicit ETAF reactive dispatcher contexts." + :group 'applications + :prefix "etaf-scheduler-") + +(defcustom etaf-scheduler-default-step-budget 128 + "Maximum dispatcher turns one context may run in one projection." + :type 'positive-integer + :group 'etaf-scheduler) + +(defvar etaf-scheduler--context-id-counter 0) +(defvar etaf-scheduler--projection-id-counter 0) + +(cl-defstruct + (etaf-scheduler-context + (:constructor etaf-scheduler-context--create)) + "Mutable queues and diagnostics for one isolated dispatcher." + id name + source-queue source-queue-tail source-set + deferred-source-queue deferred-source-queue-tail deferred-source-set + delivered-source-set + runtime-queue runtime-queue-tail runtime-set + effect-set + (active-turn-id 0) + (projection-epoch 0) + (completed-projection-epoch 0) + (depth 0) + (event-depth 0) + (busy-depth 0) + phase + fixed-point-step-budget + fault-state diagnostics + (source-enqueue-count 0) + (source-dedupe-count 0) + (source-delivery-count 0) + (subscriber-visit-count 0) + (effect-claim-count 0) + (effect-evaluation-count 0) + (runtime-enqueue-count 0) + (runtime-dedupe-count 0) + (runtime-execution-count 0) + (stale-route-drop-count 0) + (projection-fault-count 0) + (turn-count 0)) + +(cl-defstruct + (etaf-scheduler--projection + (:constructor etaf-scheduler--projection-create)) + "One cross-context reactive projection." + id all-contexts all-context-set context-before finalizers + completion-observers) + +(cl-defun etaf-scheduler-context-create (&key name fixed-point-step-budget) + "Create an isolated scheduler context named NAME. +FIXED-POINT-STEP-BUDGET is optional diagnostic metadata for its owner." + (when (and fixed-point-step-budget + (not (and (integerp fixed-point-step-budget) + (> fixed-point-step-budget 0)))) + (signal 'etaf-scheduler-error + (list :invalid-fixed-point-step-budget fixed-point-step-budget))) + (etaf-scheduler-context--create + :id (cl-incf etaf-scheduler--context-id-counter) + :name name + :source-set (make-hash-table :test #'eq) + :deferred-source-set (make-hash-table :test #'eq) + :delivered-source-set (make-hash-table :test #'eq) + :runtime-set (make-hash-table :test #'eq) + :effect-set (make-hash-table :test #'eq) + :fixed-point-step-budget + (or fixed-point-step-budget etaf-scheduler-default-step-budget))) + +(defvar etaf-scheduler-default-context + (etaf-scheduler-context-create :name 'default) + "Default scheduler context used by the compatibility facade.") + +(defvar etaf--scheduler-context nil + "Dynamically active scheduler context, or nil for the default context.") + +(defvar etaf-scheduler--active-projection nil + "Dynamically active cross-context projection.") + +(defvar etaf-scheduler-projection-observer nil + "Optional function receiving each completed projection summary.") + +(defun etaf-scheduler-current-context () + "Return the dynamically active scheduler context or the default context." + (or etaf--scheduler-context etaf-scheduler-default-context)) + +(defun etaf-scheduler-context-resolve (context) + "Return validated CONTEXT, defaulting nil to the current context." + (setq context (or context (etaf-scheduler-current-context))) + (unless (etaf-scheduler-context-p context) + (signal 'wrong-type-argument + (list 'etaf-scheduler-context-p context))) + context) + +(defun etaf-scheduler-call-with-context (context function) + "Call FUNCTION with CONTEXT as the active scheduler context." + (setq context (etaf-scheduler-context-resolve context)) + (unless (functionp function) + (signal 'wrong-type-argument (list 'functionp function))) + (cl-incf (etaf-scheduler-context-depth context)) + (let ((etaf--scheduler-context context)) + (unwind-protect + (funcall function) + (setf (etaf-scheduler-context-depth context) + (max 0 (1- (etaf-scheduler-context-depth context))))))) + +(defun etaf-scheduler-context-event-begin (context) + "Increment CONTEXT's logical event nesting depth." + (setq context (etaf-scheduler-context-resolve context)) + (cl-incf (etaf-scheduler-context-event-depth context)) + context) + +(defun etaf-scheduler-context-event-end (context) + "Decrement CONTEXT's logical event nesting depth." + (setq context (etaf-scheduler-context-resolve context)) + (setf (etaf-scheduler-context-event-depth context) + (max 0 (1- (etaf-scheduler-context-event-depth context)))) + context) + +(defun etaf-scheduler--append-context (projection context) + "Register CONTEXT with PROJECTION once." + (unless (gethash context + (etaf-scheduler--projection-all-context-set projection)) + (puthash context t + (etaf-scheduler--projection-all-context-set projection)) + (puthash context (etaf-scheduler-context-metrics context) + (etaf-scheduler--projection-context-before projection)) + (push context (etaf-scheduler--projection-all-contexts projection)))) + +(defun etaf-scheduler--append-source (context source) + "Append SOURCE to CONTEXT's FIFO in constant time." + (let ((cell (list source))) + (if (etaf-scheduler-context-source-queue-tail context) + (setcdr (etaf-scheduler-context-source-queue-tail context) cell) + (setf (etaf-scheduler-context-source-queue context) cell)) + (setf (etaf-scheduler-context-source-queue-tail context) cell))) + +(defun etaf-scheduler--append-runtime (context runtime) + "Append RUNTIME to CONTEXT's FIFO in constant time." + (let ((cell (list runtime))) + (if (etaf-scheduler-context-runtime-queue-tail context) + (setcdr (etaf-scheduler-context-runtime-queue-tail context) cell) + (setf (etaf-scheduler-context-runtime-queue context) cell)) + (setf (etaf-scheduler-context-runtime-queue-tail context) cell))) + +(defun etaf-scheduler--append-deferred-source (context source) + "Append SOURCE to CONTEXT's following-turn FIFO." + (let ((cell (list source))) + (if (etaf-scheduler-context-deferred-source-queue-tail context) + (setcdr + (etaf-scheduler-context-deferred-source-queue-tail context) cell) + (setf (etaf-scheduler-context-deferred-source-queue context) cell)) + (setf (etaf-scheduler-context-deferred-source-queue-tail context) cell))) + +(defun etaf-scheduler--enqueue-source-now + (projection context source delivery) + "Enqueue SOURCE and DELIVERY in CONTEXT under PROJECTION." + (etaf-scheduler--append-context projection context) + (let ((set (etaf-scheduler-context-source-set context)) + (deferred-set + (etaf-scheduler-context-deferred-source-set context))) + (cond + ((or (gethash source set) (gethash source deferred-set)) + (cl-incf (etaf-scheduler-context-source-dedupe-count context))) + ((and (eq (etaf-scheduler-context-phase context) 'source) + (gethash source + (etaf-scheduler-context-delivered-source-set context))) + (puthash source delivery deferred-set) + (etaf-scheduler--append-deferred-source context source) + (cl-incf (etaf-scheduler-context-source-enqueue-count context))) + (t + (puthash source delivery set) + (etaf-scheduler--append-source context source) + (cl-incf (etaf-scheduler-context-source-enqueue-count context)))))) + +(defun etaf-scheduler--enqueue-runtime-now + (projection context runtime function) + "Enqueue RUNTIME and FUNCTION in CONTEXT under PROJECTION." + (etaf-scheduler--append-context projection context) + (let ((set (etaf-scheduler-context-runtime-set context))) + (if (gethash runtime set) + (cl-incf (etaf-scheduler-context-runtime-dedupe-count context)) + (puthash runtime function set) + (etaf-scheduler--append-runtime context runtime) + (cl-incf (etaf-scheduler-context-runtime-enqueue-count context))))) + +(defun etaf-scheduler-enqueue-source (context source delivery) + "Queue DELIVERY of SOURCE once in CONTEXT for the active projection. +DELIVERY receives CONTEXT, SOURCE, and the projection epoch." + (setq context (etaf-scheduler-context-resolve context)) + (unless (functionp delivery) + (signal 'wrong-type-argument (list 'functionp delivery))) + (if etaf-scheduler--active-projection + (etaf-scheduler--enqueue-source-now + etaf-scheduler--active-projection context source delivery) + (etaf-scheduler-call-with-projection + (lambda () + (etaf-scheduler--enqueue-source-now + etaf-scheduler--active-projection context source delivery)))) + source) + +(defun etaf-scheduler-enqueue-runtime (context runtime function) + "Queue FUNCTION once for RUNTIME in CONTEXT after source delivery." + (setq context (etaf-scheduler-context-resolve context)) + (unless (functionp function) + (signal 'wrong-type-argument (list 'functionp function))) + (if etaf-scheduler--active-projection + (etaf-scheduler--enqueue-runtime-now + etaf-scheduler--active-projection context runtime function) + (etaf-scheduler-call-with-projection + (lambda () + (etaf-scheduler--enqueue-runtime-now + etaf-scheduler--active-projection context runtime function)))) + runtime) + +(defun etaf-scheduler-claim-effect (context effect) + "Return non-nil after claiming EFFECT once in CONTEXT's projection." + (setq context (etaf-scheduler-context-resolve context)) + (unless (gethash effect (etaf-scheduler-context-effect-set context)) + (puthash effect t (etaf-scheduler-context-effect-set context)) + (cl-incf (etaf-scheduler-context-effect-claim-count context)) + t)) + +(defun etaf-scheduler-record-subscriber-visit (context) + "Record one subscriber visit in scheduler CONTEXT." + (setq context (etaf-scheduler-context-resolve context)) + (when etaf-scheduler--active-projection + (etaf-scheduler--append-context + etaf-scheduler--active-projection context)) + (cl-incf (etaf-scheduler-context-subscriber-visit-count context))) + +(defun etaf-scheduler-record-effect-evaluation (context) + "Record one reactive effect evaluation in scheduler CONTEXT." + (setq context (etaf-scheduler-context-resolve context)) + (when etaf-scheduler--active-projection + (etaf-scheduler--append-context + etaf-scheduler--active-projection context)) + (cl-incf (etaf-scheduler-context-effect-evaluation-count context))) + +(defun etaf-scheduler-record-stale-route-drop (context) + "Record one stale route filtered before fan-out in CONTEXT." + (setq context (etaf-scheduler-context-resolve context)) + (when etaf-scheduler--active-projection + (etaf-scheduler--append-context + etaf-scheduler--active-projection context)) + (cl-incf (etaf-scheduler-context-stale-route-drop-count context))) + +(defun etaf-scheduler-defer-finalizer (function) + "Run FUNCTION after the active projection drains, or immediately if idle." + (unless (functionp function) + (signal 'wrong-type-argument (list 'functionp function))) + (if etaf-scheduler--active-projection + (push function + (etaf-scheduler--projection-finalizers + etaf-scheduler--active-projection)) + (funcall function)) + function) + +(defun etaf-scheduler-projection-active-p () + "Return non-nil while a scheduler projection is active." + (not (null etaf-scheduler--active-projection))) + +(defun etaf-scheduler-on-projection-complete (function) + "Call FUNCTION with the active projection summary when it completes." + (unless (functionp function) + (signal 'wrong-type-argument (list 'functionp function))) + (if etaf-scheduler--active-projection + (push function + (etaf-scheduler--projection-completion-observers + etaf-scheduler--active-projection)) + (funcall function nil)) + function) + +(defun etaf-scheduler--reset-pending (context) + "Discard CONTEXT's pending queues and dedupe membership." + (setf (etaf-scheduler-context-source-queue context) nil + (etaf-scheduler-context-source-queue-tail context) nil + (etaf-scheduler-context-deferred-source-queue context) nil + (etaf-scheduler-context-deferred-source-queue-tail context) nil + (etaf-scheduler-context-runtime-queue context) nil + (etaf-scheduler-context-runtime-queue-tail context) nil + (etaf-scheduler-context-phase context) nil) + (clrhash (etaf-scheduler-context-source-set context)) + (clrhash (etaf-scheduler-context-deferred-source-set context)) + (clrhash (etaf-scheduler-context-delivered-source-set context)) + (clrhash (etaf-scheduler-context-runtime-set context)) + context) + +(defun etaf-scheduler--promote-deferred-sources (context) + "Move CONTEXT's deferred sources into its following-turn FIFO." + (let ((deferred (etaf-scheduler-context-deferred-source-queue context)) + (deferred-set + (etaf-scheduler-context-deferred-source-set context)) + (source-set (etaf-scheduler-context-source-set context))) + (setf (etaf-scheduler-context-deferred-source-queue context) nil + (etaf-scheduler-context-deferred-source-queue-tail context) nil) + (dolist (source deferred) + (let ((delivery (gethash source deferred-set))) + (remhash source deferred-set) + (unless (gethash source source-set) + (puthash source delivery source-set) + (etaf-scheduler--append-source context source))))) + context) + +(defun etaf-scheduler--projection-contexts (projection) + "Return PROJECTION contexts in first-touch order." + (nreverse + (copy-sequence (etaf-scheduler--projection-all-contexts projection)))) + +(defun etaf-scheduler--projection-source-work-p (projection) + "Return non-nil when PROJECTION has pending source work." + (cl-some #'etaf-scheduler-context-source-queue + (etaf-scheduler--projection-contexts projection))) + +(defun etaf-scheduler--projection-work-p (projection) + "Return non-nil when PROJECTION has pending source or Runtime work." + (cl-some + (lambda (context) + (or (etaf-scheduler-context-source-queue context) + (etaf-scheduler-context-runtime-queue context))) + (etaf-scheduler--projection-contexts projection))) + +(defun etaf-scheduler--mark-turn + (context projection-id marked-contexts context-steps) + "Mark CONTEXT for PROJECTION-ID in MARKED-CONTEXTS and CONTEXT-STEPS." + (unless (gethash context marked-contexts) + (let* ((steps (1+ (gethash context context-steps 0))) + (budget (etaf-scheduler-context-fixed-point-step-budget context))) + (when (> steps budget) + (signal 'etaf-scheduler-error + (list :kind 'fixed-point-step-budget + :context-id (etaf-scheduler-context-id context) + :projection-id projection-id + :steps steps :budget budget))) + (puthash context steps context-steps)) + (puthash context t marked-contexts) + (setf (etaf-scheduler-context-projection-epoch context) projection-id) + (clrhash (etaf-scheduler-context-delivered-source-set context)) + (clrhash (etaf-scheduler-context-effect-set context)) + (cl-incf (etaf-scheduler-context-active-turn-id context)) + (cl-incf (etaf-scheduler-context-turn-count context)))) + +(defun etaf-scheduler--call-context-work (context function) + "Call FUNCTION as busy work owned by scheduler CONTEXT." + (cl-incf (etaf-scheduler-context-busy-depth context)) + (unwind-protect + (etaf-scheduler-call-with-context context function) + (setf (etaf-scheduler-context-busy-depth context) + (max 0 (1- (etaf-scheduler-context-busy-depth context)))))) + +(defun etaf-scheduler--drain-context-sources (context projection-id) + "Drain CONTEXT source FIFO for PROJECTION-ID." + (setf (etaf-scheduler-context-phase context) 'source) + (etaf-scheduler--call-context-work + context + (lambda () + (while (etaf-scheduler-context-source-queue context) + (let* ((source (pop (etaf-scheduler-context-source-queue context))) + (delivery + (gethash source + (etaf-scheduler-context-source-set context)))) + (unless (etaf-scheduler-context-source-queue context) + (setf (etaf-scheduler-context-source-queue-tail context) nil)) + (remhash source (etaf-scheduler-context-source-set context)) + (puthash source t + (etaf-scheduler-context-delivered-source-set context)) + (cl-incf (etaf-scheduler-context-source-delivery-count context)) + (when delivery + (funcall delivery context source projection-id))))))) + +(defun etaf-scheduler--detach-runtime-turns (projection) + "Detach and return PROJECTION's Runtime queues in context order." + (let (turns) + (dolist (context (etaf-scheduler--projection-contexts projection)) + (when-let* ((turn (etaf-scheduler-context-runtime-queue context))) + (setf (etaf-scheduler-context-runtime-queue context) nil + (etaf-scheduler-context-runtime-queue-tail context) nil) + (push (cons context turn) turns))) + (nreverse turns))) + +(defun etaf-scheduler--run-runtime-turn (context turn) + "Run detached Runtime TURN owned by CONTEXT." + (setf (etaf-scheduler-context-phase context) 'runtime) + (etaf-scheduler--call-context-work + context + (lambda () + (dolist (runtime turn) + (let ((function + (gethash runtime + (etaf-scheduler-context-runtime-set context)))) + (remhash runtime (etaf-scheduler-context-runtime-set context)) + (when function + (cl-incf (etaf-scheduler-context-runtime-execution-count context)) + (funcall function))))))) + +(defun etaf-scheduler--record-fault (context projection-id condition) + "Record CONTEXT failure CONDITION for PROJECTION-ID." + (setf (etaf-scheduler-context-fault-state context) (copy-tree condition)) + (cl-incf (etaf-scheduler-context-projection-fault-count context)) + (push (list :projection-epoch projection-id + :turn-id (etaf-scheduler-context-active-turn-id context) + :condition (copy-tree condition)) + (etaf-scheduler-context-diagnostics context)) + (when (> (length (etaf-scheduler-context-diagnostics context)) 64) + (setcdr (nthcdr 63 (etaf-scheduler-context-diagnostics context)) nil)) + condition) + +(defun etaf-scheduler--drain-projection (projection) + "Drain PROJECTION source-first and re-signal its first context fault." + (let ((projection-id (etaf-scheduler--projection-id projection)) + (failed-contexts (make-hash-table :test #'eq)) + (context-steps (make-hash-table :test #'eq)) + first-condition) + (while (etaf-scheduler--projection-work-p projection) + (let ((marked-contexts (make-hash-table :test #'eq))) + ;; All source propagation across every context settles before any + ;; Runtime can observe and publish the resulting reactive state. + (while (etaf-scheduler--projection-source-work-p projection) + (dolist (context (etaf-scheduler--projection-contexts projection)) + (cond + ((gethash context failed-contexts) + (etaf-scheduler--reset-pending context)) + ((etaf-scheduler-context-source-queue context) + (condition-case condition + (progn + (etaf-scheduler--mark-turn + context projection-id marked-contexts context-steps) + (etaf-scheduler--drain-context-sources + context projection-id)) + ((error quit) + (puthash context t failed-contexts) + (unless first-condition (setq first-condition condition)) + (etaf-scheduler--record-fault + context projection-id condition) + (etaf-scheduler--reset-pending context))))))) + ;; Detach every context's Runtime queue before callbacks run. A write + ;; from one callback therefore belongs to the following scheduler turn. + (dolist (entry (etaf-scheduler--detach-runtime-turns projection)) + (let ((context (car entry)) (turn (cdr entry))) + (unless (gethash context failed-contexts) + (condition-case condition + (progn + (etaf-scheduler--mark-turn + context projection-id marked-contexts context-steps) + (etaf-scheduler--run-runtime-turn context turn)) + ((error quit) + (puthash context t failed-contexts) + (unless first-condition (setq first-condition condition)) + (etaf-scheduler--record-fault + context projection-id condition) + (etaf-scheduler--reset-pending context)))))) + (dolist (context (etaf-scheduler--projection-contexts projection)) + (setf (etaf-scheduler-context-phase context) nil) + (unless (gethash context failed-contexts) + (etaf-scheduler--promote-deferred-sources context))))) + (dolist (context (etaf-scheduler--projection-contexts projection)) + (unless (gethash context failed-contexts) + (setf (etaf-scheduler-context-completed-projection-epoch context) + projection-id + (etaf-scheduler-context-fault-state context) nil))) + (when first-condition + (signal (car first-condition) (cdr first-condition))))) + +(defun etaf-scheduler--finish-projection (projection) + "Restore PROJECTION invariants and return contained finalizer errors." + (let (errors) + (dolist (context (etaf-scheduler--projection-all-contexts projection)) + (when (or (etaf-scheduler-context-source-queue context) + (etaf-scheduler-context-deferred-source-queue context) + (etaf-scheduler-context-runtime-queue context)) + (etaf-scheduler--reset-pending context)) + (setf (etaf-scheduler-context-phase context) nil) + (clrhash (etaf-scheduler-context-delivered-source-set context)) + (clrhash (etaf-scheduler-context-effect-set context))) + (let ((etaf-scheduler--active-projection nil)) + (dolist (finalizer + (nreverse (etaf-scheduler--projection-finalizers projection))) + (condition-case condition + (funcall finalizer) + ((error quit) (push condition errors))))) + (setf (etaf-scheduler--projection-finalizers projection) nil) + (nreverse errors))) + +(defun etaf-scheduler--notify-projection-complete (projection summary) + "Notify PROJECTION observers with detached SUMMARY, containing failures." + (let ((etaf-scheduler--active-projection nil) + errors) + (when (functionp etaf-scheduler-projection-observer) + (condition-case condition + (funcall etaf-scheduler-projection-observer (copy-tree summary)) + ((error quit) (push condition errors)))) + (dolist + (observer + (nreverse + (etaf-scheduler--projection-completion-observers projection))) + (condition-case condition + (funcall observer (copy-tree summary)) + ((error quit) (push condition errors)))) + (setf (etaf-scheduler--projection-completion-observers projection) nil) + (nreverse errors))) + +(defun etaf-scheduler--metric-delta (before after key) + "Return non-negative KEY delta between BEFORE and AFTER metric plists." + (max 0 (- (or (plist-get after key) 0) + (or (plist-get before key) 0)))) + +(defun etaf-scheduler--projection-summary (projection condition) + "Return detached PROJECTION metrics with optional failure CONDITION." + (let ((before-table + (etaf-scheduler--projection-context-before projection)) + context-summaries + (source-enqueues 0) (source-dedupes 0) (source-deliveries 0) + (subscriber-visits 0) (effect-claims 0) (effect-evaluations 0) + (runtime-enqueues 0) (runtime-dedupes 0) (runtime-executions 0) + (stale-route-drops 0) (turns 0) (faults 0)) + (dolist (context (etaf-scheduler--projection-contexts projection)) + (let* ((before (gethash context before-table)) + (after (etaf-scheduler-context-metrics context)) + (entry + (list + :context-id (etaf-scheduler-context-id context) + :source-enqueues + (etaf-scheduler--metric-delta + before after :source-enqueues) + :source-dedupes + (etaf-scheduler--metric-delta before after :source-dedupes) + :source-deliveries + (etaf-scheduler--metric-delta + before after :source-deliveries) + :subscriber-visits + (etaf-scheduler--metric-delta + before after :subscriber-visits) + :effect-claims + (etaf-scheduler--metric-delta before after :effect-claims) + :effect-evaluations + (etaf-scheduler--metric-delta + before after :effect-evaluations) + :runtime-enqueues + (etaf-scheduler--metric-delta before after :runtime-enqueues) + :runtime-dedupes + (etaf-scheduler--metric-delta before after :runtime-dedupes) + :runtime-executions + (etaf-scheduler--metric-delta + before after :runtime-executions) + :stale-route-drops + (etaf-scheduler--metric-delta + before after :stale-route-drops) + :turns (etaf-scheduler--metric-delta before after :turn-count) + :faults + (etaf-scheduler--metric-delta before after :fault-count)))) + (cl-incf source-enqueues (plist-get entry :source-enqueues)) + (cl-incf source-dedupes (plist-get entry :source-dedupes)) + (cl-incf source-deliveries (plist-get entry :source-deliveries)) + (cl-incf subscriber-visits (plist-get entry :subscriber-visits)) + (cl-incf effect-claims (plist-get entry :effect-claims)) + (cl-incf effect-evaluations (plist-get entry :effect-evaluations)) + (cl-incf runtime-enqueues (plist-get entry :runtime-enqueues)) + (cl-incf runtime-dedupes (plist-get entry :runtime-dedupes)) + (cl-incf runtime-executions (plist-get entry :runtime-executions)) + (cl-incf stale-route-drops (plist-get entry :stale-route-drops)) + (cl-incf turns (plist-get entry :turns)) + (cl-incf faults (plist-get entry :faults)) + (push entry context-summaries))) + (list :projection-id (etaf-scheduler--projection-id projection) + :completed-p (null condition) + :context-count (length context-summaries) + :source-enqueues source-enqueues + :source-dedupes source-dedupes + :source-deliveries source-deliveries + :subscriber-visits subscriber-visits + :effect-claims effect-claims + :effect-evaluations effect-evaluations + :runtime-enqueues runtime-enqueues + :runtime-dedupes runtime-dedupes + :runtime-executions runtime-executions + :stale-route-drops stale-route-drops + :turns turns :faults faults + :condition (and condition (copy-tree condition)) + :contexts (nreverse context-summaries)))) + +(defun etaf-scheduler-call-with-projection (function) + "Call FUNCTION and drain all contexts touched by its reactive writes. +Nested calls join the active projection. If FUNCTION signals after changing +state, queued notifications still drain before the original condition is +re-signaled." + (unless (functionp function) + (signal 'wrong-type-argument (list 'functionp function))) + (if etaf-scheduler--active-projection + (funcall function) + (let* ((projection + (etaf-scheduler--projection-create + :id (cl-incf etaf-scheduler--projection-id-counter) + :all-context-set (make-hash-table :test #'eq) + :context-before (make-hash-table :test #'eq))) + result primary-condition drain-condition summary finalizer-errors) + (let ((etaf-scheduler--active-projection projection)) + (unwind-protect + (progn + (condition-case condition + (setq result (funcall function)) + ((error quit) (setq primary-condition condition))) + (condition-case condition + (etaf-scheduler--drain-projection projection) + ((error quit) (setq drain-condition condition))) + (setq summary + (etaf-scheduler--projection-summary + projection drain-condition))) + (setq finalizer-errors + (etaf-scheduler--finish-projection projection)) + (when finalizer-errors + (setq summary + (plist-put summary :completed-p nil) + summary + (plist-put summary :finalizer-errors + (copy-tree finalizer-errors)))) + (etaf-scheduler--notify-projection-complete + projection summary))) + (cond + (primary-condition + (signal (car primary-condition) (cdr primary-condition))) + (drain-condition + (signal (car drain-condition) (cdr drain-condition))) + (t result))))) + +(defun etaf-scheduler-context-idle-p (context) + "Return non-nil when CONTEXT has no queued or running dispatcher work." + (setq context (etaf-scheduler-context-resolve context)) + (and (null (etaf-scheduler-context-source-queue context)) + (null (etaf-scheduler-context-deferred-source-queue context)) + (null (etaf-scheduler-context-runtime-queue context)) + (zerop (hash-table-count + (etaf-scheduler-context-source-set context))) + (zerop (hash-table-count + (etaf-scheduler-context-runtime-set context))) + (zerop (etaf-scheduler-context-busy-depth context)))) + +(defun etaf-scheduler-context-metrics (context) + "Return a detached metrics plist for scheduler CONTEXT." + (setq context (etaf-scheduler-context-resolve context)) + (list :context-id (etaf-scheduler-context-id context) + :projection-epoch (etaf-scheduler-context-projection-epoch context) + :completed-projection-epoch + (etaf-scheduler-context-completed-projection-epoch context) + :active-turn-id (etaf-scheduler-context-active-turn-id context) + :turn-count (etaf-scheduler-context-turn-count context) + :source-enqueues + (etaf-scheduler-context-source-enqueue-count context) + :source-dedupes + (etaf-scheduler-context-source-dedupe-count context) + :source-deliveries + (etaf-scheduler-context-source-delivery-count context) + :subscriber-visits + (etaf-scheduler-context-subscriber-visit-count context) + :effect-claims + (etaf-scheduler-context-effect-claim-count context) + :effect-evaluations + (etaf-scheduler-context-effect-evaluation-count context) + :runtime-enqueues + (etaf-scheduler-context-runtime-enqueue-count context) + :runtime-dedupes + (etaf-scheduler-context-runtime-dedupe-count context) + :runtime-executions + (etaf-scheduler-context-runtime-execution-count context) + :stale-route-drops + (etaf-scheduler-context-stale-route-drop-count context) + :fault-count + (etaf-scheduler-context-projection-fault-count context) + :fault-state + (copy-tree (etaf-scheduler-context-fault-state context)))) + +(provide 'etaf-scheduler) + +;;; etaf-scheduler.el ends here diff --git a/etaf.el b/etaf.el index 2a64efe..9c8a48f 100644 --- a/etaf.el +++ b/etaf.el @@ -30,6 +30,7 @@ (require 'etaf-view) (require 'etaf-compiler) (require 'etaf-component) +(require 'etaf-scheduler) (require 'etaf-reactive) (require 'etaf-observer) (require 'etaf-context) diff --git a/scripts/benchmark-scheduler-context.el b/scripts/benchmark-scheduler-context.el new file mode 100644 index 0000000..522d20b --- /dev/null +++ b/scripts/benchmark-scheduler-context.el @@ -0,0 +1,142 @@ +;;; benchmark-scheduler-context.el --- ETAF dispatcher microbenchmark -*- lexical-binding: t; -*- + +;; SPDX-License-Identifier: GPL-3.0-or-later + +;;; Commentary: + +;; Runs the M3a T6 fixed scheduler scenario: 40 changed sources fan out to one +;; Effect in each of four isolated contexts. Five warmups precede thirty timed +;; samples. Exact work counters and a 50 ms p95/max gate make the result useful +;; as a small dispatcher regression baseline; product GUI performance is a +;; separate end-to-end gate. + +;;; Code: + +(require 'cl-lib) +(require 'etaf) + +(defconst etaf-scheduler-benchmark-warmups 5) +(defconst etaf-scheduler-benchmark-samples 30) +(defconst etaf-scheduler-benchmark-source-count 40) +(defconst etaf-scheduler-benchmark-context-count 4) +(defconst etaf-scheduler-benchmark-limit-ms 50.0) + +(defun etaf-scheduler-benchmark--percentile (values percentile) + "Return nearest-rank PERCENTILE from numeric VALUES." + (let* ((ordered (sort (copy-sequence values) #'<)) + (rank (max 1 (ceiling (* percentile (length ordered)))))) + (nth (1- rank) ordered))) + +(defun etaf-scheduler-benchmark--metric-total (contexts key) + "Return the sum of metric KEY across scheduler CONTEXTS." + (cl-loop for context in contexts + sum (or (plist-get (etaf-scheduler-context-metrics context) key) + 0))) + +(defun etaf-scheduler-benchmark--iteration (sources iteration) + "Publish ITERATION through all reactive SOURCES once." + (etaf-reactive-call-with-batch + (lambda () + (cl-loop for source in sources + for index from 0 + do (setf (etaf-value source) + (+ (* iteration 1000) index)))))) + +;;;###autoload +(defun etaf-scheduler-benchmark-run () + "Run the fixed scheduler benchmark, print its result, and return it." + (let* ((contexts + (cl-loop for index below etaf-scheduler-benchmark-context-count + collect + (etaf-scheduler-context-create + :name (list 'benchmark index)))) + (sources + (cl-loop repeat etaf-scheduler-benchmark-source-count + collect (etaf-ref 0))) + effects durations result) + (unwind-protect + (progn + (dolist (context contexts) + (let ((effect + (etaf-scheduler-call-with-context + context + (lambda () + (etaf-reactive-effect-create + (lambda () (mapcar #'etaf-value sources))))))) + (push effect effects) + (etaf-reactive-effect-run effect))) + (dotimes (index etaf-scheduler-benchmark-warmups) + (etaf-scheduler-benchmark--iteration sources (1+ index))) + (let ((source-before + (etaf-scheduler-benchmark--metric-total + contexts :source-enqueues)) + (visit-before + (etaf-scheduler-benchmark--metric-total + contexts :subscriber-visits)) + (effect-before + (etaf-scheduler-benchmark--metric-total + contexts :effect-evaluations)) + (turn-before + (etaf-scheduler-benchmark--metric-total + contexts :turn-count))) + (dotimes (index etaf-scheduler-benchmark-samples) + (let ((started (float-time))) + (etaf-scheduler-benchmark--iteration + sources (+ etaf-scheduler-benchmark-warmups index 1)) + (push (* 1000.0 (- (float-time) started)) durations))) + (setq durations (nreverse durations)) + (let* ((expected-source-work + (* etaf-scheduler-benchmark-samples + etaf-scheduler-benchmark-source-count + etaf-scheduler-benchmark-context-count)) + (expected-context-work + (* etaf-scheduler-benchmark-samples + etaf-scheduler-benchmark-context-count)) + (source-work + (- (etaf-scheduler-benchmark--metric-total + contexts :source-enqueues) + source-before)) + (subscriber-visits + (- (etaf-scheduler-benchmark--metric-total + contexts :subscriber-visits) + visit-before)) + (effect-work + (- (etaf-scheduler-benchmark--metric-total + contexts :effect-evaluations) + effect-before)) + (turn-work + (- (etaf-scheduler-benchmark--metric-total + contexts :turn-count) + turn-before)) + (p95 + (etaf-scheduler-benchmark--percentile durations 0.95)) + (maximum (apply #'max durations))) + (unless (and (= source-work expected-source-work) + (= subscriber-visits expected-source-work) + (= effect-work expected-context-work) + (= turn-work expected-context-work)) + (error "Scheduler work counters diverged: %S" + (list source-work subscriber-visits + effect-work turn-work))) + (setq result + (list + :scenario 'scheduler-40-sources-4-contexts + :warmups etaf-scheduler-benchmark-warmups + :samples etaf-scheduler-benchmark-samples + :p95-ms p95 :max-ms maximum + :source-enqueues source-work + :subscriber-visits subscriber-visits + :effect-evaluations effect-work + :turns turn-work)) + (when (or (> p95 etaf-scheduler-benchmark-limit-ms) + (> maximum etaf-scheduler-benchmark-limit-ms)) + (error "Scheduler benchmark exceeds %.1f ms: %S" + etaf-scheduler-benchmark-limit-ms result))))) + (dolist (effect effects) (etaf--stop-effect effect))) + (prin1 result) + (terpri) + result)) + +(provide 'benchmark-scheduler-context) + +;;; benchmark-scheduler-context.el ends here diff --git a/tests/etaf-scheduler-tests.el b/tests/etaf-scheduler-tests.el new file mode 100644 index 0000000..f8a0370 --- /dev/null +++ b/tests/etaf-scheduler-tests.el @@ -0,0 +1,733 @@ +;;; etaf-scheduler-tests.el --- M3a dispatcher context gates -*- lexical-binding: t; -*- + +;; SPDX-License-Identifier: GPL-3.0-or-later + +(require 'cl-lib) +(require 'ert) +(require 'etaf) + +(define-error 'etaf-scheduler-test-condition + "ETAF scheduler test condition") +(define-error 'etaf-scheduler-test-source-condition + "ETAF scheduler source test condition") +(define-error 'etaf-scheduler-test-projection-condition + "ETAF scheduler projection test condition") +(define-error 'etaf-scheduler-test-body-condition + "ETAF scheduler body test condition") + +(etaf-define-component etaf-scheduler-test-pair (&key left right) + "Render reactive LEFT and RIGHT values." + :view + (text (expr (format "%s/%s" (etaf-value left) (etaf-value right))))) + +(etaf-define-component etaf-scheduler-test-single (&key source) + "Render one reactive SOURCE." + :view (text (expr (etaf-value source)))) + +(etaf-define-component etaf-scheduler-test-data (&key controller) + "Render CONTROLLER status, item count, and total." + :view + (text + (expr + (format "%s:%d:%d" + (etaf-value (etaf-data-status controller)) + (length (etaf-value (etaf-data-items controller))) + (etaf-value (etaf-data-total controller)))))) + +(etaf-define-component etaf-scheduler-test-data-projection-failure + (&key controller) + "Fail rendering when CONTROLLER reaches success." + :render + (let ((status (etaf-value (etaf-data-status controller)))) + (when (eq status 'success) + (signal 'etaf-scheduler-test-projection-condition '("projection"))) + (etaf-node 'text nil (list (symbol-name status))))) + +(defun etaf-scheduler-test--text (buffer-name) + "Return BUFFER-NAME text without properties." + (with-current-buffer buffer-name + (substring-no-properties (buffer-string)))) + +(defun etaf-scheduler-test--metric (context key) + "Return scheduler CONTEXT metric KEY." + (plist-get (etaf-scheduler-context-metrics context) key)) + +(defun etaf-scheduler-test--source (symbol) + "Return Lisp source text defining SYMBOL." + (let* ((loaded (symbol-file symbol 'defun)) + (source (if (and loaded (string-suffix-p ".elc" loaded)) + (substring loaded 0 -1) + loaded))) + (with-temp-buffer + (insert-file-contents source) + (buffer-string)))) + +(ert-deftest etaf-scheduler-is-the-only-dispatch-state-owner () + "Reactive semantics use the scheduler and own no process-global queues." + (let ((reactive (etaf-scheduler-test--source 'etaf--dispatch-source)) + (scheduler + (etaf-scheduler-test--source 'etaf-scheduler-context-create))) + (should-not (string-match-p "(defvar etaf--dispatch-" reactive)) + (dolist (field '("source-queue" "runtime-queue" "effect-set" + "active-turn-id" "projection-epoch" "fault-state")) + (should (string-match-p field scheduler))))) + +(ert-deftest etaf-scheduler-shared-sources-fan-out-once-per-context () + "Two changed refs schedule each isolated Runtime exactly once." + (let* ((left-buffer " *etaf-scheduler-left*") + (left-peer-buffer " *etaf-scheduler-left-peer*") + (right-buffer " *etaf-scheduler-right*") + (left-context (etaf-scheduler-context-create :name 'left)) + (right-context (etaf-scheduler-context-create :name 'right)) + (left (etaf-ref "A")) + (right (etaf-ref "1")) + (view (etaf--view-call + 'etaf-scheduler-test-pair + (list :left left :right right) nil))) + (unwind-protect + (progn + (etaf-mount left-buffer view + (list :scheduler-context left-context)) + (etaf-mount left-peer-buffer view + (list :scheduler-context left-context)) + (etaf-mount right-buffer view + (list :scheduler-context right-context)) + (let* ((left-runtime (etaf-runtime-for-buffer left-buffer)) + (left-peer-runtime + (etaf-runtime-for-buffer left-peer-buffer)) + (right-runtime (etaf-runtime-for-buffer right-buffer)) + (left-generation (etaf-runtime-generation left-runtime)) + (left-peer-generation + (etaf-runtime-generation left-peer-runtime)) + (right-generation (etaf-runtime-generation right-runtime)) + (left-source-before + (etaf-scheduler-test--metric + left-context :source-enqueues)) + (right-source-before + (etaf-scheduler-test--metric + right-context :source-enqueues)) + (left-runtime-before + (etaf-scheduler-test--metric + left-context :runtime-enqueues)) + (right-runtime-before + (etaf-scheduler-test--metric + right-context :runtime-enqueues))) + (should (eq left-context + (etaf-runtime-scheduler-context left-runtime))) + (should (eq right-context + (etaf-runtime-scheduler-context right-runtime))) + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value left) "B") + (setf (etaf-value right) "2"))) + (should (equal "B/2" (etaf-scheduler-test--text left-buffer))) + (should (equal "B/2" + (etaf-scheduler-test--text left-peer-buffer))) + (should (equal "B/2" (etaf-scheduler-test--text right-buffer))) + (should (= 1 (- (etaf-runtime-generation left-runtime) + left-generation))) + (should (= 1 (- (etaf-runtime-generation left-peer-runtime) + left-peer-generation))) + (should (= 1 (- (etaf-runtime-generation right-runtime) + right-generation))) + (should (= 2 (- (etaf-scheduler-test--metric + left-context :source-enqueues) + left-source-before))) + (should (= 2 (- (etaf-scheduler-test--metric + right-context :source-enqueues) + right-source-before))) + (should (= 2 (- (etaf-scheduler-test--metric + left-context :runtime-enqueues) + left-runtime-before))) + (should (= 1 (- (etaf-scheduler-test--metric + right-context :runtime-enqueues) + right-runtime-before))) + (should (etaf-scheduler-context-idle-p left-context)) + (should (etaf-scheduler-context-idle-p right-context)))) + (dolist (buffer-name + (list left-buffer left-peer-buffer right-buffer)) + (when-let* ((runtime (etaf-runtime-for-buffer buffer-name))) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer)))))) + +(ert-deftest etaf-scheduler-cross-context-computed-settles-before-runtime () + "A computed in one context settles before another context publishes." + (let* ((buffer-name " *etaf-scheduler-computed*") + (computed-context + (etaf-scheduler-context-create :name 'computed-owner)) + (runtime-context + (etaf-scheduler-context-create :name 'computed-consumer)) + (base (etaf-ref 1)) + (computed + (etaf-scheduler-call-with-context + computed-context + (lambda () + (etaf-computed (lambda () (* 2 (etaf-value base)))))))) + (unwind-protect + (progn + (etaf-mount + buffer-name + (etaf--view-call + 'etaf-scheduler-test-pair + (list :left base :right computed) nil) + (list :scheduler-context runtime-context)) + (let* ((runtime (etaf-runtime-for-buffer buffer-name)) + (generation (etaf-runtime-generation runtime)) + (executions + (etaf-scheduler-test--metric + runtime-context :runtime-executions))) + (should (eq computed-context + (etaf-effect-scheduler-context + (etaf-computed-effect computed)))) + (setf (etaf-value base) 2) + (should (equal "2/4" + (etaf-scheduler-test--text buffer-name))) + (should (= 1 (- (etaf-runtime-generation runtime) generation))) + (should (= 1 (- (etaf-scheduler-test--metric + runtime-context :runtime-executions) + executions))))) + (when-let* ((runtime (etaf-runtime-for-buffer buffer-name))) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer)) + (etaf--stop-effect (etaf-computed-effect computed))))) + +(ert-deftest etaf-scheduler-context-fault-does-not-suppress-peer () + "A failing context is contained until every peer receives the source." + (let* ((source (etaf-ref 0)) + (left-context (etaf-scheduler-context-create :name 'fault-left)) + (right-context (etaf-scheduler-context-create :name 'fault-right)) + (fail-p nil) + (right-value nil) + left-effect right-effect captured) + (setq left-effect + (etaf-scheduler-call-with-context + left-context + (lambda () + (etaf-reactive-effect-create + (lambda () + (prog1 (etaf-value source) + (when fail-p + (signal 'etaf-scheduler-test-condition '("left"))))))))) + (setq right-effect + (etaf-scheduler-call-with-context + right-context + (lambda () + (etaf-reactive-effect-create + (lambda () (setq right-value (etaf-value source))))))) + (unwind-protect + (progn + (etaf-reactive-effect-run left-effect) + (etaf-reactive-effect-run right-effect) + (setq fail-p t) + (condition-case condition + (setf (etaf-value source) 1) + (etaf-scheduler-test-condition (setq captured condition))) + (should captured) + (should (= right-value 1)) + (should (eq (car (etaf-scheduler-context-fault-state left-context)) + 'etaf-scheduler-test-condition)) + (should-not (etaf-scheduler-context-fault-state right-context)) + (should (etaf-scheduler-context-idle-p left-context)) + (should (etaf-scheduler-context-idle-p right-context))) + (etaf--stop-effect left-effect) + (etaf--stop-effect right-effect)))) + +(ert-deftest etaf-scheduler-body-error-precedes-drain-error-and-quit-drains () + "Body failure wins over drain failure; quit still drains queued effects." + (let* ((context (etaf-scheduler-context-create :name 'precedence)) + (source (etaf-ref 0)) + (quit-source (etaf-ref 0)) + (fail-p nil) (quit-seen nil) finalizer-ran summary + effect quit-effect captured quit-captured) + (setq effect + (etaf-scheduler-call-with-context + context + (lambda () + (etaf-reactive-effect-create + (lambda () + (prog1 (etaf-value source) + (when fail-p + (signal 'etaf-scheduler-test-condition '("drain"))))))))) + (setq quit-effect + (etaf-scheduler-call-with-context + context + (lambda () + (etaf-reactive-effect-create + (lambda () + (setq quit-seen (etaf-value quit-source))))))) + (unwind-protect + (progn + (etaf-reactive-effect-run effect) + (etaf-reactive-effect-run quit-effect) + (setq fail-p t) + (let ((etaf-scheduler-projection-observer + (lambda (value) (setq summary value)))) + (condition-case condition + (etaf-reactive-call-with-batch + (lambda () + (etaf-scheduler-defer-finalizer + (lambda () (error "contained finalizer"))) + (etaf-scheduler-defer-finalizer + (lambda () (setq finalizer-ran t))) + (setf (etaf-value source) 1) + (signal 'etaf-scheduler-test-body-condition '("body")))) + (etaf-scheduler-test-body-condition + (setq captured condition)))) + (should (equal captured + '(etaf-scheduler-test-body-condition "body"))) + (should finalizer-ran) + (should (= 1 (length (plist-get summary :finalizer-errors)))) + (should (eq (car (etaf-scheduler-context-fault-state context)) + 'etaf-scheduler-test-condition)) + (setq fail-p nil) + (condition-case condition + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value quit-source) 1) + (signal 'quit nil))) + (quit (setq quit-captured condition))) + (should (eq (car quit-captured) 'quit)) + (should (= quit-seen 1)) + (should (etaf-scheduler-context-idle-p context))) + (etaf--stop-effect effect) + (etaf--stop-effect quit-effect)))) + +(ert-deftest etaf-scheduler-reentrant-source-runs-in-following-turn () + "A same-source reentrant write is deferred to one following turn." + (let* ((context (etaf-scheduler-context-create :name 'reentrant)) + (source (etaf-ref 0)) + values turns effect) + (setq effect + (etaf-scheduler-call-with-context + context + (lambda () + (etaf-reactive-effect-create + (lambda () + (let ((value (etaf-value source))) + (push value values) + (push (etaf-scheduler-context-active-turn-id context) + turns) + (when (= value 1) + (setf (etaf-value source) 2)))))))) + (unwind-protect + (progn + (etaf-reactive-effect-run effect) + (let ((deliveries + (etaf-scheduler-test--metric + context :source-deliveries)) + (turn-count + (etaf-scheduler-test--metric context :turn-count))) + (setf (etaf-value source) 1) + (should (equal '(0 1 2) (nreverse values))) + (should (equal '(0 1 2) (nreverse turns))) + (should (= 2 (- (etaf-scheduler-test--metric + context :source-deliveries) + deliveries))) + (should (= 2 (- (etaf-scheduler-test--metric + context :turn-count) + turn-count))))) + (etaf--stop-effect effect)))) + +(ert-deftest etaf-scheduler-cross-context-cycle-is-bounded-and-reusable () + "Cross-context ping-pong fails at the exact budget and leaves idle contexts." + (let* ((left-context + (etaf-scheduler-context-create + :name 'cycle-left :fixed-point-step-budget 3)) + (right-context + (etaf-scheduler-context-create + :name 'cycle-right :fixed-point-step-budget 3)) + (left-source (etaf-ref 0)) + (right-source (etaf-ref 0)) + left-effect right-effect captured recovery-effect) + (setq left-effect + (etaf-scheduler-call-with-context + left-context + (lambda () + (etaf-reactive-effect-create + (lambda () + (let ((value (etaf-value left-source))) + (when (> value 0) + (setf (etaf-value right-source) value)))))))) + (setq right-effect + (etaf-scheduler-call-with-context + right-context + (lambda () + (etaf-reactive-effect-create + (lambda () + (let ((value (etaf-value right-source))) + (when (> value 0) + (setf (etaf-value left-source) (1+ value))))))))) + (unwind-protect + (progn + (etaf-reactive-effect-run left-effect) + (etaf-reactive-effect-run right-effect) + (condition-case condition + (setf (etaf-value left-source) 1) + (etaf-scheduler-error (setq captured condition))) + (should captured) + (should (eq 'fixed-point-step-budget + (plist-get (cdr captured) :kind))) + (should (= 4 (plist-get (cdr captured) :steps))) + (should (= 3 (plist-get (cdr captured) :budget))) + (should (etaf-scheduler-context-idle-p left-context)) + (should (etaf-scheduler-context-idle-p right-context)) + (etaf--stop-effect left-effect) + (etaf--stop-effect right-effect) + (let ((recovery-source (etaf-ref 0)) (seen nil)) + (setq recovery-effect + (etaf-scheduler-call-with-context + left-context + (lambda () + (etaf-reactive-effect-create + (lambda () + (setq seen (etaf-value recovery-source))))))) + (etaf-reactive-effect-run recovery-effect) + (setf (etaf-value recovery-source) 1) + (should (= seen 1)) + (should-not + (etaf-scheduler-context-fault-state left-context)))) + (when (and left-effect (etaf-effect-active-p left-effect)) + (etaf--stop-effect left-effect)) + (when (and right-effect (etaf-effect-active-p right-effect)) + (etaf--stop-effect right-effect)) + (when recovery-effect (etaf--stop-effect recovery-effect))))) + +(ert-deftest etaf-scheduler-data-success-is-one-turn-per-context () + "One Data success projection coalesces all changed refs in each context." + (let* ((left-buffer " *etaf-scheduler-data-left*") + (right-buffer " *etaf-scheduler-data-right*") + (left-context (etaf-scheduler-context-create :name 'data-left)) + (right-context (etaf-scheduler-context-create :name 'data-right)) + left-runtime-before right-runtime-before + (source + (etaf-data-source + :load + (lambda (_query _page _page-size) + (setq left-runtime-before + (etaf-scheduler-test--metric + left-context :runtime-enqueues) + right-runtime-before + (etaf-scheduler-test--metric + right-context :runtime-enqueues)) + (list :items '(one two three) :total 9)))) + (controller (etaf-data-controller source)) + (view (etaf--view-call + 'etaf-scheduler-test-data + (list :controller controller) nil))) + (unwind-protect + (progn + (etaf-mount left-buffer view + (list :scheduler-context left-context)) + (etaf-mount right-buffer view + (list :scheduler-context right-context)) + (etaf-data-load controller) + (should (equal "success:3:9" + (etaf-scheduler-test--text left-buffer))) + (should (equal "success:3:9" + (etaf-scheduler-test--text right-buffer))) + (should (= 1 (- (etaf-scheduler-test--metric + left-context :runtime-enqueues) + left-runtime-before))) + (should (= 1 (- (etaf-scheduler-test--metric + right-context :runtime-enqueues) + right-runtime-before)))) + (dolist (buffer-name (list left-buffer right-buffer)) + (when-let* ((runtime (etaf-runtime-for-buffer buffer-name))) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer))) + (etaf-data-stop controller)))) + +(ert-deftest etaf-scheduler-data-separates-source-and-projection-errors () + "Source failure sets Data error; render failure leaves successful refs." + (let* ((buffer-name " *etaf-scheduler-data-projection-error*") + (context (etaf-scheduler-context-create :name 'data-projection)) + (source + (etaf-data-source + :load (lambda (_query _page _page-size) + (list :items '(committed) :total 1)))) + (controller (etaf-data-controller source)) + projection-condition) + (unwind-protect + (progn + (etaf-mount + buffer-name + (etaf--view-call + 'etaf-scheduler-test-data-projection-failure + (list :controller controller) nil) + (list :scheduler-context context)) + (condition-case condition + (etaf-data-load controller) + (etaf-scheduler-test-projection-condition + (setq projection-condition condition))) + (should projection-condition) + (should (eq 'success + (etaf-value (etaf-data-status controller)))) + (should-not (etaf-value (etaf-data-error controller))) + (should (equal '(committed) + (etaf-value (etaf-data-items controller)))) + (should (equal "loading" + (etaf-scheduler-test--text buffer-name)))) + (when-let* ((runtime (etaf-runtime-for-buffer buffer-name))) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer)) + (etaf-data-stop controller))) + (let* ((source + (etaf-data-source + :load (lambda (_query _page _page-size) + (signal 'etaf-scheduler-test-source-condition + '("source"))))) + (controller (etaf-data-controller source)) + captured) + (unwind-protect + (progn + (condition-case condition + (etaf-data-load controller) + (etaf-scheduler-test-source-condition + (setq captured condition))) + (should captured) + (should (eq 'error (etaf-value (etaf-data-status controller)))) + (should (equal captured + (etaf-value (etaf-data-error controller))))) + (etaf-data-stop controller)))) + +(ert-deftest etaf-scheduler-detached-route-does-not-enter-fan-out () + "An invalidated Runtime route cannot enqueue work in its old context." + (let* ((buffer-name " *etaf-scheduler-stale-route*") + (context (etaf-scheduler-context-create :name 'stale)) + (source (etaf-ref "A"))) + (etaf-mount + buffer-name + (etaf--view-call 'etaf-scheduler-test-single (list :source source) nil) + (list :scheduler-context context)) + (let* ((runtime (etaf-runtime-for-buffer buffer-name)) + (route (etaf-runtime-route-token runtime))) + (etaf-unmount runtime) + (let ((before (etaf-scheduler-test--metric context :source-enqueues))) + (should-not (etaf-runtime-route-active-p route)) + (setf (etaf-value source) "B") + (should (= before + (etaf-scheduler-test--metric + context :source-enqueues))))) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer)))) + +(ert-deftest etaf-scheduler-stale-authority-is-filtered-before-fan-out () + "Registry, token, attaching, and detached mismatches enqueue no source." + (let* ((buffer-name " *etaf-scheduler-stale-authority*") + (context (etaf-scheduler-context-create :name 'stale-authority)) + (source (etaf-ref "A"))) + (etaf-mount + buffer-name + (etaf--view-call 'etaf-scheduler-test-single (list :source source) nil) + (list :scheduler-context context)) + (let* ((runtime (etaf-runtime-for-buffer buffer-name)) + (route (etaf-runtime-route-token runtime)) + (epoch (etaf-runtime-mount-epoch runtime)) + (authority (etaf-runtime-host-authority runtime)) + (slots (etaf-host-authority-slots authority)) + (token (etaf-runtime-route-authority-token route)) + (state (etaf-host-authority-state authority)) + (drops (etaf-scheduler-test--metric + context :stale-route-drops)) + (assert-drop + (lambda (value) + (let ((before (etaf-scheduler-test--metric + context :source-enqueues))) + (setf (etaf-value source) value) + (should (= before + (etaf-scheduler-test--metric + context :source-enqueues))))))) + (unwind-protect + (progn + (remhash epoch etaf--runtime-route-registry) + (funcall assert-drop "missing-registry") + (puthash epoch runtime etaf--runtime-route-registry) + (setf (etaf-runtime-route-authority-token route) + (make-symbol "stale-token")) + (funcall assert-drop "wrong-token") + (setf (etaf-runtime-route-authority-token route) token) + (aset slots 0 'attaching) + (funcall assert-drop "attaching") + (aset slots 0 'detached) + (funcall assert-drop "detached") + (aset slots 0 state) + (should (= 4 (- (etaf-scheduler-test--metric + context :stale-route-drops) + drops))) + (setf (etaf-value source) "live") + (should (equal "live" + (etaf-scheduler-test--text buffer-name)))) + (puthash epoch runtime etaf--runtime-route-registry) + (setf (etaf-runtime-route-authority-token route) token) + (aset slots 0 state) + (when (etaf-runtime-mounted-p runtime) + (etaf-unmount runtime)))) + (when-let* ((buffer (get-buffer buffer-name))) + (kill-buffer buffer)))) + +(ert-deftest etaf-scheduler-operation-report-exposes-cost-counters () + "Observed operations expose scheduler enqueue, dedupe, and turn costs." + (let* ((buffer-name " *etaf-scheduler-observer*") + (peer-buffer " *etaf-scheduler-observer-peer*") + (context (etaf-scheduler-context-create :name 'observer)) + (peer-context + (etaf-scheduler-context-create :name 'observer-peer)) + (left (etaf-ref "A")) + (right (etaf-ref "1")) + reports) + (unwind-protect + (progn + (etaf-mount + buffer-name + (etaf--view-call + 'etaf-scheduler-test-pair + (list :left left :right right) nil) + (list :scheduler-context context + :observer (lambda (report) (push report reports)))) + (etaf-mount + peer-buffer + (etaf--view-call + 'etaf-scheduler-test-pair + (list :left left :right right) nil) + (list :scheduler-context peer-context)) + (setq reports nil) + (let ((runtime (etaf-runtime-for-buffer buffer-name))) + (etaf-runtime-call-operation + runtime 'scheduler-test "scheduler counters" + (lambda () + (etaf-reactive-call-with-batch + (lambda () + (setf (etaf-value left) "B") + (setf (etaf-value right) "2"))))) + (let ((final (car reports))) + (should (eq 'runtime-operation (plist-get final :stage))) + (should (= (etaf-scheduler-context-id context) + (plist-get final :scheduler-context-id))) + (should (= 2 (plist-get final :scheduler-source-enqueues))) + (should (= 2 (plist-get final :scheduler-source-deliveries))) + (should (= 2 (plist-get final :scheduler-subscriber-visits))) + (should (= 1 (plist-get final :scheduler-runtime-enqueues))) + (should (= 1 (plist-get final :scheduler-runtime-dedupes))) + (should (= 1 (plist-get final :scheduler-runtime-executions))) + (should (= 1 (plist-get final :scheduler-turns))) + (let ((projection + (car (plist-get final :scheduler-projections)))) + (should (= 2 (plist-get projection :context-count))) + (should (= 4 (plist-get projection :source-enqueues))) + (should (= 4 (plist-get projection :subscriber-visits))) + (should (= 2 (plist-get projection :runtime-executions)))) + (should (> (plist-get final :scheduler-projection-epoch-after) + (plist-get final + :scheduler-projection-epoch-before)))))) + (dolist (name (list buffer-name peer-buffer)) + (when-let* ((runtime (etaf-runtime-for-buffer name))) + (when (eq name buffer-name) + (etaf-runtime-set-observer runtime nil)) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer name))) + (kill-buffer buffer)))))) + +(ert-deftest etaf-scheduler-real-flush-defers-final-report-to-projection-end () + "A direct ref write reports the complete cross-context projection." + (let* ((buffer-name " *etaf-scheduler-real-flush*") + (peer-buffer " *etaf-scheduler-real-flush-peer*") + (context (etaf-scheduler-context-create :name 'real-flush)) + (peer-context + (etaf-scheduler-context-create :name 'real-flush-peer)) + (source (etaf-ref "A")) reports) + (unwind-protect + (progn + (etaf-mount + buffer-name + (etaf--view-call + 'etaf-scheduler-test-single (list :source source) nil) + (list :scheduler-context context + :observer (lambda (report) (push report reports)))) + (etaf-mount + peer-buffer + (etaf--view-call + 'etaf-scheduler-test-single (list :source source) nil) + (list :scheduler-context peer-context)) + (setq reports nil) + (setf (etaf-value source) "B") + (let* ((final (car reports)) + (projection + (car (plist-get final :scheduler-projections)))) + (should (eq 'runtime-operation (plist-get final :stage))) + (should (= 2 (plist-get projection :context-count))) + (should (= 2 (plist-get projection :source-enqueues))) + (should (= 2 (plist-get projection :subscriber-visits))) + (should (= 2 (plist-get projection :runtime-executions))) + (should (equal "B" (etaf-scheduler-test--text buffer-name))) + (should (equal "B" (etaf-scheduler-test--text peer-buffer))))) + (dolist (name (list buffer-name peer-buffer)) + (when-let* ((runtime (etaf-runtime-for-buffer name))) + (when (equal name buffer-name) + (etaf-runtime-set-observer runtime nil)) + (etaf-unmount runtime)) + (when-let* ((buffer (get-buffer name))) + (kill-buffer buffer)))))) + +(ert-deftest etaf-scheduler-scale-matrix-scans-each-subscriber-once () + "Source-by-context fan-out has exact linear subscriber and effect costs." + (let* ((contexts + (cl-loop for index below 4 + collect (etaf-scheduler-context-create + :name (list 'scale index)))) + (sources (cl-loop repeat 40 collect (etaf-ref 0))) + effects summary + (started (float-time))) + (unwind-protect + (progn + (dolist (context contexts) + (let ((effect + (etaf-scheduler-call-with-context + context + (lambda () + (etaf-reactive-effect-create + (lambda () + (mapcar #'etaf-value sources))))))) + (push effect effects) + (etaf-reactive-effect-run effect))) + (let ((etaf-scheduler-projection-observer + (lambda (value) (setq summary value)))) + (etaf-reactive-call-with-batch + (lambda () + (cl-loop for source in sources + for value from 1 + do (setf (etaf-value source) value))))) + (should (= 4 (plist-get summary :context-count))) + (should (= 160 (plist-get summary :source-enqueues))) + (should (= 160 (plist-get summary :source-deliveries))) + (should (= 160 (plist-get summary :subscriber-visits))) + (should (= 4 (plist-get summary :effect-claims))) + (should (= 4 (plist-get summary :effect-evaluations))) + (should (= 4 (plist-get summary :turns))) + (should (< (* 1000.0 (- (float-time) started)) 1000.0))) + (dolist (effect effects) (etaf--stop-effect effect))))) + +(ert-deftest etaf-scheduler-fifo-cost-counters-are-linear () + "A bounded synthetic turn records one enqueue and one dedupe per key." + (let* ((context (etaf-scheduler-context-create :name 'cost)) + (sources (cl-loop repeat 500 collect (list (make-symbol "source")))) + (started (float-time))) + (etaf-scheduler-call-with-projection + (lambda () + (dolist (source sources) + (etaf-scheduler-enqueue-source context source #'ignore) + (etaf-scheduler-enqueue-source context source #'ignore)))) + (let ((elapsed-ms (* 1000.0 (- (float-time) started)))) + (should (= 500 (etaf-scheduler-test--metric + context :source-enqueues))) + (should (= 500 (etaf-scheduler-test--metric + context :source-dedupes))) + (should (< elapsed-ms 1000.0)) + (should (etaf-scheduler-context-idle-p context))))) + +(provide 'etaf-scheduler-tests) + +;;; etaf-scheduler-tests.el ends here diff --git a/tests/etaf-tests.el b/tests/etaf-tests.el index 3f583a2..05a1581 100644 --- a/tests/etaf-tests.el +++ b/tests/etaf-tests.el @@ -2240,24 +2240,26 @@ Event composition is a Runtime contract, not a UI-library helper contract." (ert-deftest etaf-reactive-dispatch-fifos-append-in-constant-time-order () "Keep source and Runtime dispatch order with explicit FIFO tails." - (let ((etaf--dispatch-depth 1) - (etaf--dispatch-source-queue nil) - (etaf--dispatch-source-queue-tail nil) - (etaf--dispatch-source-set (make-hash-table :test #'eq)) - (etaf--dispatch-runtime-queue nil) - (etaf--dispatch-runtime-queue-tail nil) - (etaf--dispatch-runtime-set (make-hash-table :test #'eq)) + (let ((context (etaf-scheduler-context-create :name 'fifo-test)) (sources (cl-loop repeat 100 collect (etaf-ref nil)))) - (dolist (source sources) - (etaf--dispatch-source source)) - (etaf--dispatch-source (car sources)) - (etaf-reactive-enqueue-runtime-flush 'first #'ignore) - (etaf-reactive-enqueue-runtime-flush 'second #'ignore) - (etaf-reactive-enqueue-runtime-flush 'first #'ignore) - (should (equal sources etaf--dispatch-source-queue)) - (should (eq (car etaf--dispatch-source-queue-tail) (car (last sources)))) - (should (equal '(first second) etaf--dispatch-runtime-queue)) - (should (eq 'second (car etaf--dispatch-runtime-queue-tail))))) + (etaf-scheduler-call-with-projection + (lambda () + (dolist (source sources) + (etaf-scheduler-enqueue-source context source #'ignore)) + (etaf-scheduler-enqueue-source context (car sources) #'ignore) + (etaf-scheduler-enqueue-runtime context 'first #'ignore) + (etaf-scheduler-enqueue-runtime context 'second #'ignore) + (etaf-scheduler-enqueue-runtime context 'first #'ignore) + (should (equal sources + (etaf-scheduler-context-source-queue context))) + (should (eq (car (etaf-scheduler-context-source-queue-tail context)) + (car (last sources)))) + (should (equal '(first second) + (etaf-scheduler-context-runtime-queue context))) + (should + (eq 'second + (car (etaf-scheduler-context-runtime-queue-tail context)))))) + (should (etaf-scheduler-context-idle-p context)))) (ert-deftest etaf-runtime-dirty-effect-fifo-and-priority-stay-stable () "Append dirty effects in O(1) and sort each immutable effect only once." diff --git a/tests/fixtures/etaf-m0a-condition-consumers.sexp b/tests/fixtures/etaf-m0a-condition-consumers.sexp index 70f6d22..4d94a12 100644 --- a/tests/fixtures/etaf-m0a-condition-consumers.sexp +++ b/tests/fixtures/etaf-m0a-condition-consumers.sexp @@ -8,6 +8,7 @@ (:file "etaf-performance.el" :form condition-case :conditions (error) :owner etaf-performance :policy generic-containment) (:file "etaf-reactive.el" :form condition-case :conditions (error quit) :owner etaf-reactive :policy generic-containment) (:file "etaf-reactive.el" :form condition-case :conditions (error quit) :owner etaf-reactive :policy generic-containment) + (:file "etaf-reactive.el" :form condition-case :conditions (error quit) :owner etaf-reactive :policy generic-containment) (:file "etaf-render-port.el" :form condition-case :conditions (error quit) :owner etaf-render-port :policy generic-containment) (:file "etaf-render-port.el" :form condition-case :conditions (error quit) :owner etaf-render-port :policy generic-containment) (:file "etaf-render-port.el" :form condition-case :conditions (error quit) :owner etaf-render-port :policy generic-containment) @@ -24,4 +25,11 @@ (:file "etaf-runtime.el" :form condition-case :conditions (error quit) :owner etaf-runtime :policy generic-containment) (:file "etaf-runtime.el" :form condition-case :conditions (error quit) :owner etaf-runtime :policy generic-containment) (:file "etaf-runtime.el" :form condition-case :conditions (error quit) :owner etaf-runtime :policy generic-containment) - (:file "etaf-runtime.el" :form condition-case :conditions (error quit) :owner etaf-runtime :policy generic-containment)) + (:file "etaf-runtime.el" :form condition-case :conditions (error quit) :owner etaf-runtime :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment) + (:file "etaf-scheduler.el" :form condition-case :conditions (error quit) :owner etaf-scheduler :policy generic-containment))