Sitelet https://github.com/graphprotocol/graph-node/commit/6ee05e2a6645fd3f18d141d6e58407bc892088b7
Skip to content

Commit 6ee05e2

Browse files
committed
feature(offchain): Add causality_region column to entity tables
For now this just tracks the tables that need the column and adds the column to the DDL, but still unconditionally inserts 0. Inserting the correct causality region is follow up work.
1 parent 55b045b commit 6ee05e2

24 files changed

Lines changed: 236 additions & 59 deletions

File tree

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎core/src/subgraph/registrar.rs‎

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -620,11 +620,30 @@ async fn create_subgraph_version<C: Blockchain, S: SubgraphStore>(
620620
"block" => format!("{:?}", base_block.as_ref().map(|(_,ptr)| ptr.number))
621621
);
622622

623+
// Entity types that may be touched by offchain data sources need a causality region column.
624+
let needs_causality_region = manifest
625+
.data_sources
626+
.iter()
627+
.filter_map(|ds| ds.as_offchain())
628+
.map(|ds| ds.mapping.entities.iter())
629+
.chain(
630+
manifest
631+
.templates
632+
.iter()
633+
.filter_map(|ds| ds.as_offchain())
634+
.map(|ds| ds.mapping.entities.iter()),
635+
)
636+
.flatten()
637+
.cloned()
638+
.collect();
639+
623640
// Apply the subgraph versioning and deployment operations,
624641
// creating a new subgraph deployment if one doesn't exist.
625642
let deployment = DeploymentCreate::new(raw_string, &manifest, start_block)
626643
.graft(base_block)
627-
.debug(debug_fork);
644+
.debug(debug_fork)
645+
.has_causality_region(needs_causality_region);
646+
628647
deployment_store
629648
.create_subgraph_deployment(
630649
name,

‎graph/src/components/store/mod.rs‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ mod traits;
44

55
pub use entity_cache::{EntityCache, ModificationsAndCache};
66

7+
use diesel::types::{FromSql, ToSql};
78
pub use err::StoreError;
89
use itertools::Itertools;
910
pub use traits::*;
@@ -12,13 +13,14 @@ use futures::stream::poll_fn;
1213
use futures::{Async, Poll, Stream};
1314
use graphql_parser::schema as s;
1415
use serde::{Deserialize, Serialize};
16+
use std::borrow::Borrow;
1517
use std::collections::btree_map::Entry;
1618
use std::collections::{BTreeMap, BTreeSet, HashSet};
17-
use std::fmt;
1819
use std::fmt::Display;
1920
use std::sync::atomic::{AtomicUsize, Ordering};
2021
use std::sync::{Arc, RwLock};
2122
use std::time::Duration;
23+
use std::{fmt, io};
2224

2325
use crate::blockchain::Block;
2426
use crate::data::store::scalar::Bytes;
@@ -71,6 +73,12 @@ impl<'a> From<&s::InterfaceType<'a, String>> for EntityType {
7173
}
7274
}
7375

76+
impl Borrow<str> for EntityType {
77+
fn borrow(&self) -> &str {
78+
&self.0
79+
}
80+
}
81+
7482
// This conversion should only be used in tests since it makes it too
7583
// easy to convert random strings into entity types
7684
#[cfg(debug_assertions)]
@@ -82,6 +90,22 @@ impl From<&str> for EntityType {
8290

8391
impl CheapClone for EntityType {}
8492

93+
impl FromSql<diesel::sql_types::Text, diesel::pg::Pg> for EntityType {
94+
fn from_sql(bytes: Option<&[u8]>) -> diesel::deserialize::Result<Self> {
95+
let s = <String as FromSql<_, diesel::pg::Pg>>::from_sql(bytes)?;
96+
Ok(EntityType::new(s))
97+
}
98+
}
99+
100+
impl ToSql<diesel::sql_types::Text, diesel::pg::Pg> for EntityType {
101+
fn to_sql<W: io::Write>(
102+
&self,
103+
out: &mut diesel::serialize::Output<W, diesel::pg::Pg>,
104+
) -> diesel::serialize::Result {
105+
<str as ToSql<diesel::sql_types::Text, diesel::pg::Pg>>::to_sql(self.0.as_str(), out)
106+
}
107+
}
108+
85109
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)]
86110
pub struct EntityFilterDerivative(bool);
87111

‎graph/src/data/graphql/object_or_interface.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,7 @@ impl<'a> ObjectOrInterface<'a> {
117117
ObjectOrInterface::Object(object) => Some(vec![object]),
118118
ObjectOrInterface::Interface(interface) => schema
119119
.types_for_interface()
120-
.get(&interface.into())
120+
.get(interface.name.as_str())
121121
.map(|object_types| object_types.iter().collect()),
122122
}
123123
}
@@ -131,7 +131,7 @@ impl<'a> ObjectOrInterface<'a> {
131131
) -> bool {
132132
match self {
133133
ObjectOrInterface::Object(o) => o.name == typename,
134-
ObjectOrInterface::Interface(i) => types_for_interface[&i.into()]
134+
ObjectOrInterface::Interface(i) => types_for_interface[i.name.as_str()]
135135
.iter()
136136
.any(|o| o.name == typename),
137137
}

‎graph/src/data/subgraph/schema.rs‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ use hex;
55
use lazy_static::lazy_static;
66
use rand::rngs::OsRng;
77
use rand::Rng;
8+
use std::collections::BTreeSet;
89
use std::str::FromStr;
910
use std::{fmt, fmt::Display};
1011

@@ -106,6 +107,7 @@ pub struct DeploymentCreate {
106107
pub graft_base: Option<DeploymentHash>,
107108
pub graft_block: Option<BlockPtr>,
108109
pub debug_fork: Option<DeploymentHash>,
110+
pub has_causality_region: BTreeSet<EntityType>,
109111
}
110112

111113
impl DeploymentCreate {
@@ -120,6 +122,7 @@ impl DeploymentCreate {
120122
graft_base: None,
121123
graft_block: None,
122124
debug_fork: None,
125+
has_causality_region: BTreeSet::new(),
123126
}
124127
}
125128

@@ -135,6 +138,11 @@ impl DeploymentCreate {
135138
self.debug_fork = fork;
136139
self
137140
}
141+
142+
pub fn has_causality_region(mut self, has_causality_region: BTreeSet<EntityType>) -> Self {
143+
self.has_causality_region = has_causality_region;
144+
self
145+
}
138146
}
139147

140148
/// The representation of a subgraph deployment when reading an existing
@@ -158,6 +166,7 @@ pub struct SubgraphDeploymentEntity {
158166
pub reorg_count: i32,
159167
pub current_reorg_depth: i32,
160168
pub max_reorg_depth: i32,
169+
pub has_causality_region: Vec<EntityType>,
161170
}
162171

163172
#[derive(Debug)]

‎graph/src/data_source/mod.rs‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -250,6 +250,13 @@ impl<C: Blockchain> DataSourceTemplate<C> {
250250
}
251251
}
252252

253+
pub fn as_offchain(&self) -> Option<&offchain::DataSourceTemplate> {
254+
match self {
255+
Self::Onchain(_) => None,
256+
Self::Offchain(t) => Some(&t),
257+
}
258+
}
259+
253260
pub fn into_onchain(self) -> Option<C::DataSourceTemplate> {
254261
match self {
255262
Self::Onchain(ds) => Some(ds),

‎graphql/src/introspection/resolver.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,7 @@ fn interface_type_object(
135135
description: interface_type.description.clone(),
136136
fields:
137137
field_objects(schema, type_objects, &interface_type.fields),
138-
possibleTypes: schema.types_for_interface()[&interface_type.into()]
138+
possibleTypes: schema.types_for_interface()[interface_type.name.as_str()]
139139
.iter()
140140
.map(|object_type| r::Value::String(object_type.name.to_owned()))
141141
.collect::<Vec<_>>(),

‎store/postgres/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ git-testament = "0.2.0"
3333
itertools = "0.10.5"
3434
pin-utils = "0.1"
3535
hex = "0.4.3"
36+
pretty_assertions = "1.3.0"
3637

3738
[dev-dependencies]
3839
futures = "0.3"

‎store/postgres/examples/layout.rs‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ extern crate clap;
22
extern crate graph_store_postgres;
33

44
use clap::{arg, Command};
5+
use std::collections::BTreeSet;
56
use std::process::exit;
67
use std::{fs, sync::Arc};
78

@@ -145,7 +146,7 @@ pub fn main() {
145146
);
146147
let site = Arc::new(make_dummy_site(subgraph, namespace, "anet".to_string()));
147148
let catalog = ensure(
148-
Catalog::for_tests(site.clone()),
149+
Catalog::for_tests(site.clone(), BTreeSet::new()),
149150
"Failed to construct catalog",
150151
);
151152
let layout = ensure(
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
alter table subgraphs.subgraph_deployment drop column has_causality_region;

0 commit comments

Comments
 (0)