Files
provider_imap/lib/Stores/MessageStore.php
T
2026-10-07 21:21:14 -04:00

308 lines
11 KiB
PHP

<?php
declare(strict_types=1);
/**
* SPDX-FileCopyrightText: Sebastian Krupinski <krupinski01@gmail.com>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
namespace KTXM\ProviderImap\Stores;
use KTXC\Db\DataStore;
use KTXM\ProviderImap\Providers\EntityResource;
/**
* IMAP Message Meta Store
*
* 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,
) {}
/**
* Create the collection indexes.
*
* MongoDB createIndex is idempotent when the name and specification match.
*
* @return string[]
*/
public function ensureIndexes(): array
{
$collection = $this->dataStore->selectCollection(self::COLLECTION_NAME);
return [
$collection->createIndex(
['sid' => 1, 'mailbox' => 1, 'uidValidity' => 1, 'uid' => 1],
['name' => 'messages_key', 'unique' => true]
),
$collection->createIndex(
['sid' => 1, 'mailbox' => 1, 'received' => -1],
['name' => 'messages_by_received']
),
$collection->createIndex(
['sid' => 1, 'mailbox' => 1, 'sent' => -1],
['name' => 'messages_by_sent']
),
$collection->createIndex(
['sid' => 1, 'mailbox' => 1, 'flags' => 1],
['name' => 'messages_by_flags']
),
$collection->createIndex(
['sid' => 1, 'urid' => 1],
['name' => 'messages_by_urid']
),
$collection->createIndex(
['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, int $seq): void
{
$document = $entity->toCacheMeta();
$key = [
'sid' => $serviceId,
'mailbox' => $document['mailbox'],
'uidValidity' => $uidValidity,
'uid' => $document['uid'],
];
$this->dataStore->selectCollection(self::COLLECTION_NAME)->updateOne(
$key,
[
'$set' => ['tid' => $tenantId, ...$key, ...$document, 'modSeq' => $seq],
'$setOnInsert' => ['addedSeq' => $seq],
],
['upsert' => true],
);
}
/**
* Retrieve the meta document of a message.
*/
public function fetch(string $serviceId, string $mailbox, int $uidValidity, int $uid): ?array
{
return $this->dataStore->selectCollection(self::COLLECTION_NAME)->findOne([
'sid' => $serviceId,
'mailbox' => $mailbox,
'uidValidity' => $uidValidity,
'uid' => $uid,
...self::LIVE,
]);
}
/**
* Retrieve the meta documents of several messages.
*
* @return array<int, array> keyed by UID; UIDs that are not cached are absent
*/
public function fetchMany(string $serviceId, string $mailbox, int $uidValidity, int ...$uids): array
{
if ($uids === []) {
return [];
}
$cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find([
'sid' => $serviceId,
'mailbox' => $mailbox,
'uidValidity' => $uidValidity,
'uid' => ['$in' => array_values($uids)],
...self::LIVE,
]);
$list = [];
foreach ($cursor as $document) {
$list[(int) $document['uid']] = $document;
}
return $list;
}
/**
* List the UIDs cached for a mailbox.
*
* @return int[]
*/
public function uids(string $serviceId, string $mailbox, int $uidValidity): array
{
$cursor = $this->dataStore->selectCollection(self::COLLECTION_NAME)->find(
['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, ...self::LIVE],
['projection' => ['uid' => 1]],
);
$uids = [];
foreach ($cursor as $document) {
$uids[] = (int) $document['uid'];
}
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, ...self::LIVE],
['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 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, int $seq): void
{
$this->dataStore->selectCollection(self::COLLECTION_NAME)->updateOne(
['sid' => $serviceId, 'mailbox' => $mailbox, 'uidValidity' => $uidValidity, 'uid' => $uid, ...self::LIVE],
['$set' => ['flags' => array_values($flags), 'modSeq' => $seq]],
);
}
/**
* Turn messages into tombstones: their list fields are dropped and removedSeq is set.
*/
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,
'removedSeq' => ['$lte' => $upTo],
]);
}
/**
* Delete all meta documents of a mailbox (any UIDVALIDITY), tombstones included.
*/
public function deleteByMailbox(string $serviceId, string $mailbox): void
{
$this->dataStore->selectCollection(self::COLLECTION_NAME)->deleteMany([
'sid' => $serviceId,
'mailbox' => $mailbox,
]);
}
/**
* Delete all meta documents of a service.
*/
public function deleteByService(string $serviceId): void
{
$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;
}
}