#![allow(clippy::expect_used, clippy::unwrap_used)]
mod conformance;
use std::sync::Arc;
use graph_storage::config::{GraphStorageConfig, HopStrategy};
use graph_storage::infra::engine::PgGraphEngine;
use graph_storage::infra::storage::migrations::Migrator;
use graph_storage::infra::store::PgGraphStore;
use graph_storage_sdk::models::{Direction, HopBudget, TruncationReason};
use graph_storage_sdk::plugin_api::{ExpandRequest, GraphEngineV1, GraphStoreV1, HopBackend};
use sea_orm_migration::MigratorTrait;
use testcontainers::runners::AsyncRunner as _;
use testcontainers::{ContainerAsync, ImageExt as _};
use testcontainers_modules::postgres::Postgres;
use toolkit_db::migration_runner::run_migrations_for_testing;
use toolkit_db::secure::Db;
use toolkit_db::{ConnectOpts, connect_db};
use toolkit_security::AccessScope;
use uuid::Uuid;
static STAND_PERMITS: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(0);
fn stand_permits() -> &'static tokio::sync::Semaphore {
static ONCE: std::sync::Once = std::sync::Once::new();
ONCE.call_once(|| {
let permits = std::env::var("GEARS_TEST_PG_GRAPH_STANDS")
.ok()
.and_then(|value| value.parse::<usize>().ok())
.filter(|permits| *permits > 0)
.unwrap_or(2);
STAND_PERMITS.add_permits(permits);
});
&STAND_PERMITS
}
async fn start_server() -> Result<ContainerAsync<Postgres>, testcontainers::TestcontainersError> {
let request = match graph_image() {
Some((name, tag)) => test_containers::postgres_graph()
.with_name(name)
.with_tag(tag),
None => test_containers::postgres_graph(),
};
request
.with_env_var("POSTGRES_PASSWORD", "pass")
.with_env_var("POSTGRES_USER", "user")
.with_env_var("POSTGRES_DB", "graph")
.start()
.await
}
async fn mapped_port(container: &ContainerAsync<Postgres>) -> u16 {
let mut last = None;
for _ in 0..20 {
match container.get_host_port_ipv4(5432).await {
Ok(port) => return port,
Err(error) => {
last = Some(error);
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
}
}
}
panic!("the container never published its port: {last:?}");
}
async fn connect_with_retry(dsn: &str) -> Db {
let opts = || ConnectOpts {
max_conns: Some(4),
min_conns: Some(1),
..Default::default()
};
let mut last = None;
for _ in 0..15 {
match connect_db(dsn, opts()).await {
Ok(db) => return db,
Err(error) => {
last = Some(error);
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
}
}
}
panic!("the server never accepted a connection: {last:?}");
}
struct Stand {
store: Arc<PgGraphStore>,
engine: PgGraphEngine,
db: Arc<Db>,
dsn: String,
_container: ContainerAsync<Postgres>,
_permit: tokio::sync::SemaphorePermit<'static>,
}
async fn drop_property_graph(stand: &Stand) {
use sea_orm::ConnectionTrait as _;
let raw = sea_orm::Database::connect(&stand.dsn)
.await
.expect("a plain connection for operator surgery");
raw.execute_unprepared("DROP PROPERTY GRAPH kb")
.await
.expect("the property graph is dropped");
}
fn graph_image() -> Option<(String, String)> {
let image = std::env::var("GEARS_TEST_PG_GRAPH_IMAGE").ok()?;
let (name, tag) = image.rsplit_once(':').unwrap_or((image.as_str(), "latest"));
Some((name.to_owned(), tag.to_owned()))
}
async fn stand(hop: HopStrategy) -> Option<Stand> {
if graph_image().is_none() && !test_containers::graph_lane_required() {
eprintln!(
"GEARS_TEST_PG_GRAPH_IMAGE is unset and the platform pin ({}) has no pgvector - \
skipping the SQL/PGQ lane",
test_containers::postgres_graph_tag()
);
return None;
}
let permit = stand_permits()
.acquire()
.await
.expect("the stand semaphore is never closed");
let mut started = start_server().await;
for _ in 0..4 {
if started.is_ok() {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
started = start_server().await;
}
let container = match started {
Ok(container) => container,
Err(error) => {
assert!(
!test_containers::graph_lane_required(),
"GEARS_TEST_PG_GRAPH_REQUIRED is set but PostgreSQL 19 ({}) could not start: {error}",
test_containers::postgres_graph_tag()
);
eprintln!("PostgreSQL 19 unavailable - skipping the SQL/PGQ lane: {error}");
return None;
}
};
let port = mapped_port(&container).await;
let dsn = format!("postgres://user:pass@127.0.0.1:{port}/graph");
let db = connect_with_retry(&dsn).await;
if let Err(error) = run_migrations_for_testing(&db, Migrator::migrations()).await {
let no_pgvector = error
.to_string()
.contains("extension \"vector\" is not available");
assert!(
!(no_pgvector && test_containers::graph_lane_required()),
"GEARS_TEST_PG_GRAPH_REQUIRED is set but the graph image has no pgvector: {error}"
);
if no_pgvector {
eprintln!(
"the graph image has no pgvector - skipping; set GEARS_TEST_PG_GRAPH_IMAGE \
to an image with PostgreSQL 19 and pgvector"
);
return None;
}
panic!("migrations apply: {error}");
}
let config = GraphStorageConfig {
traversal_hop: hop,
ontology_max_chain_depth: 8,
..GraphStorageConfig::default()
};
let db = Arc::new(db);
let pgq_available = graph_storage::infra::engine::probe_pgq(&db).await;
let store = Arc::new(PgGraphStore::new(
Arc::clone(&db),
config.validated().expect("the test configuration is valid"),
pgq_available,
));
let engine = PgGraphEngine::new(Arc::clone(&store));
Some(Stand {
store,
engine,
db,
dsn,
_container: container,
_permit: permit,
})
}
async fn tenant_on(stand: &Stand) -> Uuid {
let tenant = Uuid::now_v7();
let scope = AccessScope::for_tenant(tenant);
graph_storage::infra::store::ingest::ensure_meta(stand.store.as_ref(), tenant, &scope)
.await
.expect("meta rows exist");
tenant
}
async fn store_with(stand: &Stand, options: &str) -> Arc<PgGraphStore> {
let encoded = options.replace(' ', "%20").replace('=', "%3D");
let dsn = format!("{}?options={encoded}", stand.dsn);
let db = connect_db(
&dsn,
ConnectOpts {
max_conns: Some(2),
min_conns: Some(1),
..Default::default()
},
)
.await
.unwrap_or_else(|error| panic!("a second pool with `{options}` connects: {error}"));
let db = Arc::new(db);
let pgq = graph_storage::infra::engine::probe_pgq(&db).await;
Arc::new(PgGraphStore::new(
db,
GraphStorageConfig::default()
.validated()
.expect("the default configuration is valid"),
pgq,
))
}
#[tokio::test]
async fn a_filtered_vector_search_under_returns_without_iterative_scan() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let crowd = tenant_on(&stand).await;
let alone = tenant_on(&stand).await;
for (tenant, keys) in [
(
crowd,
(0..60).map(|i| format!("crowd-{i}")).collect::<Vec<_>>(),
),
(alone, vec!["the-only-one".to_owned()]),
] {
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("the ontology registers");
let nodes: Vec<_> = keys
.iter()
.map(|key| conformance::summarized(key, key, key))
.collect();
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(nodes, Vec::new()),
)
.await
.expect("the fixture commits");
}
let scope = AccessScope::for_tenant(alone);
let ctx = conformance::ctx(alone, &scope, None);
let epoch = conformance::EPOCH;
let probe = "a probe that names nothing in particular";
let strict = store_with(&stand, "-c enable_seqscan=off -c hnsw.ef_search=1").await;
let missed = conformance::search_vector(strict.as_ref(), &ctx, probe, epoch).await;
let iterative = store_with(
&stand,
"-c enable_seqscan=off -c hnsw.ef_search=1 -c hnsw.iterative_scan=relaxed_order",
)
.await;
let found = conformance::search_vector(iterative.as_ref(), &ctx, probe, epoch).await;
assert!(
missed.is_empty(),
"the small tenant's vector is not among the candidates the scan spent \
its budget on: {missed:?}"
);
assert_eq!(
found,
vec!["the-only-one".to_owned()],
"iterative scanning keeps going until the filter has something to \
return, which is the difference between a recall setting and a \
correctness one"
);
}
macro_rules! pg_case {
($name:ident, $case:path) => {
#[tokio::test]
async fn $name() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
$case(stand.store.as_ref(), tenant).await;
}
};
}
pg_case!(a_failed_batch_commits_nothing, conformance::batch_atomicity);
pg_case!(
source_generations_are_fenced_monotonically,
conformance::generation_fencing
);
pg_case!(
a_node_never_outlives_its_incident_edges,
conformance::no_orphan_edges
);
pg_case!(
a_fresh_tenant_reports_a_usable_revision,
conformance::a_fresh_tenant_reports_a_usable_revision
);
pg_case!(
materializing_a_phantom_revalidates_its_edges,
conformance::materializing_a_phantom_revalidates_its_edges
);
pg_case!(
an_edge_type_refuses_an_endpoint_it_does_not_admit,
conformance::endpoint_constraints_are_enforced
);
pg_case!(a_recorded_idempotency_key_replays, conformance::idempotency);
pg_case!(
a_keyless_retry_is_a_new_request,
conformance::a_keyless_retry_is_a_new_request
);
pg_case!(
a_document_is_retrieved_by_its_own_text,
conformance::a_document_is_retrieved_by_its_own_text
);
pg_case!(
a_declared_path_reaches_the_vector,
conformance::a_declared_path_reaches_the_vector
);
pg_case!(
a_skipped_re_ingest_preserves_the_vector,
conformance::a_skipped_re_ingest_preserves_the_vector
);
pg_case!(
a_stale_vector_stops_ranking_but_the_node_stays,
conformance::a_stale_vector_stops_ranking_but_the_node_stays
);
pg_case!(
only_the_active_epoch_ranks,
conformance::only_the_active_epoch_ranks
);
pg_case!(an_identical_batch_converges, conformance::convergent_replay);
pg_case!(
a_same_key_ingest_may_not_change_the_type,
conformance::a_same_key_ingest_may_not_change_the_type
);
pg_case!(
per_item_outcomes_follow_the_batch_order,
conformance::per_item_outcomes_follow_the_batch_order
);
pg_case!(
a_scope_and_an_idempotency_key_belong_to_their_producer,
conformance::a_scope_and_an_idempotency_key_belong_to_their_producer
);
pg_case!(
a_type_pattern_narrows_search_and_a_hop,
conformance::a_type_pattern_narrows_search_and_a_hop
);
pg_case!(
hybrid_search_fuses_both_arms,
conformance::hybrid_search_fuses_both_arms
);
pg_case!(
deleting_an_already_tombstoned_row_is_a_no_op,
conformance::deleting_an_already_tombstoned_row_is_a_no_op
);
pg_case!(
the_type_catalogue_pages_through_its_own_cursor,
conformance::the_type_catalogue_pages_through_its_own_cursor
);
pg_case!(
a_deleted_conclusion_stops_pinning_its_endpoint,
conformance::a_deleted_conclusion_stops_pinning_its_endpoint
);
pg_case!(
a_batch_that_names_one_type_twice_is_refused,
conformance::a_batch_that_names_one_type_twice_is_refused
);
pg_case!(
an_unchanged_re_ingest_embeds_nothing,
conformance::an_unchanged_re_ingest_embeds_nothing
);
pg_case!(
tombstoned_rows_are_invisible,
conformance::tombstones_are_invisible
);
pg_case!(
a_denied_row_reads_like_an_absent_one,
conformance::denied_is_indistinguishable_from_absent
);
pg_case!(
search_applies_the_scope_inside_the_statement,
conformance::search_is_scoped
);
pg_case!(
the_envelope_records_the_subject_of_each_verb,
conformance::the_envelope_records_the_subject_of_each_verb
);
pg_case!(
a_projection_row_carries_the_envelope,
conformance::a_projection_row_carries_the_envelope
);
pg_case!(
a_declared_payload_path_filters_and_orders_the_projection,
conformance::a_declared_payload_path_filters_and_orders_the_projection
);
pg_case!(
an_undeclared_payload_path_is_refused_naming_the_alternatives,
conformance::an_undeclared_payload_path_is_refused_naming_the_alternatives
);
pg_case!(
an_index_path_onto_a_non_scalar_is_refused_at_registration,
conformance::an_index_path_onto_a_non_scalar_is_refused_at_registration
);
pg_case!(
a_deeper_chain_registers_and_its_ancestor_admits_the_leaf,
conformance::a_deeper_chain_registers_and_its_ancestor_admits_the_leaf
);
#[tokio::test]
async fn readiness_reports_every_capability_and_only_some_block_service() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::readiness_reports_every_capability_and_only_some_block_service(
stand.store.as_ref(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_replacements_of_one_scope_serialize() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_replacements_of_one_scope_serialize(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn every_update_of_a_node_advances_its_version() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::every_update_of_a_node_advances_its_version(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_writers_with_one_expected_version_do_not_both_win() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_writers_with_one_expected_version_do_not_both_win(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_type_updates_do_not_share_one_revision() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_type_updates_do_not_share_one_revision(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_delete_racing_an_upsert_leaves_no_rewritten_tombstone() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::a_delete_racing_an_upsert_leaves_no_rewritten_tombstone(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_replacement_does_not_delete_what_an_ingest_moved_out_of_its_scope() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::a_replacement_does_not_delete_what_an_ingest_moved_out_of_its_scope(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_batches_naming_one_new_endpoint_both_land() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_batches_naming_one_new_endpoint_both_land(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_node_delete_racing_its_edge_delete_counts_every_edge_once() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::a_node_delete_racing_its_edge_delete_counts_every_edge_once(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_deletes_of_one_node_tombstone_it_once() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_deletes_of_one_node_tombstone_it_once(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_deleted_edge_is_revived_by_the_next_upsert() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::a_deleted_edge_is_revived_by_the_next_upsert(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn every_committed_mutation_gets_its_own_revision() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::every_committed_mutation_gets_its_own_revision(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
pg_case!(
scope_replacement_removes_an_edge_whose_endpoints_remain,
conformance::scope_replacement_removes_an_edge_whose_endpoints_remain
);
#[tokio::test]
async fn a_catalogue_scan_with_no_progress_refuses_rather_than_answering_empty() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let live = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&live, conformance::ontology_batch())
.await
.expect("the ontology registers");
let spent = graph_storage_sdk::plugin_api::StoreCtx {
budget: graph_storage_sdk::models::RemainingBudget::starting_now(std::time::Duration::ZERO),
..conformance::ctx(tenant, &scope, None)
};
let refused = stand
.store
.list_types(&spent, graph_storage_sdk::models::TypeQuery::default())
.await
.expect_err("a scan that read nothing has nothing to hand back");
assert!(
matches!(
refused,
graph_storage_sdk::plugin_api::GraphStoreError::Deadline
),
"expected a deadline refusal, got {refused:?}"
);
let page = stand
.store
.list_types(&live, graph_storage_sdk::models::TypeQuery::default())
.await
.expect("a live budget lists the catalogue");
assert!(!page.items.is_empty(), "the fixture is there to be listed");
}
pg_case!(
an_edge_is_not_taken_from_the_scope_that_declared_it,
conformance::an_edge_is_not_taken_from_the_scope_that_declared_it
);
pg_case!(
one_replacement_does_not_take_another_scopes_edges,
conformance::one_replacement_does_not_take_another_scopes_edges
);
pg_case!(
a_store_without_labels_refuses_every_label_call,
conformance::a_store_without_labels_refuses_every_label_call
);
pg_case!(
an_edge_cannot_name_a_tombstoned_endpoint,
conformance::an_edge_cannot_name_a_tombstoned_endpoint
);
pg_case!(
node_types_answers_the_live_nodes_it_is_asked_about,
conformance::node_types_answers_the_live_nodes_it_is_asked_about
);
pg_case!(
an_edge_only_scope_drops_the_edges_it_stops_declaring,
conformance::an_edge_only_scope_drops_the_edges_it_stops_declaring
);
pg_case!(
scope_replacement_removes_what_the_batch_no_longer_names,
conformance::scope_replacement_removes_what_the_batch_no_longer_names
);
pg_case!(
scope_replacement_preserves_analysis_edges_and_their_endpoints,
conformance::scope_replacement_preserves_analysis_edges_and_their_endpoints
);
pg_case!(
a_source_namespace_is_claimed_by_its_first_writer,
conformance::a_source_namespace_is_claimed_by_its_first_writer
);
pg_case!(
writing_under_another_producers_namespace_is_forbidden,
conformance::writing_under_another_producers_namespace_is_forbidden
);
pg_case!(
a_transfer_moves_the_namespace_and_records_who_moved_it,
conformance::a_transfer_moves_the_namespace_and_records_who_moved_it
);
pg_case!(
an_owned_nodes_source_field_claims_no_namespace,
conformance::an_owned_nodes_source_field_claims_no_namespace
);
pg_case!(
a_backward_compatible_change_updates_the_type_in_place,
conformance::a_backward_compatible_change_updates_the_type_in_place
);
pg_case!(
an_incompatible_change_is_refused_with_its_location,
conformance::an_incompatible_change_is_refused_with_its_location
);
pg_case!(
a_changed_schema_is_still_a_conflict_by_default,
conformance::a_changed_schema_is_still_a_conflict_by_default
);
pg_case!(
a_dry_run_reports_every_verdict_and_writes_nothing,
conformance::a_dry_run_reports_every_verdict_and_writes_nothing
);
pg_case!(
a_change_the_schemas_cannot_prove_is_admitted_when_the_rows_fit,
conformance::a_change_the_schemas_cannot_prove_is_admitted_when_the_rows_fit
);
pg_case!(
a_change_the_stored_rows_contradict_is_refused_naming_them,
conformance::a_change_the_stored_rows_contradict_is_refused_naming_them
);
pg_case!(
a_migration_moves_the_data_with_the_type,
conformance::a_migration_moves_the_data_with_the_type
);
pg_case!(
a_migration_that_leaves_rows_invalid_is_refused_naming_them,
conformance::a_migration_that_leaves_rows_invalid_is_refused_naming_them
);
pg_case!(
a_migration_without_a_schema_change_is_refused,
conformance::a_migration_without_a_schema_change_is_refused
);
pg_case!(
a_migration_stamps_its_writer_and_moves_the_version,
conformance::a_migration_stamps_its_writer_and_moves_the_version
);
pg_case!(
an_accepted_type_update_advances_the_graph_revision,
conformance::an_accepted_type_update_advances_the_graph_revision
);
pg_case!(
a_new_index_path_becomes_filterable_without_recreating_the_type,
conformance::a_new_index_path_becomes_filterable_without_recreating_the_type
);
#[tokio::test]
async fn a_payload_ordered_projection_pages_by_keyset() {
use toolkit_odata::SortDir;
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let store = stand.store.as_ref();
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
let whole = store
.project_table(
&ctx,
conformance::projection_seeded(store, &ctx, &[("payload/score", SortDir::Desc)]).await,
)
.await
.expect("projection succeeds");
let expected: Vec<String> = whole.items.iter().map(|r| r.node_key.clone()).collect();
assert_eq!(expected, vec!["t1", "t5", "t2", "t3", "t4"]);
let mut walked = Vec::new();
let mut request = conformance::projection(
&[conformance::INDEXED],
"",
&[("payload/score", SortDir::Desc)],
);
request.query = request.query.with_limit(2);
loop {
let page = store
.project_table(&ctx, request.clone())
.await
.expect("a page is served");
assert!(page.items.len() <= 2, "the page honours its limit");
walked.extend(page.items.iter().map(|r| r.node_key.clone()));
let Some(token) = page.page_info.next_cursor else {
break;
};
let cursor = toolkit_odata::CursorV1::decode(&token).expect("a CursorV1 token");
assert_eq!(cursor.s, "-payload/score,+node_key");
request.query = toolkit_odata::ODataQuery::new()
.with_limit(2)
.with_cursor(cursor);
assert!(walked.len() <= 5, "the walk terminates");
}
assert_eq!(walked, expected, "pages concatenate to the one-page answer");
}
#[tokio::test]
async fn colliding_node_keys_stay_inside_their_tenants() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let one = tenant_on(&stand).await;
let two = tenant_on(&stand).await;
conformance::tenant_isolation(stand.store.as_ref(), one, two).await;
}
#[tokio::test]
async fn the_pattern_hop_walks_the_graph() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(
vec![
conformance::node("a", "a"),
conformance::node("b", "b"),
conformance::node("c", "c"),
],
vec![conformance::edge("a", "b"), conformance::edge("b", "c")],
),
)
.await
.expect("the batch commits");
let ids = stand
.store
.resolve_node_ids(&ctx, &["a".to_owned()])
.await
.expect("resolution succeeds");
let seed = ids.first().expect("`a` resolves").1;
let response = stand
.engine
.expand(
&ctx,
ExpandRequest {
frontier: vec![seed],
direction: Direction::Outgoing,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 100,
max_edges_scanned: 1_000,
},
with_degrees: false,
},
)
.await
.expect("the pattern hop runs");
assert_eq!(
response.served_by,
HopBackend::Pattern,
"the single-statement pattern must be what answered, not the fallback"
);
let reached: Vec<String> = response.edges.iter().map(|e| e.dst.clone()).collect();
assert_eq!(reached, vec!["b".to_owned()], "one hop reaches exactly `b`");
assert!(response.truncated.is_none());
}
async fn seed_and_expand(
stand: &Stand,
direction: Direction,
) -> (Vec<i64>, Vec<String>, HopBackend) {
let tenant = tenant_on(stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(
vec![
conformance::node("hub", "hub"),
conformance::node("spoke-1", "one"),
conformance::node("spoke-2", "two"),
],
vec![
conformance::edge("hub", "spoke-1"),
conformance::edge("spoke-2", "hub"),
],
),
)
.await
.expect("the batch commits");
let seed = stand
.store
.resolve_node_ids(&ctx, &["hub".to_owned()])
.await
.expect("resolution succeeds")
.first()
.expect("`hub` resolves")
.1;
let response = stand
.engine
.expand(
&ctx,
ExpandRequest {
frontier: vec![seed],
direction,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 100,
max_edges_scanned: 1_000,
},
with_degrees: false,
},
)
.await
.expect("the hop runs");
let mut keys: Vec<String> = response
.edges
.iter()
.flat_map(|e| [e.src.clone(), e.dst.clone()])
.collect();
keys.sort();
(response.reached, keys, response.served_by)
}
#[tokio::test]
async fn both_hop_backends_return_the_same_answer() {
let Some(pgq) = stand(HopStrategy::Pgq).await else {
return;
};
let Some(two_query) = stand(HopStrategy::TwoQuery).await else {
return;
};
for direction in [Direction::Outgoing, Direction::Incoming, Direction::Either] {
let (pattern_reached, pattern_keys, pattern_backend) =
seed_and_expand(&pgq, direction).await;
let (fallback_reached, fallback_keys, fallback_backend) =
seed_and_expand(&two_query, direction).await;
assert_eq!(
(pattern_backend, fallback_backend),
(HopBackend::Pattern, HopBackend::TwoQuery),
"the {direction:?} comparison must be between two different backends"
);
assert_eq!(
pattern_reached.len(),
fallback_reached.len(),
"the backends disagree on how many nodes {direction:?} reaches"
);
assert_eq!(
pattern_keys, fallback_keys,
"the backends disagree on the {direction:?} edges"
);
}
}
#[tokio::test]
async fn a_hop_never_leaves_its_tenant() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let ours = tenant_on(&stand).await;
let theirs = tenant_on(&stand).await;
for (tenant, far) in [(ours, "ours-far"), (theirs, "theirs-far")] {
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(
vec![
conformance::node("shared-key", "start"),
conformance::node(far, far),
],
vec![conformance::edge("shared-key", far)],
),
)
.await
.expect("the batch commits");
}
let their_scope = AccessScope::for_tenant(theirs);
let their_ctx = conformance::ctx(theirs, &their_scope, None);
let theirs_view = stand
.store
.get_node(&their_ctx, &"shared-key".to_owned(), 10)
.await
.expect("the other tenant owns the same key");
assert_eq!(
theirs_view.adjacency.len(),
1,
"the trap fixture must have an edge for the leak to expose"
);
let our_scope = AccessScope::for_tenant(ours);
let our_ctx = conformance::ctx(ours, &our_scope, None);
let seed = stand
.store
.resolve_node_ids(&our_ctx, &["shared-key".to_owned()])
.await
.expect("resolution succeeds")
.first()
.expect("our key resolves")
.1;
let response = stand
.engine
.expand(
&our_ctx,
ExpandRequest {
frontier: vec![seed],
direction: Direction::Either,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 100,
max_edges_scanned: 1_000,
},
with_degrees: false,
},
)
.await
.expect("the hop runs");
let keys: Vec<String> = response.edges.iter().map(|e| e.dst.clone()).collect();
assert_eq!(
keys,
vec!["ours-far".to_owned()],
"the hop must reach only our own far node"
);
}
#[tokio::test]
async fn a_stopped_hop_reports_why() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
let response = stand
.engine
.expand(
&ctx,
ExpandRequest {
frontier: vec![1, 2, 3],
direction: Direction::Either,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 1,
max_edges_scanned: 10,
},
with_degrees: false,
},
)
.await
.expect("the hop runs");
assert_eq!(
response.truncated,
Some(TruncationReason::FrontierCap),
"a frontier over the cap must be reported, not silently trimmed"
);
}
#[tokio::test]
async fn asking_for_degrees_does_not_double_the_hop_edge_budget() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
let mut nodes = vec![conformance::node("hub", "hub")];
let mut edges = Vec::new();
for index in 0..8 {
let key = format!("spoke-{index}");
nodes.push(conformance::node(&key, &key));
edges.push(conformance::edge("hub", &key));
}
conformance::ingest_batch(stand.store.as_ref(), &ctx, conformance::batch(nodes, edges))
.await
.expect("the hub commits");
let ids = stand
.store
.resolve_node_ids(&ctx, &["hub".to_owned()])
.await
.expect("the hub resolves");
let frontier: Vec<_> = ids.into_iter().map(|(_, id)| id).collect();
let request = |with_degrees| ExpandRequest {
frontier: frontier.clone(),
direction: Direction::Either,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 1_000,
max_edges_scanned: 4,
},
with_degrees,
};
let plain = stand
.engine
.expand(&ctx, request(false))
.await
.expect("the hop runs");
let with_degrees = stand
.engine
.expand(&ctx, request(true))
.await
.expect("the hop runs with degrees");
assert_eq!(
plain.edges.len(),
with_degrees.edges.len(),
"the same budget reads the same number of edges whoever is asking"
);
assert_eq!(
with_degrees.truncated,
Some(TruncationReason::EdgeScanCap),
"and the hop still reports that its scan was cut short"
);
assert_eq!(
with_degrees.degrees.len(),
with_degrees.reached.len(),
"degrees stay index-aligned with the reached set even when unknown"
);
assert!(
with_degrees.degrees.iter().all(|degree| *degree == 0),
"the first scan spent the hop's allowance, so the degrees are unknown \
rather than bought with a second one: {:?}",
with_degrees.degrees
);
}
#[tokio::test]
async fn a_hop_that_hits_its_edge_budget_reports_it_on_both_backends() {
for hop in [HopStrategy::Pgq, HopStrategy::TwoQuery] {
let Some(stand) = stand(hop).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
let mut nodes = vec![conformance::node("hub", "hub")];
let mut edges = Vec::new();
for index in 0..6 {
let key = format!("spoke-{index}");
nodes.push(conformance::node(&key, &key));
edges.push(conformance::edge("hub", &key));
}
conformance::ingest_batch(stand.store.as_ref(), &ctx, conformance::batch(nodes, edges))
.await
.expect("the hub commits");
let ids = stand
.store
.resolve_node_ids(&ctx, &["hub".to_owned()])
.await
.expect("the hub resolves");
let frontier: Vec<_> = ids.into_iter().map(|(_, id)| id).collect();
let request = |max_edges_scanned| ExpandRequest {
frontier: frontier.clone(),
direction: Direction::Either,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 1_000,
max_edges_scanned,
},
with_degrees: false,
};
let cut = stand
.engine
.expand(&ctx, request(3))
.await
.expect("the hop runs");
assert_eq!(
cut.truncated,
Some(TruncationReason::EdgeScanCap),
"{hop:?}: a scan stopped by the edge budget must say so"
);
assert_eq!(
cut.edges.len(),
3,
"{hop:?}: the budget still bounds the work"
);
let whole = stand
.engine
.expand(&ctx, request(100))
.await
.expect("the hop runs");
assert_eq!(whole.truncated, None, "{hop:?}: an unbounded hop is whole");
assert_eq!(whole.edges.len(), 6, "{hop:?}: every edge of the hub");
}
}
#[tokio::test]
async fn a_hop_reports_the_reached_nodes_degree_only_when_asked() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(
vec![
conformance::node("root", "root"),
conformance::node("core", "core"),
conformance::node("leaf", "leaf"),
conformance::node("far-1", "far-1"),
conformance::node("far-2", "far-2"),
],
vec![
conformance::edge("root", "core"),
conformance::edge("root", "leaf"),
conformance::edge("core", "far-1"),
conformance::edge("core", "far-2"),
],
),
)
.await
.expect("the fixture commits");
let ids = stand
.store
.resolve_node_ids(&ctx, &["root".to_owned(), "core".to_owned()])
.await
.expect("the root resolves");
let by_key: std::collections::BTreeMap<String, i64> = ids.into_iter().collect();
let frontier = vec![by_key["root"]];
let request = |with_degrees| ExpandRequest {
frontier: frontier.clone(),
direction: Direction::Either,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 1_000,
max_edges_scanned: 1_000,
},
with_degrees,
};
let silent = stand
.engine
.expand(&ctx, request(false))
.await
.expect("the hop runs");
assert!(
silent.degrees.is_empty(),
"a hop that was not asked for degrees does not pay for them"
);
let ranked = stand
.engine
.expand(&ctx, request(true))
.await
.expect("the hop runs");
assert_eq!(ranked.degrees.len(), ranked.reached.len(), "index-aligned");
let degree_of_core = ranked
.reached
.iter()
.zip(&ranked.degrees)
.find(|(id, _)| **id == by_key["core"])
.map(|(_, degree)| *degree);
assert_eq!(
degree_of_core,
Some(3),
"`core` has three live edges: one to the root and two of its own"
);
let leaf_degree = ranked
.reached
.iter()
.zip(&ranked.degrees)
.filter(|(id, _)| **id != by_key["core"])
.map(|(_, degree)| *degree)
.collect::<Vec<_>>();
assert_eq!(
leaf_degree,
vec![1],
"the leaf has only the edge that reached it"
);
}
#[tokio::test]
async fn a_property_graph_lost_after_boot_is_reported_and_not_substituted_on_demand() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(
vec![conformance::node("p-a", "a"), conformance::node("p-b", "b")],
vec![conformance::edge("p-a", "p-b")],
),
)
.await
.expect("the batch commits");
let seed = stand
.store
.resolve_node_ids(&ctx, &["p-a".to_owned()])
.await
.expect("resolution succeeds")
.first()
.map_or_else(|| panic!("`p-a` resolves"), |(_, id)| *id);
assert!(
graph_storage::infra::engine::probe_pgq(stand.store.db()).await,
"the stand must start with a working property graph for this test to mean anything"
);
let hop = || ExpandRequest {
frontier: vec![seed],
direction: Direction::Outgoing,
edge_types: None,
labels: None,
budget: HopBudget {
max_frontier: 100,
max_edges_scanned: 1_000,
},
with_degrees: false,
};
let served = stand
.engine
.expand(&ctx, hop())
.await
.expect("with the property graph present the pattern answers");
assert_eq!(
served.served_by,
graph_storage_sdk::plugin_api::HopBackend::Pattern,
"precondition: the pattern hop answers before the graph is dropped"
);
assert_eq!(
sqlpgq_row(stand.store.as_ref()).await.state,
graph_storage_sdk::models::ReadinessState::Healthy
);
drop_property_graph(&stand).await;
assert!(
!graph_storage::infra::engine::probe_pgq(stand.store.db()).await,
"the probe must report the capability as absent once the graph is gone"
);
let refused = stand
.engine
.expand(&ctx, hop())
.await
.err()
.expect("`pgq` is a demand: a pattern that stopped executing is not quietly replaced");
assert!(
matches!(refused, graph_storage_sdk::plugin_api::GraphEngineError::Unavailable { ref reason } if reason.contains("traversal_hop")),
"the refusal names the setting, got {refused:?}"
);
if let graph_storage_sdk::plugin_api::GraphEngineError::Unavailable { reason } = &refused {
assert!(
reason.contains("the reason is in the gear's log")
&& !reason.contains("does not exist"),
"the refusal does not carry the server's diagnostic: {reason}"
);
}
let row = sqlpgq_row(stand.store.as_ref()).await;
assert_eq!(
row.state,
graph_storage_sdk::models::ReadinessState::Unhealthy,
"readiness reports the loss from the request that found it: {row:?}"
);
let preferring = Arc::new(PgGraphStore::new(
Arc::clone(&stand.db),
GraphStorageConfig {
traversal_hop: HopStrategy::Auto,
..GraphStorageConfig::default()
}
.validated()
.expect("the test configuration is valid"),
true,
));
assert_eq!(
sqlpgq_row(preferring.as_ref()).await.state,
graph_storage_sdk::models::ReadinessState::Healthy,
"until a request meets the loss, the store believes the probe"
);
let response = PgGraphEngine::new(Arc::clone(&preferring))
.expand(&ctx, hop())
.await
.expect("`auto` is a preference, and the fallback backend answers");
assert_eq!(
response.served_by,
graph_storage_sdk::plugin_api::HopBackend::TwoQuery
);
let reached: Vec<String> = response.edges.iter().map(|e| e.dst.clone()).collect();
assert_eq!(
reached,
vec!["p-b".to_owned()],
"the fallback backend answers the same question the pattern would have"
);
assert_eq!(
sqlpgq_row(preferring.as_ref()).await.state,
graph_storage_sdk::models::ReadinessState::Degraded,
"the loss is reported once a request has met it"
);
let probed = graph_storage::infra::engine::probe_pgq(stand.store.db()).await;
let store_for = |hop: HopStrategy| {
Arc::new(PgGraphStore::new(
Arc::clone(&stand.db),
GraphStorageConfig {
traversal_hop: hop,
..GraphStorageConfig::default()
}
.validated()
.expect("the test configuration is valid"),
probed,
))
};
let preferred = store_for(HopStrategy::Auto);
let again = PgGraphEngine::new(Arc::clone(&preferred))
.expand(&ctx, hop())
.await
.expect("`auto` is a preference, and the fallback backend answers");
assert_eq!(again.edges.len(), 1);
let row = sqlpgq_row(preferred.as_ref()).await;
assert_eq!(
row.state,
graph_storage_sdk::models::ReadinessState::Degraded,
"{row:?}"
);
let demanded = store_for(HopStrategy::Pgq);
let refused = PgGraphEngine::new(Arc::clone(&demanded))
.expand(&ctx, hop())
.await
.err()
.expect("`pgq` is a demand, and another backend is not substituted for it");
assert!(
matches!(refused, graph_storage_sdk::plugin_api::GraphEngineError::Unavailable { ref reason } if reason.contains("traversal_hop")),
"the refusal names the setting to change, got {refused:?}"
);
let row = sqlpgq_row(demanded.as_ref()).await;
assert_eq!(
row.state,
graph_storage_sdk::models::ReadinessState::Unhealthy,
"{row:?}"
);
assert!(
!graph_storage_sdk::models::Readiness::of(vec![row]).ready,
"an explicitly configured backend the server cannot provide leaves the gear not ready"
);
let chosen = store_for(HopStrategy::TwoQuery);
let row = sqlpgq_row(chosen.as_ref()).await;
assert_eq!(
row.state,
graph_storage_sdk::models::ReadinessState::Healthy,
"{row:?}"
);
}
async fn sqlpgq_row(store: &PgGraphStore) -> graph_storage_sdk::models::ComponentReadiness {
store
.probe_readiness()
.await
.into_iter()
.find(|row| row.component == graph_storage_sdk::models::SQLPGQ)
.expect("the store reports its SQL/PGQ row")
}
#[tokio::test]
async fn the_built_in_store_declines_the_snapshot_obligation() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let tenant = tenant_on(&stand).await;
let scope = AccessScope::for_tenant(tenant);
let ctx = conformance::ctx(tenant, &scope, None);
assert!(
!stand.store.capabilities().snapshots,
"the store must declare the capability absent rather than claim it"
);
stand
.store
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("ontology registers");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(vec![conformance::node("snap-before", "before")], Vec::new()),
)
.await
.expect("the first batch commits");
let snapshot = stand.store.begin_read(&ctx).await.expect("snapshot opens");
conformance::ingest_batch(
stand.store.as_ref(),
&ctx,
conformance::batch(vec![conformance::node("snap-after", "after")], Vec::new()),
)
.await
.expect("the concurrent batch commits");
let under = conformance::ctx(tenant, &scope, Some(&snapshot));
let seen = stand
.store
.resolve_node_ids(&under, &["snap-after".to_owned()])
.await
.expect("resolution succeeds");
assert!(
!seen.is_empty(),
"the built-in store does not isolate a compound read: this asserts the \
*absence* of isolation, so if it ever starts isolating, revisit \
StoreCapabilities::snapshots and DESIGN section 3.3 together"
);
assert_eq!(
snapshot.revision.revision + 1,
stand
.store
.revision(&ctx)
.await
.expect("revision reads")
.revision,
"the snapshot recorded the revision it opened at, even though it does \
not hold it"
);
stand
.store
.end_read(snapshot)
.await
.expect("snapshot closes");
}
pg_case!(
both_node_families_and_both_edge_families_round_trip,
conformance::both_node_families_and_both_edge_families_round_trip
);
pg_case!(
an_edge_read_carries_the_envelope,
conformance::an_edge_read_carries_the_envelope
);
#[tokio::test]
async fn an_edge_whose_endpoint_is_hidden_is_not_readable() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let one = tenant_on(&stand).await;
let two = tenant_on(&stand).await;
conformance::an_edge_whose_endpoint_is_hidden_is_not_readable(stand.store.as_ref(), one, two)
.await;
}
#[tokio::test]
async fn no_read_surface_answers_with_another_tenants_rows() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let one = tenant_on(&stand).await;
let two = tenant_on(&stand).await;
conformance::no_read_surface_answers_with_another_tenants_rows(stand.store.as_ref(), one, two)
.await;
}
use graph_storage::infra::store::spaces::{SpaceResolution, resolve};
fn space(model: &str) -> graph_storage_sdk::models::EmbeddingSpaceId {
graph_storage_sdk::models::EmbeddingSpaceId::new(
model,
"tokenizer-v1",
serde_json::json!({"lowercase": true}),
serde_json::json!({"mode": "mean"}),
serde_json::json!({"l2": true}),
8,
)
}
#[tokio::test]
async fn a_second_boot_with_another_provider_reports_a_mismatch_rather_than_opening_an_epoch() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let first = resolve(&stand.db, &space("model-a"))
.await
.expect("the first boot resolves");
let SpaceResolution::Active { epoch } = first else {
panic!("a first boot adopts the provider's own space, got {first:?}");
};
assert_eq!(
resolve(&stand.db, &space("model-a"))
.await
.expect("a same-provider boot resolves"),
SpaceResolution::Active { epoch },
"the same provider adopts the recorded epoch rather than opening another"
);
match resolve(&stand.db, &space("model-b"))
.await
.expect("a different-provider boot still resolves")
{
SpaceResolution::Mismatched {
recorded_identity,
recorded_epoch,
} => {
assert_eq!(recorded_epoch, epoch);
assert_eq!(recorded_identity, space("model-a").identity_hash);
}
other @ SpaceResolution::Active { .. } => {
panic!("a different provider must not be adopted silently, got {other:?}")
}
}
}
pg_case!(
an_edge_type_evolves_over_its_own_rows,
conformance::an_edge_type_evolves_over_its_own_rows
);
#[tokio::test]
async fn an_elided_insert_that_asks_for_its_row_reports_record_not_found() {
use graph_storage::infra::storage::entity::graph_meta;
use sea_orm::{ActiveValue, EntityTrait as _, sea_query::OnConflict};
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
let raw = sea_orm::Database::connect(&stand.dsn)
.await
.expect("a plain connection to the stand");
let tenant = Uuid::now_v7();
let row = || graph_meta::ActiveModel {
tenant_id: ActiveValue::Set(tenant),
key: ActiveValue::Set("elision-probe".to_owned()),
value: ActiveValue::Set(serde_json::json!(0)),
};
let do_nothing = || {
OnConflict::columns([graph_meta::Column::TenantId, graph_meta::Column::Key])
.do_nothing()
.to_owned()
};
graph_meta::Entity::insert(row())
.on_conflict(do_nothing())
.exec_with_returning(&raw)
.await
.expect("the first insert lands and returns its row");
let elided = graph_meta::Entity::insert(row())
.on_conflict(do_nothing())
.exec_with_returning(&raw)
.await;
assert!(
matches!(elided, Err(sea_orm::DbErr::RecordNotFound(_))),
"an elided insert with RETURNING must report RecordNotFound, got {elided:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_scopes_racing_to_claim_an_unowned_edge_leave_it_with_one() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_scopes_racing_to_claim_an_unowned_edge_leave_it_with_one(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}
pg_case!(
an_edge_does_not_follow_a_node_ingested_under_a_new_key,
conformance::an_edge_does_not_follow_a_node_ingested_under_a_new_key
);
pg_case!(
an_expected_version_on_an_absent_key_is_a_conflict_unless_it_is_zero,
conformance::an_expected_version_on_an_absent_key_is_a_conflict_unless_it_is_zero
);
pg_case!(
a_batch_with_one_conflicting_type_registers_none_and_names_it,
conformance::a_batch_with_one_conflicting_type_registers_none_and_names_it
);
pg_case!(
a_replacement_leaves_alone_what_does_not_carry_its_attribute,
conformance::a_replacement_leaves_alone_what_does_not_carry_its_attribute
);
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_creators_with_expected_version_zero_do_not_both_win() {
let Some(stand) = stand(HopStrategy::Pgq).await else {
return;
};
conformance::two_creators_with_expected_version_zero_do_not_both_win(
std::sync::Arc::clone(&stand.store) as std::sync::Arc<dyn GraphStoreV1>,
Uuid::now_v7(),
)
.await;
}