offlineimap rev 567

"Automatic Subversion Change Mailer" <[email protected]> Thu, 14 Aug 2003 13:33:24 -0500 (CDT)
Newsgroups gmane.mail.imap.offlineimap.subversion
Message-ID <[email protected]>
You are receiving this message because
all commits get sent to this address.

Author: jgoerzen
Date: 2003-08-14 13:33:18 -0500 (Thu, 14 Aug 2003)
New Revision: 567

Modified:
  offlineimap/branches/twisted/offlineimap/accounts.py
  offlineimap/branches/twisted/offlineimap/imapserver.py
  offlineimap/branches/twisted/offlineimap/repository/Base.py
  offlineimap/branches/twisted/offlineimap/repository/IMAP.py

Log:
Keepalive again appears to work (woohoo)


Diff:
Modified: offlineimap/branches/twisted/offlineimap/accounts.py
===================================================================
--- offlineimap/branches/twisted/offlineimap/accounts.py	2003-08-14 14:51:05 UTC (rev 566)
+++ offlineimap/branches/twisted/offlineimap/accounts.py	2003-08-14 18:33:18 UTC (rev 567)
@@ -61,6 +61,16 @@
     def getsection(self):
         return 'Account ' + self.getname()
 
+    def _getkaobjs(self):
+        kaobjs = []
+
+        if hasattr(self, 'localrepos'):
+            kaobjs.append(self.localrepos)
+        if hasattr(self, 'remoterepos'):
+            kaobjs.append(self.remoterepos)
+
+        return kaobjs
+
     def sleeper(self, maindefer = None):
         """Sleep handler.  Returns same value as UIBase.sleep:
         0 if timeout expired, 1 if there was a request to cancel the timer,
@@ -74,17 +84,11 @@
         if not self.refreshperiod:
             return rd(100)
 
-        # FIXME: handle keepalive
-        #kaobjs = []
+        kaobjs = self._getkaobjs()
 
-        #if hasattr(self, 'localrepos'):
-        #    kaobjs.append(self.localrepos)
-        #if hasattr(self, 'remoterepos'):
-        #    kaobjs.append(self.remoterepos)
+        for item in kaobjs:
+            item.startkeepalive()
 
-        #for item in kaobjs:
-        #    item.startkeepalive()
-
         firstcall = 0
         if maindefer == None:
             firstcall = 1
@@ -98,15 +102,16 @@
             return maindefer
 
     def sleeper_result(self, sleepresult, maindefer):
-        # FIXME: handle keepalive
-        #if sleepresult == 2:
-        #    # Cancel keep-alive, but don't bother terminating threads
-        #    for item in kaobjs:
-        #        item.stopkeepalive(abrupt = 1)
-        #else:
-        #    # Cancel keep-alive and wait for thread to terminate.
-        #    for item in kaobjs:
-        #        item.stopkeepalive(abrupt = 0)
+        kaobjs = self._getkaobjs()
+        deferreds = []
+        for item in kaobjs:
+            deferreds.append(item.stopkeepalive())
+        dl = defer.DeferredList(deferreds)
+        dl.addCallback(self._sleeper_result_2, sleepresult, maindefer)
+        return dl
+
+    def _sleeper_result_2(self, ignore, sleepresult, maindefer):
+        """Called after keepalives have all died."""
         print "sleepdone"
         if sleepresult == 2:
             self.acctdone(None)

Modified: offlineimap/branches/twisted/offlineimap/imapserver.py
===================================================================
--- offlineimap/branches/twisted/offlineimap/imapserver.py	2003-08-14 14:51:05 UTC (rev 566)
+++ offlineimap/branches/twisted/offlineimap/imapserver.py	2003-08-14 18:33:18 UTC (rev 567)
@@ -338,54 +338,45 @@
         self.lastowner = {}
         self.connectionlock.release()
 
-    def keepalive(self, timeout, event):
-        """Sends a NOOP to each connection recorded.   It will wait a maximum
-        of timeout seconds between doing this, and will continue to do so
-        until the Event object as passed is true.  This method is expected
-        to be invoked in a separate thread, which should be join()'d after
-        the event is set."""
+    def keepalive(self):
+        """Sends a NOOP to each connection recorded.  This will be run just
+        once.  Returns a deferred that gets callbacked when
+        all NOOPs are complete."""
         ui = UIBase.getglobalui()
-        ui.debug('imap', 'keepalive thread started')
-        while 1:
-            ui.debug('imap', 'keepalive: top of loop')
-            event.wait(timeout)
-            ui.debug('imap', 'keepalive: after wait')
-            if event.isSet():
-                ui.debug('imap', 'keepalive: event is set; exiting')
-                return
-            ui.debug('imap', 'keepalive: acquiring connectionlock')
-            self.connectionlock.acquire()
-            numconnections = len(self.assignedconnections) + \
-                             len(self.availableconnections)
-            self.connectionlock.release()
-            ui.debug('imap', 'keepalive: connectionlock released')
-            threads = []
-            imapobjs = []
-        
-            for i in range(numconnections):
-                ui.debug('imap', 'keepalive: processing connection %d of %d' % (i, numconnections))
-                imapobj = self.acquireconnection()
-                ui.debug('imap', 'keepalive: connection %d acquired' % i)
-                imapobjs.append(imapobj)
-                thr = threadutil.ExitNotifyThread(target = imapobj.noop)
-                thr.setDaemon(1)
-                thr.start()
-                threads.append(thr)
-                ui.debug('imap', 'keepalive: thread started')
+        ui.debug('imap', 'keepalive process started')
+        self.kanumconnections = self._getConnectionCount()
+        ui.debug('imap', 'keepalive found %d connections' %
+                 self.kanumconnections)
+        #self.kafinished = defer.Deferred()
+        deferreds = []
+        self.kaconnections = []
+        self.kanumdone = 0
 
-            ui.debug('imap', 'keepalive: joining threads')
+        def noopdone(ignore, i):
+            ui.debug('imap', 'keepalive: noop finished on connection %d' % i)
+            self.kanumdone += 1
 
-            for thr in threads:
-                # Make sure all the commands have completed.
-                thr.join()
+            if self.kanumdone == self.kanumconnections:
+                ui.debug('imap', 'keepalive: noop all finished')
+                for connection in self.kaconnections:
+                    self._releaseConnection(None, connection)
+                #self.kafinished.callback(None)
 
-            ui.debug('imap', 'keepalive: releasing connections')
+        def sendnoop(imapobj, i):
+            ui.debug('imap', 'keepalive: sending noop on connection %d' % i)
+            self.kaconnections.append(imapobj)
+            d = imapobj.noop()
+            d.addCallback(noopdone, i)
+        
+        for i in range(self.kanumconnections):
+            ui.debug('imap', 'keepalive: processing connection %d of %d' % (i,
+                                                                            self.kanumconnections))
+            d = self._grabConnection()
+            d.addCallback(sendnoop, i)
+            deferreds.append(d)
+        dl = defer.DeferredList(deferreds)
+        return dl
 
-            for imapobj in imapobjs:
-                self.releaseconnection(imapobj)
-
-            ui.debug('imap', 'keepalive: bottom of loop')
-
 class ConfigedIMAPServer(IMAPServer):
     """This class is designed for easier initialization given a ConfigParser
     object and an account name.  The passwordhash is used if

Modified: offlineimap/branches/twisted/offlineimap/repository/Base.py
===================================================================
--- offlineimap/branches/twisted/offlineimap/repository/Base.py	2003-08-14 14:51:05 UTC (rev 566)
+++ offlineimap/branches/twisted/offlineimap/repository/Base.py	2003-08-14 18:33:18 UTC (rev 567)
@@ -19,6 +19,7 @@
 from offlineimap import CustomConfig
 import os.path
 from twisted.internet import defer
+from offlineimap.imaputil import rd
 
 def LoadRepository(name, account, reqtype):
     from offlineimap.repository.IMAP import IMAPRepository, MappedIMAPRepository
@@ -150,5 +151,5 @@
     def stopkeepalive(self, abrupt = 0):
         """Stop keep alive.  If abrupt is 1, stop it but don't bother waiting
         for the threads to terminate."""
-        pass
+        return rd(None)
     

Modified: offlineimap/branches/twisted/offlineimap/repository/IMAP.py
===================================================================
--- offlineimap/branches/twisted/offlineimap/repository/IMAP.py	2003-08-14 14:51:05 UTC (rev 566)
+++ offlineimap/branches/twisted/offlineimap/repository/IMAP.py	2003-08-14 18:33:18 UTC (rev 567)
@@ -21,6 +21,7 @@
 from offlineimap.folder.UIDMaps import MappedIMAPFolder
 import re, types, os
 from offlineimap.imaputil import rd
+from twisted.internet import defer
 
 class IMAPRepository(BaseRepository, imaputil.AcquireMixin):
     def __init__(self, reposname, account):
@@ -48,26 +49,63 @@
                                              {'re': re})
 
     def startkeepalive(self):
+        """Returns nothing, and does so immediately."""
         keepalivetime = self.getkeepalive()
         if not keepalivetime: return
-        self.kaevent = Event()
-        self.kathread = ExitNotifyThread(target = self.imapserver.keepalive,
-                                         name = "Keep alive " + self.getname(),
-                                         args = (keepalivetime, self.kaevent))
-        self.kathread.setDaemon(1)
-        self.kathread.start()
+        
+        self.kacancel = 0
+        self._schedulekeepalive()
+        
+    def _schedulekeepalive(self):
+        keepalivetime = self.getkeepalive()
+        if not keepalivetime: return
+        if self.kacancel:
+            return
 
-    def stopkeepalive(self, abrupt = 0):
-        if not hasattr(self, 'kaevent'):
-            # Keepalive is not active.
+        from twisted.internet import reactor
+        self.katimer = reactor.callLater(keepalivetime,
+                                       self._runkeepalive)
+
+    def _runkeepalive(self):
+        if self.kacancel:
             return
+        self.kadefer = self.imapserver.keepalive()
+        self.kadefer.addCallback(lambda x: self._schedulekeepalive())
 
-        self.kaevent.set()
-        if not abrupt:
-            self.kathread.join()
-        del self.kathread
-        del self.kaevent
+    def stopkeepalive(self):
+        """Returns a deferred, whose callback will not be called until
+        we know that no keepalive is running."""
+        if not hasattr(self, 'kacancel'):
+            # Keepalive is not active.
+            return rd(None)
 
+        self.kacancel = 1
+        try:
+            self.katimer.cancel()
+        except:
+            # It's already fired (grr)
+            if hasattr(self, 'kadefer'):
+                self.kadefer.addCallback(lambda x: self._keepaliveisdone())
+                return self.kadefer
+            # If we don't have kadefer, that means that it fired but
+            # got canceled in time.
+
+        return rd(self._keepaliveisdone())
+
+    def _keepaliveisdone(self):
+        try:
+            del self.kadefer
+        except:
+            pass
+        try:
+            del self.katimer
+        except:
+            pass
+        try:
+            del self.kacancel
+        except:
+            pass
+
     def holdordropconnections(self):
         if not self.getholdconnectionopen():
             self.dropconnections()