Something went wrong. Try again.
Build Reactive Signals for Bluesky's AT Protocol Firehose in Laravel
Something went wrong. Try again.
6.9 kB · 244 lines
PHP
at dev
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245<?php
namespace SocialDept\AtpSignals\Services;
use Illuminate\Support\Facades\Log;use SocialDept\AtpSignals\Contracts\CursorStore;use SocialDept\AtpSignals\Events\SignalEvent;use SocialDept\AtpSignals\Exceptions\ConnectionException;use SocialDept\AtpSignals\Support\WebSocketConnection;
class JetstreamConsumer{ protected CursorStore $cursorStore;
protected SignalRegistry $signalRegistry;
protected EventDispatcher $eventDispatcher;
protected ?WebSocketConnection $connection = null;
protected int $reconnectAttempts = 0;
protected bool $shouldStop = false;
public function __construct( CursorStore $cursorStore, SignalRegistry $signalRegistry, EventDispatcher $eventDispatcher ) { $this->cursorStore = $cursorStore; $this->signalRegistry = $signalRegistry; $this->eventDispatcher = $eventDispatcher; }
/** * Start consuming the Jetstream. */ public function start(?int $cursor = null): void { $this->shouldStop = false;
// Get cursor from storage if not explicitly provided // null = use stored cursor, 0 = start fresh (no cursor), >0 = specific cursor if ($cursor === null) { $cursor = $this->cursorStore->get(); }
// If cursor is explicitly 0, don't send it (fresh start) $url = $this->buildWebSocketUrl($cursor > 0 ? $cursor : null);
Log::info('Signal: Starting Jetstream consumer', [ 'url' => $url, 'cursor' => $cursor > 0 ? $cursor : 'none (fresh start)', 'mode' => 'firehose', ]);
$this->connect($url); }
/** * Stop consuming the Jetstream. */ public function stop(): void { $this->shouldStop = true;
if ($this->connection) { $this->connection->close(); }
Log::info('Signal: Jetstream consumer stopped'); }
/** * Connect to the Jetstream WebSocket. */ protected function connect(string $url): void { $this->connection = new WebSocketConnection();
// Set up event handlers $this->connection ->onMessage(function (string $message) { $this->handleMessage($message); }) ->onClose(function (?int $code, ?string $reason) { $this->handleClose($code, $reason); }) ->onError(function (\Exception $e) { $this->handleError($e); });
// Connect to the WebSocket endpoint $this->connection->connect($url)->then( function () { $this->reconnectAttempts = 0; Log::info('Signal: Connected to Jetstream successfully'); }, function (\Exception $e) { Log::error('Signal: Could not connect to Jetstream', [ 'error' => $e->getMessage(), ]);
if (! $this->shouldStop) { $this->attemptReconnect(); } } );
// Run the event loop (blocking) $this->connection->run(); }
/** * Handle incoming WebSocket message. */ protected function handleMessage(string $message): void { try { $data = json_decode($message, true);
if (! $data) { Log::warning('Signal: Failed to decode message');
return; }
$event = SignalEvent::fromArray($data);
// Update cursor $this->cursorStore->set($event->timeUs);
// Dispatch to matching signals $this->eventDispatcher->dispatch($event);
} catch (\Exception $e) { Log::error('Signal: Error handling message', [ 'error' => $e->getMessage(), 'trace' => $e->getTraceAsString(), ]); } }
/** * Handle WebSocket connection close. */ protected function handleClose(?int $code, ?string $reason): void { Log::warning('Signal: Connection closed', [ 'code' => $code, 'reason' => $reason ?: 'none', 'reconnect_attempts' => $this->reconnectAttempts, ]);
// Attempt reconnection if enabled if (! $this->shouldStop) { $this->attemptReconnect(); } }
/** * Handle WebSocket connection error. */ protected function handleError(\Exception $error): void { Log::error('Signal: Connection error', [ 'error' => $error->getMessage(), 'trace' => $error->getTraceAsString(), ]); }
/** * Attempt to reconnect to the Jetstream with exponential backoff. */ protected function attemptReconnect(): void { $maxAttempts = config('signal.connection.reconnect_attempts', 5);
if ($this->reconnectAttempts >= $maxAttempts) { Log::error('Signal: Max reconnection attempts reached');
throw new ConnectionException('Failed to reconnect to Jetstream after '.$maxAttempts.' attempts'); }
$this->reconnectAttempts++;
// Calculate exponential backoff delay $baseDelay = config('signal.connection.reconnect_delay', 5); $maxDelay = config('signal.connection.max_reconnect_delay', 60);
$delay = min( $baseDelay * (2 ** ($this->reconnectAttempts - 1)), $maxDelay );
Log::info('Signal: Attempting to reconnect', [ 'attempt' => $this->reconnectAttempts, 'max_attempts' => $maxAttempts, 'delay' => $delay, ]);
sleep($delay);
$cursor = $this->cursorStore->get(); $url = $this->buildWebSocketUrl($cursor);
$this->connect($url); }
/** * Build the WebSocket URL with optional cursor and collection filters. */ protected function buildWebSocketUrl(?int $cursor = null): string { $baseUrl = config('signal.websocket_url', 'wss://jetstream2.us-east.bsky.network'); $url = rtrim($baseUrl, '/').'/subscribe';
$params = [];
// Add cursor parameter if provided if ($cursor !== null) { $params[] = 'cursor='.$cursor; }
// Add collection filters from all registered signals $collections = $this->signalRegistry->all() ->flatMap(fn ($signal) => $signal->collections() ?? []) ->unique() ->filter() ->values();
if ($collections->isNotEmpty()) { foreach ($collections as $collection) { $params[] = 'wantedCollections='.urlencode($collection); } }
if (! empty($params)) { $url .= '?'.implode('&', $params); }
return $url; }}