perf: O(1) bounded worker-pool FIFO queue (cons+reverse two-list)
ober
e3a79f24dd37ed3d8eb5186c9cca6cfaf111b691
--- a/lib/secmon/server/limits.sls +++ b/lib/secmon/server/limits.sls @@ -141,31 +141,44 @@ (let ([row (find-row (a-rows admission) peer)]) (if row (vector-ref row 1) 0))))) - ;; #(fixed-pool workers max-queue mutex condition queue closed threads) + ;; #(fixed-pool workers max-queue mutex condition front back count closed + ;; threads) + ;; Two-list amortized FIFO: submit conses onto BACK in O(1); take pops from + ;; FRONT, refilling it by reversing BACK once drained. COUNT tracks the queued + ;; job total so the bounded-queue admission check is O(1), not (length queue). (define (make-fixed-worker-pool workers max-queue) (unless (and (integer? workers) (> workers 0) (integer? max-queue) (> max-queue 0)) (error 'make-fixed-worker-pool "invalid worker/queue bounds")) (vector 'fixed-pool workers max-queue (make-mutex) (make-condition) - '() #f '())) + '() '() 0 #f '())) (define (p-workers p) (vector-ref p 1)) (define (p-max-queue p) (vector-ref p 2)) (define (p-mutex p) (vector-ref p 3)) (define (p-condition p) (vector-ref p 4)) - (define (p-queue p) (vector-ref p 5)) - (define (p-queue-set! p x) (vector-set! p 5 x)) - (define (p-closed? p) (vector-ref p 6)) - (define (p-closed-set! p x) (vector-set! p 6 x)) + (define (p-front p) (vector-ref p 5)) + (define (p-front-set! p x) (vector-set! p 5 x)) + (define (p-back p) (vector-ref p 6)) + (define (p-back-set! p x) (vector-set! p 6 x)) + (define (p-count p) (vector-ref p 7)) + (define (p-count-set! p x) (vector-set! p 7 x)) + (define (p-closed? p) (vector-ref p 8)) + (define (p-closed-set! p x) (vector-set! p 8 x)) (define (pool-take! pool) (mutex-acquire (p-mutex pool)) (let loop () (cond - [(pair? (p-queue pool)) - (let ([job (car (p-queue pool))]) - (p-queue-set! pool (cdr (p-queue pool))) + [(pair? (p-front pool)) + (let ([job (car (p-front pool))]) + (p-front-set! pool (cdr (p-front pool))) + (p-count-set! pool (- (p-count pool) 1)) (mutex-release (p-mutex pool)) job)] + [(pair? (p-back pool)) + (p-front-set! pool (reverse (p-back pool))) + (p-back-set! pool '()) + (loop)] [(p-closed? pool) (mutex-release (p-mutex pool)) #f] @@ -176,7 +189,7 @@ (define (fixed-worker-pool-start! pool) (let loop ([n (p-workers pool)] [threads '()]) (if (= n 0) - (begin (vector-set! pool 7 threads) pool) + (begin (vector-set! pool 9 threads) pool) (loop (- n 1) (cons (fork-thread (lambda () @@ -190,9 +203,10 @@ (define (fixed-worker-pool-submit! pool job) (mutex-acquire (p-mutex pool)) (let ([accepted (and (not (p-closed? pool)) - (< (length (p-queue pool)) (p-max-queue pool)))]) + (< (p-count pool) (p-max-queue pool)))]) (when accepted - (p-queue-set! pool (append (p-queue pool) (list job))) + (p-back-set! pool (cons job (p-back pool))) + (p-count-set! pool (+ (p-count pool) 1)) (condition-signal (p-condition pool))) (mutex-release (p-mutex pool)) accepted)) --- a/tests/limits-test.ss +++ b/tests/limits-test.ss @@ -39,7 +39,44 @@ #f) #t) +;; Worker pool: bounded admission is O(1) (count-tracked) and the FIFO order +;; survives the cons+reverse two-list queue. +(let ([pool (make-fixed-worker-pool 1 4)]) + (check "queue slot 1 accepted" (fixed-worker-pool-submit! pool (lambda () (void))) #t) + (check "queue slot 2 accepted" (fixed-worker-pool-submit! pool (lambda () (void))) #t) + (check "queue slot 3 accepted" (fixed-worker-pool-submit! pool (lambda () (void))) #t) + (check "queue slot 4 accepted" (fixed-worker-pool-submit! pool (lambda () (void))) #t) + (check "queue full rejects fifth submit" + (fixed-worker-pool-submit! pool (lambda () (void))) #f) + (fixed-worker-pool-stop! pool)) + +(let ([pool (make-fixed-worker-pool 1 8)] + [mtx (make-mutex)] + [order (box '())] + [n 5]) + (fixed-worker-pool-start! pool) + (let submit-loop ([i 0]) + (when (< i n) + (fixed-worker-pool-submit! pool + (let ([id i]) + (lambda () + (mutex-acquire mtx) + (set-box! order (cons id (unbox order))) + (mutex-release mtx)))) + (submit-loop (+ i 1)))) + (let wait-loop ([tries 0]) + (mutex-acquire mtx) + (let ([seen (length (unbox order))]) + (mutex-release mtx) + (cond + [(>= seen n) (void)] + [(>= tries 500) (void)] + [else (sleep (make-time 'time-duration 10000000 0)) + (wait-loop (+ tries 1))]))) + (check "worker pool drains jobs FIFO" (reverse (unbox order)) '(0 1 2 3 4)) + (fixed-worker-pool-stop! pool)) + (if (= failures 0) - (begin (display "limits-test: all admission checks passed") (newline)) + (begin (display "limits-test: all admission and worker-pool checks passed") (newline)) (begin (display (format "limits-test: ~a failures" failures)) (newline) (exit 1)))