Files
provider_imap/lib/Providers/CachedService.php
T
2026-10-07 22:20:13 -04:00

491 lines
18 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\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
{
return iterator_to_array($this->entityFetchStream(...$identifiers), true);
}
/**
* Fetch messages from the cache; messages that are not cached are fetched from the
* server and cached, mailboxes never harmonized are read from the server.
*/
public function entityFetchStream(EntityIdentifierInterface ...$identifiers): Generator
{
$byMailbox = [];
foreach ($identifiers as $identifier) {
if ($identifier->provider() !== $this->provider() || (string) $identifier->service() !== $this->serviceId()) {
throw new \InvalidArgumentException('Entity identifier does not belong to this service: ' . $identifier);
}
$byMailbox[(string) $identifier->collection()][] = (int) $identifier->entity();
}
foreach ($byMailbox as $mailbox => $uids) {
$uidValidity = $this->mailboxStore->state($this->serviceId(), (string) $mailbox)['uidValidity'];
$entities = $uidValidity === null
? $this->liveMail()->entityFetch((string) $mailbox, ...$uids)
: $this->hydrate((string) $mailbox, (int) $uidValidity, $uids);
foreach ($entities as $entity) {
yield $entity->urn() => $entity;
}
}
}
/**
* 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<int, EntityResource>
*/
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;
}
if ($entity->getProperties()->getIncompleteSections() !== []) {
$this->completeSections($mailbox, $uidValidity, $entity);
}
$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];
}
}
}
}
/**
* Fetch the full content of sections cut off at ingest and store the completed message.
*
* Content is immutable, so the meta document and change sequence are untouched. A
* failure leaves the cut off text in place; it is tried again on the next read.
*/
private function completeSections(string $mailbox, int $uidValidity, EntityResource $entity): void
{
$properties = $entity->getProperties();
$uid = (int) $entity->identifier();
try {
$sections = $this->liveMail()->messageSections($mailbox, $uid, ...$properties->getIncompleteSections());
if ($sections === []) {
return;
}
foreach ($sections as $partId => $content) {
$properties->completeSection((string) $partId, $content);
}
$this->fileStore->write($this->tenantId(), $this->serviceId(), $mailbox, $uidValidity, $uid, $entity->toCacheContent());
} catch (\Throwable) {
// keep the cut off text; retried on the next read
}
}
/**
* 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());
}
}