feat: implement content cache

Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
2026-09-29 19:56:30 -04:00
parent d28800a04e
commit 0250b5a847
10 changed files with 511 additions and 8 deletions
+13
View File
@@ -16,6 +16,7 @@ final class Message
* @param list<MessageAddress> $bcc * @param list<MessageAddress> $bcc
* @param array<string, string> $bodySections * @param array<string, string> $bodySections
* @param list<string> $references message ids from the References header, without angle brackets * @param list<string> $references message ids from the References header, without angle brackets
* @param list<string> $truncatedSections part ids of text sections cut off by a partial fetch
*/ */
public function __construct( public function __construct(
private readonly int $sequence, private readonly int $sequence,
@@ -37,6 +38,7 @@ final class Message
private readonly ?MessagePart $bodyStructure, private readonly ?MessagePart $bodyStructure,
private readonly array $bodySections, private readonly array $bodySections,
private readonly array $references = [], private readonly array $references = [],
private readonly array $truncatedSections = [],
) {} ) {}
public function sequence(): int public function sequence(): int
@@ -166,6 +168,16 @@ final class Message
return $this->bodySections; return $this->bodySections;
} }
/**
* Part ids of text sections that were cut off by a partial fetch.
*
* @return list<string>
*/
public function truncatedSections(): array
{
return $this->truncatedSections;
}
public function bodyRaw(): ?string public function bodyRaw(): ?string
{ {
return $this->bodySections[''] ?? null; return $this->bodySections[''] ?? null;
@@ -196,6 +208,7 @@ final class Message
$bodyStructure, $bodyStructure,
$bodySections, $bodySections,
$this->references, $this->references,
$this->truncatedSections,
); );
} }
} }
@@ -33,7 +33,8 @@ final class FetchMessageParser
$uid = self::toInt($attributes['UID'] ?? null, 'FETCH response is missing UID: ' . $raw); $uid = self::toInt($attributes['UID'] ?? null, 'FETCH response is missing UID: ' . $raw);
$envelope = is_array($attributes['ENVELOPE'] ?? null) ? $attributes['ENVELOPE'] : null; $envelope = is_array($attributes['ENVELOPE'] ?? null) ? $attributes['ENVELOPE'] : null;
$bodyStructure = isset($attributes['BODYSTRUCTURE']) ? self::parseBodyPart($attributes['BODYSTRUCTURE'], '') : null; $bodyStructure = isset($attributes['BODYSTRUCTURE']) ? self::parseBodyPart($attributes['BODYSTRUCTURE'], '') : null;
$bodySections = self::parseBodySections($attributes, $bodyStructure); $truncatedSections = [];
$bodySections = self::parseBodySections($attributes, $bodyStructure, $truncatedSections);
$headers = self::parseFetchedHeaders($attributes); $headers = self::parseFetchedHeaders($attributes);
return new Message( return new Message(
@@ -56,6 +57,7 @@ final class FetchMessageParser
$bodyStructure, $bodyStructure,
$bodySections, $bodySections,
self::extractReferences($headers), self::extractReferences($headers),
$truncatedSections,
); );
} }
@@ -543,7 +545,7 @@ final class FetchMessageParser
* @param array<string, mixed> $attributes * @param array<string, mixed> $attributes
* @return array<string, string> * @return array<string, string>
*/ */
private static function parseBodySections(array $attributes, ?MessagePart $bodyStructure = null): array private static function parseBodySections(array $attributes, ?MessagePart $bodyStructure = null, array &$truncated = []): array
{ {
$sections = []; $sections = [];
@@ -570,7 +572,7 @@ final class FetchMessageParser
} }
if ($bodyStructure->isMultipart()) { if ($bodyStructure->isMultipart()) {
$derivedSections = self::sectionsFromBodyText($sections['TEXT'], $bodyStructure); $derivedSections = self::sectionsFromBodyText($sections['TEXT'], $bodyStructure, $truncated);
unset($sections['TEXT']); unset($sections['TEXT']);
foreach ($derivedSections as $section => $content) { foreach ($derivedSections as $section => $content) {
@@ -581,6 +583,11 @@ final class FetchMessageParser
} }
if (str_starts_with($bodyStructure->mimeType(), 'text/')) { if (str_starts_with($bodyStructure->mimeType(), 'text/')) {
// a single-part body has no closing boundary; a partial fetch shows as fewer octets than declared
$declaredSize = $bodyStructure->size();
if ($declaredSize !== null && strlen($sections['TEXT']) < $declaredSize) {
$truncated[] = $bodyStructure->partId();
}
$sections[$bodyStructure->partId()] ??= $sections['TEXT']; $sections[$bodyStructure->partId()] ??= $sections['TEXT'];
unset($sections['TEXT']); unset($sections['TEXT']);
} }
@@ -616,7 +623,7 @@ final class FetchMessageParser
/** /**
* @return array<string, string> * @return array<string, string>
*/ */
private static function sectionsFromBodyText(string $content, MessagePart $part): array private static function sectionsFromBodyText(string $content, MessagePart $part, array &$truncatedParts = []): array
{ {
if ($part->isMultipart()) { if ($part->isMultipart()) {
$boundary = $part->parameters()['boundary'] ?? ''; $boundary = $part->parameters()['boundary'] ?? '';
@@ -633,7 +640,7 @@ final class FetchMessageParser
} }
$segmentTruncated = $truncated && $index === $lastIndex; $segmentTruncated = $truncated && $index === $lastIndex;
foreach (self::sectionsFromMimeEntity($segments[$index], $childPart, $segmentTruncated) as $section => $childContent) { foreach (self::sectionsFromMimeEntity($segments[$index], $childPart, $segmentTruncated, $truncatedParts) as $section => $childContent) {
$sections[$section] = $childContent; $sections[$section] = $childContent;
} }
} }
@@ -651,7 +658,7 @@ final class FetchMessageParser
/** /**
* @return array<string, string> * @return array<string, string>
*/ */
private static function sectionsFromMimeEntity(string $content, MessagePart $part, bool $truncated = false): array private static function sectionsFromMimeEntity(string $content, MessagePart $part, bool $truncated = false, array &$truncatedParts = []): array
{ {
// a truncated entity cut off inside its headers has no usable body // a truncated entity cut off inside its headers has no usable body
if ($truncated && !str_contains($content, "\r\n\r\n") && !str_contains($content, "\n\n")) { if ($truncated && !str_contains($content, "\r\n\r\n") && !str_contains($content, "\n\n")) {
@@ -661,13 +668,17 @@ final class FetchMessageParser
[, $body] = self::splitMimeEntity($content); [, $body] = self::splitMimeEntity($content);
if ($part->isMultipart()) { if ($part->isMultipart()) {
return self::sectionsFromBodyText($body, $part); return self::sectionsFromBodyText($body, $part, $truncatedParts);
} }
if (!str_starts_with($part->mimeType(), 'text/')) { if (!str_starts_with($part->mimeType(), 'text/')) {
return []; return [];
} }
if ($truncated) {
$truncatedParts[] = $part->partId();
}
return [$part->partId() => $body]; return [$part->partId() => $body];
} }
+4 -1
View File
@@ -105,6 +105,7 @@ class EntityResource extends EntityMutableAbstract {
return [ return [
'schemaVersion' => self::CACHE_SCHEMA_VERSION, 'schemaVersion' => self::CACHE_SCHEMA_VERSION,
'created' => $this->data['created'] ?? null, 'created' => $this->data['created'] ?? null,
'incomplete' => $this->getProperties()->getIncompleteSections(),
'properties' => $this->getProperties()->toCacheContent(), 'properties' => $this->getProperties()->toCacheContent(),
]; ];
} }
@@ -124,7 +125,9 @@ class EntityResource extends EntityMutableAbstract {
$this->data['created'] = $content['created']; $this->data['created'] = $content['created'];
} }
$this->getProperties()->fromCacheContent($content['properties'] ?? []); $this->getProperties()
->fromCacheContent($content['properties'] ?? [])
->setIncompleteSections(...array_map('strval', $content['incomplete'] ?? []));
return $this; return $this;
} }
+24
View File
@@ -20,12 +20,20 @@ use KTXF\Mail\Object\MessagePropertiesMutableAbstract;
*/ */
class MessageProperties extends MessagePropertiesMutableAbstract { class MessageProperties extends MessagePropertiesMutableAbstract {
/**
* Part ids of text sections that were cut off when fetched (cache only, not part of the API shape)
*
* @var list<string>
*/
private array $incompleteSections = [];
/** /**
* Convert IMAP data to mail message properties object. * Convert IMAP data to mail message properties object.
*/ */
public function fromImap(Message $message): static public function fromImap(Message $message): static
{ {
$this->data[static::PROPERTY_SIZE] = $message->size(); $this->data[static::PROPERTY_SIZE] = $message->size();
$this->incompleteSections = $message->truncatedSections();
if ($message->messageId() !== null) { if ($message->messageId() !== null) {
$this->data[static::PROPERTY_URID] = $message->messageId(); $this->data[static::PROPERTY_URID] = $message->messageId();
@@ -111,6 +119,22 @@ class MessageProperties extends MessagePropertiesMutableAbstract {
// ── Cache (meta store / content store) ─────────────────────────────────── // ── Cache (meta store / content store) ───────────────────────────────────
/**
* Part ids of text sections that were cut off when fetched and must be completed on open.
*
* @return list<string>
*/
public function getIncompleteSections(): array
{
return $this->incompleteSections;
}
public function setIncompleteSections(string ...$partIds): static
{
$this->incompleteSections = array_values(array_unique($partIds));
return $this;
}
/** /**
* Serialise the fields the meta store needs for list, filter and sort. * Serialise the fields the meta store needs for list, filter and sort.
* *
+55
View File
@@ -0,0 +1,55 @@
<?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 KTXM\ProviderImap\Providers\EntityResource;
use KTXM\ProviderImap\Stores\MessageFileStore;
use KTXM\ProviderImap\Stores\MessageStore;
/**
* Writes fetched messages into the cache (content store + meta store).
*
* Shared by harmonization and by the cold list path, which caches the page it streams.
*/
class MessageIngestor
{
/** Maximum octets of BODY[TEXT] fetched per message for the cache; text past it is completed on open */
public const BODY_TEXT_LIMIT = 262144; // 256 KB
public function __construct(
private readonly MessageFileStore $fileStore,
private readonly MessageStore $messageStore,
) {}
/**
* 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.
*
* @return int[] UIDs that were cached
*/
public function ingest(string $tenantId, string $serviceId, int $uidValidity, EntityResource ...$entities): array
{
$ingested = [];
foreach ($entities as $entity) {
$mailbox = (string) $entity->collection();
$uid = (int) $entity->identifier();
$this->fileStore->write($tenantId, $serviceId, $mailbox, $uidValidity, $uid, $entity->toCacheContent());
$this->messageStore->upsert($tenantId, $serviceId, $uidValidity, $entity);
$ingested[] = $uid;
}
return $ingested;
}
}
+182
View File
@@ -0,0 +1,182 @@
<?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 DI\Attribute\Inject;
use FilesystemIterator;
use InvalidArgumentException;
use JsonException;
use RecursiveDirectoryIterator;
use RecursiveIteratorIterator;
use RuntimeException;
/**
* IMAP Message Content Store
*
* Holds the immutable content of cached messages (message.json) on the filesystem:
*
* {rootDir}/storage/{tid}/provider_imap/{sid}/{sha1(mailbox)}/{uidValidity}/{uid}/message.json
*
* Mailbox names are hashed because they contain delimiters and modified UTF-7.
* Files are written atomically (temporary file + rename).
*/
class MessageFileStore
{
private const CONTENT_FILE = 'message.json';
private const FOLDER_PERMISSIONS = 0750;
private const FILE_PERMISSIONS = 0640;
public function __construct(
#[Inject('rootDir')] private readonly string $rootDir,
) {}
/**
* Write the content of a message, replacing any existing content.
*/
public function write(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int $uid, array $content): void
{
$folder = $this->messagePath($tenantId, $serviceId, $mailbox, $uidValidity, $uid);
if (!is_dir($folder) && !@mkdir($folder, self::FOLDER_PERMISSIONS, true) && !is_dir($folder)) {
throw new RuntimeException("Failed to create message folder: {$folder}");
}
try {
$json = json_encode($content, JSON_THROW_ON_ERROR | JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES | JSON_INVALID_UTF8_SUBSTITUTE);
} catch (JsonException $e) {
throw new RuntimeException('Failed to encode message content: ' . $e->getMessage(), 0, $e);
}
$file = $folder . DIRECTORY_SEPARATOR . self::CONTENT_FILE;
$temporary = $file . '.' . bin2hex(random_bytes(6)) . '.tmp';
if (@file_put_contents($temporary, $json, LOCK_EX) !== strlen($json)) {
@unlink($temporary);
throw new RuntimeException("Failed to write message content: {$file}");
}
@chmod($temporary, self::FILE_PERMISSIONS);
if (!@rename($temporary, $file)) {
@unlink($temporary);
throw new RuntimeException("Failed to replace message content: {$file}");
}
}
/**
* Read the content of a message.
*
* @return array|null null when the message is not cached or its content is unreadable
*/
public function read(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int $uid): ?array
{
$file = $this->messagePath($tenantId, $serviceId, $mailbox, $uidValidity, $uid) . DIRECTORY_SEPARATOR . self::CONTENT_FILE;
if (!is_file($file)) {
return null;
}
$json = @file_get_contents($file);
if ($json === false) {
return null;
}
try {
$content = json_decode($json, true, 512, JSON_THROW_ON_ERROR);
} catch (JsonException) {
return null;
}
return is_array($content) ? $content : null;
}
/**
* Delete the cached files of messages.
*/
public function delete(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int ...$uids): void
{
foreach ($uids as $uid) {
$this->removeTree($this->messagePath($tenantId, $serviceId, $mailbox, $uidValidity, $uid));
}
}
/**
* Delete the cached files of a mailbox (all UIDVALIDITY generations).
*/
public function deleteByMailbox(string $tenantId, string $serviceId, string $mailbox): void
{
$this->removeTree($this->mailboxPath($tenantId, $serviceId, $mailbox));
}
/**
* Delete the cached files of a service.
*/
public function deleteByService(string $tenantId, string $serviceId): void
{
$this->removeTree($this->servicePath($tenantId, $serviceId));
}
private function servicePath(string $tenantId, string $serviceId): string
{
return implode(DIRECTORY_SEPARATOR, [
rtrim($this->rootDir, DIRECTORY_SEPARATOR),
'storage',
self::pathSegment($tenantId, 'tenant'),
'provider_imap',
self::pathSegment($serviceId, 'service'),
]);
}
private function mailboxPath(string $tenantId, string $serviceId, string $mailbox): string
{
return $this->servicePath($tenantId, $serviceId) . DIRECTORY_SEPARATOR . sha1($mailbox);
}
private function messagePath(string $tenantId, string $serviceId, string $mailbox, int $uidValidity, int $uid): string
{
if ($uidValidity < 0 || $uid <= 0) {
throw new InvalidArgumentException('Invalid UIDVALIDITY or UID');
}
return $this->mailboxPath($tenantId, $serviceId, $mailbox)
. DIRECTORY_SEPARATOR . $uidValidity
. DIRECTORY_SEPARATOR . $uid;
}
/**
* Guard identifiers used as path segments against traversal.
*/
private static function pathSegment(string $value, string $label): string
{
if (preg_match('/^[A-Za-z0-9_-]+$/', $value) !== 1) {
throw new InvalidArgumentException("Invalid {$label} identifier for cache path");
}
return $value;
}
private function removeTree(string $path): void
{
if (!file_exists($path)) {
return;
}
if (!is_dir($path) || is_link($path)) {
@unlink($path);
return;
}
$items = new RecursiveIteratorIterator(
new RecursiveDirectoryIterator($path, FilesystemIterator::SKIP_DOTS),
RecursiveIteratorIterator::CHILD_FIRST,
);
foreach ($items as $item) {
$item->isDir() && !$item->isLink() ? @rmdir($item->getPathname()) : @unlink($item->getPathname());
}
@rmdir($path);
}
}
+14
View File
@@ -81,6 +81,20 @@ final class CacheSerializationTest extends TestCase
$this->assertSame(['flagged' => true], $restored->getProperties()->getFlags()); $this->assertSame(['flagged' => true], $restored->getProperties()->getFlags());
} }
public function testIncompleteSectionsRoundTripThroughContentButNotTheApi(): void
{
$entity = $this->entity();
$entity->getProperties()->setIncompleteSections('1');
$content = json_decode(json_encode($entity->toCacheContent()), true);
$restored = (new EntityResource('imap', 'svc'))->fromCacheContent($content);
$this->assertSame(['1'], $content['incomplete']);
$this->assertArrayNotHasKey('incomplete', $content['properties']);
$this->assertSame(['1'], $restored->getProperties()->getIncompleteSections());
$this->assertStringNotContainsString('incomplete', json_encode($restored));
}
public function testContentWithOtherSchemaVersionIsRejected(): void public function testContentWithOtherSchemaVersionIsRejected(): void
{ {
$content = $this->entity()->toCacheContent(); $content = $this->entity()->toCacheContent();
+34
View File
@@ -78,6 +78,40 @@ final class FetchMessageParserTest extends TestCase
$this->assertSame(['1' => 'Hello wor'], $message->bodySections()); $this->assertSame(['1' => 'Hello wor'], $message->bodySections());
} }
public function testTruncatedSectionsAreReported(): void
{
$complete = FetchMessageParser::parse($this->fetchResponse(
self::MIXED_STRUCTURE,
'BODY[TEXT]',
self::ALTERNATIVE_PREFIX . "<p>Hello</p>\r\n--alt--\r\n--mix--\r\n",
));
$cutInHtml = FetchMessageParser::parse($this->fetchResponse(
self::MIXED_STRUCTURE,
'BODY[TEXT]<0>',
self::ALTERNATIVE_PREFIX . '<p>Hello wor',
));
$cutInAttachment = FetchMessageParser::parse($this->fetchResponse(
self::MIXED_STRUCTURE,
'BODY[TEXT]<0>',
self::ALTERNATIVE_PREFIX . "<p>Hello</p>\r\n--alt--\r\n--mix\r\nContent-Type: application/pdf\r\n\r\nJVBE",
));
$this->assertSame([], $complete->truncatedSections());
$this->assertSame(['1.2'], $cutInHtml->truncatedSections());
$this->assertSame([], $cutInAttachment->truncatedSections());
}
public function testTruncatedSinglePartBodyIsReportedBySize(): void
{
$structure = '("TEXT" "PLAIN" ("CHARSET" "UTF-8") NIL NIL "7BIT" 11 1)';
$partial = FetchMessageParser::parse($this->fetchResponse($structure, 'BODY[TEXT]<0>', 'Hello'));
$complete = FetchMessageParser::parse($this->fetchResponse($structure, 'BODY[TEXT]', 'Hello world'));
$this->assertSame(['1'], $partial->truncatedSections());
$this->assertSame([], $complete->truncatedSections());
}
private function fetchResponse(string $bodyStructure, string $section, string $body): string private function fetchResponse(string $bodyStructure, string $section, string $body): string
{ {
return sprintf( return sprintf(
+103
View File
@@ -0,0 +1,103 @@
<?php
declare(strict_types=1);
namespace KTXT\ProviderImap\Tests\Unit;
use FilesystemIterator;
use InvalidArgumentException;
use KTXM\ProviderImap\Stores\MessageFileStore;
use PHPUnit\Framework\TestCase;
use RecursiveDirectoryIterator;
use RecursiveIteratorIterator;
final class MessageFileStoreTest extends TestCase
{
private string $root;
private MessageFileStore $store;
protected function setUp(): void
{
$this->root = sys_get_temp_dir() . '/ktrix-imap-cache-' . bin2hex(random_bytes(4));
$this->store = new MessageFileStore($this->root);
}
protected function tearDown(): void
{
if (!is_dir($this->root)) {
return;
}
$items = new RecursiveIteratorIterator(
new RecursiveDirectoryIterator($this->root, FilesystemIterator::SKIP_DOTS),
RecursiveIteratorIterator::CHILD_FIRST,
);
foreach ($items as $item) {
$item->isDir() ? rmdir($item->getPathname()) : unlink($item->getPathname());
}
rmdir($this->root);
}
public function testWriteThenReadReturnsContent(): void
{
$content = ['schemaVersion' => 1, 'properties' => ['subject' => 'Grüße / ✓']];
$this->store->write('tenant', 'svc', 'INBOX/Reports', 7, 42, $content);
$this->assertSame($content, $this->store->read('tenant', 'svc', 'INBOX/Reports', 7, 42));
$this->assertFileExists(sprintf(
'%s/storage/tenant/provider_imap/svc/%s/7/42/message.json',
$this->root,
sha1('INBOX/Reports'),
));
}
public function testWriteReplacesAndLeavesNoTemporaryFiles(): void
{
$this->store->write('tenant', 'svc', 'INBOX', 7, 42, ['v' => 1]);
$this->store->write('tenant', 'svc', 'INBOX', 7, 42, ['v' => 2]);
$this->assertSame(['v' => 2], $this->store->read('tenant', 'svc', 'INBOX', 7, 42));
$folder = sprintf('%s/storage/tenant/provider_imap/svc/%s/7/42', $this->root, sha1('INBOX'));
$this->assertSame(['message.json'], array_values(array_diff(scandir($folder), ['.', '..'])));
}
public function testReadOfMissingOrCorruptContentReturnsNull(): void
{
$this->assertNull($this->store->read('tenant', 'svc', 'INBOX', 7, 42));
$this->store->write('tenant', 'svc', 'INBOX', 7, 42, ['v' => 1]);
file_put_contents(
sprintf('%s/storage/tenant/provider_imap/svc/%s/7/42/message.json', $this->root, sha1('INBOX')),
'{truncated',
);
$this->assertNull($this->store->read('tenant', 'svc', 'INBOX', 7, 42));
}
public function testDeleteScopes(): void
{
$this->store->write('tenant', 'svc', 'INBOX', 7, 1, ['v' => 1]);
$this->store->write('tenant', 'svc', 'INBOX', 7, 2, ['v' => 2]);
$this->store->write('tenant', 'svc', 'Sent', 3, 1, ['v' => 3]);
$this->store->write('tenant', 'other', 'INBOX', 7, 1, ['v' => 4]);
$this->store->delete('tenant', 'svc', 'INBOX', 7, 1);
$this->assertNull($this->store->read('tenant', 'svc', 'INBOX', 7, 1));
$this->assertNotNull($this->store->read('tenant', 'svc', 'INBOX', 7, 2));
$this->store->deleteByMailbox('tenant', 'svc', 'INBOX');
$this->assertNull($this->store->read('tenant', 'svc', 'INBOX', 7, 2));
$this->assertNotNull($this->store->read('tenant', 'svc', 'Sent', 3, 1));
$this->store->deleteByService('tenant', 'svc');
$this->assertNull($this->store->read('tenant', 'svc', 'Sent', 3, 1));
$this->assertDirectoryDoesNotExist($this->root . '/storage/tenant/provider_imap/svc');
$this->assertNotNull($this->store->read('tenant', 'other', 'INBOX', 7, 1));
}
public function testIdentifiersThatCouldEscapeTheRootAreRejected(): void
{
$this->expectException(InvalidArgumentException::class);
$this->store->write('tenant', '../escape', 'INBOX', 7, 1, []);
}
}
+64
View File
@@ -0,0 +1,64 @@
<?php
declare(strict_types=1);
namespace KTXT\ProviderImap\Tests\Unit;
use KTXM\ProviderImap\Providers\EntityResource;
use KTXM\ProviderImap\Providers\MessageProperties;
use KTXM\ProviderImap\Service\Cache\MessageIngestor;
use KTXM\ProviderImap\Stores\MessageFileStore;
use KTXM\ProviderImap\Stores\MessageStore;
use PHPUnit\Framework\TestCase;
use RuntimeException;
final class MessageIngestorTest extends TestCase
{
public function testWritesContentBeforeMetaForEachMessage(): void
{
$calls = [];
$fileStore = $this->createStub(MessageFileStore::class);
$fileStore->method('write')->willReturnCallback(
function (string $tid, string $sid, string $mailbox, int $uidValidity, int $uid, array $content) use (&$calls): void {
$calls[] = ['file', $tid, $sid, $mailbox, $uidValidity, $uid, $content['schemaVersion']];
},
);
$messageStore = $this->createStub(MessageStore::class);
$messageStore->method('upsert')->willReturnCallback(
function (string $tid, string $sid, int $uidValidity, EntityResource $entity) use (&$calls): void {
$calls[] = ['meta', $tid, $sid, $entity->collection(), $uidValidity, $entity->identifier()];
},
);
$ingested = (new MessageIngestor($fileStore, $messageStore))
->ingest('tenant', 'svc', 7, $this->entity(41), $this->entity(42));
$this->assertSame([41, 42], $ingested);
$this->assertSame([
['file', 'tenant', 'svc', 'INBOX', 7, 41, EntityResource::CACHE_SCHEMA_VERSION],
['meta', 'tenant', 'svc', 'INBOX', 7, 41],
['file', 'tenant', 'svc', 'INBOX', 7, 42, EntityResource::CACHE_SCHEMA_VERSION],
['meta', 'tenant', 'svc', 'INBOX', 7, 42],
], $calls);
}
public function testMetaIsNotWrittenWhenContentFails(): void
{
$fileStore = $this->createStub(MessageFileStore::class);
$fileStore->method('write')->willThrowException(new RuntimeException('disk full'));
$messageStore = $this->createMock(MessageStore::class);
$messageStore->expects($this->never())->method('upsert');
$this->expectException(RuntimeException::class);
(new MessageIngestor($fileStore, $messageStore))->ingest('tenant', 'svc', 7, $this->entity(41));
}
private function entity(int $uid): EntityResource
{
$properties = new MessageProperties([]);
$properties->setSubject('Message ' . $uid);
return (new EntityResource('imap', 'svc'))->fromMutation('INBOX', $uid, $properties);
}
}