* SPDX-License-Identifier: AGPL-3.0-or-later */ namespace KTXM\ProviderImap\Providers; use Generator; use KTXF\Mail\Collection\CollectionBaseInterface; use KTXF\Mail\Collection\CollectionPropertiesBaseInterface; use KTXF\Mail\Object\AddressInterface; use KTXF\Mail\Object\MessagePropertiesMutableInterface; use KTXF\Mail\Submission\EntitySubmitResult; use KTXF\Resource\BinaryResource; use KTXF\Resource\Delta\Delta; use KTXF\Resource\Filter\IFilter; use KTXF\Resource\Identifier\CollectionIdentifier; use KTXF\Resource\Identifier\EntityIdentifier; use KTXF\Resource\Identifier\EntityIdentifierInterface; use KTXF\Resource\Range\IRange; use KTXF\Resource\Range\IRangeTally; use KTXF\Resource\Range\RangeAnchorType; use KTXF\Resource\Sort\ISort; use KTXM\ProviderImap\Service\Cache\HarmonizationService; use KTXM\ProviderImap\Service\Cache\MessageDeltaService; use KTXM\ProviderImap\Service\Cache\MessageIngestor; use KTXM\ProviderImap\Service\Cache\MessageQueryBuilder; use KTXM\ProviderImap\Service\Live\LiveMailService; use KTXM\ProviderImap\Stores\MailboxStore; use KTXM\ProviderImap\Stores\MessageFileStore; use KTXM\ProviderImap\Stores\MessageStore; use UnexpectedValueException; /** * IMAP mail service that serves reads from the cache. * * Holds a LiveService for the same account: writes and uncached operations go * through it (its rules stay in one place), reads come from the cache and * harmonize with the server when stale. Extends ServiceBase rather than * LiveService, so every mail method is written out here. */ class CachedService extends ServiceBase { /** Seconds after which a mailbox (or the mailbox list) is harmonized again on access */ public const FRESHNESS_WINDOW = 60; /** Messages hydrated from the cache per meta store query */ private const HYDRATE_BATCH_SIZE = 100; private ?LiveMailService $liveMail = null; private ?LiveService $live = null; public function __construct( private readonly MailboxStore $mailboxStore, private readonly MessageStore $messageStore, private readonly MessageFileStore $fileStore, private readonly MessageIngestor $ingestor, private readonly MessageQueryBuilder $queries, private readonly HarmonizationService $harmonizer, private readonly MessageDeltaService $deltas, ) {} // ── Collections (cache) ────────────────────────────────────────────────── /** * Unfiltered lists come from the cache; filtered lists (e.g. role lookups) go to the server, * so filter behaviour stays identical to live mode. */ public function collectionList(string|int|null $location, ?IFilter $filter = null, ?ISort $sort = null): array { if ($location !== null || $filter !== null) { return $this->live()->collectionList($location, $filter, $sort); } $this->harmonizeMailboxesIfStale(); $list = []; foreach ($this->mailboxStore->list($this->serviceId()) as $name => $document) { $list[(string) $name] = $this->collectionFromCache($document); } return $list; } public function collectionExtant(string|int ...$identifiers): array { $this->harmonizeMailboxesIfStale(); $cached = $this->mailboxStore->list($this->serviceId()); $list = []; foreach ($identifiers as $identifier) { $list[(string) $identifier] = isset($cached[(string) $identifier]); } return $list; } public function collectionFetch(string|int $identifier): ?CollectionResource { $this->harmonizeMailboxesIfStale(); $document = $this->mailboxStore->fetch($this->serviceId(), (string) $identifier); return $document === null ? null : $this->collectionFromCache($document); } // ── Collections (server; cache maintenance follows in step 6) ──────────── public function collectionCreate(CollectionIdentifier|null $target, CollectionPropertiesBaseInterface $properties, array $options = []): CollectionBaseInterface { return $this->live()->collectionCreate($target, $properties, $options); } public function collectionUpdate(CollectionIdentifier $target, CollectionPropertiesBaseInterface $properties): CollectionBaseInterface { return $this->live()->collectionUpdate($target, $properties); } public function collectionDelete(CollectionIdentifier $target, bool $force = false): CollectionBaseInterface | true { return $this->live()->collectionDelete($target, $force); } public function collectionMove(CollectionIdentifier $target, CollectionIdentifier $source): CollectionBaseInterface { return $this->live()->collectionMove($target, $source); } // ── Entities: reads ────────────────────────────────────────────────────── public function entityListBulk(string|int $collection, ?IFilter $filter = null, ?ISort $sort = null, ?IRange $range = null, ?array $properties = null): array { return iterator_to_array($this->entityListStream($collection, $filter, $sort, $range, $properties), true); } /** * List messages of a mailbox. * * - mailbox not harmonized yet: streamed from the server, each message cached as it passes * - harmonized: UIDs from the meta store (filter / sort translated, range applied like the * live service), content from message.json; a filter on body / full text has the server * find the UIDs (IMAP SEARCH) and only the content comes from the cache */ public function entityListStream(string|int $collection, ?IFilter $filter = null, ?ISort $sort = null, ?IRange $range = null, ?array $properties = null): Generator { $mailbox = (string) $collection; $state = $this->mailboxStore->state($this->serviceId(), $mailbox); if (!$state['harmonizationComplete'] || $state['uidValidity'] === null) { yield from $this->listFromServer($mailbox, $filter, $sort, $range); return; } $uidValidity = (int) $state['uidValidity']; if ($this->queries->needsServer($filter)) { $uids = $this->liveMail()->entityFind($mailbox, $filter, $sort, $range); } else { $order = $this->queries->sort($sort); $uids = self::applyRange( $this->messageStore->query($this->serviceId(), $mailbox, $uidValidity, $this->queries->filter($filter), $order['sort'], $order['collation']), $range, ); } foreach ($this->hydrate($mailbox, $uidValidity, $uids) as $entity) { yield $entity->urn() => $entity; } } public function entityFetchBulk(EntityIdentifierInterface ...$identifiers): array { // served from the cache in 5.2c return $this->live()->entityFetchBulk(...$identifiers); } public function entityFetchStream(EntityIdentifierInterface ...$identifiers): Generator { // served from the cache in 5.2c return $this->live()->entityFetchStream(...$identifiers); } /** * Changes since a signature, harmonizing the mailbox first when stale. * * Harmonization never waits on another run's lock, and a server that cannot be * reached does not fail the request: the delta is answered from what is cached. */ public function entityDelta(string|int $collection, string $signature, string $detail = 'ids'): Delta { $mailbox = (string) $collection; if ($this->isStale($this->mailboxStore->state($this->serviceId(), $mailbox)['harmonizedAt'])) { try { $this->harmonizer()->harmonizeMessages($mailbox); } catch (\Throwable) { // answer from the cache } } return $this->deltas->delta($this->serviceId(), $mailbox, $signature); } public function entityExtant(string|int $collection, string|int ...$identifiers): array { return $this->live()->entityExtant($collection, ...$identifiers); } public function entityDownload(EntityIdentifierInterface $target, array|null $part): BinaryResource { return $this->live()->entityDownload($target, $part); } // ── Entities: writes (server; cache maintenance follows in step 6) ─────── public function entitySubmit(AddressInterface $sender, EntityIdentifierInterface|null $source = null, MessagePropertiesMutableInterface|null $message = null): EntitySubmitResult { return $this->live()->entitySubmit($sender, $source, $message); } public function entityCreate(CollectionIdentifier $target, MessagePropertiesMutableInterface $properties, array $options = []): EntityResource { return $this->live()->entityCreate($target, $properties, $options); } public function entityModify(EntityIdentifier $target, MessagePropertiesMutableInterface $properties): EntityResource { return $this->live()->entityModify($target, $properties); } public function entityPatch(MessagePropertiesMutableInterface $properties, EntityIdentifier ...$targets): array { return $this->live()->entityPatch($properties, ...$targets); } public function entityDelete(EntityIdentifier ...$targets): array { return $this->live()->entityDelete(...$targets); } public function entityMove(CollectionIdentifier $target, EntityIdentifier ...$sources): array { return $this->live()->entityMove($target, ...$sources); } public function entityCopy(CollectionIdentifier $target, EntityIdentifier ...$sources): array { return $this->live()->entityCopy($target, ...$sources); } // ── Internals ──────────────────────────────────────────────────────────── /** * Stream a list from the server, caching each message as it passes (cold mailbox). */ private function listFromServer(string $mailbox, ?IFilter $filter, ?ISort $sort, ?IRange $range): Generator { $uidValidity = $this->cacheGeneration($mailbox); foreach ($this->liveMail()->entityList($mailbox, $filter, $sort, $range) as $entity) { if ($uidValidity !== null) { $this->ingestQuietly($uidValidity, $entity); } yield $entity->urn() => $entity; } } /** * Entities for UIDs in the given order: cached content + meta flags, with messages * missing from the cache (or written with another schema version) fetched from the * server and cached. * * @param int[] $uids * @return Generator */ private function hydrate(string $mailbox, int $uidValidity, array $uids): Generator { foreach (array_chunk($uids, self::HYDRATE_BATCH_SIZE) as $batch) { $metas = $this->messageStore->fetchMany($this->serviceId(), $mailbox, $uidValidity, ...$batch); $entities = []; $missing = []; foreach ($batch as $uid) { $entity = isset($metas[$uid]) ? $this->entityFromCache($mailbox, $uidValidity, $metas[$uid]) : null; if ($entity === null) { $missing[] = $uid; continue; } $entities[$uid] = $entity; } if ($missing !== []) { foreach ($this->liveMail()->entityFetch($mailbox, ...$missing) as $uid => $entity) { $this->ingestQuietly($uidValidity, $entity); $entities[$uid] = $entity; } } foreach ($batch as $uid) { if (isset($entities[$uid])) { yield $uid => $entities[$uid]; } } } } /** * An entity from its meta document and message.json; null when the content is missing or outdated. */ private function entityFromCache(string $mailbox, int $uidValidity, array $meta): ?EntityResource { $content = $this->fileStore->read($this->tenantId(), $this->serviceId(), $mailbox, $uidValidity, (int) $meta['uid']); if ($content === null) { return null; } try { return $this->entityFresh()->fromCacheContent($content)->fromCacheMeta($meta); } catch (UnexpectedValueException) { return null; } } /** * UIDVALIDITY to cache a not yet harmonized mailbox under; null when it cannot be cached. */ private function cacheGeneration(string $mailbox): ?int { // the mailbox document holds the change sequence, so it has to exist before ingesting if ($this->mailboxStore->fetch($this->serviceId(), $mailbox) === null) { try { $this->harmonizer()->harmonizeMailboxes(); } catch (\Throwable) { return null; } if ($this->mailboxStore->fetch($this->serviceId(), $mailbox) === null) { return null; } } $state = $this->mailboxStore->state($this->serviceId(), $mailbox); if ($state['uidValidity'] !== null) { return (int) $state['uidValidity']; } return $this->liveMail()->mailboxFetch($mailbox)?->uidValidity(); } /** * Cache a message; a failing cache write must not fail the read that triggered it. */ private function ingestQuietly(int $uidValidity, EntityResource $entity): void { try { $this->ingestor->ingest($this->tenantId(), $this->serviceId(), $uidValidity, $entity); } catch (\Throwable) { // the next harmonization caches it } } /** * Apply a list range to sorted UIDs, as the live service does: absolute skips `position` * messages, relative starts at the UID given as `position`. * * @param int[] $uids * @return int[] */ private static function applyRange(array $uids, ?IRange $range): array { if (!$range instanceof IRangeTally) { return array_values($uids); } $tally = max(0, $range->getTally()); if ($tally === 0) { return []; } if ($range->getAnchor() === RangeAnchorType::ABSOLUTE) { $start = max(0, (int) $range->getPosition()); } else { $index = array_search((int) $range->getPosition(), $uids, true); $start = $index === false ? 0 : $index; } return array_values(array_slice($uids, $start, $tally)); } /** * Build a collection from its cached document; its signature is the delta signature of the mailbox. */ private function collectionFromCache(array $document): CollectionResource { if (isset($document['uidValidity'])) { $document['signature'] = MessageDeltaService::signature((int) $document['uidValidity'], (int) ($document['changeSeq'] ?? 0)); } return $this->collectionFresh()->fromCacheMeta($document); } private function harmonizeMailboxesIfStale(): void { if ($this->isStale($this->mailboxStore->listedAt($this->serviceId()))) { $this->harmonizer()->harmonizeMailboxes(); } } private function isStale(?int $timestamp): bool { return $timestamp === null || time() - $timestamp >= self::FRESHNESS_WINDOW; } private function serviceId(): string { return (string) $this->identifier(); } private function tenantId(): string { return (string) $this->tenantIdentifier(); } /** * One IMAP connection per service object, shared by the LiveService and harmonization. */ protected function liveMail(): LiveMailService { return $this->liveMail ??= new LiveMailService($this); } private function live(): LiveService { return $this->live ??= (new LiveService($this->liveMail()))->fromStore($this->toStore()); } private function harmonizer(): HarmonizationService { return $this->harmonizer->for($this, $this->liveMail()); } }