| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
| Expand Up | @@ -509,18 +509,63 @@ private function queueOutgoing(Request|Notification|Response|Error $message, arr | |||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Consume (get and clear) all outgoing messages for a session. | ||||||
| * Consume (get and clear) outgoing messages for a session. | ||||||
| * | ||||||
| * With response ids only the matching responses are taken, the rest stays queued. | ||||||
| * A null id matches an error without an id. | ||||||
| * | ||||||
| * @param list<int|string|null>|null $responseIds | ||||||
| * | ||||||
| * @return array<int, array{message: string, context: array<string, mixed>}> | ||||||
| */ | ||||||
| public function consumeOutgoingMessages(Uuid $sessionId): array | ||||||
| public function consumeOutgoingMessages(Uuid $sessionId, ?array $responseIds = null): array | ||||||
| { | ||||||
| $session = $this->sessionManager->createWithId($sessionId); | ||||||
| /** @var array<int, array{message: string, context: array<string, mixed>}> $queue */ | ||||||
| $queue = $session->get(self::SESSION_OUTGOING_QUEUE, []); | ||||||
| $session->set(self::SESSION_OUTGOING_QUEUE, []); | ||||||
| $session->save(); | ||||||
|
|
||||||
| return $queue; | ||||||
| if (null === $responseIds) { | ||||||
| $session->set(self::SESSION_OUTGOING_QUEUE, []); | ||||||
| $session->save(); | ||||||
|
|
||||||
| return $queue; | ||||||
| } | ||||||
|
|
||||||
| $consumed = []; | ||||||
| $remaining = []; | ||||||
| foreach ($queue as $message) { | ||||||
| if (self::isResponseTo($message, $responseIds)) { | ||||||
| $consumed[] = $message; | ||||||
| } else { | ||||||
| $remaining[] = $message; | ||||||
| } | ||||||
| } | ||||||
|
|
||||||
| if ([] !== $consumed) { | ||||||
| $session->set(self::SESSION_OUTGOING_QUEUE, $remaining); | ||||||
| $session->save(); | ||||||
| } | ||||||
|
|
||||||
| return $consumed; | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * @param array{message: string, context: array<string, mixed>} $message | ||||||
| * @param list<int|string|null> $responseIds | ||||||
| */ | ||||||
| private static function isResponseTo(array $message, array $responseIds): bool | ||||||
|
Comment thread
Copy link
Copy Markdown
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality
Suggested change
Sorry, something went wrong.
All reactions
|
||||||
| { | ||||||
| if ('response' !== ($message['context']['type'] ?? null)) { | ||||||
| return false; | ||||||
| } | ||||||
|
|
||||||
| try { | ||||||
| $decoded = json_decode($message['message'], true, flags: \JSON_THROW_ON_ERROR); | ||||||
| } catch (\JsonException) { | ||||||
| return false; | ||||||
| } | ||||||
|
|
||||||
| return \is_array($decoded) && \in_array($decoded['id'] ?? null, $responseIds, true); | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| Expand Down | ||||||
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
| Expand Up | @@ -74,6 +74,10 @@ class StreamableHttpTransport extends BaseTransport implements StatelessAwareTra | |||||
| private ?string $immediateResponse = null; | ||||||
| private ?int $immediateStatusCode = null; | ||||||
|
|
||||||
| /** @var list<int|string|null> */ | ||||||
| private array $expectedResponseIds = []; | ||||||
| private bool $batchRequest = false; | ||||||
|
|
||||||
| /** @var list<MiddlewareInterface>|null null until {@see self::listen()} resolves the defaults */ | ||||||
| private ?array $middleware; | ||||||
|
|
||||||
| Expand Down Expand Up | @@ -186,6 +190,9 @@ protected function handleOptionsRequest(): ResponseInterface | |||||
| */ | ||||||
| protected function handlePostRequest(string $body): ResponseInterface | ||||||
| { | ||||||
| // Concurrent requests of one session share the outgoing queue. | ||||||
| [$this->expectedResponseIds, $this->batchRequest] = self::expectedResponses($body); | ||||||
|
|
||||||
| $this->handleMessage($body, $this->sessionId); | ||||||
|
|
||||||
| // Consume the immediate response exactly once, so a transport instance | ||||||
| Expand Down Expand Up | @@ -223,15 +230,15 @@ protected function handleDeleteRequest(): ResponseInterface | |||||
|
|
||||||
| protected function createJsonResponse(): ResponseInterface | ||||||
| { | ||||||
| $outgoingMessages = $this->getOutgoingMessages($this->sessionId); | ||||||
| $outgoingMessages = $this->getOutgoingResponses($this->sessionId, $this->expectedResponseIds); | ||||||
|
|
||||||
| if (empty($outgoingMessages)) { | ||||||
| return $this->responseFactory->createResponse(202) | ||||||
| ->withHeader('Content-Type', 'application/json'); | ||||||
| } | ||||||
|
|
||||||
| $messages = array_column($outgoingMessages, 'message'); | ||||||
| $responseBody = 1 === \count($messages) ? $messages[0] : '['.implode(',', $messages).']'; | ||||||
| $responseBody = $this->batchRequest ? '['.implode(',', $messages).']' : $messages[0]; | ||||||
|
|
||||||
| $response = $this->responseFactory->createResponse(200) | ||||||
| ->withHeader('Content-Type', 'application/json') | ||||||
| Expand Down Expand Up | @@ -461,6 +468,39 @@ private function handleModernRequest(ServerRequestInterface $request, string $bo | |||||
| return $this->responder->respond($this->stateless->handle($body, self::headers($request))); | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * Ids the body expects responses for (null for a message without a usable id), and whether it is a batch. | ||||||
| * | ||||||
| * @return array{list<int|string|null>, bool} | ||||||
| */ | ||||||
| private static function expectedResponses(string $body): array | ||||||
|
Comment thread
Copy link
Copy Markdown
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality
Suggested change
Sorry, something went wrong.
All reactions
|
||||||
| { | ||||||
| try { | ||||||
| $data = json_decode($body, true, flags: \JSON_THROW_ON_ERROR); | ||||||
| } catch (\JsonException) { | ||||||
| // Parse errors are sent directly, not queued. | ||||||
| return [[], false]; | ||||||
| } | ||||||
|
|
||||||
| if (!\is_array($data) || [] === $data) { | ||||||
| return [[null], false]; | ||||||
| } | ||||||
|
|
||||||
| $batch = array_is_list($data); | ||||||
| $ids = []; | ||||||
|
|
||||||
| foreach ($batch ? $data : [$data] as $message) { | ||||||
| if (\is_array($message) && (isset($message['result']) || isset($message['error']))) { | ||||||
| continue; | ||||||
| } | ||||||
|
|
||||||
| $id = \is_array($message) ? ($message['id'] ?? null) : null; | ||||||
| $ids[] = \is_int($id) || \is_string($id) ? $id : null; | ||||||
| } | ||||||
|
|
||||||
| return [$ids, $batch]; | ||||||
| } | ||||||
|
|
||||||
| /** | ||||||
| * @return array<string, string> | ||||||
| */ | ||||||
| Expand Down | ||||||
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
| Expand Up | @@ -84,11 +84,12 @@ public function onMessage(callable $listener): void; | |||||
| public function onSessionEnd(callable $listener): void; | ||||||
|
|
||||||
| /** | ||||||
| * Set a provider function to retrieve all queued outgoing messages. | ||||||
| * Set a provider function to retrieve queued outgoing messages. | ||||||
| * | ||||||
| * The transport calls this to retrieve all queued messages for a session. | ||||||
| * The transport calls this to retrieve queued messages for a session. When response ids | ||||||
| * are passed, only the responses to those requests are returned. | ||||||
| * | ||||||
| * @param callable(Uuid $sessionId): array<int, array{message: string, context: array<string, mixed>}> $provider | ||||||
| * @param callable(Uuid $sessionId, list<int|string|null>|null $responseIds=): array<int, array{message: string, context: array<string, mixed>}> $provider | ||||||
|
Comment thread
Copy link
Copy Markdown
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality
Suggested change
Sorry, something went wrong.
All reactions
|
||||||
| */ | ||||||
| public function setOutgoingMessagesProvider(callable $provider): void; | ||||||
|
|
||||||
| Expand Down | ||||||
| Back | FazBrowse Home | New Git URL |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Qualitythis make sense as solution but introduces an issue of potentially piling up $remaining responses in the memory, right?
Sorry, something went wrong.
Uh oh!
There was an error while loading. Please reload this page.