feat: harmonize mailbox messages

Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
2026-10-05 20:34:02 -04:00
parent b02ca63208
commit 65fdb23fa4
12 changed files with 879 additions and 43 deletions
+14
View File
@@ -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;
+2 -11
View File
@@ -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));
}
}
+20 -1
View File
@@ -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,
);
}
}
+17 -13
View File
@@ -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<string>
*/
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) ───────────────────────────────────
/**
+143
View File
@@ -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.
*/
+78 -17
View File
@@ -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<int, list<string>> 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;
}
}
+19
View File
@@ -145,6 +145,25 @@ class MessageStore
return $uids;
}
/**
* List the flags of every message cached for a mailbox.
*
* @return array<int, list<string>> 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.
*