mod actions;
mod render;
use std::path::PathBuf;
use clap::{Args as ClapArgs, Subcommand, ValueEnum};
use quicknode_sdk::streams::{FilterLanguage, StreamDataset, StreamRegion, StreamStatus};
use crate::context::Ctx;
use crate::errors::CliError;
#[derive(Debug, ClapArgs)]
#[command(after_help = "Examples:\n \
qn stream list --limit 20\n \
qn stream create --name blocks --network ethereum-mainnet --dataset block \\\n \
--start 24691804 --end=-1 --region usa-east --webhook https://hook.example.com\n \
qn stream create --stream-config-file stream.json\n \
qn stream pause s-1234")]
pub struct Args {
#[command(subcommand)]
pub cmd: StreamCmd,
}
#[derive(Debug, Subcommand)]
pub enum StreamCmd {
#[command(visible_alias = "ls")]
List(ListArgs),
Create(Box<CreateArgs>),
Show {
#[arg(value_name = "STREAM_ID")]
id: String,
},
Update(UpdateArgs),
Delete {
#[arg(value_name = "STREAM_ID")]
id: String,
},
Activate {
#[arg(value_name = "STREAM_ID")]
id: String,
},
Pause {
#[arg(value_name = "STREAM_ID")]
id: String,
},
TestFilter(TestFilterArgs),
EnabledCount {
#[arg(long)]
r#type: Option<String>,
},
}
#[derive(Debug, ClapArgs)]
pub struct ListArgs {
#[arg(long)]
pub limit: Option<i64>,
#[arg(long)]
pub offset: Option<i64>,
#[arg(long)]
pub order_by: Option<String>,
#[arg(long, value_parser = ["asc", "desc"])]
pub order_direction: Option<String>,
#[arg(long)]
pub stream_type: Option<String>,
}
#[derive(Debug, ClapArgs)]
pub struct CreateArgs {
#[arg(long, conflicts_with_all = ["name", "network", "dataset", "start", "end", "region", "plan", "webhook"])]
pub stream_config_file: Option<PathBuf>,
#[arg(long, required_unless_present = "stream_config_file")]
pub name: Option<String>,
#[arg(long, required_unless_present = "stream_config_file")]
pub network: Option<String>,
#[arg(long, value_enum, required_unless_present = "stream_config_file")]
pub dataset: Option<DatasetArg>,
#[arg(long, required_unless_present = "stream_config_file")]
pub start: Option<i64>,
#[arg(
long,
allow_hyphen_values = true,
required_unless_present = "stream_config_file"
)]
pub end: Option<i64>,
#[arg(long, value_enum, required_unless_present = "stream_config_file")]
pub region: Option<RegionArg>,
#[arg(long)]
pub plan: Option<String>,
#[arg(long, required_unless_present = "stream_config_file")]
pub webhook: Option<String>,
#[arg(long)]
pub webhook_security_token: Option<String>,
#[arg(long, default_value = "none")]
pub compression: String,
#[arg(long)]
pub batch_size: Option<i64>,
#[arg(long, value_parser = clap::builder::BoolishValueParser::new())]
pub fix_block_reorgs: Option<bool>,
#[arg(long)]
pub elastic_batch_enabled: Option<bool>,
#[arg(long)]
pub notification_email: Option<String>,
#[arg(long, value_enum)]
pub status: Option<StatusArg>,
#[arg(long, conflicts_with = "filter_file")]
pub filter: Option<String>,
#[arg(long)]
pub filter_file: Option<PathBuf>,
#[arg(long, value_enum)]
pub filter_language: Option<FilterLanguageArg>,
#[arg(long)]
pub threshold_fetch_buffer: Option<i64>,
#[arg(long)]
pub keep_distance_from_tip: Option<i64>,
}
#[derive(Debug, ClapArgs)]
pub struct UpdateArgs {
#[arg(value_name = "STREAM_ID")]
pub id: String,
#[arg(long)]
pub name: Option<String>,
#[arg(long, value_enum)]
pub status: Option<StatusArg>,
#[arg(long)]
pub notification_email: Option<String>,
}
#[derive(Debug, ClapArgs)]
pub struct TestFilterArgs {
#[arg(long)]
pub network: String,
#[arg(long, value_enum)]
pub dataset: DatasetArg,
#[arg(long)]
pub block: String,
#[arg(long, conflicts_with = "filter_file")]
pub filter: Option<String>,
#[arg(long)]
pub filter_file: Option<PathBuf>,
#[arg(long, value_enum)]
pub filter_language: Option<FilterLanguageArg>,
}
#[derive(Debug, Clone, Copy, ValueEnum)]
pub enum RegionArg {
UsaEast,
EuropeCentral,
AsiaEast,
}
impl From<RegionArg> for StreamRegion {
fn from(r: RegionArg) -> Self {
match r {
RegionArg::UsaEast => StreamRegion::UsaEast,
RegionArg::EuropeCentral => StreamRegion::EuropeCentral,
RegionArg::AsiaEast => StreamRegion::AsiaEast,
}
}
}
#[derive(Debug, Clone, Copy, ValueEnum)]
pub enum DatasetArg {
Block,
BlockWithReceipts,
Transactions,
Logs,
Receipts,
TraceBlocks,
DebugTraces,
BlockWithReceiptsDebugTrace,
BlockWithReceiptsTraceBlock,
BlobSidecars,
ProgramsWithLogs,
Ledger,
Events,
Orders,
Trades,
BookUpdates,
Twap,
WriterActions,
}
impl From<DatasetArg> for StreamDataset {
fn from(d: DatasetArg) -> Self {
match d {
DatasetArg::Block => StreamDataset::Block,
DatasetArg::BlockWithReceipts => StreamDataset::BlockWithReceipts,
DatasetArg::Transactions => StreamDataset::Transactions,
DatasetArg::Logs => StreamDataset::Logs,
DatasetArg::Receipts => StreamDataset::Receipts,
DatasetArg::TraceBlocks => StreamDataset::TraceBlocks,
DatasetArg::DebugTraces => StreamDataset::DebugTraces,
DatasetArg::BlockWithReceiptsDebugTrace => StreamDataset::BlockWithReceiptsDebugTrace,
DatasetArg::BlockWithReceiptsTraceBlock => StreamDataset::BlockWithReceiptsTraceBlock,
DatasetArg::BlobSidecars => StreamDataset::BlobSidecars,
DatasetArg::ProgramsWithLogs => StreamDataset::ProgramsWithLogs,
DatasetArg::Ledger => StreamDataset::Ledger,
DatasetArg::Events => StreamDataset::Events,
DatasetArg::Orders => StreamDataset::Orders,
DatasetArg::Trades => StreamDataset::Trades,
DatasetArg::BookUpdates => StreamDataset::BookUpdates,
DatasetArg::Twap => StreamDataset::Twap,
DatasetArg::WriterActions => StreamDataset::WriterActions,
}
}
}
#[derive(Debug, Clone, Copy, ValueEnum)]
pub enum StatusArg {
Active,
Paused,
Terminated,
Completed,
Blocked,
}
impl From<StatusArg> for StreamStatus {
fn from(s: StatusArg) -> Self {
match s {
StatusArg::Active => StreamStatus::Active,
StatusArg::Paused => StreamStatus::Paused,
StatusArg::Terminated => StreamStatus::Terminated,
StatusArg::Completed => StreamStatus::Completed,
StatusArg::Blocked => StreamStatus::Blocked,
}
}
}
#[derive(Debug, Clone, Copy, ValueEnum)]
pub enum FilterLanguageArg {
Javascript,
Go,
Wasm,
}
impl From<FilterLanguageArg> for FilterLanguage {
fn from(f: FilterLanguageArg) -> Self {
match f {
FilterLanguageArg::Javascript => FilterLanguage::Javascript,
FilterLanguageArg::Go => FilterLanguage::Go,
FilterLanguageArg::Wasm => FilterLanguage::Wasm,
}
}
}
pub async fn run(args: Args, ctx: Ctx) -> Result<(), CliError> {
match args.cmd {
StreamCmd::List(a) => actions::list(a, ctx).await,
StreamCmd::Create(a) => actions::create(*a, ctx).await,
StreamCmd::Show { id } => actions::show(&id, ctx).await,
StreamCmd::Update(a) => actions::update(a, ctx).await,
StreamCmd::Delete { id } => actions::delete(&id, ctx).await,
StreamCmd::Activate { id } => actions::activate(&id, ctx).await,
StreamCmd::Pause { id } => actions::pause(&id, ctx).await,
StreamCmd::TestFilter(a) => actions::test_filter(a, ctx).await,
StreamCmd::EnabledCount { r#type } => actions::enabled_count(r#type, ctx).await,
}
}