#![allow(clippy::expect_used, clippy::unwrap_used)]
#[allow(dead_code)]
mod conformance;
mod support;
use std::sync::Arc;
use graph_storage::config::GraphStorageConfig;
use graph_storage::domain::error::DomainError;
use graph_storage::infra::fake_store::FakeGraphStore;
use graph_storage_sdk::models::{
NeighborhoodRequest, NodeSpec, SearchMode, SearchRequest, TraverseRequest, TruncationReason,
TypeQuery,
};
use support::Harness;
#[tokio::test]
async fn the_ontology_registers_and_reads_back_through_the_service() {
let harness = Harness::allowed();
let ctx = harness.ctx();
let registered = harness
.services
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("the ontology registers");
assert!(
registered.len() >= 2,
"every submitted type is reported: {}",
registered.len()
);
let record = harness
.services
.get_type(&ctx, &conformance::OWNED.to_owned())
.await
.expect("the producer type reads back");
assert_eq!(record.type_id, conformance::OWNED);
assert_eq!(
record.effective_traits.family.as_deref(),
Some("owned"),
"traits are merged down the chain, not read off the leaf"
);
let page = harness
.services
.list_types(&ctx, TypeQuery::default())
.await
.expect("types list");
assert!(
page.items
.iter()
.any(|item| item.type_id == conformance::OWNED),
"the list carries what was registered"
);
}
#[tokio::test]
async fn a_schema_outside_its_declared_chain_is_refused_by_the_service() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let orphan = graph_storage_sdk::models::TypeRegistration {
type_id: "gts.cf.core.graph.node.v1~cf.core.graph.owned_node.v1~acme.gs._.orphan.v1~"
.to_owned(),
schema: serde_json::json!({
"$id": "gts://gts.cf.core.graph.node.v1~cf.core.graph.owned_node.v1~acme.gs._.orphan.v1~",
"$schema": "http://json-schema.org/draft-07/schema#",
"type": "object",
"allOf": [{ "$ref": "gts://gts.nobody.registered.this.v1~" }]
}),
};
let error = harness
.services
.register_types(&ctx, vec![orphan])
.await
.expect_err("a schema that cannot compile is refused at registration");
let rendered = error.to_string();
assert!(
matches!(
error,
DomainError::Validation { .. } | DomainError::InvalidArgument { .. }
),
"expected a validation failure, got {rendered}"
);
}
#[tokio::test]
async fn an_ingest_goes_through_authorization_admission_and_the_coordinator() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let outcome = harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("a", "first"),
conformance::node("b", "second"),
],
vec![conformance::edge("a", "b")],
),
)
.await
.expect("the batch commits");
assert_eq!(outcome.counts.nodes_inserted, 2);
assert_eq!(outcome.counts.edges_inserted, 1);
assert!(outcome.revision.revision > 0, "the revision advanced");
let view = harness
.services
.get_node(&ctx, &"a".to_owned(), None)
.await
.expect("the node reads back");
assert!(
view.has_embedding,
"the service ran the batch through the embedding coordinator, not around it"
);
}
#[tokio::test]
async fn a_batch_over_the_node_bound_is_refused_before_any_validation() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let too_many = (0..=harness.services.config().ingest_max_nodes)
.map(|i| conformance::node(&format!("n{i}"), "x"))
.collect();
let error = harness
.services
.ingest(&ctx, conformance::batch(too_many, Vec::new()))
.await
.expect_err("the bound is enforced");
match error {
DomainError::LimitExceeded { what } => assert!(
what.contains("ingest_max_nodes"),
"the refusal names the bound it enforced: {what}"
),
other => panic!("expected a limit refusal, got {other}"),
}
}
#[tokio::test]
async fn an_unregistered_type_fails_the_item_not_the_request() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let mut node = conformance::node("unknown-type", "x");
node.type_id =
"gts.cf.core.graph.node.v1~cf.core.graph.owned_node.v1~acme.gs._.nope.v1~".to_owned();
let error = harness
.services
.ingest(&ctx, conformance::batch(vec![node], Vec::new()))
.await
.expect_err("an unregistered type is refused");
match error {
DomainError::Validation { items } => {
assert_eq!(items.len(), 1, "one item, one error: {items:?}");
}
other => panic!("expected per-item validation errors, got {other}"),
}
}
#[tokio::test]
async fn every_read_path_answers_through_the_service() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("read-a", "findable alpha"),
conformance::node("read-b", "findable beta"),
],
vec![conformance::edge("read-a", "read-b")],
),
)
.await
.expect("the batch commits");
let node = harness
.services
.get_node(&ctx, &"read-a".to_owned(), Some(5))
.await
.expect("the node reads");
assert_eq!(node.adjacency.len(), 1);
let edge = harness
.services
.get_edge(&ctx, &node.adjacency[0].edge_key)
.await
.expect("the edge reads");
assert_eq!((edge.src.as_str(), edge.dst.as_str()), ("read-a", "read-b"));
let page = harness
.services
.project_nodes(&ctx, &[], toolkit_odata::ODataQuery::default())
.await
.expect("the projection answers");
assert_eq!(page.items.len(), 2, "both rows: {:?}", page.items);
let hits = harness
.services
.search(
&ctx,
SearchRequest {
mode: SearchMode::Lexical,
query: Some("findable".to_owned()),
arm_limit: 10,
limit: 10,
type_patterns: Vec::new(),
},
)
.await
.expect("search answers");
assert_eq!(hits.hits.len(), 2, "both documents rank: {:?}", hits.hits);
let revision = harness
.services
.revision(&ctx)
.await
.expect("revision reads");
assert!(revision.revision > 0);
}
#[tokio::test]
async fn a_read_bound_is_refused_rather_than_clamped() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let over = harness.services.config().node_read_max_adjacency + 1;
let error = harness
.services
.get_node(&ctx, &"anything".to_owned(), Some(over))
.await
.expect_err("an adjacency limit above the ceiling is refused");
assert!(
matches!(error, DomainError::LimitExceeded { .. }),
"expected a limit refusal, got {error}"
);
let error = harness
.services
.search(
&ctx,
SearchRequest {
mode: SearchMode::Hybrid,
query: None,
arm_limit: 10,
limit: 10,
type_patterns: Vec::new(),
},
)
.await
.expect_err("a search without text is refused");
assert!(
matches!(error, DomainError::LimitCombination { .. }),
"expected a limit-combination refusal, got {error}"
);
}
#[tokio::test]
async fn a_traversal_walks_hops_and_a_neighborhood_answers_from_a_root() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("hop-a", "a"),
conformance::node("hop-b", "b"),
conformance::node("hop-c", "c"),
],
vec![
conformance::edge("hop-a", "hop-b"),
conformance::edge("hop-b", "hop-c"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["hop-a".to_owned()],
depth: 2,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect("the traversal answers");
let mut keys: Vec<String> = walked.nodes.iter().map(|n| n.node_key.clone()).collect();
keys.sort();
assert_eq!(
keys,
vec!["hop-a".to_owned(), "hop-b".to_owned(), "hop-c".to_owned()],
"two hops reach the whole chain"
);
assert_eq!(
walked.seeds,
vec!["hop-a".to_owned()],
"the answer says which seeds it actually started from"
);
let partly = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec![
"hop-a".to_owned(),
"hop-a".to_owned(),
"never-existed".to_owned(),
],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect("the traversal answers");
assert_eq!(
partly.seeds,
vec!["hop-a".to_owned()],
"duplicates collapse and an unknown seed is absent, not an error"
);
let around = harness
.services
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "hop-b".to_owned(),
depth: 1,
node_budget: Some(10),
include_phantoms: false,
},
)
.await
.expect("the neighborhood answers");
assert_eq!(
around.nodes.len(),
3,
"one hop around the middle reaches both sides: {:?}",
around.nodes.iter().map(|n| &n.node_key).collect::<Vec<_>>()
);
}
#[tokio::test]
async fn a_traversal_outside_its_bounds_is_refused() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let over_depth = harness.services.config().traversal_max_depth + 1;
for (request, what) in [
(
TraverseRequest {
seeds: Vec::new(),
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: None,
},
"a traversal with no seed",
),
(
TraverseRequest {
seeds: vec!["x".to_owned()],
depth: over_depth,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: None,
},
"a depth above the ceiling",
),
] {
assert!(
harness.services.traverse(&ctx, request).await.is_err(),
"{what} must be refused"
);
}
}
#[tokio::test]
async fn a_delete_tombstones_the_node_and_its_edge_together() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("del-a", "a"),
conformance::node("del-b", "b"),
],
vec![conformance::edge("del-a", "del-b")],
),
)
.await
.expect("the batch commits");
let outcome = harness
.services
.delete_node(&ctx, &"del-a".to_owned())
.await
.expect("the delete succeeds");
assert_eq!(outcome.tombstoned_nodes, 1);
assert_eq!(
outcome.tombstoned_edges, 1,
"an incident edge follows the node in the same transaction"
);
assert!(
harness
.services
.get_node(&ctx, &"del-a".to_owned(), None)
.await
.is_err(),
"a tombstoned node is absent from the read"
);
}
#[tokio::test]
async fn readiness_answers_without_a_caller() {
let harness = Harness::allowed();
let readiness = harness.services.readiness().await;
assert!(
!readiness.components.is_empty(),
"readiness names its components"
);
}
#[tokio::test]
async fn an_embedding_space_mismatch_is_unhealthy_and_leaves_the_gear_ready() {
use graph_storage::domain::embedding::{EmbeddingCoordinator, SpaceState};
use graph_storage_sdk::models::{EMBEDDING_SPACE, ReadinessState};
let blocked = EmbeddingCoordinator::new(conformance::provider(), SpaceState::Blocked, 8 * 1024);
let harness = Harness::with_coordinator(Arc::new(support::AllowInOwnTenant), blocked);
let readiness = harness.services.readiness().await;
let space = readiness
.components
.iter()
.find(|row| row.component == EMBEDDING_SPACE)
.expect("the embedding space has a row");
assert_eq!(space.state, ReadinessState::Unhealthy, "{space:?}");
assert!(
space
.blocked
.as_deref()
.is_some_and(|what| what.contains("EMBEDDING_SPACE_MISMATCH")),
"the row names what it refuses: {space:?}"
);
assert!(
readiness.ready,
"a space mismatch blocks vector search, not the gear: {readiness:?}"
);
let healthy = Harness::allowed().services.readiness().await;
let row = healthy
.components
.iter()
.find(|row| row.component == EMBEDDING_SPACE)
.expect("the embedding space has a row");
assert_eq!(row.state, ReadinessState::Healthy, "{row:?}");
}
#[tokio::test]
async fn a_readiness_row_does_not_repeat_what_a_failing_dependency_said() {
use async_trait::async_trait;
use graph_storage::domain::embedding::{EmbeddingCoordinator, SpaceState};
use graph_storage::infra::embedding::fake::FakeEmbeddingProvider;
use graph_storage_sdk::models::{EMBEDDING_PROVIDER, EmbeddingSpaceId, ReadinessState};
use graph_storage_sdk::plugin_api::{
EmbedRequest, EmbedResponse, EmbeddingProviderError, EmbeddingProviderV1,
};
const INTERNAL: &str = "https://embeddings.internal.example:8443/v1/embeddings";
struct Unreachable(FakeEmbeddingProvider);
#[async_trait]
impl EmbeddingProviderV1 for Unreachable {
fn embedding_space(&self) -> &EmbeddingSpaceId {
self.0.embedding_space()
}
fn dimension(&self) -> u32 {
self.0.dimension()
}
async fn embed(&self, req: EmbedRequest) -> Result<EmbedResponse, EmbeddingProviderError> {
self.0.embed(req).await
}
async fn health(&self) -> Result<(), EmbeddingProviderError> {
Err(EmbeddingProviderError::Unavailable {
reason: format!("{INTERNAL}: connection refused"),
})
}
}
let provider = Arc::new(Unreachable(FakeEmbeddingProvider::new(
conformance::DIMENSION,
)));
let coordinator = EmbeddingCoordinator::new(
provider,
SpaceState::Active {
epoch: conformance::EPOCH,
},
8 * 1024,
);
let harness = Harness::with_coordinator(Arc::new(support::AllowInOwnTenant), coordinator);
let readiness = harness.services.readiness().await;
let row = readiness
.components
.iter()
.find(|row| row.component == EMBEDDING_PROVIDER)
.expect("the provider has a row");
assert_eq!(row.state, ReadinessState::Degraded, "{row:?}");
let said = format!("{readiness:?}");
assert!(
!said.contains("embeddings.internal.example") && !said.contains("connection refused"),
"readiness repeats what the provider said: {row:?}"
);
}
#[tokio::test]
async fn a_type_preview_needs_read_and_the_change_it_previews_still_needs_admin() {
use graph_storage::infra::fake_store::FakeGraphStore;
let store = Arc::new(FakeGraphStore::new());
let administrator = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
GraphStorageConfig::default(),
);
let admin_ctx = administrator.ctx();
administrator.seed_ontology(&admin_ctx).await;
let mut producer = Harness::configured_over(
store,
Arc::new(support::ReadOnly),
GraphStorageConfig::default(),
);
producer.tenant = administrator.tenant;
let producer_ctx = producer.ctx();
let candidate = conformance::ontology_batch();
let verdicts = producer
.services
.type_compatibility(&producer_ctx, candidate.clone(), true, Vec::new())
.await
.expect("a preview is a read, and read is what this caller holds");
assert!(!verdicts.is_empty(), "the preview answers for each type");
let refused = producer
.services
.register_types(&producer_ctx, candidate)
.await
.expect_err("making the change is still administration");
assert!(
matches!(refused, DomainError::AccessDenied),
"the refusal is a permission refusal, got {refused}"
);
}
#[tokio::test]
async fn the_namespace_surface_lists_and_transfers() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
assert!(
harness
.services
.list_source_namespaces(&ctx)
.await
.expect("the list answers")
.is_empty(),
"nothing is claimed before anything is written"
);
let assigned = harness
.services
.transfer_source_namespace(&ctx, "unclaimed", "mirror-gear")
.await
.expect("an unclaimed namespace can be pre-assigned");
assert_eq!(assigned.owner_principal, "mirror-gear");
assert_eq!(assigned.previous_owner, None, "there was no previous owner");
assert!(
assigned.transferred_by.is_some(),
"the administrative act records who performed it"
);
let moved = harness
.services
.transfer_source_namespace(&ctx, "unclaimed", "other-gear")
.await
.expect("and then handed on");
assert_eq!(moved.previous_owner.as_deref(), Some("mirror-gear"));
let listed = harness
.services
.list_source_namespaces(&ctx)
.await
.expect("the list answers");
assert_eq!(listed.len(), 1, "the boundary is visible: {listed:?}");
assert_eq!(listed[0].owner_principal, "other-gear");
for stranded in ["other-gear ", " other-gear", "other-gear\n"] {
let error = harness
.services
.transfer_source_namespace(&ctx, "unclaimed", stranded)
.await
.expect_err("a principal no writer could equal is not a transfer");
assert!(
matches!(error, DomainError::InvalidQuery { .. }),
"{stranded:?} must be refused as invalid, got {error:?}"
);
}
let after = harness
.services
.list_source_namespaces(&ctx)
.await
.expect("the list answers");
assert_eq!(
after[0].owner_principal, "other-gear",
"a refused transfer leaves the owner where it was"
);
}
#[tokio::test]
async fn an_unauthorized_transfer_answers_the_same_whatever_it_was_given() {
let harness = Harness::denied();
let ctx = harness.ctx();
let mut refusals = Vec::new();
for (namespace, principal) in [
("github", "mirror-gear"), ("git\nhub", "mirror-gear"), (" github ", "mirror-gear"), ("github", "mirror-gear "), ] {
let refused = harness
.services
.transfer_source_namespace(&ctx, namespace, principal)
.await
.expect_err("a denied caller never transfers anything");
refusals.push(std::mem::discriminant(&refused));
}
assert!(
refusals.windows(2).all(|pair| pair[0] == pair[1]),
"every refusal must be the same variant, whatever the input looked like"
);
}
#[tokio::test]
async fn a_denying_pdp_stops_every_surface_before_the_store() {
let harness = Harness::denied();
let ctx = harness.ctx();
let refusals = [
harness
.services
.register_types(&ctx, conformance::ontology_batch())
.await
.err(),
harness
.services
.ingest(&ctx, conformance::batch(Vec::new(), Vec::new()))
.await
.err(),
harness
.services
.get_node(&ctx, &"x".to_owned(), None)
.await
.err(),
harness.services.get_edge(&ctx, &"x".to_owned()).await.err(),
harness
.services
.project_nodes(&ctx, &[], toolkit_odata::ODataQuery::default())
.await
.err(),
harness.services.revision(&ctx).await.err(),
harness
.services
.delete_node(&ctx, &"x".to_owned())
.await
.err(),
];
for refusal in refusals {
assert!(
matches!(refusal, Some(DomainError::AccessDenied)),
"a denied caller must be refused by the PEP, got {refusal:?}"
);
}
}
#[tokio::test]
async fn the_local_client_answers_like_the_service_and_is_bounded_like_it() {
use graph_storage::domain::local_client::GraphStorageLocalClient;
use graph_storage_sdk::GraphStorageClientV1;
let harness = Harness::allowed();
let ctx = harness.ctx();
let client = GraphStorageLocalClient::new(std::sync::Arc::clone(&harness.services));
client
.register_types(&ctx, conformance::ontology_batch())
.await
.expect("the ontology registers in-process");
let outcome = client
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("local-a", "findable one"),
conformance::node("local-b", "findable two"),
],
vec![conformance::edge("local-a", "local-b")],
),
)
.await
.expect("the batch commits in-process");
assert_eq!(outcome.counts.nodes_inserted, 2);
let node = client
.get_node(&ctx, &"local-a".to_owned(), None)
.await
.expect("the node reads");
assert_eq!(node.adjacency.len(), 1);
let edge_key = node.adjacency[0].edge_key.clone();
assert_eq!(
client
.get_type(&ctx, &conformance::OWNED.to_owned())
.await
.expect("the type reads")
.type_id,
conformance::OWNED
);
assert!(
!client
.list_types(&ctx, TypeQuery::default())
.await
.expect("types list")
.items
.is_empty()
);
assert_eq!(
client
.project_nodes(&ctx, &[], toolkit_odata::ODataQuery::default())
.await
.expect("the projection answers")
.items
.len(),
2
);
assert_eq!(
client
.search(
&ctx,
SearchRequest {
mode: SearchMode::Lexical,
query: Some("findable".to_owned()),
arm_limit: 10,
limit: 10,
type_patterns: Vec::new(),
},
)
.await
.expect("search answers")
.hits
.len(),
2
);
assert_eq!(
client
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["local-a".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect("the traversal answers")
.nodes
.len(),
2
);
assert!(
!client
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "local-a".to_owned(),
depth: 1,
node_budget: Some(10),
include_phantoms: false,
},
)
.await
.expect("the neighborhood answers")
.nodes
.is_empty()
);
assert!(
client
.revision(&ctx)
.await
.expect("revision reads")
.revision
> 0,
"the in-process path reports the same revision surface"
);
let over = harness.services.config().node_read_max_adjacency + 1;
let refused = client
.get_node(&ctx, &"local-a".to_owned(), Some(over))
.await
.expect_err("the in-process path is bounded too");
assert_eq!(
refused.status_code(),
400,
"the same refusal, rendered through the same mapping: {refused:?}"
);
let deleted = client
.delete_edge(&ctx, &edge_key)
.await
.expect("the edge is deleted in-process");
assert_eq!(deleted.tombstoned_edges, 1);
assert_eq!(
client
.delete_node(&ctx, &"local-a".to_owned())
.await
.expect("the node is deleted in-process")
.tombstoned_nodes,
1
);
}
#[tokio::test]
async fn no_new_work_starts_after_the_budget_is_spent() {
let expired = GraphStorageConfig {
deadline_interactive_secs: 0,
..GraphStorageConfig::default()
};
assert!(
expired.validate().is_err(),
"a zero deadline is not a configuration anyone can set"
);
let pdp = Arc::new(support::CountingPdp::default());
let harness = Harness::configured(
Arc::clone(&pdp) as Arc<dyn authz_resolver_sdk::api::AuthZResolverApi>,
expired,
);
let ctx = harness.ctx();
let refused = harness
.services
.ingest(
&ctx,
conformance::batch(vec![conformance::node("late-a", "a")], Vec::new()),
)
.await
.expect_err("an ingest does not start under a spent budget");
assert!(
matches!(refused, DomainError::Deadline),
"expected a deadline refusal, got {refused}"
);
let refused = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["late-a".to_owned()],
depth: 2,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect_err("a compound read does not start under a spent budget");
assert!(
matches!(refused, DomainError::Deadline),
"expected a deadline refusal, got {refused}"
);
let search = harness.services.search(
&ctx,
SearchRequest {
mode: SearchMode::Lexical,
query: Some("anything".to_owned()),
arm_limit: 10,
limit: 10,
type_patterns: Vec::new(),
},
);
let reads: Vec<(&str, DomainError)> = vec![
(
"get_node",
harness
.services
.get_node(&ctx, &"late-a".to_owned(), Some(10))
.await
.expect_err("get_node does not start"),
),
("search", search.await.expect_err("search does not start")),
(
"list_types",
harness
.services
.list_types(&ctx, TypeQuery::default())
.await
.expect_err("the catalogue does not start"),
),
(
"project_nodes",
harness
.services
.project_nodes(&ctx, &[], toolkit_odata::ODataQuery::default())
.await
.expect_err("the projection does not start"),
),
(
"revision",
harness
.services
.revision(&ctx)
.await
.expect_err("the revision read does not start"),
),
(
"neighborhood",
harness
.services
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "late-a".to_owned(),
depth: 1,
node_budget: Some(10),
include_phantoms: false,
},
)
.await
.expect_err("the neighborhood does not start"),
),
];
for (what, error) in reads {
assert!(
matches!(error, DomainError::Deadline),
"{what} must refuse a spent deadline, got {error}"
);
}
assert_eq!(
pdp.calls(),
0,
"the policy decision point must not be reached by a request that is already over \
its deadline"
);
}
#[tokio::test]
async fn a_store_without_snapshots_is_not_asked_for_one() {
let store = Arc::new(FakeGraphStore::declining_snapshots(0));
let harness = Harness::over(store, Arc::new(support::AllowInOwnTenant));
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![conformance::node("s-a", "a"), conformance::node("s-b", "b")],
vec![conformance::edge("s-a", "s-b")],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["s-a".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect("the traversal answers without a snapshot");
assert_eq!(walked.nodes.len(), 2, "the walk still works: {walked:?}");
assert!(
walked.consistent_snapshot,
"nothing was written while it ran, so the arms agree even without a snapshot"
);
}
#[tokio::test]
async fn a_walk_that_spans_a_commit_is_reported_as_inconsistent() {
let store = Arc::new(FakeGraphStore::declining_snapshots(1));
let harness = Harness::over(store, Arc::new(support::AllowInOwnTenant));
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(vec![conformance::node("d-a", "a")], Vec::new()),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["d-a".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect("the traversal answers");
assert!(
!walked.consistent_snapshot,
"the revision moved under the walk, and the answer says so"
);
}
#[tokio::test]
async fn a_traversal_stops_at_the_byte_budget_and_reports_it() {
let small = GraphStorageConfig {
response_max_bytes: 4 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "z".repeat(1_500);
let fat = |key: &str| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
};
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![fat("b-a"), fat("b-b"), fat("b-c"), fat("b-d")],
vec![
conformance::edge("b-a", "b-b"),
conformance::edge("b-b", "b-c"),
conformance::edge("b-c", "b-d"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["b-a".to_owned()],
depth: 3,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(100),
},
)
.await
.expect("the traversal answers");
assert!(
walked.nodes.len() < 4,
"the budget cut the answer short: {} nodes",
walked.nodes.len()
);
assert_eq!(
walked.truncated,
Some(TruncationReason::ResponseBytes),
"and the cut is reported as a byte budget rather than a node budget"
);
let returned: std::collections::BTreeSet<&str> =
walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
for edge in &walked.edges {
assert!(
returned.contains(edge.src.as_str()) && returned.contains(edge.dst.as_str()),
"the edge filter follows the byte cut: {edge:?} against {returned:?}"
);
}
}
#[tokio::test]
async fn an_oversized_projection_page_is_refused_rather_than_trimmed() {
let small = GraphStorageConfig {
response_max_bytes: 4 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "p".repeat(1_500);
let fat = |key: &str| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
};
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![fat("p-a"), fat("p-b"), fat("p-c"), fat("p-d")],
Vec::new(),
),
)
.await
.expect("the batch commits");
let refused = harness
.services
.project_nodes(&ctx, &[], toolkit_odata::ODataQuery::default())
.await
.expect_err("the page hydrates past the budget");
assert!(
matches!(refused, DomainError::LimitExceeded { .. }),
"expected a bound refusal, got {refused}"
);
assert!(
refused.to_string().contains("response_max_bytes") && refused.to_string().contains("$top"),
"the refusal names the bound and what to do about it: {refused}"
);
}
#[tokio::test]
async fn an_oversized_hit_list_is_cut_and_reported() {
let small = GraphStorageConfig {
response_max_bytes: 1_024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let long = "s".repeat(400);
let nodes: Vec<NodeSpec> = (0..6)
.map(|i| conformance::node(&format!("hit-{i}"), &format!("{long}-{i}")))
.collect();
harness
.services
.ingest(&ctx, conformance::batch(nodes, Vec::new()))
.await
.expect("the batch commits");
let found = harness
.services
.search(
&ctx,
SearchRequest {
mode: SearchMode::Lexical,
query: Some(long.clone()),
arm_limit: 50,
limit: 50,
type_patterns: Vec::new(),
},
)
.await
.expect("the search answers");
assert!(
found.hits.len() < 6,
"the budget cut the list: {} hits",
found.hits.len()
);
assert_eq!(
found.truncated,
Some(TruncationReason::ResponseBytes),
"and a short list says why, since a small graph looks the same"
);
}
#[tokio::test]
async fn the_byte_budget_charges_the_edges_too() {
const LEAVES: usize = 8;
let small = GraphStorageConfig {
response_max_bytes: 1_536,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let mut nodes = vec![conformance::node("hub", "hub")];
let mut edges = Vec::new();
for index in 0..LEAVES {
let leaf = format!("leaf-{index}");
edges.push(conformance::edge("hub", &leaf));
nodes.push(conformance::node(&leaf, &leaf));
}
harness
.services
.ingest(&ctx, conformance::batch(nodes, edges))
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["hub".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(100),
},
)
.await
.expect("the traversal answers");
let returned: std::collections::BTreeSet<&str> =
walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
for seed in &walked.seeds {
assert!(
returned.contains(seed.as_str()),
"a seed the answer names is a seed the answer contains: {seed} not in {returned:?}"
);
}
assert_eq!(
walked.nodes.len(),
LEAVES + 1,
"every node fits, so the nodes are not what the budget cut: {returned:?}"
);
assert!(
walked.edges.len() < LEAVES,
"the edges are what it cut, which is only possible if they were \
charged: {} of {LEAVES}",
walked.edges.len()
);
assert_eq!(
walked.truncated,
Some(TruncationReason::ResponseBytes),
"and the cut is reported"
);
for edge in &walked.edges {
assert!(
returned.contains(edge.src.as_str()) && returned.contains(edge.dst.as_str()),
"an edge names two returned nodes: {edge:?}"
);
}
}
#[tokio::test]
async fn seeds_that_cannot_fit_the_budget_are_refused() {
let tiny = GraphStorageConfig {
response_max_bytes: 2 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), tiny);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "w".repeat(1_500);
let fat = |key: &str| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
};
harness
.services
.ingest(
&ctx,
conformance::batch(vec![fat("s-1"), fat("s-2")], Vec::new()),
)
.await
.expect("the batch commits");
let refused = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["s-1".to_owned(), "s-2".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(10),
},
)
.await
.expect_err("two seeds larger than the budget cannot be answered");
assert!(
matches!(refused, DomainError::LimitExceeded { .. }),
"expected a bound refusal, got {refused}"
);
assert!(
refused.to_string().contains("the seeds alone"),
"the refusal says which part did not fit: {refused}"
);
}
#[tokio::test]
async fn a_seed_is_not_evicted_by_its_own_bytes_counted_twice() {
const SEEDS: usize = 6;
let measured = GraphStorageConfig {
response_max_bytes: 12 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), measured);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "z".repeat(1_500);
let fat = |key: &str| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
};
let keys: Vec<String> = (0..SEEDS).map(|index| format!("seed-{index}")).collect();
harness
.services
.ingest(
&ctx,
conformance::batch(keys.iter().map(|k| fat(k)).collect(), Vec::new()),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: keys.clone(),
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(100),
},
)
.await
.expect("the seeds fit the budget, so the traversal answers");
let returned: std::collections::BTreeSet<&str> =
walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
for seed in &keys {
assert!(
returned.contains(seed.as_str()),
"every seed that fit must be in the answer: {seed} missing from {returned:?}"
);
}
assert_eq!(
walked.nodes.len(),
SEEDS,
"the seeds are the whole answer here: {returned:?}"
);
}
#[tokio::test]
async fn a_traversal_reads_what_its_budget_can_hold_and_not_the_whole_walk() {
const LEAVES: usize = 100;
let store = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let budgeted = GraphStorageConfig {
response_max_bytes: 12 * 1024,
item_max_bytes: 2 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
budgeted,
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "q".repeat(1_500);
let mut nodes = vec![conformance::node("hub", "hub")];
let mut edges = Vec::new();
for index in 0..LEAVES {
let leaf = format!("fat-{index}");
edges.push(conformance::edge("hub", &leaf));
nodes.push(NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(&leaf, &leaf)
});
}
harness
.services
.ingest(&ctx, conformance::batch(nodes, edges))
.await
.expect("the batch commits");
let before = store.rows_hydrated();
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["hub".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(1_000),
},
)
.await
.expect("the traversal answers");
let read = store.rows_hydrated() - before;
assert_eq!(
walked.truncated,
Some(TruncationReason::ResponseBytes),
"the budget is what cut the answer, or this proves nothing about it"
);
let returned = walked.nodes.len() as u64;
assert!(
returned > 1 && returned < (LEAVES as u64),
"the answer is a handful of leaves, not none and not all: {returned}"
);
let piece = budgeted_piece(12 * 1024, 2 * 1024);
assert!(
read <= returned + piece,
"read {read} rows to return {returned} out of a walk of {}; the reads \
should stop within one piece ({piece}) of the budget",
LEAVES + 1
);
}
#[tokio::test]
async fn a_seed_set_over_the_budget_is_refused_before_it_is_read_in_full() {
const SEEDS: usize = 40;
let store = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let budgeted = GraphStorageConfig {
response_max_bytes: 12 * 1024,
item_max_bytes: 2 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
budgeted,
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "q".repeat(1_500);
let keys: Vec<String> = (0..SEEDS)
.map(|index| format!("fat-seed-{index}"))
.collect();
let nodes = keys
.iter()
.map(|key| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
})
.collect();
harness
.services
.ingest(&ctx, conformance::batch(nodes, Vec::new()))
.await
.expect("the batch commits");
let before = store.rows_hydrated();
let refused = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: keys,
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: Vec::new(),
max_nodes: Some(1_000),
},
)
.await
.expect_err("forty seeds of 1.5 KB do not fit a 12 KB budget");
let read = store.rows_hydrated() - before - SEEDS as u64;
assert!(
matches!(refused, DomainError::LimitExceeded { .. }),
"the refusal is a bound: {refused}"
);
let holds: u64 = (12_u64 * 1024).div_euclid(1_500);
assert!(
read <= holds + 1,
"read {read} of {SEEDS} seed rows to refuse a budget that holds {holds}"
);
}
#[tokio::test]
async fn a_filtered_walk_near_its_budget_does_not_hydrate_row_by_row() {
const SEEDS: usize = 6;
const PHANTOMS: usize = 50;
let store = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let budgeted = GraphStorageConfig {
response_max_bytes: 12 * 1024,
item_max_bytes: 2 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
budgeted,
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "q".repeat(1_500);
let seeds: Vec<String> = (0..SEEDS).map(|index| format!("near-{index}")).collect();
let nodes = seeds
.iter()
.map(|key| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
})
.collect();
let edges = (0..PHANTOMS)
.map(|index| conformance::edge(&seeds[0], &format!("ghost-{index}")))
.collect();
harness
.services
.ingest(&ctx, conformance::batch(nodes, edges))
.await
.expect("the batch commits");
let before = store.hydrate_calls();
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: seeds.clone(),
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: vec![conformance::OWNED.to_owned()],
max_nodes: Some(1_000),
},
)
.await
.expect("the traversal answers");
let calls = store.hydrate_calls() - before;
let mut returned: Vec<String> = walked.nodes.iter().map(|n| n.node_key.clone()).collect();
returned.sort();
let mut expected = seeds;
expected.sort();
assert_eq!(
returned, expected,
"the answer is the seeds; every phantom is filtered"
);
assert!(
calls <= 3,
"{calls} hydrate calls for a walk of {SEEDS} seeds and {PHANTOMS} filtered neighbours"
);
}
#[tokio::test]
async fn a_type_filter_drops_the_edges_of_the_nodes_it_filters() {
let store = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let harness = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
GraphStorageConfig::default(),
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("seed", "seed"),
conformance::node("kept", "kept"),
],
vec![
conformance::edge("seed", "kept"),
conformance::edge("seed", "ghost"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["seed".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: vec![conformance::OWNED.to_owned()],
max_nodes: Some(1_000),
},
)
.await
.expect("the traversal answers");
let mut returned: Vec<&str> = walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
returned.sort_unstable();
assert_eq!(returned, ["kept", "seed"], "the phantom is filtered");
let edges: Vec<(&str, &str)> = walked
.edges
.iter()
.map(|edge| (edge.src.as_str(), edge.dst.as_str()))
.collect();
assert_eq!(
edges,
[("seed", "kept")],
"the edge to the filtered phantom goes with it; only the edge between two \
returned nodes stays"
);
}
#[tokio::test]
async fn a_filtered_walk_survives_a_store_without_node_types_and_not_one_that_fails_them() {
async fn walk(
store: graph_storage::infra::fake_store::FakeGraphStore,
) -> Result<Vec<String>, DomainError> {
let store = Arc::new(store);
let harness = Harness::configured_over(
Arc::clone(&store),
Arc::new(support::AllowInOwnTenant),
GraphStorageConfig::default(),
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("seed", "seed"),
conformance::node("kept", "kept"),
],
vec![
conformance::edge("seed", "kept"),
conformance::edge("seed", "ghost"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["seed".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: vec![conformance::OWNED.to_owned()],
max_nodes: Some(1_000),
},
)
.await?;
let mut keys: Vec<String> = walked.nodes.into_iter().map(|n| n.node_key).collect();
keys.sort_unstable();
Ok(keys)
}
let answered = walk(graph_storage::infra::fake_store::FakeGraphStore::without_node_types())
.await
.expect("a store that cannot say the types is hydrated and filtered instead");
assert_eq!(
answered,
["kept", "seed"],
"the filter still applies, after hydration"
);
let failed = walk(graph_storage::infra::fake_store::FakeGraphStore::failing_node_types())
.await
.expect_err("a store that fails to answer fails the walk");
assert!(
failed.to_string().contains("node_types is failing"),
"the caller is told the store's failure, not served a slower answer: {failed}"
);
}
#[tokio::test]
async fn a_store_that_inherits_the_default_node_types_is_hydrated_and_filtered() {
let fake = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let harness = Harness::configured_over_store(
Arc::new(support::without_node_types::StoreWithoutNodeTypes(
Arc::clone(&fake),
)),
fake,
Arc::new(support::AllowInOwnTenant),
GraphStorageConfig::default(),
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("seed", "seed"),
conformance::node("kept", "kept"),
],
vec![
conformance::edge("seed", "kept"),
conformance::edge("seed", "ghost"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.traverse(
&ctx,
TraverseRequest {
seeds: vec!["seed".to_owned()],
depth: 1,
edge_type_patterns: Vec::new(),
node_type_patterns: vec![conformance::OWNED.to_owned()],
max_nodes: Some(1_000),
},
)
.await
.expect("the default answer is `Unsupported`, and the walk hydrates and filters");
let mut keys: Vec<&str> = walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
keys.sort_unstable();
assert_eq!(keys, ["kept", "seed"], "the filter applies after hydration");
}
#[tokio::test]
async fn a_store_that_inherits_the_default_node_types_still_excludes_phantoms() {
use graph_storage_sdk::models::NeighborhoodRequest;
let fake = Arc::new(graph_storage::infra::fake_store::FakeGraphStore::new());
let harness = Harness::configured_over_store(
Arc::new(support::without_node_types::StoreWithoutNodeTypes(
Arc::clone(&fake),
)),
fake,
Arc::new(support::AllowInOwnTenant),
GraphStorageConfig::default(),
);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![
conformance::node("seed", "seed"),
conformance::node("kept", "kept"),
],
vec![
conformance::edge("seed", "kept"),
conformance::edge("seed", "ghost"),
],
),
)
.await
.expect("the batch commits");
let walked = harness
.services
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "seed".to_owned(),
depth: 1,
node_budget: None,
include_phantoms: false,
},
)
.await
.expect("the neighborhood answers through the inherited default");
let mut keys: Vec<&str> = walked.nodes.iter().map(|n| n.node_key.as_str()).collect();
keys.sort_unstable();
assert_eq!(
keys,
["kept", "seed"],
"the phantom is dropped by the toggle alone; nothing else filters here"
);
let edges: Vec<(&str, &str)> = walked
.edges
.iter()
.map(|edge| (edge.src.as_str(), edge.dst.as_str()))
.collect();
assert_eq!(
edges,
[("seed", "kept")],
"no edge names the excluded phantom"
);
}
fn budgeted_piece(budget: u64, item_ceiling: u64) -> u64 {
budget.div_euclid(item_ceiling).max(1)
}
#[tokio::test]
async fn a_batch_is_bounded_by_its_total_size_and_not_only_its_counts() {
let small = GraphStorageConfig {
ingest_max_bytes: 16 * 1024,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let filler = "x".repeat(2 * 1024);
let fat = |key: &str| NodeSpec {
payload: Some(serde_json::json!({ "note": filler })),
..conformance::node(key, key)
};
harness
.services
.ingest(&ctx, conformance::batch(vec![fat("one")], Vec::new()))
.await
.expect("one item of this size is ordinary");
let many: Vec<NodeSpec> = (0..32).map(|i| fat(&format!("many-{i}"))).collect();
let refused = harness
.services
.ingest(&ctx, conformance::batch(many, Vec::new()))
.await
.expect_err("the sum of them is not");
assert!(
matches!(refused, DomainError::LimitExceeded { .. }),
"expected a bound refusal, got {refused}"
);
assert!(
refused.to_string().contains("ingest_max_bytes"),
"the refusal names the bound it hit: {refused}"
);
}
#[tokio::test]
async fn an_element_is_bounded_as_a_whole() {
let small = GraphStorageConfig {
item_max_bytes: 4_096,
payload_max_bytes: 4_096,
..GraphStorageConfig::default()
};
let harness = Harness::configured(Arc::new(support::AllowInOwnTenant), small);
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let node = NodeSpec {
payload: Some(serde_json::json!({ "note": "y".repeat(4_000) })),
..conformance::node("whole", &"n".repeat(200))
};
let refused = harness
.services
.ingest(&ctx, conformance::batch(vec![node], Vec::new()))
.await
.expect_err("the element as a whole is over the ceiling");
assert!(
refused.to_string().contains("item_max_bytes"),
"the refusal names the bound it hit: {refused}"
);
}
#[tokio::test]
async fn the_caller_controlled_strings_are_bounded() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let bound = harness.services.config().identifier_max_bytes as usize;
let long = "k".repeat(bound + 1);
let refused = harness
.services
.ingest(
&ctx,
conformance::batch(vec![conformance::node(&long, "fine")], vec![]),
)
.await
.expect_err("an oversized node_key is refused");
assert!(
matches!(refused, DomainError::LimitExceeded { .. }),
"expected a bound refusal, got {refused}"
);
let refused = harness
.services
.ingest(
&ctx,
conformance::batch(vec![conformance::node("fine", &long)], vec![]),
)
.await
.expect_err("an oversized name is refused");
assert!(matches!(refused, DomainError::LimitExceeded { .. }));
let mut batch = conformance::batch(vec![conformance::node("fine", "fine")], vec![]);
batch.idempotency_key = Some(long.clone());
let refused = harness
.services
.ingest(&ctx, batch)
.await
.expect_err("an oversized idempotency_key is refused");
assert!(matches!(refused, DomainError::LimitExceeded { .. }));
let over = harness.services.config().projection_max_page + 1;
let refused = harness
.services
.list_types(
&ctx,
TypeQuery {
top: Some(over),
..TypeQuery::default()
},
)
.await
.expect_err("an oversized catalogue page is refused");
assert!(matches!(refused, DomainError::LimitExceeded { .. }));
let query = "q".repeat(harness.services.config().search_query_max_bytes as usize + 1);
let refused = harness
.services
.search(
&ctx,
SearchRequest {
mode: SearchMode::Lexical,
query: Some(query),
arm_limit: 10,
limit: 10,
type_patterns: Vec::new(),
},
)
.await
.expect_err("an oversized query is refused");
assert!(matches!(refused, DomainError::LimitExceeded { .. }));
harness
.services
.ingest(
&ctx,
conformance::batch(
vec![conformance::node("ordinary", "an ordinary name")],
vec![],
),
)
.await
.expect("an ordinary node commits");
}
#[tokio::test]
async fn a_budgeted_neighborhood_keeps_the_best_connected_neighbours() {
let harness = Harness::allowed();
let ctx = harness.ctx();
harness.seed_ontology(&ctx).await;
let mut nodes = vec![conformance::node("hub", "hub")];
for leaf in ["leaf-1", "leaf-2", "leaf-3"] {
nodes.push(conformance::node(leaf, leaf));
}
for core in ["core-1", "core-2", "core-3"] {
nodes.push(conformance::node(core, core));
}
for far in ["far-1", "far-2", "far-3", "far-4", "far-5", "far-6"] {
nodes.push(conformance::node(far, far));
}
let mut edges = Vec::new();
for neighbour in ["leaf-1", "leaf-2", "leaf-3", "core-1", "core-2", "core-3"] {
edges.push(conformance::edge("hub", neighbour));
}
for (core, far) in [
("core-1", "far-1"),
("core-1", "far-2"),
("core-2", "far-3"),
("core-2", "far-4"),
("core-3", "far-5"),
("core-3", "far-6"),
] {
edges.push(conformance::edge(core, far));
}
harness
.services
.ingest(&ctx, conformance::batch(nodes, edges))
.await
.expect("the hub commits");
let around = harness
.services
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "hub".to_owned(),
depth: 1,
node_budget: Some(4),
include_phantoms: false,
},
)
.await
.expect("the neighborhood answers");
let mut kept: Vec<String> = around.nodes.iter().map(|n| n.node_key.clone()).collect();
kept.sort();
assert_eq!(
kept,
vec![
"core-1".to_owned(),
"core-2".to_owned(),
"core-3".to_owned(),
"hub".to_owned(),
],
"the budget keeps the root and the three connected neighbours, not the leaves"
);
assert!(
around.truncated.is_some(),
"a truncated neighborhood says so"
);
let whole = harness
.services
.neighborhood(
&ctx,
NeighborhoodRequest {
root: "hub".to_owned(),
depth: 1,
node_budget: Some(50),
include_phantoms: false,
},
)
.await
.expect("the neighborhood answers");
assert_eq!(whole.nodes.len(), 7, "the root and all six neighbours");
assert!(whole.truncated.is_none());
}
#[tokio::test]
async fn every_caller_supplied_identifier_is_bounded_before_it_reaches_the_store() {
let harness = Harness::allowed();
let ctx = harness.ctx();
let limit = GraphStorageConfig::default().identifier_max_bytes as usize;
let long = "k".repeat(limit + 1);
let fits = "k".repeat(limit);
let traverse = |seed: &str, edge: &str, node: &str| TraverseRequest {
seeds: vec![seed.to_owned()],
depth: 1,
edge_type_patterns: if edge.is_empty() {
Vec::new()
} else {
vec![edge.to_owned()]
},
node_type_patterns: if node.is_empty() {
Vec::new()
} else {
vec![node.to_owned()]
},
max_nodes: Some(10),
};
let neighborhood = |root: &str| NeighborhoodRequest {
root: root.to_owned(),
depth: 1,
node_budget: Some(10),
include_phantoms: false,
};
let search = |pattern: &str| SearchRequest {
mode: SearchMode::Lexical,
query: Some("anything".to_owned()),
arm_limit: 10,
limit: 10,
type_patterns: vec![pattern.to_owned()],
};
let type_query = |pattern: &str| TypeQuery {
kind: None,
pattern: Some(pattern.to_owned()),
top: Some(10),
cursor: None,
};
let s = &harness.services;
for (what, over, at) in [
(
"traverse seed",
s.traverse(&ctx, traverse(&long, "", "")).await.err(),
s.traverse(&ctx, traverse(&fits, "", "")).await.err(),
),
(
"traverse edge pattern",
s.traverse(&ctx, traverse("k", &long, "")).await.err(),
s.traverse(&ctx, traverse("k", &fits, "")).await.err(),
),
(
"traverse node pattern",
s.traverse(&ctx, traverse("k", "", &long)).await.err(),
s.traverse(&ctx, traverse("k", "", &fits)).await.err(),
),
(
"neighborhood root",
s.neighborhood(&ctx, neighborhood(&long)).await.err(),
s.neighborhood(&ctx, neighborhood(&fits)).await.err(),
),
(
"search type pattern",
s.search(&ctx, search(&long)).await.err(),
s.search(&ctx, search(&fits)).await.err(),
),
(
"type catalogue pattern",
s.list_types(&ctx, type_query(&long)).await.err(),
s.list_types(&ctx, type_query(&fits)).await.err(),
),
(
"get_type",
s.get_type(&ctx, &long).await.err(),
s.get_type(&ctx, &fits).await.err(),
),
(
"get_node",
s.get_node(&ctx, &long, None).await.err(),
s.get_node(&ctx, &fits, None).await.err(),
),
(
"get_edge",
s.get_edge(&ctx, &long).await.err(),
s.get_edge(&ctx, &fits).await.err(),
),
(
"delete_node",
s.delete_node(&ctx, &long).await.err(),
s.delete_node(&ctx, &fits).await.err(),
),
(
"delete_edge",
s.delete_edge(&ctx, &long).await.err(),
s.delete_edge(&ctx, &fits).await.err(),
),
(
"namespace transfer",
s.transfer_source_namespace(&ctx, &long, "owner")
.await
.err(),
s.transfer_source_namespace(&ctx, &fits, "owner")
.await
.err(),
),
(
"namespace owner",
s.transfer_source_namespace(&ctx, "github", &long)
.await
.err(),
s.transfer_source_namespace(&ctx, "github", &fits)
.await
.err(),
),
] {
assert!(
matches!(over, Some(DomainError::LimitExceeded { .. })),
"{what}: an identifier over the bound must be refused as one, got {over:?}"
);
assert!(
!matches!(at, Some(DomainError::LimitExceeded { .. })),
"{what}: an identifier at the bound must pass admission, got {at:?}"
);
}
}