diff --git a/README.md b/README.md index 37665cf..1f6bec2 100644 --- a/README.md +++ b/README.md @@ -92,6 +92,8 @@ локально битые — `BatchRejection(index, detail)` без HTTP; ответ сервера 202 **и** 422 («ни одного не принято») парсится одинаково — `index` серверного отклонения ремапится на исходную позицию; 5xx/транспорт — исключение целиком. - `status(string $eventId): EventStatus` — uuid-проверка, 404 → `NotFoundException`. +- `waitForStatus(string $eventId, array $statuses = ['done', 'failed'], float $timeoutSeconds = 30, float $intervalSeconds = 1): EventStatus` — поллинг до целевого статуса (для приёмки/тестов, не бизнес-кода); морг транспорта не прерывает ожидание; не дождались — `StatusTimeoutException`. +- `send(..., ?string $userId = null)` / `emit(..., ?string $userId = null)` — конвенция `payload.user_id` сразу параметром (перебивает payload'овский). - `health(): array`, `ready(): array` — `/api/healthz` и `/api/readyz`, без ключа. ## Исключения (`GNexus\Synapse\Exception\`) @@ -106,6 +108,7 @@ | `NotFoundException` | 404 | Статус чужого/несуществующего события | | `ValidationException` | 422 / null | Незарегистрированный тип или битый конверт (null = локально) | | `ServerException` | 5xx | Synapse упал | +| `StatusTimeoutException` | — | `waitForStatus` не дождался целевого статуса за timeout | | `InvalidWebhookException` | — | приём s2s: подпись/freshness/не-JSON — не прошло проверку `WebhookVerifier` | Ловить — одну `SynapseException`. diff --git a/examples/plain-php/smoke.php b/examples/plain-php/smoke.php index 781c737..d34d7d3 100644 --- a/examples/plain-php/smoke.php +++ b/examples/plain-php/smoke.php @@ -22,6 +22,8 @@ // уникальный на прогон: дедуп-окно сервера 24 ч $dedup = 'client-smoke-' . time(); +// источник теста: из env (SYNAPSE_DEFAULT_SOURCE, см. заголовок) — как и в реальном сервисе +$source = getenv('SYNAPSE_DEFAULT_SOURCE') ?: 'libtest'; $config = SynapseConfig::fromGlobals(); @@ -35,11 +37,11 @@ printf(" health: %s\n", json_encode($client->health(), JSON_UNESCAPED_UNICODE)); -$event = $client->send('libtest', 'ping', 'done', payload: ['user_id' => 'mcp'], dedupKey: $dedup); +$event = $client->send($source, 'ping', 'done', userId: 'mcp', dedupKey: $dedup); printf(" send: id=%s status=%s\n", $event->id, $event->status); assert($event->status === 'queued' && !$event->deduplicated); -$repeat = $client->send('libtest', 'ping', 'done', payload: ['user_id' => 'mcp'], dedupKey: $dedup); +$repeat = $client->send($source, 'ping', 'done', userId: 'mcp', dedupKey: $dedup); printf(" repeat: id=%s deduplicated=%s\n", $repeat->id, $repeat->deduplicated ? 'true' : 'false'); assert($repeat->deduplicated && $repeat->id === $event->id); // дедуп вернул первое событие @@ -78,6 +80,6 @@ new HttpFactory(), new HttpFactory(), ); -printf(" emit к неподнятому хосту вернул: %s (null — ок)\n", json_encode($emitClient->emit('libtest', 'a', 'b'))); +printf(" emit к неподнятому хосту вернул: %s (null — ок)\n", json_encode($emitClient->emit($source, 'a', 'b'))); echo " ALL OK\n"; \ No newline at end of file diff --git a/src/Exception/StatusTimeoutException.php b/src/Exception/StatusTimeoutException.php new file mode 100644 index 0000000..bed60cc --- /dev/null +++ b/src/Exception/StatusTimeoutException.php @@ -0,0 +1,10 @@ +buildEnvelope( - $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, + $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, $userId, ); if (isset($envelope['payload']) && is_array($envelope['payload'])) { EnvelopeBuilder::warnBadUserId($envelope['payload'], $this->logger); @@ -97,11 +99,12 @@ ?string $dedupKey = null, ?int $ttlSeconds = null, ?\DateTimeInterface $scheduledAt = null, + ?string $userId = null, bool $throw = false, ): ?SentEvent { try { return $this->send( - $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, + $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, $userId, ); } catch (SynapseException $ex) { if ($throw) { @@ -220,6 +223,47 @@ return $this->diagnostics('/api/readyz'); } + /** + * Поллинг до целевого статуса или таймаут — для приёмки/тестов, не + * бизнес-кода. Морг транспорта не прерывает ожидание до дедлайна. + * Не достигли → StatusTimeoutException с последним известным статусом. + * + * @param list $statuses + */ + public function waitForStatus( + string $eventId, + array $statuses = ['done', 'failed'], + float $timeoutSeconds = 30.0, + float $intervalSeconds = 1.0, + ): EventStatus { + $deadline = microtime(true) + $timeoutSeconds; + $lastStatus = null; + $sleepMicro = (int) round(max($intervalSeconds, 0.01) * 1_000_000); + while (true) { + try { + $st = $this->status($eventId); + $lastStatus = $st->status; + if (in_array($st->status, $statuses, true)) { + return $st; + } + } catch (TransportException) { + // Synapse моргнул — держим дедлайн + } + if (microtime(true) >= $deadline) { + throw new StatusTimeoutException( + sprintf( + 'событие %s не в [%s] за %s c (последний статус: %s)', + $eventId, + implode(', ', $statuses), + $timeoutSeconds, + $lastStatus ?? 'нет ответа', + ), + ); + } + usleep($sleepMicro); + } + } + // -- internal -------------------------------------------------------- /** @@ -248,7 +292,13 @@ ?string $dedupKey, ?int $ttlSeconds, ?\DateTimeInterface $scheduledAt, + ?string $userId = null, ): array { + // user_id — first-class конвенция «о ком событие» (docs/05): + // явный параметр перебивает payload["user_id"], если тот был + if ($userId !== null && $userId !== '') { + $payload = [...($payload ?? []), 'user_id' => $userId]; + } return EnvelopeBuilder::build( $this->resolveSource($source), $subject, diff --git a/src/Webhook/WebhookVerifier.php b/src/Webhook/WebhookVerifier.php index 50beee6..5eb9a8c 100644 --- a/src/Webhook/WebhookVerifier.php +++ b/src/Webhook/WebhookVerifier.php @@ -31,6 +31,8 @@ final class WebhookVerifier { public const SIGNATURE_HEADER = 'x-gnexus-signature'; + public const EVENT_TYPE_HEADER = 'x-gnexus-event-type'; // ".." + public const SOURCE_HEADER = 'x-synapse-source'; public const DEFAULT_MAX_SKEW = 300; // секунд (зеркало gnexus-auth) /** private __construct: только статика. */ @@ -48,6 +50,36 @@ } /** + * `".."` из `X-Gnexus-Event-Type` → + * `list{string, string, string}`; null — нет/битый заголовок. + * + * @return list{string, string, string}|null + */ + public static function parseEventType(?string $header): ?array + { + if ($header === null || $header === '') { + return null; + } + $parts = explode('.', $header); + if (count($parts) !== 3 || in_array('', $parts, true)) { + return null; + } + + return $parts; + } + + /** Тело ответа 2xx приёмника: ['received' => true, 'event_id' => ...]. */ + public static function ack(?string $eventId = null): array + { + $body = ['received' => true]; + if ($eventId !== null && $eventId !== '') { + $body['event_id'] = $eventId; + } + + return $body; + } + + /** * Проверить подпись и распарсить доставку; возвращает конверт (assoc-array). * * @param array> $headers произвольный массив заголовков diff --git a/tests/SynapseClientTest.php b/tests/SynapseClientTest.php index e8aba38..e2df4b2 100644 --- a/tests/SynapseClientTest.php +++ b/tests/SynapseClientTest.php @@ -10,6 +10,7 @@ use GNexus\Synapse\Exception\ForbiddenException; use GNexus\Synapse\Exception\NotFoundException; use GNexus\Synapse\Exception\ServerException; +use GNexus\Synapse\Exception\StatusTimeoutException; use GNexus\Synapse\Exception\SynapseException; use GNexus\Synapse\Exception\TransportException; use GNexus\Synapse\Exception\ValidationException; @@ -246,6 +247,68 @@ $noKey->send('s', 'x', 'y'); } + // -- userId + waitForStatus (v0.1.2) -------------------------------- + + public function testUserIdMergesIntoPayloadAndOverrides(): void + { + $body = null; + $client = $this->clientWith(static function ($r) use (&$body) { + $body = json_decode((string) $r->getBody(), true, flags: JSON_THROW_ON_ERROR); + return new Response(202, [], json_encode([ + 'ok' => true, 'id' => 'e-1', 'status' => 'queued', 'deduplicated' => false, + ])); + }); + $client->send('bugtrail', 'task', 'created', 'normal', ['title' => 'x', 'user_id' => 'oops'], userId: 'u-42'); + self::assertSame(['title' => 'x', 'user_id' => 'u-42'], $body['payload']); + } + + public function testUserIdNullLeavesPayloadOut(): void + { + $body = null; + $client = $this->clientWith(static function ($r) use (&$body) { + $body = json_decode((string) $r->getBody(), true, flags: JSON_THROW_ON_ERROR); + return new Response(202, [], json_encode([ + 'ok' => true, 'id' => 'e-1', 'status' => 'queued', 'deduplicated' => false, + ])); + }); + $client->emit('bugtrail', 'task', 'created'); + self::assertArrayNotHasKey('payload', $body); + } + + public function testWaitForStatusPollsUntilTerminal(): void + { + $calls = 0; + $client = $this->clientWith(static function () use (&$calls) { + ++$calls; + $status = $calls < 3 ? 'queued' : 'done'; + return new Response(200, [], json_encode([ + 'id' => '11111111-1111-1111-1111-111111111111', 'source' => 'libtest', + 'subject' => 'ping', 'action' => 'done', 'priority' => 'normal', + 'status' => $status, 'created_at' => '2026-10-03T12:00:00+00:00', + 'deliveries' => [], + ])); + }); + $st = $client->waitForStatus('11111111-1111-1111-1111-111111111111', ['done'], 5, 0.01); + self::assertSame('done', $st->status); + self::assertSame(3, $calls); + } + + public function testWaitForStatusTimeout(): void + { + $client = $this->clientWith(static fn () => new Response(200, [], json_encode([ + 'id' => '11111111-1111-1111-1111-111111111111', 'source' => 'libtest', + 'subject' => 'ping', 'action' => 'done', 'priority' => 'normal', + 'status' => 'queued', 'created_at' => '2026-10-03T12:00:00+00:00', + 'deliveries' => [], + ]))); + try { + $client->waitForStatus('11111111-1111-1111-1111-111111111111', ['done'], 0.08, 0.01); + self::fail('таймаут не сработал'); + } catch (StatusTimeoutException $ex) { + self::assertStringContainsString('queued', $ex->detail); + } + } + // -- helpers -------------------------------------------------------- private static function config(): SynapseConfig diff --git a/tests/WebhookVerifierTest.php b/tests/WebhookVerifierTest.php index e7b401f..93c9e2b 100644 --- a/tests/WebhookVerifierTest.php +++ b/tests/WebhookVerifierTest.php @@ -134,4 +134,26 @@ self::assertStringContainsString('не JSON-объект', $ex->detail); } } + + // -- хелперы контекста и ответа ------------------------------------ + + public function testParseEventType(): void + { + self::assertSame( + ['monitoring', 'container', 'down'], + WebhookVerifier::parseEventType('monitoring.container.down'), + ); + self::assertNull(WebhookVerifier::parseEventType('a.b')); + self::assertNull(WebhookVerifier::parseEventType('..')); + self::assertNull(WebhookVerifier::parseEventType(null)); + } + + public function testAck(): void + { + self::assertSame(['received' => true], WebhookVerifier::ack()); + self::assertSame( + ['received' => true, 'event_id' => self::EVENT_ID], + WebhookVerifier::ack(self::EVENT_ID), + ); + } } \ No newline at end of file