diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..d548adf --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +/vendor/ +.phpunit.cache/ +.phpunit.result.cache +/composer.lock \ No newline at end of file diff --git a/README.md b/README.md new file mode 100644 index 0000000..6bedcba --- /dev/null +++ b/README.md @@ -0,0 +1,129 @@ +# gnexus-synapse-client-php + +Тонкий PHP-клиент Synapse (централизованный хаб уведомлений Gnexus) для сервисов экосистемы: +gnexus-auth (Laravel), Navi-инстансы и другие PHP-сервисы. + +Клиент берёт на себя весь повторяемый шаблон: env-конфиг, сборку и локальную валидацию конверта, +типизированные исключения по HTTP-статусам, fire-and-forget `emit`, batch, статус события. +Ретраев и локальных очередей нет — Synapse принимает событие с ответом 202 мгновенно, всё +остальное делает его воркер. + +Транспорт инжектный (PSR-18 + PSR-17, например Guzzle — в `require-dev` для тестов и примеров, +жёсткой зависимости нет). + +## Установка + +В `composer.json` сервиса: + +```json +{ + "repositories": [ + { + "type": "vcs", + "url": "https://git.gnexus.space/git/root/gn-synapse-client-php.git" + } + ] +} +``` + +```bash +composer require gnexus/synapse-client:^0.1.0 +``` + +## Quickstart + +```php +use GNexus\Synapse\Config\SynapseConfig; +use GNexus\Synapse\SynapseClient; +use GuzzleHttp\Client; +use GuzzleHttp\Psr7\HttpFactory; + +$client = new SynapseClient( + SynapseConfig::fromGlobals(), // SYNAPSE_URL, SYNAPSE_API_KEY, ... + new Client(['timeout' => 10.0]), // PSR-18 (таймаут задаёт потребитель) + new HttpFactory(), // PSR-17 RequestFactory + new HttpFactory(), // PSR-17 StreamFactory +); + +// Обычная отправка: событие обязательно дойти — ошибка бросает исключение +$event = $client->send('bugtrail', 'test', 'failed', 'high', + ['user_id' => 'u-42'], dedupKey: "failed-42"); +$event->id; // uuid +$event->status; // 'queued' — 202, воркер заберёт позже + +// Fire-and-forget: упавший Synapse не должен ломать основной путь сервиса +$client->emit(null, 'user', 'password_changed', payload: ['user_id' => $sub]); +// → SentEvent|null; SynapseException и транспортные сбои ловятся, +// пишутся в $logger->warning('synapse emit failed: ...', ...) и глотаются. + +// Статус и доставки +$status = $client->status($event->id); +foreach ($status->deliveries as $d) { /* channel, target, status, attempts, error */ } +``` + +## Конфигурация (env) + +| Переменная | Обязательна | По умолчанию | Что делает | +|---|---|---|---| +| `SYNAPSE_URL` | да | — | База Synapse (`http://localhost:8013`) | +| `SYNAPSE_API_KEY` | да (для send/emit) | — | Ключ источника (`syn_...`) | +| `SYNAPSE_TIMEOUT` | нет | `10.0` | Секунды (используется как `fromGlobals()->timeoutSeconds`) | +| `SYNAPSE_DEFAULT_SOURCE` | нет | — | `source` по умолчанию: `$client->send(null, ...)` | + +`SynapseConfig` валидирует всё на конструировании (fail fast): пустой URL, схема не http(s), +таймаут ≤ 0, ключ-пустая-строка → `ConfigurationException`. Send/sendBatch при отсутствии +ключа тоже бросает `ConfigurationException` **до** HTTP. + +`default_source` — смягчение антивспуфинга сервера: 403 получают вызовы, где `source` конверта +не совпадает с источником ключа; с `SYNAPSE_DEFAULT_SOURCE` это ошибка конфигурации, а не кода. + +## Методы + +- `send(...): SentEvent` — конверт валидируется локально (без HTTP): имена `^[a-z0-9]([a-z0-9._-]*[a-z0-9])?$` + 1..64, priority ∈ `low|normal|high|critical`, `ttlSeconds` 1..604800, `dedupKey` ≤255 (`''` → null), + payload — assoc-массив. `normal` priority в конверт не попадает, null-поля не шлются + (`scheduledAt` → UTC `DATE_ATOM`). Ошибки валидации — локальный `ValidationException` + (statusCode = null: это ошибка кода, не Synapse). +- `emit(..., bool $throw = false): ?SentEvent` — то же, но любые `SynapseException`/транспортные + ошибки ловятся → `$logger->warning(...)` → null (`throw: true` → перебрасывает). +- `sendBatch(array $envelopes): BatchResult` — каждый элемент прогоняется через `EnvelopeBuilder::build()`; + локально битые — `BatchRejection(index, detail)` без HTTP; ответ сервера 202 **и** 422 («ни одного + не принято») парсится одинаково — `index` серверного отклонения ремапится на исходную позицию; 5xx/транспорт — исключение целиком. +- `status(string $eventId): EventStatus` — uuid-проверка, 404 → `NotFoundException`. +- `health(): array`, `ready(): array` — `/api/healthz` и `/api/readyz`, без ключа. + +## Исключения (`GNexus\Synapse\Exception\`) + +| Класс | Код | Когда | +|---|---|---| +| `SynapseException` | — | База: `.detail`, `.statusCode` | +| `ConfigurationException` | — | Битый конфиг / нет ключа (до HTTP) | +| `TransportException` | — | Сеть упала (`.transportError`) | +| `AuthException` | 401 | Ключ нет/отозван/источник архивен | +| `ForbiddenException` | 403 | `source` ≠ источник ключа (антивспуфинг) | +| `NotFoundException` | 404 | Статус чужого/несуществующего события | +| `ValidationException` | 422 / null | Незарегистрированный тип или битый конверт (null = локально) | +| `ServerException` | 5xx | Synapse упал | + +Ловить — одну `SynapseException`. + +## Оговорки + +- **Дедуп best-effort** в окне 24ч: повтор возвращает первое событие с `deduplicated: true`. + Не строить бизнес-логику на отсутствии дублей. +- **Инжектный PSR-18 без таймаута висит вечно** (какой бы ни был `timeoutSeconds`) — в Guzzle и + аналогах явно задавайте таймаут (Laravel-binding с таймаутом — `examples/laravel/README.md`). +- `scheduled_at` — резерв MVP: принимается, не отправляется. +- `payload.user_id` — непустая строка (uuid gnexus-auth); остальное — предупреждение в лог (сервер + запишет address-доставки со статусом `skipped`). + +## Разработка + +```bash +composer install +vendor/bin/phpunit +``` + +Интеграционный смок против живого стека — `examples/plain-php/smoke.php`. + +Лицензия: проприетарная (внутри экосистемы Gnexus). \ No newline at end of file diff --git a/composer.json b/composer.json new file mode 100644 index 0000000..59cdec0 --- /dev/null +++ b/composer.json @@ -0,0 +1,32 @@ +{ + "name": "gnexus/synapse-client", + "description": "Framework-agnostic PHP client library for gn-synapse event ingestion", + "keywords": ["gnexus", "synapse", "notifications", "ingestion", "events"], + "type": "library", + "license": "proprietary", + "require": { + "php": "^8.3", + "ext-json": "*", + "psr/http-client": "^1.0", + "psr/http-factory": "^1.0", + "psr/http-message": "^2.0", + "psr/log": "^3.0" + }, + "require-dev": { + "phpunit/phpunit": "^11.5", + "guzzlehttp/guzzle": "^7.9", + "guzzlehttp/psr7": "^2.7" + }, + "autoload": { + "psr-4": { + "GNexus\\Synapse\\": "src/" + } + }, + "autoload-dev": { + "psr-4": { + "GNexus\\Synapse\\Tests\\": "tests/" + } + }, + "minimum-stability": "stable", + "prefer-stable": true +} \ No newline at end of file diff --git a/examples/laravel/README.md b/examples/laravel/README.md new file mode 100644 index 0000000..78dd329 --- /dev/null +++ b/examples/laravel/README.md @@ -0,0 +1,47 @@ +# Подключение Synapse-клиента в Laravel-сервис + +PSR-18 транспорт — тот, что уже есть в сервисе (в Laravel обычно Guzzle +через `guzzlehttp/guzzle`). Таймаут берётся из конфига клиента. + +`app/Providers/AppServiceProvider.php`: + +```php +use GNexus\Synapse\Config\SynapseConfig; +use GNexus\Synapse\SynapseClient; +use GuzzleHttp\Client; +use GuzzleHttp\Psr7\HttpFactory; + +public function register(): void +{ + $this->app->singleton(SynapseConfig::class, fn () => SynapseConfig::fromGlobals()); + + $this->app->singleton(SynapseClient::class, function ($app) { + $config = $app->make(SynapseConfig::class); + return new SynapseClient( + $config, + new Client(['timeout' => $config->timeoutSeconds]), + new HttpFactory(), // RequestFactoryInterface + new HttpFactory(), // StreamFactoryInterface + ); + }); +} +``` + +`.env` (значения — в vault/env сервиса, ключ печатается при выдаче один раз): + +``` +SYNAPSE_URL=http://synapse.gnexus.space:8013 +SYNAPSE_API_KEY=syn_... +SYNAPSE_DEFAULT_SOURCE=gnexus-auth +``` + +Использование (в job/event listener): + +```php +$synapse = app(SynapseClient::class); +$synapse->emit(source: null, subject: 'user', action: 'password_changed', + payload: ['user_id' => $user->uuid], dedupKey: "pw-{$user->uuid}"); +``` + +`emit` не бросает исключения: упавший Synapse не должен ломать выход +пользователя. Для критичных событий — `send` с `try/catch SynapseException`. \ No newline at end of file diff --git a/examples/plain-php/smoke.php b/examples/plain-php/smoke.php new file mode 100644 index 0000000..781c737 --- /dev/null +++ b/examples/plain-php/smoke.php @@ -0,0 +1,83 @@ + $config->timeoutSeconds]), + new HttpFactory(), + new HttpFactory(), +); + +printf(" health: %s\n", json_encode($client->health(), JSON_UNESCAPED_UNICODE)); + +$event = $client->send('libtest', 'ping', 'done', payload: ['user_id' => '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); +printf(" repeat: id=%s deduplicated=%s\n", $repeat->id, $repeat->deduplicated ? 'true' : 'false'); +assert($repeat->deduplicated && $repeat->id === $event->id); // дедуп вернул первое событие + +$status = $client->status($event->id); +printf(" status: %s; доставки:\n", $status->status); +foreach ($status->deliveries as $d) { + printf( + " %s%s: %s ×%d%s\n", + $d->channel, + $d->target !== null ? ' → ' . $d->target : '', + $d->status, + $d->attempts, + $d->error !== null ? ' (' . $d->error . ')' : '', + ); +} + +try { + $client->send('wrong-source', 'ping', 'done'); // анти-спуфинг сервера + fwrite(STDERR, " ОШИБКА: 403 не получен\n"); + exit(1); +} catch (SynapseException $ex) { + printf(" %d пойман (анти-спуфинг): %s\n", $ex->statusCode, $ex->detail); +} + +try { + $client->status('00000000-0000-0000-0000-000000000000'); +} catch (NotFoundException) { + echo " 404 пойман (чужой/несуществующий id)\n"; +} + +// fire-and-forget: даже на лежащем Synapse не бросает +$downConfig = new SynapseConfig('http://127.0.0.1:9', $config->apiKey, 1.0); +$emitClient = new SynapseClient( + $downConfig, + new Client(['timeout' => 1.0]), + new HttpFactory(), + new HttpFactory(), +); +printf(" emit к неподнятому хосту вернул: %s (null — ок)\n", json_encode($emitClient->emit('libtest', 'a', 'b'))); + +echo " ALL OK\n"; \ No newline at end of file diff --git a/phpunit.xml.dist b/phpunit.xml.dist new file mode 100644 index 0000000..f2878c1 --- /dev/null +++ b/phpunit.xml.dist @@ -0,0 +1,14 @@ + + + + + tests + + + \ No newline at end of file diff --git a/src/Config/SynapseConfig.php b/src/Config/SynapseConfig.php new file mode 100644 index 0000000..1eddd6b --- /dev/null +++ b/src/Config/SynapseConfig.php @@ -0,0 +1,66 @@ + 0"); + } + if ($apiKey !== null && $apiKey === '') { + throw new ConfigurationException('api_key пустая строка — укажите ключ или null'); + } + } + + /** Фолбэк в env: fromGlobals() или из синтетического массива (тесты). */ + public static function fromGlobals(array $env = []): self + { + $env = $env === [] ? $_SERVER + $_ENV : $env; + $timeout = $env[self::ENV_TIMEOUT] ?? ''; + return new self( + baseUrl: $env[self::ENV_URL] ?? '', + apiKey: isset($env[self::ENV_API_KEY]) && $env[self::ENV_API_KEY] !== '' + ? $env[self::ENV_API_KEY] + : null, + timeoutSeconds: $timeout === '' ? 10.0 : (float) $timeout, + defaultSource: $env[self::ENV_DEFAULT_SOURCE] ?? '', + ); + } + + public function url(string $path): string + { + // путь — с начальным '/' + return rtrim($this->baseUrl, '/') . $path; + } + + public function source(): string + { + return $this->defaultSource; + } +} \ No newline at end of file diff --git a/src/DTO/BatchRejection.php b/src/DTO/BatchRejection.php new file mode 100644 index 0000000..0320f3c --- /dev/null +++ b/src/DTO/BatchRejection.php @@ -0,0 +1,17 @@ + $results + */ + public function __construct( + public array $results, + ) {} + + /** @return list */ + public function accepted(): array + { + return array_values(array_filter( + $this->results, + static fn ($r) => $r instanceof SentEvent, + )); + } + + /** @return list */ + public function rejected(): array + { + return array_values(array_filter( + $this->results, + static fn ($r) => $r instanceof BatchRejection, + )); + } + + public function allAccepted(): bool + { + return $this->rejected() === []; + } +} \ No newline at end of file diff --git a/src/DTO/Delivery.php b/src/DTO/Delivery.php new file mode 100644 index 0000000..3aa71e6 --- /dev/null +++ b/src/DTO/Delivery.php @@ -0,0 +1,18 @@ + $deliveries + */ + public function __construct( + public string $id, + public string $source, + public string $subject, + public string $action, + public string $priority, + public string $status, + public \DateTimeImmutable $createdAt, + public ?\DateTimeImmutable $expiresAt, + public array $deliveries = [], + ) {} +} \ No newline at end of file diff --git a/src/DTO/SentEvent.php b/src/DTO/SentEvent.php new file mode 100644 index 0000000..c541bc0 --- /dev/null +++ b/src/DTO/SentEvent.php @@ -0,0 +1,15 @@ += 1 + public const DEDUP_MAX = 255; + + /** + * Собрать конверт v1; лишнего в нём нет (сервер extra="forbid"). + * + * @param array|null $payload JSON-объект + * @return array + */ + public static function build( + string $source, + string $subject, + string $action, + string $priority = 'normal', + ?array $payload = null, + ?string $dedupKey = null, + ?int $ttlSeconds = null, + ?\DateTimeInterface $scheduledAt = null, + ): array { + $errors = []; + foreach (['source' => $source, 'subject' => $subject, 'action' => $action] as $kind => $name) { + $err = self::nameError($kind, $name); + if ($err !== null) { + $errors[] = $err; + } + } + if (!in_array($priority, self::PRIORITY_VALUES, true)) { + $errors[] = sprintf( + "недопустимый priority='%s': разрешены %s", + $priority, + implode(', ', self::PRIORITY_VALUES), + ); + } + if ($dedupKey === '') { + $dedupKey = null; // пустая строка = не задаём + } elseif ($dedupKey !== null && strlen($dedupKey) > self::DEDUP_MAX) { + $errors[] = sprintf('dedup_key длиннее %d символов', self::DEDUP_MAX); + } + if ($ttlSeconds !== null && ($ttlSeconds < 1 || $ttlSeconds > self::TTL_MAX)) { + $errors[] = sprintf('ttlSeconds=%d: допустимо 1..%d', $ttlSeconds, self::TTL_MAX); + } + if ($errors !== []) { + throw new ValidationException(implode('; ', $errors)); + } + + $envelope = [ + 'source' => $source, + 'subject' => $subject, + 'action' => $action, + ]; + if ($priority !== 'normal') { + $envelope['priority'] = $priority; // серверный дефолт не шлём + } + if ($payload !== null) { + if ($payload === [] || array_is_list($payload)) { + throw new ValidationException('payload должен быть словарем (JSON-объектом)'); + } + $envelope['payload'] = $payload; + } + if ($dedupKey !== null) { + $envelope['dedup_key'] = $dedupKey; + } + if ($ttlSeconds !== null) { + $envelope['ttl_seconds'] = $ttlSeconds; + } + if ($scheduledAt !== null) { + $envelope['scheduled_at'] = $scheduledAt + ->setTimezone(new \DateTimeZone('UTC')) + ->format(DATE_ATOM); + } + return $envelope; + } + + /** + * Конвенция payload.user_id (= sub gnexus-auth) — «о ком событие». + * Сервер не валидирует payload: не-строчный/пустой user_id запишет адресные + * доставки как skipped. Warning, не ошибка. + * + * @param array $payload + */ + public static function warnBadUserId(array $payload, LoggerInterface $logger): void + { + $uid = $payload['user_id'] ?? null; + if ((is_string($uid) && $uid !== '') || (is_int($uid) && !is_bool($uid))) { + return; + } + $logger->warning( + 'payload.user_id должен быть непустой строкой (или числом); получено {uid} — адресные доставки будут skipped', + ['uid' => $uid], + ); + } + + private static function nameError(string $kind, string $value): ?string + { + if ($value === '') { + return sprintf("недопустимое %s='': строка 1..%d символов", $kind, self::NAME_MAX); + } + if (strlen($value) > self::NAME_MAX) { + return sprintf("недопустимое %s='%s': длина > %d", $kind, $value, self::NAME_MAX); + } + if (preg_match(self::NAME_PATTERN, $value) !== 1) { + return sprintf( + "недопустимое %s='%s': не совпадает с шаблоном %s", + $kind, + $value, + self::NAME_PATTERN, + ); + } + return null; + } +} \ No newline at end of file diff --git a/src/Exception/AuthException.php b/src/Exception/AuthException.php new file mode 100644 index 0000000..e4c1fce --- /dev/null +++ b/src/Exception/AuthException.php @@ -0,0 +1,8 @@ + AuthException::class, + 403 => ForbiddenException::class, + 404 => NotFoundException::class, + 422 => ValidationException::class, + ]; + + private LoggerInterface $logger; + + public function __construct( + private readonly SynapseConfig $config, + private readonly ClientInterface $httpClient, + private readonly RequestFactoryInterface $requestFactory, + private readonly StreamFactoryInterface $streamFactory, + ?LoggerInterface $logger = null, + ) { + $this->logger = $logger ?? new NullLogger(); + } + + /** + * POST /api/v1/events → 202. Ошибки — типизированные исключения. + * + * @param array $payload + */ + public function send( + ?string $source = null, + string $subject = '', + string $action = '', + string $priority = 'normal', + array $payload = [], + ?string $dedupKey = null, + ?int $ttlSeconds = null, + ?\DateTimeInterface $scheduledAt = null, + ): SentEvent { + $envelope = $this->buildEnvelope( + $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, + ); + if (isset($envelope['payload']) && is_array($envelope['payload'])) { + EnvelopeBuilder::warnBadUserId($envelope['payload'], $this->logger); + } + return $this->parseSentEvent($this->post($envelope, self::API_PREFIX . '/events')); + } + + /** + * Fire-and-forget: ловит SynapseException, пишет warning, возвращает null. + * + * @param array $payload + */ + public function emit( + ?string $source = null, + string $subject = '', + string $action = '', + string $priority = 'normal', + array $payload = [], + ?string $dedupKey = null, + ?int $ttlSeconds = null, + ?\DateTimeInterface $scheduledAt = null, + bool $throw = false, + ): ?SentEvent { + try { + return $this->send( + $source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, + ); + } catch (SynapseException $ex) { + if ($throw) { + throw $ex; + } + $this->logger->warning('synapse emit failed: ' . $ex->detail); + return null; + } + } + + /** + * Голый массив конвертов; сервер отвечает per-index (202 и 422 парсим + * одинаково: rejected — не исключение, частичный успех норма). Локально + * битые конверты не посылаются — сразу BatchRejection на своей позиции. + * Транспортный сбой/5xx — исключение: частичный успех неизвестен. + * + * @param list> $events + */ + public function sendBatch(array $events): BatchResult + { + $results = []; + $positions = []; + $payload = []; + foreach (array_values($events) as $i => $raw) { + try { + $rawSource = $raw['source'] ?? null; + $rawDedup = $raw['dedup_key'] ?? null; + $rawTtl = $raw['ttl_seconds'] ?? null; + $rawScheduled = $raw['scheduled_at'] ?? null; + $payload[] = EnvelopeBuilder::build( + is_string($rawSource) && $rawSource !== '' ? $rawSource : $this->resolveSource(), + (string) ($raw['subject'] ?? ''), + (string) ($raw['action'] ?? ''), + is_string($raw['priority'] ?? null) ? $raw['priority'] : 'normal', + is_array($raw['payload'] ?? null) ? $raw['payload'] : null, + is_string($rawDedup) ? $rawDedup : null, + is_int($rawTtl) ? $rawTtl : null, + $rawScheduled instanceof \DateTimeInterface ? $rawScheduled : null, + ); + $positions[] = $i; + } catch (SynapseException $ex) { // Configuration + Validation локально + $results[$i] = new BatchRejection($i, $ex->detail); + } + } + + if ($payload !== []) { + $resp = $this->post($payload, self::API_PREFIX . '/events/batch'); + $body = $this->jsonBody($resp); + $serverResults = is_array($body['results'] ?? null) ? $body['results'] : []; + foreach ($positions as $pos => $origIndex) { + $item = $serverResults[$pos] ?? []; + if (!is_array($item)) { + $results[$origIndex] = new BatchRejection($origIndex, 'битый элемент ответа'); + } elseif (($item['ok'] ?? null) === true) { + $results[$origIndex] = new SentEvent( + (string) ($item['id'] ?? ''), + (string) ($item['status'] ?? ''), + (bool) ($item['deduplicated'] ?? false), + ); + } else { + $results[$origIndex] = new BatchRejection( + $origIndex, + (string) ($item['detail'] ?? ''), + ); + } + } + } + ksort($results); + return new BatchResult(array_values($results)); + } + + public function status(string $eventId): EventStatus + { + $resp = $this->request('GET', self::API_PREFIX . '/events/' . rawurlencode($eventId)); + $body = $this->jsonBody($resp); + if ($resp->getStatusCode() !== 200) { + throw $this->errorFor($resp, $body); + } + $deliveries = []; + foreach (is_array($body['deliveries'] ?? null) ? $body['deliveries'] : [] as $d) { + $deliveries[] = new Delivery( + (string) ($d['channel'] ?? ''), + isset($d['target']) && $d['target'] !== null ? (string) $d['target'] : null, + (string) ($d['status'] ?? ''), + (int) ($d['attempts'] ?? 0), + isset($d['error']) && $d['error'] !== null ? (string) $d['error'] : null, + isset($d['rendered_message']) && $d['rendered_message'] !== null + ? (string) $d['rendered_message'] + : null, + ); + } + return new EventStatus( + (string) ($body['id'] ?? $eventId), + (string) ($body['source'] ?? ''), + (string) ($body['subject'] ?? ''), + (string) ($body['action'] ?? ''), + (string) ($body['priority'] ?? ''), + (string) ($body['status'] ?? ''), + new \DateTimeImmutable((string) ($body['created_at'] ?? 'now')), + isset($body['expires_at']) && $body['expires_at'] !== null + ? new \DateTimeImmutable((string) $body['expires_at']) + : null, + $deliveries, + ); + } + + /** GET /api/healthz (без ключа) — смок/мониторинг. */ + public function health(): array + { + return $this->diagnostics('/api/healthz'); + } + + /** GET /api/readyz (без ключа). */ + public function ready(): array + { + return $this->diagnostics('/api/readyz'); + } + + // -- internal -------------------------------------------------------- + + /** + * @return array + */ + private function diagnostics(string $path): array + { + $resp = $this->request('GET', $path); + $body = $this->jsonBody($resp); + if ($resp->getStatusCode() !== 200) { + throw $this->errorFor($resp, $body); + } + return $body; + } + + /** + * @param array|null $payload + * @return array + */ + private function buildEnvelope( + ?string $source, + string $subject, + string $action, + string $priority, + ?array $payload, + ?string $dedupKey, + ?int $ttlSeconds, + ?\DateTimeInterface $scheduledAt, + ): array { + return EnvelopeBuilder::build( + $this->resolveSource($source), + $subject, + $action, + $priority, + $payload !== null && $payload !== [] ? $payload : null, + $dedupKey, + $ttlSeconds, + $scheduledAt, + ); + } + + private function resolveSource(?string $source = null): string + { + $source = $source ?? $this->config->source(); + if ($source === '') { + throw new ConfigurationException( + 'не задан source: передайте source или SYNAPSE_DEFAULT_SOURCE (' . + SynapseConfig::ENV_DEFAULT_SOURCE . ')', + ); + } + return $source; + } + + /** + * @param array|list> $body + */ + private function post(array $body, string $path): ResponseInterface + { + if ($this->config->apiKey === null) { + throw new ConfigurationException( + 'не задан api_key: передайте ключ или ' . SynapseConfig::ENV_API_KEY, + ); + } + return $this->request('POST', $path, json_encode($body, JSON_THROW_ON_ERROR)); + } + + private function request(string $method, string $path, ?string $json = null): ResponseInterface + { + $request = $this->requestFactory->createRequest($method, $this->config->url($path)) + ->withHeader('Content-Type', 'application/json') + ->withHeader('User-Agent', self::USER_AGENT); + if ($this->config->apiKey !== null) { + $request = $request->withHeader('Authorization', 'Bearer ' . $this->config->apiKey); + } + if ($json !== null) { + $request = $request->withBody($this->streamFactory->createStream($json)); + } + try { + return $this->httpClient->sendRequest($request); + } catch (NetworkExceptionInterface $ex) { + throw new TransportException('Synapse недоступен: ' . $ex->getMessage(), $ex); + } catch (ClientExceptionInterface $ex) { + throw new TransportException('Транспортная ошибка: ' . $ex->getMessage(), $ex); + } + } + + /** + * @return array + */ + private function jsonBody(ResponseInterface $resp): array + { + $raw = (string) $resp->getBody(); + if ($raw === '') { + return []; + } + try { + $decoded = json_decode($raw, true, 512, JSON_THROW_ON_ERROR); + } catch (\JsonException) { + return []; + } + return is_array($decoded) ? $decoded : []; + } + + private function errorFor(ResponseInterface $resp, array $body): SynapseException + { + $code = $resp->getStatusCode(); + if (isset($body['detail']) && is_scalar($body['detail'])) { + $detail = (string) $body['detail']; + } else { + $detail = sprintf('HTTP %d: %s', $code, substr((string) $resp->getBody(), 0, 300)); + } + if ($code >= 500) { + return new ServerException($detail, $code); + } + $cls = self::STATUS_ERRORS[$code] ?? SynapseException::class; + return new $cls($detail, $code); + } + + private function parseSentEvent(ResponseInterface $resp): SentEvent + { + $body = $this->jsonBody($resp); + if ($resp->getStatusCode() !== 202) { + throw $this->errorFor($resp, $body); + } + return new SentEvent( + (string) ($body['id'] ?? ''), + (string) ($body['status'] ?? ''), + (bool) ($body['deduplicated'] ?? false), + ); + } +} \ No newline at end of file diff --git a/tests/EnvelopeBuilderTest.php b/tests/EnvelopeBuilderTest.php new file mode 100644 index 0000000..1626d6d --- /dev/null +++ b/tests/EnvelopeBuilderTest.php @@ -0,0 +1,109 @@ + 's', 'subject' => 'x', 'action' => 'y'], + EnvelopeBuilder::build('s', 'x', 'y'), + ); + // сервер extra="forbid": None/пустых полей в конверте нет + $env = EnvelopeBuilder::build('s', 'x', 'y', payload: null, dedupKey: '', ttlSeconds: null); + self::assertSame(['source' => 's', 'subject' => 'x', 'action' => 'y'], $env); + } + + public function testBuildsFullEnvelope(): void + { + $ts = new \DateTimeImmutable('2026-10-03 12:00:00', new \DateTimeZone('UTC')); + $env = EnvelopeBuilder::build('s', 'x', 'y', 'high', ['user_id' => '42'], 'd1', 60, $ts); + self::assertSame([ + 'source' => 's', 'subject' => 'x', 'action' => 'y', 'priority' => 'high', + 'payload' => ['user_id' => '42'], 'dedup_key' => 'd1', 'ttl_seconds' => 60, + 'scheduled_at' => '2026-10-03T12:00:00+00:00', + ], $env); + // дефолтный priority не шлётся + self::assertSame( + ['source' => 's', 'subject' => 'x', 'action' => 'y'], + EnvelopeBuilder::build('s', 'x', 'y', 'normal'), + ); + } + + public static function badNamesProvider(): iterable + { + yield 'uppercase' => ['Bugtrail']; + yield 'space' => ['a b']; + yield 'leading dash' => ['-lead']; + yield 'trailing dash' => ['trail-']; + yield 'empty' => ['']; + yield 'too long' => [str_repeat('a', 65)]; + yield 'cyrillic' => ['спецсимволы']; + } + + #[DataProvider('badNamesProvider')] + public function testRejectsBadNames(string $bad): void + { + foreach ([[$bad, 'x', 'y'], ['s', $bad, 'y'], ['s', 'x', $bad]] as $args) { + try { + EnvelopeBuilder::build(...$args); + self::fail("ожидалась ValidationException для '$args[0]', '$args[1]', '$args[2]'"); + } catch (ValidationException $ex) { + self::assertIsString($ex->detail); + } + } + } + + public function testRejectsPriorityTtlDedup(): void + { + try { + EnvelopeBuilder::build('s', 'x', 'y', 'urgent'); + self::fail('ожидалась ValidationException'); + } catch (ValidationException $ex) { + self::assertStringContainsString('priority', $ex->detail); + } + foreach ([0, -5, 604801] as $ttl) { + try { + EnvelopeBuilder::build('s', 'x', 'y', ttlSeconds: $ttl); + self::fail("ожидалась ValidationException для ttlSeconds=$ttl"); + } catch (ValidationException $ex) { + self::assertStringContainsString('ttlSeconds', $ex->detail); + } + } + self::assertSame(604800, EnvelopeBuilder::build('s', 'x', 'y', ttlSeconds: 604800)['ttl_seconds']); + try { + EnvelopeBuilder::build('s', 'x', 'y', dedupKey: str_repeat('d', 256)); + self::fail('ожидалась ValidationException'); + } catch (ValidationException $ex) { + self::assertStringContainsString('dedup_key', $ex->detail); + } + } + + public function testRejectsListPayload(): void + { + $this->expectException(ValidationException::class); + $this->expectExceptionMessage('словар'); + EnvelopeBuilder::build('s', 'x', 'y', payload: ['не', 'словарь']); + } + + public function testWarnsOnBadUserId(): void + { + $logs = new FakeLogger(); + EnvelopeBuilder::warnBadUserId(['user_id' => '42'], $logs); + EnvelopeBuilder::warnBadUserId(['user_id' => 42], $logs); + self::assertCount(0, $logs->warnings()); + EnvelopeBuilder::warnBadUserId(['user_id' => ''], $logs); + EnvelopeBuilder::warnBadUserId(['user_id' => null], $logs); + EnvelopeBuilder::warnBadUserId([], $logs); + self::assertCount(3, $logs->warnings()); + } +} \ No newline at end of file diff --git a/tests/FakeHttpClient.php b/tests/FakeHttpClient.php new file mode 100644 index 0000000..3f74a10 --- /dev/null +++ b/tests/FakeHttpClient.php @@ -0,0 +1,71 @@ + */ + public array $requests = []; + + /** @var \Closure(RequestInterface): \Psr\Http\Message\ResponseInterface */ + private \Closure $handler; + + /** @param callable(RequestInterface): \Psr\Http\Message\ResponseInterface $handler */ + public function __construct(callable $handler) + { + $this->handler = \Closure::fromCallable($handler); + } + + public function sendRequest(RequestInterface $request): ResponseInterface + { + $this->requests[] = $request; + return ($this->handler)($request); + } + + /** Фейк-клиент, бросающий транспортную ошибку PSR-18. */ + public static function failing(\Throwable $ex): self + { + return new self(static function () use ($ex): never { + throw $ex; + }); + } +} + +/** Транспортная ошибка, как её бросают реальные PSR-18 клиенты (Guzzle). */ +final class FakeNetworkException extends \RuntimeException implements \Psr\Http\Client\NetworkExceptionInterface +{ + public function __construct(string $message, private readonly RequestInterface $request) + { + parent::__construct($message); + } + + public function getRequest(): RequestInterface + { + return $this->request; + } +} + +/** PSR-17 фейк: строит guzzlehttp/psr7 запросы/стримы (dev-зависимость). */ +final class FakePsr17 +{ + public static function requestFactory(): \Psr\Http\Message\RequestFactoryInterface + { + return new \GuzzleHttp\Psr7\HttpFactory(); + } + + public static function streamFactory(): \Psr\Http\Message\StreamFactoryInterface + { + return new \GuzzleHttp\Psr7\HttpFactory(); + } +} \ No newline at end of file diff --git a/tests/SynapseClientTest.php b/tests/SynapseClientTest.php new file mode 100644 index 0000000..e8aba38 --- /dev/null +++ b/tests/SynapseClientTest.php @@ -0,0 +1,266 @@ + */ + public array $records = []; + + public function emergency(string|\Stringable $message, array $context = []): void { $this->log('emergency', $message, $context); } + public function alert(string|\Stringable $message, array $context = []): void { $this->log('alert', $message, $context); } + public function critical(string|\Stringable $message, array $context = []): void { $this->log('critical', $message, $context); } + public function error(string|\Stringable $message, array $context = []): void { $this->log('error', $message, $context); } + public function warning(string|\Stringable $message, array $context = []): void { $this->log('warning', $message, $context); } + public function notice(string|\Stringable $message, array $context = []): void { $this->log('notice', $message, $context); } + public function info(string|\Stringable $message, array $context = []): void { $this->log('info', $message, $context); } + public function debug(string|\Stringable $message, array $context = []): void { $this->log('debug', $message, $context); } + public function log(mixed $level, string|\Stringable $message, array $context = []): void + { + $this->records[] = [$level, (string) $message]; + } + + public function warnings(): array + { + return array_values(array_filter($this->records, static fn ($r) => $r[0] === 'warning')); + } +} + +final class SynapseClientTest extends TestCase +{ + // -- send --------------------------------------------------------- + + public function testSendSuccessBuildsRequest(): void + { + $req = null; + $client = $this->clientWith(static function ($r) use (&$req) { + $req = $r; + return new Response(202, [], json_encode([ + 'ok' => true, 'id' => 'e-1', 'status' => 'queued', 'deduplicated' => false, + ])); + }); + $event = $client->send('bugtrail', 'test', 'failed', 'high', ['user_id' => '42'], 'd1'); + self::assertSame('e-1', $event->id); + self::assertSame('queued', $event->status); + self::assertFalse($event->deduplicated); + self::assertSame('/api/v1/events', $req->getUri()->getPath()); + self::assertSame('Bearer syn_test', $req->getHeaderLine('authorization')); + self::assertStringStartsWith('gnexus-synapse-php/', $req->getHeaderLine('user-agent')); + $body = json_decode((string) $req->getBody(), true, flags: JSON_THROW_ON_ERROR); + self::assertSame([ + 'source' => 'bugtrail', 'subject' => 'test', 'action' => 'failed', + 'priority' => 'high', 'payload' => ['user_id' => '42'], 'dedup_key' => 'd1', + ], $body); + } + + public function testSendMapsStatusErrors(): void + { + $cases = [ + 401 => [AuthException::class, 'API-ключ неизвестен или отозван'], + 403 => [ForbiddenException::class, 'не совпадает с источником ключа'], + 422 => [ValidationException::class, 'Тип не зарегистрирован'], + 500 => [ServerException::class, 'boom'], + 418 => [SynapseException::class, 'teapot'], + ]; + foreach ($cases as $code => [$cls, $detail]) { + $client = $this->clientWith(static fn () => new Response($code, [], json_encode(['detail' => $detail]))); + try { + $client->send('s', 'x', 'y'); + self::fail("ожидалось исключение для HTTP $code"); + } catch (SynapseException $ex) { + self::assertSame($cls, get_class($ex)); + self::assertSame($detail, $ex->detail); + self::assertSame($code, $ex->statusCode); + } + } + } + + public function testSendTransportErrorWrapped(): void + { + $cause = new FakeNetworkException('connection refused', new Request('POST', 'http://x')); + $client = new SynapseClient( + self::config(), + FakeHttpClient::failing($cause), + FakePsr17::requestFactory(), + FakePsr17::streamFactory(), + ); + try { + $client->send('s', 'x', 'y'); + self::fail('ожидался TransportException'); + } catch (TransportException $ex) { + self::assertStringContainsString('Synapse недоступен', $ex->detail); + self::assertSame($cause, $ex->transportError); + } + } + + public function testLocalValidationNoHttp(): void + { + $called = 0; + $client = $this->clientWith(static function () use (&$called) { + ++$called; + return new Response(202, [], '{}'); + }); + try { + $client->send('Wrong Source', 'x', 'y'); + self::fail('ожидалась ValidationException'); + } catch (ValidationException $ex) { + self::assertNull($ex->statusCode); // локальная, не сервер + } + self::assertSame(0, $called); + } + + // -- emit --------------------------------------------------------- + + public function testEmitSwallowsAndLogs(): void + { + $logs = new FakeLogger(); + $client = $this->clientWith( + static fn () => new Response(503, [], json_encode(['detail' => 'нет воркеров'])), + $logs, + ); + self::assertNull($client->emit('s', 'x', 'y')); + self::assertCount(1, $logs->warnings()); + self::assertStringContainsString('нет воркеров', $logs->warnings()[0][1]); + } + + public function testEmitThrowTrueRethrows(): void + { + $client = $this->clientWith(static fn () => new Response(503, [], json_encode(['detail' => 'бум']))); + $this->expectException(ServerException::class); + $client->emit('s', 'x', 'y', throw: true); + } + + // -- batch -------------------------------------------------------- + + public function testBatchPartialAndIndexRemap(): void + { + $client = $this->clientWith(static fn () => new Response(202, [], json_encode([ + 'results' => [ + ['ok' => true, 'id' => 'a', 'status' => 'queued', 'deduplicated' => false], + ['ok' => false, 'index' => 1, 'detail' => 'Тип не зарегистрирован'], + ], + ]))); + $result = $client->sendBatch([ + ['source' => 's', 'subject' => 'x', 'action' => 'y1'], + ['subject' => 'x', 'action' => 'y'], // нет source — локальный reject + ['source' => 's', 'subject' => 'x', 'action' => 'y2'], // сервер отверг + ]); + self::assertFalse($result->allAccepted()); + self::assertSame(['a'], array_map(static fn ($e) => $e->id, $result->accepted())); + $rejected = $result->rejected(); + self::assertCount(2, $rejected); + self::assertSame(1, $rejected[0]->index); + self::assertStringContainsString('source', $rejected[0]->detail); + self::assertSame(2, $rejected[1]->index); + self::assertSame('Тип не зарегистрирован', $rejected[1]->detail); + } + + public function testBatchAllRejectedByServer(): void + { + $client = $this->clientWith(static fn () => new Response(422, [], json_encode([ + 'results' => [ + ['ok' => false, 'index' => 0, 'detail' => 'Тип не зарегистрирован'], + ['ok' => false, 'index' => 1, 'detail' => 'бум'], + ], + ]))); + $result = $client->sendBatch([ + ['source' => 's', 'subject' => 'x', 'action' => 'y'], + ['source' => 's', 'subject' => 'x', 'action' => 'z'], + ]); + self::assertFalse($result->allAccepted()); + self::assertCount(2, $result->rejected()); + } + + public function testBatchEmptyNoHttp(): void + { + $client = $this->clientWith(static function (): never { + throw new \RuntimeException('HTTP не нужен'); + }); + $result = $client->sendBatch([]); + self::assertSame([], $result->results); + self::assertTrue($result->allAccepted()); + } + + // -- status / health ---------------------------------------------- + + public function testStatusParsesAnd404(): void + { + $eid = '11111111-1111-1111-1111-111111111111'; + $other = '22222222-2222-2222-2222-222222222222'; + $client = $this->clientWith(static function ($req) use ($eid, $other) { + if (str_ends_with($req->getUri()->getPath(), $other)) { + return new Response(404, [], json_encode(['detail' => 'Событие не найдено'])); + } + return new Response(200, [], json_encode([ + 'id' => $eid, 'source' => 'bugtrail', 'subject' => 'test', 'action' => 'failed', + 'priority' => 'high', 'status' => 'delivered', + 'created_at' => '2026-10-03T12:00:00+00:00', 'expires_at' => null, + 'deliveries' => [ + ['channel' => 'internal_log', 'target' => null, 'status' => 'delivered', + 'attempts' => 1, 'error' => null, 'rendered_message' => 'тест упал'], + ['channel' => 'telegram', 'target' => 'qa-chat', 'status' => 'retrying', + 'attempts' => 2, 'error' => '429', 'rendered_message' => null], + ], + ])); + }); + $st = $client->status($eid); + self::assertSame('delivered', $st->status); + self::assertSame(2, $st->deliveries[1]->attempts); + self::assertSame('retrying', $st->deliveries[1]->status); + self::assertNull($st->expiresAt); + $this->expectException(NotFoundException::class); + $client->status($other); + } + + public function testHealthAndNoKeyConfigError(): void + { + $client = $this->clientWith(static fn () => new Response(200, [], json_encode(['status' => 'ok']))); + self::assertSame(['status' => 'ok'], $client->health()); + self::assertSame(['status' => 'ok'], $client->ready()); + // ключа нет — send поднимает ConfigurationException до HTTP + $noKey = new SynapseClient( + new SynapseConfig('http://synapse.test', null), + FakeHttpClient::failing(new FakeNetworkException('не должен дойти', new Request('POST', 'http://x'))), + FakePsr17::requestFactory(), + FakePsr17::streamFactory(), + ); + $this->expectException(ConfigurationException::class); + $noKey->send('s', 'x', 'y'); + } + + // -- helpers -------------------------------------------------------- + + private static function config(): SynapseConfig + { + return new SynapseConfig('http://synapse.test', 'syn_test'); + } + + private function clientWith(callable $handler, ?LoggerInterface $logger = null): SynapseClient + { + return new SynapseClient( + self::config(), + new FakeHttpClient($handler), + FakePsr17::requestFactory(), + FakePsr17::streamFactory(), + $logger, + ); + } +} \ No newline at end of file diff --git a/tests/SynapseConfigTest.php b/tests/SynapseConfigTest.php new file mode 100644 index 0000000..b7dffb0 --- /dev/null +++ b/tests/SynapseConfigTest.php @@ -0,0 +1,63 @@ +url('/api/v1/events')); + self::assertSame('syn_abc', $cfg->apiKey); + self::assertSame(10.0, $cfg->timeoutSeconds); + self::assertSame('', $cfg->source()); + } + + public function testRejectsEmptyUrl(): void + { + $this->expectException(ConfigurationException::class); + $this->expectExceptionMessage(SynapseConfig::ENV_URL); + new SynapseConfig(''); + } + + public function testRejectsUrlWithoutScheme(): void + { + $this->expectException(ConfigurationException::class); + new SynapseConfig('synapse:8013'); // без схемы http(s) + } + + public function testRejectsZeroTimeout(): void + { + $this->expectException(ConfigurationException::class); + new SynapseConfig('http://s', timeoutSeconds: 0); + } + + public function testFromGlobalsReadsEnv(): void + { + $cfg = SynapseConfig::fromGlobals([ + SynapseConfig::ENV_URL => 'http://synapse:8013/', + SynapseConfig::ENV_API_KEY => 'syn_abc', + SynapseConfig::ENV_TIMEOUT => '3.5', + SynapseConfig::ENV_DEFAULT_SOURCE => 'bugtrail', + ]); + self::assertSame('http://synapse:8013/api/v1', $cfg->url('/api/v1')); + self::assertSame('syn_abc', $cfg->apiKey); + self::assertSame(3.5, $cfg->timeoutSeconds); + self::assertSame('bugtrail', $cfg->source()); + } + + public function testFromGlobalsDefaultsWithoutKey(): void + { + // url обязателен: пустой env — ConfigurationException, fail fast + $cfg = SynapseConfig::fromGlobals([SynapseConfig::ENV_URL => 'http://synapse.test']); + self::assertNull($cfg->apiKey); // ключа нет — send поднимет ConfigurationException, + self::assertSame(10.0, $cfg->timeoutSeconds); // health() же работает + self::assertSame('', $cfg->source()); + } +} \ No newline at end of file diff --git a/tests/bootstrap.php b/tests/bootstrap.php new file mode 100644 index 0000000..ae62546 --- /dev/null +++ b/tests/bootstrap.php @@ -0,0 +1,5 @@ +