Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .sampo/changesets/quiet-otters-degrade.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
cargo/posthog-rs: patch
---

Do not panic when the HTTP client cannot be built. A container without CA certificates makes the reqwest builder fail, which stopped the calling program. The client now logs a warning, disables itself, and lets the program continue. The feature flag pollers degrade the same way: they log a warning and stay stopped instead of panicking.
39 changes: 25 additions & 14 deletions src/client/async_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,8 +46,8 @@ fn is_retryable_feature_flags_error(err: &reqwest::Error) -> bool {

use super::common::{
already_reported, build_dedup_key, extract_flag_details, flag_called_event,
flag_event_dedup_cache, local_record, remote_record_from_detail, report_flags_error,
DetailedFlagsResponse, FlagEventDedupCache,
flag_event_dedup_cache, http_client_or_disable, local_record, remote_record_from_detail,
report_flags_error, DetailedFlagsResponse, FlagEventDedupCache,
};
use super::transport::{Completion, Control, TransportHandle};
use super::{CaptureSummary, ClientOptions};
Expand All @@ -57,7 +57,7 @@ use reqwest::header::CONTENT_ENCODING;
/// A [`Client`] facilitates interactions with the PostHog API over HTTP.
pub struct Client {
options: ClientOptions,
client: HttpClient,
client: Option<HttpClient>,
local_evaluator: Option<LocalEvaluator>,
_flag_poller: Option<AsyncFlagPoller>,
flag_event_host: OnceLock<Arc<dyn FeatureFlagEvaluationsHost>>,
Expand Down Expand Up @@ -126,12 +126,16 @@ impl FeatureFlagEvaluationsHost for AsyncFlagEventHost {
///
/// This constructor is available with the default `async-client` feature and
/// must be awaited. Passing a blank API key creates a disabled client.
/// If the HTTP client cannot be built, logs a warning and creates a disabled
/// client instead.
pub async fn client<C: Into<ClientOptions>>(options: C) -> Client {
let options = options.into().sanitize();
let client = HttpClient::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build()
.unwrap(); // Unwrap here is as safe as `HttpClient::new`
let mut options = options.into().sanitize();
let client = http_client_or_disable(
HttpClient::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build(),
&mut options,
);
Comment thread
posthog[bot] marked this conversation as resolved.

let (local_evaluator, flag_poller) =
if options.enable_local_evaluation && !options.is_disabled() {
Expand Down Expand Up @@ -178,6 +182,13 @@ pub async fn client<C: Into<ClientOptions>>(options: C) -> Client {
}

impl Client {
/// The HTTP client, or [`Error::Connection`] if initialization failed.
fn http(&self) -> Result<&HttpClient, Error> {
self.client.as_ref().ok_or_else(|| {
Error::Connection("HTTP client is unavailable; PostHog client is disabled".to_string())
})
}

/// Capture the provided event, sending it to PostHog.
///
/// # Parameters
Expand Down Expand Up @@ -580,7 +591,7 @@ impl Client {
)?;

let step = match self
.client
.http()?
.post(&prep.url)
.headers(headers)
.body(body)
Expand Down Expand Up @@ -646,7 +657,7 @@ impl Client {
let mut attempt: u32 = 1;
loop {
let mut request = self
.client
.http()?
.post(&prep.url)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -994,7 +1005,7 @@ impl Client {

let distinct_id = payload.get("distinct_id").and_then(|v| v.as_str());
let response = match self
.client
.http()?
.post(&flags_endpoint)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -1260,7 +1271,7 @@ impl Client {
let mut attempt = 1;
loop {
let request = self
.client
.http()?
.post(flags_endpoint)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -1497,7 +1508,7 @@ mod minimal_gate_tests {
let options = ClientOptions::from(("phc_test", "http://localhost:0"));
let client = Client {
options,
client: HttpClient::builder().build().unwrap(),
client: Some(HttpClient::builder().build().unwrap()),
local_evaluator: Some(LocalEvaluator::new(cache)),
_flag_poller: None,
flag_event_host: OnceLock::new(),
Expand Down Expand Up @@ -1597,7 +1608,7 @@ mod local_payload_tests {
let options = ClientOptions::from(("phc_test", "http://localhost:0"));
let client = Client {
options,
client: HttpClient::builder().build().unwrap(),
client: Some(HttpClient::builder().build().unwrap()),
local_evaluator: Some(LocalEvaluator::new(cache)),
_flag_poller: None,
flag_event_host: OnceLock::new(),
Expand Down
39 changes: 25 additions & 14 deletions src/client/blocking.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,8 @@ fn is_retryable_feature_flags_error(err: &reqwest::Error) -> bool {

use super::common::{
already_reported, build_dedup_key, extract_flag_details, flag_called_event,
flag_event_dedup_cache, local_record, remote_record_from_detail, report_flags_error,
DetailedFlagsResponse, FlagEventDedupCache,
flag_event_dedup_cache, http_client_or_disable, local_record, remote_record_from_detail,
report_flags_error, DetailedFlagsResponse, FlagEventDedupCache,
};
use super::transport::{Completion, Control, TransportHandle};
use super::{CaptureSummary, ClientOptions};
Expand All @@ -60,7 +60,7 @@ use reqwest::header::CONTENT_ENCODING;
/// A [`Client`] facilitates interactions with the PostHog API over HTTP.
pub struct Client {
options: ClientOptions,
client: HttpClient,
client: Option<HttpClient>,
local_evaluator: Option<LocalEvaluator>,
_flag_poller: Option<FlagPoller>,
flag_event_host: OnceLock<Arc<dyn FeatureFlagEvaluationsHost>>,
Expand Down Expand Up @@ -130,12 +130,16 @@ impl FeatureFlagEvaluationsHost for BlockingFlagEventHost {
///
/// Passing a blank API key creates a disabled client. Enable the default
/// `async-client` feature to use the async client instead.
/// If the HTTP client cannot be built, logs a warning and creates a disabled
/// client instead.
pub fn client<C: Into<ClientOptions>>(options: C) -> Client {
let options = options.into().sanitize();
let client = HttpClient::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build()
.unwrap(); // Unwrap here is as safe as `HttpClient::new`
let mut options = options.into().sanitize();
let client = http_client_or_disable(
HttpClient::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build(),
&mut options,
);

let (local_evaluator, flag_poller) =
if options.enable_local_evaluation && !options.is_disabled() {
Expand Down Expand Up @@ -182,6 +186,13 @@ pub fn client<C: Into<ClientOptions>>(options: C) -> Client {
}

impl Client {
/// The HTTP client, or [`Error::Connection`] if initialization failed.
fn http(&self) -> Result<&HttpClient, Error> {
self.client.as_ref().ok_or_else(|| {
Error::Connection("HTTP client is unavailable; PostHog client is disabled".to_string())
})
}

/// Capture the provided event, sending it to PostHog.
///
/// # Parameters
Expand Down Expand Up @@ -563,7 +574,7 @@ impl Client {
)?;

let step = match self
.client
.http()?
.post(&prep.url)
.headers(headers)
.body(body)
Expand Down Expand Up @@ -627,7 +638,7 @@ impl Client {
let mut attempt: u32 = 1;
loop {
let mut request = self
.client
.http()?
.post(&prep.url)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -963,7 +974,7 @@ impl Client {

let distinct_id = payload.get("distinct_id").and_then(|v| v.as_str());
let response = match self
.client
.http()?
.post(&flags_endpoint)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -1228,7 +1239,7 @@ impl Client {
let mut attempt = 1;
loop {
let request = self
.client
.http()?
.post(flags_endpoint)
.header(CONTENT_TYPE, "application/json")
.header(USER_AGENT, get_default_user_agent())
Expand Down Expand Up @@ -1420,7 +1431,7 @@ mod minimal_gate_tests {
let options = ClientOptions::from(("phc_test", "http://localhost:0"));
let client = Client {
options,
client: HttpClient::builder().build().unwrap(),
client: Some(HttpClient::builder().build().unwrap()),
local_evaluator: Some(LocalEvaluator::new(cache)),
_flag_poller: None,
flag_event_host: OnceLock::new(),
Expand Down Expand Up @@ -1519,7 +1530,7 @@ mod local_payload_tests {
let options = ClientOptions::from(("phc_test", "http://localhost:0"));
let client = Client {
options,
client: HttpClient::builder().build().unwrap(),
client: Some(HttpClient::builder().build().unwrap()),
local_evaluator: Some(LocalEvaluator::new(cache)),
_flag_poller: None,
flag_event_host: OnceLock::new(),
Expand Down
37 changes: 36 additions & 1 deletion src/client/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,15 @@ use std::sync::{Mutex, OnceLock};

use crate::client::BeforeSendHook;
use crate::client::CaptureDefaults;
use crate::client::ClientOptions;
use crate::client::FlagsFailure;
use crate::client::OnErrorHook;
use crate::client::PostHogError;
use crate::feature_flag_evaluations::{EvaluatedFlagRecord, FlagCalledEventParams};
use crate::feature_flags::{FeatureFlagsResponse, FlagDetail, FlagMetadata, FlagValue};
use crate::Error;
use crate::Event;
use tracing::error;
use tracing::{error, warn};

/// Cap on the number of `distinct_id` entries in the `$feature_flag_called`
/// dedup cache. On overflow the entire map is reset (matches the JS SDK).
Expand Down Expand Up @@ -336,6 +337,23 @@ fn normalize_payload(payload: serde_json::Value) -> serde_json::Value {
}
}

/// Disable the client when HTTP initialization fails, for example when the
/// system has no CA certificates.
pub(super) fn http_client_or_disable<C, E: std::fmt::Debug>(
result: Result<C, E>,
options: &mut ClientOptions,
) -> Option<C> {
match result {
Ok(client) => Some(client),
Err(err) => {
Comment thread
posthog[bot] marked this conversation as resolved.
// Debug includes the underlying builder error.
warn!("Failed to build HTTP client ({err:?}); disabling PostHog client");
options.disabled = true;
None
}
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down Expand Up @@ -493,4 +511,21 @@ mod tests {
apply_before_send_hooks(&options.before_send, Event::new("test", "user-1")).is_none()
);
}

#[test]
fn http_client_build_failure_disables_client() {
let mut options = crate::ClientOptionsBuilder::default()
.api_key("phc_test".to_string())
.build()
.unwrap();
assert!(!options.is_disabled());

let client = http_client_or_disable(
Err::<(), _>("failed to load native root certificates"),
&mut options,
);

assert!(client.is_none());
assert!(options.is_disabled());
}
}
52 changes: 37 additions & 15 deletions src/client/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -447,7 +447,14 @@ fn run_worker(
// any sane teardown budget.
let shutdown_timeout =
Duration::from_millis(options.shutdown_timeout_ms).min(MAX_SHUTDOWN_TIMEOUT);
let mut pipeline = Pipeline::new(&options, Arc::clone(&clock), len);
// Without an HTTP client nothing can ever be delivered, so stop the worker
// instead of running a pipeline that cannot send. Producers keep going: any
// queued completion is signalled on the way out, and once the channel is
// disconnected an enqueue releases its reserved slot.
let Some(mut pipeline) = Pipeline::new(&options, Arc::clone(&clock), Arc::clone(&len)) else {
drain_pending_completions(&rx, &len);
return;
};

let mut buffer: Vec<Event> = Vec::new();
let mut buffer_since: Option<Instant> = None;
Expand Down Expand Up @@ -673,6 +680,25 @@ fn drain_historical(
}
}

/// Build the worker's blocking HTTP client.
///
/// `build()` is fallible — it sets up the TLS trust store, and the blocking
/// client also spawns its own thread and runtime — and `Default` is not a safe
/// fallback: it calls `Client::new`, which panics on exactly those failures.
/// Return `None` instead so the worker can stop without panicking.
fn build_http(options: &ClientOptions) -> Option<reqwest::blocking::Client> {
match reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build()
{
Ok(http) => Some(http),
Err(e) => {
warn!("posthog-rs: failed to build the transport HTTP client: {e:?}");
None
}
}
}

// ===========================================================================
// V1 pipeline
// ===========================================================================
Expand Down Expand Up @@ -705,22 +731,20 @@ struct Pipeline {

#[cfg(feature = "capture-v1")]
impl Pipeline {
fn new(options: &ClientOptions, clock: Arc<dyn Clock>, len: Arc<AtomicUsize>) -> Self {
let http = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build()
.unwrap_or_default();
/// `None` when the HTTP client cannot be built, so nothing could be sent.
fn new(options: &ClientOptions, clock: Arc<dyn Clock>, len: Arc<AtomicUsize>) -> Option<Self> {
let http = build_http(options)?;
let url = options
.endpoints()
.build_custom_url(super::v1_capture::V1_CAPTURE_PATH);
Self {
Some(Self {
http,
options: options.clone(),
url,
clock,
len,
retries: VecDeque::new(),
}
})
}

fn send_batch(
Expand Down Expand Up @@ -985,22 +1009,20 @@ struct Pipeline {

#[cfg(not(feature = "capture-v1"))]
impl Pipeline {
fn new(options: &ClientOptions, clock: Arc<dyn Clock>, len: Arc<AtomicUsize>) -> Self {
let http = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(options.request_timeout_seconds))
.build()
.unwrap_or_default();
/// `None` when the HTTP client cannot be built, so nothing could be sent.
fn new(options: &ClientOptions, clock: Arc<dyn Clock>, len: Arc<AtomicUsize>) -> Option<Self> {
let http = build_http(options)?;
let url_base = options
.endpoints()
.build_url(crate::endpoints::Endpoint::Batch);
Self {
Some(Self {
http,
options: options.clone(),
url_base,
clock,
len,
retries: VecDeque::new(),
}
})
}

fn send_batch(
Expand Down
Loading