diff --git a/store/postgres/src/relational_queries.rs b/store/postgres/src/relational_queries.rs index d87b950adc1..60ada6e1622 100644 --- a/store/postgres/src/relational_queries.rs +++ b/store/postgres/src/relational_queries.rs @@ -2241,7 +2241,9 @@ impl<'a> QueryFragment for FindDerivedQuery<'a> { // For truly gigantic `excluded_keys` lists, this will be slow, and // we should rewrite this query to use a CTE or a temp table to hold // the excluded keys. - out.push_sql(" != any("); + // `id <> all(..)` excludes every listed key; `id != any(..)` + // would hold for any row once two or more keys are listed + out.push_sql(" <> all("); self.excluded_keys.push_bind_param(&mut out)?; out.push_sql(") and "); } diff --git a/store/postgres/src/writable.rs b/store/postgres/src/writable.rs index 633688aeb1a..3ff7f3ea3d6 100644 --- a/store/postgres/src/writable.rs +++ b/store/postgres/src/writable.rs @@ -1344,12 +1344,14 @@ impl Queue { ) -> impl Iterator)> + 'a { batch .effective_ops(&derived_query.entity_type, at) - .filter_map(|op| match op { + .map(|op| match op { EntityOp::Write { key, entity } if is_related(derived_query, entity) => { - Some((key.clone(), Some(entity.clone()))) + (key.clone(), Some(entity.clone())) } - EntityOp::Write { .. } => None, - EntityOp::Remove { key } => Some((key.clone(), None)), + // A newer queued version supersedes the database row even + // when it is no longer related: exclude its key + EntityOp::Write { key, .. } => (key.clone(), None), + EntityOp::Remove { key } => (key.clone(), None), }) } diff --git a/store/test-store/tests/postgres/writable.rs b/store/test-store/tests/postgres/writable.rs index 7c17e8f201f..5fa3ad34ebf 100644 --- a/store/test-store/tests/postgres/writable.rs +++ b/store/test-store/tests/postgres/writable.rs @@ -41,6 +41,14 @@ const SCHEMA_GQL: &str = " id: String!, value: String! } + type Parent @entity { + id: ID!, + children: [Child!]! @derivedFrom(field: \"parent\") + } + type Child @entity { + id: ID!, + parent: Parent! + } type PoolCreated @entity(immutable: true) { id: Bytes!, token0: Bytes!, @@ -69,6 +77,8 @@ lazy_static! { .expect("Failed to parse user schema"); static ref COUNTER_TYPE: EntityType = TEST_SUBGRAPH_SCHEMA.entity_type(COUNTER).unwrap(); static ref COUNTER2_TYPE: EntityType = TEST_SUBGRAPH_SCHEMA.entity_type(COUNTER2).unwrap(); + static ref PARENT_TYPE: EntityType = TEST_SUBGRAPH_SCHEMA.entity_type("Parent").unwrap(); + static ref CHILD_TYPE: EntityType = TEST_SUBGRAPH_SCHEMA.entity_type("Child").unwrap(); } /// Inserts test data into the store. @@ -316,6 +326,108 @@ fn get_derived_nobatch() { get_with_pending(false, count_get_derived); } +fn set_parent(id: &str, vid: i64) -> EntityOperation { + EntityOperation::Set { + key: PARENT_TYPE.parse_key(id).unwrap(), + data: entity! { TEST_SUBGRAPH_SCHEMA => id: id, vid: vid }, + } +} + +fn set_child(id: &str, parent: &str, vid: i64) -> EntityOperation { + EntityOperation::Set { + key: CHILD_TYPE.parse_key(id).unwrap(), + data: entity! { TEST_SUBGRAPH_SCHEMA => id: id, parent: parent, vid: vid }, + } +} + +fn remove_child(id: &str) -> EntityOperation { + EntityOperation::Remove { + key: CHILD_TYPE.parse_key(id).unwrap(), + } +} + +/// The ids of the children of `parent` according to `get_derived` +async fn children_of(writable: &Arc, parent: &str) -> Vec { + let query = DerivedEntityQuery { + entity_type: CHILD_TYPE.clone(), + entity_field: Word::from("parent"), + value: PARENT_TYPE.parse_id(parent).unwrap(), + causality_region: CausalityRegion::ONCHAIN, + }; + let map = writable.get_derived(&query).await.unwrap(); + let mut ids: Vec = map.keys().map(|k| k.entity_id.to_string()).collect(); + ids.sort(); + ids +} + +/// `get_derived` while the changes of a block are still in the write +/// queue must give the same answer as after they are written. Block 1 +/// (parents `p0`, `p1`; children `w`, `x` under `p0`) is written to the +/// database; block 2 is empty and absorbs a writer that may already be +/// past its pause point; block 3 (`queued`) is held in the queue. +fn derived_with_pending(queued: Vec, expected: &'static [&'static str]) { + run_test(move |store, writable, _, deployment| async move { + let subgraph_store = store.subgraph_store(); + let base = vec![ + set_parent("p0", 1), + set_parent("p1", 2), + set_child("w", "p0", 3), + set_child("x", "p0", 4), + ]; + transact_entity_operations(&subgraph_store, &deployment, block_pointer(1), base) + .await + .unwrap(); + pause_writer(&deployment).await; + transact_entity_operations(&subgraph_store, &deployment, block_pointer(2), vec![]) + .await + .unwrap(); + transact_entity_operations(&subgraph_store, &deployment, block_pointer(3), queued) + .await + .unwrap(); + let head = subgraph_store + .least_block_ptr(&deployment.hash) + .await + .unwrap() + .map(|ptr| ptr.number); + assert!( + matches!(head, Some(1) | Some(2)), + "block 3 must still be queued, database head is {head:?}" + ); + let expected: Vec = expected.iter().map(|s| s.to_string()).collect(); + assert_eq!(expected, children_of(&writable, "p0").await, "while queued"); + + writable.flush().await.unwrap(); + assert_eq!(expected, children_of(&writable, "p0").await, "after flush"); + }) +} + +/// Control: one queued removal is excluded (worked before the fix) +#[test] +fn get_derived_pending_one_removal() { + derived_with_pending(vec![remove_child("x")], &["w"]); +} + +/// Two queued removals: with `!= any(..)` nothing was excluded +#[test] +fn get_derived_pending_two_removals() { + derived_with_pending(vec![remove_child("w"), remove_child("x")], &[]); +} + +/// A queued removal plus a queued creation: two excluded keys +#[test] +fn get_derived_pending_removal_and_creation() { + derived_with_pending( + vec![remove_child("x"), set_child("y", "p0", 5)], + &["w", "y"], + ); +} + +/// A queued write that moves `x` to another parent must exclude `x` +#[test] +fn get_derived_pending_parent_change() { + derived_with_pending(vec![set_child("x", "p1", 5)], &["w"]); +} + #[test] fn restart() { run_test(|store, writable, _, deployment| async move {