feat: node list back end

Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
2026-09-04 22:12:27 -04:00
parent 93aa1b0f24
commit 7d18bcab1a
2 changed files with 214 additions and 98 deletions
+134 -98
View File
@@ -163,56 +163,13 @@ class DefaultController extends ControllerAbstract {
'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)
};
}
// ==================== Identifier Helpers ====================
/**
* Parse a required collection target identifier (provider:service:collection)
*/
private function collectionTarget(array $data, string $key = 'target'): CollectionIdentifier {
if (!isset($data[$key])) {
throw new InvalidArgumentException('Missing parameter: ' . $key);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException("Invalid parameter: $key must be a string");
}
$identifier = ResourceIdentifier::fromString($data[$key]);
if (!$identifier instanceof CollectionIdentifier) {
throw new InvalidArgumentException(self::ERR_TARGET_COLLECTION);
}
return $identifier;
}
/**
* Parse an optional collection target identifier (absent = root)
*/
private function collectionTargetOptional(array $data, string $key = 'target'): ?CollectionIdentifier {
if (!isset($data[$key]) || $data[$key] === null || $data[$key] === '') {
return null;
}
return $this->collectionTarget($data, $key);
}
/**
* Parse a required entity target identifier (provider:service:collection:entity)
*/
private function entityTarget(array $data, string $key = 'target'): EntityIdentifier {
if (!isset($data[$key])) {
throw new InvalidArgumentException('Missing parameter: ' . $key);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException("Invalid parameter: $key must be a string");
}
$identifier = ResourceIdentifier::fromString($data[$key]);
if (!$identifier instanceof EntityIdentifier) {
throw new InvalidArgumentException(self::ERR_TARGET_ENTITY);
}
return $identifier;
}
// ==================== Provider Operations ====================
private function providerList(string $tenantId, string $userId, array $data): mixed {
@@ -620,55 +577,6 @@ class DefaultController extends ControllerAbstract {
);
}
/**
* 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): \Generator {
// Peek the first event: an expected-total marker rides on the start frame.
$expected = null;
$items->rewind();
if ($items->valid() && $items->current() instanceof ExpectedTotal) {
$expected = $items->current()->expectedTotal();
$items->next();
}
$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' => $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];
}
private function entityFetch(string $tenantId, string $userId, array $data): mixed {
if (!isset($data['targets'])) {
throw new InvalidArgumentException(self::ERR_MISSING_TARGETS);
@@ -780,9 +688,6 @@ class DefaultController extends ControllerAbstract {
return $this->entityRelocate($tenantId, $userId, $data, 'entityCopy');
}
/**
* Shared request handling for entity move/copy operations
*/
private function entityRelocate(string $tenantId, string $userId, array $data, string $method): mixed {
if (!isset($data['sources'])) {
throw new InvalidArgumentException(self::ERR_MISSING_SOURCES);
@@ -898,4 +803,135 @@ class DefaultController extends ControllerAbstract {
];
}
// ==================== Node Operations ====================
private function nodeList(string $tenantId, string $userId, array $data, int $version, string $transaction): StreamedNdJsonResponse {
if (isset($data['sources'])) {
if (!is_array($data['sources'])) {
throw new InvalidArgumentException(self::ERR_INVALID_SOURCES);
}
$sources = ResourceIdentifiers::fromArray($data['sources']);
foreach ($sources as $source) {
if (!$source instanceof ServiceIdentifier && !$source instanceof CollectionIdentifier) {
throw new InvalidArgumentException('Invalid parameter: sources must contain provider:service or provider:service:collection identifiers');
}
}
} else {
$sources = null;
}
$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 ====================
private function collectionTarget(array $data, string $key = 'target'): CollectionIdentifier {
if (!isset($data[$key])) {
throw new InvalidArgumentException('Missing parameter: ' . $key);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException("Invalid parameter: $key must be a string");
}
$identifier = ResourceIdentifier::fromString($data[$key]);
if (!$identifier instanceof CollectionIdentifier) {
throw new InvalidArgumentException(self::ERR_TARGET_COLLECTION);
}
return $identifier;
}
private function collectionTargetOptional(array $data, string $key = 'target'): ?CollectionIdentifier {
if (!isset($data[$key]) || $data[$key] === null || $data[$key] === '') {
return null;
}
return $this->collectionTarget($data, $key);
}
private function entityTarget(array $data, string $key = 'target'): EntityIdentifier {
if (!isset($data[$key])) {
throw new InvalidArgumentException('Missing parameter: ' . $key);
}
if (!is_string($data[$key])) {
throw new InvalidArgumentException("Invalid parameter: $key must be a string");
}
$identifier = ResourceIdentifier::fromString($data[$key]);
if (!$identifier instanceof EntityIdentifier) {
throw new InvalidArgumentException(self::ERR_TARGET_ENTITY);
}
return $identifier;
}
/**
* 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];
}
}