$identifier, 'action' => $action, 'source' => $source, 'revision' => 1, 'disposition' => 'staged', 'sender' => $sender, 'message' => $message, 'attachments' => [], 'remote' => [ 'status' => 'dirty', 'entity' => null, 'error' => null, ], ]; // 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; } $sourceIndentifier = ResourceIdentifier::fromString($source); if ($sourceIndentifier === null) { throw new InvalidArgumentException('Invalid source identifier'); } $entities = $this->mailManager->entityFetchBulk($tenantId, $userId, $sourceIndentifier); if (count($entities) === 0) { throw new InvalidArgumentException('Source message not found'); } // retrieve the attachments from the source message $attachments = reset($entities)->getProperties()->getAttachments(); // For each attachment, we download the content and stage it in the composition store foreach ($attachments as $attachment) { $contentId = $attachment->getContentId(); $isInline = $attachment->getDisposition() === 'inline'; // target part $targetPart = [ 'partId' => $attachment->getId(), 'blobId' => $attachment->getBlobId(), 'cid' => $contentId, 'cId' => $contentId, ]; // retrieve the attachment content as a stream $data = $this->mailManager->entityDownload($tenantId, $userId, $sourceIndentifier, $targetPart); // stage the attachment in the composition store $meta = $this->compositionStore->attachmentStageFromStream( tenantId: $tenantId, userId: $userId, compositionId: $identifier, attachmentId: UUID::v4(), data: $data, ); // supplement the attachment metadata with source information $meta['origin'] = 'message'; $meta['source'] = $source; $meta['partId'] = $attachment->getId(); $meta['blobId'] = $attachment->getBlobId(); $meta['cid'] = $contentId; $meta['contentId'] = $contentId; $meta['disposition'] = $attachment->getDisposition(); $meta['inline'] = $isInline; $meta['name'] = $attachment->getName() ?? $meta['name']; $meta['type'] = $attachment->getType() ?? $meta['type']; $meta['size'] = $attachment->getSize() ?? $meta['size']; $snapshot['attachments'][$meta['identifier']] = $meta; } return $snapshot; }; // Execute the mutation and atomically persist the snapshot it returns. $snapshot = $this->compositionStore->compositionSave( $tenantId, $userId, $identifier, $saveComposition, ) ?? $snapshot; $event = new CompositionSavedEvent($tenantId, $userId, $identifier, (int)$snapshot['revision']); $this->events->dispatch($event); $response = $snapshot; unset($response['remote']); return $response; } public function patch(string $tenantId, string $userId, string $identifier, array $data): array { $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 (!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'] ?? []; $composed = $this->markRemoteDirty($composed); $result = [ 'identifier' => $identifier, 'disposition' => 'patched', ]; return $composed; }; // Execute the mutation and atomically persist the snapshot it returns. $composed = $this->compositionStore->compositionSave( $tenantId, $userId, $identifier, $patchComposition, ); if (($result['disposition'] ?? null) === 'patched' && $composed !== null) { $event = new CompositionSavedEvent($tenantId, $userId, $identifier, (int)($composed['revision'] ?? 0)); $this->events->dispatch($event); } return $result; } public function save(string $tenantId, string $userId, string $identifier, array $data): array { // Construct the snapshot mutation that will run while the composition is locked. $saveComposition = function (?array $composed) use ($data): ?array { if ($composed === null) { return null; } if (isset($data['revision']) && is_int($data['revision'])) { $composed['revision'] = $data['revision']; } $composed['sender'] = $data['sender']; $composed['message'] = $data['message']; return $this->markRemoteDirty($composed); }; // Execute the mutation and atomically persist the snapshot it returns. $composed = $this->compositionStore->compositionSave( $tenantId, $userId, $identifier, $saveComposition, ); if ($composed === null) { return [ 'identifier' => $identifier, 'disposition' => 'error', 'error' => [ 'type' => 'composition_not_found', 'message' => 'Composition not found', ], ]; } $synchronized = $this->synchronize($tenantId, $userId, $identifier); if ($synchronized['disposition'] !== 'saved') { return [ 'identifier' => $identifier, 'disposition' => 'error', 'error' => [ 'type' => 'composition_save_failed', 'message' => (string)($synchronized['error'] ?? 'Draft could not be saved to the server'), ], ]; } $this->compositionStore->compositionDiscard($tenantId, $userId, $identifier); return ['identifier' => $identifier, 'disposition' => 'saved']; } public function discard(string $tenantId, string $userId, string $identifier): array { $composed = $this->compositionStore->compositionFetch($tenantId, $userId, $identifier); if ($composed === null) { return ['identifier' => $identifier, 'disposition' => 'discarded']; } try { $remoteEntity = $this->remoteEntityIdentifier($composed['remote']['entity'] ?? null); if ($remoteEntity !== null) { $outcomes = $this->mailManager->entityDelete($tenantId, $userId, $remoteEntity); $outcome = $outcomes[(string)$remoteEntity] ?? reset($outcomes); if (!is_array($outcome) || ($outcome['disposition'] ?? 'error') === 'error') { return [ 'identifier' => $identifier, 'disposition' => 'error', 'error' => [ 'type' => 'composition_discard_failed', 'message' => is_array($outcome) ? (string)($outcome['error'] ?? 'Remote draft could not be discarded') : 'Remote draft could not be discarded', ], ]; } } } catch (Throwable $throwable) { return [ 'identifier' => $identifier, 'disposition' => 'error', 'error' => [ 'type' => 'composition_discard_failed', 'message' => $throwable->getMessage(), ], ]; } $this->compositionStore->compositionDiscard($tenantId, $userId, $identifier); return ['identifier' => $identifier, 'disposition' => 'discarded']; } public function synchronize(string $tenantId, string $userId, string $identifier): array { $result = [ 'identifier' => $identifier, 'disposition' => 'staged', ]; // Construct the synchronization operation that will run while the composition is locked. $synchronizeComposition = function (?array $composed) use ($tenantId, $userId, $identifier, &$result): ?array { if ($composed === null) { $result['error'] = 'Composition not found'; return null; } $remote = isset($composed['remote']) && is_array($composed['remote']) ? $composed['remote'] : ['status' => 'dirty', 'entity' => null, 'error' => null]; if (($remote['status'] ?? 'dirty') === 'synced') { $result['disposition'] = 'saved'; return null; } try { $sender = isset($composed['sender']) && is_array($composed['sender']) ? $composed['sender'] : []; $senderAddress = (string)($sender['address'] ?? ''); if ($senderAddress === '' || !filter_var($senderAddress, FILTER_VALIDATE_EMAIL)) { throw new RuntimeException('Composition sender is invalid'); } $service = $this->mailManager->serviceFindByAddress($tenantId, $userId, $senderAddress); if ($service === null || $service->getEnabled() === false) { throw new RuntimeException("No enabled mail service handles '$senderAddress'"); } if (!$service instanceof ServiceEntityMutableInterface) { throw new RuntimeException("Mail service '{$service->identifier()}' does not support draft mutations"); } $target = $this->resolveDraftCollection($service); $attachments = isset($composed['attachments']) && is_array($composed['attachments']) ? array_values($composed['attachments']) : []; $properties = $this->buildMessageProperties( service: $service, sender: Address::fromArray($sender), message: is_array($composed['message'] ?? null) ? $composed['message'] : [], attachments: $attachments, composed: $composed, tenantId: $tenantId, userId: $userId, compositionId: $identifier, ); $properties->setFlag('draft', true); $properties->setFlag('seen', true); $remoteEntity = $this->remoteEntityIdentifier($remote['entity'] ?? null); if ($remoteEntity === null) { $entity = $this->mailManager->entityCreate($tenantId, $userId, $target, $properties); } else { if ($remoteEntity->provider() !== $service->provider() || (string)$remoteEntity->service() !== (string)$service->identifier()) { throw new RuntimeException('Changing the draft service is not supported yet'); } $entity = $this->mailManager->entityModify($tenantId, $userId, $remoteEntity, $properties); } if ($entity->collection() === null || $entity->identifier() === null) { throw new RuntimeException('Provider returned an incomplete draft identifier'); } $composed['remote'] = [ ...$remote, 'status' => 'synced', 'entity' => (string)new EntityIdentifier( $entity->provider(), (string)$entity->service(), (string)$entity->collection(), (string)$entity->identifier(), ), 'error' => null, ]; $result['disposition'] = 'saved'; } catch (Throwable $throwable) { $composed['remote'] = [ ...$remote, 'status' => 'failed', 'entity' => $remote['entity'] ?? null, 'error' => $throwable->getMessage(), ]; $result['error'] = $throwable->getMessage(); } return $composed; }; // Synchronize and atomically persist the resulting remote state. $this->compositionStore->compositionSave( $tenantId, $userId, $identifier, $synchronizeComposition, ); return $result; } public function send(string $tenantId, string $userId, string $identifier, array $sender, array $message, array $attachments): array { // 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 $this->markRemoteDirty($composed); }; // Execute the mutation and atomically persist the snapshot it returns. $composed = $this->compositionStore->compositionSave( $tenantId, $userId, $identifier, $saveComposition, ); if ($composed === null) { return [ 'disposition' => 'error', 'error' => [ 'type' => 'composition_not_found', 'message' => 'Composition not found', ], ]; } // validate the sender address before attempting to send the message if ($sender['address'] === '' || !filter_var($sender['address'], FILTER_VALIDATE_EMAIL)) { $result['disposition'] = 'error'; $result['error'] = [ 'type' => 'composition_invalid_sender', 'message' => 'Invalid sender address', ]; return $result; } $senderAddress = $sender['address']; $senderObject = Address::fromArray($sender); // resolve submit-capable service for sender $service = $this->mailManager->serviceFindByAddress($tenantId, $userId, $senderAddress); if ($service === null || $service->getEnabled() === false) { return [ 'disposition' => 'error', 'error' => [ 'type' => 'service_not_found', 'message' => "Service not found for sender '$senderAddress' or service is disabled", ], ]; } if ($service instanceof ServiceEntitySubmitInterface === false) { return [ 'disposition' => 'error', 'error' => [ 'type' => 'service_not_supported', 'message' => "Service '{$service->identifier()}' does not support entity submission", ], ]; } $source = null; $properties = $this->buildMessageProperties( service: $service, sender: $senderObject, message: $message, attachments: $attachments, composed: $composed, tenantId: $tenantId, userId: $userId, compositionId: $identifier, ); $sendResult = $service->entitySubmit($senderObject, $source, $properties); if ($sendResult->disposition === EntitySubmitResult::DISPOSITION_ERROR) { return [ 'identifier' => $identifier, 'disposition' => 'error', 'error' => [ 'type' => 'service_submission_error', 'message' => $sendResult->errorMessage ?? 'An unknown error occurred during submission', ], ]; } return [ 'identifier' => $identifier, 'disposition' => 'sent' ]; } private function buildMessageProperties( ServiceEntitySubmitInterface|ServiceEntityMutableInterface $service, AddressInterface $sender, array $message, array $attachments, array $composed, string $tenantId, string $userId, string $compositionId, ): MessagePropertiesMutableInterface { $properties = $service->entityFresh()->getProperties(); $bodyTextPlain = isset($message['body']['text']) ? (string)$message['body']['text'] : ''; $bodyTextHtml = isset($message['body']['html']) ? (string)$message['body']['html'] : ''; $properties->setFrom($sender); $to = $this->mapAddresses($message['to'] ?? []); if ($to !== []) { $properties->setTo(...$to); } $cc = $this->mapAddresses($message['cc'] ?? []); if ($cc !== []) { $properties->setCc(...$cc); } $bcc = $this->mapAddresses($message['bcc'] ?? []); if ($bcc !== []) { $properties->setBcc(...$bcc); } $replyTo = $this->mapAddresses($message['replyTo'] ?? []); if ($replyTo !== []) { $properties->setReplyTo(...$replyTo); } $properties->setSubject((string)($message['subject'] ?? '')); $properties->setBodyTextPlain($bodyTextPlain); $properties->setBodyTextHtml($bodyTextHtml); if (isset($message['flags']) && is_array($message['flags'])) { $properties->setFlags($message['flags']); } $attachmentObjects = []; foreach ($attachments as $attachment) { if (!isset($attachment['identifier']) || !is_string($attachment['identifier']) || $attachment['identifier'] === '') { continue; } $composedAttachment = $composed['attachments'][$attachment['identifier']] ?? null; if (!is_array($composedAttachment)) { continue; } $attachmentObjects[] = MessagePart::fromArray([ 'partId' => $attachment['identifier'], 'blobId' => $attachment['blobId'] ?? $composedAttachment['blobId'] ?? null, 'size' => $composedAttachment['size'] ?? null, 'name' => $composedAttachment['name'] ?? 'unknown.bin', 'type' => $composedAttachment['type'] ?? 'application/octet-stream', 'disposition' => (($attachment['inline'] ?? $composedAttachment['inline'] ?? false) === true) ? 'inline' : 'attachment', 'content' => $this->compositionStore->attachmentFetchData($tenantId, $userId, $compositionId, $attachment['identifier']), 'cid' => $attachment['contentId'] ?? $attachment['cid'] ?? $composedAttachment['contentId'] ?? $composedAttachment['cid'] ?? null, ]); } if ($attachmentObjects !== []) { $properties->setAttachments(...$attachmentObjects); } return $properties; } public function attachmentAdd(string $tenantId, string $userId, string $composition, array ...$attachments): array { $result = [ 'composition' => $composition, 'attachments' => [], ]; $failed = []; $compositionChanged = false; // Construct the snapshot mutation that will run while the composition is locked. $addAttachments = function (?array $composed) use ($tenantId, $userId, $composition, $attachments, &$result, &$failed, &$compositionChanged): ?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); } $deviceAttachmentsChanged = $uploads !== []; $documentAttachmentsChanged = count($failed) < count($documents); if ($deviceAttachmentsChanged || $documentAttachmentsChanged) { $composed = $this->markRemoteDirty($composed); $compositionChanged = true; } return $composed; }; // Execute the mutation and atomically persist the snapshot it returns. $composed = $this->compositionStore->compositionSave( $tenantId, $userId, $composition, $addAttachments, ); if ($composed === null) { return $result; } // construct the result $result['attachments'] = $composed['attachments']; $result['failed'] = $failed; $result['disposition'] = match (true) { $failed === [] => 'added', $result['attachments'] !== [] => 'partial', default => 'error', }; if ($result['disposition'] === 'error') { $result['error'] = [ 'type' => 'attachment_add_failed', 'message' => 'None of the requested attachments could be added', ]; } if ($compositionChanged) { $event = new CompositionSavedEvent($tenantId, $userId, $composition, (int)($composed['revision'] ?? 0)); $this->events->dispatch($event); } return $result; } /** * Stage device-uploaded (base64 payload) attachment entries */ private function attachmentAddFromDevice(string $tenantId, string $userId, string $composition, array $uploads, array &$composed): void { foreach ($uploads as $attachment) { $identifier = (isset($attachment['identifier']) && is_string($attachment['identifier']) && $attachment['identifier'] !== '') ? $attachment['identifier'] : UUID::v4(); $decoded = base64_decode((string) ($attachment['data'] ?? ''), true); if ($decoded === false) { throw new InvalidArgumentException('Attachment payload is not valid base64'); } $data = new BinaryResource( $attachment['name'] ?? 'unknown.bin', $attachment['type'] ?? 'application/octet-stream', $this->stringToGenerator($decoded), ); $meta = $this->compositionStore->attachmentStageFromStream( tenantId: $tenantId, userId: $userId, compositionId: $composition, attachmentId: $identifier, data: $data, ); $meta['origin'] = 'device'; $composed['attachments'][$meta['identifier']] = $meta; } } private function resolveDocumentsManager(): ?DocumentsManager { $module = $this->moduleManager->fetch(self::DOCUMENTS_MODULE_HANDLE); if ($module === null || !$module->enabled()) { return null; } if (!$this->container->has(DocumentsManager::class)) { return null; } return $this->container->get(DocumentsManager::class); } /** * Stage documents-sourced attachment entries */ private function attachmentAddFromDocuments(string $tenantId, string $userId, string $composition, array $sourced, array &$composed): array { $failed = []; // The documents module is an optional dependency, so we must resolve it // lazily and handle the case where it is not installed or enabled. $documentsManager = $this->resolveDocumentsManager(); if ($documentsManager === null) { foreach ($sourced as $attachment) { $failed[] = [ 'identifier' => (string) ($attachment['source'] ?? ''), 'message' => 'Documents module is not installed or enabled', ]; } return $failed; } $requested = []; foreach ($sourced as $attachment) { $rawIdentifier = (string) $attachment['source']; try { $parsed = ResourceIdentifier::fromString($rawIdentifier); } catch (InvalidArgumentException) { $failed[] = ['identifier' => $rawIdentifier, 'message' => 'Invalid document identifier']; continue; } if (!$parsed instanceof EntityIdentifier) { $failed[] = ['identifier' => $rawIdentifier, 'message' => 'Invalid document identifier']; continue; } $requested[] = ['attachment' => $attachment, 'identifier' => $parsed]; } if ($requested === []) { return $failed; } $entities = $documentsManager->entityFetchBulk($tenantId, $userId, ...array_column($requested, 'identifier')); $entitiesByUrn = []; foreach ($entities as $entity) { $entitiesByUrn[$entity->urn()] = $entity; } foreach ($requested as ['attachment' => $attachment, 'identifier' => $entityIdentifier]) { $urn = (string) $entityIdentifier; $entity = $entitiesByUrn[$urn] ?? null; if ($entity === null) { $failed[] = ['identifier' => $urn, 'message' => 'Document not found or access denied']; continue; } $resource = $documentsManager->entityReadStream($tenantId, $userId, $entityIdentifier); if ($resource === null) { $failed[] = ['identifier' => $urn, 'message' => 'Document content could not be read']; continue; } $identifier = (isset($attachment['identifier']) && is_string($attachment['identifier']) && $attachment['identifier'] !== '') ? $attachment['identifier'] : UUID::v4(); $properties = $entity->getProperties(); $data = new BinaryResource($properties->getLabel(), $properties->getMime(), $this->streamToGenerator($resource)); $meta = $this->compositionStore->attachmentStageFromStream( tenantId: $tenantId, userId: $userId, compositionId: $composition, attachmentId: $identifier, data: $data, ); $meta['origin'] = 'documents'; $meta['source'] = $urn; $composed['attachments'][$meta['identifier']] = $meta; } return $failed; } public function attachmentRemove(string $tenantId, string $userId, string $composition, string $identifier): array { $result = [ 'composition' => $composition, 'identifier' => $identifier, ]; // 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]); $composed = $this->markRemoteDirty($composed); $result['disposition'] = 'removed'; return $composed; }; // Execute the mutation and atomically persist the snapshot it returns. $composed = $this->compositionStore->compositionSave( $tenantId, $userId, $composition, $removeAttachment, ); if (($result['disposition'] ?? null) === 'removed' && $composed !== null) { $event = new CompositionSavedEvent($tenantId, $userId, $composition, (int)($composed['revision'] ?? 0)); $this->events->dispatch($event); } return $result; } private function markRemoteDirty(array $composed): array { $remote = isset($composed['remote']) && is_array($composed['remote']) ? $composed['remote'] : []; $remote['status'] = 'dirty'; $remote['entity'] ??= null; $remote['error'] = null; $composed['remote'] = $remote; return $composed; } private function resolveDraftCollection(ServiceBaseInterface $service): CollectionIdentifier { $configuredTarget = $service->getAuxiliary()['draftTarget'] ?? null; if ($configuredTarget !== null && (string)$configuredTarget !== '') { return new CollectionIdentifier( $service->provider(), (string)$service->identifier(), (string)$configuredTarget, ); } $filter = $service->collectionListFilter(); $filter->condition('role', CollectionRoles::Drafts->value); $collections = $service->collectionList('', $filter); $collection = reset($collections); if (!$collection instanceof CollectionBaseInterface || $collection->identifier() === null) { throw new RuntimeException("Mail service '{$service->identifier()}' has no Drafts collection"); } return new CollectionIdentifier( $service->provider(), (string)$service->identifier(), (string)$collection->identifier(), ); } private function remoteEntityIdentifier(mixed $value): ?EntityIdentifier { if (!is_string($value) || $value === '') { return null; } $identifier = ResourceIdentifier::fromString($value); if (!$identifier instanceof EntityIdentifier) { throw new RuntimeException("Invalid remote draft identifier '$value'"); } return $identifier; } /** * @param array $entries * @return array */ private function mapAddresses(array $entries): array { $addresses = []; foreach ($entries as $entry) { if (is_array($entry)) { $address = Address::fromArray($entry); if ($address->getAddress() !== '') { $addresses[] = $address; } continue; } if (is_string($entry) && $entry !== '') { $address = Address::fromString($entry); if ($address->getAddress() !== '') { $addresses[] = $address; } } } return $addresses; } /** * Wrap an in-memory byte string in a single-chunk Generator. */ private function stringToGenerator(string $bytes): \Generator { yield $bytes; } /** * Wrap a raw PHP stream resource in a Generator yielding chunks, closing it once exhausted. */ private function streamToGenerator($resource): \Generator { while (!feof($resource)) { yield fread($resource, 65536); } fclose($resource); } }