Skip to content
Merged
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 2 additions & 1 deletion chain/substreams/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ dirs-next = "2.0"
anyhow = "1.0"
tiny-keccak = "1.5.0"
hex = "0.4.3"
semver = "1.0.14"
semver = "1.0.12"
base64 = "0.13.0"

itertools = "0.10.5"

Expand Down
21 changes: 4 additions & 17 deletions chain/substreams/examples/substreams.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,7 @@ use graph::blockchain::block_stream::BlockStreamEvent;
use graph::blockchain::substreams_block_stream::SubstreamsBlockStream;
use graph::prelude::{info, tokio, DeploymentHash, Registry};
use graph::tokio_stream::StreamExt;
use graph::{
env::env_var,
firehose::FirehoseEndpoint,
log::logger,
substreams::{self},
};
use graph::{env::env_var, firehose::FirehoseEndpoint, log::logger, substreams};
use graph_chain_substreams::mapper::Mapper;
use graph_core::MetricsRegistry;
use prost::Message;
Expand All @@ -27,7 +22,7 @@ async fn main() -> Result<(), Error> {

let endpoint = env_var(
"SUBSTREAMS_ENDPOINT",
"https://api-dev.streamingfast.io".to_string(),
"https://api.streamingfast.io".to_string(),
);

let package_file = env_var("SUBSTREAMS_PACKAGE", "".to_string());
Expand Down Expand Up @@ -79,17 +74,9 @@ async fn main() -> Result<(), Error> {
Ok(block_stream_event) => match block_stream_event {
BlockStreamEvent::Revert(_, _) => {}
BlockStreamEvent::ProcessBlock(block_with_trigger, _) => {
let changes = block_with_trigger.block;
for change in changes.entity_changes {
info!(&logger, "----- Entity -----");
info!(
&logger,
"name: {} operation: {}", change.entity, change.operation
);
for change in block_with_trigger.block.changes.entity_changes {
for field in change.fields {
info!(&logger, "field: {}, type: {}", field.name, field.value_type);
info!(&logger, "new value: {}", hex::encode(field.new_value));
info!(&logger, "old value: {}", hex::encode(field.old_value));
info!(&logger, "field: {:?}", field);
}
}
}
Expand Down
46 changes: 25 additions & 21 deletions chain/substreams/proto/codec.proto
Original file line number Diff line number Diff line change
Expand Up @@ -2,20 +2,16 @@ syntax = "proto3";

package substreams.entity.v1;

message EntitiesChanges {
bytes block_id = 1;
uint64 block_number = 2;
bytes prev_block_id = 3;
uint64 prev_block_number = 4;
repeated EntityChange entityChanges = 5;
message EntityChanges {
repeated EntityChange entity_changes = 5;
}

message EntityChange {
string entity = 1;
bytes id = 2;
string id = 2;
uint64 ordinal = 3;
enum Operation {
UNSET = 0; // Protobuf default should not be used, this is used so that the consume can ensure that the value was actually specified
UNSET = 0; // Protobuf default should not be used, this is used so that the consume can ensure that the value was actually specified
CREATE = 1;
UPDATE = 2;
DELETE = 3;
Expand All @@ -24,19 +20,27 @@ message EntityChange {
repeated Field fields = 5;
}

message Value {
oneof typed {
int32 int32 = 1;
string bigdecimal = 2;
string bigint = 3;
string string = 4;
bytes bytes = 5;
bool bool = 6;

//reserved 7 to 9; // For future types

Array array = 10;
}
}

message Array {
repeated Value value = 1;
}

message Field {
string name = 1;
enum Type {
UNSET = 0; // Protobuf default should not be used, this is used so that the consume can ensure that the value was actually specified
BIGDECIMAL = 1;
BIGINT = 2;
INT = 3; // int32
BYTES = 4;
STRING = 5;
}
Type value_type = 2;
bytes new_value = 3;
bool new_value_null = 4;
bytes old_value = 5;
bool old_value_null = 6;
optional Value new_value = 3;
optional Value old_value = 5;
}
22 changes: 13 additions & 9 deletions chain/substreams/src/chain.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use crate::{data_source::*, Block, TriggerData, TriggerFilter, TriggersAdapter};
use crate::{data_source::*, EntityChanges, TriggerData, TriggerFilter, TriggersAdapter};
use anyhow::Error;
use core::fmt;
use graph::firehose::FirehoseEndpoints;
Expand All @@ -17,19 +17,23 @@ use graph::{
};
use std::{str::FromStr, sync::Arc};

#[derive(Default, Debug, Clone)]
pub struct Block {
pub hash: BlockHash,
pub number: BlockNumber,
pub changes: EntityChanges,
}

impl blockchain::Block for Block {
fn ptr(&self) -> BlockPtr {
return BlockPtr {
hash: BlockHash(Box::from(self.block_id.clone())),
number: self.block_number as i32,
};
BlockPtr {
hash: self.hash.clone(),
number: self.number,
}
}

fn parent_ptr(&self) -> Option<BlockPtr> {
Some(BlockPtr {
hash: BlockHash(Box::from(self.prev_block_id.clone())),
number: self.prev_block_number as i32,
})
None
}
}

Expand Down
3 changes: 1 addition & 2 deletions chain/substreams/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,9 +8,8 @@ pub mod mapper;

pub use block_stream::BlockStreamBuilder;
pub use chain::*;
pub use codec::EntitiesChanges as Block;
pub use codec::EntityChanges;
pub use data_source::*;
pub use trigger::*;

pub use codec::field::Type as FieldType;
pub use codec::Field;
61 changes: 44 additions & 17 deletions chain/substreams/src/mapper.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
use crate::{Block, Chain, TriggerData};
use crate::{Block, Chain, EntityChanges, TriggerData};
use graph::blockchain::block_stream::SubstreamsError::{
MultipleModuleOutputError, UnexpectedStoreDeltaOutput,
};
use graph::blockchain::block_stream::{
BlockStreamEvent, BlockWithTriggers, FirehoseCursor, SubstreamsError, SubstreamsMapper,
};
use graph::prelude::{async_trait, BlockNumber, BlockPtr, Logger};
use graph::prelude::{async_trait, BlockHash, BlockNumber, BlockPtr, Logger};
use graph::substreams::module_output::Data;
use graph::substreams::{BlockScopedData, ForkStep};
use graph::substreams::{BlockScopedData, Clock, ForkStep};
use prost::Message;

pub struct Mapper {}
Expand All @@ -19,28 +19,50 @@ impl SubstreamsMapper<Chain> for Mapper {
_logger: &Logger,
block_scoped_data: &BlockScopedData,
) -> Result<Option<BlockStreamEvent<Chain>>, SubstreamsError> {
let step = ForkStep::from_i32(block_scoped_data.step).unwrap_or_else(|| {
let BlockScopedData {
outputs,
clock,
step,
cursor: _,
} = block_scoped_data;

let step = ForkStep::from_i32(*step).unwrap_or_else(|| {
panic!(
"unknown step i32 value {}, maybe you forgot update & re-regenerate the protobuf definitions?",
block_scoped_data.step
step
)
});

if block_scoped_data.outputs.len() == 0 {
if outputs.len() == 0 {
return Ok(None);
}

if block_scoped_data.outputs.len() > 1 {
return Err(MultipleModuleOutputError());
if outputs.len() > 1 {
return Err(MultipleModuleOutputError);
}

//todo: handle step
let module_output = &block_scoped_data.outputs[0];
let cursor = &block_scoped_data.cursor;

match module_output.data.as_ref().unwrap() {
Data::MapOutput(msg) => {
let changes: Block = Message::decode(msg.value.as_slice()).unwrap();
let clock = match clock {
Some(clock) => clock,
None => return Err(SubstreamsError::MissingClockError),
};

let Clock {
id: hash,
number,
timestamp: _,
} = clock;

let hash: BlockHash = hash.as_str().try_into()?;
let number: BlockNumber = *number as BlockNumber;

match module_output.data.as_ref() {
Some(Data::MapOutput(msg)) => {
let changes: EntityChanges = Message::decode(msg.value.as_slice())
.map_err(SubstreamsError::DecodingError)?;

use ForkStep::*;
match step {
Expand All @@ -52,14 +74,18 @@ impl SubstreamsMapper<Chain> for Mapper {

// TODO(filipe): Fix once either trigger data can be empty
// or we move the changes into trigger data.
BlockWithTriggers::new(changes, vec![TriggerData {}]),
BlockWithTriggers::new(
Block {
hash,
number,
changes,
},
vec![TriggerData {}],
),
FirehoseCursor::from(cursor.clone()),
))),
StepUndo => {
let parent_ptr = BlockPtr {
hash: changes.prev_block_id.clone().into(),
number: changes.prev_block_number as BlockNumber,
};
let parent_ptr = BlockPtr { hash, number };

Ok(Some(BlockStreamEvent::Revert(
parent_ptr,
Expand All @@ -71,7 +97,8 @@ impl SubstreamsMapper<Chain> for Mapper {
}
}
}
Data::StoreDeltas(_) => Err(UnexpectedStoreDeltaOutput()),
Some(Data::StoreDeltas(_)) => Err(UnexpectedStoreDeltaOutput),
_ => Err(SubstreamsError::ModuleOutputNotPresentOrUnexpected),
}
}
}
75 changes: 39 additions & 36 deletions chain/substreams/src/protobuf/substreams.entity.v1.rs
Original file line number Diff line number Diff line change
@@ -1,22 +1,14 @@
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct EntitiesChanges {
#[prost(bytes="vec", tag="1")]
pub block_id: ::prost::alloc::vec::Vec<u8>,
#[prost(uint64, tag="2")]
pub block_number: u64,
#[prost(bytes="vec", tag="3")]
pub prev_block_id: ::prost::alloc::vec::Vec<u8>,
#[prost(uint64, tag="4")]
pub prev_block_number: u64,
pub struct EntityChanges {
#[prost(message, repeated, tag="5")]
pub entity_changes: ::prost::alloc::vec::Vec<EntityChange>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct EntityChange {
#[prost(string, tag="1")]
pub entity: ::prost::alloc::string::String,
#[prost(bytes="vec", tag="2")]
pub id: ::prost::alloc::vec::Vec<u8>,
#[prost(string, tag="2")]
pub id: ::prost::alloc::string::String,
#[prost(uint64, tag="3")]
pub ordinal: u64,
#[prost(enumeration="entity_change::Operation", tag="4")]
Expand All @@ -37,32 +29,43 @@ pub mod entity_change {
}
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Value {
#[prost(oneof="value::Typed", tags="1, 2, 3, 4, 5, 6, 10")]
pub typed: ::core::option::Option<value::Typed>,
}
/// Nested message and enum types in `Value`.
pub mod value {
#[derive(Clone, PartialEq, ::prost::Oneof)]
pub enum Typed {
#[prost(int32, tag="1")]
Int32(i32),
#[prost(string, tag="2")]
Bigdecimal(::prost::alloc::string::String),
#[prost(string, tag="3")]
Bigint(::prost::alloc::string::String),
#[prost(string, tag="4")]
String(::prost::alloc::string::String),
#[prost(bytes, tag="5")]
Bytes(::prost::alloc::vec::Vec<u8>),
#[prost(bool, tag="6")]
Bool(bool),
//reserved 7 to 9; // For future types

#[prost(message, tag="10")]
Array(super::Array),
}
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Array {
#[prost(message, repeated, tag="1")]
pub value: ::prost::alloc::vec::Vec<Value>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Field {
#[prost(string, tag="1")]
pub name: ::prost::alloc::string::String,
#[prost(enumeration="field::Type", tag="2")]
pub value_type: i32,
#[prost(bytes="vec", tag="3")]
pub new_value: ::prost::alloc::vec::Vec<u8>,
#[prost(bool, tag="4")]
pub new_value_null: bool,
#[prost(bytes="vec", tag="5")]
pub old_value: ::prost::alloc::vec::Vec<u8>,
#[prost(bool, tag="6")]
pub old_value_null: bool,
}
/// Nested message and enum types in `Field`.
pub mod field {
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
#[repr(i32)]
pub enum Type {
/// Protobuf default should not be used, this is used so that the consume can ensure that the value was actually specified
Unset = 0,
Bigdecimal = 1,
Bigint = 2,
/// int32
Int = 3,
Bytes = 4,
String = 5,
}
#[prost(message, optional, tag="3")]
pub new_value: ::core::option::Option<Value>,
#[prost(message, optional, tag="5")]
pub old_value: ::core::option::Option<Value>,
}
Loading