IMAP: allow only one sync at a time, divided header fetching.

* CheckMailboxesCommand, and FetchHeadersCommand now inherit from SyncCommand
  which will prevent new CheckMailboxesCommand to be enqueued.
* FetchHeadersCommand now only retrieves up to kMaxFetchEntries headers at
  once. This gets the same TODO about scaling as the limit in
  CheckMailboxesCommand when fetching the flags/UIDs. Since we already read
  all new UIDs at that point, we could easily do better there, though.
This commit is contained in:
Axel Dörfler
2015-01-06 15:26:10 +01:00
parent aebdd0c14b
commit 29e5da6f20
2 changed files with 86 additions and 11 deletions
@@ -64,6 +64,11 @@ public:
return fWorker._EnqueueCommand(command); return fWorker._EnqueueCommand(command);
} }
void SyncCommandDone()
{
fWorker._SyncCommandDone();
}
void Quit() void Quit()
{ {
fWorker.fStopped = true; fWorker.fStopped = true;
@@ -84,6 +89,10 @@ public:
}; };
class SyncCommand : public WorkerCommand {
};
class QuitCommand : public WorkerCommand { class QuitCommand : public WorkerCommand {
public: public:
QuitCommand() QuitCommand()
@@ -112,7 +121,7 @@ public:
}; };
class FetchHeadersCommand : public WorkerCommand, public IMAP::FetchListener { class FetchHeadersCommand : public SyncCommand, public IMAP::FetchListener {
public: public:
FetchHeadersCommand(IMAPFolder& folder, IMAPMailbox& mailbox, FetchHeadersCommand(IMAPFolder& folder, IMAPMailbox& mailbox,
uint32 from, uint32 to) uint32 from, uint32 to)
@@ -133,9 +142,16 @@ public:
if (status != B_OK) if (status != B_OK)
return status; return status;
// TODO: this does not scale that well. Over time, the holes in the
// UIDs might become really large
uint32 to = fTo;
if (to - fFrom >= kMaxFetchEntries)
to = fFrom + kMaxFetchEntries - 1;
// TODO: trigger download of mails for all messages below the // TODO: trigger download of mails for all messages below the
// body fetch limit // body fetch limit
IMAP::FetchCommand fetch(fFrom, fTo, printf("IMAP: fetch headers from %lu to %lu\n", fFrom, to);
IMAP::FetchCommand fetch(fFrom, to,
IMAP::kFetchHeader | IMAP::kFetchFlags); IMAP::kFetchHeader | IMAP::kFetchFlags);
fetch.SetListener(this); fetch.SetListener(this);
@@ -143,9 +159,15 @@ public:
if (status != B_OK) if (status != B_OK)
return status; return status;
fFrom = to + 1;
return B_OK; return B_OK;
} }
virtual bool IsDone() const
{
return fFrom >= fTo;
}
virtual bool FetchData(uint32 fetchFlags, BDataIO& stream, size_t& length) virtual bool FetchData(uint32 fetchFlags, BDataIO& stream, size_t& length)
{ {
fFetchStatus = fFolder.StoreMessage(fFile, fetchFlags, stream, fFetchStatus = fFolder.StoreMessage(fFile, fetchFlags, stream,
@@ -170,10 +192,11 @@ private:
}; };
class CheckMailboxesCommand : public WorkerCommand { class CheckMailboxesCommand : public SyncCommand {
public: public:
CheckMailboxesCommand() CheckMailboxesCommand(IMAPConnectionWorker& worker)
: :
fWorker(worker),
fFolders(5, false), fFolders(5, false),
fState(INIT), fState(INIT),
fFolder(NULL), fFolder(NULL),
@@ -188,8 +211,10 @@ public:
if (fState == INIT) { if (fState == INIT) {
// Collect folders // Collect folders
status_t status = WorkerPrivate(worker).AddFolders(fFolders); status_t status = WorkerPrivate(worker).AddFolders(fFolders);
if (status != B_OK || fFolders.IsEmpty()) if (status != B_OK || fFolders.IsEmpty()) {
fState = DONE;
return status; return status;
}
fState = SELECT; fState = SELECT;
} }
@@ -228,8 +253,8 @@ public:
// UIDs might become really large // UIDs might become really large
uint32 from = fLastUID; uint32 from = fLastUID;
uint32 to = fNextUID; uint32 to = fNextUID;
if (to - from > kMaxFetchEntries) if (to - from >= kMaxFetchEntries)
to = from + kMaxFetchEntries; to = from + kMaxFetchEntries - 1;
printf("IMAP: get entries from %lu to %lu\n", from, to); printf("IMAP: get entries from %lu to %lu\n", from, to);
// TODO: we don't really need the flags at this point at all // TODO: we don't really need the flags at this point at all
@@ -248,13 +273,13 @@ public:
fTotalEntries += entries.size(); fTotalEntries += entries.size();
fMailboxEntries += entries.size(); fMailboxEntries += entries.size();
fLastUID = to; fLastUID = to + 1;
if (to == fNextUID) { if (to == fNextUID) {
if (fMailboxEntries > 0) { if (fMailboxEntries > 0) {
// Add pending command to fetch the message headers // Add pending command to fetch the message headers
WorkerCommand* command = new FetchHeadersCommand(*fFolder, WorkerCommand* command = new FetchHeadersCommand(*fFolder,
*fMailbox, fFirstUID, fLastUID); *fMailbox, fFirstUID, fNextUID);
if (!fFetchCommands.AddItem(command)) if (!fFetchCommands.AddItem(command))
delete command; delete command;
} }
@@ -278,6 +303,7 @@ private:
DONE DONE
}; };
IMAPConnectionWorker& fWorker;
BObjectList<IMAPFolder> fFolders; BObjectList<IMAPFolder> fFolders;
State fState; State fState;
IMAPFolder* fFolder; IMAPFolder* fFolder;
@@ -292,6 +318,38 @@ private:
}; };
struct CommandDelete
{
inline void operator()(WorkerCommand* command)
{
delete command;
}
};
/*! An auto deleter similar to ObjectDeleter that called SyncCommandDone()
for all SyncCommands.
*/
struct CommandDeleter : BPrivate::AutoDeleter<WorkerCommand, CommandDelete>
{
CommandDeleter(IMAPConnectionWorker& worker, WorkerCommand* command)
:
BPrivate::AutoDeleter<WorkerCommand, CommandDelete>(command),
fWorker(worker)
{
}
~CommandDeleter()
{
if (dynamic_cast<SyncCommand*>(fObject) != NULL)
WorkerPrivate(fWorker).SyncCommandDone();
}
private:
IMAPConnectionWorker& fWorker;
};
// #pragma mark - // #pragma mark -
@@ -400,8 +458,13 @@ IMAPConnectionWorker::EnqueueCheckSubscribedFolders()
status_t status_t
IMAPConnectionWorker::EnqueueCheckMailboxes() IMAPConnectionWorker::EnqueueCheckMailboxes()
{ {
// Do not schedule checking mailboxes again if we're still working on
// those.
if (fSyncPending > 0)
return B_OK;
printf("IMAP: worker %p: enqueue check mailboxes\n", this); printf("IMAP: worker %p: enqueue check mailboxes\n", this);
return _EnqueueCommand(new CheckMailboxesCommand()); return _EnqueueCommand(new CheckMailboxesCommand(*this));
} }
@@ -450,7 +513,7 @@ IMAPConnectionWorker::_Worker()
if (command == NULL) if (command == NULL)
continue; continue;
ObjectDeleter<WorkerCommand> deleter(command); CommandDeleter deleter(*this, command);
status_t status = _Connect(); status_t status = _Connect();
if (status != B_OK) if (status != B_OK)
@@ -484,6 +547,9 @@ IMAPConnectionWorker::_EnqueueCommand(WorkerCommand* command)
return B_NO_MEMORY; return B_NO_MEMORY;
} }
if (dynamic_cast<SyncCommand*>(command) != NULL)
fSyncPending++;
locker.Unlock(); locker.Unlock();
release_sem(fPendingCommandsSemaphore); release_sem(fPendingCommandsSemaphore);
return B_OK; return B_OK;
@@ -534,6 +600,13 @@ IMAPConnectionWorker::_MailboxFor(IMAPFolder& folder)
} }
void
IMAPConnectionWorker::_SyncCommandDone()
{
fSyncPending--;
}
status_t status_t
IMAPConnectionWorker::_Connect() IMAPConnectionWorker::_Connect()
{ {
@@ -61,6 +61,7 @@ private:
uint32* _nextUID); uint32* _nextUID);
IMAPMailbox* _MailboxFor(IMAPFolder& folder); IMAPMailbox* _MailboxFor(IMAPFolder& folder);
IMAPFolder* _Selected() const { return fSelectedBox; } IMAPFolder* _Selected() const { return fSelectedBox; }
void _SyncCommandDone();
status_t _Connect(); status_t _Connect();
void _Disconnect(); void _Disconnect();
@@ -76,6 +77,7 @@ private:
IMAP::Protocol fProtocol; IMAP::Protocol fProtocol;
sem_id fPendingCommandsSemaphore; sem_id fPendingCommandsSemaphore;
WorkerCommandList fPendingCommands; WorkerCommandList fPendingCommands;
uint32 fSyncPending;
IMAP::ExistsHandler fExistsHandler; IMAP::ExistsHandler fExistsHandler;
IMAP::ExpungeHandler fExpungeHandler; IMAP::ExpungeHandler fExpungeHandler;