generated from Nodarx/template
feat: add CachedService with cached collections and delta
Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
@@ -0,0 +1,260 @@
|
||||
<?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\Sort\ISort;
|
||||
use KTXM\ProviderImap\Service\Cache\HarmonizationService;
|
||||
use KTXM\ProviderImap\Service\Cache\MessageDeltaService;
|
||||
use KTXM\ProviderImap\Service\Live\LiveMailService;
|
||||
use KTXM\ProviderImap\Stores\MailboxStore;
|
||||
|
||||
/**
|
||||
* 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;
|
||||
|
||||
private ?LiveMailService $liveMail = null;
|
||||
private ?LiveService $live = null;
|
||||
|
||||
public function __construct(
|
||||
private readonly MailboxStore $mailboxStore,
|
||||
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
|
||||
{
|
||||
// served from the cache in 5.2b
|
||||
return $this->live()->entityListBulk($collection, $filter, $sort, $range, $properties);
|
||||
}
|
||||
|
||||
public function entityListStream(string|int $collection, ?IFilter $filter = null, ?ISort $sort = null, ?IRange $range = null, ?array $properties = null): Generator
|
||||
{
|
||||
// served from the cache in 5.2b
|
||||
return $this->live()->entityListStream($collection, $filter, $sort, $range, $properties);
|
||||
}
|
||||
|
||||
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 ────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* 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();
|
||||
}
|
||||
|
||||
/**
|
||||
* 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());
|
||||
}
|
||||
}
|
||||
@@ -55,6 +55,16 @@ class LiveService extends ServiceBase
|
||||
{
|
||||
private LiveMailService $mailService;
|
||||
|
||||
/**
|
||||
* @param LiveMailService|null $mailService share an existing IMAP connection (e.g. CachedService); created on demand otherwise
|
||||
*/
|
||||
public function __construct(?LiveMailService $mailService = null)
|
||||
{
|
||||
if ($mailService !== null) {
|
||||
$this->mailService = $mailService;
|
||||
}
|
||||
}
|
||||
|
||||
public function collectionList(string|int|null $location, ?IFilter $filter = null, ?ISort $sort = null): array
|
||||
{
|
||||
$this->initialize();
|
||||
|
||||
Reference in New Issue
Block a user