SVN: ZODB/trunk/src/ZEO/ The storage server is now multi-threaded.

Jim Fulton <[email protected]>
Newsgroups gmane.comp.python.zope.zodb.cvs
Message-ID <[email protected]>
Log message for revision 108694:
  The storage server is now multi-threaded.
  

Changed:
  U   ZODB/trunk/src/ZEO/StorageServer.py
  U   ZODB/trunk/src/ZEO/zrpc/connection.py

-=-
Modified: ZODB/trunk/src/ZEO/StorageServer.py
===================================================================
--- ZODB/trunk/src/ZEO/StorageServer.py	2010-02-01 16:17:32 UTC (rev 108693)
+++ ZODB/trunk/src/ZEO/StorageServer.py	2010-02-01 19:12:15 UTC (rev 108694)
@@ -1346,7 +1346,7 @@
         self.rpc.callAsync('endVerify')
 
     def invalidateTransaction(self, tid, args):
-        self.rpc.callAsyncNoPoll('invalidateTransaction', tid, args)
+        self.rpc.callAsync('invalidateTransaction', tid, args)
 
     def serialnos(self, arg):
         self.rpc.callAsyncNoPoll('serialnos', arg)
@@ -1372,11 +1372,11 @@
 class ClientStub308(ClientStub):
 
     def invalidateTransaction(self, tid, args):
-        self.rpc.callAsyncNoPoll(
-            'invalidateTransaction', tid, [(arg, '') for arg in args])
+        ClientStub.invalidateTransaction(
+            self, tid, [(arg, '') for arg in args])
 
     def invalidateVerify(self, oid):
-        self.rpc.callAsync('invalidateVerify', (oid, ''))
+        ClientStub.invalidateVerify(self, (oid, ''))
 
 class ZEOStorage308Adapter:
 

Modified: ZODB/trunk/src/ZEO/zrpc/connection.py
===================================================================
--- ZODB/trunk/src/ZEO/zrpc/connection.py	2010-02-01 16:17:32 UTC (rev 108693)
+++ ZODB/trunk/src/ZEO/zrpc/connection.py	2010-02-01 19:12:15 UTC (rev 108694)
@@ -560,15 +560,18 @@
     # Exception types that should not be logged:
     unlogged_exception_types = (ZODB.POSException.POSKeyError, )
 
-    # Servers use a shared server trigger that uses the asyncore socket map
-    trigger = ZEO.zrpc.trigger.trigger()
-    call_from_thread = trigger.pull_trigger
-
     def __init__(self, sock, addr, obj, mgr):
         self.mgr = mgr
-        Connection.__init__(self, sock, addr, obj, 'S')
+        map = {}
+        Connection.__init__(self, sock, addr, obj, 'S', map=map)
         self.marshal = ServerMarshaller()
+        self.trigger = ZEO.zrpc.trigger.trigger(map)
+        self.call_from_thread = self.trigger.pull_trigger
 
+        t = threading.Thread(target=server_loop, args=(map,))
+        t.setDaemon(True)
+        t.start()
+
     def handshake(self):
         # Send the server's preferred protocol to the client.
         self.message_output(self.current_protocol)
@@ -601,6 +604,13 @@
 
     poll = smac.SizedMessageAsyncConnection.handle_write
 
+def server_loop(map):
+    while len(map) > 1:
+        asyncore.poll(30.0, map)
+
+    for o in map.values():
+        o.close()
+
 class ManagedClientConnection(Connection):
     """Client-side Connection subclass."""
     __super_init = Connection.__init__
@@ -714,10 +724,6 @@
 
         self.trigger.pull_trigger()
 
-        # Delay used when we call asyncore.poll() directly.
-        # Start with a 1 msec delay, double until 1 sec.
-        delay = 0.001
-
         self.replies_cond.acquire()
         try:
             while 1:
lmpx.com only provides a reader for public news (NNTP) servers. It is not affiliated with the servers or forums shown here and is not responsible for the content of articles, which is written by their respective authors.