<?php
declare(strict_types=1);
namespace GNexus\Synapse;
use GNexus\Synapse\Config\SynapseConfig;
use GNexus\Synapse\DTO\BatchRejection;
use GNexus\Synapse\DTO\BatchResult;
use GNexus\Synapse\DTO\Delivery;
use GNexus\Synapse\DTO\EventStatus;
use GNexus\Synapse\DTO\SentEvent;
use GNexus\Synapse\Envelope\EnvelopeBuilder;
use GNexus\Synapse\Exception\AuthException;
use GNexus\Synapse\Exception\ConfigurationException;
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;
use Psr\Http\Client\ClientExceptionInterface;
use Psr\Http\Client\ClientInterface;
use Psr\Http\Client\NetworkExceptionInterface;
use Psr\Http\Message\RequestFactoryInterface;
use Psr\Http\Message\RequestInterface;
use Psr\Http\Message\ResponseInterface;
use Psr\Http\Message\StreamFactoryInterface;
use Psr\Log\LoggerInterface;
use Psr\Log\NullLogger;
/**
* Тонкий клиент Ingestion API v1 хаба уведомлений Gnexus Synapse.
* Один POST на send; никаких ретраев и локальных очередей — очередь у
* Synapse своя (приём отвечает мгновенно, 202). Транспорт инжектный
* (PSR-18 + PSR-17), таймаут настраивает потребитель в своём http-клиенте.
*/
final class SynapseClient
{
public const USER_AGENT = 'gnexus-synapse-php/0.1.2';
private const API_PREFIX = '/api/v1';
private const STATUS_ERRORS = [
401 => 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<string, mixed> $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,
?string $userId = null,
): SentEvent {
$envelope = $this->buildEnvelope(
$source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, $userId,
);
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<string, mixed> $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,
?string $userId = null,
bool $throw = false,
): ?SentEvent {
try {
return $this->send(
$source, $subject, $action, $priority, $payload, $dedupKey, $ttlSeconds, $scheduledAt, $userId,
);
} 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<array<string, mixed>> $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');
}
/**
* Поллинг до целевого статуса или таймаут — для приёмки/тестов, не
* бизнес-кода. Морг транспорта не прерывает ожидание до дедлайна.
* Не достигли → StatusTimeoutException с последним известным статусом.
*
* @param list<string> $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 --------------------------------------------------------
/**
* @return array<string, mixed>
*/
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<string, mixed>|null $payload
* @return array<string, mixed>
*/
private function buildEnvelope(
?string $source,
string $subject,
string $action,
string $priority,
?array $payload,
?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,
$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<string, mixed>|list<array<string, mixed>> $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<string, mixed>
*/
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),
);
}
}