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
119 changes: 59 additions & 60 deletions src/reducers/asset_holders_by_asset_id.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
use pallas::ledger::traverse::MultiEraOutput;
use pallas::ledger::traverse::{MultiEraBlock, OutputRef, Subject};
use pallas::ledger::traverse::{Asset, MultiEraOutput};
use pallas::ledger::traverse::{MultiEraBlock, OutputRef};
use serde::Deserialize;

use pallas::crypto::hash::Hash;
use crate::{crosscut, model, prelude::*};
use pallas::crypto::hash::Hash;

use crate::crosscut::epochs::block_epoch;
use std::str::FromStr;
Expand All @@ -18,33 +18,38 @@ pub struct Config {
pub key_prefix: Option<String>,
pub filter: Option<crosscut::filters::Predicate>,
pub aggr_by: Option<AggrType>,
pub policy_ids_hex: Option<Vec<String>>, // if specified only those policy ids as hex will be taken into account, if not all policy ids will be indexed

/// Policies to match
///
/// If specified only those policy ids as hex will be taken into account, if
/// not all policy ids will be indexed.
pub policy_ids_hex: Option<Vec<String>>,
}

pub struct Reducer {
config: Config,
policy: crosscut::policies::RuntimePolicy,
chain: crosscut::ChainWellKnownInfo,
policy_ids: Option<Vec<Hash<28>>>
policy_ids: Option<Vec<Hash<28>>>,
}

impl Reducer {
fn config_key(&self, asset_id: String, epoch_no: u64) -> String {
fn config_key(&self, subject: String, epoch_no: u64) -> String {
let def_key_prefix = "asset_holders_by_asset_id";

match &self.config.aggr_by {
Some(aggr_type) if matches!(aggr_type, AggrType::Epoch) => {
return match &self.config.key_prefix {
Some(prefix) => format!("{}.{}.{}", prefix, asset_id, epoch_no),
None => format!("{}.{}", def_key_prefix.to_string(), asset_id),
Some(prefix) => format!("{}.{}.{}", prefix, subject, epoch_no),
None => format!("{}.{}", def_key_prefix.to_string(), subject),
};
},
}
_ => {
return match &self.config.key_prefix {
Some(prefix) => format!("{}.{}", prefix, asset_id),
None => format!("{}.{}", def_key_prefix.to_string(), asset_id),
Some(prefix) => format!("{}.{}", prefix, subject),
None => format!("{}.{}", def_key_prefix.to_string(), subject),
};
},
}
};
}

Expand All @@ -71,26 +76,21 @@ impl Reducer {

let address = utxo.address().map(|addr| addr.to_string()).or_panic()?;

for asset in utxo.assets().iter() {
let sub = &asset.subject;
let quantity = &asset.quantity;

let delta = *quantity as i64 * (-1);
for asset in utxo.assets() {
match asset {
Asset::NativeAsset(policy_id, _, quantity) => {
if self.is_policy_id_accepted(&policy_id) {
let subject = asset.subject();
let key = self.config_key(subject, epoch_no);
let delta = quantity as i64 * (-1);

match sub {
Subject::NativeAsset(policy_id, asset_name) => {
if self.is_policy_id_accepted(policy_id) {
let asset_id = format!("{}{}", policy_id, asset_name);
let crdt =
model::CRDTCommand::SortedSetRemove(key, address.to_string(), delta);

let key = self.config_key(asset_id, epoch_no);

let crdt = model::CRDTCommand::SortedSetRemove(key, address.to_string(), delta);

output.send(gasket::messaging::Message::from(crdt))?;
}

}
_ => {},
_ => (),
};
}

Expand All @@ -103,28 +103,26 @@ impl Reducer {
epoch_no: u64,
output: &mut super::OutputPort,
) -> Result<(), gasket::error::Error> {
let address = tx_output.address().map(|addr| addr.to_string()).or_panic()?;
let address = tx_output
.address()
.map(|addr| addr.to_string())
.or_panic()?;

for asset in tx_output.assets() {
match asset {
Asset::NativeAsset(policy_id, _, quantity) => {
if self.is_policy_id_accepted(&policy_id) {
let subject = asset.subject();
let key = self.config_key(subject, epoch_no);
let delta = quantity as i64;

let crdt =
model::CRDTCommand::SortedSetAdd(key, address.to_string(), delta);

for asset in tx_output.assets().iter() {
let sub = &asset.subject;
let quantity = &asset.quantity;

let delta = *quantity as i64;

match sub {
Subject::NativeAsset(policy_id, asset_name) => {
if self.is_policy_id_accepted(policy_id) {

let asset_id = format!("{}{}", policy_id, asset_name);

let key = self.config_key(asset_id, epoch_no);

let crdt = model::CRDTCommand::SortedSetAdd(key, address.to_string(), delta);

output.send(gasket::messaging::Message::from(crdt))?;
}
}
_ => {},
_ => {}
};
}

Expand Down Expand Up @@ -156,23 +154,22 @@ impl Reducer {
}

impl Config {
pub fn plugin(self,
chain: &crosscut::ChainWellKnownInfo,
policy: &crosscut::policies::RuntimePolicy,
) -> super::Reducer {

let policy_ids: Option<Vec<Hash<28>>> = match &self.policy_ids_hex {
Some(pids) => {
let ps = pids
pub fn plugin(
self,
chain: &crosscut::ChainWellKnownInfo,
policy: &crosscut::policies::RuntimePolicy,
) -> super::Reducer {
let policy_ids: Option<Vec<Hash<28>>> = match &self.policy_ids_hex {
Some(pids) => {
let ps = pids
.iter()
.map(|pid| Hash::<28>::from_str(pid)
.expect("invalid policy_id"))
.map(|pid| Hash::<28>::from_str(pid).expect("invalid policy_id"))
.collect();

Some(ps)
},
None => None,
};
Some(ps)
}
None => None,
};

let reducer = Reducer {
config: self,
Expand All @@ -186,5 +183,7 @@ impl Config {
}

// How to query
// 127.0.0.1:6379> ZRANGEBYSCORE "asset_holders_by_asset_id.5d9d887de76a2c9d057b3e5d34d5411f7f8dc4d54f0c06e8ed2eb4a9494e4459" 1 +inf
// 127.0.0.1:6379> ZRANGEBYSCORE
// "asset_holders_by_asset_id.
// 5d9d887de76a2c9d057b3e5d34d5411f7f8dc4d54f0c06e8ed2eb4a9494e4459" 1 +inf
// 1) "addr1q8lmu79hgm3sppz8dta3aftf0cwh2v2eja56wqvzqy4jj0zjt7qgvj7saxdxve35c4ehuxuam4czlz9fw6ls7zr4as9s609d7u"
22 changes: 22 additions & 0 deletions testdrive/custom/adahandle.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
[source]
type = "N2N"
address = "preprod-node.world.dev.cardano.org:30000"

[[reducers]]
type = "AddressByAdaHandle"
key_prefix = "AddressByAdaHandle"
policy_id_hex = "f0ff48bbb7bbe9d59a40f1ce90e9e9d0ff5002ec48f232b49ca0fb9a"

[storage]
type = "Redis"
connection_params = "redis://127.0.0.1:6379"

[chain]
type = "PreProd"

[intersect]
type = "Point"
value = [
8261225,
"103ff3cfec6e388db803fc10dbdecb3026b4a51382008d01ff3f9121be6fa6e4",
]
2 changes: 1 addition & 1 deletion testdrive/custom/start.sh
Original file line number Diff line number Diff line change
@@ -1 +1 @@
RUST_LOG=info cargo run --all-features --bin scrolls -- daemon --console tui --config ./daemon.toml
RUST_LOG=info cargo run --all-features --bin scrolls -- daemon --console plain`` --config ./adahandle.toml