diff --git a/python_tests/test_edges.py b/python_tests/test_edges.py index c0d6965..ecb63f1 100644 --- a/python_tests/test_edges.py +++ b/python_tests/test_edges.py @@ -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 diff --git a/python_tests/test_event_vocabulary.py b/python_tests/test_event_vocabulary.py index fa84091..ed65ffb 100644 --- a/python_tests/test_event_vocabulary.py +++ b/python_tests/test_event_vocabulary.py @@ -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 = [ @@ -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: @@ -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 @@ -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): diff --git a/python_tests/test_filter_resources.py b/python_tests/test_filter_resources.py index b280f0a..5142ee8 100644 --- a/python_tests/test_filter_resources.py +++ b/python_tests/test_filter_resources.py @@ -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, @@ -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 @@ -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): diff --git a/python_tests/test_subscriptions.py b/python_tests/test_subscriptions.py index 0278cca..2582a6a 100644 --- a/python_tests/test_subscriptions.py +++ b/python_tests/test_subscriptions.py @@ -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 @@ -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") @@ -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") @@ -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. @@ -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. @@ -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 diff --git a/python_tests/test_units.py b/python_tests/test_units.py index 42ed5f7..8ae6cd1 100644 --- a/python_tests/test_units.py +++ b/python_tests/test_units.py @@ -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] diff --git a/src/assets/tests.rs b/src/assets/tests.rs index c374069..4926c29 100644 --- a/src/assets/tests.rs +++ b/src/assets/tests.rs @@ -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; @@ -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"); @@ -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"); @@ -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"); diff --git a/src/events/tests.rs b/src/events/tests.rs index 61c8268..b088a1d 100644 --- a/src/events/tests.rs +++ b/src/events/tests.rs @@ -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] @@ -793,14 +797,20 @@ mod vocabulary { ) -> Result<(), Box> { 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" @@ -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. diff --git a/src/functions/test.rs b/src/functions/test.rs index 1c8e42c..d30bc97 100644 --- a/src/functions/test.rs +++ b/src/functions/test.rs @@ -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"); @@ -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}; diff --git a/src/labels/test.rs b/src/labels/test.rs index 6df7ec9..427cdc9 100644 --- a/src/labels/test.rs +++ b/src/labels/test.rs @@ -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> { let api = create_api_service(); @@ -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> { use crate::relations::RelForm; use crate::resources::Resource; diff --git a/src/resources/label_update_tests.rs b/src/resources/label_update_tests.rs index f729ec7..831ae28 100644 --- a/src/resources/label_update_tests.rs +++ b/src/resources/label_update_tests.rs @@ -2,7 +2,7 @@ //! //! Two groups: //! - offline serde tests that lock the request wire-format (set/add/remove, node identity); -//! - `#[ignore]` live tests that exercise the backend's label-update semantics, in particular the +//! - live tests that exercise the backend's label-update semantics, in particular the //! privileged **type-labels** (`ASSET`/`DATASET`/`POLICY`/`TIMESERIES`/`FUNCTION`), which //! `TypeLabels.applyLabelUpdate` on the server forces to stay exactly the node's own type — no //! update may add, remove, or swap it — while ordinary labels follow `base(set|current)+add-remove`. @@ -108,8 +108,7 @@ mod tests { } // ---------------------------------------------------------------------------------------- - // Live: label-update semantics. Ignored by default (needs a backend + mutates state). - // Run with `cargo test label_update -- --ignored --nocapture`. + // Live: label-update semantics. // ---------------------------------------------------------------------------------------- /// Sorted label list from an update/read response's first node. @@ -182,7 +181,6 @@ mod tests { /// The intrinsic type-label (ASSET here) is preserved no matter what an update tries. #[tokio::test] - #[ignore] async fn special_type_label_is_preserved() -> Result<(), Box> { let api = create_api_service(); let ext = "sdk_lblupd_type"; @@ -224,7 +222,6 @@ mod tests { /// set / add / remove each applied on their own to ordinary (non-type) labels. #[tokio::test] - #[ignore] async fn set_add_remove_individually() -> Result<(), Box> { let api = create_api_service(); let ext = "sdk_lblupd_basic"; @@ -259,7 +256,6 @@ mod tests { /// `set_and_delta_are_mutually_exclusive` test), so there is no combined-`set`+`add`/`remove` /// request to exercise live — the type system no longer lets one be built. #[tokio::test] - #[ignore] async fn add_then_remove_in_one_delta() -> Result<(), Box> { let api = create_api_service(); let ext = "sdk_lblupd_delta"; diff --git a/src/resources/tests.rs b/src/resources/tests.rs index 45be52d..b4689fe 100644 --- a/src/resources/tests.rs +++ b/src/resources/tests.rs @@ -1040,7 +1040,6 @@ async fn plain_listing_is_typed_capped_and_uncursored() -> Result<(), Box Result<(), ResponseError> { let api = create_api_service(); diff --git a/src/subscriptions/test.rs b/src/subscriptions/test.rs index 97972e4..b167563 100644 --- a/src/subscriptions/test.rs +++ b/src/subscriptions/test.rs @@ -282,20 +282,11 @@ mod tests { result } - // Regression test for the subscription-bound timeseries delete path. - // - // Deleting a timeseries that is still referenced by a subscription must surface a *terminal* - // 400 naming the blocking subscription — not an opaque 500. The backend guard - // (ResourceService.checkForSubscriptions) throws ResourceDeleteException; the fix that maps - // it to 400 lives in datahub-platform on branch - // `bugfix-500-on-subscription-bound-timeseries-delete`. - // - // #[ignore]d until that backend fix ships: against an unpatched backend this delete returns an - // empty 500 (which the SDK may also buffer/retry as a 5xx), so the assertions below would fail. - // Un-ignore once the backend change is deployed. Run with `cargo test -- --ignored`. + // Deleting a timeseries that a subscription still references is refused with a 409 + // `referenced` problem naming the subscription in `blockedBy`. It used to surface as an + // empty 500. #[tokio::test] - #[ignore] - async fn test_delete_timeseries_bound_to_subscription_returns_400( + async fn test_delete_timeseries_bound_to_subscription_is_refused( ) -> Result<(), Box> { let api_service = create_api_service(); let ts_ext = unique_id("sub_guard_ts"); @@ -318,24 +309,20 @@ mod tests { api_service.subscriptions.create(&sub).await?; let mut sub_cleanup = cleanup_subscriptions(vec![sub_ext.clone()]); - // Deleting the bound timeseries must be refused with a terminal 400, not a 500. let ts_ids = DataWrapper::from_vec(vec![IdAndExtId::from_external_id(&ts_ext)]); match api_service.time_series.delete(&ts_ids).await { Ok(_) => panic!( "expected the delete to be refused while a subscription references the timeseries" ), Err(e) => { - assert_eq!( - e.get_status(), - StatusCode::BAD_REQUEST, - "expected 400, got {}: {}", - e.get_status(), - e.get_message() - ); + assert_eq!(e.get_status(), StatusCode::CONFLICT, "{}", e.get_message()); + let problem = e.problem().expect("the refusal should be a problem document"); + assert_eq!(problem.slug(), Some("referenced"), "{}", e.get_message()); + let blockers = problem.blocked_by(); assert!( - e.get_message().to_lowercase().contains("subscription"), - "error should name the blocking subscription, got: {}", - e.get_message() + blockers.iter().any(|b| b.get("subscriptionExternalId") + .and_then(|v| v.as_str()) == Some(sub_ext.as_str())), + "blockedBy should name {sub_ext}: {blockers:?}" ); } } diff --git a/src/tests.rs b/src/tests.rs index 0e399ff..0f393b2 100644 --- a/src/tests.rs +++ b/src/tests.rs @@ -582,7 +582,6 @@ pub mod cleanup { /// the blocked-runtime half faithfully while turning a hang into a failed assertion — /// a deadlock would otherwise just stall the suite with no output. #[tokio::test] - #[ignore] async fn a_guard_tears_down_after_the_token_cache_was_cleared() { // `create_default` reads process env; only `create_api_service` loads `.env`. dotenv::dotenv().ok(); diff --git a/src/timeseries/test.rs b/src/timeseries/test.rs index 3332f93..2b5008f 100644 --- a/src/timeseries/test.rs +++ b/src/timeseries/test.rs @@ -1589,7 +1589,6 @@ mod tests { /// The only live coverage of that retry: every other path resolves before it sends. Needs a /// backend that serves the binary endpoint. #[tokio::test] - #[ignore] async fn test_insert_datapoints_binary_re_resolves_a_recreated_series() -> Result<(), Box> { let api_service = create_api_service(); let ext_id = unique_id("ts_binary_recreated");