From 65fdb23fa4551cd2c98bc461f57d0394378016d6 Mon Sep 17 00:00:00 2001 From: Sebastian Krupinski Date: Mon, 5 Oct 2026 20:34:02 -0400 Subject: [PATCH] feat: harmonize mailbox messages Signed-off-by: Sebastian Krupinski --- lib/Client/Mailbox.php | 14 + lib/Client/Protocol/Command/ListCommand.php | 13 +- lib/Client/Protocol/Command/SelectCommand.php | 21 +- lib/Providers/MessageProperties.php | 30 +- lib/Service/Cache/HarmonizationService.php | 143 +++++++++ lib/Service/Live/LiveMailService.php | 95 +++++- lib/Stores/MessageStore.php | 19 ++ tests/php/Unit/CollectionFetchTest.php | 95 ++++++ tests/php/Unit/HarmonizationServiceTest.php | 4 +- tests/php/Unit/MailboxSelectionTest.php | 116 +++++++ tests/php/Unit/MessageHarmonizationTest.php | 303 ++++++++++++++++++ tests/php/Unit/SelectCommandTest.php | 69 ++++ 12 files changed, 879 insertions(+), 43 deletions(-) create mode 100644 tests/php/Unit/CollectionFetchTest.php create mode 100644 tests/php/Unit/MailboxSelectionTest.php create mode 100644 tests/php/Unit/MessageHarmonizationTest.php create mode 100644 tests/php/Unit/SelectCommandTest.php diff --git a/lib/Client/Mailbox.php b/lib/Client/Mailbox.php index 0b3119d..7e7576e 100644 --- a/lib/Client/Mailbox.php +++ b/lib/Client/Mailbox.php @@ -22,6 +22,8 @@ final class Mailbox private readonly int $recent = 0, private readonly array $flags = [], private readonly bool $readOnly = true, + private readonly ?int $uidNext = null, + private readonly ?int $highestModSeq = null, ) {} public function fromStatus(MailboxStatusResult $status): self @@ -38,6 +40,8 @@ final class Mailbox $this->recent, $this->flags, $this->readOnly, + $items['UIDNEXT'] ?? $this->uidNext, + $items['HIGHESTMODSEQ'] ?? $this->highestModSeq, ); } @@ -64,6 +68,16 @@ final class Mailbox return $this->uidValidity; } + public function uidNext(): ?int + { + return $this->uidNext; + } + + public function highestModSeq(): ?int + { + return $this->highestModSeq; + } + public function messages(): int { return $this->messages; diff --git a/lib/Client/Protocol/Command/ListCommand.php b/lib/Client/Protocol/Command/ListCommand.php index 95a0ca4..795318e 100644 --- a/lib/Client/Protocol/Command/ListCommand.php +++ b/lib/Client/Protocol/Command/ListCommand.php @@ -13,6 +13,7 @@ use KTXM\ProviderImap\Client\ImapException; use KTXM\ProviderImap\Client\Protocol\Command\Argument\ListReturnOptions; use KTXM\ProviderImap\Client\Protocol\Command\Argument\ListSelectionOptions; use KTXM\ProviderImap\Client\Mailbox; +use KTXM\ProviderImap\Client\Result\MailboxStatusResult; use KTXM\ProviderImap\Client\Protocol\RequestFrame; use KTXM\ProviderImap\Client\Protocol\Response\TaggedResponse; use KTXM\ProviderImap\Client\Protocol\Response\UntaggedResponse; @@ -144,16 +145,6 @@ final class ListCommand implements CommandInterface */ private function applyStatus(Mailbox $mailbox, array $status): Mailbox { - return new Mailbox( - $mailbox->name(), - $mailbox->delimiter(), - $mailbox->attributes(), - $status['MESSAGES'] ?? $mailbox->messages(), - $status['UNSEEN'] ?? $mailbox->unread(), - $mailbox->uidValidity(), - $mailbox->recent(), - $mailbox->flags(), - $mailbox->readOnly(), - ); + return $mailbox->fromStatus(new MailboxStatusResult($mailbox->name(), $status)); } } \ No newline at end of file diff --git a/lib/Client/Protocol/Command/SelectCommand.php b/lib/Client/Protocol/Command/SelectCommand.php index f6d718a..02b0390 100644 --- a/lib/Client/Protocol/Command/SelectCommand.php +++ b/lib/Client/Protocol/Command/SelectCommand.php @@ -53,6 +53,9 @@ final class SelectCommand implements CommandInterface { $exists = 0; $recent = 0; + $uidValidity = null; + $uidNext = null; + $highestModSeq = null; $flags = []; $readOnly = $this->readOnly; @@ -70,6 +73,16 @@ final class SelectCommand implements CommandInterface continue; } + // response codes, e.g. "* OK [UIDVALIDITY 3857529045] UIDs valid" + if (preg_match('/^\*\s+OK\s+\[(UIDVALIDITY|UIDNEXT|HIGHESTMODSEQ)\s+(\d+)\]/i', $raw, $matches)) { + match (strtoupper($matches[1])) { + 'UIDVALIDITY' => $uidValidity = (int) $matches[2], + 'UIDNEXT' => $uidNext = (int) $matches[2], + 'HIGHESTMODSEQ' => $highestModSeq = (int) $matches[2], + }; + continue; + } + if ($response->label() === 'FLAGS' && preg_match('/\(([^)]*)\)/', $response->payload(), $matches)) { $flags = $this->parseFlags($matches[1]); continue; @@ -77,6 +90,10 @@ final class SelectCommand implements CommandInterface } if ($response instanceof TaggedResponse) { + // a failed SELECT leaves no mailbox selected (RFC 3501 6.3.1) + if ($response->status() !== 'OK') { + $context->setSelectedMailbox(null); + } CompletionChecker::assertSuccess($this->name(), $response); if (str_contains(strtoupper($response->text()), 'READ-ONLY')) { @@ -92,10 +109,12 @@ final class SelectCommand implements CommandInterface [], $exists, 0, - null, + $uidValidity, $recent, $flags, $readOnly, + $uidNext, + $highestModSeq, ); } } diff --git a/lib/Providers/MessageProperties.php b/lib/Providers/MessageProperties.php index b3d8df7..5047e2b 100644 --- a/lib/Providers/MessageProperties.php +++ b/lib/Providers/MessageProperties.php @@ -100,23 +100,27 @@ class MessageProperties extends MessagePropertiesMutableAbstract { } } - $this->data[static::PROPERTY_FLAGS] = []; - foreach ($message->flags() as $flag) { - $flag = ltrim($flag, '\\'); - $normalized = match (strtolower($flag)) { - 'seen' => 'seen', - 'flagged' => 'flagged', - 'answered' => 'answered', - 'draft' => 'draft', - 'deleted' => 'deleted', - default => strtolower($flag), - }; - $this->data[static::PROPERTY_FLAGS][$normalized] = true; - } + $this->data[static::PROPERTY_FLAGS] = array_fill_keys(self::normalizeFlags($message->flags()), true); return $this; } + /** + * Normalise IMAP flags (e.g. "\\Seen", "$Forwarded") to flag names (e.g. "seen", "$forwarded"). + * + * @param string[] $flags + * @return list + */ + public static function normalizeFlags(array $flags): array + { + $normalized = []; + foreach ($flags as $flag) { + $normalized[] = strtolower(ltrim($flag, '\\')); + } + + return array_values(array_unique($normalized)); + } + // ── Cache (meta store / content store) ─────────────────────────────────── /** diff --git a/lib/Service/Cache/HarmonizationService.php b/lib/Service/Cache/HarmonizationService.php index f08618b..cf5f93c 100644 --- a/lib/Service/Cache/HarmonizationService.php +++ b/lib/Service/Cache/HarmonizationService.php @@ -9,13 +9,17 @@ declare(strict_types=1); namespace KTXM\ProviderImap\Service\Cache; +use KTXM\ProviderImap\Client\Protocol\Command\Argument\FetchOptions; use KTXM\ProviderImap\Providers\CollectionResource; +use KTXM\ProviderImap\Providers\EntityResource; +use KTXM\ProviderImap\Providers\MessageProperties; use KTXM\ProviderImap\Providers\Service; use KTXM\ProviderImap\Service\Live\LiveMailService; use KTXM\ProviderImap\Stores\MailboxStore; use KTXM\ProviderImap\Stores\MessageFileStore; use KTXM\ProviderImap\Stores\MessageStore; use LogicException; +use RuntimeException; /** * Brings the cache of one service in line with its IMAP server. @@ -25,6 +29,12 @@ use LogicException; */ class HarmonizationService { + /** Messages fetched per UID FETCH when ingesting; each is cached as it streams in */ + public const INGEST_BATCH_SIZE = 200; + + /** Lifetime of a mailbox harmonization lock in seconds; long enough for a large first run */ + public const LOCK_TTL = 900; + private ?Service $service = null; private ?LiveMailService $live = null; @@ -32,6 +42,7 @@ class HarmonizationService private readonly MailboxStore $mailboxStore, private readonly MessageStore $messageStore, private readonly MessageFileStore $fileStore, + private readonly MessageIngestor $ingestor, ) {} /** @@ -87,6 +98,138 @@ class HarmonizationService return $result; } + /** + * Harmonize the messages of one mailbox with the server. + * + * Ingests messages that are not cached yet (newest first), updates changed + * flags and removes messages expunged on the server. A changed UIDVALIDITY + * discards the mailbox's cache first. Never waits: when another run holds the + * mailbox lock, nothing is done. + * + * @return array{status: 'harmonized'|'skipped', reset: bool, added: int, updated: int, removed: int, complete: bool} + */ + public function harmonizeMessages(string $mailbox): array + { + [$service, $live] = $this->bound(); + $tenantId = (string) $service->tenantIdentifier(); + $serviceId = (string) $service->identifier(); + + $result = ['status' => 'skipped', 'reset' => false, 'added' => 0, 'updated' => 0, 'removed' => 0, 'complete' => false]; + + // the lock lives on the mailbox document, so it has to exist first + if ($this->mailboxStore->fetch($serviceId, $mailbox) === null) { + $this->harmonizeMailboxes(); + if ($this->mailboxStore->fetch($serviceId, $mailbox) === null) { + throw new RuntimeException("Mailbox not found on server: {$mailbox}"); + } + } + + $owner = bin2hex(random_bytes(8)); + if (!$this->mailboxStore->acquireLock($serviceId, $mailbox, $owner, self::LOCK_TTL)) { + return $result; + } + + try { + $state = $this->mailboxStore->state($serviceId, $mailbox); + $selected = $live->collectionFetch($mailbox) + ?? throw new RuntimeException("Mailbox not found on server: {$mailbox}"); + $uidValidity = $selected->uidValidity() + ?? throw new RuntimeException("Server did not report UIDVALIDITY for mailbox: {$mailbox}"); + + // a new UIDVALIDITY invalidates every cached UID of the mailbox + if ($state['uidValidity'] !== null && (int) $state['uidValidity'] !== $uidValidity) { + $this->messageStore->deleteByMailbox($serviceId, $mailbox); + $this->fileStore->deleteByMailbox($tenantId, $serviceId, $mailbox); + $result['reset'] = true; + } + + $cachedFlags = $this->messageStore->flags($serviceId, $mailbox, $uidValidity); + + // stream every UID and its flags from the server: cached messages get their flags + // compared (and are dropped from $cachedFlags), unknown UIDs are collected as new. + // No other IMAP command may run until the stream is consumed. + $new = []; + foreach ($live->entityFlags($mailbox) as $uid => $flags) { + if (!isset($cachedFlags[$uid])) { + $new[] = $uid; + continue; + } + + $flags = MessageProperties::normalizeFlags($flags); + if (!self::sameFlags($flags, $cachedFlags[$uid])) { + $this->messageStore->updateFlags($serviceId, $mailbox, $uidValidity, $uid, $flags); + $result['updated']++; + } + unset($cachedFlags[$uid]); + } + + // new: newest first so recent mail is available soonest + rsort($new); + $missing = []; + foreach (array_chunk($new, self::INGEST_BATCH_SIZE) as $batch) { + $ingested = $this->ingestBatch($tenantId, $serviceId, $mailbox, $uidValidity, $batch); + $result['added'] += count($ingested); + array_push($missing, ...array_diff($batch, $ingested)); + } + + // expunged: cached UIDs the server did not list; meta documents first, then files + $expunged = array_keys($cachedFlags); + if ($expunged !== []) { + $this->messageStore->delete($serviceId, $mailbox, $uidValidity, ...$expunged); + $this->fileStore->delete($tenantId, $serviceId, $mailbox, $uidValidity, ...$expunged); + $result['removed'] = count($expunged); + } + + // UIDs that could not be fetched (e.g. expunged meanwhile) are retried next run + $result['complete'] = $missing === []; + $result['status'] = 'harmonized'; + + $this->mailboxStore->updateState($serviceId, $mailbox, [ + 'uidValidity' => $uidValidity, + 'uidNext' => $selected->uidNext(), + 'highestModSeq' => $selected->highestModSeq(), + 'harmonizedAt' => time(), + 'harmonizationComplete' => $result['complete'], + ]); + } finally { + $this->mailboxStore->releaseLock($serviceId, $mailbox, $owner); + } + + return $result; + } + + /** + * Fetch a batch of messages and cache each one as it streams in. + * + * @param int[] $uids + * @return int[] UIDs that were cached + */ + private function ingestBatch(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, array $uids): array + { + [$service, $live] = $this->bound(); + $options = FetchOptions::message()->withBodyText(MessageIngestor::BODY_TEXT_LIMIT); + + $ingested = []; + foreach ($live->entityFetch($mailbox, $options, ...$uids) as $message) { + $entity = (new EntityResource($service->provider(), $serviceId))->fromImap($message, $mailbox); + array_push($ingested, ...$this->ingestor->ingest($tenantId, $serviceId, $uidValidity, $entity)); + } + + return $ingested; + } + + /** + * @param string[] $left + * @param string[] $right + */ + private static function sameFlags(array $left, array $right): bool + { + sort($left); + sort($right); + + return $left === $right; + } + /** * Remove a mailbox and everything cached for it: meta documents first, then files, then the mailbox. */ diff --git a/lib/Service/Live/LiveMailService.php b/lib/Service/Live/LiveMailService.php index 294ec77..9702c9e 100644 --- a/lib/Service/Live/LiveMailService.php +++ b/lib/Service/Live/LiveMailService.php @@ -65,6 +65,8 @@ class LiveMailService private ?ImapClient $imapClient = null; private ?SmtpClient $smtpClient = null; + /** Result of the last SELECT on the current connection, reused while still selected */ + private ?Mailbox $selected = null; public function __construct( private readonly Service $service, @@ -200,16 +202,25 @@ class LiveMailService */ public function collectionFetch(string $identifier): ?Mailbox { + // LIST-STATUS (RFC 5819) returns the status with the LIST response in one round trip + $listStatus = $this->imapClient()->hasCapability('LIST-STATUS'); + $command = $listStatus + ? new ListCommand('', $identifier, null, ListReturnOptions::status(...self::DEFAULT_MAILBOX_STATUS_ITEMS)) + : new ListCommand('', $identifier); + // retrieve mailbox from remote - $mailbox = iterator_to_array($this->imapClient()->perform(new ListCommand('', $identifier, null, ListReturnOptions::status(...self::DEFAULT_MAILBOX_STATUS_ITEMS)))); + $mailbox = iterator_to_array($this->imapClient()->perform($command)); if (empty($mailbox)) { return null; } $mailbox = reset($mailbox); - // enrich with STATUS - $status = $this->imapClient()->perform(new StatusCommand($mailbox->name(), self::DEFAULT_MAILBOX_STATUS_ITEMS)); - $mailbox = $mailbox->fromStatus($status); - + + // enrich with STATUS when LIST could not provide it + if (!$listStatus && $mailbox->isSelectable()) { + $status = $this->imapClient()->perform(new StatusCommand($mailbox->name(), self::DEFAULT_MAILBOX_STATUS_ITEMS)); + $mailbox = $mailbox->fromStatus($status); + } + return $mailbox; } @@ -270,7 +281,7 @@ class LiveMailService $nativeFilter = $this->buildEntitySearchCriteria($filter); $nativeSort = $sort !== null ? $this->entitySortCriteria($sort) : []; - $this->imapClient()->perform(new SelectCommand($collection, true)); + $this->select($collection); $rfc5258 = $this->imapClient()->hasCapability('SORT'); $uids = []; @@ -298,6 +309,31 @@ class LiveMailService return $this->entityApplyRange($uids, $range); } + /** + * Stream the flags of every message in a mailbox. + * + * The stream also lists every UID that exists in the mailbox. + * + * @return Generator> IMAP flags keyed by UID + */ + public function entityFlags(string $collection): Generator + { + // fresh SELECT: the message count decides whether "1:*" may be sent + $mailbox = $this->select($collection, refresh: true); + + // "1:*" on an empty mailbox is rejected by some servers + if ($mailbox->messages() === 0) { + return; + } + + foreach ($this->imapClient()->perform(new FetchManyCommand( + MessageTarget::uid('1:*'), + FetchOptions::of('FLAGS'), + )) as $message) { + yield $message->uid() => $message->flags(); + } + } + /** * Determine which of the given UIDs exist in a mailbox. * @@ -309,7 +345,7 @@ class LiveMailService return []; } - $this->imapClient()->perform(new SelectCommand($collection, true)); + $this->select($collection); return $this->imapClient()->perform(new SearchCommand( SearchCriteriaBuilder::create()->uid(SequenceSet::items(...$uids)), @@ -329,9 +365,10 @@ class LiveMailService // fast path: fetch all messages without filtering, sorting or pagination if ($filter === null && $sort === null && $range === null) { - $mailbox = $this->imapClient()->perform(new SelectCommand($collection, true)); + $mailbox = $this->select($collection, refresh: true); - if ($mailbox === null) { + // "1:*" on an empty mailbox is rejected by some servers + if ($mailbox->messages() === 0) { return []; } @@ -366,7 +403,7 @@ class LiveMailService } $options ??= FetchOptions::message()->withBodyText(); - $this->imapClient()->perform(new SelectCommand($collection, true)); + $this->select($collection); $request = new FetchManyCommand( MessageTarget::uid(SequenceSet::items(...array_values($uids))), @@ -391,7 +428,7 @@ class LiveMailService */ public function entityDownload(string $collection, int $uid, ?string $partId = null): BinaryResource { - $this->imapClient()->perform(new SelectCommand($collection, true)); + $this->select($collection); $encoding = null; @@ -550,7 +587,7 @@ class LiveMailService return; } - $this->imapClient()->perform(new SelectCommand($collection, false)); + $this->select($collection, readOnly: false); $this->imapClient()->perform(new StoreCommand( MessageTarget::uid(SequenceSet::items(...array_values($uids))), $flags, @@ -569,7 +606,7 @@ class LiveMailService $target = MessageTarget::uid(SequenceSet::items(...array_values($uids))); - $this->imapClient()->perform(new SelectCommand($collection, false)); + $this->select($collection, readOnly: false); $this->imapClient()->perform(new StoreCommand($target, ['\\Deleted'], '+')); $this->imapClient()->perform(new ExpungeCommand($target)); @@ -583,7 +620,7 @@ class LiveMailService return; } - $this->imapClient()->perform(new SelectCommand($collection, false)); + $this->select($collection, readOnly: false); $flagsToAdd = $this->normalizeFlags($flagsToAdd); $flagsToRemove = $this->normalizeFlags($flagsToRemove); @@ -615,13 +652,13 @@ class LiveMailService // if MOVE is supported, use it; otherwise, fall back to COPY + EXPUNGE if ($rfc6851) { - $this->imapClient()->perform(new SelectCommand($sourceCollection, false)); + $this->select($sourceCollection, readOnly: false); $response = $this->imapClient()->perform(new MoveCommand( MessageTarget::uid(SequenceSet::items(...array_values($uids))), $targetCollection, )); } else { - $this->imapClient()->perform(new SelectCommand($sourceCollection, false)); + $this->select($sourceCollection, readOnly: false); $response = $this->imapClient()->perform(new CopyCommand( MessageTarget::uid(SequenceSet::items(...array_values($uids))), $targetCollection, @@ -657,7 +694,7 @@ class LiveMailService return []; } - $this->imapClient()->perform(new SelectCommand($sourceCollection, false)); + $this->select($sourceCollection, readOnly: false); $response = $this->imapClient()->perform(new CopyCommand( MessageTarget::uid(SequenceSet::items(...array_values($uids))), $targetCollection, @@ -1143,4 +1180,28 @@ class LiveMailService } return $normalized; } + + /** + * Select a mailbox, reusing the current selection when possible. + * + * The current selection is reused when the same mailbox is still selected on the + * session and its access mode suffices (read-write covers read-only). A refresh + * forces a new SELECT, e.g. when an up-to-date message count is needed. + */ + private function select(string $collection, bool $readOnly = true, bool $refresh = false): Mailbox + { + if (!$refresh + && $this->selected !== null + && $this->selected->name() === $collection + && $this->imapClient()->session()->selectedMailbox() === $collection + && ($readOnly || !$this->selected->readOnly()) + ) { + return $this->selected; + } + + $this->selected = null; + $this->selected = $this->imapClient()->perform(new SelectCommand($collection, $readOnly)); + + return $this->selected; + } } diff --git a/lib/Stores/MessageStore.php b/lib/Stores/MessageStore.php index da10dcf..660b5e7 100644 --- a/lib/Stores/MessageStore.php +++ b/lib/Stores/MessageStore.php @@ -145,6 +145,25 @@ class MessageStore return $uids; } + /** + * List the flags of every message cached for a mailbox. + * + * @return array> set flags keyed by UID + */ + public function flags(string $serviceId, string $mailbox, int $uidValidity): array + { + $cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find( + ['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity], + ['projection' => ['uid' => 1, 'flags' => 1]], + ); + + $flags = []; + foreach ($cursor as $document) { + $flags[(int) $document['uid']] = array_values(array_map('strval', (array) ($document['flags'] ?? []))); + } + return $flags; + } + /** * Replace the flags of a message. * diff --git a/tests/php/Unit/CollectionFetchTest.php b/tests/php/Unit/CollectionFetchTest.php new file mode 100644 index 0000000..b2c0f76 --- /dev/null +++ b/tests/php/Unit/CollectionFetchTest.php @@ -0,0 +1,95 @@ +service(listStatus: true)->collectionFetch('INBOX'); + + $this->assertCount(1, $this->commands); + $this->assertStringContainsString('RETURN (STATUS', $this->commands[0]); + $this->assertSame(3, $mailbox?->messages()); + $this->assertSame(7, $mailbox?->uidValidity()); + $this->assertSame(10, $mailbox?->uidNext()); + } + + public function testWithoutListStatusFallsBackToStatus(): void + { + $mailbox = $this->service(listStatus: false)->collectionFetch('INBOX'); + + $this->assertCount(2, $this->commands); + $this->assertStringNotContainsString('RETURN', $this->commands[0]); + $this->assertStringStartsWith('STATUS', $this->commands[1]); + $this->assertSame(3, $mailbox?->messages()); + $this->assertSame(7, $mailbox?->uidValidity()); + $this->assertSame(10, $mailbox?->uidNext()); + } + + public function testUnknownMailboxIsNull(): void + { + $this->assertNull($this->service(listStatus: false, exists: false)->collectionFetch('Missing')); + } + + private function service(bool $listStatus, bool $exists = true): LiveMailService + { + $lines = ["* PREAUTH Ready\r\n"]; + $connection = $this->createStub(ConnectionInterface::class); + $connection->method('readLine')->willReturnCallback(static function () use (&$lines): string { + return array_shift($lines) ?? throw new \RuntimeException('Unexpected response read'); + }); + $connection->method('write')->willReturnCallback(function (string $wire) use (&$lines, $listStatus, $exists): void { + [$tag, $command] = explode(' ', trim($wire), 2); + $operation = explode(' ', $command)[0]; + if ($operation === 'CAPABILITY') { + $lines[] = '* CAPABILITY IMAP4rev1' . ($listStatus ? ' LIST-STATUS' : '') . "\r\n"; + } else { + $this->commands[] = $command; + } + if ($operation === 'LIST') { + if (!$listStatus && str_contains($command, 'RETURN')) { + $lines[] = "$tag BAD Unknown RETURN option\r\n"; + return; + } + if ($exists) { + $lines[] = "* LIST (\\HasNoChildren) \"/\" \"INBOX\"\r\n"; + if ($listStatus) { + $lines[] = "* STATUS \"INBOX\" (MESSAGES 3 UNSEEN 1 UIDNEXT 10 UIDVALIDITY 7)\r\n"; + } + } + } + if ($operation === 'STATUS') { + $lines[] = "* STATUS \"INBOX\" (MESSAGES 3 UNSEEN 1 UIDNEXT 10 UIDVALIDITY 7)\r\n"; + } + $lines[] = "$tag OK Completed\r\n"; + }); + $factory = $this->createStub(ConnectionFactoryInterface::class); + $factory->method('create')->willReturn($connection); + + $client = new Client($factory); + $client->connect(new ConnectionConfig('localhost')); + + return new class($this->createStub(Service::class), $client) extends LiveMailService { + public function __construct(Service $service, private readonly Client $client) + { + parent::__construct($service); + } + + public function imapClient(): Client + { + return $this->client; + } + }; + } +} diff --git a/tests/php/Unit/HarmonizationServiceTest.php b/tests/php/Unit/HarmonizationServiceTest.php index 789523f..ffd0fa0 100644 --- a/tests/php/Unit/HarmonizationServiceTest.php +++ b/tests/php/Unit/HarmonizationServiceTest.php @@ -13,6 +13,7 @@ use KTXM\ProviderImap\Client\Mailbox; use KTXM\ProviderImap\Providers\CollectionResource; use KTXM\ProviderImap\Providers\Service; use KTXM\ProviderImap\Service\Cache\HarmonizationService; +use KTXM\ProviderImap\Service\Cache\MessageIngestor; use KTXM\ProviderImap\Service\Live\LiveMailService; use KTXM\ProviderImap\Stores\MailboxStore; use KTXM\ProviderImap\Stores\MessageFileStore; @@ -63,6 +64,7 @@ final class HarmonizationServiceTest extends TestCase $this->createStub(MailboxStore::class), $this->createStub(MessageStore::class), $this->createStub(MessageFileStore::class), + $this->createStub(MessageIngestor::class), ); $this->expectException(LogicException::class); @@ -148,7 +150,7 @@ final class HarmonizationServiceTest extends TestCase $live = new HarmonizationServiceTestLiveStub($service); $live->mailboxes = array_map(static fn (string $name): Mailbox => new Mailbox($name, '/', []), $remote); - return (new HarmonizationService($mailboxes, $messages, $files))->for($service, $live); + return (new HarmonizationService($mailboxes, $messages, $files, $this->createStub(MessageIngestor::class)))->for($service, $live); } private function dataStore(Collection $collection): DataStore diff --git a/tests/php/Unit/MailboxSelectionTest.php b/tests/php/Unit/MailboxSelectionTest.php new file mode 100644 index 0000000..aeb297e --- /dev/null +++ b/tests/php/Unit/MailboxSelectionTest.php @@ -0,0 +1,116 @@ +service(); + + $service->entityExtant('INBOX', 1); + $service->entityExtant('INBOX', 2); + + $this->assertSame(1, $this->sent('EXAMINE "INBOX"')); + } + + public function testOtherMailboxIsSelected(): void + { + $service = $this->service(); + + $service->entityExtant('INBOX', 1); + $service->entityExtant('Sent', 1); + $service->entityExtant('INBOX', 1); + + $this->assertSame(2, $this->sent('EXAMINE "INBOX"')); + $this->assertSame(1, $this->sent('EXAMINE "Sent"')); + } + + public function testWriteAfterReadOnlySelectionSelectsReadWrite(): void + { + $service = $this->service(); + + $service->entityExtant('INBOX', 1); + $service->entityPatch('INBOX', ['\\Seen'], [], 1); + $service->entityExtant('INBOX', 1); + + $this->assertSame(1, $this->sent('EXAMINE "INBOX"')); + $this->assertSame(1, $this->sent('SELECT "INBOX"')); + } + + public function testFailedSelectionIsNotReused(): void + { + $service = $this->service(); + $service->entityExtant('INBOX', 1); + + $this->failing = ['EXAMINE "Gone"']; + try { + $service->entityExtant('Gone', 1); + $this->fail('Expected the selection to fail'); + } catch (ImapException) { + } + + $service->entityExtant('INBOX', 1); + + $this->assertSame(2, $this->sent('EXAMINE "INBOX"')); + } + + private function sent(string $command): int + { + return count(array_filter($this->commands, static fn (string $sent): bool => $sent === $command)); + } + + private function service(): LiveMailService + { + $lines = ["* PREAUTH Ready\r\n"]; + $connection = $this->createStub(ConnectionInterface::class); + $connection->method('readLine')->willReturnCallback(static function () use (&$lines): string { + return array_shift($lines) ?? throw new \RuntimeException('Unexpected response read'); + }); + $connection->method('write')->willReturnCallback(function (string $wire) use (&$lines): void { + [$tag, $command] = explode(' ', trim($wire), 2); + $operation = explode(' ', $command)[0]; + if ($operation === 'CAPABILITY') { + $lines[] = "* CAPABILITY IMAP4rev1\r\n"; + } else { + $this->commands[] = $command; + } + if (in_array($command, $this->failing, true)) { + $lines[] = "$tag NO Mailbox does not exist\r\n"; + return; + } + if (in_array($operation, ['SELECT', 'EXAMINE'], true)) { + $lines[] = "* 3 EXISTS\r\n"; + } + $lines[] = "$tag OK Completed\r\n"; + }); + $factory = $this->createStub(ConnectionFactoryInterface::class); + $factory->method('create')->willReturn($connection); + + $client = new Client($factory); + $client->connect(new ConnectionConfig('localhost')); + + return new class($this->createStub(Service::class), $client) extends LiveMailService { + public function __construct(Service $service, private readonly Client $client) + { + parent::__construct($service); + } + + public function imapClient(): Client + { + return $this->client; + } + }; + } +} diff --git a/tests/php/Unit/MessageHarmonizationTest.php b/tests/php/Unit/MessageHarmonizationTest.php new file mode 100644 index 0000000..8cc3e35 --- /dev/null +++ b/tests/php/Unit/MessageHarmonizationTest.php @@ -0,0 +1,303 @@ +createStub(DataStore::class); + $this->mailboxes = new FakeMailboxStore($dataStore); + $this->messages = new FakeMessageStore($dataStore); + $this->files = new FakeMessageFileStore('/nonexistent'); + + $service = (new Service())->fromStore(['tid' => 'tenant', 'sid' => 'svc']); + $this->live = new MessageHarmonizationLiveStub($service); + $this->live->uidValidity = 7; + $this->live->uidNext = 100; + + $this->harmonizer = (new HarmonizationService( + $this->mailboxes, + $this->messages, + $this->files, + new MessageIngestor($this->files, $this->messages), + ))->for($service, $this->live); + + $this->mailboxes->upsert('tenant', 'svc', (new CollectionResource('imap', 'svc'))->fromImap(new Mailbox('INBOX', '/', []))); + } + + public function testFirstRunIngestsEverythingNewestFirst(): void + { + $this->live->flags = [1 => ['\\Seen'], 2 => [], 3 => ['\\Flagged']]; + + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertSame(['status' => 'harmonized', 'reset' => false, 'added' => 3, 'updated' => 0, 'removed' => 0, 'complete' => true], $result); + $this->assertSame([[3, 2, 1]], $this->live->fetchBatches); + $this->assertSame([1 => ['seen'], 2 => [], 3 => ['flagged']], $this->messages->flags('svc', 'INBOX', 7)); + $this->assertSame([1, 2, 3], $this->files->uids(7)); + $this->assertSame(262144, $this->live->bodyTextLimit); + + $state = $this->mailboxes->state('svc', 'INBOX'); + $this->assertSame(7, $state['uidValidity']); + $this->assertSame(100, $state['uidNext']); + $this->assertTrue($state['harmonizationComplete']); + $this->assertIsInt($state['harmonizedAt']); + $this->assertNull($this->mailboxes->lockOwner); + } + + public function testIncrementalRunAddsUpdatesAndRemoves(): void + { + $this->live->flags = [1 => ['\\Seen'], 2 => []]; + $this->harmonizer->harmonizeMessages('INBOX'); + + $this->live->flags = [2 => ['\\Seen', '$Forwarded'], 3 => []]; + $this->live->fetchBatches = []; + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertSame(['status' => 'harmonized', 'reset' => false, 'added' => 1, 'updated' => 1, 'removed' => 1, 'complete' => true], $result); + $this->assertSame([[3]], $this->live->fetchBatches); + $this->assertSame([2 => ['seen', '$forwarded'], 3 => []], $this->messages->flags('svc', 'INBOX', 7)); + $this->assertSame([2, 3], $this->files->uids(7)); + } + + public function testChangedUidValidityDiscardsTheCache(): void + { + $this->live->flags = [1 => [], 2 => []]; + $this->harmonizer->harmonizeMessages('INBOX'); + + $this->live->uidValidity = 8; + $this->live->flags = [1 => []]; + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertTrue($result['reset']); + $this->assertSame(1, $result['added']); + $this->assertSame([], $this->messages->flags('svc', 'INBOX', 7)); + $this->assertSame([1 => []], $this->messages->flags('svc', 'INBOX', 8)); + $this->assertSame([], $this->files->uids(7)); + $this->assertSame(8, $this->mailboxes->state('svc', 'INBOX')['uidValidity']); + } + + public function testHeldLockSkipsTheRun(): void + { + $this->live->flags = [1 => []]; + $this->mailboxes->lockOwner = 'someone-else'; + + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertSame('skipped', $result['status']); + $this->assertSame(0, $this->live->selects); + $this->assertSame('someone-else', $this->mailboxes->lockOwner); + } + + public function testMessagesThatCannotBeFetchedLeaveTheMailboxIncomplete(): void + { + $this->live->flags = [1 => [], 2 => []]; + $this->live->unfetchable = [2]; + + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertSame(1, $result['added']); + $this->assertFalse($result['complete']); + $this->assertFalse($this->mailboxes->state('svc', 'INBOX')['harmonizationComplete']); + } + + public function testNewMessagesAreFetchedInBatches(): void + { + $this->live->flags = array_fill_keys(range(1, 450), []); + + $result = $this->harmonizer->harmonizeMessages('INBOX'); + + $this->assertSame(450, $result['added']); + $this->assertSame([200, 200, 50], array_map('count', $this->live->fetchBatches)); + $this->assertSame(450, $this->live->fetchBatches[0][0]); + $this->assertSame(1, $this->live->fetchBatches[2][49]); + } + + public function testLockIsReleasedWhenTheRunFails(): void + { + $this->live->uidValidity = null; + + try { + $this->harmonizer->harmonizeMessages('INBOX'); + $this->fail('Expected missing UIDVALIDITY to fail the run'); + } catch (\RuntimeException) { + } + + $this->assertNull($this->mailboxes->lockOwner); + } +} + +final class MessageHarmonizationLiveStub extends LiveMailService +{ + public ?int $uidValidity = null; + public ?int $uidNext = null; + /** @var array> */ + public array $flags = []; + /** @var int[] */ + public array $unfetchable = []; + public array $fetchBatches = []; + public int $selects = 0; + public ?int $bodyTextLimit = null; + + public function collectionFetch(string $identifier): ?Mailbox + { + $this->selects++; + return new Mailbox($identifier, '/', [], count($this->flags), 0, $this->uidValidity, 0, [], true, $this->uidNext); + } + + public function entityFlags(string $collection): Generator + { + yield from $this->flags; + } + + public function entityFetch(string $collection, ?FetchOptions $options = null, int ...$uids): Generator + { + $this->fetchBatches[] = $uids; + if (preg_match('/BODY\.PEEK\[TEXT\]<0\.(\d+)>/', (string) $options?->toCommand(), $matches) === 1) { + $this->bodyTextLimit = (int) $matches[1]; + } + + foreach ($uids as $uid) { + if (in_array($uid, $this->unfetchable, true)) { + continue; + } + $flags = implode(' ', $this->flags[$uid] ?? []); + yield $uid => FetchMessageParser::parse("* {$uid} FETCH (UID {$uid} FLAGS ({$flags}))"); + } + } +} + +final class FakeMailboxStore extends MailboxStore +{ + /** @var array */ + public array $documents = []; + public ?string $lockOwner = null; + + public function upsert(string $tenantId, string $serviceId, CollectionResource $collection): void + { + $name = (string) $collection->identifier(); + $this->documents[$name] = ($this->documents[$name] ?? self::STATE_DEFAULTS) + ['name' => $name]; + } + + public function fetch(string $serviceId, string $name): ?array + { + return $this->documents[$name] ?? null; + } + + public function list(string $serviceId): array + { + return $this->documents; + } + + public function updateState(string $serviceId, string $name, array $state): void + { + $this->documents[$name] = array_replace($this->documents[$name], array_intersect_key($state, self::STATE_DEFAULTS)); + } + + public function acquireLock(string $serviceId, string $name, string $owner, int $ttl): bool + { + if ($this->lockOwner !== null && $this->lockOwner !== $owner) { + return false; + } + $this->lockOwner = $owner; + return true; + } + + public function releaseLock(string $serviceId, string $name, string $owner): void + { + if ($this->lockOwner === $owner) { + $this->lockOwner = null; + } + } +} + +final class FakeMessageStore extends MessageStore +{ + /** @var array>> flags by UIDVALIDITY and UID */ + private array $documents = []; + + public function upsert(string $tenantId, string $serviceId, int $uidValidity, EntityResource $entity): void + { + $meta = $entity->toCacheMeta(); + $this->documents[$uidValidity][$meta['uid']] = $meta['flags']; + ksort($this->documents[$uidValidity]); + } + + public function flags(string $serviceId, string $mailbox, int $uidValidity): array + { + return $this->documents[$uidValidity] ?? []; + } + + public function updateFlags(string $serviceId, string $mailbox, int $uidValidity, int $uid, array $flags): void + { + $this->documents[$uidValidity][$uid] = array_values($flags); + } + + public function delete(string $serviceId, string $mailbox, int $uidValidity, int ...$uids): void + { + foreach ($uids as $uid) { + unset($this->documents[$uidValidity][$uid]); + } + } + + public function deleteByMailbox(string $serviceId, string $mailbox): void + { + $this->documents = []; + } +} + +final class FakeMessageFileStore extends MessageFileStore +{ + /** @var array> content by UIDVALIDITY and UID */ + private array $files = []; + + public function write(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int $uid, array $content): void + { + $this->files[$uidValidity][$uid] = $content; + } + + public function delete(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int ...$uids): void + { + foreach ($uids as $uid) { + unset($this->files[$uidValidity][$uid]); + } + } + + public function deleteByMailbox(string $tenantId, string $serviceId, string $mailbox): void + { + $this->files = []; + } + + /** @return int[] */ + public function uids(int $uidValidity): array + { + $uids = array_keys($this->files[$uidValidity] ?? []); + sort($uids); + return $uids; + } +} diff --git a/tests/php/Unit/SelectCommandTest.php b/tests/php/Unit/SelectCommandTest.php new file mode 100644 index 0000000..225f580 --- /dev/null +++ b/tests/php/Unit/SelectCommandTest.php @@ -0,0 +1,69 @@ +client([ + "* FLAGS (\\Answered \\Flagged \\Deleted \\Seen \\Draft)\r\n", + "* OK [PERMANENTFLAGS ()] Read-only mailbox\r\n", + "* 172 EXISTS\r\n", + "* 1 RECENT\r\n", + "* OK [UIDVALIDITY 3857529045] UIDs valid\r\n", + "* OK [UIDNEXT 4392] Predicted next UID\r\n", + "* OK [HIGHESTMODSEQ 715194045007] Highest\r\n", + ])->perform(new SelectCommand('INBOX', true)); + + $this->assertSame(172, $mailbox->messages()); + $this->assertSame(1, $mailbox->recent()); + $this->assertSame(3857529045, $mailbox->uidValidity()); + $this->assertSame(4392, $mailbox->uidNext()); + $this->assertSame(715194045007, $mailbox->highestModSeq()); + } + + public function testMissingResponseCodesAreNull(): void + { + $mailbox = $this->client(["* 0 EXISTS\r\n"])->perform(new SelectCommand('INBOX', true)); + + $this->assertNull($mailbox->uidValidity()); + $this->assertNull($mailbox->uidNext()); + $this->assertNull($mailbox->highestModSeq()); + } + + /** + * @param string[] $untagged untagged lines the server sends before completing SELECT + */ + private function client(array $untagged): Client + { + $lines = ["* PREAUTH Ready\r\n"]; + $connection = $this->createStub(ConnectionInterface::class); + $connection->method('readLine')->willReturnCallback(static function () use (&$lines): string { + return array_shift($lines) ?? throw new \RuntimeException('Unexpected response read'); + }); + $connection->method('write')->willReturnCallback(static function (string $wire) use (&$lines, $untagged): void { + [$tag, $command] = explode(' ', trim($wire), 2); + if (str_starts_with($command, 'CAPABILITY')) { + $lines[] = "* CAPABILITY IMAP4rev1\r\n"; + } + if (str_starts_with($command, 'EXAMINE') || str_starts_with($command, 'SELECT')) { + array_push($lines, ...$untagged); + } + $lines[] = "$tag OK Completed\r\n"; + }); + $factory = $this->createStub(ConnectionFactoryInterface::class); + $factory->method('create')->willReturn($connection); + + $client = new Client($factory); + $client->connect(new ConnectionConfig('localhost')); + return $client; + } +}