From 37e749b8a57dcad813e8084ddc7a773b1dd1b730 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Thu, 17 Sep 2026 09:41:36 +0200 Subject: [PATCH 1/2] fix: binary ingest builds against ResponseError::content_type #123 was branched before #128 gave ResponseError a content_type field and merged after it, so main stopped compiling at the three places the binary path builds one. All three are raised by the SDK with no server response behind them (a 204 that fails to deserialize, a value refused locally, a series /byids cannot find), so each carries None. Co-Authored-By: Claude Opus 5 (1M context) Signed-off-by: jgjesdal --- src/generic.rs | 1 + src/timeseries/binary.rs | 2 ++ 2 files changed, 3 insertions(+) diff --git a/src/generic.rs b/src/generic.rs index 56ed379..feddc72 100644 --- a/src/generic.rs +++ b/src/generic.rs @@ -865,6 +865,7 @@ pub trait ApiServiceProvider { ResponseError { status: response.status(), message: err.to_string(), + content_type: None, } }); } diff --git a/src/timeseries/binary.rs b/src/timeseries/binary.rs index 0ebea83..fee9834 100644 --- a/src/timeseries/binary.rs +++ b/src/timeseries/binary.rs @@ -737,6 +737,7 @@ fn unprocessable(message: String) -> ResponseError { ResponseError { status: StatusCode::UNPROCESSABLE_ENTITY, message, + content_type: None, } } @@ -931,6 +932,7 @@ impl TimeSeriesService { return Err(ResponseError { status: StatusCode::NOT_FOUND, message: format!("Could not find following timeseries: {}", missing.join(", ")), + content_type: None, }); } Ok(resolved) From a52fdd4717575a78b4491161edf356ef6621a353 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Thu, 17 Sep 2026 09:41:44 +0200 Subject: [PATCH 2/2] fix(binary): decide the stale-series retry from the problem document MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit is_stale_series_rejection searched the raw body for "unknown-timeseries" or "external-id-mismatch", the idiom from before ResponseError::problem() existed. That also matched the SDK's own client-side 404, which names the series it could not resolve, so an external id spelling one of those tokens earned a pointless second attempt. It now reads the problem: the api answers every binary refusal with the one type datapoint-block-rejected and puts the cause in a `reason` extension, so the slug alone does not discriminate here the way it does on every other endpoint. README and AGENTS.md say so, and AGENTS.md drops a stale "not in the Python bindings yet" — #123 added them. The retry had no live coverage because the SDK resolves every series before sending. test_insert_datapoints_binary_re_resolves_a_recreated_series deletes and recreates a series under the same external id, so the server refuses the cached id; it fails if unknown-timeseries stops matching. Co-Authored-By: Claude Opus 5 (1M context) Signed-off-by: jgjesdal --- AGENTS.md | 2 +- README.md | 5 +++- src/timeseries/binary.rs | 51 ++++++++++++++++++++++++++++++++++++---- src/timeseries/test.rs | 48 +++++++++++++++++++++++++++++++++++++ 4 files changed, 100 insertions(+), 6 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 719c041..f24e78c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -125,7 +125,7 @@ When a datapoint/event send can't get through, ingestion spools to a segmented, ### Binary datapoint ingest (`src/timeseries/binary.rs`) -`TimeSeriesService::insert_datapoints_binary` is the second ingest path, to `POST /timeseries/data/binary`: the same `DatapointsCollection` input, resolved through `/timeseries/byids` (cached per service instance; needs read access on the dataset), checked against each series' value type locally, cut into Arrow IPC frames at the contract's caps, zstd-compressed per frame (mandatory: level 1, 3 or 9, default 9) and posted through `execute_post_bytes_request`. `FrameWriter` builds one frame and is public; `cut_into_writers` and `pack_requests` hold the caps. The byte layout is the platform's `binary_datapoints_format.md`. The arrow-rs crates (`arrow-array`, `arrow-schema`, `arrow-ipc`) exist for this path and are the seed of the Arrow read path. Not in the Python bindings yet. The ignored `timeseries::tests::test_datapoints_binary` is the live twin of `test_datapoints` and needs a backend that serves the endpoint; the writer's own tests in `binary.rs` run offline. +`TimeSeriesService::insert_datapoints_binary` is the second ingest path, to `POST /timeseries/data/binary`: the same `DatapointsCollection` input, resolved through `/timeseries/byids` (cached per service instance; needs read access on the dataset), checked against each series' value type locally, cut into Arrow IPC frames at the contract's caps, zstd-compressed per frame (mandatory: level 1, 3 or 9, default 9) and posted through `execute_post_bytes_request`. The api answers every refusal on this path with the single problem type `datapoint-block-rejected` and puts the actual cause in a `reason` extension, so here — unlike everywhere else — the slug is not the discriminant; `is_stale_series_rejection` reads `reason` to decide the one re-resolve-and-retry. `FrameWriter` builds one frame and is public; `cut_into_writers` and `pack_requests` hold the caps. The byte layout is the platform's `binary_datapoints_format.md`. The arrow-rs crates (`arrow-array`, `arrow-schema`, `arrow-ipc`) exist for this path and are the seed of the Arrow read path. The ignored `timeseries::tests::test_datapoints_binary` is the live twin of `test_datapoints` and needs a backend that serves the endpoint; the writer's own tests in `binary.rs` run offline. ### The `ApiServiceProvider` trait (`src/generic.rs`) diff --git a/README.md b/README.md index 9ae94e0..f821812 100644 --- a/README.md +++ b/README.md @@ -128,7 +128,10 @@ frame, 32 frames per request), compressed with zstd (level 9 by default, 1 and 3 choices) and posted. A 204 means every frame was accepted. A series that does not exist is a 404 before anything is sent, a value that does not fit its -type is a 422, and a 429 or a 5xx is retried. Resolving by external id goes through +type is a 422, and a 429 or a 5xx is retried. Every refusal from the server is a problem +document of type `datapoint-block-rejected`, so branch on its `reason` extension +(`unknown-timeseries`, `value-type-mismatch`, `too-many-in-flight`, …) rather than on +`problem_slug()`. Resolving by external id goes through `/timeseries/byids`, so the caller needs read access on the dataset as well as write access. The durable spool does not cover this path. `binary::FrameWriter` is public for producers that build frames themselves. diff --git a/src/timeseries/binary.rs b/src/timeseries/binary.rs index fee9834..1905a0c 100644 --- a/src/timeseries/binary.rs +++ b/src/timeseries/binary.rs @@ -744,10 +744,15 @@ fn unprocessable(message: String) -> ResponseError { /// The server rejected the request because the series named in a frame no longer match what the /// client resolved: unknown after a delete, or renamed since. fn is_stale_series_rejection(error: &ResponseError) -> bool { - let status = error.get_status(); - (status == StatusCode::NOT_FOUND || status == StatusCode::UNPROCESSABLE_ENTITY) - && (error.message.contains("unknown-timeseries") - || error.message.contains("external-id-mismatch")) + let Some(problem) = error.problem() else { + return false; + }; + // Every binary refusal shares this one type; what went wrong is in `reason`. + problem.slug() == Some("datapoint-block-rejected") + && matches!( + problem.extensions.get("reason").and_then(|reason| reason.as_str()), + Some("unknown-timeseries" | "external-id-mismatch") + ) } impl TimeSeriesService { @@ -966,6 +971,44 @@ mod tests { (batches, schema) } + #[test] + fn only_a_stale_series_problem_rebuilds_the_request() { + let error = |status: StatusCode, body: &str| ResponseError { + status, + message: body.to_string(), + content_type: Some("application/problem+json".to_string()), + }; + let rejected = |reason: &str| { + format!( + r#"{{"type":"https://intellistream.ai/errors/datapoint-block-rejected","title":"Datapoint block rejected","reason":"{reason}","timeseriesIds":[3]}}"# + ) + }; + + assert!(is_stale_series_rejection(&error( + StatusCode::NOT_FOUND, + &rejected("unknown-timeseries") + ))); + assert!(is_stale_series_rejection(&error( + StatusCode::UNPROCESSABLE_ENTITY, + &rejected("external-id-mismatch") + ))); + assert!(!is_stale_series_rejection(&error( + StatusCode::UNPROCESSABLE_ENTITY, + &rejected("value-type-mismatch") + ))); + assert!(!is_stale_series_rejection(&error( + StatusCode::NOT_FOUND, + r#"{"type":"https://intellistream.ai/errors/not-found","reason":"unknown-timeseries"}"# + ))); + // The SDK's own 404 names the series it could not resolve, and an external id can spell a + // reason: matching on the body text retried this. + assert!(!is_stale_series_rejection(&ResponseError { + status: StatusCode::NOT_FOUND, + message: "Could not find following timeseries: unknown-timeseries-pump".to_string(), + content_type: None, + })); + } + #[test] fn float_frame_is_sorted_deduplicated_and_readable_by_arrow() { let mut writer = FrameWriter::new(DatapointValueType::Float); diff --git a/src/timeseries/test.rs b/src/timeseries/test.rs index 4512a75..8866c29 100644 --- a/src/timeseries/test.rs +++ b/src/timeseries/test.rs @@ -1573,6 +1573,54 @@ mod tests { Ok(()) } + /// Recreating a series under the same external id gives it a new id, so the server refuses the + /// one this service cached as `unknown-timeseries`; the call must re-resolve and succeed. + #[tokio::test] + async fn test_insert_datapoints_binary_re_resolves_a_recreated_series() -> Result<(), Box> { + let api_service = create_api_service(); + let ext_id = unique_id("ts_binary_recreated"); + let mut ts_cleanup = cleanup_timeseries(vec![ext_id.clone()]); + + let mut ts_collection = DataWrapper::new(); + ts_collection.add_item( + TimeSeries::builder() + .set_external_id(&ext_id) + .set_name(&ext_id) + .set_unit("celsius") + .set_value_type("float") + .clone(), + ); + let mut data_request: DataWrapper> = DataWrapper::new(); + let mut dp_collection = DatapointsCollection::from_external_id(&ext_id); + dp_collection.datapoints = vec![ + DatapointString::from_datetime(Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap(), "42.0"), + ]; + data_request.add_item(dp_collection); + + api_service.time_series.create(&ts_collection).await.expect("could not create the series"); + api_service + .time_series + .insert_datapoints_binary(&data_request, &BinaryIngestOptions::default()) + .await + .expect("the first binary insert failed"); + + delete_timeseries(&api_service, &[&ext_id]).await; + api_service.time_series.create(&ts_collection).await.expect("could not recreate the series"); + + match api_service + .time_series + .insert_datapoints_binary(&data_request, &BinaryIngestOptions::default()) + .await + { + Ok(r) => assert_eq!(r.get_http_status_code().unwrap(), StatusCode::NO_CONTENT.as_u16()), + Err(e) => panic!("the cached id was not re-resolved: {}: {}", e.get_status(), e.get_message()), + } + + delete_timeseries(&api_service, &[&ext_id]).await; + ts_cleanup.disarm(); + Ok(()) + } + fn validate_data_insertion(result: Result, ResponseError>) { match result { Ok(r) => {