Code Coverage |
||||||||||
Lines |
Functions and Methods |
Classes and Traits |
||||||||
| Total | |
93.85% |
61 / 65 |
|
85.71% |
12 / 14 |
CRAP | |
0.00% |
0 / 1 |
| DataHolder | |
93.85% |
61 / 65 |
|
85.71% |
12 / 14 |
24.13 | |
0.00% |
0 / 1 |
| __construct | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| create | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
1 | |||
| update | |
100.00% |
8 / 8 |
|
100.00% |
1 / 1 |
2 | |||
| change | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
3 | |||
| publish | |
100.00% |
12 / 12 |
|
100.00% |
1 / 1 |
5 | |||
| accept | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| reject | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| acknowledge | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| revoke | |
0.00% |
0 / 3 |
|
0.00% |
0 / 1 |
2 | |||
| subscribe | |
100.00% |
4 / 4 |
|
100.00% |
1 / 1 |
2 | |||
| announce | |
100.00% |
5 / 5 |
|
100.00% |
1 / 1 |
1 | |||
| forget | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
3 | |||
| actionRequests | |
0.00% |
0 / 1 |
|
0.00% |
0 / 1 |
2 | |||
| require | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| 1 | <?php |
| 2 | |
| 3 | declare(strict_types=1); |
| 4 | |
| 5 | namespace LambdaTwelve\OneRecord\Server; |
| 6 | |
| 7 | use InvalidArgumentException; |
| 8 | use LambdaTwelve\OneRecord\Api\ActionRequest; |
| 9 | use LambdaTwelve\OneRecord\Api\Error; |
| 10 | use LambdaTwelve\OneRecord\Api\Notification; |
| 11 | use LambdaTwelve\OneRecord\Api\NotificationEventType; |
| 12 | use LambdaTwelve\OneRecord\Api\RequestStatus; |
| 13 | use LambdaTwelve\OneRecord\Api\Subscription; |
| 14 | use LambdaTwelve\OneRecord\Change\Change; |
| 15 | use LambdaTwelve\OneRecord\Change\ChangeBuilder; |
| 16 | use LambdaTwelve\OneRecord\Model\IriMinter; |
| 17 | use LambdaTwelve\OneRecord\Model\LocalGraph; |
| 18 | use LambdaTwelve\OneRecord\Model\LogisticsObject; |
| 19 | use LambdaTwelve\OneRecord\Model\ResolvedGraph; |
| 20 | use LambdaTwelve\OneRecord\Rdf\Iri; |
| 21 | use LambdaTwelve\OneRecord\Server\Event\LogisticsObjectCreated; |
| 22 | use LambdaTwelve\OneRecord\Server\Notification\Fanout; |
| 23 | use LambdaTwelve\OneRecord\Server\Spi\OutboundNotification; |
| 24 | use LambdaTwelve\OneRecord\Server\Spi\StoredObject; |
| 25 | |
| 26 | /** |
| 27 | * The holder's side of the API, in PHP: publish objects, decide action |
| 28 | * requests, share and forget. This is what a host calls from its own code; |
| 29 | * partners reach the same state through the HTTP endpoints. |
| 30 | * |
| 31 | * Changes the holder makes to its own objects are recorded as accepted |
| 32 | * change requests in the audit trail, as the spec asks for a complete |
| 33 | * history whatever the source of a change. |
| 34 | */ |
| 35 | final class DataHolder |
| 36 | { |
| 37 | private readonly ActionRequests $requests; |
| 38 | |
| 39 | public function __construct(private readonly Services $services) |
| 40 | { |
| 41 | $this->requests = new ActionRequests($services); |
| 42 | } |
| 43 | |
| 44 | /** |
| 45 | * Create an object that does not exist yet. |
| 46 | */ |
| 47 | public function create(LogisticsObject $object): StoredObject |
| 48 | { |
| 49 | return $this->services->unitOfWork->run(function () use ($object): StoredObject { |
| 50 | $stored = $this->services->objects->create($object->withEmbeddedIds($this->services->embeddedIds), $this->services->clock->now()); |
| 51 | Deprecations::log($stored->object->graph, $stored->object->iri, $this->services->vocabulary, $this->services->logger); |
| 52 | $this->services->dispatcher->dispatch(new LogisticsObjectCreated($stored, $this->services->config->dataHolder)); |
| 53 | (new Fanout($this->services))->logisticsObjectCreated($stored); |
| 54 | |
| 55 | return $stored; |
| 56 | }); |
| 57 | } |
| 58 | |
| 59 | /** |
| 60 | * Bring a stored object to this version: computed as a change, recorded as |
| 61 | * an accepted change request, applied as a new revision. Null when nothing |
| 62 | * differs. |
| 63 | */ |
| 64 | public function update(LogisticsObject $object, ?string $description = null): ?ActionRequest |
| 65 | { |
| 66 | return $this->services->unitOfWork->run(function () use ($object, $description): ?ActionRequest { |
| 67 | $current = $this->services->objects->latest($object->iri) |
| 68 | ?? throw new InvalidArgumentException(\sprintf('"%s" does not exist; create it first.', $object->iri->value)); |
| 69 | $change = (new ChangeBuilder($this->services->vocabulary))->diff($current->object, $object, $current->revision, $description); |
| 70 | if ($change === null) { |
| 71 | return null; |
| 72 | } |
| 73 | |
| 74 | return $this->change($change); |
| 75 | }); |
| 76 | } |
| 77 | |
| 78 | /** |
| 79 | * Apply a change of the holder's own making: created and accepted at once. |
| 80 | * |
| 81 | * @throws ChangeFailed when the change was not applied (it failed, or a creation listener rejected it); the request is kept only if the unit of work commits it |
| 82 | */ |
| 83 | public function change(Change $change): ActionRequest |
| 84 | { |
| 85 | return $this->services->unitOfWork->run(function () use ($change): ActionRequest { |
| 86 | $request = $this->requests->create($change, $this->services->config->dataHolder); |
| 87 | // A creation listener may already have decided it (R7-001); deciding twice is an illegal transition. |
| 88 | $decided = $request->status === RequestStatus::Pending ? $this->requests->accept($request, $this->services->config->dataHolder) : $request; |
| 89 | if ($decided->status !== RequestStatus::Accepted) { |
| 90 | // Returning a failed or rejected request read as success to callers checking for null (AR-027). |
| 91 | throw new ChangeFailed($decided); |
| 92 | } |
| 93 | |
| 94 | return $decided; |
| 95 | }); |
| 96 | } |
| 97 | |
| 98 | /** |
| 99 | * Create or update every object of a graph, idempotently: unchanged objects |
| 100 | * produce nothing, changed ones a new revision. Pass the URIs of objects |
| 101 | * already published so they keep them. A change that cannot be applied |
| 102 | * throws ChangeFailed out of the unit of work, so with a transactional |
| 103 | * UnitOfWork nothing of the graph is kept; without one, objects handled |
| 104 | * before the failure stay published. |
| 105 | * |
| 106 | * @param array<string, Iri> $existing local key => URI |
| 107 | * @throws ChangeFailed |
| 108 | */ |
| 109 | public function publish(LocalGraph|ResolvedGraph $graph, ?IriMinter $minter = null, array $existing = []): PublishResult |
| 110 | { |
| 111 | return $this->services->unitOfWork->run(function () use ($graph, $minter, $existing): PublishResult { |
| 112 | if ($graph instanceof LocalGraph) { |
| 113 | $graph = $graph->resolve($minter ?? throw new InvalidArgumentException('A LocalGraph needs an IriMinter to publish.'), $existing); |
| 114 | } |
| 115 | $outcomes = []; |
| 116 | foreach ($graph->objects as $key => $object) { |
| 117 | if (!$this->services->objects->exists($object->iri)) { |
| 118 | $this->create($object); |
| 119 | $outcomes[$key] = PublishResult::CREATED; |
| 120 | continue; |
| 121 | } |
| 122 | $outcomes[$key] = $this->update($object, 'Republished by the data holder') === null ? PublishResult::UNCHANGED : PublishResult::UPDATED; |
| 123 | } |
| 124 | |
| 125 | return new PublishResult($graph, $outcomes); |
| 126 | }); |
| 127 | } |
| 128 | |
| 129 | public function accept(Iri $request): ActionRequest |
| 130 | { |
| 131 | return $this->services->unitOfWork->run(function () use ($request): ActionRequest { |
| 132 | return $this->requests->accept($this->require($request), $this->services->config->dataHolder); |
| 133 | }); |
| 134 | } |
| 135 | |
| 136 | /** |
| 137 | * @param list<Error> $errors |
| 138 | */ |
| 139 | public function reject(Iri $request, array $errors = []): ActionRequest |
| 140 | { |
| 141 | return $this->services->unitOfWork->run(function () use ($request, $errors): ActionRequest { |
| 142 | return $this->requests->reject($this->require($request), $this->services->config->dataHolder, $errors); |
| 143 | }); |
| 144 | } |
| 145 | |
| 146 | public function acknowledge(Iri $request): ActionRequest |
| 147 | { |
| 148 | return $this->services->unitOfWork->run(function () use ($request): ActionRequest { |
| 149 | return $this->requests->acknowledge($this->require($request), $this->services->config->dataHolder); |
| 150 | }); |
| 151 | } |
| 152 | |
| 153 | public function revoke(Iri $request): ActionRequest |
| 154 | { |
| 155 | return $this->services->unitOfWork->run(function () use ($request): ActionRequest { |
| 156 | return $this->requests->revoke($this->require($request), $this->services->config->dataHolder); |
| 157 | }); |
| 158 | } |
| 159 | |
| 160 | /** |
| 161 | * Publisher-initiated subscription (spec: "Get Subscription information as |
| 162 | * Publisher"): the subscription a partner answered with becomes an |
| 163 | * accepted SubscriptionRequest, so notifications can reference it and the |
| 164 | * partner can revoke it. |
| 165 | */ |
| 166 | public function subscribe(Subscription $subscription): ActionRequest |
| 167 | { |
| 168 | return $this->services->unitOfWork->run(function () use ($subscription): ActionRequest { |
| 169 | $request = $this->requests->create($subscription, $subscription->subscriber); |
| 170 | |
| 171 | // A creation listener may already have decided it (R8-001); its decision stands and the |
| 172 | // caller reads the status it got, instead of this method deciding a second time. |
| 173 | return $request->status === RequestStatus::Pending ? $this->requests->accept($request, $this->services->config->dataHolder) : $request; |
| 174 | }); |
| 175 | } |
| 176 | |
| 177 | /** |
| 178 | * Tell a partner an object exists (LOGISTICS_OBJECT_AVAILABLE): the spec's |
| 179 | * way of sharing a URI so the partner can read or subscribe to it. |
| 180 | */ |
| 181 | public function announce(Iri $object, Iri $recipient): void |
| 182 | { |
| 183 | $this->services->unitOfWork->run(function () use ($object, $recipient): void { |
| 184 | $stored = $this->services->objects->latest($object) ?? throw new InvalidArgumentException(\sprintf('"%s" does not exist.', $object->value)); |
| 185 | $notification = new Notification(NotificationEventType::LogisticsObjectAvailable, $object, $stored->object->mostSpecificType($this->services->vocabulary)); |
| 186 | $this->services->outbox->enqueue(new OutboundNotification($recipient, $notification, $this->services->clock->now(), $this->services->ids->next())); |
| 187 | }); |
| 188 | } |
| 189 | |
| 190 | /** |
| 191 | * Forget an object: every revision is erased, and by default its events |
| 192 | * and the grants on it too. The API has no delete, so this is the host's |
| 193 | * data-protection operation. Keep events (false) where status history |
| 194 | * must outlive the object; action requests are never erased, they are |
| 195 | * other parties' history too. The host's own access rules must stop |
| 196 | * answering for the URI as well. |
| 197 | */ |
| 198 | public function forget(Iri $object, bool $events = true, bool $grants = true): void |
| 199 | { |
| 200 | $this->services->unitOfWork->run(function () use ($object, $events, $grants): void { |
| 201 | $this->services->objects->erase($object); |
| 202 | if ($events) { |
| 203 | $this->services->events->eraseFor($object); |
| 204 | } |
| 205 | if ($grants) { |
| 206 | $this->services->delegations->eraseFor($object); |
| 207 | } |
| 208 | }); |
| 209 | } |
| 210 | |
| 211 | public function actionRequests(): ActionRequests |
| 212 | { |
| 213 | return $this->requests; |
| 214 | } |
| 215 | |
| 216 | private function require(Iri $iri): ActionRequest |
| 217 | { |
| 218 | return $this->requests->get($iri) ?? throw new InvalidArgumentException(\sprintf('No action request "%s".', $iri->value)); |
| 219 | } |
| 220 | } |