diff --git a/lib/Console/HarmonizeCommand.php b/lib/Console/HarmonizeCommand.php
index aa80f9c..6dad3d8 100644
--- a/lib/Console/HarmonizeCommand.php
+++ b/lib/Console/HarmonizeCommand.php
@@ -13,6 +13,7 @@ use KTXC\Context\TenantContext;
use KTXM\ProviderImap\Providers\Provider;
use KTXM\ProviderImap\Providers\Service;
use KTXM\ProviderImap\Service\Cache\HarmonizationService;
+use KTXM\ProviderImap\Service\Cache\MessageDeltaService;
use KTXM\ProviderImap\Stores\MailboxStore;
use KTXM\ProviderImap\Stores\MessageStore;
use Symfony\Component\Console\Attribute\AsCommand;
@@ -141,9 +142,10 @@ class HarmonizeCommand extends Command
foreach ($messages as $name => $outcome) {
if (is_string($outcome)) {
$failed++;
- $rows[] = [$name, 'failed', '', '', '', '', $outcome];
+ $rows[] = [$name, 'failed', '', '', '', '', '', $outcome];
continue;
}
+ $state = $this->mailboxStore->state($serviceId, (string) $name);
$rows[] = [
$name,
$outcome['status'],
@@ -151,10 +153,13 @@ class HarmonizeCommand extends Command
$outcome['updated'],
$outcome['removed'],
$outcome['complete'] ? 'yes' : 'no',
+ $state['uidValidity'] !== null
+ ? MessageDeltaService::signature((int) $state['uidValidity'], (int) $state['changeSeq'])
+ : '',
$outcome['reset'] ? 'UIDVALIDITY changed, cache reset' : '',
];
}
- $io->table(['Mailbox', 'Status', 'Added', 'Updated', 'Removed', 'Complete', 'Note'], $rows);
+ $io->table(['Mailbox', 'Status', 'Added', 'Updated', 'Removed', 'Complete', 'Signature', 'Note'], $rows);
$io->writeln(sprintf(
'Cache: storage/%s/provider_imap/%s, MongoDB provider_imap_mail_mailboxes / provider_imap_mail_messages',
diff --git a/lib/Service/Cache/HarmonizationService.php b/lib/Service/Cache/HarmonizationService.php
index a78b45d..44017a0 100644
--- a/lib/Service/Cache/HarmonizationService.php
+++ b/lib/Service/Cache/HarmonizationService.php
@@ -184,7 +184,8 @@ class HarmonizationService
$flags = MessageProperties::normalizeFlags($flags);
if (!self::sameFlags($flags, $cachedFlags[$uid])) {
- $this->messageStore->updateFlags($serviceId, $mailbox, $uidValidity, $uid, $flags);
+ $seq = $this->mailboxStore->reserveSequence($serviceId, $mailbox);
+ $this->messageStore->updateFlags($serviceId, $mailbox, $uidValidity, $uid, $flags, $seq);
$result['updated']++;
}
unset($cachedFlags[$uid]);
@@ -199,10 +200,11 @@ class HarmonizationService
array_push($missing, ...array_diff($batch, $ingested));
}
- // expunged: cached UIDs the server did not list; meta documents first, then files
+ // expunged: cached UIDs the server did not list; tombstones first, then files
$expunged = array_keys($cachedFlags);
if ($expunged !== []) {
- $this->messageStore->delete($serviceId, $mailbox, $uidValidity, ...$expunged);
+ $seq = $this->mailboxStore->reserveSequence($serviceId, $mailbox);
+ $this->messageStore->tombstone($serviceId, $mailbox, $uidValidity, $seq, ...$expunged);
$this->fileStore->delete($tenantId, $serviceId, $mailbox, $uidValidity, ...$expunged);
$result['removed'] = count($expunged);
}
diff --git a/lib/Service/Cache/MessageDeltaService.php b/lib/Service/Cache/MessageDeltaService.php
new file mode 100644
index 0000000..a7bb29e
--- /dev/null
+++ b/lib/Service/Cache/MessageDeltaService.php
@@ -0,0 +1,88 @@
+
+ * SPDX-License-Identifier: AGPL-3.0-or-later
+ */
+
+namespace KTXM\ProviderImap\Service\Cache;
+
+use KTXF\Resource\Delta\Delta;
+use KTXF\Resource\Delta\DeltaCollection;
+use KTXM\ProviderImap\Stores\MailboxStore;
+use KTXM\ProviderImap\Stores\MessageStore;
+
+/**
+ * Answers "what changed since this signature?" for a cached mailbox.
+ *
+ * A signature is `:`. Read-only: harmonizing before
+ * answering is the caller's decision.
+ */
+class MessageDeltaService
+{
+ public function __construct(
+ private readonly MailboxStore $mailboxStore,
+ private readonly MessageStore $messageStore,
+ ) {}
+
+ /**
+ * Changes of a mailbox since a signature.
+ *
+ * - empty signature: the current signature and no changes, so the client gets a
+ * starting point; while the initial harmonization is incomplete the signature is
+ * `:0`, so additions of the run in progress are not skipped
+ * - signature of another UIDVALIDITY, older than the purged tombstones, or not
+ * understood: a reset, answered as changes since 0 (every message is an addition)
+ * - otherwise: additions, modifications and deletions after the signature
+ *
+ * A mailbox that is not cached (or never harmonized) yields an empty delta.
+ */
+ public function delta(string $serviceId, string $mailbox, string $signature): Delta
+ {
+ $state = $this->mailboxStore->state($serviceId, $mailbox);
+ if ($state['uidValidity'] === null) {
+ return new Delta();
+ }
+
+ $uidValidity = (int) $state['uidValidity'];
+ $current = self::signature($uidValidity, (int) $state['changeSeq']);
+
+ if ($signature === '') {
+ return new Delta(signature: $state['harmonizationComplete'] ? $current : self::signature($uidValidity, 0));
+ }
+
+ $since = self::since($signature, $uidValidity, (int) $state['purgedSeq'], (int) $state['changeSeq']) ?? 0;
+ $changes = $this->messageStore->changes($serviceId, $mailbox, $uidValidity, $since);
+
+ return new Delta(
+ new DeltaCollection(array_map('strval', $changes['additions'])),
+ new DeltaCollection(array_map('strval', $changes['modifications'])),
+ new DeltaCollection(array_map('strval', $changes['deletions'])),
+ $current,
+ );
+ }
+
+ public static function signature(int $uidValidity, int $changeSeq): string
+ {
+ return $uidValidity . ':' . $changeSeq;
+ }
+
+ /**
+ * The change sequence a signature refers to, or null when it needs a reset.
+ */
+ private static function since(string $signature, int $uidValidity, int $purgedSeq, int $changeSeq): ?int
+ {
+ if (preg_match('/^(\d+):(\d+)$/', $signature, $matches) !== 1) {
+ return null;
+ }
+
+ $since = (int) $matches[2];
+ if ((int) $matches[1] !== $uidValidity || $since < $purgedSeq || $since > $changeSeq) {
+ return null;
+ }
+
+ return $since;
+ }
+}
diff --git a/lib/Service/Cache/MessageIngestor.php b/lib/Service/Cache/MessageIngestor.php
index f27e625..4bfb946 100644
--- a/lib/Service/Cache/MessageIngestor.php
+++ b/lib/Service/Cache/MessageIngestor.php
@@ -10,6 +10,7 @@ declare(strict_types=1);
namespace KTXM\ProviderImap\Service\Cache;
use KTXM\ProviderImap\Providers\EntityResource;
+use KTXM\ProviderImap\Stores\MailboxStore;
use KTXM\ProviderImap\Stores\MessageFileStore;
use KTXM\ProviderImap\Stores\MessageStore;
@@ -26,28 +27,37 @@ class MessageIngestor
public function __construct(
private readonly MessageFileStore $fileStore,
private readonly MessageStore $messageStore,
+ private readonly MailboxStore $mailboxStore,
) {}
/**
* Cache messages of one mailbox generation.
*
* The content is written before the meta document, so an interruption leaves at
- * most an orphaned file, never a meta document without content.
+ * most an orphaned file, never a meta document without content. Each message is
+ * stamped with its own change sequence, reserved per mailbox in one step.
*
* @return int[] UIDs that were cached
*/
public function ingest(string $tenantId, string $serviceId, int $uidValidity, EntityResource ...$entities): array
{
- $ingested = [];
-
+ $byMailbox = [];
foreach ($entities as $entity) {
- $mailbox = (string) $entity->collection();
- $uid = (int) $entity->identifier();
+ $byMailbox[(string) $entity->collection()][] = $entity;
+ }
- $this->fileStore->write($tenantId, $serviceId, $mailbox, $uidValidity, $uid, $entity->toCacheContent());
- $this->messageStore->upsert($tenantId, $serviceId, $uidValidity, $entity);
+ $ingested = [];
+ foreach ($byMailbox as $mailbox => $group) {
+ $seq = $this->mailboxStore->reserveSequence($serviceId, (string) $mailbox, count($group));
- $ingested[] = $uid;
+ foreach ($group as $entity) {
+ $uid = (int) $entity->identifier();
+
+ $this->fileStore->write($tenantId, $serviceId, (string) $mailbox, $uidValidity, $uid, $entity->toCacheContent());
+ $this->messageStore->upsert($tenantId, $serviceId, $uidValidity, $entity, $seq++);
+
+ $ingested[] = $uid;
+ }
}
return $ingested;
diff --git a/lib/Stores/MailboxStore.php b/lib/Stores/MailboxStore.php
index 552e911..46604d0 100644
--- a/lib/Stores/MailboxStore.php
+++ b/lib/Stores/MailboxStore.php
@@ -10,6 +10,8 @@ declare(strict_types=1);
namespace KTXM\ProviderImap\Stores;
use KTXC\Db\DataStore;
+use MongoDB\Operation\FindOneAndUpdate;
+use RuntimeException;
use KTXM\ProviderImap\Providers\CollectionResource;
/**
@@ -29,6 +31,7 @@ class MailboxStore
'uidNext' => null,
'highestModSeq' => null,
'changeSeq' => 0,
+ 'purgedSeq' => 0,
'harmonizedAt' => null,
'harmonizationComplete' => false,
];
@@ -112,7 +115,7 @@ class MailboxStore
/**
* Harmonization state of a mailbox; defaults when the mailbox is not cached.
*
- * @return array{uidValidity: ?int, uidNext: ?int, highestModSeq: ?int, changeSeq: int, harmonizedAt: ?int, harmonizationComplete: bool}
+ * @return array{uidValidity: ?int, uidNext: ?int, highestModSeq: ?int, changeSeq: int, purgedSeq: int, harmonizedAt: ?int, harmonizationComplete: bool}
*/
public function state(string $serviceId, string $name): array
{
@@ -137,6 +140,34 @@ class MailboxStore
);
}
+ /**
+ * Reserve a range of change sequence numbers for a mailbox.
+ *
+ * Atomically advances the mailbox's changeSeq by $count and returns the first
+ * number of the reserved range (range = first .. first + count - 1).
+ *
+ * @throws RuntimeException when the mailbox is not cached
+ */
+ public function reserveSequence(string $serviceId, string $name, int $count = 1): int
+ {
+ $count = max(1, $count);
+
+ $document = $this->dataStore->selectCollection(self::COLLECTION_NAME)->getMongoCollection()->findOneAndUpdate(
+ ['sid' => $serviceId, 'name' => $name],
+ ['$inc' => ['changeSeq' => $count]],
+ [
+ 'projection' => ['changeSeq' => 1],
+ 'returnDocument' => FindOneAndUpdate::RETURN_DOCUMENT_AFTER,
+ ],
+ );
+
+ if ($document === null) {
+ throw new RuntimeException("Mailbox is not cached: {$name}");
+ }
+
+ return (int) ((array) $document)['changeSeq'] - $count + 1;
+ }
+
/**
* Take the harmonization lock of a mailbox when it is free or expired.
*
diff --git a/lib/Stores/MessageStore.php b/lib/Stores/MessageStore.php
index 660b5e7..0a6d3e4 100644
--- a/lib/Stores/MessageStore.php
+++ b/lib/Stores/MessageStore.php
@@ -18,11 +18,27 @@ use KTXM\ProviderImap\Providers\EntityResource;
* One MongoDB document per cached message in `provider_imap_mail_messages`,
* keyed by (sid, mailbox, uidValidity, uid). Holds only what list, filter and
* sort need; message content lives in the content store (message.json).
+ *
+ * Changes are tracked on the documents themselves (as IMAP CONDSTORE does):
+ * `addedSeq` and `modSeq` come from the mailbox's change sequence, and an
+ * expunged message stays behind as a tombstone with `removedSeq` until it is
+ * purged. Every read except changes() skips tombstones.
*/
class MessageStore
{
protected const COLLECTION_NAME = 'provider_imap_mail_messages';
+ /** Filter that excludes tombstones */
+ private const LIVE = ['removedSeq' => null];
+
+ /** Fields dropped when a message becomes a tombstone */
+ private const TOMBSTONE_UNSET = [
+ 'created' => '', 'received' => '', 'sent' => '', 'size' => '',
+ 'subject' => '', 'from' => '', 'to' => '', 'cc' => '',
+ 'urid' => '', 'inReplyTo' => '', 'references' => '',
+ 'flags' => '', 'hasAttachments' => '', 'preview' => '', 'blobs' => '',
+ ];
+
public function __construct(
protected readonly DataStore $dataStore,
) {}
@@ -63,15 +79,28 @@ class MessageStore
['tid' => 1, 'sid' => 1],
['name' => 'messages_by_tenant_service']
),
+ $collection->createIndex(
+ ['sid' => 1, 'mailbox' => 1, 'uidValidity' => 1, 'addedSeq' => 1],
+ ['name' => 'messages_by_added_seq']
+ ),
+ $collection->createIndex(
+ ['sid' => 1, 'mailbox' => 1, 'uidValidity' => 1, 'modSeq' => 1],
+ ['name' => 'messages_by_mod_seq']
+ ),
+ $collection->createIndex(
+ ['sid' => 1, 'mailbox' => 1, 'uidValidity' => 1, 'removedSeq' => 1],
+ ['name' => 'messages_by_removed_seq', 'partialFilterExpression' => ['removedSeq' => ['$exists' => true]]]
+ ),
];
}
/**
* Insert or replace the meta document of a message.
*
+ * `$seq` becomes the document's modSeq, and its addedSeq when the document is new.
* Fields owned by other writers (e.g. blobs) are left untouched.
*/
- public function upsert(string $tenantId, string $serviceId, int $uidValidity, EntityResource $entity): void
+ public function upsert(string $tenantId, string $serviceId, int $uidValidity, EntityResource $entity, int $seq): void
{
$document = $entity->toCacheMeta();
$key = [
@@ -83,7 +112,10 @@ class MessageStore
$this->dataStore->selectCollection(self::COLLECTION_NAME)->updateOne(
$key,
- ['$set' => ['tid' => $tenantId, ...$key, ...$document]],
+ [
+ '$set' => ['tid' => $tenantId, ...$key, ...$document, 'modSeq' => $seq],
+ '$setOnInsert' => ['addedSeq' => $seq],
+ ],
['upsert' => true],
);
}
@@ -98,6 +130,7 @@ class MessageStore
'mailbox' => $mailbox,
'uidValidity' => $uidValidity,
'uid' => $uid,
+ ...self::LIVE,
]);
}
@@ -117,6 +150,7 @@ class MessageStore
'mailbox' => $mailbox,
'uidValidity' => $uidValidity,
'uid' => ['$in' => array_values($uids)],
+ ...self::LIVE,
]);
$list = [];
@@ -134,7 +168,7 @@ class MessageStore
public function uids(string $serviceId, string $mailbox, int $uidValidity): array
{
$cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find(
- ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity],
+ ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, ...self::LIVE],
['projection' => ['uid' => 1]],
);
@@ -153,7 +187,7 @@ class MessageStore
public function flags(string $serviceId, string $mailbox, int $uidValidity): array
{
$cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find(
- ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity],
+ ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, ...self::LIVE],
['projection' => ['uid' => 1, 'flags' => 1]],
);
@@ -165,37 +199,78 @@ class MessageStore
}
/**
- * Replace the flags of a message.
+ * Replace the flags of a message and stamp it with the change sequence.
*
* @param string[] $flags set flags, e.g. ['seen', 'flagged']
*/
- public function updateFlags(string $serviceId, string $mailbox, int $uidValidity, int $uid, array $flags): void
+ public function updateFlags(string $serviceId, string $mailbox, int $uidValidity, int $uid, array $flags, int $seq): void
{
$this->dataStore->selectCollection(self::COLLECTION_NAME)->updateOne(
- ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, 'uid' => $uid],
- ['$set' => ['flags' => array_values($flags)]],
+ ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, 'uid' => $uid, ...self::LIVE],
+ ['$set' => ['flags' => array_values($flags), 'modSeq' => $seq]],
);
}
/**
- * Delete the meta documents of messages.
+ * Turn messages into tombstones: their list fields are dropped and removedSeq is set.
*/
- public function delete(string $serviceId, string $mailbox, int $uidValidity, int ...$uids): void
+ public function tombstone(string $serviceId, string $mailbox, int $uidValidity, int $seq, int ...$uids): void
{
if ($uids === []) {
return;
}
+ $this->dataStore->selectCollection(self::COLLECTION_NAME)->updateMany(
+ [
+ 'sid' => $serviceId,
+ 'mailbox' => $mailbox,
+ 'uidValidity' => $uidValidity,
+ 'uid' => ['$in' => array_values($uids)],
+ ...self::LIVE,
+ ],
+ [
+ '$set' => ['removedSeq' => $seq, 'modSeq' => $seq],
+ '$unset' => self::TOMBSTONE_UNSET,
+ ],
+ );
+ }
+
+ /**
+ * UIDs added, modified and removed after a change sequence.
+ *
+ * A message added after $since is an addition even if it was modified later;
+ * a message both added and removed after $since is not reported at all.
+ *
+ * @return array{additions: int[], modifications: int[], deletions: int[]}
+ */
+ public function changes(string $serviceId, string $mailbox, int $uidValidity, int $since): array
+ {
+ $key = ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity];
+
+ return [
+ 'additions' => $this->findUids([...$key, ...self::LIVE, 'addedSeq' => ['$gt' => $since]]),
+ 'modifications' => $this->findUids([...$key, ...self::LIVE, 'addedSeq' => ['$lte' => $since], 'modSeq' => ['$gt' => $since]]),
+ 'deletions' => $this->findUids([...$key, 'addedSeq' => ['$lte' => $since], 'removedSeq' => ['$gt' => $since]]),
+ ];
+ }
+
+ /**
+ * Delete tombstones removed at or before a change sequence.
+ *
+ * The caller records $upTo as the mailbox's purgedSeq: older signatures can no
+ * longer be answered with deletions and get a reset instead.
+ */
+ public function purgeTombstones(string $serviceId, string $mailbox, int $upTo): void
+ {
$this->dataStore->selectCollection(self::COLLECTION_NAME)->deleteMany([
- 'sid' => $serviceId,
- 'mailbox' => $mailbox,
- 'uidValidity' => $uidValidity,
- 'uid' => ['$in' => array_values($uids)],
+ 'sid' => $serviceId,
+ 'mailbox' => $mailbox,
+ 'removedSeq' => ['$lte' => $upTo],
]);
}
/**
- * Delete all meta documents of a mailbox (any UIDVALIDITY).
+ * Delete all meta documents of a mailbox (any UIDVALIDITY), tombstones included.
*/
public function deleteByMailbox(string $serviceId, string $mailbox): void
{
@@ -212,4 +287,21 @@ class MessageStore
{
$this->dataStore->selectCollection(self::COLLECTION_NAME)->deleteMany(['sid' => $serviceId]);
}
+
+ /**
+ * @return int[]
+ */
+ private function findUids(array $filter): array
+ {
+ $cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find(
+ $filter,
+ ['projection' => ['uid' => 1], 'sort' => ['uid' => 1]],
+ );
+
+ $uids = [];
+ foreach ($cursor as $document) {
+ $uids[] = (int) $document['uid'];
+ }
+ return $uids;
+ }
}
diff --git a/tests/php/Unit/CacheSerializationTest.php b/tests/php/Unit/CacheSerializationTest.php
index 283ea15..701f8a8 100644
--- a/tests/php/Unit/CacheSerializationTest.php
+++ b/tests/php/Unit/CacheSerializationTest.php
@@ -128,16 +128,18 @@ final class CacheSerializationTest extends TestCase
->with(
['sid' => 'svc', 'mailbox' => 'INBOX', 'uidValidity' => 7, 'uid' => 42],
$this->callback(function (array $update): bool {
- $this->assertSame(['$set'], array_keys($update));
+ $this->assertSame(['$set', '$setOnInsert'], array_keys($update));
$this->assertSame('tenant', $update['$set']['tid']);
$this->assertSame(7, $update['$set']['uidValidity']);
+ $this->assertSame(5, $update['$set']['modSeq']);
+ $this->assertSame(['addedSeq' => 5], $update['$setOnInsert']);
$this->assertArrayNotHasKey('blobs', $update['$set']);
return true;
}),
['upsert' => true],
);
- (new MessageStore($this->dataStore($collection)))->upsert('tenant', 'svc', 7, $this->entity());
+ (new MessageStore($this->dataStore($collection)))->upsert('tenant', 'svc', 7, $this->entity(), 5);
}
public function testMailboxStoreUpsertLeavesHarmonizationState(): void
diff --git a/tests/php/Unit/MessageHarmonizationTest.php b/tests/php/Unit/MessageHarmonizationTest.php
index c87256a..e10b80e 100644
--- a/tests/php/Unit/MessageHarmonizationTest.php
+++ b/tests/php/Unit/MessageHarmonizationTest.php
@@ -5,6 +5,7 @@ declare(strict_types=1);
namespace KTXT\ProviderImap\Tests\Unit;
use Generator;
+use KTXF\Resource\Delta\Delta;
use KTXF\Resource\Filter\IFilter;
use KTXF\Resource\Sort\ISort;
use KTXC\Db\DataStore;
@@ -15,6 +16,7 @@ use KTXM\ProviderImap\Providers\CollectionResource;
use KTXM\ProviderImap\Providers\EntityResource;
use KTXM\ProviderImap\Providers\Service;
use KTXM\ProviderImap\Service\Cache\HarmonizationService;
+use KTXM\ProviderImap\Service\Cache\MessageDeltaService;
use KTXM\ProviderImap\Service\Cache\MessageIngestor;
use KTXM\ProviderImap\Service\Live\LiveMailService;
use KTXM\ProviderImap\Stores\MailboxStore;
@@ -29,6 +31,7 @@ final class MessageHarmonizationTest extends TestCase
private FakeMessageFileStore $files;
private MessageHarmonizationLiveStub $live;
private HarmonizationService $harmonizer;
+ private MessageDeltaService $deltas;
protected function setUp(): void
{
@@ -46,8 +49,9 @@ final class MessageHarmonizationTest extends TestCase
$this->mailboxes,
$this->messages,
$this->files,
- new MessageIngestor($this->files, $this->messages),
+ new MessageIngestor($this->files, $this->messages, $this->mailboxes),
))->for($service, $this->live);
+ $this->deltas = new MessageDeltaService($this->mailboxes, $this->messages);
$this->mailboxes->upsert('tenant', 'svc', (new CollectionResource('imap', 'svc'))->fromImap(new Mailbox('INBOX', '/', [])));
}
@@ -159,6 +163,87 @@ final class MessageHarmonizationTest extends TestCase
$this->assertNull($this->mailboxes->lockOwner);
}
+ public function testEmptySignatureReturnsTheCurrentSignature(): void
+ {
+ $this->live->flags = [1 => [], 2 => [], 3 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+
+ $delta = $this->deltas->delta('svc', 'INBOX', '');
+
+ $this->assertSame('7:3', $delta->signature);
+ $this->assertSame([[], [], []], $this->changeLists($delta));
+ }
+
+ public function testEmptySignatureWhileIncompleteStartsFromZero(): void
+ {
+ $this->live->flags = [1 => [], 2 => []];
+ $this->live->unfetchable = [2];
+ $this->harmonizer->harmonizeMessages('INBOX');
+
+ $delta = $this->deltas->delta('svc', 'INBOX', '');
+ $this->assertSame('7:0', $delta->signature);
+
+ $this->assertSame([['1'], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', '7:0')));
+ }
+
+ public function testDeltaReportsAdditionsModificationsAndDeletions(): void
+ {
+ $this->live->flags = [1 => [], 2 => [], 3 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $signature = $this->deltas->delta('svc', 'INBOX', '')->signature;
+
+ $this->live->flags = [2 => ['\\Seen'], 3 => [], 4 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $delta = $this->deltas->delta('svc', 'INBOX', $signature);
+
+ $this->assertSame([['4'], ['2'], ['1']], $this->changeLists($delta));
+ $this->assertSame([[], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', $delta->signature)));
+ }
+
+ public function testMessagesAddedSinceTheSignatureAreOnlyAdditions(): void
+ {
+ $this->live->flags = [1 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $signature = $this->deltas->delta('svc', 'INBOX', '')->signature;
+
+ $this->live->flags = [1 => [], 2 => [], 3 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $this->live->flags = [1 => [], 2 => ['\\Flagged']];
+ $this->harmonizer->harmonizeMessages('INBOX');
+
+ // 2: added then flagged -> addition; 3: added then expunged -> not reported
+ $this->assertSame([['2'], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', $signature)));
+ }
+
+ public function testSignatureThatCannotBeAnsweredIsReset(): void
+ {
+ $this->live->flags = [1 => [], 2 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $old = $this->deltas->delta('svc', 'INBOX', '')->signature;
+
+ $this->assertSame([['1', '2'], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', 'garbage')));
+ $this->assertSame([['1', '2'], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', '7:999')));
+
+ $this->mailboxes->updateState('svc', 'INBOX', ['purgedSeq' => 2]);
+ $this->assertSame([['1', '2'], [], []], $this->changeLists($this->deltas->delta('svc', 'INBOX', '7:1')));
+
+ $this->live->uidValidity = 8;
+ $this->live->flags = [5 => []];
+ $this->harmonizer->harmonizeMessages('INBOX');
+ $delta = $this->deltas->delta('svc', 'INBOX', $old);
+
+ $this->assertSame([['5'], [], []], $this->changeLists($delta));
+ $this->assertStringStartsWith('8:', $delta->signature);
+ }
+
+ public function testMailboxThatIsNotCachedHasAnEmptyDelta(): void
+ {
+ $delta = $this->deltas->delta('svc', 'Unknown', '7:1');
+
+ $this->assertSame('', $delta->signature);
+ $this->assertSame([[], [], []], $this->changeLists($delta));
+ }
+
public function testLockIsReleasedWhenTheRunFails(): void
{
$this->live->uidValidity = null;
@@ -171,6 +256,18 @@ final class MessageHarmonizationTest extends TestCase
$this->assertNull($this->mailboxes->lockOwner);
}
+
+ /**
+ * @return array{0: string[], 1: string[], 2: string[]}
+ */
+ private function changeLists(Delta $delta): array
+ {
+ return [
+ json_decode(json_encode($delta->additions), true),
+ json_decode(json_encode($delta->modifications), true),
+ json_decode(json_encode($delta->deletions), true),
+ ];
+ }
}
final class MessageHarmonizationLiveStub extends LiveMailService
@@ -272,37 +369,80 @@ final class FakeMailboxStore extends MailboxStore
$this->lockOwner = null;
}
}
+
+ public function reserveSequence(string $serviceId, string $name, int $count = 1): int
+ {
+ $this->documents[$name]['changeSeq'] += $count;
+ return $this->documents[$name]['changeSeq'] - $count + 1;
+ }
+
+ public function state(string $serviceId, string $name): array
+ {
+ return array_replace(self::STATE_DEFAULTS, array_intersect_key($this->documents[$name] ?? [], self::STATE_DEFAULTS));
+ }
}
final class FakeMessageStore extends MessageStore
{
- /** @var array>> flags by UIDVALIDITY and UID */
+ /** @var array, addedSeq: int, modSeq: int, removedSeq: ?int}>> by UIDVALIDITY and UID */
private array $documents = [];
- public function upsert(string $tenantId, string $serviceId, int $uidValidity, EntityResource $entity): void
+ public function upsert(string $tenantId, string $serviceId, int $uidValidity, EntityResource $entity, int $seq): void
{
$meta = $entity->toCacheMeta();
- $this->documents[$uidValidity][$meta['uid']] = $meta['flags'];
+ $existing = $this->documents[$uidValidity][$meta['uid']] ?? null;
+ $this->documents[$uidValidity][$meta['uid']] = [
+ 'flags' => $meta['flags'],
+ 'addedSeq' => $existing['addedSeq'] ?? $seq,
+ 'modSeq' => $seq,
+ 'removedSeq' => null,
+ ];
ksort($this->documents[$uidValidity]);
}
public function flags(string $serviceId, string $mailbox, int $uidValidity): array
{
- return $this->documents[$uidValidity] ?? [];
+ $flags = [];
+ foreach ($this->documents[$uidValidity] ?? [] as $uid => $document) {
+ if ($document['removedSeq'] === null) {
+ $flags[$uid] = $document['flags'];
+ }
+ }
+ return $flags;
}
- public function updateFlags(string $serviceId, string $mailbox, int $uidValidity, int $uid, array $flags): void
+ public function updateFlags(string $serviceId, string $mailbox, int $uidValidity, int $uid, array $flags, int $seq): void
{
- $this->documents[$uidValidity][$uid] = array_values($flags);
+ $this->documents[$uidValidity][$uid]['flags'] = array_values($flags);
+ $this->documents[$uidValidity][$uid]['modSeq'] = $seq;
}
- public function delete(string $serviceId, string $mailbox, int $uidValidity, int ...$uids): void
+ public function tombstone(string $serviceId, string $mailbox, int $uidValidity, int $seq, int ...$uids): void
{
foreach ($uids as $uid) {
- unset($this->documents[$uidValidity][$uid]);
+ $this->documents[$uidValidity][$uid]['removedSeq'] = $seq;
+ $this->documents[$uidValidity][$uid]['modSeq'] = $seq;
+ $this->documents[$uidValidity][$uid]['flags'] = [];
}
}
+ public function changes(string $serviceId, string $mailbox, int $uidValidity, int $since): array
+ {
+ $changes = ['additions' => [], 'modifications' => [], 'deletions' => []];
+ foreach ($this->documents[$uidValidity] ?? [] as $uid => $document) {
+ if ($document['removedSeq'] !== null) {
+ if ($document['addedSeq'] <= $since && $document['removedSeq'] > $since) {
+ $changes['deletions'][] = $uid;
+ }
+ } elseif ($document['addedSeq'] > $since) {
+ $changes['additions'][] = $uid;
+ } elseif ($document['modSeq'] > $since) {
+ $changes['modifications'][] = $uid;
+ }
+ }
+ return $changes;
+ }
+
public function deleteByMailbox(string $serviceId, string $mailbox): void
{
$this->documents = [];
diff --git a/tests/php/Unit/MessageIngestorTest.php b/tests/php/Unit/MessageIngestorTest.php
index 5b69f93..fd92f15 100644
--- a/tests/php/Unit/MessageIngestorTest.php
+++ b/tests/php/Unit/MessageIngestorTest.php
@@ -7,6 +7,7 @@ namespace KTXT\ProviderImap\Tests\Unit;
use KTXM\ProviderImap\Providers\EntityResource;
use KTXM\ProviderImap\Providers\MessageProperties;
use KTXM\ProviderImap\Service\Cache\MessageIngestor;
+use KTXM\ProviderImap\Stores\MailboxStore;
use KTXM\ProviderImap\Stores\MessageFileStore;
use KTXM\ProviderImap\Stores\MessageStore;
use PHPUnit\Framework\TestCase;
@@ -26,20 +27,22 @@ final class MessageIngestorTest extends TestCase
);
$messageStore = $this->createStub(MessageStore::class);
$messageStore->method('upsert')->willReturnCallback(
- function (string $tid, string $sid, int $uidValidity, EntityResource $entity) use (&$calls): void {
- $calls[] = ['meta', $tid, $sid, $entity->collection(), $uidValidity, $entity->identifier()];
+ function (string $tid, string $sid, int $uidValidity, EntityResource $entity, int $seq) use (&$calls): void {
+ $calls[] = ['meta', $tid, $sid, $entity->collection(), $uidValidity, $entity->identifier(), $seq];
},
);
+ $mailboxStore = $this->createStub(MailboxStore::class);
+ $mailboxStore->method('reserveSequence')->willReturn(10);
- $ingested = (new MessageIngestor($fileStore, $messageStore))
+ $ingested = (new MessageIngestor($fileStore, $messageStore, $mailboxStore))
->ingest('tenant', 'svc', 7, $this->entity(41), $this->entity(42));
$this->assertSame([41, 42], $ingested);
$this->assertSame([
['file', 'tenant', 'svc', 'INBOX', 7, 41, EntityResource::CACHE_SCHEMA_VERSION],
- ['meta', 'tenant', 'svc', 'INBOX', 7, 41],
+ ['meta', 'tenant', 'svc', 'INBOX', 7, 41, 10],
['file', 'tenant', 'svc', 'INBOX', 7, 42, EntityResource::CACHE_SCHEMA_VERSION],
- ['meta', 'tenant', 'svc', 'INBOX', 7, 42],
+ ['meta', 'tenant', 'svc', 'INBOX', 7, 42, 11],
], $calls);
}
@@ -51,7 +54,7 @@ final class MessageIngestorTest extends TestCase
$messageStore->expects($this->never())->method('upsert');
$this->expectException(RuntimeException::class);
- (new MessageIngestor($fileStore, $messageStore))->ingest('tenant', 'svc', 7, $this->entity(41));
+ (new MessageIngestor($fileStore, $messageStore, $this->createStub(MailboxStore::class)))->ingest('tenant', 'svc', 7, $this->entity(41));
}
private function entity(int $uid): EntityResource
diff --git a/tests/php/Unit/MessageStoreChangesTest.php b/tests/php/Unit/MessageStoreChangesTest.php
new file mode 100644
index 0000000..a2b5d5f
--- /dev/null
+++ b/tests/php/Unit/MessageStoreChangesTest.php
@@ -0,0 +1,92 @@
+createStub(Collection::class);
+ $collection->method('find')->willReturnCallback(function (array $filter) use (&$filters): Cursor {
+ $filters[] = $filter;
+ return new Cursor(new ArrayIterator([['uid' => count($filters)]]));
+ });
+
+ $changes = (new MessageStore($this->dataStore($collection)))->changes('svc', 'INBOX', 7, 10);
+
+ $key = ['sid' => 'svc', 'mailbox' => 'INBOX', 'uidValidity' => 7];
+ $this->assertSame([
+ [...$key, 'removedSeq' => null, 'addedSeq' => ['$gt' => 10]],
+ [...$key, 'removedSeq' => null, 'addedSeq' => ['$lte' => 10], 'modSeq' => ['$gt' => 10]],
+ [...$key, 'addedSeq' => ['$lte' => 10], 'removedSeq' => ['$gt' => 10]],
+ ], $filters);
+ $this->assertSame(['additions' => [1], 'modifications' => [2], 'deletions' => [3]], $changes);
+ }
+
+ public function testTombstoneKeepsKeyAndSequencesOnly(): void
+ {
+ $collection = $this->createMock(Collection::class);
+ $collection->expects($this->once())
+ ->method('updateMany')
+ ->with(
+ ['sid' => 'svc', 'mailbox' => 'INBOX', 'uidValidity' => 7, 'uid' => ['$in' => [3, 4]], 'removedSeq' => null],
+ $this->callback(function (array $update): bool {
+ $this->assertSame(['removedSeq' => 12, 'modSeq' => 12], $update['$set']);
+ $this->assertArrayHasKey('subject', $update['$unset']);
+ $this->assertArrayHasKey('flags', $update['$unset']);
+ $this->assertArrayNotHasKey('addedSeq', $update['$unset']);
+ return true;
+ }),
+ );
+
+ (new MessageStore($this->dataStore($collection)))->tombstone('svc', 'INBOX', 7, 12, 3, 4);
+ }
+
+ public function testLiveReadsSkipTombstones(): void
+ {
+ $collection = $this->createMock(Collection::class);
+ $collection->expects($this->once())
+ ->method('find')
+ ->with($this->callback(fn (array $filter): bool => array_key_exists('removedSeq', $filter) && $filter['removedSeq'] === null))
+ ->willReturn(new Cursor(new ArrayIterator([])));
+
+ (new MessageStore($this->dataStore($collection)))->flags('svc', 'INBOX', 7);
+ }
+
+ public function testReserveSequenceReturnsTheFirstNumberOfTheRange(): void
+ {
+ $mongo = $this->createMock(MongoCollection::class);
+ $mongo->expects($this->once())
+ ->method('findOneAndUpdate')
+ ->with(
+ ['sid' => 'svc', 'name' => 'INBOX'],
+ ['$inc' => ['changeSeq' => 3]],
+ $this->callback(fn (array $options): bool => $options['returnDocument'] === FindOneAndUpdate::RETURN_DOCUMENT_AFTER),
+ )
+ ->willReturn(['changeSeq' => 12]);
+ $collection = $this->createStub(Collection::class);
+ $collection->method('getMongoCollection')->willReturn($mongo);
+
+ $this->assertSame(10, (new MailboxStore($this->dataStore($collection)))->reserveSequence('svc', 'INBOX', 3));
+ }
+
+ private function dataStore(Collection $collection): DataStore
+ {
+ $dataStore = $this->createStub(DataStore::class);
+ $dataStore->method('selectCollection')->willReturn($collection);
+ return $dataStore;
+ }
+}