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
8 changes: 1 addition & 7 deletions python_tests/test_edges.py
Original file line number Diff line number Diff line change
Expand Up @@ -186,13 +186,7 @@ def test_create_between_existing_resources(sync_client):
relationship_type="SDK_TEST_LINK", from_external_id=a, to_external_id=b
)
try:
try:
created = sync_client.edges.create([form])
except intellistream_datahub_sdk.DataHubException as exc:
# POST /edges/create is newer than some backends; skip rather than fail there.
if "405" in str(exc):
pytest.skip("this backend has no POST /edges/create")
raise
created = sync_client.edges.create([form])

assert len(created) == 1
edge_id = created[0].id
Expand Down
52 changes: 35 additions & 17 deletions python_tests/test_event_vocabulary.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,13 @@
tenant actually use?" for the four categorical event fields, so they back filter dropdowns and
autocompletes. They return plain strings, not events.
"""
from datetime import datetime, timezone

import intellistream_datahub_sdk
import pytest

from fixtures import async_client, sync_client
from fixtures import TEST_PREFIX, async_client, sync_client, unique_id
from polling import poll_until


DIMENSIONS = [
Expand All @@ -17,6 +20,30 @@
intellistream_datahub_sdk.EventDimension.SOURCE,
]

# Fixed, not unique: a dimension row outlives the events that put it there, so a fresh name per
# run would add a vocabulary entry per run.
SEEDED_TYPES = [f"{TEST_PREFIX}vocab_type_a", f"{TEST_PREFIX}vocab_type_b"]


@pytest.fixture(scope="module")
def seeded_types(sync_client):
events = [
intellistream_datahub_sdk.Event(
type=t, external_id=unique_id("vocab_event"), event_time=datetime.now(timezone.utc)
)
for t in SEEDED_TYPES
]
sync_client.events.create(events)
try:
found = poll_until(
lambda: sync_client.events.search_types(f"{TEST_PREFIX}vocab_type"),
lambda found: set(SEEDED_TYPES) <= set(found),
)
assert set(SEEDED_TYPES) <= set(found), found
yield SEEDED_TYPES
finally:
sync_client.events.delete(events)


def test_list_dimensions_are_distinct(sync_client):
for dimension in DIMENSIONS:
Expand Down Expand Up @@ -46,10 +73,7 @@ def test_named_helpers_match_the_generic_form(sync_client):
)


def test_limit_caps_results(sync_client):
values = sync_client.events.list_types()
if len(values) < 2:
pytest.skip("tenant has fewer than two distinct event types")
def test_limit_caps_results(sync_client, seeded_types):
assert len(sync_client.events.list_types(limit=1)) == 1


Expand All @@ -59,22 +83,16 @@ def test_limit_is_clamped_not_rejected(sync_client):
assert isinstance(sync_client.events.list_types(limit=0), list)


def test_search_is_case_insensitive_substring(sync_client):
values = sync_client.events.list_types()
if not values:
pytest.skip("tenant has no events with a type")
sample = values[0]
def test_search_is_case_insensitive_substring(sync_client, seeded_types):
sample = seeded_types[0]

assert sample in sync_client.events.search_types(sample)

flipped = sample.upper() if sample.islower() else sample.lower()
assert sample in sync_client.events.search_types(flipped)

if len(sample) > 2:
assert sample in sync_client.events.search_types(sample[:-1])
assert sample in sync_client.events.search_types(sample.upper())
assert sample in sync_client.events.search_types(sample[:-1])

# Search filters the vocabulary; it never invents values.
assert set(sync_client.events.search_types(sample)) <= set(values)
hits = sync_client.events.search_types(sample[:-1])
assert all(sample[:-1] in v.lower() for v in hits), hits


def test_no_match_is_empty_not_an_error(sync_client):
Expand Down
36 changes: 17 additions & 19 deletions python_tests/test_filter_resources.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
import intellistream_datahub_sdk
from intellistream_datahub_sdk import DataHubException

from fixtures import async_client, sync_client, unique_id # noqa: F401 (fixtures)
from fixtures import async_client, make_resource, sync_client, unique_id # noqa: F401 (fixtures)
from filter_fixtures import ( # noqa: F401 (fixtures)
datasets,
prefix,
Expand Down Expand Up @@ -269,20 +269,20 @@ def test_every_node_carries_its_type_as_a_label(sync_client, resource_corpus, pr
assert "TIMESERIES" in (node.labels or []), f"{node.external_id} has labels {node.labels}"


@pytest.mark.parametrize("type_label,node_type", [
("ASSET", "asset"),
("DATASET", "dataset"),
("TIMESERIES", "timeseries"),
("FUNCTION", "function"),
("POLICY", "policy"),
@pytest.mark.parametrize("type_label,build", [
("ASSET", lambda ext: intellistream_datahub_sdk.Asset(external_id=ext, name=ext, is_root=True)),
("DATASET", lambda ext: intellistream_datahub_sdk.Dataset(external_id=ext, name=ext)),
("TIMESERIES", lambda ext: intellistream_datahub_sdk.TimeSeries(
external_id=ext, name=ext, value_type="float", unit="a.u")),
("FUNCTION", lambda ext: intellistream_datahub_sdk.Function(external_id=ext, name=ext)),
("POLICY", lambda ext: intellistream_datahub_sdk.Policy(
external_id=ext, name=ext, type="HAS_REQUIREMENT", value=True)),
])
def test_every_type_label_is_matchable(sync_client, type_label, node_type):
def test_every_type_label_is_matchable(sync_client, make_resource, type_label, build):
"""A node reports its type-label on every read, so filtering by it must find that node.

Asserted tenant-wide rather than against this run's corpus: a type-label is one shared row per
name, so whether it matches is a property of that row, not of any node a fixture can create.
The failure mode is silent either way — an empty result is indistinguishable from "nothing is
tagged that way" — which is what makes it worth pinning.
The failure mode is silent — an empty result is indistinguishable from "nothing is tagged that
way" — which is what makes it worth pinning.

Four of these five were broken by hash drift until V37: commit 5c22b485 dropped the line writing
`label.hash` with XXH64 and left `Label.setName`'s XXH3 as the only writer, so every row written
Expand All @@ -292,15 +292,13 @@ def test_every_type_label_is_matchable(sync_client, type_label, node_type):
POLICY was a different fault and stayed red after V37: there was no POLICY row in the label
table at all, so a policy node reported `labels: ['POLICY']` from the denormalised
`node.labels` column while the filter's join on the label hash had nothing to join to — the
label visible and unsearchable, with no row to repair. Fixed server-side; it is a plain case
here now rather than a strict xfail.
label visible and unsearchable, with no row to repair.
"""
of_type = sync_client.resources.filter(node_type=[node_type], limit=1000)
if not of_type:
pytest.skip(f"no {node_type} nodes in this tenant to match")
ext = unique_id(f"type_label_{type_label.lower()}")
make_resource([build(ext)])

by_label = sync_client.resources.filter(labels=[type_label], limit=1000)
assert by_label, f"{len(of_type)} {node_type} nodes exist but none match labels=[{type_label}]"
by_label = sync_client.resources.filter(labels=[type_label], external_id=[ext])
assert externals(by_label) == {ext}, f"the new node does not match labels=[{type_label}]"


def test_garbled_criteria_match_nothing_without_erroring(flt, resource_corpus):
Expand Down
15 changes: 1 addition & 14 deletions python_tests/test_subscriptions.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,7 @@
"""Tests for the Python subscriptions module.

Mirrors src/subscriptions/test.rs: a CRUD round-trip and a listen end-to-end test. The
end-to-end listen test needs the backend's Pulsar fan-out consumer running, so it is gated
behind RUN_LISTEN_TESTS=1 (matching the Rust `#[ignore]`).
Mirrors src/subscriptions/test.rs: a CRUD round-trip and a listen end-to-end test.
"""
import os
import time
from datetime import datetime, timezone

Expand Down Expand Up @@ -126,12 +123,6 @@ def test_ws_datapoint_as_float():

# --- Listen end-to-end -----------------------------------------------------------------

# Skipped by default — needs the backend's Pulsar consumer running so REST datapoint writes
# fan out to the subscription topic. Set RUN_LISTEN_TESTS=1 to enable.
listen_enabled = os.environ.get("RUN_LISTEN_TESTS") == "1"


@pytest.mark.skipif(not listen_enabled, reason="set RUN_LISTEN_TESTS=1 to run live listen tests")
def test_listen_end_to_end(sync_client):
ts_ext = unique_id("listen_ts")
sub_ext = unique_id("listen")
Expand Down Expand Up @@ -205,7 +196,6 @@ def test_listen_end_to_end(sync_client):
pass


@pytest.mark.skipif(not listen_enabled, reason="set RUN_LISTEN_TESTS=1 to run live listen tests")
def test_listen_context_manager_closes_cleanly(sync_client):
ts_ext = unique_id("ctx_ts")
sub_ext = unique_id("ctx")
Expand Down Expand Up @@ -234,7 +224,6 @@ def test_listen_context_manager_closes_cleanly(sync_client):
sync_client.timeseries.delete([ts])


@pytest.mark.skipif(not listen_enabled, reason="set RUN_LISTEN_TESTS=1 to run live listen tests")
def test_listen_fans_out_all_bound_timeseries(sync_client):
"""One subscription bound to several timeseries must fan out datapoints from ALL of them.

Expand Down Expand Up @@ -303,7 +292,6 @@ def test_listen_fans_out_all_bound_timeseries(sync_client):
)


@pytest.mark.skipif(not listen_enabled, reason="set RUN_LISTEN_TESTS=1 to run live listen tests")
def test_listen_refused_subscription_surfaces_as_error(sync_client):
"""A subscription the server refuses is raised to the caller as an exception, not swallowed.

Expand All @@ -329,7 +317,6 @@ def test_listen_refused_subscription_surfaces_as_error(sync_client):
pass


@pytest.mark.skipif(not listen_enabled, reason="set RUN_LISTEN_TESTS=1 to run live listen tests")
def test_listen_partial_refusal_keeps_valid_subscription(sync_client):
"""A refused subscription on a multiplexed connection surfaces as an error but does NOT tear the
socket down — the valid subscription on the same connection keeps delivering. This is the core
Expand Down
3 changes: 1 addition & 2 deletions python_tests/test_units.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,7 @@
@pytest.fixture(scope="module")
def some_unit(sync_client):
units = sync_client.units.list()
if not units:
pytest.skip("backend has no units configured")
assert units, "the units catalogue is seeded by the api's migrations"
return units[0]


Expand Down
6 changes: 0 additions & 6 deletions src/assets/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -113,9 +113,6 @@ fn search_body_omits_an_absent_filter() {
}

/// Live round-trip over the whole `/assets` surface.
///
/// `#[ignore]` like every test here that needs a backend; run with
/// `cargo test assets:: -- --ignored --nocapture`.
mod live {
use super::*;
use crate::create_api_service;
Expand All @@ -125,7 +122,6 @@ mod live {
use crate::tests::ids::unique_id;

#[tokio::test]
#[ignore]
async fn assets_full_roundtrip() {
let api = create_api_service();
let ext_id = unique_id("asset");
Expand Down Expand Up @@ -238,7 +234,6 @@ mod live {
/// A node that exists but is not an asset is reported as missing, not as a type error — so a
/// 404 here does not tell you whether the id exists.
#[tokio::test]
#[ignore]
async fn a_non_asset_id_is_reported_as_missing() {
let api = create_api_service();
let ext_id = unique_id("fn_not_asset");
Expand Down Expand Up @@ -269,7 +264,6 @@ mod live {
/// `geoLocation` set through the shared update form reaches the asset — the one update field
/// that means anything on exactly one node type.
#[tokio::test]
#[ignore]
async fn geolocation_is_updatable() {
let api = create_api_service();
let ext_id = unique_id("asset_geo");
Expand Down
34 changes: 22 additions & 12 deletions src/events/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -752,7 +752,11 @@ mod related_resources_serde {

mod vocabulary {
use crate::create_api_service;
use crate::events::EventDimension;
use crate::events::{Event, EventDimension};
use crate::tests::cleanup::cleanup_events;
use crate::tests::ids::{unique_id, TEST_PREFIX};
use crate::tests::polling::poll_until;
use chrono::Utc;

/// The four list endpoints: distinct values, alphabetical, honouring `limit`.
#[tokio::test]
Expand Down Expand Up @@ -793,14 +797,20 @@ mod vocabulary {
) -> Result<(), Box<dyn std::error::Error>> {
let api = create_api_service();

// Pick a real value to search for so the test doesn't depend on any particular tenant data.
let types = api.events.list_types(None).await?;
let Some(sample) = types.get_items().first().cloned() else {
println!("SKIP test_search_dimensions: this tenant has no events with a type");
return Ok(());
};
// Fixed, not unique: a dimension row outlives the events that put it there, so a fresh
// name per run would add a vocabulary entry per run.
let sample = format!("{TEST_PREFIX}vocab_type_a");
let ext_id = unique_id("vocab_event");
let _cleanup = cleanup_events(vec![ext_id.clone()]);
api.events
.create(&Event::new(ext_id, sample.clone(), Utc::now()))
.await?;

let exact = api.events.search_types(&sample, None).await?;
let exact = poll_until(
|| async { api.events.search_types(&sample, None).await },
|r| r.as_ref().is_ok_and(|dw| dw.get_items().contains(&sample)),
)
.await?;
assert!(
exact.get_items().contains(&sample),
"searching for {sample:?} should find it"
Expand Down Expand Up @@ -831,10 +841,10 @@ mod vocabulary {
);
}

// Every result is a subset of the full list — search filters, it does not invent values.
let all = types.get_items();
for found in exact.get_items() {
assert!(all.contains(&found), "{found:?} is not in the full list");
// Search filters the vocabulary, it does not invent values.
let part = &sample[..sample.len() - 1];
for found in api.events.search_types(part, None).await?.get_items() {
assert!(found.to_lowercase().contains(part), "{found:?} does not contain {part:?}");
}

// No match is an empty list, not an error.
Expand Down
8 changes: 0 additions & 8 deletions src/functions/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,7 @@ mod tests {
use crate::tests::ids::unique_id;

/// The whole `/functions` surface: create, list, get by id, by_ids, update, delete.
///
/// `#[ignore]` because it needs a live backend; run with
/// `cargo test functions:: -- --ignored`. It needs nothing beyond that — the `forecast-ema`
/// model template this comment used to name went away with the functions feature itself (see
/// the server's "Remove functions feature. revert to simple metadata store"), and a function is
/// now a plain node.
#[tokio::test]
#[ignore]
async fn functions_full_roundtrip() {
let api = create_api_service();
let ext_id = unique_id("fn");
Expand Down Expand Up @@ -98,7 +91,6 @@ mod tests {
/// `/functions/byids`, `/filter` and `/search` — the server-side reads that replaced the
/// client-side walk over the listing.
#[tokio::test]
#[ignore]
async fn functions_byids_filter_and_search() {
use crate::filters::NodeFilter;
use crate::functions::{FunctionFilter, FunctionFilterForm};
Expand Down
5 changes: 1 addition & 4 deletions src/labels/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,8 @@ mod tests {
}

// Live end-to-end exercise of the whole label lifecycle, including the
// delete-while-in-use error. Ignored by default: needs a configured backend (.env) and
// mutates tenant state. Run with `cargo test labels -- --ignored --nocapture`.
// delete-while-in-use error.
#[tokio::test]
#[ignore]
async fn test_label_lifecycle() -> Result<(), Box<dyn std::error::Error>> {
let api = create_api_service();

Expand Down Expand Up @@ -103,7 +101,6 @@ mod tests {
// Delete-while-in-use: create a label, attach it to a resource, and confirm the delete is
// rejected with a 400 whose body names the blocking resource. Ignored by default.
#[tokio::test]
#[ignore]
async fn test_delete_label_in_use_reports_blocker() -> Result<(), Box<dyn std::error::Error>> {
use crate::relations::RelForm;
use crate::resources::Resource;
Expand Down
Loading
Loading