Skip to content
Closed
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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<DatapointString>` 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<DatapointString>` 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`)

Expand Down
5 changes: 4 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions src/generic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -865,6 +865,7 @@ pub trait ApiServiceProvider {
ResponseError {
status: response.status(),
message: err.to_string(),
content_type: None,
}
});
}
Expand Down
53 changes: 49 additions & 4 deletions src/timeseries/binary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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);
Expand Down
48 changes: 48 additions & 0 deletions src/timeseries/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn std::error::Error>> {
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<DatapointsCollection<DatapointString>> = 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<DataWrapper<String>, ResponseError>) {
match result {
Ok(r) => {
Expand Down
Loading