diff --git a/AGENTS.md b/AGENTS.md index 17acc56..1d4c522 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -10,9 +10,17 @@ cargo test # substring match on test name cargo test -- --ignored # run tests marked #[ignore] (e.g. long-running datapoint tests) cargo test ::tests:: # e.g. `events::tests::test_events_full` cargo test -- --nocapture # show println! from tests (the SDK prints response bodies) +cargo test --release # ALWAYS --release for anything timed (see below) ./run_python_tests.sh # Python-bindings suite (rebuilds the PyO3 module first — see below) ``` +**Never time anything under a plain `cargo test`.** It builds with the `dev` profile at +`opt-level = 0`. The `bench_json_vs_binary` comparison first reported Rust as *slower than the +Java SDK* on both ingest paths that way, 339k against 720k points per second on JSON and 1.0M +against 3.8M on binary; the whole result was the missing `--release`. That test now panics +rather than run unoptimised, and the same applies to the PyO3 module, which needs +`maturin develop --release` before any Python timing means anything. + Most tests are integration tests that call a live backend via `create_api_service()`. They read configuration from a local `.env` file (gitignored). Required: - `BASE_URL` — backend root, e.g. `http://localhost:8081` @@ -115,6 +123,10 @@ Synchronous mirror of the async API behind the `blocking` cargo feature — the When a datapoint/event send can't get through, ingestion spools to a segmented, zstd-compressed NDJSON log on disk and flushes automatically on a later ingest call. Invariants to preserve: memory use is bounded by a single segment (plain append-only active segment, zstd-sealed at ~50 MiB rollover via temp file + atomic rename, drained oldest-first one segment at a time); bounded by time retention (whole segments past the window dropped, expired records skipped on read) and a size cap (oldest segment deleted); a torn trailing line from an unclean shutdown is skipped on read. Each on-disk line is `\t`; the spool is content-agnostic. +### 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. + ### The `ApiServiceProvider` trait (`src/generic.rs`) Every subservice implements `ApiServiceProvider`, which owns the HTTP plumbing: token acquisition, `execute_get_request`, `execute_post_request`, `execute_file_upload_request`, `execute_get_stream_request`. Subservice methods should go through these helpers rather than calling `reqwest` directly — a few early methods (e.g. `TimeSeriesService::list`) still bypass the trait and should be migrated when touched. diff --git a/Cargo.toml b/Cargo.toml index 932774f..2d07d25 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,6 +42,9 @@ geojson = "1" # Reading (never verifying) the payload of a JWT the SDK already holds, to explain an # otherwise-unexplained 401 — see `auth_diagnostics`. base64 = "0.22" +arrow-array = "59.3" +arrow-schema = "59.3" +arrow-ipc = "59.3" #[lib] diff --git a/README.md b/README.md index 1fe7073..9ae94e0 100644 --- a/README.md +++ b/README.md @@ -117,6 +117,22 @@ sealed at a ~50 MiB rollover and drained one segment at a time — so even a mul spool never loads into memory, and a torn trailing line from an unclean shutdown is skipped on read. +## Binary datapoint ingest + +`time_series.insert_datapoints_binary(&collections, &BinaryIngestOptions::default())` takes +the same collections as `insert_datapoints` and sends them through +`POST /timeseries/data/binary`: each series is resolved once to its id and value type and +cached, values are checked against that type locally, sorted and de-duplicated, cut into Arrow +IPC frames at the contract's caps (100 000 points per numeric frame, 10 000 per text or mixed +frame, 32 frames per request), compressed with zstd (level 9 by default, 1 and 3 are the other +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 +`/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. + ## Python bindings `datahub_python_bindings/` wraps this SDK as the Python package `intellistream-datahub-sdk` (import name diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index 2c983b4..18ffbd5 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -635,6 +635,18 @@ class TimeSeriesServiceSync: """ def insert_datapoints(self, input: list[DatapointsCollectionString]) -> list[str]: ... + def insert_datapoints_binary( + self, + input: list[DatapointsCollectionString], + zstd_level: int | None = None, + ) -> list[str]: ... + def insert_from_lists_binary( + self, + timestamps: list[datetime.datetime], + values: list[float], + ts: Identifiable, + zstd_level: int | None = None, + ) -> list[str]: ... def insert_from_lists( self, timestamps: list[datetime.datetime], diff --git a/datahub_python_bindings/src/timeseries/mod.rs b/datahub_python_bindings/src/timeseries/mod.rs index d0255a9..d4de4f0 100644 --- a/datahub_python_bindings/src/timeseries/mod.rs +++ b/datahub_python_bindings/src/timeseries/mod.rs @@ -469,6 +469,55 @@ pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> { Ok(()) } +/// Options for the binary ingest bindings. zstd level 1, 3 or 9; 9 is the default, as in the +/// Rust and Java SDKs, because the client is what pays for it. +pub(crate) fn binary_options( + zstd_level: Option, +) -> intellistream_datahub_sdk::timeseries::BinaryIngestOptions { + use intellistream_datahub_sdk::timeseries::BinaryIngestOptions; + match zstd_level { + Some(level) => BinaryIngestOptions::new().zstd_level(level), + None => BinaryIngestOptions::default(), + } +} + +/// Parallel timestamp and value sequences for one series, which is the shape a DataFrame's +/// index and one of its columns arrive in, turned into the wire collection. +pub(crate) fn lists_to_collection( + timestamps: Vec>, + values: Vec, + ts: Identifiable, +) -> PyResult> { + if timestamps.len() != values.len() { + return Err(PyValueError::new_err(format!( + "timestamps and values must be the same length, got {} and {}", + timestamps.len(), + values.len() + ))); + } + let datapoints: Vec = timestamps + .into_iter() + .zip(values) + .map(|(timestamp, value)| { + Ok(DatapointString { + timestamp: crate::datetime::py_datetime_to_utc(×tamp)? + .timestamp_millis() + .to_string(), + value: value.to_string(), + }) + }) + .collect::>>()?; + let id_collection = ts.id_collection(); + Ok(DatapointsCollection { + datapoints, + next_cursor: None, + id: id_collection.id, + external_id: id_collection.external_id, + unit: None, + unit_external_id: None, + }) +} + /// Reverse lookup: the events that reference this timeseries. (The other direction — /// `Event.related_resource_nodes()` — resolves an event's resources.) Available only on /// timeseriess returned by the API; calling on a locally-constructed one raises. diff --git a/datahub_python_bindings/src/timeseries/sync_service.rs b/datahub_python_bindings/src/timeseries/sync_service.rs index 4bd8a25..dad14d9 100644 --- a/datahub_python_bindings/src/timeseries/sync_service.rs +++ b/datahub_python_bindings/src/timeseries/sync_service.rs @@ -218,6 +218,53 @@ impl PyTimeSeriesServiceSync { Ok(result.get_items().clone()) }) } + /// `POST /timeseries/data/binary`: the same collections as `insert_datapoints`, sent as + /// zstd-compressed Arrow frames. `zstd_level` is 1, 3 or 9 and defaults to 9. + #[pyo3(signature = (input, zstd_level=None))] + fn insert_datapoints_binary<'py>( + &self, + py: Python<'py>, + input: Vec, + zstd_level: Option, + ) -> PyResult> { + let service = self.api_service.clone(); + let vec: Vec> = + input.into_iter().map(|item| item.into()).collect(); + let wrapper = DataWrapper::>::from_vec(vec); + let options = binary_options(zstd_level); + py.detach(|| { + let result = self + .runtime + .block_on(service.time_series.insert_datapoints_binary(&wrapper, &options)) + .map_err(|e| crate::datahub_err(e))?; + Ok(result.get_items().clone()) + }) + } + + /// The binary twin of `insert_from_lists`: parallel timestamp and value sequences for one + /// series, which is the shape a DataFrame column pair arrives in. + #[pyo3(signature = (timestamps, values, ts, zstd_level=None))] + fn insert_from_lists_binary<'py>( + &self, + py: Python<'py>, + timestamps: Vec>, + values: Vec, + ts: Identifiable, + zstd_level: Option, + ) -> PyResult> { + let service = self.api_service.clone(); + let collection = lists_to_collection(timestamps, values, ts)?; + let wrapper = DataWrapper::>::from_vec(vec![collection]); + let options = binary_options(zstd_level); + py.detach(|| { + let result = self + .runtime + .block_on(service.time_series.insert_datapoints_binary(&wrapper, &options)) + .map_err(|e| crate::datahub_err(e))?; + Ok(result.get_items().clone()) + }) + } + fn insert_from_lists<'py>( &self, py: Python<'py>, diff --git a/python_tests/bench_json_vs_binary.py b/python_tests/bench_json_vs_binary.py new file mode 100644 index 0000000..d354d19 --- /dev/null +++ b/python_tests/bench_json_vs_binary.py @@ -0,0 +1,166 @@ +# JSON against binary datapoint ingest, through the Python SDK. +# +# Not a pytest: it needs a live platform and takes minutes, so it is a script you run when you +# want the numbers. The Java equivalent lives in the platform repo's datahub-e2e module and +# reports the same columns, so the two are comparable. +# +# python python_tests/bench_json_vs_binary.py --points 10000000 +# +# Reads the same .env the SDK does. The api must have its daily quota and rate limiter off, or +# it will refuse a run of this size: +# -Ddatahub.limits.quota.enabled=false -Ddatahub.limits.rate.enabled=false +"""Measure the two ingest paths from Python and print a comparison.""" + +from __future__ import annotations + +import argparse +import datetime as dt +import math +import os +import resource +import time +import urllib.request + +from intellistream_datahub_sdk import DataHubClient, TimeSeries + +START = dt.datetime(2025, 1, 1, tzinfo=dt.timezone.utc) + + +def generate(count: int, offset: int, series_index: int) -> tuple[list[dt.datetime], list[float]]: + """A slow sine plus noise, one signal per series. + + Identical series would let zstd compress the repetition across them and report a wire size + no real fleet of sensors would produce. + """ + base = 150.0 + series_index * 0.7 + phase = series_index * 0.37 + timestamps = [] + values = [] + for i in range(count): + t = offset + i + timestamps.append(START + dt.timedelta(seconds=t)) + # Deterministic, so two runs generate the same bytes. + z = (t + series_index * 6364136223846793005) * 6364136223846793005 & 0xFFFFFFFFFFFFFFFF + z ^= z >> 33 + noise = ((z >> 40) / float(1 << 24)) - 0.5 + values.append(base + 20.0 * math.sin(t / 600.0 + phase) + noise) + return timestamps, values + + +def clickhouse_count(ids: list[int]) -> int: + url = ( + os.environ.get("CLICKHOUSE_URL", "http://localhost:18123") + + "/?user=" + os.environ.get("CLICKHOUSE_USER", "foobar") + + "&password=" + os.environ.get("CLICKHOUSE_PASSWORD", "changeme") + ) + db = os.environ.get("CLICKHOUSE_DB", "foo") + sql = ( + f"SELECT count() FROM {db}.datapoints_float WHERE timeseries_id IN " + f"({','.join(str(i) for i in ids)})" + ) + with urllib.request.urlopen(url, sql.encode()) as response: + return int(response.read().decode().strip()) + + +def peak_rss_mb() -> float: + # ru_maxrss is kilobytes on Linux. + return resource.getrusage(resource.RUSAGE_SELF).ru_maxrss / 1024.0 + + +def run(client, label: str, binary: bool, points: int, series: int, chunk: int) -> dict: + run_id = f"pybench_{int(time.time())}_{'bin' if binary else 'json'}" + external_ids = [f"{run_id}_{i}" for i in range(series)] + internal_ids = [] + for external_id in external_ids: + created = client.timeseries.create( + [TimeSeries(name=external_id, external_id=external_id, value_type="float", unit="celsius")] + ) + internal_ids.append(created[0].id) + + per_series = chunk // series + latencies = [] + sent = 0 + offset = 0 + started = time.perf_counter() + while sent < points: + this_chunk = min(chunk, points - sent) + per_series = max(1, this_chunk // series) + for index, external_id in enumerate(external_ids): + timestamps, values = generate(per_series, offset, index) + call = time.perf_counter() + if binary: + client.timeseries.insert_from_lists_binary(timestamps, values, external_id) + else: + client.timeseries.insert_from_lists(timestamps, values, external_id) + latencies.append(time.perf_counter() - call) + sent += per_series * series + offset += per_series + elapsed = time.perf_counter() - started + print(f" {label}: {sent:,} / {points:,} points, {sent / elapsed:,.0f} pts/s", flush=True) + ingest_seconds = time.perf_counter() - started + + settle_start = time.perf_counter() + deadline = settle_start + 1800 + while time.perf_counter() < deadline: + if clickhouse_count(internal_ids) >= sent: + break + time.sleep(1) + else: + raise SystemExit(f"{label}: only {clickhouse_count(internal_ids)} of {sent} rows became readable") + settle_seconds = time.perf_counter() - settle_start + + client.timeseries.delete(external_ids) + latencies.sort() + return { + "path": label, + "points": sent, + "ingest_seconds": ingest_seconds, + "settle_seconds": settle_seconds, + "points_per_second": sent / ingest_seconds, + # Wall clock includes generating the points in interpreted Python, which dominates it. + # This is the same points divided by the time actually spent inside the SDK calls, so + # it says what the transport did rather than what the loop above did. + "points_per_second_in_call": sent / sum(latencies), + "requests": len(latencies), + "latency_mean_ms": sum(latencies) / len(latencies) * 1000, + "latency_p50_ms": latencies[len(latencies) // 2] * 1000, + "latency_p99_ms": latencies[min(len(latencies) - 1, int(len(latencies) * 0.99))] * 1000, + "client_peak_rss_mb": peak_rss_mb(), + } + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--points", type=int, default=10_000_000) + parser.add_argument("--series", type=int, default=100) + parser.add_argument("--chunk", type=int, default=1_000_000) + args = parser.parse_args() + + env_file = os.environ.get("DATAHUB_ENV_FILE", ".env") + client = DataHubClient.from_envfile(env_file) + print(f"\n=== {args.points:,} points across {args.series:,} series, float ===", flush=True) + results = [ + run(client, "JSON", False, args.points, args.series, args.chunk), + run(client, "binary", True, args.points, args.series, args.chunk), + ] + + rows = [ + ("ingest wall time (s)", "{:.1f}", "ingest_seconds"), + ("points per second", "{:,.0f}", "points_per_second"), + ("points per second in-call", "{:,.0f}", "points_per_second_in_call"), + ("settle to readable (s)", "{:.1f}", "settle_seconds"), + ("requests", "{:,.0f}", "requests"), + ("latency mean (ms)", "{:,.0f}", "latency_mean_ms"), + ("latency p50 (ms)", "{:,.0f}", "latency_p50_ms"), + ("latency p99 (ms)", "{:,.0f}", "latency_p99_ms"), + ("client peak RSS (MB)", "{:,.0f}", "client_peak_rss_mb"), + ] + print(f"\n=== datapoint ingest from Python: JSON against binary ===") + print(f"{args.points:,} points across {args.series:,} series, float\n") + print("{:<26}{:>18}{:>18}".format("metric", *[r["path"] for r in results])) + for label, fmt, key in rows: + print("{:<26}{:>18}{:>18}".format(label, *[fmt.format(r[key]) for r in results])) + + +if __name__ == "__main__": + main() diff --git a/src/blocking.rs b/src/blocking.rs index 9a6d28f..3989aae 100644 --- a/src/blocking.rs +++ b/src/blocking.rs @@ -45,7 +45,7 @@ use crate::relations::{EdgeProxy, RelForm, RelTypeForm, RelationshipType}; use crate::resources::{ RelatedResourcesForm, Resource, ResourceFilter, ResourceNetwork, ResourceUpdate, }; -use crate::timeseries::{TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection}; +use crate::timeseries::{BinaryIngestOptions, TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection}; use crate::unit::Unit; /// Generate blocking methods that delegate to the same-named async method on one of @@ -168,6 +168,7 @@ impl TimeSeriesService { fn search_by_query(query: &str) -> Result, ResponseError>; fn insert_datapoint(id: Option, external_id: Option, timestamp: DateTime, value: String) -> Result, ResponseError>; fn insert_datapoints(json: &mut DataWrapper>) -> Result, ResponseError>; + fn insert_datapoints_binary(json: &DataWrapper>, options: &BinaryIngestOptions) -> Result, ResponseError>; fn retrieve_datapoints(json: &DataWrapper) -> Result>, ResponseError>; fn delete_datapoints(json: &DataWrapper) -> Result, ResponseError>; fn retrieve_latest_datapoint(json: &DataWrapper) -> Result>, ResponseError>; @@ -177,6 +178,11 @@ impl TimeSeriesService { pub fn buffered_count(&self) -> u64 { self.api.time_series.buffered_count() } + + /// Already synchronous on the async service; passed through directly. + pub fn evict_binary_series_cache(&self) { + self.api.time_series.evict_binary_series_cache() + } } /// Blocking counterpart of [`crate::ResourceService`]. diff --git a/src/generic.rs b/src/generic.rs index f324b6e..0a1bbbd 100644 --- a/src/generic.rs +++ b/src/generic.rs @@ -834,6 +834,44 @@ pub trait ApiServiceProvider { } } + /// `POST` a binary body under its own media type: the datapoint frames of + /// `/timeseries/data/binary`. The api's 204 becomes an empty wrapper, as in the JSON helper; + /// a rejection carries the api's `problem+json` text. + async fn execute_post_bytes_request( + &self, + path: &str, + body: Vec, + content_type: &str, + ) -> Result { + let token = self.get_token().await?; + let response = self + .get_api_service() + .http_client + .post(path) + .header(http::header::CONTENT_TYPE, content_type) + .header(http::header::ACCEPT, "application/json, application/problem+json") + .body(body) + .bearer_auth(token.clone()) + .send() + .await + .map_err(|err| { + eprintln!("HTTP request failed: {}", err); + ResponseError::from_err(err) + })?; + if response.status() == 204 { + return T::deserialize_and_set_status("", response.status().as_u16()).map_err(|err| { + ResponseError { + status: response.status(), + message: err.to_string(), + } + }); + } + match process_response::(response, path).await { + Ok(value) => Ok(value), + Err(e) => Err(self.on_request_error(e, &token).await), + } + } + /// `GET` an endpoint that answers with bytes rather than JSON (currently only /// `/files/download/{id}`). /// diff --git a/src/timeseries/binary.rs b/src/timeseries/binary.rs new file mode 100644 index 0000000..0ebea83 --- /dev/null +++ b/src/timeseries/binary.rs @@ -0,0 +1,1170 @@ +//! The binary datapoint contract behind `POST /timeseries/data/binary`. +//! +//! A request body is one or more *frames*. Each frame carries the datapoints of one value type as +//! an Arrow IPC stream in the schema the server's ClickHouse table wants, compressed with zstd by +//! the client and forwarded compressed all the way to the consumer, behind a 28-byte envelope and +//! a directory naming the series inside. The layout is the platform's `binary_datapoints_format.md`. +//! +//! [`FrameWriter`] builds one frame. [`TimeSeriesService::insert_datapoints_binary`] does the rest: +//! resolves series to their ids and value types (cached), cuts frames at the caps, compresses them +//! in parallel, packs them into requests and posts them. + +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; + +use arrow_array::{ + ArrayRef, Decimal128Array, Float32Array, Float64Array, Int64Array, RecordBatch, StringArray, + TimestampMillisecondArray, +}; +use arrow_ipc::writer::StreamWriter; +use arrow_schema::{DataType, Field, Schema, SchemaRef, TimeUnit}; +use futures::future::join_all; +use oauth2::http::StatusCode; + +use crate::generic::{ + ApiServiceProvider, DataWrapper, DatapointString, DatapointsCollection, IdAndExtId, +}; +use crate::http::ResponseError; +use crate::timeseries::TimeSeriesService; + +/// The request media type. Anything else is a 415. +pub const MEDIA_TYPE: &str = "application/vnd.intellistream.datapoint-block"; +pub const MAGIC: &[u8; 4] = b"DHDP"; +pub const VERSION: u8 = 1; +pub const CODEC_ARROW_IPC: u8 = 1; +pub const COMPRESSION_ZSTD: u8 = 1; +/// Bytes before the series directory. +pub const HEADER_BYTES: usize = 28; + +/// Decompressed bytes per frame; a frame is one Pulsar message on the server. +pub const MAX_FRAME_RAW_BYTES: usize = 4 * 1024 * 1024; +pub const MAX_FRAMES_PER_REQUEST: usize = 32; +/// Decompressed bytes per request, summed over its frames. +pub const MAX_REQUEST_RAW_BYTES: usize = 64 * 1024 * 1024; +pub const MAX_SERIES_PER_FRAME: usize = 10_000; +pub const MAX_EXTERNAL_ID_BYTES: usize = 1024; +pub const MAX_VALUE_CHARS: usize = 64; +pub const MAX_VALUE_BYTES: usize = 256; +/// Rows per frame for the numeric value types, the JSON path's per-collection cap. +pub const NUMERIC_ROWS_PER_FRAME: usize = 100_000; +/// Rows per frame for `text` and `mixed`. +pub const TEXT_ROWS_PER_FRAME: usize = 10_000; + +const NUMERIC_UNSCALED_MAX: i128 = 999_999_999_999_999_999; +const DECIMAL32_UNSCALED_MAX: i128 = 999_999_999; +/// Headroom under the raw cap: arrow-rs pads every buffer to 64 bytes, which the estimate ignores. +const RAW_BYTES_MARGIN: usize = 64 * 1024; + +/// A timeseries value type, numbered as the frame envelope carries it. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum DatapointValueType { + Bigint = 1, + Float = 2, + Numeric = 3, + Text = 4, + Decimal32 = 5, + Mixed = 6, + Float32 = 7, +} + +impl DatapointValueType { + /// The name a timeseries read reports in `valueType`, matched case-insensitively. + pub fn from_name(name: &str) -> Option { + Some(match name.to_ascii_lowercase().as_str() { + "bigint" => Self::Bigint, + "float" => Self::Float, + "numeric" => Self::Numeric, + "text" => Self::Text, + "decimal32" => Self::Decimal32, + "mixed" => Self::Mixed, + "float32" => Self::Float32, + _ => return None, + }) + } + + pub fn id(self) -> u8 { + self as u8 + } + + pub fn carries_text(self) -> bool { + matches!(self, Self::Text | Self::Mixed) + } + + /// Rows one frame of this type may carry. + pub fn max_rows(self) -> usize { + if self.carries_text() { + TEXT_ROWS_PER_FRAME + } else { + NUMERIC_ROWS_PER_FRAME + } + } + + /// The one Arrow schema a frame of this type is accepted in. + pub fn schema(self) -> SchemaRef { + let id = Field::new("timeseries_id", DataType::Int64, false); + let timestamp = Field::new( + "timestamp", + DataType::Timestamp(TimeUnit::Millisecond, Some("UTC".into())), + false, + ); + let fields = match self { + Self::Bigint => vec![id, timestamp, Field::new("value", DataType::Int64, false)], + Self::Float => vec![id, timestamp, Field::new("value", DataType::Float64, false)], + Self::Float32 => vec![id, timestamp, Field::new("value", DataType::Float32, false)], + Self::Numeric => vec![id, timestamp, Field::new("value", DataType::Decimal128(18, 6), false)], + Self::Decimal32 => vec![id, timestamp, Field::new("value", DataType::Decimal128(9, 4), false)], + Self::Text => vec![id, timestamp, Field::new("value", DataType::Utf8, false)], + Self::Mixed => vec![ + id, + timestamp, + Field::new("value_numeric", DataType::Float64, true), + Field::new("value_text", DataType::Utf8, true), + ], + }; + Arc::new(Schema::new(fields)) + } + + fn bytes_per_row(self) -> usize { + match self { + Self::Bigint | Self::Float => 24, + Self::Float32 => 20, + Self::Numeric | Self::Decimal32 => 32, + Self::Text => 20, + Self::Mixed => 28, + } + } +} + +/// How [`TimeSeriesService::insert_datapoints_binary`] compresses and retries. +#[derive(Debug, Clone)] +pub struct BinaryIngestOptions { + /// zstd level per frame: 1, 3 or 9. Compression is mandatory, so there is no way to turn it off. + pub zstd_level: i32, + /// Retries of one request after a 429 or a 5xx, one second apart per attempt. + pub max_retries: u32, + /// Requests in flight at once when one call packs into more than one, which happens above + /// 3.2 million points. Below that a call is a single request and this has no effect. + pub request_concurrency: usize, +} + +impl Default for BinaryIngestOptions { + fn default() -> Self { + BinaryIngestOptions { + zstd_level: 9, + max_retries: 3, + request_concurrency: 4, + } + } +} + +impl BinaryIngestOptions { + pub fn new() -> Self { + Self::default() + } + + pub fn zstd_level(mut self, level: i32) -> Self { + self.zstd_level = level; + self + } + + pub fn max_retries(mut self, retries: u32) -> Self { + self.max_retries = retries; + self + } + + pub fn request_concurrency(mut self, requests: usize) -> Self { + self.request_concurrency = requests; + self + } + + fn validate(&self) -> Result<(), ResponseError> { + if ![1, 3, 9].contains(&self.zstd_level) { + return Err(ResponseError::bad_request(format!( + "zstd level {} is not one of 1, 3 or 9", + self.zstd_level + ))); + } + Ok(()) + } +} + +/// A series as the binary path needs it: its id, the external id the frame directory names it by, +/// and the value type that picks the frame it goes into. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ResolvedSeries { + pub id: u64, + pub external_id: String, + pub value_type: DatapointValueType, +} + +/// One finished frame, ready to post as is or packed with others. +#[derive(Debug, Clone)] +pub struct Frame { + pub bytes: Vec, + pub rows: usize, + pub series: usize, + /// Payload bytes before compression; what the request cap counts. + pub raw_len: usize, + /// Decimal32 values that were clamped to the type's range. + pub clamped: usize, +} + +#[derive(Debug, Clone)] +enum Cell { + I64(i64), + F64(f64), + F32(f32), + Scaled(i128), + Text(String), + MixedNumeric(f64), + MixedText(String), +} + +#[derive(Debug, Clone)] +struct Row { + id: i64, + timestamp: i64, + cell: Cell, +} + +/// Builds one frame for one value type: rows in any order in, a sorted, de-duplicated, compressed +/// frame out. Values arrive in the JSON contract's string form and are checked the way the JSON +/// path checks them on the server, so both paths store the same bytes for the same input. +#[derive(Debug, Clone)] +pub struct FrameWriter { + value_type: DatapointValueType, + series: HashMap, + rows: Vec, + text_bytes: usize, + clamped: usize, +} + +impl FrameWriter { + pub fn new(value_type: DatapointValueType) -> Self { + FrameWriter { + value_type, + series: HashMap::new(), + rows: Vec::new(), + text_bytes: 0, + clamped: 0, + } + } + + pub fn value_type(&self) -> DatapointValueType { + self.value_type + } + + /// Names a series the frame will carry; every id added must be named before [`build`](Self::build). + pub fn series(&mut self, id: u64, external_id: &str) -> Result<(), String> { + if self.series.contains_key(&(id as i64)) { + return Ok(()); + } + if external_id.is_empty() { + return Err(format!("external id for series {id} is empty")); + } + if external_id.len() > MAX_EXTERNAL_ID_BYTES { + return Err(format!( + "external id for series {id} exceeds {MAX_EXTERNAL_ID_BYTES} bytes" + )); + } + self.series.insert(id as i64, external_id.to_string()); + Ok(()) + } + + pub fn has_series(&self, id: u64) -> bool { + self.series.contains_key(&(id as i64)) + } + + pub fn row_count(&self) -> usize { + self.rows.len() + } + + pub fn series_count(&self) -> usize { + self.series.len() + } + + /// Decimal32 values clamped so far, for the caller to warn about. + pub fn clamped_count(&self) -> usize { + self.clamped + } + + /// Payload bytes the frame would need before compression, for cutting frames before building. + pub fn estimated_raw_bytes(&self) -> usize { + 2048 + self.value_type.bytes_per_row() * self.rows.len() + self.text_bytes + } + + /// Whether another row for `id` would push the frame past a cap. + pub fn is_full_for(&self, id: u64) -> bool { + self.rows.len() >= self.value_type.max_rows() + || self.estimated_raw_bytes() + MAX_VALUE_BYTES + RAW_BYTES_MARGIN > MAX_FRAME_RAW_BYTES + || (!self.has_series(id) && self.series.len() >= MAX_SERIES_PER_FRAME) + } + + /// Adds a value in the JSON contract's string form, parsed by the frame's value type. + pub fn add(&mut self, id: u64, timestamp_ms: i64, value: &str) -> Result<(), String> { + let cell = match self.value_type { + DatapointValueType::Bigint => Cell::I64( + value + .parse::() + .map_err(|_| format!("{value:?} is not a bigint"))?, + ), + DatapointValueType::Float => Cell::F64( + value + .parse::() + .map_err(|_| format!("{value:?} is not a float"))?, + ), + DatapointValueType::Float32 => Cell::F32( + value + .parse::() + .map_err(|_| format!("{value:?} is not a float32"))?, + ), + DatapointValueType::Numeric => { + let unscaled = parse_scaled(value, 6) + .ok_or_else(|| format!("{value:?} is not a decimal number"))?; + if unscaled > NUMERIC_UNSCALED_MAX || unscaled < -NUMERIC_UNSCALED_MAX { + return Err(format!("numeric value {value} does not fit Decimal(18, 6)")); + } + Cell::Scaled(unscaled) + } + DatapointValueType::Decimal32 => { + let unscaled = parse_scaled(value, 4) + .ok_or_else(|| format!("{value:?} is not a decimal number"))?; + let bounded = unscaled.clamp(-DECIMAL32_UNSCALED_MAX, DECIMAL32_UNSCALED_MAX); + if bounded != unscaled { + self.clamped += 1; + } + Cell::Scaled(bounded) + } + DatapointValueType::Text => Cell::Text(checked_text(value)?), + DatapointValueType::Mixed => match decimal_parts(value) { + Some(_) => Cell::MixedNumeric( + value + .parse::() + .map_err(|_| format!("{value:?} is not a float"))?, + ), + None => Cell::MixedText(checked_text(value)?), + }, + }; + self.push(id, timestamp_ms, cell); + Ok(()) + } + + pub fn add_bigint(&mut self, id: u64, timestamp_ms: i64, value: i64) -> Result<(), String> { + self.expect(DatapointValueType::Bigint)?; + self.push(id, timestamp_ms, Cell::I64(value)); + Ok(()) + } + + pub fn add_float(&mut self, id: u64, timestamp_ms: i64, value: f64) -> Result<(), String> { + self.expect(DatapointValueType::Float)?; + self.push(id, timestamp_ms, Cell::F64(value)); + Ok(()) + } + + pub fn add_float32(&mut self, id: u64, timestamp_ms: i64, value: f32) -> Result<(), String> { + self.expect(DatapointValueType::Float32)?; + self.push(id, timestamp_ms, Cell::F32(value)); + Ok(()) + } + + pub fn add_text(&mut self, id: u64, timestamp_ms: i64, value: &str) -> Result<(), String> { + self.expect(DatapointValueType::Text)?; + let text = checked_text(value)?; + self.push(id, timestamp_ms, Cell::Text(text)); + Ok(()) + } + + fn expect(&self, expected: DatapointValueType) -> Result<(), String> { + if self.value_type != expected { + return Err(format!( + "writer is for {:?}, not {:?}", + self.value_type, expected + )); + } + Ok(()) + } + + fn push(&mut self, id: u64, timestamp: i64, cell: Cell) { + if let Cell::Text(t) | Cell::MixedText(t) = &cell { + self.text_bytes += t.len(); + } + self.rows.push(Row { + id: id as i64, + timestamp, + cell, + }); + } + + fn key(&self, index: usize) -> (i64, i64) { + let row = &self.rows[index]; + (row.id, row.timestamp) + } + + /// Sorts by (id, timestamp), keeps the last value of a repeated pair, checks the caps, writes + /// the Arrow stream, compresses it and returns the complete frame. + pub fn build(&self, zstd_level: i32) -> Result { + if self.rows.is_empty() { + return Err("frame has no rows".to_string()); + } + let mut order: Vec = (0..self.rows.len()).collect(); + order.sort_by(|&a, &b| self.key(a).cmp(&self.key(b))); + let mut kept = Vec::with_capacity(order.len()); + for (i, &index) in order.iter().enumerate() { + let last_of_key = i + 1 == order.len() || self.key(index) != self.key(order[i + 1]); + if last_of_key { + kept.push(index); + } + } + let rows = kept.len(); + let max_rows = self.value_type.max_rows(); + if rows > max_rows { + return Err(format!( + "frame has {rows} rows, the cap for {:?} is {max_rows}; split it", + self.value_type + )); + } + let mut directory: Vec<(i64, &str)> = Vec::new(); + for &index in &kept { + let id = self.rows[index].id; + if directory.last().map_or(true, |(last, _)| *last != id) { + let external_id = self + .series + .get(&id) + .ok_or_else(|| format!("no external id registered for series {id}"))?; + directory.push((id, external_id.as_str())); + } + } + if directory.len() > MAX_SERIES_PER_FRAME { + return Err(format!( + "frame has {} series, the cap is {MAX_SERIES_PER_FRAME}; split it", + directory.len() + )); + } + + let raw = self.write_ipc(&kept)?; + if raw.len() > MAX_FRAME_RAW_BYTES { + return Err(format!( + "frame payload is {} bytes, the cap is {MAX_FRAME_RAW_BYTES}; split it", + raw.len() + )); + } + let compressed = + zstd::bulk::compress(&raw, zstd_level).map_err(|e| format!("zstd failed: {e}"))?; + + let directory_len: usize = directory + .iter() + .map(|(_, external_id)| 8 + varint_size(external_id.len()) + external_id.len()) + .sum(); + let mut frame = Vec::with_capacity(HEADER_BYTES + directory_len + compressed.len()); + frame.extend_from_slice(MAGIC); + frame.push(VERSION); + frame.push(self.value_type.id()); + frame.push(CODEC_ARROW_IPC); + frame.push(COMPRESSION_ZSTD); + for n in [rows, directory.len(), directory_len, compressed.len(), raw.len()] { + frame.extend_from_slice(&(n as u32).to_le_bytes()); + } + for (id, external_id) in &directory { + frame.extend_from_slice(&id.to_le_bytes()); + write_varint(&mut frame, external_id.len()); + frame.extend_from_slice(external_id.as_bytes()); + } + frame.extend_from_slice(&compressed); + Ok(Frame { + bytes: frame, + rows, + series: directory.len(), + raw_len: raw.len(), + clamped: self.clamped, + }) + } + + fn write_ipc(&self, kept: &[usize]) -> Result, String> { + let ids: Vec = kept.iter().map(|&i| self.rows[i].id).collect(); + let timestamps: Vec = kept.iter().map(|&i| self.rows[i].timestamp).collect(); + let mut columns: Vec = vec![ + Arc::new(Int64Array::from(ids)), + Arc::new(TimestampMillisecondArray::from(timestamps).with_timezone("UTC")), + ]; + let cells = kept.iter().map(|&i| &self.rows[i].cell); + match self.value_type { + DatapointValueType::Bigint => { + let values: Vec = cells + .map(|c| match c { + Cell::I64(v) => *v, + _ => unreachable!("bigint frame holds bigint cells"), + }) + .collect(); + columns.push(Arc::new(Int64Array::from(values))); + } + DatapointValueType::Float => { + let values: Vec = cells + .map(|c| match c { + Cell::F64(v) => *v, + _ => unreachable!("float frame holds float cells"), + }) + .collect(); + columns.push(Arc::new(Float64Array::from(values))); + } + DatapointValueType::Float32 => { + let values: Vec = cells + .map(|c| match c { + Cell::F32(v) => *v, + _ => unreachable!("float32 frame holds float32 cells"), + }) + .collect(); + columns.push(Arc::new(Float32Array::from(values))); + } + DatapointValueType::Numeric | DatapointValueType::Decimal32 => { + let values: Vec = cells + .map(|c| match c { + Cell::Scaled(v) => *v, + _ => unreachable!("decimal frame holds scaled cells"), + }) + .collect(); + let (precision, scale) = if self.value_type == DatapointValueType::Numeric { + (18, 6) + } else { + (9, 4) + }; + let array = Decimal128Array::from(values) + .with_precision_and_scale(precision, scale) + .map_err(|e| e.to_string())?; + columns.push(Arc::new(array)); + } + DatapointValueType::Text => { + let values: Vec<&str> = cells + .map(|c| match c { + Cell::Text(v) => v.as_str(), + _ => unreachable!("text frame holds text cells"), + }) + .collect(); + columns.push(Arc::new(StringArray::from(values))); + } + DatapointValueType::Mixed => { + let mut numeric: Vec> = Vec::with_capacity(kept.len()); + let mut text: Vec> = Vec::with_capacity(kept.len()); + for cell in cells { + match cell { + Cell::MixedNumeric(v) => { + numeric.push(Some(*v)); + text.push(None); + } + Cell::MixedText(v) => { + numeric.push(None); + text.push(Some(v.as_str())); + } + _ => unreachable!("mixed frame holds mixed cells"), + } + } + columns.push(Arc::new(Float64Array::from(numeric))); + columns.push(Arc::new(StringArray::from(text))); + } + } + let schema = self.value_type.schema(); + let batch = RecordBatch::try_new(schema.clone(), columns).map_err(|e| e.to_string())?; + let mut raw = Vec::with_capacity(self.estimated_raw_bytes()); + let mut writer = + StreamWriter::try_new(&mut raw, schema.as_ref()).map_err(|e| e.to_string())?; + writer.write(&batch).map_err(|e| e.to_string())?; + writer.finish().map_err(|e| e.to_string())?; + drop(writer); + Ok(raw) + } +} + +/// The server counts the character limit in UTF-16 units, so this does too. +fn checked_text(value: &str) -> Result { + if value.is_empty() { + return Err("text value is empty".to_string()); + } + let chars = value.encode_utf16().count(); + if chars > MAX_VALUE_CHARS { + return Err(format!( + "text value is {chars} characters, allowed {MAX_VALUE_CHARS}" + )); + } + if value.len() > MAX_VALUE_BYTES { + return Err(format!( + "text value is {} bytes, allowed {MAX_VALUE_BYTES}", + value.len() + )); + } + Ok(value.to_string()) +} + +/// Splits a decimal literal (`[+-]digits[.digits][e[+-]digits]`, the grammar `BigDecimal` takes) +/// into its sign, its digit string and the number of those digits before the decimal point. +fn decimal_parts(value: &str) -> Option<(bool, String, i32)> { + let (negative, rest) = match value.as_bytes().first()? { + b'-' => (true, &value[1..]), + b'+' => (false, &value[1..]), + _ => (false, value), + }; + let (mantissa, exponent) = match rest.find(['e', 'E']) { + Some(at) => (&rest[..at], rest[at + 1..].parse::().ok()?), + None => (rest, 0), + }; + let (integer, fraction) = match mantissa.find('.') { + Some(at) => (&mantissa[..at], &mantissa[at + 1..]), + None => (mantissa, ""), + }; + if integer.is_empty() && fraction.is_empty() { + return None; + } + if !integer.bytes().all(|b| b.is_ascii_digit()) || !fraction.bytes().all(|b| b.is_ascii_digit()) { + return None; + } + let point = i32::try_from(integer.len()).ok()?.checked_add(exponent)?; + Some((negative, format!("{integer}{fraction}"), point)) +} + +/// Parses a decimal literal to its unscaled value at `scale` decimals, rounded half-up. +fn parse_scaled(value: &str, scale: u32) -> Option { + let (negative, digits, point) = decimal_parts(value)?; + let total = i32::try_from(digits.len()).ok()?; + // Digits to keep: everything before the point plus `scale` after it. + let keep = point.checked_add(scale as i32)?; + let magnitude: i128 = if keep >= total { + let shift = u32::try_from(keep - total).ok()?; + if digits.len() as u32 + shift > 38 { + return None; + } + digits.parse::().ok()?.checked_mul(10i128.checked_pow(shift)?)? + } else if keep < 0 { + 0 + } else { + let keep = keep as usize; + if keep > 38 { + return None; + } + let head: i128 = if keep == 0 { 0 } else { digits[..keep].parse().ok()? }; + let round_up = digits.as_bytes()[keep] >= b'5'; + if round_up { + head + 1 + } else { + head + } + }; + Some(if negative { -magnitude } else { magnitude }) +} + +fn varint_size(mut value: usize) -> usize { + let mut n = 1; + while value >= 0x80 { + value >>= 7; + n += 1; + } + n +} + +fn write_varint(out: &mut Vec, mut value: usize) { + while value >= 0x80 { + out.push((value as u8 & 0x7F) | 0x80); + value >>= 7; + } + out.push(value as u8); +} + +/// Series the binary path has resolved, per service instance. Dropped as a whole when the server +/// says a series is unknown or renamed, so the next call re-reads every series it names. +#[derive(Debug, Default)] +pub(crate) struct SeriesCache { + by_external_id: HashMap, + by_id: HashMap, +} + +/// Groups the collections' datapoints into writers, one value type at a time, cutting a new +/// writer whenever the open one would pass a cap. +pub(crate) fn cut_into_writers( + items: &[(ResolvedSeries, &[DatapointString])], +) -> Result, ResponseError> { + let mut open: HashMap = HashMap::new(); + let mut done = Vec::new(); + for (series, datapoints) in items { + let value_type = series.value_type; + for dp in datapoints.iter() { + let timestamp = dp.timestamp.parse::().map_err(|_| { + ResponseError::bad_request(format!( + "timestamp {:?} of series {} is not epoch milliseconds", + dp.timestamp, series.external_id + )) + })?; + let writer = open + .entry(value_type) + .or_insert_with(|| FrameWriter::new(value_type)); + if writer.is_full_for(series.id) { + let full = std::mem::replace(writer, FrameWriter::new(value_type)); + done.push(full); + } + writer + .series(series.id, &series.external_id) + .map_err(unprocessable)?; + writer + .add(series.id, timestamp, &dp.value) + .map_err(|message| unprocessable(format!("series {}: {message}", series.external_id)))?; + } + } + done.extend(open.into_values().filter(|w| w.row_count() > 0)); + Ok(done) +} + +/// Concatenates frames into request bodies under the per-request caps, in order. +pub(crate) fn pack_requests(frames: Vec) -> Vec> { + let mut requests = Vec::new(); + let mut body = Vec::new(); + let mut count = 0; + let mut raw = 0; + for frame in frames { + if count > 0 + && (count == MAX_FRAMES_PER_REQUEST || raw + frame.raw_len > MAX_REQUEST_RAW_BYTES) + { + requests.push(std::mem::take(&mut body)); + count = 0; + raw = 0; + } + body.extend_from_slice(&frame.bytes); + count += 1; + raw += frame.raw_len; + } + if count > 0 { + requests.push(body); + } + requests +} + +fn unprocessable(message: String) -> ResponseError { + ResponseError { + status: StatusCode::UNPROCESSABLE_ENTITY, + message, + } +} + +/// 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")) +} + +impl TimeSeriesService { + /// Inserts datapoints through `POST /timeseries/data/binary`, the high-throughput path. + /// + /// The same input as [`insert_datapoints`](Self::insert_datapoints), sent as zstd-compressed + /// Arrow frames instead of JSON: each series is resolved once to its id and value type (through + /// `/timeseries/byids`, which needs read access on its dataset, and cached on this service), + /// the values are checked against that type here, sorted and de-duplicated per series, cut into + /// frames at the contract's caps, compressed in parallel and posted in as many requests as the + /// caps allow. A 204 means every frame was accepted; a request is all-or-nothing on the server. + /// + /// Errors before any request: a series that does not exist, or that the caller cannot read, is + /// a 404 naming it; a value that does not fit its series' type, or a text value over the + /// limits, is a 422. A 429 or a 5xx is retried per [`BinaryIngestOptions::max_retries`]. A + /// request the server refuses because a series was deleted or renamed after it was cached is + /// rebuilt once after re-resolving. The durable spool does not apply to this path. + pub async fn insert_datapoints_binary( + &self, + json: &DataWrapper>, + options: &BinaryIngestOptions, + ) -> Result, ResponseError> { + options.validate()?; + match self.insert_binary_once(json, options).await { + Err(error) if is_stale_series_rejection(&error) => { + self.evict_binary_series_cache(); + self.insert_binary_once(json, options).await + } + other => other, + } + } + + /// Forgets every series the binary path has resolved, so the next call re-reads them. + pub fn evict_binary_series_cache(&self) { + let mut cache = self.binary_series.lock().unwrap(); + cache.by_external_id.clear(); + cache.by_id.clear(); + } + + async fn insert_binary_once( + &self, + json: &DataWrapper>, + options: &BinaryIngestOptions, + ) -> Result, ResponseError> { + let collections = json.get_items(); + let resolved = self.resolve_series(collections).await?; + let items: Vec<(ResolvedSeries, &[DatapointString])> = resolved + .into_iter() + .zip(collections.iter()) + .map(|(series, collection)| (series, collection.datapoints.as_slice())) + .collect(); + let writers = cut_into_writers(&items)?; + + let level = options.zstd_level; + let builds = writers + .into_iter() + .map(|writer| tokio::task::spawn_blocking(move || writer.build(level))); + let mut frames = Vec::new(); + for built in join_all(builds).await { + let frame = built + .map_err(|e| ResponseError::bad_request(format!("frame build failed: {e}")))? + .map_err(unprocessable)?; + frames.push(frame); + } + + // Concurrently, not one after another. A call of up to 3.2M points packs into a single + // request and this changes nothing, but a larger one becomes several and the api + // validates them independently, so there is no reason to serialise them. Bounded by + // `request_concurrency` so a very large call cannot open an unbounded number of + // connections. + let path = format!("{}/data/binary", self.base_url); + let bodies = pack_requests(frames); + let limit = options.request_concurrency.max(1); + for window in bodies.chunks(limit) { + let sends = window + .iter() + .map(|body| self.post_frames(&path, body.clone(), options.max_retries)); + for outcome in join_all(sends).await { + outcome?; + } + } + let mut result = DataWrapper::new(); + result.set_http_status_code(204); + Ok(result) + } + + async fn post_frames( + &self, + path: &str, + body: Vec, + max_retries: u32, + ) -> Result<(), ResponseError> { + let mut attempt = 0u32; + loop { + match self + .execute_post_bytes_request::>(path, body.clone(), MEDIA_TYPE) + .await + { + Ok(_) => return Ok(()), + Err(error) + if attempt < max_retries + && (error.get_status() == StatusCode::TOO_MANY_REQUESTS + || error.get_status().is_server_error()) => + { + attempt += 1; + tokio::time::sleep(std::time::Duration::from_secs(attempt as u64)).await; + } + Err(error) => return Err(error), + } + } + } + + /// One [`ResolvedSeries`] per collection, in order, from the cache or `/timeseries/byids`. + async fn resolve_series( + &self, + collections: &[DatapointsCollection], + ) -> Result, ResponseError> { + let mut lookups: Vec = Vec::new(); + let mut seen: HashSet = HashSet::new(); + { + let cache = self.binary_series.lock().unwrap(); + for collection in collections { + match (collection.id, &collection.external_id) { + (Some(id), _) => { + if !cache.by_id.contains_key(&id) && seen.insert(format!("id:{id}")) { + lookups.push(IdAndExtId::from_id(id)); + } + } + (None, Some(external_id)) => { + if !cache.by_external_id.contains_key(external_id) + && seen.insert(format!("ext:{external_id}")) + { + lookups.push(IdAndExtId::from_external_id(external_id)); + } + } + (None, None) => { + return Err(ResponseError::bad_request( + "a datapoint collection names neither id nor externalId".to_string(), + )) + } + } + } + } + for chunk in lookups.chunks(MAX_SERIES_PER_FRAME) { + let found = self.by_ids(&DataWrapper::from_vec(chunk.to_vec())).await?; + let mut cache = self.binary_series.lock().unwrap(); + for ts in found.get_items() { + let value_type = ts.value_type.as_deref().and_then(DatapointValueType::from_name); + if let (Some(id), Some(value_type)) = (ts.id, value_type) { + let series = ResolvedSeries { + id, + external_id: ts.external_id.clone(), + value_type, + }; + cache.by_external_id.insert(series.external_id.clone(), series.clone()); + cache.by_id.insert(id, series); + } + } + } + + let cache = self.binary_series.lock().unwrap(); + let mut resolved = Vec::with_capacity(collections.len()); + let mut missing = Vec::new(); + for collection in collections { + let hit = match (collection.id, &collection.external_id) { + (Some(id), _) => cache.by_id.get(&id), + (None, Some(external_id)) => cache.by_external_id.get(external_id), + (None, None) => None, + }; + match hit { + Some(series) => resolved.push(series.clone()), + None => missing.push( + collection + .id + .map(|id| id.to_string()) + .or_else(|| collection.external_id.clone()) + .unwrap_or_default(), + ), + } + } + if !missing.is_empty() { + return Err(ResponseError { + status: StatusCode::NOT_FOUND, + message: format!("Could not find following timeseries: {}", missing.join(", ")), + }); + } + Ok(resolved) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow_array::Array; + use arrow_ipc::reader::StreamReader; + use std::io::Cursor; + + fn u32_at(frame: &[u8], offset: usize) -> u32 { + u32::from_le_bytes(frame[offset..offset + 4].try_into().unwrap()) + } + + /// Decompresses the payload and reads the stream back with arrow-rs. + fn read_back(frame: &Frame) -> (Vec, SchemaRef) { + let bytes = &frame.bytes; + let directory_len = u32_at(bytes, 16) as usize; + let payload_len = u32_at(bytes, 20) as usize; + let raw_len = u32_at(bytes, 24) as usize; + let payload = &bytes[HEADER_BYTES + directory_len..]; + assert_eq!(payload.len(), payload_len); + let raw = zstd::bulk::decompress(payload, raw_len).unwrap(); + assert_eq!(raw.len(), raw_len); + let reader = StreamReader::try_new(Cursor::new(raw), None).unwrap(); + let schema = reader.schema(); + let batches: Vec = reader.map(|b| b.unwrap()).collect(); + (batches, schema) + } + + #[test] + fn float_frame_is_sorted_deduplicated_and_readable_by_arrow() { + let mut writer = FrameWriter::new(DatapointValueType::Float); + writer.series(7, "pump_b").unwrap(); + writer.series(3, "pump_a").unwrap(); + writer.add(7, 2_000, "2.5").unwrap(); + writer.add(3, 1_000, "1.0").unwrap(); + writer.add(7, 1_000, "2.0").unwrap(); + writer.add(3, 1_000, "1.5").unwrap(); // same key as the second row: the last one wins + let frame = writer.build(3).unwrap(); + + let bytes = &frame.bytes; + assert_eq!(&bytes[0..4], MAGIC); + assert_eq!(bytes[4], VERSION); + assert_eq!(bytes[5], DatapointValueType::Float.id()); + assert_eq!(bytes[6], CODEC_ARROW_IPC); + assert_eq!(bytes[7], COMPRESSION_ZSTD); + assert_eq!(u32_at(bytes, 8), 3, "rows after dedupe"); + assert_eq!(u32_at(bytes, 12), 2, "series"); + assert_eq!(frame.rows, 3); + assert_eq!(frame.series, 2); + + // Directory: ascending ids, each an i64 then a varint length and the external id. + let directory_len = u32_at(bytes, 16) as usize; + let directory = &bytes[HEADER_BYTES..HEADER_BYTES + directory_len]; + let mut expected = Vec::new(); + expected.extend_from_slice(&3i64.to_le_bytes()); + expected.push(6); + expected.extend_from_slice(b"pump_a"); + expected.extend_from_slice(&7i64.to_le_bytes()); + expected.push(6); + expected.extend_from_slice(b"pump_b"); + assert_eq!(directory, expected.as_slice()); + + let (batches, schema) = read_back(&frame); + assert_eq!(schema, DatapointValueType::Float.schema()); + assert_eq!(batches.len(), 1); + let batch = &batches[0]; + let ids = batch.column(0).as_any().downcast_ref::().unwrap(); + let ts = batch.column(1).as_any().downcast_ref::().unwrap(); + let values = batch.column(2).as_any().downcast_ref::().unwrap(); + assert_eq!(ids.values(), &[3, 7, 7]); + assert_eq!(ts.values(), &[1_000, 1_000, 2_000]); + assert_eq!(values.values(), &[1.5, 2.0, 2.5]); + } + + #[test] + fn every_value_type_writes_its_canonical_schema() { + for value_type in [ + DatapointValueType::Bigint, + DatapointValueType::Float, + DatapointValueType::Float32, + DatapointValueType::Numeric, + DatapointValueType::Decimal32, + DatapointValueType::Text, + DatapointValueType::Mixed, + ] { + let mut writer = FrameWriter::new(value_type); + writer.series(1, "s").unwrap(); + let value = match value_type { + DatapointValueType::Text => "on", + DatapointValueType::Mixed => "off", + _ => "42", + }; + writer.add(1, 0, value).unwrap(); + let frame = writer.build(1).unwrap(); + assert_eq!(frame.bytes[5], value_type.id()); + let (batches, schema) = read_back(&frame); + assert_eq!(schema, value_type.schema(), "{value_type:?}"); + assert_eq!(batches[0].num_rows(), 1); + } + } + + #[test] + fn decimals_are_scaled_half_up_and_decimal32_is_clamped() { + assert_eq!(parse_scaled("123.456789", 6), Some(123_456_789)); + assert_eq!(parse_scaled("0.0000005", 6), Some(1)); + assert_eq!(parse_scaled("0.0000004", 6), Some(0)); + assert_eq!(parse_scaled("-1.5", 0), Some(-2)); + assert_eq!(parse_scaled("1e3", 2), Some(100_000)); + assert_eq!(parse_scaled("12.34E-1", 4), Some(12_340)); + assert_eq!(parse_scaled(".5", 1), Some(5)); + assert_eq!(parse_scaled("5.", 1), Some(50)); + assert_eq!(parse_scaled("", 6), None); + assert_eq!(parse_scaled("NaN", 6), None); + assert_eq!(parse_scaled("1.2.3", 6), None); + + let mut numeric = FrameWriter::new(DatapointValueType::Numeric); + numeric.series(1, "s").unwrap(); + assert!(numeric.add(1, 0, "1000000000000.5").is_err(), "outside Decimal(18, 6)"); + + let mut decimal32 = FrameWriter::new(DatapointValueType::Decimal32); + decimal32.series(1, "s").unwrap(); + decimal32.add(1, 0, "123456.78").unwrap(); + decimal32.add(1, 1, "-99999.99995").unwrap(); + decimal32.add(1, 2, "99999.9999").unwrap(); + assert_eq!(decimal32.clamped_count(), 2); + let frame = decimal32.build(1).unwrap(); + let (batches, _) = read_back(&frame); + let values = batches[0].column(2).as_any().downcast_ref::().unwrap(); + assert_eq!(values.values(), &[999_999_999, -999_999_999, 999_999_999]); + } + + #[test] + fn mixed_rows_set_exactly_one_side() { + let mut writer = FrameWriter::new(DatapointValueType::Mixed); + writer.series(1, "s").unwrap(); + writer.add(1, 0, "12.5").unwrap(); + writer.add(1, 1, "open").unwrap(); + writer.add(1, 2, "-3e2").unwrap(); + writer.add(1, 3, "NaN").unwrap(); // not a decimal literal, so text + let frame = writer.build(1).unwrap(); + let (batches, _) = read_back(&frame); + let numeric = batches[0].column(2).as_any().downcast_ref::().unwrap(); + let text = batches[0].column(3).as_any().downcast_ref::().unwrap(); + assert_eq!(numeric.null_count(), 2); + assert_eq!(text.null_count(), 2); + assert_eq!(numeric.value(0), 12.5); + assert!(numeric.is_null(1)); + assert_eq!(text.value(1), "open"); + assert_eq!(numeric.value(2), -300.0); + assert_eq!(text.value(3), "NaN"); + } + + #[test] + fn text_limits_and_series_registration_are_enforced() { + let mut writer = FrameWriter::new(DatapointValueType::Text); + writer.series(1, "s").unwrap(); + assert!(writer.add(1, 0, "").is_err()); + assert!(writer.add(1, 0, &"x".repeat(65)).is_err()); + // 33 astral characters are 66 UTF-16 units, over the limit as the server counts it. + assert!(writer.add(1, 0, &"\u{1F600}".repeat(33)).unwrap_err().contains("characters")); + writer.add(1, 0, &"\u{20ac}".repeat(64)).unwrap(); + writer.add(1, 0, &"x".repeat(64)).unwrap(); + writer.add(2, 0, "unregistered").unwrap(); + assert!(writer.build(1).unwrap_err().contains("series 2")); + + let mut typed = FrameWriter::new(DatapointValueType::Float); + assert!(typed.add_bigint(1, 0, 1).is_err()); + assert!(typed.add(1, 0, "abc").is_err()); + assert!(FrameWriter::new(DatapointValueType::Float).build(1).is_err(), "no rows"); + } + + #[test] + fn frames_are_cut_at_the_row_cap_and_packed_under_the_request_cap() { + let series = ResolvedSeries { + id: 5, + external_id: "s".to_string(), + value_type: DatapointValueType::Float, + }; + let datapoints: Vec = (0..250_000) + .map(|i| DatapointString::new(&i.to_string(), "1.0")) + .collect(); + let writers = cut_into_writers(&[(series.clone(), datapoints.as_slice())]).unwrap(); + let rows: Vec = writers.iter().map(FrameWriter::row_count).collect(); + assert_eq!(rows, vec![100_000, 100_000, 50_000]); + + let bad = vec![DatapointString::new("2025-01-01T00:00:00Z", "1.0")]; + let error = cut_into_writers(&[(series.clone(), bad.as_slice())]).unwrap_err(); + assert_eq!(error.get_status(), StatusCode::BAD_REQUEST); + + let unfit = vec![DatapointString::new("0", "warm")]; + let error = cut_into_writers(&[(series, unfit.as_slice())]).unwrap_err(); + assert_eq!(error.get_status(), StatusCode::UNPROCESSABLE_ENTITY); + + let frame = |raw_len: usize| Frame { + bytes: vec![0xAB], + rows: 1, + series: 1, + raw_len, + clamped: 0, + }; + let by_count = pack_requests((0..33).map(|_| frame(1)).collect()); + assert_eq!(by_count.iter().map(Vec::len).collect::>(), vec![32, 1]); + let by_bytes = pack_requests(vec![ + frame(MAX_REQUEST_RAW_BYTES / 2), + frame(MAX_REQUEST_RAW_BYTES / 2), + frame(1), + ]); + assert_eq!(by_bytes.iter().map(Vec::len).collect::>(), vec![2, 1]); + } + + #[test] + fn varints_are_unsigned_leb128() { + let mut out = Vec::new(); + write_varint(&mut out, 0); + write_varint(&mut out, 127); + write_varint(&mut out, 128); + write_varint(&mut out, 300); + assert_eq!(out, vec![0, 127, 0x80, 0x01, 0xAC, 0x02]); + assert_eq!(varint_size(127), 1); + assert_eq!(varint_size(128), 2); + assert_eq!(varint_size(16_384), 3); + } + + #[test] + fn options_accept_only_the_three_levels() { + assert!(BinaryIngestOptions::default().validate().is_ok()); + assert_eq!(BinaryIngestOptions::default().zstd_level, 9); + assert!(BinaryIngestOptions::new().zstd_level(3).validate().is_ok()); + assert!(BinaryIngestOptions::new().zstd_level(0).validate().is_err()); + assert!(BinaryIngestOptions::new().zstd_level(22).validate().is_err()); + } +} diff --git a/src/timeseries/mod.rs b/src/timeseries/mod.rs index d71fe66..6a199f2 100644 --- a/src/timeseries/mod.rs +++ b/src/timeseries/mod.rs @@ -1,5 +1,8 @@ +pub mod binary; mod test; +pub use binary::{BinaryIngestOptions, DatapointValueType, Frame, FrameWriter, ResolvedSeries}; + use crate::buffer::DurableSpool; use crate::datahub::DataHubConfig; use crate::fields::{Field, MapField}; @@ -35,6 +38,8 @@ pub struct TimeSeriesService { base_url: String, // Durable spool for datapoint ingestion (lazily opened on first buffered send; None if off). spool: Mutex>, + // Series the binary path has resolved to id and value type. + binary_series: Mutex, } impl TimeSeriesService { @@ -44,6 +49,7 @@ impl TimeSeriesService { api_service, base_url, spool: Mutex::new(None), + binary_series: Mutex::new(binary::SeriesCache::default()), } } @@ -369,8 +375,12 @@ impl TimeSeriesService { .datapoints .extend(orig_dp_collection.datapoints.clone()); } - new_json.add_item(new_dp_collection.clone()); - total_datapoints = total_datapoints - new_dp_collection.datapoints.len(); + // Moved, not cloned. This used to deep-copy the collection here and the + // whole request body again below, so every datapoint was copied twice on + // its way out, and a DatapointString is two heap Strings. + let moved_points = new_dp_collection.datapoints.len(); + new_json.add_item(new_dp_collection); + total_datapoints = total_datapoints - moved_points; println!("Total datapoints left: {}", total_datapoints); } @@ -383,8 +393,7 @@ impl TimeSeriesService { new_total_datapoints ); - let new_json_clone = new_json.clone(); - new_request_bodies.push(new_json_clone); + new_request_bodies.push(new_json); } } // Now create futures after all request bodies are created diff --git a/src/timeseries/test.rs b/src/timeseries/test.rs index 3955299..4512a75 100644 --- a/src/timeseries/test.rs +++ b/src/timeseries/test.rs @@ -12,6 +12,7 @@ mod tests { use crate::generic::{DataWrapper, Datapoint, DatapointString, DatapointsCollection, DeleteFilter, IdAndExtId, RetrieveFilter}; use crate::http::ResponseError; use crate::timeseries::{TimeSeries, TimeSeriesFilter, TimeSeriesFilterForm, TimeSeriesUpdate, TimeSeriesUpdateCollection, TimeSeriesUpdateFields}; + use crate::timeseries::binary::BinaryIngestOptions; use crate::tests::cleanup::cleanup_timeseries; use crate::tests::ids::{unique_id, unique_token}; use crate::tests::polling::poll_until_for; @@ -564,6 +565,397 @@ mod tests { Ok(()) } + + /// The binary twin of `test_datapoints`: the same sixty days of one-second values, sent + /// through `POST /timeseries/data/binary` as zstd-compressed Arrow frames, then read back + /// through the same validations. Needs a backend that serves the binary endpoint. + #[tokio::test] + #[ignore] + async fn test_datapoints_binary() -> Result<(), Box> { + // Fixed for the same reason as test_datapoints's 6540, and two away from it so the two + // ignored tests can run in one process without addressing each other's series. + let unique_id: u64 = 6542; + let api_service = create_api_service(); + let new_ts_ext_id = format!("rust_sdk_test_{id}_ts", id = unique_id); + + // Fixed id, so a run that died before its teardown left this behind. + delete_timeseries(&api_service, &[&new_ts_ext_id]).await; + + let mut ts_collection = DataWrapper::new(); + let ts = TimeSeries::builder() + .set_external_id(new_ts_ext_id.as_str()) + .set_name(format!("Rust SDK Test {id} TimeSeries", id = unique_id).as_str()) + .set_description("This is test timeseries generated by rust sdk test code.") + .set_unit("celsius") + .set_value_type("float") + .clone(); + ts_collection.add_item(ts); + let created = api_service.time_series.create(&ts_collection).await; + let mut ts_cleanup = cleanup_timeseries(vec![new_ts_ext_id.clone()]); + match created { + Ok(timeseries) => assert_eq!(timeseries.length(), 1), + Err(e) => panic!("could not create the timeseries: {:?}", e.get_message()), + } + + println!("Prepare datapoints..."); + let mut data_request: DataWrapper> = DataWrapper::new(); + let mut dp_collection = DatapointsCollection::from_external_id(new_ts_ext_id.as_str()); + let datetime = Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap(); + dp_collection.datapoints = create_daily_datapoints(datetime); + let last = dp_collection.datapoints.last().unwrap().clone(); + let inserted_points = dp_collection.datapoints.len(); + data_request.add_item(dp_collection); + + println!("Start binary datapoint insert of {inserted_points} points!"); + let started = std::time::Instant::now(); + let result = api_service + .time_series + .insert_datapoints_binary(&data_request, &BinaryIngestOptions::default()) + .await; + match result { + Ok(r) => assert_eq!(r.get_http_status_code().unwrap(), StatusCode::NO_CONTENT.as_u16()), + Err(e) => panic!("binary insert failed with {}: {}", e.get_status(), e.get_message()), + } + println!("Binary insert took {:?}", started.elapsed()); + + // The binary path refreshes the latest-value cache from each series' last row as the + // request is accepted, so this is readable before ClickHouse has merged anything. + let id_collection = DataWrapper::from_vec(vec![IdAndExtId::from_external_id(&new_ts_ext_id)]); + let latest = api_service.time_series.retrieve_latest_datapoint(&id_collection).await?; + let latest_dp = latest.get_items().first().unwrap().datapoints.first().unwrap(); + assert_eq!(latest_dp.timestamp.timestamp_millis(), last.timestamp.parse::().unwrap()); + // Compared within a few ULP, not bit-exactly. The value crosses three systems as decimal + // text (this client, the api's cache, the read path's own rendering), and one of these + // fixture values carries all 17 significant digits a f64 can hold, which came back one + // ULP off. What this asserts is that the latest-value cache holds the right point. + let want = last.value.parse::().unwrap(); + let got = latest_dp.value.unwrap(); + assert!( + (got - want).abs() <= want.abs() * 4.0 * f64::EPSILON, + "latest value {got} is not {want}" + ); + + // Wait for the ClickHouse insert+merge to expose every datapoint, then validate. + poll_datapoint_count(&api_service, &new_ts_ext_id, 100000).await; + validate_datapoints(&api_service, vec![new_ts_ext_id.clone()]).await; + + // validate_daily_avg checks each daily average against the source file, which says the + // binary path stored the right values. Comparing avg/min/max with a JSON-ingested twin + // on top says the two paths store the *same* values, which is what this test is for. + println!("Validate aggregates against the JSON path..."); + let json_ext_id = format!("{new_ts_ext_id}_json"); + let mut json_ts_collection = DataWrapper::new(); + json_ts_collection.add_item( + TimeSeries::builder() + .set_external_id(json_ext_id.as_str()) + .set_name(json_ext_id.as_str()) + .set_unit("celsius") + .set_value_type("float") + .clone(), + ); + api_service.time_series.create(&json_ts_collection).await + .expect("could not create the JSON comparison series"); + let mut json_cleanup = cleanup_timeseries(vec![json_ext_id.clone()]); + + let mut json_request: DataWrapper> = DataWrapper::new(); + let mut json_dps = DatapointsCollection::from_external_id(json_ext_id.as_str()); + json_dps.datapoints = create_daily_datapoints(datetime); + json_request.add_item(json_dps); + api_service.time_series.insert_datapoints(&mut json_request).await + .expect("JSON insert failed"); + + // Both series must be complete before the aggregates can be compared. poll_datapoint_count + // reads with a 100k limit, so it cannot see past the first 100k of 5.18M and returns long + // before the series has landed; comparing then comes back unequal because one side is + // still filling, which looks exactly like a path that stores different values. + poll_all_points(&api_service, &new_ts_ext_id, inserted_points).await; + poll_all_points(&api_service, &json_ext_id, inserted_points).await; + + validate_daily_avg(&api_service, vec![new_ts_ext_id.clone()]).await; + + let binary_aggs = daily_aggregates(&api_service, &new_ts_ext_id).await; + let json_aggs = daily_aggregates(&api_service, &json_ext_id).await; + assert!(!binary_aggs.is_empty(), "no daily buckets came back"); + assert_eq!( + binary_aggs, json_aggs, + "the binary path aggregates differently from the JSON path for identical input" + ); + + println!("Validate raw datapoints with a cursor walk..."); + // Not validate_raw_datapoints_with_cursor: that one asserts a fixed final page size + // measured against a series that had been written to more than once, so it only holds + // for whatever the JSON test's series happens to contain. Here the count is known, so + // the walk asserts the total and the page shape instead. + let walked = walk_all_datapoints(&api_service, &new_ts_ext_id).await; + assert_eq!(walked, inserted_points, "cursor walk returned {walked} of {inserted_points} points"); + + println!("Delete datapoints"); + validate_deleted_datapoints(&api_service, new_ts_ext_id.clone()).await; + + delete_timeseries(&api_service, &[&new_ts_ext_id, &json_ext_id]).await; + ts_cleanup.disarm(); // explicit delete succeeded; skip the drop teardown + json_cleanup.disarm(); + + Ok(()) + } + + /// JSON against binary ingest from Rust, reporting the same columns as the Java benchmark in + /// the platform's `datahub-e2e` module and the Python one in `python_tests`, so the three + /// clients can be compared. + /// + /// Sized by `DATAHUB_BENCH_POINTS` (default 10 million) across `DATAHUB_BENCH_SERIES` + /// series. The api must have its daily quota and rate limiter off for a run of any size: + /// `-Ddatahub.limits.quota.enabled=false -Ddatahub.limits.rate.enabled=false`. + /// + /// **Run it with `--release`.** `cargo test` builds with the `dev` profile at + /// `opt-level = 0`, and the first time this was measured that way it reported Rust as + /// slower than Java at both paths: 339k against 720k points per second on JSON, and 1.0M + /// against 3.8M on binary. The assertion below refuses to run without optimisation rather + /// than print numbers that mean nothing. + #[tokio::test] + #[ignore] + async fn bench_json_vs_binary() -> Result<(), Box> { + // debug_assertions is on in the dev profile and off in release, which is the cheapest + // reliable way to tell which one built this. + if cfg!(debug_assertions) { + panic!("built without optimisation; run with `cargo test --release`, or this \ + measures opt-level 0 rather than the SDK"); + } + let total: usize = std::env::var("DATAHUB_BENCH_POINTS") + .ok().and_then(|v| v.parse().ok()).unwrap_or(10_000_000); + let series_count: usize = std::env::var("DATAHUB_BENCH_SERIES") + .ok().and_then(|v| v.parse().ok()).unwrap_or(100); + let chunk: usize = std::env::var("DATAHUB_BENCH_CHUNK") + .ok().and_then(|v| v.parse().ok()).unwrap_or(1_000_000); + let api_service = create_api_service(); + println!("\n=== {total} points across {series_count} series, float32 ==="); + + // Which paths to run. A sweep over request sizes only needs the binary one, and the JSON + // path is slow enough that running it every time makes the sweep the expensive part. + let paths = std::env::var("DATAHUB_BENCH_PATHS").unwrap_or_else(|_| "both".to_string()); + let wanted: Vec = match paths.as_str() { + "binary" => vec![true], + "json" => vec![false], + _ => vec![false, true], + }; + + let mut results = Vec::new(); + for binary in wanted { + let label = if binary { "binary" } else { "JSON" }; + let run_id = unique_id(if binary { "bench_bin" } else { "bench_json" }); + let external_ids: Vec = + (0..series_count).map(|i| format!("{run_id}_{i}")).collect(); + + let mut ts_collection = DataWrapper::new(); + for external_id in &external_ids { + ts_collection.add_item( + TimeSeries::builder() + .set_external_id(external_id) + .set_name(external_id) + .set_unit("celsius") + .set_value_type("float32") + .clone(), + ); + } + api_service.time_series.create(&ts_collection).await + .expect("could not create the benchmark series"); + let mut cleanup = cleanup_timeseries(external_ids.clone()); + + let per_series = chunk / series_count; + let mut latencies: Vec = Vec::new(); + let mut sent = 0usize; + let mut offset = 0i64; + let started = std::time::Instant::now(); + while sent < total { + let this_chunk = std::cmp::min(chunk, total - sent); + let per = std::cmp::max(1, this_chunk / series_count); + let mut request: DataWrapper> = DataWrapper::new(); + for (index, external_id) in external_ids.iter().enumerate() { + let mut collection = DatapointsCollection::from_external_id(external_id); + collection.datapoints = bench_points(offset, per, index); + request.add_item(collection); + } + let call = std::time::Instant::now(); + let concurrency: usize = std::env::var("DATAHUB_BENCH_REQ_CONCURRENCY") + .ok().and_then(|v| v.parse().ok()).unwrap_or(4); + if binary { + api_service.time_series + .insert_datapoints_binary( + &request, + &BinaryIngestOptions::default().request_concurrency(concurrency), + ) + .await + .unwrap_or_else(|e| panic!("{label} insert failed: {}", e.get_message())); + } else { + let mut json_request = request.clone(); + api_service.time_series + .insert_datapoints(&mut json_request) + .await + .unwrap_or_else(|e| panic!("{label} insert failed: {}", e.get_message())); + } + latencies.push(call.elapsed().as_secs_f64() * 1000.0); + sent += per * series_count; + offset += per as i64; + let elapsed = started.elapsed().as_secs_f64(); + println!(" {label}: {sent} / {total} points, {:.0} pts/s", sent as f64 / elapsed); + } + let ingest_seconds = started.elapsed().as_secs_f64(); + + latencies.sort_by(|a, b| a.partial_cmp(b).unwrap()); + let mean = latencies.iter().sum::() / latencies.len() as f64; + let p50 = latencies[latencies.len() / 2]; + let p99 = latencies[std::cmp::min(latencies.len() - 1, latencies.len() * 99 / 100)]; + results.push(( + label, + sent, + ingest_seconds, + sent as f64 / ingest_seconds, + latencies.len(), + mean, + p50, + p99, + )); + + delete_timeseries(&api_service, &external_ids.iter().map(|s| s.as_str()).collect::>()).await; + cleanup.disarm(); + } + + println!("\n=== datapoint ingest from Rust ==="); + println!("{total} points across {series_count} series, {chunk} per call, float32\n"); + print!("{:<28}", "metric"); + for r in &results { print!("{:>18}", r.0); } + println!(); + let row = |label: &str, f: &dyn Fn(&(&str, usize, f64, f64, usize, f64, f64, f64)) -> String| { + print!("{label:<28}"); + for r in &results { print!("{:>18}", f(r)); } + println!(); + }; + row("ingest wall time (s)", &|r| format!("{:.1}", r.2)); + row("points per second", &|r| format!("{:.0}", r.3)); + row("points per second in-call", &|r| format!("{:.0}", r.1 as f64 / (r.5 / 1000.0 * r.4 as f64))); + row("sdk calls", &|r| format!("{}", r.4)); + row("latency mean (ms)", &|r| format!("{:.0}", r.5)); + row("latency p50 (ms)", &|r| format!("{:.0}", r.6)); + row("latency p99 (ms)", &|r| format!("{:.0}", r.7)); + Ok(()) + } + + /// A slow sine plus noise, one signal per series. Identical series would let zstd compress + /// the repetition across them and report a wire size no real fleet of sensors produces. + fn bench_points(offset_seconds: i64, count: usize, series_index: usize) -> Vec { + let base = 150.0 + series_index as f64 * 0.7; + let phase = series_index as f64 * 0.37; + let mut points = Vec::with_capacity(count); + for i in 0..count { + let t = offset_seconds + i as i64; + let mut z = (t as u64) + .wrapping_add((series_index as u64).wrapping_mul(0x5851_F42D_4C95_7F2D)) + .wrapping_mul(0x9E37_79B9_7F4A_7C15); + z ^= z >> 30; + z = z.wrapping_mul(0xBF58_476D_1CE4_E5B9); + z ^= z >> 27; + let noise = ((z >> 40) as f64 / (1u64 << 24) as f64) - 0.5; + let value = (base + 20.0 * (t as f64 / 600.0 + phase).sin() + noise) as f32; + let timestamp = 1_735_689_600_000i64 + t * 1000; + points.push(DatapointString::new(×tamp.to_string(), &value.to_string())); + } + points + } + + /// Waits until the whole series is readable, by walking it. Needed because the cheap count + /// read is capped at its own limit and cannot tell "100k so far" from "all of it". + async fn poll_all_points(api_service: &Arc, external_id: &str, want: usize) { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300); + loop { + let have = walk_all_datapoints(api_service, external_id).await; + if have >= want { + return; + } + if std::time::Instant::now() > deadline { + panic!("{external_id} reached only {have} of {want} points before the deadline"); + } + tokio::time::sleep(std::time::Duration::from_secs(2)).await; + } + } + + fn truncate_10(x: f64) -> f64 { + // Clickhouse will have rounding errors using for example avg(), so we truncate the returned + // values to mitigate this + let multiplier = 10f64.powf(10.0); + (x * multiplier).floor() / multiplier + } + + /// Daily avg/min/max for the whole window, truncated so two paths that stored the same + /// values compare equal without depending on how many digits the read path prints. + async fn daily_aggregates( + api_service: &Arc, + external_id: &str, + ) -> Vec<(i64, f64, f64, f64)> { + let mut data_request: DataWrapper = DataWrapper::new(); + let mut rf = RetrieveFilter::new(); + rf.set_external_id(external_id); + rf.set_start(Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap()); + rf.set_end(Utc.with_ymd_and_hms(2025, 3, 2, 0, 0, 0).unwrap()); + rf.set_aggregates(vec!["avg".to_string(), "min".to_string(), "max".to_string()]); + rf.set_granularity("1d"); + data_request.add_item(rf); + let response = api_service + .time_series + .retrieve_datapoints(&data_request) + .await + .expect("aggregate read failed"); + response + .get_items() + .first() + .map(|item| { + item.datapoints + .iter() + .map(|dp| { + ( + dp.timestamp().timestamp_millis(), + truncate_10(dp.average().unwrap_or(f64::NAN)), + truncate_10(dp.min().unwrap_or(f64::NAN)), + truncate_10(dp.max().unwrap_or(f64::NAN)), + ) + }) + .collect() + }) + .unwrap_or_default() + } + + /// Pages the whole series with the keyset cursor and returns how many points came back. + /// Every page but the last must be full, which is what proves the cursor is not skipping. + async fn walk_all_datapoints(api_service: &Arc, external_id: &str) -> usize { + let mut data_request: DataWrapper = DataWrapper::new(); + let mut rf = RetrieveFilter::new(); + rf.set_external_id(external_id); + rf.set_start(Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap()); + rf.set_end(Utc.with_ymd_and_hms(2025, 3, 2, 0, 0, 0).unwrap()); + // Limit 0, not an explicit page size: an explicit limit caps the whole result and comes + // back without a cursor, so the walk would stop after one page and read as a series that + // only ever received its first hundred thousand points. + rf.set_limit(0); + data_request.add_item(rf); + + let mut total = 0usize; + let mut cursor: Option = None; + loop { + let mut request = data_request.clone(); + request.get_items_mut().first_mut().unwrap().cursor = cursor.clone(); + let response = api_service + .time_series + .retrieve_datapoints(&request) + .await + .expect("cursor page failed"); + let page = response.get_items().first().expect("no series in the page"); + total += page.datapoints.len(); + cursor = page.next_cursor.clone(); + if cursor.is_none() { + break; + } + } + total + } // total is 9 354 000 /// Poll a series until `want` datapoints are readable from the start of 2025. @@ -887,15 +1279,9 @@ mod tests { let result = api_service.time_series.retrieve_datapoints(&data_request).await; match result { Ok(r) => { - if let Some(first_item) = r.get_items().first() { - if let Some(external_id) = &first_item.external_id { - if external_id == "rust_sdk_test_6540_ts" { - assert_eq!(r.get_items().first().unwrap().datapoints.len(), 59); - } else { - assert_eq!(r.get_items().first().unwrap().datapoints.len(), 25); - } - } - } + // Every caller inserts the full sixty days from 2025-01-01 and reads back to + // 2025-03-01: 59 daily buckets, whichever series carries them. + assert_eq!(r.get_items().first().unwrap().datapoints.len(), 59); let start_date = Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap(); let end_date = Utc.with_ymd_and_hms(2025, 3, 1, 0, 0, 0).unwrap(); @@ -1155,6 +1541,38 @@ mod tests { Ok(()) } + /// The binary path resolves every series before it builds a frame, so a missing one is + /// refused here, with the JSON path's wording, and no frame is ever sent. + #[tokio::test] + async fn test_insert_datapoints_binary_missing_timeseries_returns_not_found() -> Result<(), Box> { + let api_service = create_api_service(); + + let missing_ext_id = unique_id("ts_missing"); + let mut data_request: DataWrapper> = DataWrapper::new(); + let mut dp_collection = DatapointsCollection::from_external_id(&missing_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); + + let result = api_service + .time_series + .insert_datapoints_binary(&data_request, &BinaryIngestOptions::default()) + .await; + match result { + Ok(_) => panic!("Expected 404 Not Found for non-existent timeseries"), + Err(e) => { + assert_eq!(e.get_status(), StatusCode::NOT_FOUND); + let msg = e.get_message(); + assert!( + msg.contains("Could not find following timeseries"), + "unexpected error body: {msg}" + ); + } + } + Ok(()) + } + fn validate_data_insertion(result: Result, ResponseError>) { match result { Ok(r) => {