From 501a4910c56a97a702e16a85c74f60c873435875 Mon Sep 17 00:00:00 2001 From: Nick Gerakines <12125+ngerakines@users.noreply.github.com> Date: Thu, 7 Nov 2024 08:57:51 -0500 Subject: [PATCH] feature: compression configuration option Signed-off-by: Nick Gerakines <12125+ngerakines@users.noreply.github.com> --- README.md | 3 +- dev-server.sh | 3 +- docs/playbook-enable-compression.md | 34 +++++++++++++++ src/bin/supercell.rs | 1 + src/config.rs | 30 ++++++++++++- src/consumer.rs | 66 +++++++++++++++++------------ 6 files changed, 106 insertions(+), 31 deletions(-) create mode 100644 docs/playbook-enable-compression.md diff --git a/README.md b/README.md index 395cfd3..fb095d5 100644 --- a/README.md +++ b/README.md @@ -12,7 +12,8 @@ The following environment variables are used: * `EXTERNAL_BASE` - The hostname of the feed generator. * `DATABASE_URL` - The URL of the database to use. * `JETSTREAM_HOSTNAME` - The hostname of the JetStream server to consume events from. -* `ZSTD_DICTIONARY` - The path to the ZSTD dictionary to use. +* `COMPRESSION` - Use zstd compression. Default `false`. +* `ZSTD_DICTIONARY` - The path to the ZSTD dictionary to use. Required when compression is enabled. * `CONSUMER_TASK_ENABLE` - Whether or not to enable the consumer tasks. Default `true`. * `VMC_TASK_ENABLE` - Whether or not to enable the VMC (verification method cache) tasks. Default `true`. * `PLC_HOSTNAME` - The hostname of the PLC server to use for VMC tasks. Default `plc.directory`. diff --git a/dev-server.sh b/dev-server.sh index e28e259..7caabd3 100755 --- a/dev-server.sh +++ b/dev-server.sh @@ -4,9 +4,10 @@ export HTTP_PORT=4050 export EXTERNAL_BASE=feeds.smokesignal.events export DATABASE_URL=sqlite://development.db export JETSTREAM_HOSTNAME=jetstream1.us-east.bsky.network -export ZSTD_DICTIONARY=$(pwd)/jetstream_zstd_dictionary export CONSUMER_TASK_ENABLE=true export FEEDS=$(pwd)/config.yml +# export COMPRESSION=true +# export ZSTD_DICTIONARY=$(pwd)/jetstream_zstd_dictionary touch development.db sqlx migrate run --database-url sqlite://development.db diff --git a/docs/playbook-enable-compression.md b/docs/playbook-enable-compression.md new file mode 100644 index 0000000..65216af --- /dev/null +++ b/docs/playbook-enable-compression.md @@ -0,0 +1,34 @@ +# Playbook: Enable compression + +Jetstream supports optional zstd compression for streamed events. This feature can reduce bandwidth usage by up to 50% with minimal performance impact, but at the cost of a more complicated deployment and additional maintenance steps. + +[zstd](https://github.com/facebook/zstd) is a dictionary-based compression algorithm that is optimized for real-time compression and decompression. + +## Configuration + +Compression is enabled by setting the `COMPRESSION` environment variable to `true`. When enabled, the `ZSTD_DICTIONARY` environment variable must be set to the path of the ZSTD dictionary to use. + +The dictionary file can be downloaded from [github.com/bluesky-social/jetstream/blob/main/pkg/models/zstd_dictionary](https://github.com/bluesky-social/jetstream/blob/main/pkg/models/zstd_dictionary). + +## FAQ + +### Why is compression disabled by default? + +The benefits of compression are not guaranteed and depend on the data being compressed. For most supercell deployments, the impact on CPU and memory is minimal, and the benefits of reduced bandwidth is significant. However, this feature is maturing and may not be suitable for all deployments. + +### Why is a custom dictionary required? + +Zstd uses a dictionary to improve compression performance. There is no default dictionary, so a custom dictionary must be provided. + +The dictionary is occasionally rebuilt to improve compression performance. The dictionary is built from a sample of the data that is being compressed, so the dictionary is specific to the data that is being compressed. + +### Dictionary mismatch + +This error occurs when compression is enabled and the dictionary is invalid or does not match the dictionary that was used to compress the data. + +To resolve, ensure the `ZSTD_DICTIONARY` environment variable is set to the correct path and that the file at that path is the same as the one from the jetstream repository. + +### Destination buffer is too small + +This error occurs when the buffer used to decompress the data is too small to contain the decompressed event. + diff --git a/src/bin/supercell.rs b/src/bin/supercell.rs index f5b48dd..1f79be1 100644 --- a/src/bin/supercell.rs +++ b/src/bin/supercell.rs @@ -102,6 +102,7 @@ async fn main() -> Result<()> { if task_enable { let consumer_task_config = ConsumerTaskConfig { user_agent: inner_config.user_agent.clone(), + compression: *inner_config.compression.as_ref(), zstd_dictionary_location: inner_config.zstd_dictionary.clone(), jetstream_hostname: inner_config.jetstream_hostname.clone(), feeds: inner_config.feeds.clone(), diff --git a/src/config.rs b/src/config.rs index c937888..7dec8e0 100644 --- a/src/config.rs +++ b/src/config.rs @@ -45,6 +45,9 @@ pub struct CertificateBundles(Vec); #[derive(Clone)] pub struct TaskEnable(bool); +#[derive(Clone)] +pub struct Compression(bool); + #[derive(Clone)] pub struct Config { pub version: String, @@ -59,6 +62,7 @@ pub struct Config { pub zstd_dictionary: String, pub jetstream_hostname: String, pub feeds: Feeds, + pub compression: Compression, } impl Config { @@ -72,7 +76,14 @@ impl Config { optional_env("CERTIFICATE_BUNDLES").try_into()?; let jetstream_hostname = require_env("JETSTREAM_HOSTNAME")?; - let zstd_dictionary = require_env("ZSTD_DICTIONARY")?; + + let compression: Compression = default_env("COMPRESSION", "false").try_into()?; + + let zstd_dictionary = if compression.0 { + require_env("ZSTD_DICTIONARY")? + } else { + "".to_string() + }; let consumer_task_enable: TaskEnable = default_env("CONSUMER_TASK_ENABLE", "true").try_into()?; @@ -103,6 +114,7 @@ impl Config { jetstream_hostname, zstd_dictionary, feeds, + compression, }) } } @@ -186,6 +198,22 @@ impl TryFrom for TaskEnable { } } +impl AsRef for Compression { + fn as_ref(&self) -> &bool { + &self.0 + } +} + +impl TryFrom for Compression { + type Error = anyhow::Error; + fn try_from(value: String) -> Result { + let value = value.parse::().map_err(|err| { + anyhow::Error::new(err).context(anyhow!("parsing compression into bool failed")) + })?; + Ok(Self(value)) + } +} + impl TryFrom for Feeds { type Error = anyhow::Error; fn try_from(value: String) -> Result { diff --git a/src/consumer.rs b/src/consumer.rs index b222f1e..87fe9f1 100644 --- a/src/consumer.rs +++ b/src/consumer.rs @@ -1,6 +1,6 @@ use std::str::FromStr; -use anyhow::{Context, Result}; +use anyhow::{anyhow, Context, Result}; use futures_util::SinkExt; use futures_util::StreamExt; use http::HeaderValue; @@ -22,6 +22,7 @@ const MAX_MESSAGE_SIZE: usize = 25000; #[derive(Clone)] pub struct ConsumerTaskConfig { pub user_agent: String, + pub compression: bool, pub zstd_dictionary_location: String, pub jetstream_hostname: String, pub feeds: config::Feeds, @@ -56,16 +57,14 @@ impl ConsumerTask { let last_time_us = consumer_control_get(&self.pool, &self.config.jetstream_hostname).await?; - // mkdir -p data/ && curl -o data/zstd_dictionary https://github.com/bluesky-social/jetstream/raw/refs/heads/main/pkg/models/zstd_dictionary - let data: Vec = std::fs::read(self.config.zstd_dictionary_location.clone()) - .context("unable to load zstd dictionary")?; - let uri = Uri::from_str(&format!( - "wss://{}/subscribe?compress=true&requireHello=true", - self.config.jetstream_hostname + "wss://{}/subscribe?compress={}&requireHello=true", + self.config.jetstream_hostname, self.config.compression )) .context("invalid jetstream URL")?; + tracing::debug!(uri = ?uri, "connecting to jetstream"); + let (mut client, _) = ClientBuilder::from_uri(uri) .add_header( http::header::USER_AGENT, @@ -89,8 +88,16 @@ impl ConsumerTask { .await .map_err(|err| anyhow::Error::msg(err).context("cannot send update"))?; - let mut decompressor = zstd::bulk::Decompressor::with_dictionary(&data) - .map_err(|err| anyhow::Error::msg(err).context("cannot create decompressor"))?; + let mut decompressor = if self.config.compression { + // mkdir -p data/ && curl -o data/zstd_dictionary https://github.com/bluesky-social/jetstream/raw/refs/heads/main/pkg/models/zstd_dictionary + let data: Vec = std::fs::read(self.config.zstd_dictionary_location.clone()) + .context("unable to load zstd dictionary")?; + zstd::bulk::Decompressor::with_dictionary(&data) + .map_err(|err| anyhow::Error::msg(err).context("cannot create decompressor"))? + } else { + zstd::bulk::Decompressor::new() + .map_err(|err| anyhow::Error::msg(err).context("cannot create decompressor"))? + }; let interval = std::time::Duration::from_secs(120); let sleeper = sleep(interval); @@ -120,29 +127,30 @@ impl ConsumerTask { } let item = item.unwrap(); - if !item.is_binary() { - tracing::warn!("message from jetstream is not binary"); - continue; - } - let payload = item.into_payload(); - - let decoded = decompressor.decompress(&payload, MAX_MESSAGE_SIZE * 3); - if let Err(err) = decoded { - let length = payload.len(); - tracing::error!(error = ?err, length = ?length, "error processing jetstream message"); - continue; - } - let decoded = decoded.unwrap(); + let event = if self.config.compression { + if !item.is_binary() { + tracing::debug!("compression enabled but message from jetstream is not binary"); + continue; + } + let payload = item.into_payload(); - let event = serde_json::from_slice::(&decoded); + let decoded = decompressor.decompress(&payload, MAX_MESSAGE_SIZE * 3); + if let Err(err) = decoded { + tracing::debug!(err = ?err, "cannot decompress message"); + continue; + } + let decoded = decoded.unwrap(); + serde_json::from_slice::(&decoded) + } else { + if !item.is_text() { + tracing::debug!("compression enabled but message from jetstream is not binary"); + continue; + } + serde_json::from_str(item.as_text().ok_or(anyhow!("cannot convert message to text"))?) + }; if let Err(err) = event { tracing::error!(error = ?err, "error processing jetstream message"); - #[cfg(debug_assertions)] - { - println!("{:?}", std::str::from_utf8(&decoded)); - } - continue; } let event = event.unwrap(); @@ -160,6 +168,8 @@ impl ConsumerTask { } let event_value = event_value.unwrap(); + tracing::trace!(event = ?event, "received event"); + for feed_matcher in self.feed_matchers.0.iter() { if feed_matcher.matches(&event_value) { tracing::debug!(feed_id = ?feed_matcher.feed, "matched event"); -- 2.51.2