generated from Nodarx/template
feat: track cache changes with sequence numbers and tombstones
Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
@@ -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, '<error>failed</error>', '', '', '', '', $outcome];
|
||||
$rows[] = [$name, '<error>failed</error>', '', '', '', '', '', $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: <info>storage/%s/provider_imap/%s</info>, MongoDB <info>provider_imap_mail_mailboxes</info> / <info>provider_imap_mail_messages</info>',
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
/**
|
||||
* SPDX-FileCopyrightText: Sebastian Krupinski <krupinski01@gmail.com>
|
||||
* 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 `<uidValidity>:<changeSeq>`. 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
|
||||
* `<uidValidity>: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;
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
|
||||
+107
-15
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user