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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,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`. Every refusal on this path is a problem document typed for its status — `invalid-frame` (400), `unknown-timeseries` (404), `request-too-large` (413), `unsupported-media-type` (415), `value-type-mismatch` and `external-id-mismatch` (422), `too-many-in-flight` (429) — with the kebab-case sub-case in a `reason` extension beside it. It was one type, `datapoint-block-rejected`, answering with all six statuses until the api split it (platform #120); the SDK matches the new slugs only, and nothing here recognises the old one — the binary path has never been in a release, so no caller can be on the other side of that change. `is_stale_series_rejection` reads the slug to decide the one re-resolve-and-retry, and reads it from the problem rather than the body text because the SDK's own pre-flight 404 names the series it could not resolve — an external id spelling `unknown-timeseries` used to trigger the 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. 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`. Every refusal on this path is a problem document typed for its status — `invalid-frame` (400), `unknown-timeseries` (404), `request-too-large` (413), `unsupported-media-type` (415), `value-type-mismatch` and `external-id-mismatch` (422), `too-many-in-flight` (429) — with the kebab-case sub-case in a `reason` extension beside it. It was one type, `datapoint-block-rejected`, answering with all six statuses until the api split it (platform #120); the SDK matches the new slugs only, and nothing here recognises the old one — the binary path has never been in a release, so no caller can be on the other side of that change. `is_stale_series_rejection` reads the slug to decide the one re-resolve-and-retry, and reads it from the problem rather than the body text because the SDK's own pre-flight 404 names the series it could not resolve — an external id spelling `unknown-timeseries` used to trigger the 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. In Python it is `insert_datapoints_binary` and `insert_from_lists_binary` on the **sync** client only, each taking an optional `zstd_level`; the async client has no binary methods. 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. `python_tests/test_timeseries_binary.py` round-trips the Python binding through the JSON read path.

### The `ApiServiceProvider` trait (`src/generic.rs`)

Expand Down
112 changes: 112 additions & 0 deletions python_tests/test_timeseries_binary.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
"""Binary datapoint ingest (`POST /timeseries/data/binary`) through the Python bindings.

What goes in through the binary path must read back through the ordinary JSON read path, point for
point. The Rust side covers the frame layout offline (`src/timeseries/binary.rs`) and the
multi-million-point load in the ignored `test_datapoints_binary`; this is the binding surface.
"""
import datetime

import pytest

import intellistream_datahub_sdk
from intellistream_datahub_sdk import DataHubException
from polling import poll_until
from python_tests.fixtures import * # noqa: F401,F403

START = datetime.datetime(2024, 3, 1, tzinfo=datetime.timezone.utc)
TIMESTAMPS = [START + datetime.timedelta(minutes=i) for i in range(50)]


def _read_all(client, ts, expected):
rf = intellistream_datahub_sdk.RetrieveFilter(
ts=ts,
start=START - datetime.timedelta(days=1),
end=START + datetime.timedelta(days=1),
)

def fetch():
collections = client.timeseries.retrieve_datapoints(rf)
return collections[0].get_datapoints() if collections else []

return poll_until(fetch, lambda dps: len(dps) >= expected)


@pytest.mark.parametrize("zstd_level", [None, 1, 3, 9])
def test_insert_from_lists_binary_round_trips_floats(sync_client, make_ts, zstd_level):
ts = make_ts(value_type="float")
# Out of order on purpose: the frame writer sorts, the read must still be in time order.
values = [i * 1.25 - 20.0 for i in range(len(TIMESTAMPS))]
shuffled = list(zip(TIMESTAMPS, values))[::-1]

result = sync_client.timeseries.insert_from_lists_binary(
[t for t, _ in shuffled], [v for _, v in shuffled], ts, zstd_level=zstd_level
)
assert result == []

dps = _read_all(sync_client, ts, len(TIMESTAMPS))
assert [dp.timestamp for dp in dps] == TIMESTAMPS
assert [dp.value for dp in dps] == pytest.approx(values)


def test_insert_datapoints_binary_round_trips_bigints(sync_client, make_ts):
ts = make_ts(value_type="bigint")
values = [(-1) ** i * i * 1_000_003 for i in range(len(TIMESTAMPS))]
data = [
intellistream_datahub_sdk.DatapointString.from_int(t, v)
for t, v in zip(TIMESTAMPS, values)
]
collection = intellistream_datahub_sdk.DatapointsCollectionString(datapoints=data, ts=ts)

assert sync_client.timeseries.insert_datapoints_binary([collection]) == []

dps = _read_all(sync_client, ts, len(TIMESTAMPS))
assert [dp.timestamp for dp in dps] == TIMESTAMPS
assert [dp.value for dp in dps] == values


def test_insert_datapoints_binary_spans_several_series(sync_client, make_ts):
first, second = make_ts(value_type="float"), make_ts(value_type="bigint")
collections = [
intellistream_datahub_sdk.DatapointsCollectionString(
datapoints=[intellistream_datahub_sdk.DatapointString.from_float(t, 0.5) for t in TIMESTAMPS],
ts=first,
),
intellistream_datahub_sdk.DatapointsCollectionString(
datapoints=[intellistream_datahub_sdk.DatapointString.from_int(t, 7) for t in TIMESTAMPS[:10]],
ts=second,
),
]

assert sync_client.timeseries.insert_datapoints_binary(collections) == []

assert [dp.value for dp in _read_all(sync_client, first, 50)] == [0.5] * 50
assert [dp.value for dp in _read_all(sync_client, second, 10)] == [7] * 10


def test_binary_insert_duplicate_timestamps_keep_one_point(sync_client, make_ts):
ts = make_ts(value_type="float")
at = TIMESTAMPS[0]

sync_client.timeseries.insert_from_lists_binary([at, at, TIMESTAMPS[1]], [1.0, 2.0, 3.0], ts)

dps = _read_all(sync_client, ts, 2)
assert [dp.timestamp for dp in dps] == [at, TIMESTAMPS[1]]


def test_binary_insert_refuses_an_unknown_zstd_level(sync_client, make_ts):
ts = make_ts(value_type="float")
with pytest.raises(DataHubException, match="zstd level 2") as excinfo:
sync_client.timeseries.insert_from_lists_binary([START], [1.0], ts, zstd_level=2)
assert excinfo.value.status_code == 400


def test_binary_insert_into_a_missing_series_is_not_found(sync_client):
missing = unique_id("ts_missing")
with pytest.raises(DataHubException, match="Could not find following timeseries") as excinfo:
sync_client.timeseries.insert_from_lists_binary([START], [1.0], missing)
assert excinfo.value.status_code == 404


def test_insert_from_lists_binary_rejects_mismatched_lengths(sync_client):
with pytest.raises(ValueError, match="got 2 and 1"):
sync_client.timeseries.insert_from_lists_binary([START, START], [1.0], "never_sent")
Loading