Files
documents_manager/lib/Controllers/DefaultController.php
T
2026-09-05 20:05:08 -04:00

942 lines
36 KiB
PHP

<?php
declare(strict_types=1);
/**
* SPDX-FileCopyrightText: Sebastian Krupinski <krupinski01@gmail.com>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
namespace KTXM\DocumentsManager\Controllers;
use InvalidArgumentException;
use KTXC\Http\Response\JsonResponse;
use KTXC\Http\Response\Response;
use KTXC\Http\Response\StreamedNdJsonResponse;
use KTXC\Context\IdentityContextInterface;
use KTXC\Context\TenantContextInterface;
use KTXF\Controller\ControllerAbstract;
use KTXF\Json\JsonSerializable;
use KTXF\Resource\Identifier\CollectionIdentifier;
use KTXF\Resource\Identifier\EntityIdentifier;
use KTXF\Resource\Identifier\ResourceIdentifier;
use KTXF\Resource\Identifier\ResourceIdentifiers;
use KTXF\Resource\Identifier\ServiceIdentifier;
use KTXF\Routing\Attributes\AuthenticatedRoute;
use KTXM\DocumentsManager\Manager;
use KTXM\DocumentsManager\Stream\ExpectedTotal;
use Psr\Log\LoggerInterface;
use Throwable;
class DefaultController extends ControllerAbstract {
private const ERR_MISSING_PROVIDER = 'Missing parameter: provider';
private const ERR_MISSING_IDENTIFIER = 'Missing parameter: identifier';
private const ERR_MISSING_SERVICE = 'Missing parameter: service';
private const ERR_MISSING_DATA = 'Missing parameter: data';
private const ERR_MISSING_SOURCES = 'Missing parameter: sources';
private const ERR_MISSING_SOURCE = 'Missing parameter: source';
private const ERR_MISSING_TARGET = 'Missing parameter: target';
private const ERR_MISSING_TARGETS = 'Missing parameter: targets';
private const ERR_INVALID_OPERATION = 'Invalid operation: ';
private const ERR_INVALID_PROVIDER = 'Invalid parameter: provider must be a string';
private const ERR_INVALID_SERVICE = 'Invalid parameter: service must be a string';
private const ERR_INVALID_IDENTIFIER = 'Invalid parameter: identifier must be a string';
private const ERR_INVALID_SOURCES = 'Invalid parameter: sources must be an array';
private const ERR_INVALID_SOURCE = 'Invalid parameter: source must be a string';
private const ERR_INVALID_TARGET = 'Invalid parameter: target must be a string';
private const ERR_INVALID_TARGETS = 'Invalid parameter: targets must be an array';
private const ERR_INVALID_DATA = 'Invalid parameter: data must be an array';
private const ERR_TARGET_COLLECTION = 'Invalid parameter: target must be provider:service:collection';
private const ERR_TARGET_ENTITY = 'Invalid parameter: target must be provider:service:collection:entity';
public function __construct(
private readonly TenantContextInterface $tenantContext,
private readonly IdentityContextInterface $identityContext,
private readonly Manager $manager,
private readonly LoggerInterface $logger
) {}
/**
* Main API endpoint for documents operations
*
* Single operation:
* {
* "version": 1,
* "transaction": "tx-1",
* "operation": "entity.create",
* "data": {...}
* }
*
* @return Response
*/
#[AuthenticatedRoute('/v1', name: 'documents.manager.v1', methods: ['POST'])]
public function index(
int $version,
string $transaction,
string|null $operation = null,
array|null $data = null,
string|null $user = null
): Response {
// authorize request
$tenantId = $this->tenantContext->identifier();
$userId = $this->identityContext->identifier();
try {
if ($operation !== null) {
$result = $this->processOperation($tenantId, $userId, $operation, $data ?? [], $version, $transaction);
if ($result instanceof Response) {
return $result;
}
return new JsonResponse([
'version' => $version,
'transaction' => $transaction,
'operation' => $operation,
'status' => 'success',
'data' => $result
], JsonResponse::HTTP_OK);
}
throw new InvalidArgumentException('Operation must be provided');
} catch (Throwable $t) {
$this->logger->error('Error processing request', ['exception' => $t]);
return new JsonResponse([
'version' => $version,
'transaction' => $transaction,
'operation' => $operation,
'status' => 'error',
'data' => [
'code' => $t->getCode(),
'message' => $t->getMessage()
]
], JsonResponse::HTTP_INTERNAL_SERVER_ERROR);
}
}
/**
* Process a single operation
*/
private function processOperation(string $tenantId, string $userId, string $operation, array $data, int $version = 1, string $transaction = ''): mixed {
return match ($operation) {
// Provider operations
'provider.list' => $this->providerList($tenantId, $userId, $data),
'provider.fetch' => $this->providerFetch($tenantId, $userId, $data),
'provider.extant' => $this->providerExtant($tenantId, $userId, $data),
// Service operations
'service.list' => $this->serviceList($tenantId, $userId, $data),
'service.fetch' => $this->serviceFetch($tenantId, $userId, $data),
'service.extant' => $this->serviceExtant($tenantId, $userId, $data),
'service.create' => $this->serviceCreate($tenantId, $userId, $data),
'service.update' => $this->serviceUpdate($tenantId, $userId, $data),
'service.delete' => $this->serviceDelete($tenantId, $userId, $data),
'service.test' => $this->serviceTest($tenantId, $userId, $data),
// Collection operations
'collection.list' => $this->collectionList($tenantId, $userId, $data),
'collection.fetch' => $this->collectionFetch($tenantId, $userId, $data),
'collection.extant' => $this->collectionExtant($tenantId, $userId, $data),
'collection.create' => $this->collectionCreate($tenantId, $userId, $data),
'collection.update' => $this->collectionUpdate($tenantId, $userId, $data),
'collection.delete' => $this->collectionDelete($tenantId, $userId, $data),
'collection.copy' => $this->collectionCopy($tenantId, $userId, $data),
'collection.move' => $this->collectionMove($tenantId, $userId, $data),
// Entity operations
'entity.listBulk' => $this->entityListBulk($tenantId, $userId, $data),
'entity.listStream' => $this->entityListStream($tenantId, $userId, $data, $version, $transaction),
'entity.fetch' => $this->entityFetch($tenantId, $userId, $data),
'entity.extant' => $this->entityExtant($tenantId, $userId, $data),
'entity.delta' => $this->entityDelta($tenantId, $userId, $data),
'entity.create' => $this->entityCreate($tenantId, $userId, $data),
'entity.update' => $this->entityUpdate($tenantId, $userId, $data),
'entity.delete' => $this->entityDelete($tenantId, $userId, $data),
'entity.copy' => $this->entityCopy($tenantId, $userId, $data),
'entity.move' => $this->entityMove($tenantId, $userId, $data),
'entity.read' => $this->entityRead($tenantId, $userId, $data),
'entity.read.chunk' => $this->entityReadChunk($tenantId, $userId, $data),
'entity.write' => $this->entityWrite($tenantId, $userId, $data),
'entity.write.chunk' => $this->entityWriteChunk($tenantId, $userId, $data),
// Node operations (unified collections + entities)
'node.list' => $this->nodeList($tenantId, $userId, $data, $version, $transaction),
default => throw new InvalidArgumentException(self::ERR_INVALID_OPERATION . $operation)
};
}
// ==================== Provider Operations ====================
private function providerList(string $tenantId, string $userId, array $data): mixed {
if (isset($data['targets'])) {
if (!is_array($data['targets'])) {
throw new InvalidArgumentException(self::ERR_INVALID_TARGETS);
}
foreach ($data['targets'] as $target) {
if (!is_string($target)) {
throw new InvalidArgumentException(self::ERR_INVALID_TARGETS);
}
}
}
return $this->manager->providerList($tenantId, $userId, $data['targets'] ?? []);
}
private function providerFetch(string $tenantId, string $userId, array $data): mixed {
if (!isset($data['target'])) {
throw new InvalidArgumentException(self::ERR_MISSING_TARGET);
}
if (!is_string($data['target'])) {
throw new InvalidArgumentException(self::ERR_INVALID_TARGET);
}
return $this->manager->providerFetch($tenantId, $userId, $data['target']);
}
private function providerExtant(string $tenantId, string $userId, array $data): mixed {
if (!isset($data['targets'])) {
throw new InvalidArgumentException(self::ERR_MISSING_TARGETS);
}
foreach ($data['targets'] as $target) {
if (!is_string($target)) {
throw new InvalidArgumentException(self::ERR_INVALID_TARGETS);
}
}
return $this->manager->providerExtant($tenantId, $userId, $data['targets']);
}
// ==================== Service Operations =====================
private function serviceList(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->optionalIdentifiers(
$data,
'targets',
ServiceIdentifier::class,
'Invalid parameter: targets must contain provider:service, provider:service:collection, or provider:service:collection:entity identifiers'
);
// perform operation
return $this->manager->serviceList($tenantId, $userId, $targets);
}
private function serviceFetch(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$provider = $this->requireString($data, 'provider', self::ERR_MISSING_PROVIDER, self::ERR_INVALID_PROVIDER);
$identifier = $this->requireString($data, 'identifier', self::ERR_MISSING_IDENTIFIER, self::ERR_INVALID_IDENTIFIER);
// perform operation
return $this->manager->serviceFetch($tenantId, $userId, $provider, $identifier);
}
private function serviceExtant(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->requireIdentifiers(
$data,
'targets',
ServiceIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service identifiers'
);
// perform operation
return $this->manager->serviceExtant($tenantId, $userId, $targets);
}
private function serviceCreate(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$provider = $this->requireString($data, 'provider', self::ERR_MISSING_PROVIDER, self::ERR_INVALID_PROVIDER);
$properties = $this->requireArray($data, 'data', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
// perform operation
return $this->manager->serviceCreate(
$tenantId,
$userId,
$provider,
$properties
);
}
private function serviceUpdate(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$provider = $this->requireString($data, 'provider', self::ERR_MISSING_PROVIDER, self::ERR_INVALID_PROVIDER);
$identifier = $this->requireString($data, 'identifier', self::ERR_MISSING_IDENTIFIER, self::ERR_INVALID_IDENTIFIER);
$properties = $this->requireArray($data, 'data', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
if (isset($data['delta']) && !is_bool($data['delta'])) {
throw new InvalidArgumentException('Invalid parameter: delta must be a boolean');
}
// perform operation
return $this->manager->serviceUpdate(
$tenantId,
$userId,
$provider,
$identifier,
$properties,
$data['delta'] ?? false,
);
}
private function serviceDelete(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$provider = $this->requireString($data, 'provider', self::ERR_MISSING_PROVIDER, self::ERR_INVALID_PROVIDER);
$identifier = $this->requireString($data, 'identifier', self::ERR_MISSING_IDENTIFIER, self::ERR_INVALID_IDENTIFIER);
// perform operation
return $this->manager->serviceDelete(
$tenantId,
$userId,
$provider,
$identifier
);
}
private function serviceTest(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$provider = $this->requireString($data, 'provider', self::ERR_MISSING_PROVIDER, self::ERR_INVALID_PROVIDER);
if (!isset($data['identifier']) && !isset($data['location']) && !isset($data['identity'])) {
throw new InvalidArgumentException('Either a service identifier or location and identity must be provided for service test');
}
// perform operation
return $this->manager->serviceTest(
$tenantId,
$userId,
$provider,
$data['identifier'] ?? null,
$data['location'] ?? null,
$data['identity'] ?? null,
);
}
// ==================== Collection Operations ====================
private function collectionList(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$sources = $this->optionalIdentifiers(
$data,
'sources',
ServiceIdentifier::class,
'Invalid parameter: sources must contain provider:service or provider:service:collection identifiers'
);
$filter = $data['filter'] ?? null;
$sort = $data['sort'] ?? null;
// perform operation
return $this->manager->collectionList($tenantId, $userId, $sources, $filter, $sort);
}
private function collectionFetch(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targetIdentifiers = $this->requireIdentifiers(
$data,
'targets',
CollectionIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
self::ERR_TARGET_COLLECTION,
EntityIdentifier::class
);
// perform operation
$list = [];
foreach ($targetIdentifiers as $targetIdentifier) {
$collection = $this->manager->collectionFetch($tenantId, $userId, $targetIdentifier);
if ($collection !== null) {
$list[(string) $targetIdentifier] = $collection;
}
}
return (object) $list;
}
private function collectionExtant(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$sources = $this->requireIdentifiers(
$data,
'targets',
CollectionIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service:collection identifiers'
);
// perform operation
return $this->manager->collectionExtant($tenantId, $userId, $sources);
}
private function collectionCreate(string $tenantId, string $userId, array $data): mixed {
if (!isset($data['target'])) {
throw new InvalidArgumentException(self::ERR_MISSING_TARGET);
}
if (!is_string($data['target'])) {
throw new InvalidArgumentException(self::ERR_INVALID_TARGET);
}
$targetIdentifier = ResourceIdentifier::fromString($data['target']);
if (!$targetIdentifier instanceof ServiceIdentifier || $targetIdentifier instanceof EntityIdentifier) {
throw new InvalidArgumentException('Invalid parameter: target must be provider:service or provider:service:collection');
}
$properties = $this->requireArray($data, 'properties', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
// perform operation
return $this->manager->collectionCreate($tenantId, $userId, $targetIdentifier, $properties, $data['options'] ?? []);
}
private function collectionUpdate(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET,
self::ERR_TARGET_COLLECTION
);
$properties = $this->requireArray($data, 'properties', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
// perform operation
return $this->manager->collectionUpdate($tenantId, $userId, $target, $properties);
}
private function collectionDelete(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET,
self::ERR_TARGET_COLLECTION
);
// perform operation
$result = $this->manager->collectionDelete($tenantId, $userId, $target, $data['options'] ?? []);
if (is_bool($result)) {
return [
'disposition' => 'deleted'
];
}
if ($result instanceof JsonSerializable) {
return [
'disposition' => 'moved',
'mutation' => $result
];
}
return $result;
}
private function collectionCopy(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$sourceIdentifier = $this->requireIdentifier(
$data,
'source',
CollectionIdentifier::class,
self::ERR_MISSING_SOURCE,
self::ERR_INVALID_SOURCE
);
$targetIdentifier = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
// perform operation
return $this->manager->collectionCopy($tenantId, $userId, $targetIdentifier, $sourceIdentifier);
}
private function collectionMove(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$sourceIdentifier = $this->requireIdentifier(
$data,
'source',
CollectionIdentifier::class,
self::ERR_MISSING_SOURCE,
self::ERR_INVALID_SOURCE
);
$targetIdentifier = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
// perform operation
return $this->manager->collectionMove($tenantId, $userId, $targetIdentifier, $sourceIdentifier);
}
// ==================== Entity Operations ====================
private function entityListBulk(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$sources = $this->optionalIdentifiers(
$data,
'sources',
ServiceIdentifier::class,
'Invalid parameter: sources must contain provider:service or provider:service:collection identifiers',
null,
self::ERR_INVALID_SOURCES
);
$filter = $data['filter'] ?? null;
$sort = $data['sort'] ?? null;
$range = $data['range'] ?? null;
// perform operation
return $this->manager->entityListBulk($tenantId, $userId, $sources, $filter, $sort, $range);
}
private function entityListStream(string $tenantId, string $userId, array $data, int $version, string $transaction): StreamedNdJsonResponse {
// Validate parameters
$sources = $this->optionalIdentifiers(
$data,
'sources',
ServiceIdentifier::class,
'Invalid parameter: sources must contain provider:service or provider:service:collection identifiers',
null,
self::ERR_INVALID_SOURCES
);
$filter = $data['filter'] ?? null;
$sort = $data['sort'] ?? null;
$range = $data['range'] ?? null;
// perform operation
$entities = $this->manager->entityListStream($tenantId, $userId, $sources, $filter, $sort, $range);
return new StreamedNdJsonResponse(
$this->streamEnvelope($entities, $version, $transaction),
1,
200,
['Content-Type' => 'application/json'],
);
}
private function entityFetch(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->requireIdentifiers(
$data,
'targets',
EntityIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service:collection:entity identifiers'
);
// perform operation
return (object) $this->manager->entityFetchBulk($tenantId, $userId, ...$targets->all());
}
private function entityExtant(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->requireIdentifiers(
$data,
'targets',
CollectionIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service:collection or provider:service:collection:entity identifiers'
);
// perform operation
return $this->manager->entityExtant($tenantId, $userId, $targets);
}
private function entityDelta(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->requireIdentifiers(
$data,
'targets',
CollectionIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service:collection or provider:service:collection:signature identifiers'
);
// perform operation
return $this->manager->entityDelta($tenantId, $userId, $targets);
}
private function entityCreate(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$properties = $this->requireArray($data, 'properties', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
$options = $data['options'] ?? [];
// perform operation
return $this->manager->entityCreate($tenantId, $userId, $target, $properties, $options);
}
private function entityUpdate(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$properties = $this->requireArray($data, 'properties', self::ERR_MISSING_DATA, self::ERR_INVALID_DATA);
// perform operation
return $this->manager->entityModify($tenantId, $userId, $target, $properties);
}
private function entityDelete(string $tenantId, string $userId, array $data): mixed {
// Validate parameters
$targets = $this->requireIdentifiers(
$data,
'targets',
EntityIdentifier::class,
self::ERR_MISSING_TARGETS,
self::ERR_INVALID_TARGETS,
'Invalid parameter: targets must contain provider:service:collection:entity identifiers'
);
// perform operation
return $this->manager->entityDelete($tenantId, $userId, ...$targets->all());
}
private function entityMove(string $tenantId, string $userId, array $data): mixed {
return $this->entityMopy($tenantId, $userId, $data, 'entityMove');
}
private function entityCopy(string $tenantId, string $userId, array $data): mixed {
return $this->entityMopy($tenantId, $userId, $data, 'entityCopy');
}
private function entityMopy(string $tenantId, string $userId, array $data, string $method): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
CollectionIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$sources = $this->requireIdentifiers(
$data,
'sources',
EntityIdentifier::class,
self::ERR_MISSING_SOURCES,
self::ERR_INVALID_SOURCES,
'Invalid parameter: sources must contain provider:service:collection:entity identifiers'
);
// perform operation
return $this->manager->{$method}($tenantId, $userId, $target, ...$sources->all());
}
// ==================== Entity Content Operations ====================
private function entityRead(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
EntityIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
// perform operation
$content = $this->manager->entityRead($tenantId, $userId, $target);
return [
'content' => $content !== null ? base64_encode($content) : null,
'encoding' => 'base64'
];
}
private function entityReadChunk(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
EntityIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$offset = $this->requireInt($data, 'offset');
$length = $this->requireInt($data, 'length');
// perform operation
$chunk = $this->manager->entityReadChunk(
$tenantId,
$userId,
$target,
$offset,
$length
);
return [
'content' => $chunk !== null ? base64_encode($chunk) : null,
'encoding' => 'base64',
'offset' => $offset,
'length' => $chunk !== null ? strlen($chunk) : 0,
];
}
private function entityWrite(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
EntityIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$content = $this->requireContent($data);
// perform operation
$bytesWritten = $this->manager->entityWrite($tenantId, $userId, $target, $content);
return [
'bytesWritten' => $bytesWritten
];
}
private function entityWriteChunk(string $tenantId, string $userId, array $data = []): mixed {
// Validate parameters
$target = $this->requireIdentifier(
$data,
'target',
EntityIdentifier::class,
self::ERR_MISSING_TARGET,
self::ERR_INVALID_TARGET
);
$offset = $this->requireInt($data, 'offset');
$content = $this->requireContent($data);
// perform operation
$bytesWritten = $this->manager->entityWriteChunk(
$tenantId,
$userId,
$target,
$offset,
$content
);
return [
'bytesWritten' => $bytesWritten,
'offset' => $offset,
];
}
// ==================== Node Operations ====================
private function nodeList(string $tenantId, string $userId, array $data, int $version, string $transaction): StreamedNdJsonResponse {
$sources = $this->requireIdentifiers(
$data,
'sources',
ServiceIdentifier::class,
self::ERR_MISSING_SOURCES,
self::ERR_INVALID_SOURCES,
self::ERR_INVALID_SOURCES,
EntityIdentifier::class
);
$filter = $data['filter'] ?? null;
$sort = $data['sort'] ?? null;
$range = $data['range'] ?? null;
$nodes = $this->manager->nodeList($tenantId, $userId, $sources, $filter, $sort, $range);
return new StreamedNdJsonResponse(
$this->streamEnvelope($nodes, $version, $transaction, static function (\JsonSerializable $node): array {
$data = $node->jsonSerialize();
$data['@type'] = $node instanceof \KTXF\Documents\Collection\CollectionBaseInterface
? 'document:collection'
: 'document:entity';
return $data;
}),
1,
200,
['Content-Type' => 'application/x-ndjson'],
);
}
// ==================== Helper Methods ====================
/**
* Extract and validate a single identifier by instance class
*
* @template T of ResourceIdentifier
* @param class-string<T> $class
* @return T
*/
private function requireIdentifier(array $data, string $key, string $class, string $missingError, string $invalidTypeError, ?string $invalidClassError = null): ResourceIdentifier {
if (!isset($data[$key])) {
throw new InvalidArgumentException($missingError);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException($invalidTypeError);
}
$identifier = ResourceIdentifier::fromString($data[$key]);
if (!$identifier instanceof $class) {
throw new InvalidArgumentException($invalidClassError ?? $invalidTypeError);
}
return $identifier;
}
/**
* Extract and validate a list of identifiers by instance class
*/
private function requireIdentifiers(array $data, string $key, string $class, string $missingError, string $invalidArrayError, string $invalidItemError, ?string $excludeClass = null): ResourceIdentifiers {
if (!isset($data[$key])) {
throw new InvalidArgumentException($missingError);
}
if (!is_array($data[$key])) {
throw new InvalidArgumentException($invalidArrayError);
}
$identifiers = ResourceIdentifiers::fromArray($data[$key]);
foreach ($identifiers as $identifier) {
if (!$identifier instanceof $class || ($excludeClass !== null && $identifier instanceof $excludeClass)) {
throw new InvalidArgumentException($invalidItemError);
}
}
return $identifiers;
}
/**
* Extract and validate a list of identifiers by instance class, returning null if the key is not present
*/
private function optionalIdentifiers(array $data, string $key, string $class, string $invalidItemError, ?string $excludeClass = null, ?string $invalidArrayError = null): ?ResourceIdentifiers {
if (!isset($data[$key])) {
return null;
}
if (!is_array($data[$key])) {
if ($invalidArrayError !== null) {
throw new InvalidArgumentException($invalidArrayError);
}
return null;
}
$identifiers = ResourceIdentifiers::fromArray($data[$key]);
foreach ($identifiers as $identifier) {
if (!$identifier instanceof $class || ($excludeClass !== null && $identifier instanceof $excludeClass)) {
throw new InvalidArgumentException($invalidItemError);
}
}
return $identifiers;
}
/**
* Extract and validate a required string parameter, throwing when it is missing or not a string.
*/
private function requireString(array $data, string $key, string $missingError, string $invalidError): string {
if (!isset($data[$key])) {
throw new InvalidArgumentException($missingError);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException($invalidError);
}
return $data[$key];
}
/**
* Extract and validate a required integer parameter, throwing when it is missing or not an int.
*/
private function requireInt(array $data, string $key, ?string $missingError = null): int {
if (!isset($data[$key]) || !is_int($data[$key])) {
throw new InvalidArgumentException($missingError ?? "Missing parameter: $key");
}
return $data[$key];
}
/**
* Extract and validate a required array parameter, throwing when it is missing or not an array.
*/
private function requireArray(array $data, string $key, string $missingError, string $invalidError): array {
if (!isset($data[$key])) {
throw new InvalidArgumentException($missingError);
}
if (!is_array($data[$key])) {
throw new InvalidArgumentException($invalidError);
}
return $data[$key];
}
/**
* Extract and validate the required 'content' parameter, base64-decoding it when 'encoding' is 'base64'.
*/
private function requireContent(array $data): string {
if (!isset($data['content'])) {
throw new InvalidArgumentException(self::ERR_MISSING_DATA);
}
$content = $data['content'];
if (isset($data['encoding']) && $data['encoding'] === 'base64') {
$content = base64_decode($content);
if ($content === false) {
throw new InvalidArgumentException('Invalid base64 encoded content');
}
}
return $content;
}
/**
* Wrap a generator of JsonSerializable domain objects in the canonical NDJSON
* stream envelope shared by every streaming operation:
*
* control:start {version, transaction, total?} — total? = expected count
* data {data} — one per domain object
* error {message} — on failure, then stop
* control:end {total} — total = objects emitted
*
* If the generator leads with an {@see ExpectedTotal} event, its value is
* folded into the start frame's `total` (the progress denominator) rather
* than emitted as a data frame.
*
* @param \Generator<\JsonSerializable> $items
*/
private function streamEnvelope(\Generator $items, int $version, string $transaction, ?callable $serialize = null): \Generator {
// Peek the first event: an expected-total marker rides on the start frame.
$expected = null;
try {
$items->rewind();
if ($items->valid() && $items->current() instanceof ExpectedTotal) {
$expected = $items->current()->expectedTotal();
$items->next();
}
} catch (\Throwable $t) {
$this->logger->error('Error starting stream', ['exception' => $t]);
yield ['type' => 'error', 'message' => $t->getMessage()];
return;
}
$start = ['type' => 'control', 'status' => 'start', 'version' => $version, 'transaction' => $transaction];
if ($expected !== null) {
$start['total'] = $expected;
}
yield $start;
$total = 0;
try {
for (; $items->valid(); $items->next()) {
$item = $items->current();
if (!$item instanceof \JsonSerializable) {
continue;
}
yield ['type' => 'data', 'data' => $serialize !== null ? $serialize($item) : $item->jsonSerialize()];
$total++;
}
} catch (\Throwable $t) {
$this->logger->error('Error streaming response', ['exception' => $t]);
yield ['type' => 'error', 'message' => $t->getMessage()];
return;
}
yield ['type' => 'control', 'status' => 'end', 'total' => $total];
}
}