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 Ursetto 
Date: 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

|#