Code Coverage |
||||||||||
Lines |
Functions and Methods |
Classes and Traits |
||||||||
| Total | |
84.09% |
74 / 88 |
|
84.62% |
11 / 13 |
CRAP | |
0.00% |
0 / 1 |
| ActionRequests | |
84.09% |
74 / 88 |
|
84.62% |
11 / 13 |
31.16 | |
0.00% |
0 / 1 |
| __construct | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| create | |
100.00% |
8 / 8 |
|
100.00% |
1 / 1 |
2 | |||
| get | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| accept | |
100.00% |
18 / 18 |
|
100.00% |
1 / 1 |
5 | |||
| reject | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
1 | |||
| acknowledge | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
1 | |||
| revoke | |
100.00% |
9 / 9 |
|
100.00% |
1 / 1 |
3 | |||
| fail | |
0.00% |
0 / 7 |
|
0.00% |
0 / 1 |
2 | |||
| applyChange | |
69.57% |
16 / 23 |
|
0.00% |
0 / 1 |
9.80 | |||
| assertTransition | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
2 | |||
| current | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| transition | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| announce | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| 1 | <?php |
| 2 | |
| 3 | declare(strict_types=1); |
| 4 | |
| 5 | namespace LambdaTwelve\OneRecord\Server; |
| 6 | |
| 7 | use LambdaTwelve\OneRecord\Api\AccessDelegation; |
| 8 | use LambdaTwelve\OneRecord\Api\ActionRequest; |
| 9 | use LambdaTwelve\OneRecord\Api\ActionRequestType; |
| 10 | use LambdaTwelve\OneRecord\Api\Error; |
| 11 | use LambdaTwelve\OneRecord\Api\Notification; |
| 12 | use LambdaTwelve\OneRecord\Api\NotificationEventType; |
| 13 | use LambdaTwelve\OneRecord\Api\RequestStatus; |
| 14 | use LambdaTwelve\OneRecord\Api\Subscription; |
| 15 | use LambdaTwelve\OneRecord\Api\Verification; |
| 16 | use LambdaTwelve\OneRecord\Change\Change; |
| 17 | use LambdaTwelve\OneRecord\Change\ChangeApplier; |
| 18 | use LambdaTwelve\OneRecord\Change\ChangeRejected; |
| 19 | use LambdaTwelve\OneRecord\Rdf\Iri; |
| 20 | use LambdaTwelve\OneRecord\Server\Event\ActionRequestCreated; |
| 21 | use LambdaTwelve\OneRecord\Server\Event\ActionRequestStatusChanged; |
| 22 | use LambdaTwelve\OneRecord\Server\Event\LogisticsObjectRevised; |
| 23 | use LambdaTwelve\OneRecord\Server\Notification\Fanout; |
| 24 | use LambdaTwelve\OneRecord\Server\Spi\Grant; |
| 25 | use LambdaTwelve\OneRecord\Server\Spi\OutboundNotification; |
| 26 | use LambdaTwelve\OneRecord\Server\Spi\StoreException; |
| 27 | |
| 28 | /** |
| 29 | * The action-request lifecycle the spec describes, shared by the HTTP |
| 30 | * endpoints and the host's PHP API: creating requests, and moving them |
| 31 | * through accept / reject / acknowledge / revoke with the consequences each |
| 32 | * step has (applying a change, granting access, notifying). |
| 33 | */ |
| 34 | final class ActionRequests |
| 35 | { |
| 36 | public function __construct(private readonly Services $services) {} |
| 37 | |
| 38 | public function create(Change|Subscription|AccessDelegation|Verification $payload, Iri $requestedBy): ActionRequest |
| 39 | { |
| 40 | return $this->services->unitOfWork->run(function () use ($payload, $requestedBy): ActionRequest { |
| 41 | $request = ActionRequest::create($this->services->config->actionRequestIri($this->services->ids->next()), $payload, $requestedBy, $this->services->clock->now()); |
| 42 | $this->services->actionRequests->save($request); |
| 43 | // The Pending notification is queued before any listener runs: a listener that decides the |
| 44 | // request synchronously must not get its decision's notification ahead of this one (R7-001). |
| 45 | if ($request->notifyRequestStatusChange()) { |
| 46 | (new Fanout($this->services))->actionRequestStatusChanged($request); |
| 47 | } |
| 48 | $this->services->dispatcher->dispatch(new ActionRequestCreated($request)); |
| 49 | |
| 50 | // A listener may have decided already; callers get what is stored, not the Pending snapshot. |
| 51 | return $this->services->actionRequests->get($request->iri) ?? $request; |
| 52 | }); |
| 53 | } |
| 54 | |
| 55 | public function get(Iri $iri): ?ActionRequest |
| 56 | { |
| 57 | return $this->services->actionRequests->get($iri); |
| 58 | } |
| 59 | |
| 60 | /** |
| 61 | * Accept: a change is applied and becomes a new revision (or the request |
| 62 | * fails with the errors recorded), a subscription becomes active, an |
| 63 | * access delegation becomes grants. Other pending changes written against |
| 64 | * the same revision are rejected, as the spec requires. |
| 65 | */ |
| 66 | public function accept(ActionRequest $request, Iri $by): ActionRequest |
| 67 | { |
| 68 | return $this->services->unitOfWork->run(function () use ($request, $by): ActionRequest { |
| 69 | $request = $this->current($request); |
| 70 | $this->assertTransition($request, RequestStatus::Accepted); |
| 71 | $now = $this->services->clock->now(); |
| 72 | $accepted = $request->withStatus(RequestStatus::Accepted, $now, $by); |
| 73 | |
| 74 | // The compare-and-set on the status is the decision, and it succeeds exactly once; every |
| 75 | // side effect (grants, a revision, notifications) comes after it, so a decision that lost |
| 76 | // the race writes nothing even under a host without a transactional unit of work. The |
| 77 | // event and the status notification come last, once the consequences are in place, so a |
| 78 | // listener reading the policy or the object sees the state after the decision. |
| 79 | $previous = $request->status; |
| 80 | $this->transition($accepted, $previous); |
| 81 | $final = $accepted; |
| 82 | if ($request->payload instanceof Change) { |
| 83 | $final = $this->applyChange($accepted, $request->payload, $by); |
| 84 | } elseif ($request->payload instanceof AccessDelegation) { |
| 85 | foreach ($request->payload->delegates as $delegate) { |
| 86 | foreach ($request->payload->logisticsObjects as $object) { |
| 87 | $this->services->delegations->grant(new Grant($delegate, $object, $request->payload->permissions, $request->payload->expiresAt, $request->iri)); |
| 88 | $stored = $this->services->objects->latest($object); |
| 89 | $this->services->outbox->enqueue(new OutboundNotification($delegate, new Notification(NotificationEventType::LogisticsObjectAccessGranted, $object, $stored?->object->mostSpecificType($this->services->vocabulary), $request->iri), $now, $this->services->ids->next())); |
| 90 | } |
| 91 | } |
| 92 | } |
| 93 | return $this->announce($final, $previous); |
| 94 | }); |
| 95 | } |
| 96 | |
| 97 | /** |
| 98 | * @param list<Error> $errors |
| 99 | */ |
| 100 | public function reject(ActionRequest $request, Iri $by, array $errors = []): ActionRequest |
| 101 | { |
| 102 | return $this->services->unitOfWork->run(function () use ($request, $by, $errors): ActionRequest { |
| 103 | $request = $this->current($request); |
| 104 | $this->assertTransition($request, RequestStatus::Rejected); |
| 105 | $rejected = $request->withStatus(RequestStatus::Rejected, $this->services->clock->now(), $by, $errors); |
| 106 | $this->transition($rejected, $request->status); |
| 107 | |
| 108 | return $this->announce($rejected, $request->status); |
| 109 | }); |
| 110 | } |
| 111 | |
| 112 | public function acknowledge(ActionRequest $request, Iri $by): ActionRequest |
| 113 | { |
| 114 | return $this->services->unitOfWork->run(function () use ($request, $by): ActionRequest { |
| 115 | $request = $this->current($request); |
| 116 | $this->assertTransition($request, RequestStatus::Acknowledged); |
| 117 | $acknowledged = $request->withStatus(RequestStatus::Acknowledged, $this->services->clock->now(), $by); |
| 118 | $this->transition($acknowledged, $request->status); |
| 119 | |
| 120 | return $this->announce($acknowledged, $request->status); |
| 121 | }); |
| 122 | } |
| 123 | |
| 124 | public function revoke(ActionRequest $request, Iri $by): ActionRequest |
| 125 | { |
| 126 | return $this->services->unitOfWork->run(function () use ($request, $by): ActionRequest { |
| 127 | $request = $this->current($request); |
| 128 | $this->assertTransition($request, RequestStatus::Revoked); |
| 129 | $revoked = $request->withStatus(RequestStatus::Revoked, $this->services->clock->now(), $by); |
| 130 | $this->transition($revoked, $request->status); |
| 131 | if ($request->type === ActionRequestType::AccessDelegation && $request->status === RequestStatus::Accepted) { |
| 132 | $this->services->delegations->revokeFrom($request->iri); |
| 133 | } |
| 134 | |
| 135 | return $this->announce($revoked, $request->status); |
| 136 | }); |
| 137 | } |
| 138 | |
| 139 | /** |
| 140 | * @param list<Error> $errors |
| 141 | */ |
| 142 | public function fail(ActionRequest $request, array $errors): ActionRequest |
| 143 | { |
| 144 | return $this->services->unitOfWork->run(function () use ($request, $errors): ActionRequest { |
| 145 | $request = $this->current($request); |
| 146 | $this->assertTransition($request, RequestStatus::Failed); |
| 147 | $failed = $request->withStatus(RequestStatus::Failed, $this->services->clock->now(), null, $errors); |
| 148 | $this->transition($failed, $request->status); |
| 149 | |
| 150 | return $this->announce($failed, $request->status); |
| 151 | }); |
| 152 | } |
| 153 | |
| 154 | /** |
| 155 | * Called once the request is stored as Accepted: applying the change may |
| 156 | * still fail, which moves the request on to Failed. |
| 157 | */ |
| 158 | private function applyChange(ActionRequest $accepted, Change $change, Iri $by): ActionRequest |
| 159 | { |
| 160 | $current = $this->services->objects->latest($change->logisticsObject); |
| 161 | if ($current === null) { |
| 162 | $failed = $accepted->withStatus(RequestStatus::Failed, $this->services->clock->now(), null, [Error::of('Resource not found', '404', 'The logistics object no longer exists.', null, $change->logisticsObject->value)]); |
| 163 | $this->transition($failed, RequestStatus::Accepted); |
| 164 | |
| 165 | return $failed; |
| 166 | } |
| 167 | $applier = new ChangeApplier($this->services->vocabulary, $this->services->embeddedIds); |
| 168 | try { |
| 169 | $result = $applier->apply($current->object, $current->revision, $change); |
| 170 | $stored = $this->services->objects->saveRevision($result->object, $current->revision, $this->services->clock->now()); |
| 171 | } catch (ChangeRejected $e) { |
| 172 | $failed = $accepted->withStatus(RequestStatus::Failed, $this->services->clock->now(), null, $e->errors); |
| 173 | $this->transition($failed, RequestStatus::Accepted); |
| 174 | |
| 175 | return $failed; |
| 176 | } catch (StoreException $e) { |
| 177 | $failed = $accepted->withStatus(RequestStatus::Failed, $this->services->clock->now(), null, [Error::of('Conflict with Logistics Object revision number', '409', $e->getMessage(), null, $change->logisticsObject->value)]); |
| 178 | $this->transition($failed, RequestStatus::Accepted); |
| 179 | |
| 180 | return $failed; |
| 181 | } |
| 182 | |
| 183 | Deprecations::log($stored->object->graph, $stored->object->iri, $this->services->vocabulary, $this->services->logger); |
| 184 | $this->services->dispatcher->dispatch(new LogisticsObjectRevised($stored, $accepted->iri, $result->changedProperties)); |
| 185 | (new Fanout($this->services))->logisticsObjectUpdated($stored, $result->changedProperties, $accepted->iri); |
| 186 | |
| 187 | // Every other pending change written against the revision just replaced is now stale. |
| 188 | foreach ($this->services->actionRequests->pendingChanges($change->logisticsObject) as $other) { |
| 189 | if (!$other->iri->equals($accepted->iri) && $other->payload instanceof Change && $other->payload->revision === $change->revision) { |
| 190 | $this->reject($other, $by, [Error::of('Conflict with Logistics Object revision number', '409', \sprintf('Another change to revision %d was accepted first; the object is now at revision %d.', $change->revision, $stored->revision), null, $change->logisticsObject->value)]); |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | return $accepted; |
| 195 | } |
| 196 | |
| 197 | private function assertTransition(ActionRequest $request, RequestStatus $next): void |
| 198 | { |
| 199 | if (!$request->canTransitionTo($next)) { |
| 200 | throw new IllegalTransition($request, $next); |
| 201 | } |
| 202 | } |
| 203 | |
| 204 | /** |
| 205 | * The stored state of a request, whatever snapshot the caller holds: a |
| 206 | * decision made on a stale copy would otherwise undo a newer one (AR-002). |
| 207 | */ |
| 208 | private function current(ActionRequest $request): ActionRequest |
| 209 | { |
| 210 | return $this->services->actionRequests->get($request->iri) ?? throw StoreException::notFound($request->iri); |
| 211 | } |
| 212 | |
| 213 | /** |
| 214 | * The compare-and-set that is the decision: it succeeds for exactly one worker. |
| 215 | */ |
| 216 | private function transition(ActionRequest $request, RequestStatus $previous): void |
| 217 | { |
| 218 | $this->services->actionRequests->transition($request, $previous); |
| 219 | } |
| 220 | |
| 221 | /** |
| 222 | * Tells the requester and then the host about a decision whose side |
| 223 | * effects are already in place; $previous is the status the request had |
| 224 | * before the decision, whatever intermediate state the store recorded on |
| 225 | * the way. The notification is queued before the event, so a listener that |
| 226 | * advances the request again cannot put its notification ahead of this |
| 227 | * one; and since a listener may advance it, the stored request is what is |
| 228 | * returned (R7-001). |
| 229 | */ |
| 230 | private function announce(ActionRequest $request, RequestStatus $previous): ActionRequest |
| 231 | { |
| 232 | (new Fanout($this->services))->actionRequestStatusChanged($request); |
| 233 | $this->services->dispatcher->dispatch(new ActionRequestStatusChanged($request, $previous)); |
| 234 | |
| 235 | return $this->services->actionRequests->get($request->iri) ?? $request; |
| 236 | } |
| 237 | } |