Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion store/postgres/src/relational_queries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2241,7 +2241,9 @@ impl<'a> QueryFragment<Pg> 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 ");
}
Expand Down
10 changes: 6 additions & 4 deletions store/postgres/src/writable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1344,12 +1344,14 @@ impl Queue {
) -> impl Iterator<Item = (EntityKey, Option<Entity>)> + '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),
})
}

Expand Down
112 changes: 112 additions & 0 deletions store/test-store/tests/postgres/writable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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<dyn WritableStore>, parent: &str) -> Vec<String> {
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<String> = 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<EntityOperation>, 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<String> = 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 {
Expand Down