feat: Atomically read, mutate, and replace a composition snapshot.

Signed-off-by: Sebastian <krupinski01@gmail.com>
This commit is contained in:
2026-08-17 23:20:35 -04:00
parent c95e503251
commit 8fb3086792
4 changed files with 423 additions and 121 deletions
+148 -94
View File
@@ -48,9 +48,12 @@ class CompositionManager {
'attachments' => [],
];
$this->compositionStore->compositionSave($tenantId, $userId, $identifier, $snapshot);
// Construct the snapshot mutation that will run while the composition is locked.
$saveComposition = function () use ($tenantId, $userId, $identifier, $action, $source, $snapshot): array {
if ($action !== 'forward') {
return $snapshot;
}
if ($action === 'forward') {
$sourceIndentifier = ResourceIdentifier::fromString($source);
if ($sourceIndentifier === null) {
throw new InvalidArgumentException('Invalid source identifier');
@@ -96,46 +99,63 @@ class CompositionManager {
$meta['size'] = $attachment->getSize() ?? $meta['size'];
$snapshot['attachments'][$meta['identifier']] = $meta;
}
// update the snapshot with the staged attachments
$this->compositionStore->compositionSave($tenantId, $userId, $identifier, $snapshot);
}
return $snapshot;
};
// Execute the mutation and atomically persist the snapshot it returns.
$snapshot = $this->compositionStore->compositionSave(
$tenantId,
$userId,
$identifier,
$saveComposition,
) ?? $snapshot;
return $snapshot;
}
public function patch(string $tenantId, string $userId, string $identifier, array $data): array {
$composed = $this->compositionStore->compositionFetch($tenantId, $userId, $identifier);
$result = [];
// Construct the snapshot mutation that will run while the composition is locked.
$patchComposition = function (?array $composed) use ($data, $identifier, &$result): ?array {
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
];
return null;
}
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
if (!isset($data['revision']) || !is_int($data['revision']) || $data['revision'] < $composed['revision']) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'revision_mismatch',
'message' => 'Revision mismatch',
'currentRevision' => $composed['revision'],
];
return null;
}
$composed['revision'] = $data['revision'];
$composed['sender'] = $data['sender'] ?? [];
$composed['message'] = $data['message'] ?? [];
$result = [
'identifier' => $identifier,
'disposition' => 'patched',
];
return $result;
}
if (!isset($data['revision']) || !is_int($data['revision']) || $data['revision'] < $composed['revision']) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'revision_mismatch',
'message' => 'Revision mismatch',
'currentRevision' => $composed['revision'],
];
return $result;
}
return $composed;
};
// Apply the patch data to the composed message
$composed['revision'] = $data['revision'];
$composed['sender'] = $data['sender'] ?? [];
$composed['message'] = $data['message'] ?? [];
// Execute the mutation and atomically persist the snapshot it returns.
$this->compositionStore->compositionSave(
$tenantId,
$userId,
$identifier,
$patchComposition,
);
$this->compositionStore->compositionSave($tenantId, $userId, $identifier, $composed);
$result = [
'identifier' => $identifier,
'disposition' => 'patched',
];
return $result;
}
@@ -147,22 +167,35 @@ class CompositionManager {
}
public function send(string $tenantId, string $userId, string $identifier, array $sender, array $message, array $attachments): array {
$composed = $this->compositionStore->compositionFetch($tenantId, $userId, $identifier);
// Construct the snapshot mutation that will run while the composition is locked.
$saveComposition = function (?array $composed) use ($sender, $message): ?array {
if ($composed === null) {
return null;
}
$composed['sender'] = $sender;
$composed['message'] = $message;
return $composed;
};
// Execute the mutation and atomically persist the snapshot it returns.
$composed = $this->compositionStore->compositionSave(
$tenantId,
$userId,
$identifier,
$saveComposition,
);
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
return [
'disposition' => 'error',
'error' => [
'type' => 'composition_not_found',
'message' => 'Composition not found',
],
];
return $result;
}
// store the message and sender information before sending
$composed['sender'] = $sender;
$composed['message'] = $message;
$this->compositionStore->compositionSave($tenantId, $userId, $identifier, $composed);
// validate the sender address before attempting to send the message
if ($sender['address'] === '' || !filter_var($sender['address'], FILTER_VALIDATE_EMAIL)) {
$result['disposition'] = 'error';
@@ -306,38 +339,49 @@ class CompositionManager {
'composition' => $composition,
'attachments' => [],
];
// fetch the composition to ensure it exists
$composed = $this->compositionStore->compositionFetch($tenantId, $userId, $composition);
$failed = [];
// Construct the snapshot mutation that will run while the composition is locked.
$addAttachments = function (?array $composed) use ($tenantId, $userId, $composition, $attachments, &$result, &$failed): ?array {
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
];
return null;
}
$uploads = [];
$documents = [];
foreach ($attachments as $attachment) {
if (isset($attachment['origin']) && $attachment['origin'] === 'documents') {
$documents[] = $attachment;
} elseif (isset($attachment['origin']) && $attachment['origin'] === 'device') {
$uploads[] = $attachment;
}
}
if ($uploads !== []) {
$this->attachmentAddFromDevice($tenantId, $userId, $composition, $uploads, $composed);
}
if ($documents !== []) {
$failed = $this->attachmentAddFromDocuments($tenantId, $userId, $composition, $documents, $composed);
}
return $composed;
};
// Execute the mutation and atomically persist the snapshot it returns.
$composed = $this->compositionStore->compositionSave(
$tenantId,
$userId,
$composition,
$addAttachments,
);
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
];
return $result;
}
// separate attachments by type
$uploads = [];
$documents = [];
foreach ($attachments as $attachment) {
if (isset($attachment['origin']) && $attachment['origin'] === 'documents') {
$documents[] = $attachment;
} elseif (isset($attachment['origin']) && $attachment['origin'] === 'device') {
$uploads[] = $attachment;
}
}
// process uploaded attachments
if ($uploads !== []) {
$this->attachmentAddFromDevice($tenantId, $userId, $composition, $uploads, $composed);
}
// process attachments from the documents module
if ($documents !== []) {
$failed = $documents !== []
? $this->attachmentAddFromDocuments($tenantId, $userId, $composition, $documents, $composed)
: [];
}
// save the updated composition with new attachments
$this->compositionStore->compositionSave($tenantId, $userId, $composition, $composed);
// construct the result
$result['attachments'] = $composed['attachments'];
$result['failed'] = $failed;
@@ -486,29 +530,39 @@ class CompositionManager {
'composition' => $composition,
'identifier' => $identifier,
];
$composed = $this->compositionStore->compositionFetch($tenantId, $userId, $composition);
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
];
return $result;
}
if (!isset($composed['attachments'][$identifier])) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'attachment_not_found',
'message' => 'Attachment not found',
];
return $result;
}
unset($composed['attachments'][$identifier]);
$this->compositionStore->compositionSave($tenantId, $userId, $composition, $composed);
// Construct the snapshot mutation that will run while the composition is locked.
$removeAttachment = function (?array $composed) use ($identifier, &$result): ?array {
if ($composed === null) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'composition_not_found',
'message' => 'Composition not found',
];
return null;
}
if (!isset($composed['attachments'][$identifier])) {
$result['disposition'] = 'error';
$result['error'] = [
'type' => 'attachment_not_found',
'message' => 'Attachment not found',
];
return null;
}
unset($composed['attachments'][$identifier]);
$result['disposition'] = 'removed';
return $composed;
};
// Execute the mutation and atomically persist the snapshot it returns.
$this->compositionStore->compositionSave(
$tenantId,
$userId,
$composition,
$removeAttachment,
);
$result['disposition'] = 'removed';
return $result;
}
@@ -555,4 +609,4 @@ class CompositionManager {
}
fclose($resource);
}
}
}
+113 -25
View File
@@ -6,6 +6,8 @@ namespace KTXM\Mail\Stores;
use DI\Attribute\Inject;
use KTXF\Resource\BinaryResource;
use RuntimeException;
use UnexpectedValueException;
final class CompositionStore {
@@ -20,38 +22,118 @@ final class CompositionStore {
}
public function compositionFetch(string $tenantId, string $userId, string $draftId): ?array {
$draftDir = $this->draftDir($tenantId, $userId, $draftId);
if (!is_dir($draftDir)) {
return null;
}
$messagePath = $draftDir . '/' . self::COMPOSITION_FILENAME;
if (!file_exists($messagePath)) {
return null;
}
$decoded = json_decode((string)file_get_contents($messagePath), true);
return is_array($decoded) ? $decoded : null;
return $this->withCompositionLock(
$tenantId,
$userId,
$draftId,
LOCK_SH,
fn(): ?array => $this->compositionRead($tenantId, $userId, $draftId),
);
}
public function compositionSave(string $tenantId, string $userId, string $draftId, array $snapshot): array {
$draftDir = $this->draftDir($tenantId, $userId, $draftId);
if (!is_dir($draftDir)) {
mkdir($draftDir, 0755, true);
}
/**
* Atomically read, mutate, and replace a composition snapshot.
*
* Returning null from the callback leaves the stored snapshot unchanged.
* The callback receives null when the composition does not exist.
*
* @param callable(?array): ?array $update
*/
public function compositionSave(string $tenantId, string $userId, string $draftId, callable $update): ?array {
return $this->withCompositionLock(
$tenantId,
$userId,
$draftId,
LOCK_EX,
function () use ($tenantId, $userId, $draftId, $update): ?array {
$snapshot = $update($this->compositionRead($tenantId, $userId, $draftId));
if ($snapshot === null) {
return null;
}
if (!is_array($snapshot)) {
throw new UnexpectedValueException('Composition save callback must return an array or null');
}
file_put_contents($draftDir . '/' . self::COMPOSITION_FILENAME, json_encode($snapshot, JSON_PRETTY_PRINT | JSON_UNESCAPED_SLASHES));
return $snapshot;
$json = json_encode($snapshot, JSON_PRETTY_PRINT | JSON_UNESCAPED_SLASHES | JSON_THROW_ON_ERROR);
$draftDir = $this->draftDir($tenantId, $userId, $draftId);
$this->ensureDirectory($draftDir);
$temporaryPath = tempnam($draftDir, '.composition-');
if ($temporaryPath === false) {
throw new RuntimeException("Unable to create a temporary composition file for '$draftId'");
}
try {
$written = file_put_contents($temporaryPath, $json);
if ($written !== strlen($json)) {
throw new RuntimeException("Unable to write composition '$draftId'");
}
if (!rename($temporaryPath, $draftDir . '/' . self::COMPOSITION_FILENAME)) {
throw new RuntimeException("Unable to replace composition '$draftId'");
}
} finally {
if (is_file($temporaryPath)) {
unlink($temporaryPath);
}
}
return $snapshot;
},
);
}
public function compositionDiscard(string $tenantId, string $userId, string $draftId): bool {
$draftDir = $this->draftDir($tenantId, $userId, $draftId);
if (!is_dir($draftDir)) {
return false;
return $this->withCompositionLock(
$tenantId,
$userId,
$draftId,
LOCK_EX,
function () use ($tenantId, $userId, $draftId): bool {
$draftDir = $this->draftDir($tenantId, $userId, $draftId);
if (!is_dir($draftDir)) {
return false;
}
$this->deleteDir($draftDir);
return true;
},
);
}
private function compositionRead(string $tenantId, string $userId, string $draftId): ?array {
$messagePath = $this->draftDir($tenantId, $userId, $draftId) . '/' . self::COMPOSITION_FILENAME;
if (!is_file($messagePath)) {
return null;
}
$this->deleteDir($draftDir);
return true;
$contents = file_get_contents($messagePath);
if ($contents === false) {
throw new RuntimeException("Unable to read composition '$draftId'");
}
$decoded = json_decode($contents, true);
return is_array($decoded) ? $decoded : null;
}
private function withCompositionLock(string $tenantId, string $userId, string $draftId, int $operation, callable $callback): mixed {
$lockDir = $this->storagePath . '/.locks';
$this->ensureDirectory($lockDir);
$lockId = hash('sha256', $tenantId . "\0" . $userId . "\0" . $draftId);
$handle = fopen($lockDir . '/' . $lockId . '.lock', 'c+b');
if ($handle === false) {
throw new RuntimeException("Unable to open composition lock for '$draftId'");
}
try {
if (!flock($handle, $operation)) {
throw new RuntimeException("Unable to lock composition '$draftId'");
}
return $callback();
} finally {
flock($handle, LOCK_UN);
fclose($handle);
}
}
public function attachmentStageFromStream(string $tenantId, string $userId, string $compositionId, string $attachmentId, BinaryResource $data): array {
@@ -112,6 +194,12 @@ final class CompositionStore {
return $this->storagePath . '/' . $tenantId . '/' . $userId . '/' . $draftId;
}
private function ensureDirectory(string $path): void {
if (!is_dir($path) && !mkdir($path, 0755, true) && !is_dir($path)) {
throw new RuntimeException("Unable to create directory '$path'");
}
}
private function deleteDir(string $path): void {
if (!is_dir($path)) {
return;
@@ -133,4 +221,4 @@ final class CompositionStore {
rmdir($path);
}
}
}
+9 -2
View File
@@ -190,10 +190,17 @@ final class CompositionManagerTest extends TestCase {
private function stageEmptyComposition(): string {
$compositionId = 'draft-' . uniqid();
$this->compositionStore->compositionSave(self::TENANT_ID, self::USER_ID, $compositionId, [
$saveComposition = static fn(): array => [
'identifier' => $compositionId,
'attachments' => [],
]);
];
$this->compositionStore->compositionSave(
self::TENANT_ID,
self::USER_ID,
$compositionId,
$saveComposition,
);
return $compositionId;
}
+153
View File
@@ -0,0 +1,153 @@
<?php
declare(strict_types=1);
namespace KTXT\Mail\Tests\Unit;
use JsonException;
use KTXM\Mail\Stores\CompositionStore;
use PHPUnit\Framework\TestCase;
final class CompositionStoreTest extends TestCase {
private const TENANT_ID = 'tenant-1';
private const USER_ID = 'user-1';
private const COMPOSITION_ID = 'draft-1';
private string $rootDir;
private CompositionStore $store;
protected function setUp(): void {
$this->rootDir = sys_get_temp_dir() . '/ktrix-mail-store-test-' . uniqid('', true);
mkdir($this->rootDir, 0755, true);
$this->store = new CompositionStore($this->rootDir);
}
protected function tearDown(): void {
$this->deleteDir($this->rootDir);
}
public function testFailedReplacementPreservesPreviousSnapshot(): void {
$initial = ['identifier' => self::COMPOSITION_ID, 'revision' => 1];
$this->saveSnapshot($initial);
$resource = fopen('php://memory', 'r');
try {
$saveInvalidSnapshot = static fn(): array => [
'identifier' => self::COMPOSITION_ID,
'invalid' => $resource,
];
$this->store->compositionSave(
self::TENANT_ID,
self::USER_ID,
self::COMPOSITION_ID,
$saveInvalidSnapshot,
);
$this->fail('Saving a value that cannot be encoded as JSON should fail');
} catch (JsonException) {
$this->assertSame(
$initial,
$this->store->compositionFetch(self::TENANT_ID, self::USER_ID, self::COMPOSITION_ID),
);
} finally {
fclose($resource);
}
}
public function testUpdateCanLeaveSnapshotUnchanged(): void {
$initial = ['identifier' => self::COMPOSITION_ID, 'revision' => 1];
$this->saveSnapshot($initial);
$leaveUnchanged = static fn(?array $snapshot): ?array => null;
$updated = $this->store->compositionSave(
self::TENANT_ID,
self::USER_ID,
self::COMPOSITION_ID,
$leaveUnchanged,
);
$this->assertNull($updated);
$this->assertSame(
$initial,
$this->store->compositionFetch(self::TENANT_ID, self::USER_ID, self::COMPOSITION_ID),
);
}
public function testConcurrentUpdatesDoNotLoseChanges(): void {
if (!function_exists('pcntl_fork')) {
$this->markTestSkipped('The pcntl extension is required for the concurrency test');
}
$this->saveSnapshot([
'identifier' => self::COMPOSITION_ID,
'updates' => 0,
]);
$childPid = pcntl_fork();
if ($childPid === -1) {
$this->fail('Unable to fork the concurrency test process');
}
if ($childPid === 0) {
$this->incrementAfterDelay();
exit(0);
}
$this->incrementAfterDelay();
pcntl_waitpid($childPid, $status);
$this->assertTrue(pcntl_wifexited($status));
$this->assertSame(0, pcntl_wexitstatus($status));
$this->assertSame(
2,
$this->store->compositionFetch(self::TENANT_ID, self::USER_ID, self::COMPOSITION_ID)['updates'],
);
}
private function incrementAfterDelay(): void {
$incrementSnapshot = static function (?array $snapshot): array {
usleep(200_000);
$snapshot['updates']++;
return $snapshot;
};
$this->store->compositionSave(
self::TENANT_ID,
self::USER_ID,
self::COMPOSITION_ID,
$incrementSnapshot,
);
}
private function saveSnapshot(array $snapshot): void {
$saveSnapshot = static fn(): array => $snapshot;
$this->store->compositionSave(
self::TENANT_ID,
self::USER_ID,
self::COMPOSITION_ID,
$saveSnapshot,
);
}
private function deleteDir(string $path): void {
if (!is_dir($path)) {
return;
}
foreach (scandir($path) ?: [] as $entry) {
if ($entry === '.' || $entry === '..') {
continue;
}
$child = $path . '/' . $entry;
if (is_dir($child)) {
$this->deleteDir($child);
continue;
}
unlink($child);
}
rmdir($path);
}
}