commit ee6824621640a553408d4cc97c8549b3a788305f
parent b7fdacad8d43a68daabd6c05082939b57548e83d
Author: Antoine A <>
Date: Fri, 24 Apr 2026 10:38:01 +0200
nexus: finish submit and improve testbench
Diffstat:
3 files changed, 90 insertions(+), 36 deletions(-)
diff --git a/src/bin/testbench.rs b/src/bin/testbench.rs
@@ -29,7 +29,7 @@ use libeufin::{
run,
};
use owo_colors::OwoColorize as _;
-use reedline::{Prompt, Reedline, Signal};
+use reedline::{FileBackedHistory, Prompt, Reedline, Signal};
use taler_common::{config::Config, log::taler_logger, types::payto::TransferIbanPayto};
use tracing::Level;
use tracing_subscriber::util::SubscriberInitExt as _;
@@ -52,16 +52,19 @@ pub struct TestbenchCmd {
#[command(name = "shell", no_binary_name = true)]
/// Run integration tests on banks provider
pub enum NexusCmd {
- /// Reset EBICS keys
ResetKeys,
- /// Reset DB
ResetDb,
- // Initiate a new transaction
Tx,
- /// Fetch all documents
- Fetch,
- Submit,
+ Fetch {
+ #[arg(trailing_var_arg = true, allow_hyphen_values = true)]
+ raw_args: Vec<String>,
+ },
+ Submit {
+ #[arg(trailing_var_arg = true, allow_hyphen_values = true)]
+ raw_args: Vec<String>,
+ },
List {
+ #[arg(trailing_var_arg = true, allow_hyphen_values = true)]
raw_args: Vec<String>,
},
Wss,
@@ -97,7 +100,12 @@ pub async fn nexus_cmd(cfg: &Config, cmd: &str) -> bool {
let args = std::iter::once("libeufin_nexus").chain(parts.iter().map(|it| it.as_str()));
match libeufin::Args::try_parse_from(args) {
- Ok(cmd) => check(run(cfg.clone(), cmd.cmd).await),
+ Ok(cmd) => {
+ tokio::select! {
+ res = run(cfg.clone(), cmd.cmd) => check(res),
+ _ = tokio::signal::ctrl_c() => false
+ }
+ }
Err(e) => {
println!("Error: {}", e);
false
@@ -177,7 +185,18 @@ async fn main() -> anyhow::Result<()> {
cmd.platform
);
- let mut line_editor = Reedline::create();
+ let history = Box::new(
+ FileBackedHistory::with_file(
+ 1000,
+ match cmd.component {
+ Component::Ebisync => ".ebisync_history",
+ Component::Nexus => ".nexus_history",
+ }
+ .into(),
+ )
+ .expect("Error configuring history with file"),
+ );
+ let mut line_editor = Reedline::create().with_history(history);
let prompt = BenchPrompt {
prompt: format!("{:?} {}", cmd.component, cmd.platform),
};
@@ -207,8 +226,7 @@ async fn main() -> anyhow::Result<()> {
_ => todo!("{}", cfg.currency),
};
let payto = TransferIbanPayto::from_str(payto).unwrap();
- let log_flags = format!("--debug-ebics testbench/test/{}", cmd.platform);
- let transient_flag = format!("{log_flags} --transient");
+ let ebics_log = format!("--debug-ebics testbench/test/{}", cmd.platform);
loop {
// Automatic setup
{
@@ -222,7 +240,7 @@ async fn main() -> anyhow::Result<()> {
|| bank.map(|it| !it.accepted).unwrap_or(true)
{
step("Run EBICS setup");
- if !nexus_cmd(&cfg.cfg, &format!("ebics-setup {log_flags}")).await {
+ if !nexus_cmd(&cfg.cfg, &format!("ebics-setup {ebics_log}")).await {
if let Some(settings) = settings {
let client =
load_client_keys(ebics.client_priv_keys_path.as_ref()).unwrap();
@@ -245,7 +263,6 @@ async fn main() -> anyhow::Result<()> {
}
}
let Signal::Success(buf) = line_editor.read_line(&prompt).unwrap() else {
- err("^C");
break;
};
match NexusCmd::try_parse_from(buf.split_whitespace()) {
@@ -253,11 +270,25 @@ async fn main() -> anyhow::Result<()> {
NexusCmd::ResetDb => {
nexus_cmd(&cfg.cfg, "dbinit -r").await;
}
- NexusCmd::Fetch => {
- nexus_cmd(&cfg.cfg, &format!("ebics-fetch {transient_flag}")).await;
+ NexusCmd::Fetch { raw_args } => {
+ nexus_cmd(
+ &cfg.cfg,
+ &format!(
+ "ebics-fetch {ebics_log} {}",
+ shlex::try_join(raw_args.iter().map(|it| it.as_str())).unwrap()
+ ),
+ )
+ .await;
}
- NexusCmd::Submit => {
- nexus_cmd(&cfg.cfg, &format!("ebics-submit {transient_flag}")).await;
+ NexusCmd::Submit { raw_args } => {
+ nexus_cmd(
+ &cfg.cfg,
+ &format!(
+ "ebics-submit {ebics_log} {}",
+ shlex::try_join(raw_args.iter().map(|it| it.as_str())).unwrap()
+ ),
+ )
+ .await;
}
NexusCmd::Tx => {
nexus_cmd(
@@ -287,10 +318,10 @@ async fn main() -> anyhow::Result<()> {
std::fs::remove_file(&ebics.bank_pub_keys_path)?;
}
NexusCmd::TxCheck => {
- nexus_cmd(&cfg.cfg, &format!("testing tx-check {log_flags}")).await;
+ nexus_cmd(&cfg.cfg, &format!("testing tx-check {ebics_log}")).await;
}
NexusCmd::Wss => {
- nexus_cmd(&cfg.cfg, &format!("testing wss {log_flags}")).await;
+ nexus_cmd(&cfg.cfg, &format!("testing wss {ebics_log}")).await;
}
NexusCmd::Exit => return Ok(()),
},
diff --git a/src/db/initiated.rs b/src/db/initiated.rs
@@ -308,13 +308,13 @@ pub async fn order_step(db: &PgPool, order_id: &str, msg: &str) -> sqlx::Result<
"
UPDATE initiated_outgoing_batches
SET status = 'pending', status_msg = $1
- WHERE order_id = $1 AND {PENDING}
+ WHERE order_id = $2 AND {PENDING}
RETURNING initiated_outgoing_batch_id
"
))
.bind(msg)
.bind(order_id)
- .try_map(|r: PgRow| r.try_get_u32(0))
+ .try_map(|r: PgRow| r.try_get_u64(0))
.fetch_optional(&mut *tx)
.await?;
if let Some(batch_id) = batch_id {
@@ -347,7 +347,7 @@ pub async fn order_success(db: &PgPool, order_id: &str) -> sqlx::Result<Option<S
"
))
.bind(order_id)
- .try_map(|r: PgRow| Ok((r.try_get_u32(0)?, r.try_get(1)?)))
+ .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?)))
.fetch_optional(&mut *tx)
.await?;
if let Some((batch_id, _)) = &res {
@@ -534,24 +534,24 @@ mod test {
// Check batch status
let msg_id: String = sqlx::query(
"
- SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=?
+ SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=$1
"
).bind(batch_id as i64)
.try_map(|r: PgRow| {
let msg_id: String = r.try_get("message_id")?;
- assert_eq!((batch_status, batch_msg), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}");
+ assert_eq!((batch_status, Some(batch_msg).filter(|it| !it.is_empty())), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}");
Ok(msg_id)
}).fetch_one(&db).await.unwrap();
// Check tx status
sqlx::query(
"
- SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=?
+ SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=$1
"
).bind(batch_id as i64).try_map(|r: PgRow| {
let end_to_end_id: &str = r.try_get("end_to_end_id")?;
let expected = match end_to_end_id {
- "TX" => (tx_status, tx_msg),
- "TX_SETTLED" => (settled_status, settled_msg),
+ "TX" => (tx_status, Some(tx_msg).filter(|it| !it.is_empty())),
+ "TX_SETTLED" => (settled_status, Some(settled_msg).filter(|it| !it.is_empty())),
_ =>panic!("Unexpected tx $endToEndId")
};
assert_eq!(expected,
@@ -587,10 +587,7 @@ mod test {
check_order_tx(order_id, status, msg, status).await;
};
- async fn test<F>(db: &PgPool, lambda: impl FnOnce(u64) -> F)
- where
- F: Future<Output = ()>,
- {
+ async fn test(db: &PgPool, lambda: impl AsyncFnOnce(u64)) {
// Reset DB
sqlx::query("DELETE FROM initiated_outgoing_transactions")
.execute(db)
@@ -644,7 +641,8 @@ mod test {
.find(|it| it.msg_id == "BATCH")
.unwrap()
.id,
- );
+ )
+ .await;
// Check witness status is unaltered
sqlx::query(
"
@@ -709,7 +707,7 @@ mod test {
.1
.as_deref()
);
- check_order("ORDER", permanent_failure, "late msg").await;
+ check_order("ORDER", permanent_failure, "failure msg").await;
order_step(&db, "ORDER", "late msg").await.unwrap();
check_order("ORDER", permanent_failure, "failure msg").await;
})
diff --git a/src/lib.rs b/src/lib.rs
@@ -350,6 +350,7 @@ pub async fn ebics_submit(
.iter()
.map(|tx| {
let creditor = FullIbanPayto::from_str(tx.creditor.as_ref().as_str()).unwrap();
+ // TODO handle missing name ?
Pain001Tx {
creditor,
amount: tx.amount,
@@ -365,10 +366,12 @@ pub async fn ebics_submit(
let submit_all = async || -> anyhow::Result<()> {
let standard = cfg.ebics()?.dialect.standard();
+
+ // Find a supported debit order
let mut instant_order = standard.instant_direct_debit();
let debit_order = standard.direct_debit();
- // Create batch if nescessary
+ // Create batch if necessary
batch_initiated(
db,
&Timestamp::now(),
@@ -376,6 +379,8 @@ pub async fn ebics_submit(
submit_cfg.require_ack,
)
.await?;
+
+ // Send submittable batches
for batch in initiated_submittable(db, &cfg.currency).await? {
debug!(target: "ebics-submit", "Submitting batch {}", batch.msg_id);
let res = async {
@@ -419,10 +424,30 @@ pub async fn ebics_submit(
Ok(())
};
if transient {
- debug!(target: "ebics-submit", "Transient mode: submitting what found and returning.");
+ debug!(target: "ebics-submit", "Transient mode: submitting what found and returning");
submit_all().await
} else {
- todo!()
+ debug!(target: "ebics-submit", "Running with a frequency of {}", submit_cfg.frequency_raw);
+ loop {
+ let now = Timestamp::now();
+ let success = match submit_all().await {
+ Ok(_) => true,
+ Err(e) => {
+ error!(target: "ebics-submit", "{e}");
+ false
+ }
+ };
+ if let Err(e) = update_task_status(db, SUBMIT_TASK_KEY, &now, success).await {
+ warn!(target: "ebics-submit", "{e}");
+ }
+ tokio::time::sleep(Duration::from_millis(
+ Timestamp::now()
+ .duration_until(now + submit_cfg.frequency)
+ .abs()
+ .as_millis() as u64,
+ ))
+ .await;
+ }
}
}