diff --git a/CHANGELOG.md b/CHANGELOG.md new file mode 100644 index 0000000..90ea6a3 --- /dev/null +++ b/CHANGELOG.md @@ -0,0 +1,34 @@ +# Changelog + +All notable changes to this project will be documented in this file. + +The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), +and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). + +## [1.1.0] - 2025-01-29 + +### Added +- Automatic retry logic with exponential backoff for all connections +- Three new configuration fields in `JetstreamConfig`: + - `max_backoff_seconds: Int` - Maximum wait time between retries (default: 60) + - `log_connection_events: Bool` - Log connection state changes (default: True) + - `log_retry_attempts: Bool` - Log detailed retry information (default: True) + +### Changed +- `start_consumer()` now automatically retries failed connections and handles disconnections +- Enhanced error handling distinguishes between harmless timeouts and real connection failures + +### Fixed +- Connection failures no longer cause application to stop +- Harmless 60-second timeouts no longer trigger unnecessary reconnections +- WebSocket disconnections are handled gracefully with automatic reconnection + +## [1.0.0] - 2024-10-28 + +### Added +- Initial release +- WebSocket consumer for AT Protocol Jetstream events +- Support for collection and DID filtering +- Zstd compression support +- Cursor-based replay +- Event parsing for Commit, Identity, and Account events diff --git a/README.md b/README.md index 4492bb2..2ea3051 100644 --- a/README.md +++ b/README.md @@ -9,29 +9,16 @@ A Gleam WebSocket consumer for AT Protocol Jetstream events. gleam add goose@1 ``` -## Example +## Quick Start ```gleam import goose import gleam/io -import gleam/option pub fn main() { - // Create a default configuration + // Use default config with automatic retry logic let config = goose.default_config() - // Or configure with custom options - let config = goose.JetstreamConfig( - endpoint: "wss://jetstream2.us-east.bsky.network/subscribe", - wanted_collections: ["app.bsky.feed.post", "app.bsky.feed.like"], - wanted_dids: [], - cursor: option.None, - max_message_size_bytes: option.None, - compress: False, - require_hello: False, - ) - - // Start consuming events goose.start_consumer(config, fn(json_event) { let event = goose.parse_event(json_event) @@ -53,8 +40,35 @@ pub fn main() { } ``` +### Custom Configuration + +```gleam +import goose +import gleam/option + +pub fn main() { + let config = goose.JetstreamConfig( + endpoint: "wss://jetstream2.us-east.bsky.network/subscribe", + wanted_collections: ["app.bsky.feed.post", "app.bsky.feed.like"], + wanted_dids: [], + cursor: option.None, + max_message_size_bytes: option.None, + compress: False, + require_hello: False, + // Retry configuration + max_backoff_seconds: 60, // Max wait between retries + log_connection_events: True, // Log connects/disconnects + log_retry_attempts: False, // Skip verbose retry logs + ) + + goose.start_consumer(config, handle_event) +} +``` + ## Configuration Options +**Note:** Goose automatically handles connection failures with exponential backoff retry logic (1s, 2s, 4s, 8s, 16s, 32s, up to max). All connections automatically retry on failure, reconnect on disconnection, and distinguish between harmless timeouts and real errors. + ### `wanted_collections` An array of Collection NSIDs to filter which records you receive (default: empty = all collections) @@ -131,6 +145,41 @@ Pause replay/live-tail until server receives a `SubscriberOptionsUpdatePayload` require_hello: True ``` +### `max_backoff_seconds` +Maximum wait time in seconds between retry attempts + +- Uses exponential backoff: 1s, 2s, 4s, 8s, 16s, 32s, then capped at this value +- Default: `60` + +**Example:** +```gleam +max_backoff_seconds: 120 // Allow up to 2 minute waits between retries +``` + +### `log_connection_events` +Whether to log connection state changes (connected, disconnected) + +- Set to `True` to log important connection events +- Default: `True` +- Recommended: `True` for production (know when disconnects happen) + +**Example:** +```gleam +log_connection_events: True +``` + +### `log_retry_attempts` +Whether to log detailed retry attempt information + +- Set to `True` to log attempt numbers, errors, and backoff times +- Default: `True` +- Recommended: `False` for production (reduces log noise) + +**Example:** +```gleam +log_retry_attempts: False // Production: skip verbose retry logs +``` + ## Full Configuration Example ```gleam @@ -145,6 +194,9 @@ let config = goose.JetstreamConfig( max_message_size_bytes: option.Some(2097152), // 2MB compress: True, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: False, ) goose.start_consumer(config, handle_event) diff --git a/example/src/example.gleam b/example/src/example.gleam index fa58352..9703940 100644 --- a/example/src/example.gleam +++ b/example/src/example.gleam @@ -13,13 +13,16 @@ pub fn main() -> Nil { max_message_size_bytes: option.None, compress: True, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) io.println("Starting Jetstream consumer...") io.println("Connected to: " <> config.endpoint) io.println("Listening for all events...\n") - // Start consuming and log all events + // Start consuming and log all events (automatically retries on failure) goose.start_consumer(config, fn(json_event) { let event = goose.parse_event(json_event) diff --git a/gleam.toml b/gleam.toml index 172c24a..7610902 100644 --- a/gleam.toml +++ b/gleam.toml @@ -1,5 +1,5 @@ name = "goose" -version = "1.0.0" +version = "1.1.0" description = "A Gleam WebSocket consumer for AT Protocol Jetstream events" licences = ["Apache-2.0"] repository = { type = "github", user = "bigmoves", repo = "goose" } diff --git a/src/goose.gleam b/src/goose.gleam index 944de47..6548c7b 100644 --- a/src/goose.gleam +++ b/src/goose.gleam @@ -1,6 +1,8 @@ import gleam/dynamic.{type Dynamic} import gleam/dynamic/decode +import gleam/erlang/atom import gleam/erlang/process.{type Pid} +import gleam/int import gleam/io import gleam/json import gleam/list @@ -44,10 +46,17 @@ pub type JetstreamConfig { max_message_size_bytes: Option(Int), compress: Bool, require_hello: Bool, + /// Maximum backoff time in seconds for retry logic (default: 60) + max_backoff_seconds: Int, + /// Whether to log connection events (connected, disconnected) (default: True) + log_connection_events: Bool, + /// Whether to log retry attempts and errors (default: True) + log_retry_attempts: Bool, ) } /// Create a default configuration for US East endpoint +/// Includes automatic retry with exponential backoff (1s, 2s, 4s, 8s, 16s, 32s, capped at 60s) pub fn default_config() -> JetstreamConfig { JetstreamConfig( endpoint: "wss://jetstream2.us-east.bsky.network/subscribe", @@ -57,6 +66,9 @@ pub fn default_config() -> JetstreamConfig { max_message_size_bytes: option.None, compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) } @@ -127,10 +139,35 @@ pub fn connect( compress: Bool, ) -> Result(Pid, Dynamic) -/// Start consuming the Jetstream feed +/// Start consuming the Jetstream feed with automatic retry logic +/// +/// Handles connection failures gracefully with exponential backoff and automatic reconnection. +/// The retry behavior is configured through the JetstreamConfig fields: +/// - max_backoff_seconds: Maximum wait time between retries +/// - log_connection_events: Log connects/disconnects +/// - log_retry_attempts: Log retry attempts and errors +/// +/// Example: +/// ```gleam +/// let config = goose.default_config() +/// +/// goose.start_consumer(config, fn(event_json) { +/// // Handle event +/// io.println(event_json) +/// }) +/// ``` pub fn start_consumer( config: JetstreamConfig, on_event: fn(String) -> Nil, +) -> Nil { + start_with_retry_internal(config, on_event, 0) +} + +/// Internal function to handle connection with retry +fn start_with_retry_internal( + config: JetstreamConfig, + on_event: fn(String) -> Nil, + retry_count: Int, ) -> Nil { let url = build_url(config) let self = process.self() @@ -138,33 +175,115 @@ pub fn start_consumer( case result { Ok(_conn_pid) -> { - receive_loop(on_event) + case config.log_connection_events { + True -> io.println("Connected to Jetstream successfully") + False -> Nil + } + // Start receiving with retry support + receive_with_retry(config, on_event) } Error(err) -> { - io.println("Failed to connect to Jetstream") - io.println_error(string.inspect(err)) + // Connection failed, calculate backoff and retry + let backoff_seconds = calculate_backoff(retry_count, config) + case config.log_retry_attempts { + True -> { + io.println( + "Failed to connect to Jetstream (attempt " + <> int.to_string(retry_count + 1) + <> "): " + <> string.inspect(err), + ) + io.println( + "Retrying in " <> int.to_string(backoff_seconds) <> " seconds...", + ) + } + False -> Nil + } + + // Sleep for backoff period + process.sleep(backoff_seconds * 1000) + + // Retry connection + start_with_retry_internal(config, on_event, retry_count + 1) } } } -/// Receive loop for WebSocket messages -fn receive_loop(on_event: fn(String) -> Nil) -> Nil { - // Call Erlang to receive one message +/// Calculate exponential backoff with configurable maximum +fn calculate_backoff(retry_count: Int, config: JetstreamConfig) -> Int { + let backoff = case retry_count { + 0 -> 1 + 1 -> 2 + 2 -> 4 + 3 -> 8 + 4 -> 16 + 5 -> 32 + _ -> config.max_backoff_seconds + } + + // Cap at max_backoff_seconds + case backoff > config.max_backoff_seconds { + True -> config.max_backoff_seconds + False -> backoff + } +} + +/// Receive messages with retry logic +fn receive_with_retry( + config: JetstreamConfig, + on_event: fn(String) -> Nil, +) -> Nil { case receive_ws_message() { Ok(text) -> { on_event(text) - receive_loop(on_event) + receive_with_retry(config, on_event) } - Error(_) -> { - // Timeout or error, continue loop - receive_loop(on_event) + Error(error_dynamic) -> { + // Decode error type + let atm = atom.cast_from_dynamic(error_dynamic) + let error_type = atom.to_string(atm) + + case error_type { + "timeout" -> { + // No messages in 60s, connection is alive - continue + receive_with_retry(config, on_event) + } + "closed" -> { + // Connection closed - log if configured + case config.log_connection_events { + True -> io.println("Jetstream connection closed, reconnecting...") + False -> Nil + } + start_with_retry_internal(config, on_event, 0) + } + "connection_error" -> { + // Connection error - log if configured + case config.log_connection_events { + True -> io.println("Jetstream connection error, reconnecting...") + False -> Nil + } + start_with_retry_internal(config, on_event, 0) + } + _ -> { + // Unknown error - log if retry logging is enabled + case config.log_retry_attempts { + True -> { + io.println("Unknown Jetstream error: " <> error_type) + io.println("Reconnecting...") + } + False -> Nil + } + start_with_retry_internal(config, on_event, 0) + } + } } } } /// Receive a WebSocket message from the message queue +/// Returns Ok(text) for messages, or Error with one of: timeout, closed, connection_error @external(erlang, "goose_ffi", "receive_ws_message") -fn receive_ws_message() -> Result(String, Nil) +fn receive_ws_message() -> Result(String, Dynamic) /// Parse a JSON event string into a JetstreamEvent pub fn parse_event(json_string: String) -> JetstreamEvent { diff --git a/src/goose_ffi.erl b/src/goose_ffi.erl index f1c9e5a..3ededf7 100644 --- a/src/goose_ffi.erl +++ b/src/goose_ffi.erl @@ -11,13 +11,13 @@ receive_ws_message() -> %% Ignore binary messages, try again receive_ws_message(); {ws_closed, _Reason} -> - {error, nil}; + {error, closed}; {ws_error, _Reason} -> - {error, nil}; + {error, connection_error}; _Other -> %% Ignore unexpected messages receive_ws_message() after 60000 -> - %% Timeout - return error to continue loop - {error, nil} + %% Timeout - connection is still alive, just no messages + {error, timeout} end. diff --git a/test/goose_test.gleam b/test/goose_test.gleam index 5ffb1e2..e336557 100644 --- a/test/goose_test.gleam +++ b/test/goose_test.gleam @@ -26,6 +26,9 @@ pub fn build_url_with_collections_test() { max_message_size_bytes: option.None, compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -44,6 +47,9 @@ pub fn build_url_with_dids_test() { max_message_size_bytes: option.None, compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -62,6 +68,9 @@ pub fn build_url_with_both_test() { max_message_size_bytes: option.None, compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -80,6 +89,9 @@ pub fn build_url_with_cursor_test() { max_message_size_bytes: option.None, compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -98,6 +110,9 @@ pub fn build_url_with_max_size_test() { max_message_size_bytes: option.Some(1_048_576), compress: False, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -116,6 +131,9 @@ pub fn build_url_with_compress_test() { max_message_size_bytes: option.None, compress: True, require_hello: False, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -134,6 +152,9 @@ pub fn build_url_with_require_hello_test() { max_message_size_bytes: option.None, compress: False, require_hello: True, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config) @@ -152,6 +173,9 @@ pub fn build_url_with_all_options_test() { max_message_size_bytes: option.Some(2_097_152), compress: True, require_hello: True, + max_backoff_seconds: 60, + log_connection_events: True, + log_retry_attempts: True, ) let url = goose.build_url(config)