Welcome to the CHICKEN Scheme pasting service
first stab at polling pasted by zbigniew the justified on Mon Jul 11 09:18:47 2011
From 5edb8a1c3094e5804507375e345123acfff325d1 Mon Sep 17 00:00:00 2001 From: Jim UrsettoDate: Mon, 11 Jul 2011 02:17:22 -0500 Subject: [PATCH] revamp polling with opaque poll item object --- zmq.scm | 156 ++++++++++++++++++++++++++++++++------------------------------- 1 files changed, 80 insertions(+), 76 deletions(-) diff --git a/zmq.scm b/zmq.scm index c226478..891c477 100644 --- a/zmq.scm +++ b/zmq.scm @@ -5,8 +5,11 @@ make-socket socket? close-socket bind-socket connect-socket socket-option-set! socket-option socket-fd send-message receive-message receive-message* - make-poll-item poll poll-item-socket - poll-item-fd poll-item-in? poll-item-out? poll-item-error?) + poll poll-items + for-each-poll-item + ;; make-poll-item poll-item-socket + ;; poll-item-fd poll-item-in? poll-item-out? poll-item-error? + ) (import (except chicken errno) scheme foreign data-structures) (use lolevel foreigners srfi-1 srfi-18 srfi-13) @@ -351,79 +354,80 @@ ;; polling -(define %make-poll-item make-poll-item) - -(define (make-poll-item socket/fd #!key in out) - (let ((item (%make-poll-item (make-foreign-poll-item) - (and (socket? socket/fd) socket/fd) - in out))) - (if (socket? socket/fd) - (%poll-item-socket-set! (poll-item-pointer item) (socket-pointer socket/fd)) - (%poll-item-fd-set! (poll-item-pointer item) socket/fd)) - - (%poll-item-events-set! (poll-item-pointer item) - (bitwise-ior (if in zmq/pollin 0) - (if out zmq/pollout 0))) - - (%poll-item-revents-set! (poll-item-pointer item) 0) - - (set-finalizer! item (lambda (i) - (free-foreign-poll-item (poll-item-pointer i)))))) - -(define (poll-item-fd item) - (%poll-item-fd (poll-item-pointer item))) - -(define (poll-item-revents item) - (%poll-item-revents (poll-item-pointer item))) - -(define (poll-item-in? item) - (not (zero? (bitwise-and zmq/pollin (poll-item-revents item))))) - -(define (poll-item-out? item) - (not (zero? (bitwise-and zmq/pollout (poll-item-revents item))))) - -(define (poll-item-error? item) - (not (zero? (bitwise-and zmq/pollerr (poll-item-revents item))))) - +(define _zmq_pollitem_size (foreign-value "sizeof(zmq_pollitem_t)" int)) +(define %poll-items-init/socket + (foreign-lambda* void ((scheme-pointer p) (int i) (socket s) (short events)) + "zmq_pollitem_t *pi = (zmq_pollitem_t *)p+i;" + "pi->socket = s; pi->fd = 0; pi->events = events; pi->revents = 0;")) +(define %poll-items-init/fd + (foreign-lambda* void ((scheme-pointer p) (int i) (int fd) (short events)) + "zmq_pollitem_t *pi = (zmq_pollitem_t *)p+i;" + "pi->socket = 0; pi->fd = fd; pi->events = events; pi->revents = 0;")) +(define %poll-items-in? + (foreign-lambda* bool ((scheme-pointer p) (int i)) + "return(((zmq_pollitem_t *)p+i)->revents & ZMQ_POLLIN);")) +(define %poll-items-out? + (foreign-lambda* bool ((scheme-pointer p) (int i)) + "return(((zmq_pollitem_t *)p+i)->revents & ZMQ_POLLOUT);")) (define %poll-sockets - (foreign-safe-lambda* int - ((scheme-object poll_item_ref) - (unsigned-int length) - (long timeout)) - "zmq_pollitem_t items[length]; - zmq_pollitem_t *item_ptrs[length]; - int i; - - for (i = 0; i < length; i++) { - C_save(C_fix(i)); - item_ptrs[i] = (zmq_pollitem_t *)C_pointer_address(C_callback(poll_item_ref, 1)); - } - - for (i = 0; i < length; i++) { - items[i] = *item_ptrs[i]; - } - - int rc = zmq_poll(items, length, timeout); - - if (rc != -1) { - for (i = 0; i < length; i++) { - (*item_ptrs[i]).revents = items[i].revents; - } - } - - C_return(rc);")) - -(define (poll poll-items timeout/block) - (if (null? poll-items) - (error 'poll "null list passed for poll-items") - (let ((result (%poll-sockets (lambda (i) - (poll-item-pointer (list-ref poll-items i))) - (length poll-items) - (case timeout/block - ((#f) 0) - ((#t) -1) - (else timeout/block))))) - (if (= result -1) - (zmq-error 'poll) - result)))) + (foreign-lambda* int ((scheme-pointer p) + (unsigned-int length) (long timeout)) + "return(zmq_poll((zmq_pollitem_t *)p, length, timeout));")) + +(define-record poll-items store sockets) ;; blob, vector +(define (poll-items-length items) + (vector-length (poll-items-sockets items))) +(define (poll-items-socket items i) ;; can also be used for fds + (vector-ref (poll-items-sockets items) i)) +(define (poll-items-in? items i) + (when (or (< i 0) (>= i (poll-items-length items))) + (error 'poll-items-out? "index out of range" i items)) + (%poll-items-in? (poll-items-store items) i)) +(define (poll-items-out? items i) + (when (or (< i 0) (>= i (poll-items-length items))) + (error 'poll-items-out? "index out of range" i items)) + (%poll-items-out? (poll-items-store items) i)) + +(define (poll-items in out) + (let* ((ilen (length in)) + (olen (length out)) + (len (+ ilen olen))) + (when (= 0 len) + (error 'poll-items "no items to poll")) + (let ((store (make-blob (* _zmq_pollitem_size len)))) + (let loop ((i 0) (in in)) + (unless (null? in) + (let ((s (car in))) + (if (socket? s) + (%poll-items-init/socket store i (socket-pointer s) zmq/pollin) + (%poll-items-init/fd store i s zmq/pollin))) + (loop (+ i 1) (cdr in)))) + (let loop ((i ilen) (out out)) + (unless (null? out) + (let ((s (car out))) + (if (socket? s) + (%poll-items-init/socket store i (socket-pointer s) zmq/pollout) + (%poll-items-init/fd store i s zmq/pollout))) + (loop (+ i 1) (cdr out)))) + (make-poll-items store (list->vector (append in out)))))) + +(define (poll items timeout/block) + (let ((result (%poll-sockets (poll-items-store items) + (poll-items-length items) + (case timeout/block + ((#f) 0) + ((#t) -1) + (else timeout/block))))) + (if (= result -1) + (zmq-error 'poll) + result))) + +(define (for-each-poll-item items in out) + (do ((i 0 (+ i 1))) + ((= i (poll-items-length items))) + (cond ((poll-items-in? items i) + (in (poll-items-socket items i))) + ((poll-items-out? items i) + (out (poll-items-socket items i)))))) + ) -- 1.7.4.1
example poll client (mspoller) pasted by zbigniew the ancient on Mon Jul 11 09:20:53 2011
(use zmq srfi-18) (define wx (make-socket 'sub)) (define vent (make-socket 'pull)) (socket-option-set! wx 'subscribe "60047 ") (connect-socket wx "tcp://localhost:5556") (connect-socket vent "tcp://localhost:5557") (define items (poll-items (list wx vent) '())) (let loop () (poll items #t) (for-each-poll-item items (lambda (s) (print "received msg: " (receive-message s))) identity) (thread-sleep! 0.25) ;; just in case (actually to verify that messages interleave) (loop)) #| received msg: 60047 118 19 received msg: 0 received msg: 74 received msg: 60047 17 51 received msg: 81 received msg: 60047 112 18 received msg: 96 received msg: 83 received msg: 60047 131 29 received msg: 73 received msg: 88 ... |#
alternative poll client pasted by zbigniew the gravid on Mon Jul 11 21:46:48 2011
(use zmq srfi-18) (define wx (make-socket 'sub)) (define vent (make-socket 'pull)) (socket-option-set! wx 'subscribe "60047 ") (connect-socket wx "tcp://localhost:5556") (connect-socket vent "tcp://localhost:5557") (define twx (make-thread (lambda () (let loop () (print "received weather msg: " (receive-message* wx)) (loop))))) (define tvent (make-thread (lambda () (receive-message* vent) ;; sync with start of batch (let loop () (let ((msg (receive-message* vent))) (print "received worker msg: " msg) (thread-sleep! (/ (string->number msg) 1000)) (loop)))))) (thread-start! twx) (thread-start! tvent) (thread-join! twx) (thread-join! tvent) #| received weather msg: 60047 7 51 received weather msg: 60047 109 54 received weather msg: 60047 25 56 received weather msg: 60047 100 42 received weather msg: 60047 60 11 received worker msg: 35 received worker msg: 35 received worker msg: 63 received worker msg: 39 received worker msg: 48 received worker msg: 75 received worker msg: 91 received weather msg: 60047 91 52 received weather msg: 60047 115 54 received worker msg: 4 received worker msg: 37 received worker msg: 68 received worker msg: 59 received worker msg: 84 ... |#
task worker w/ control channel (taskwork2) pasted by zbigniew the gritty on Fri Jul 15 00:05:59 2011
(use zmq srfi-18) (define ctrl (make-socket 'sub)) (define vent (make-socket 'pull)) (define sink (make-socket 'push)) (socket-option-set! ctrl 'subscribe "KILL") (connect-socket ctrl "tcp://localhost:5559") (connect-socket vent "tcp://localhost:5557") (connect-socket sink "tcp://localhost:5558") (define tctrl (make-thread (lambda () (receive-message* ctrl)))) (define tvent (make-thread (lambda () (let loop ((n 0)) ;; include zero msg (let ((msg (receive-message* vent))) (display ".") (flush-output) (thread-sleep! (/ (string->number msg) 1000)) (send-message sink "") (loop (+ n 1))))))) (thread-start! tctrl) (thread-start! tvent) (thread-join! tctrl) ;; finish once we receive KILL msg on ctrl
router broker (rrbroker) pasted by zbigniew interrupted on Fri Jul 15 00:07:32 2011
(use zmq srfi-18) (define front (make-socket 'xrep)) ;; FIXME: should be 'router (define back (make-socket 'xreq)) ;; FIXME: should be 'dealer (bind-socket front "tcp://*:5559") (bind-socket back "tcp://*:5560") (define-syntax forever (syntax-rules () ((forever e0 e1 ...) (let loop () e0 e1 ... (loop))))) (define (route) (forever (let more () (let ((msg (receive-message* front)) (more? (socket-option front 'rcvmore))) (send-message back msg send-more: more?) (when more? (more)))))) (define (deal) (forever (let more () (let ((msg (receive-message* back)) (more? (socket-option back 'rcvmore))) (send-message front msg send-more: more?) (when more? (more)))))) (let ((rt (make-thread route)) (dt (make-thread deal))) (thread-start! rt) (thread-start! dt) (thread-join! rt) (thread-join! dt))
router-dealer custom routing example (rtdealer) added by zbigniew the dominar on Fri Jul 15 00:19:57 2011
(use zmq random-bsd srfi-18) ;; Custom routing Router to Dealer (ROUTER to DEALER) ;; While this example runs in a single process, that is just to make ;; it easier to start and stop the example. Each thread has its own ;; context and conceptually acts as a separate process. (define (make-worker id) ;; ID: Socket identity (string) (lambda () (let* ((context (make-context 1)) (worker (make-socket 'xreq context))) ;; FIXME: 'dealer (socket-option-set! worker 'identity id) (connect-socket worker "ipc://routing.ipc") (let loop ((total 0)) (if (string=? (receive-message* worker) "END") (print id " received: " total) (loop (+ total 1)))) (close-socket worker) (terminate-context context)))) ;;; main (define client (make-socket 'xrep)) ;; FIXME: 'router (bind-socket client "ipc://routing.ipc") (define A (thread-start! (make-worker "A"))) (define B (thread-start! (make-worker "B"))) ;; Wait for threads to connect, since otherwise the messages ;; we send won't be routable. (thread-sleep! 1) (do ((i 0 (+ i 1))) ((>= i 1000)) ;; Send 1000 tasks scattered to A twice as often as B (send-message client (if (> (random 3) 0) "A" "B") send-more: #t) (send-message client "This is the workload")) (send-message client "A" send-more: #t) (send-message client "END") (send-message client "B" send-more: #t) (send-message client "END") ;;; finish (close-socket client) (terminate-context (zmq-default-context)) (for-each thread-join! (list A B)) #| $ ./rtdealer B received: 331 A received: 669 $ ./rtdealer B received: 324 A received: 676 $ ./rtdealer B received: 343 A received: 657 |#