From eea9b5630a87837dc9a6e7d1187935932c6e9330 Mon Sep 17 00:00:00 2001 From: Sebastian Krupinski Date: Wed, 7 Oct 2026 21:21:14 -0400 Subject: [PATCH] feat: track cache changes with sequence numbers and tombstones Signed-off-by: Sebastian Krupinski --- lib/Console/HarmonizeCommand.php | 9 +- lib/Service/Cache/HarmonizationService.php | 8 +- lib/Service/Cache/MessageDeltaService.php | 88 +++++++++++ lib/Service/Cache/MessageIngestor.php | 26 +++- lib/Stores/MailboxStore.php | 33 +++- lib/Stores/MessageStore.php | 122 +++++++++++++-- tests/php/Unit/CacheSerializationTest.php | 6 +- tests/php/Unit/MessageHarmonizationTest.php | 158 ++++++++++++++++++-- tests/php/Unit/MessageIngestorTest.php | 15 +- tests/php/Unit/MessageStoreChangesTest.php | 92 ++++++++++++ 10 files changed, 511 insertions(+), 46 deletions(-) create mode 100644 lib/Service/Cache/MessageDeltaService.php create mode 100644 tests/php/Unit/MessageStoreChangesTest.php 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; + } +}