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/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..1905a0c 100644 --- a/src/timeseries/binary.rs +++ b/src/timeseries/binary.rs @@ -737,16 +737,22 @@ fn unprocessable(message: String) -> ResponseError { ResponseError { status: StatusCode::UNPROCESSABLE_ENTITY, message, + content_type: None, } } /// 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 { @@ -931,6 +937,7 @@ impl TimeSeriesService { return Err(ResponseError { status: StatusCode::NOT_FOUND, message: format!("Could not find following timeseries: {}", missing.join(", ")), + content_type: None, }); } Ok(resolved) @@ -964,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) => {