libeufin

Integration and sandbox testing for FinTech APIs and data formats
Log | Files | Refs | Submodules | README | LICENSE

commit 4c2ffb9cc4011bdabd62bb918268b8d6b88dcc05
parent d94adf8e06df23cd3c6dee58b13198fbbefea345
Author: Antoine A <>
Date:   Fri, 24 Apr 2026 10:38:01 +0200

nexus: add tx-check testing cmd and improve BTS logic

Diffstat:
Msrc/bin/testbench.rs | 4++++
Msrc/db/list.rs | 8+-------
Msrc/ebics/administrative.rs | 21++++++++++++---------
Msrc/ebics/bts.rs | 234++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Msrc/ebics/key_management.rs | 12+++++-------
Msrc/ebics/mod.rs | 337++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Msrc/lib.rs | 5+----
Msrc/testing.rs | 44++++++++++++++++++++++++++++++++++----------
Msrc/worker.rs | 16++++++++--------
9 files changed, 532 insertions(+), 149 deletions(-)

diff --git a/src/bin/testbench.rs b/src/bin/testbench.rs @@ -65,6 +65,7 @@ pub enum NexusCmd { raw_args: Vec<String>, }, Wss, + TxCheck, Exit, } @@ -285,6 +286,9 @@ 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; + } NexusCmd::Wss => { nexus_cmd(&cfg.cfg, &format!("testing wss {log_flags}")).await; } diff --git a/src/db/list.rs b/src/db/list.rs @@ -149,13 +149,7 @@ pub async fn incoming( debtor: r.try_get("debit_payto")?, bounced: r.try_get("bounced")?, talerable: match r.try_get::<Option<CompactString>, _>("type")? { - None => { - if let Some(pending) = pending_pub { - Some(format!("pending mapped by {pending}")) - } else { - None - } - } + None => pending_pub.map(|pending| format!("pending mapped by {pending}")), Some(ty) => Some(format!( "{ty} {}{map}", r.try_get::<EddsaPublicKey, _>("metadata")? diff --git a/src/ebics/administrative.rs b/src/ebics/administrative.rs @@ -109,18 +109,21 @@ pub fn hev_msg(cfg: &EbicsHostCfg) -> String { pub fn parse_hev(xml: &[u8]) -> xml::Result<EbicsResponse<Box<[VersionNumber]>>> { Xml::parse(xml, "ebicsHEVResponse", |root| { + let s = root.one("SystemReturnCode")?; Ok(EbicsResponse { - technical_code: root.one("SystemReturnCode").one("ReturnCode").parse()?, + technical_code: s.one("ReturnCode").parse()?, + technical_text: s.one("ReportText").parse()?, bank_code: EbicsReturnCode::EBICS_OK, - content: root - .many("VersionNumber") - .map(|n| { - Ok(VersionNumber { - number: n.parse()?, - schema: n.attr("ProtocolVersion")?.into(), + content: Some( + root.many("VersionNumber") + .map(|n| { + Ok(VersionNumber { + number: n.parse()?, + schema: n.attr("ProtocolVersion")?.into(), + }) }) - }) - .collect::<xml::Result<_>>()?, + .collect::<xml::Result<_>>()?, + ), }) }) } diff --git a/src/ebics/bts.rs b/src/ebics/bts.rs @@ -27,6 +27,7 @@ use crate::{ crypto::ebics_pub_key_hash, ebics::{ EbicsResponse, PreparedUploadData, + ebics_code::EbicsReturnCode, order::{BTF, Order}, }, keys::{BankPubKeysFile, ClientPriKeysFile}, @@ -104,7 +105,7 @@ fn service(w: &mut XmlWriter, service: &BTF) { ) } -pub fn download_init( +pub fn d_init( cfg: &EbicsHostCfg, bank: &BankPubKeysFile, client: &ClientPriKeysFile, @@ -153,7 +154,7 @@ pub fn download_init( }) } -pub fn download_transfer( +pub fn d_transfer( cfg: &EbicsHostCfg, client: &ClientPriKeysFile, order: &Order, @@ -179,7 +180,7 @@ pub fn download_transfer( }) } -pub fn download_receipt( +pub fn receipt( cfg: &EbicsHostCfg, client: &ClientPriKeysFile, order: &Order, @@ -207,7 +208,7 @@ pub fn download_receipt( }) } -pub fn upload_init( +pub fn u_init( cfg: &EbicsHostCfg, bank: &BankPubKeysFile, client: &ClientPriKeysFile, @@ -260,7 +261,7 @@ pub fn upload_init( }) } -pub fn upload_transfer( +pub fn u_transfer( cfg: &EbicsHostCfg, client: &ClientPriKeysFile, order: &Order, @@ -295,41 +296,210 @@ pub struct DataEncryptionInfo { pub bank_pub_digest: Vec<u8>, } -pub struct BTSResponse { - pub tx_id: Option<CompactString>, - pub order_id: Option<CompactString>, - pub data_encryption_info: Option<DataEncryptionInfo>, - pub segment: Option<Vec<u8>>, - pub segment_number: Option<usize>, - pub nb_segments: Option<usize>, +fn expect_phase(n: Xml<'_>, phase: &str) -> xml::Result<()> { + let n = n.one("TransactionPhase")?; + if n.text() != phase { + Err(n.parse_err(format_args!("Expected phase '{phase}' got '{}'", n.text()))) + } else { + Ok(()) + } } -pub fn parse_bts(xml: &[u8]) -> xml::Result<EbicsResponse<BTSResponse>> { +pub struct DInit { + pub tx_id: CompactString, + pub data_encryption_info: DataEncryptionInfo, + pub segment: Vec<u8>, + pub nb_segments: usize, +} + +pub fn parse_d_init(xml: &[u8]) -> xml::Result<EbicsResponse<DInit>> { Xml::parse(xml, "ebicsResponse", |root| { let header = root.one_signed("header")?; let st = header.one("static")?; let mutable = header.one("mutable")?; let body = root.one("body")?; - let data = body.opt("DataTransfer")?; + + let bank_code: EbicsReturnCode = body.one_signed("ReturnCode").parse()?; + let technical_code: EbicsReturnCode = mutable.one("ReturnCode").parse()?; + let technical_text = mutable.one("ReportText").parse()?; + + if technical_code.is_error() || bank_code.is_error() { + return Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: None, + }); + } + + expect_phase(mutable, "Initialisation")?; + + let data: Xml<'_> = body.one("DataTransfer")?; + let enc_info = data.one_signed("DataEncryptionInfo")?; Ok(EbicsResponse { - technical_code: mutable.one("ReturnCode").parse()?, - bank_code: body.one_signed("ReturnCode").parse()?, - content: BTSResponse { - tx_id: st.opt("TransactionID").parse()?, - order_id: mutable.opt("OrderID").parse()?, - data_encryption_info: data - .opt_signed("DataEncryptionInfo")? - .map(|n| { - Ok(DataEncryptionInfo { - tx_key: n.one("TransactionKey").b64()?, - bank_pub_digest: n.one("EncryptionPubKeyDigest").b64()?, - }) - }) - .transpose()?, - segment: data.map(|it| it.one("OrderData").b64()).transpose()?, - segment_number: mutable.opt("SegmentNumber").parse()?, - nb_segments: st.opt("NumSegments").parse()?, - }, + technical_code, + bank_code, + technical_text, + content: Some(DInit { + tx_id: st.one("TransactionID").parse()?, + data_encryption_info: DataEncryptionInfo { + tx_key: enc_info.one("TransactionKey").b64()?, + bank_pub_digest: enc_info.one("EncryptionPubKeyDigest").b64()?, + }, + segment: data.one("OrderData").b64()?, + nb_segments: st.one("NumSegments").parse()?, + }), + }) + }) +} + +pub struct DTransfer { + pub tx_id: CompactString, + pub segment: Vec<u8>, + pub nb_segments: usize, +} + +pub fn parse_d_transfer(xml: &[u8]) -> xml::Result<EbicsResponse<DTransfer>> { + Xml::parse(xml, "ebicsResponse", |root| { + let header = root.one_signed("header")?; + let st = header.one("static")?; + let mutable = header.one("mutable")?; + let body = root.one("body")?; + + let bank_code: EbicsReturnCode = body.one_signed("ReturnCode").parse()?; + let technical_code: EbicsReturnCode = mutable.one("ReturnCode").parse()?; + let technical_text = mutable.one("ReportText").parse()?; + + if technical_code.is_error() || bank_code.is_error() { + return Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: None, + }); + } + + expect_phase(mutable, "Transfer")?; + + Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: Some(DTransfer { + tx_id: st.one("TransactionID").parse()?, + segment: body.one("DataTransfer").one("OrderData").b64()?, + nb_segments: st.one("NumSegments").parse()?, + }), + }) + }) +} + +pub struct Receipt { + pub tx_id: CompactString, +} + +pub fn parse_receipt(xml: &[u8]) -> xml::Result<EbicsResponse<Receipt>> { + Xml::parse(xml, "ebicsResponse", |root| { + let header = root.one_signed("header")?; + let st = header.one("static")?; + let mutable = header.one("mutable")?; + let body = root.one("body")?; + + let bank_code: EbicsReturnCode = body.one_signed("ReturnCode").parse()?; + let technical_code: EbicsReturnCode = mutable.one("ReturnCode").parse()?; + let technical_text = mutable.one("ReportText").parse()?; + + if technical_code.is_error() || bank_code.is_error() { + return Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: None, + }); + } + + expect_phase(mutable, "Receipt")?; + + Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: Some(Receipt { + tx_id: st.one("TransactionID").parse()?, + }), + }) + }) +} + +pub struct U { + pub tx_id: CompactString, + pub order_id: CompactString, +} + +pub fn parse_u_init(xml: &[u8]) -> xml::Result<EbicsResponse<U>> { + Xml::parse(xml, "ebicsResponse", |root| { + let header = root.one_signed("header")?; + let st = header.one("static")?; + let mutable = header.one("mutable")?; + let body = root.one("body")?; + + let bank_code: EbicsReturnCode = body.one_signed("ReturnCode").parse()?; + let technical_code: EbicsReturnCode = mutable.one("ReturnCode").parse()?; + let technical_text = mutable.one("ReportText").parse()?; + + if technical_code.is_error() || bank_code.is_error() { + return Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: None, + }); + } + + expect_phase(mutable, "Initialisation")?; + + Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: Some(U { + order_id: mutable.one("OrderID").parse()?, + tx_id: st.one("TransactionID").parse()?, + }), + }) + }) +} + +pub fn parse_u_transfer(xml: &[u8]) -> xml::Result<EbicsResponse<U>> { + Xml::parse(xml, "ebicsResponse", |root| { + let header = root.one_signed("header")?; + let st = header.one("static")?; + let mutable = header.one("mutable")?; + let body = root.one("body")?; + + let bank_code: EbicsReturnCode = body.one_signed("ReturnCode").parse()?; + let technical_code: EbicsReturnCode = mutable.one("ReturnCode").parse()?; + let technical_text = mutable.one("ReportText").parse()?; + + if technical_code.is_error() || bank_code.is_error() { + return Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: None, + }); + } + + expect_phase(mutable, "Transfer")?; + + Ok(EbicsResponse { + technical_code, + bank_code, + technical_text, + content: Some(U { + order_id: mutable.one("OrderID").parse()?, + tx_id: st.one("TransactionID").parse()?, + }), }) }) } diff --git a/src/ebics/key_management.rs b/src/ebics/key_management.rs @@ -255,14 +255,12 @@ impl EbicsClient { let res = self.post_to_bank(signed, &ctx).await?; Xml::parse(&res, "ebicsKeyManagementResponse", |root| { let body = root.one("body")?; + let mutable = root.one_signed("header").one("mutable")?; Ok(EbicsResponse { - technical_code: root - .one_signed("header") - .one("mutable") - .one("ReturnCode") - .parse()?, + technical_code: mutable.one("ReturnCode").parse()?, + technical_text: mutable.one("ReportText").parse()?, bank_code: body.one_signed("ReturnCode").parse()?, - content: if let Some(data) = body.opt("DataTransfer")? { + content: Some(if let Some(data) = body.opt("DataTransfer")? { let info = data.one_signed("DataEncryptionInfo")?; let info = DataEncryptionInfo { tx_key: info.one("TransactionKey").b64()?, @@ -273,7 +271,7 @@ impl EbicsClient { Some(decoded) } else { None - }, + }), }) }) .ctx(&ctx) diff --git a/src/ebics/mod.rs b/src/ebics/mod.rs @@ -24,6 +24,7 @@ use base64::{Engine, prelude::BASE64_STANDARD}; use compact_str::CompactString; use flate2::write::ZlibDecoder; use jiff::Timestamp; +use rand::{RngExt as _, distr::Alphanumeric}; use reqwest::{ Client, StatusCode, header::{CONTENT_TYPE, HeaderValue}, @@ -42,8 +43,9 @@ use crate::{ ebics::{ administrative::{HAA, VersionNumber, hev_msg, parse_haa, parse_hev}, bts::{ - BTSResponse, DataEncryptionInfo, download_init, download_receipt, download_transfer, - parse_bts, upload_init, upload_transfer, + DInit, DTransfer, DataEncryptionInfo, U, d_init, d_transfer, parse_d_init, + parse_d_transfer, parse_receipt, parse_u_init, parse_u_transfer, receipt, u_init, + u_transfer, }, ebics_code::EbicsReturnCode, logger::EbicsLogger, @@ -87,6 +89,7 @@ pub struct EbicsCtx<'a> { pub now: Timestamp, pub order: Cow<'a, Order>, pub phase: Option<Phase>, + pub tx_id: Option<CompactString>, } impl<'a> EbicsCtx<'a> { @@ -95,12 +98,62 @@ impl<'a> EbicsCtx<'a> { now: Timestamp::now(), order: Cow::Borrowed(order), phase: None, + tx_id: None, } } - pub fn with_phase(self, phase: Phase) -> Self { + pub fn init(self) -> Self { Self { - phase: Some(phase), + phase: Some(Phase::Init), + tx_id: None, + ..self + } + } + + pub fn interrupt(self, id: &str) -> Self { + Self { + phase: Some(Phase::Interrupt), + tx_id: Some( + self.tx_id + .filter(|it| it != id) + .unwrap_or_else(|| id.into()), + ), + ..self + } + } + + pub fn transfer(self, id: &str, segment: usize) -> Self { + Self { + phase: Some(Phase::Transfer(segment)), + tx_id: Some( + self.tx_id + .filter(|it| it != id) + .unwrap_or_else(|| id.into()), + ), + ..self + } + } + + pub fn process(self, id: &str) -> Self { + Self { + phase: Some(Phase::Process), + tx_id: Some( + self.tx_id + .filter(|it| it != id) + .unwrap_or_else(|| id.into()), + ), + ..self + } + } + + pub fn receipt(self, id: &str) -> Self { + Self { + phase: Some(Phase::Receipt), + tx_id: Some( + self.tx_id + .filter(|it| it != id) + .unwrap_or_else(|| id.into()), + ), ..self } } @@ -108,11 +161,19 @@ impl<'a> EbicsCtx<'a> { impl std::fmt::Display for EbicsCtx<'_> { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let Self { order, phase, .. } = self; + let Self { + order, + phase, + tx_id, + .. + } = self; write!(f, "{order}")?; if let Some(phase) = phase { write!(f, " {phase}")?; } + if let Some(tx_id) = tx_id { + write!(f, " {tx_id}")?; + } Ok(()) } } @@ -126,7 +187,7 @@ pub struct EbicsError { impl std::fmt::Display for EbicsError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { let Self { ctx, kind } = self; - write!(f, "{ctx} {kind}") + write!(f, "{ctx}: {kind}") } } @@ -188,6 +249,7 @@ impl EbicsErrKind { ctx: Box::new(EbicsCtx { now: ctx.now, phase: ctx.phase, + tx_id: ctx.tx_id.clone(), order: Cow::Owned(ctx.order.as_ref().clone()), }), kind: self, @@ -197,18 +259,22 @@ impl EbicsErrKind { pub struct EbicsResponse<T> { pub technical_code: EbicsReturnCode, pub bank_code: EbicsReturnCode, - pub content: T, + pub technical_text: String, + pub content: Option<T>, } impl<T> EbicsResponse<T> { fn ok_or_fail(self) -> Result<T, EbicsErrKind> { - if self.technical_code.is_error() || self.bank_code.is_error() { + if let Some(content) = self.content + && !self.technical_code.is_error() + && !self.bank_code.is_error() + { + Ok(content) + } else { Err(EbicsErrKind::Code { technical: self.technical_code, bank: self.bank_code, }) - } else { - Ok(self.content) } } } @@ -249,22 +315,20 @@ impl EbicsClient { } /** POST an EBICS BTS request [xmlReq] using [client] returning a validated and parsed XML response */ - async fn post_bts(&self, xml: String, ctx: &EbicsCtx<'_>) -> Result<BTSResponse, EbicsError> { + pub async fn post_bts<T>( + &self, + xml: String, + ctx: &EbicsCtx<'_>, + parse: impl FnOnce(&[u8]) -> xml::Result<EbicsResponse<T>>, + ) -> Result<T, EbicsError> { let xml = self.post_to_bank(xml, ctx).await?; - // TODO verify ebics - let res = parse_bts(&xml).ctx(ctx)?; - // TODO phase in logs ? + // TODO verify ebics signature + let res = parse(&xml).ctx(ctx)?; trace!(target: "ebics", - "{ctx}{}: {} {}", - std::fmt::from_fn(|f| { - if let Some(tx_id) = &res.content.tx_id { - write!(f, " {tx_id}") - } else { - Ok(()) - } - }), + "{ctx}: {} {} - {}", res.technical_code, - res.bank_code + res.bank_code, + res.technical_text ); res.ok_or_fail().ctx(ctx) } @@ -324,10 +388,10 @@ impl EbicsClient { })); // Close interrupted - ctx = ctx.with_phase(Phase::Interrupt); while let Some(tx_id) = ebics_first(db).await.ctx(&ctx)? { - let xml = download_receipt(&self.cfg, client, order, &tx_id, false); - if let Err(e) = self.post_bts(xml, &ctx).await { + let ctx = EbicsCtx::new(order).interrupt(&tx_id); + let xml = receipt(&self.cfg, client, order, &tx_id, false); + if let Err(e) = self.post_bts(xml, &ctx, parse_d_init).await { if !matches!( e.kind, // Transaction already closed or expired - EBICS protocol error @@ -347,51 +411,35 @@ impl EbicsClient { } // Init phase - ctx = ctx.with_phase(Phase::Init); - let xml = download_init(&self.cfg, bank, client, order, range); - let BTSResponse { + ctx = ctx.init(); + let xml = d_init(&self.cfg, bank, client, order, range); + let DInit { tx_id, nb_segments, segment, data_encryption_info, - .. - } = self.post_bts(xml, &ctx).await?; - // TODO DAO add - let (tx_id, nb_segments, segment, encr_info) = ( - tx_id - .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) - .ctx(&ctx)?, - nb_segments - .ok_or_else(|| EbicsErrKind::Custom("missing num segments".into())) - .ctx(&ctx)?, - segment - .ok_or_else(|| EbicsErrKind::Custom("missing OrderData".into())) - .ctx(&ctx)?, - data_encryption_info - .ok_or_else(|| EbicsErrKind::Custom("missing EncryptionInfo".into())) - .ctx(&ctx)?, - ); + } = self.post_bts(xml, &ctx, parse_d_init).await?; ebics_register(db, &tx_id).await.ctx(&ctx)?; // Transfer phase let mut segments = vec![segment]; for segment_nb in 2..=nb_segments { - ctx = ctx.with_phase(Phase::Transfer(segment_nb)); - let xml = download_transfer(&self.cfg, client, order, nb_segments, segment_nb, &tx_id); - let BTSResponse { segment, .. } = self.post_bts(xml, &ctx).await?; - segments.push(segment.unwrap()); // TODO error + ctx = ctx.transfer(&tx_id, segment_nb); + let xml = d_transfer(&self.cfg, client, order, nb_segments, segment_nb, &tx_id); + let DTransfer { segment, .. } = self.post_bts(xml, &ctx, parse_d_transfer).await?; + segments.push(segment); } // Processing phase - ctx = ctx.with_phase(Phase::Process); - let payload = decrypt_and_decompress_payload(&client.enc, encr_info, segments); + ctx = ctx.process(&tx_id); + let payload = decrypt_and_decompress_payload(&client.enc, data_encryption_info, segments); self.logger.log_payload(&ctx, &payload, order.file_type())?; let res = processing(payload).await.ctx(&ctx); // Receipt phase - ctx = ctx.with_phase(Phase::Receipt); - let xml = download_receipt(&self.cfg, client, order, &tx_id, res.is_ok() && !peek); - if self.post_bts(xml, &ctx).await.is_ok() { + ctx = ctx.receipt(&tx_id); + let xml = receipt(&self.cfg, client, order, &tx_id, res.is_ok() && !peek); + if self.post_bts(xml, &ctx, parse_receipt).await.is_ok() { ebics_remove(db, &tx_id).await.ok(); } @@ -419,26 +467,15 @@ impl EbicsClient { let payload = prepare_upload_payload(&self.cfg, client, bank, payload); // Init phase - ctx = ctx.with_phase(Phase::Init); - let xml = upload_init(&self.cfg, bank, client, order, &payload); - let BTSResponse { - tx_id, order_id, .. - } = self.post_bts(xml, &ctx).await?; - - let (tx_id, order_id) = ( - tx_id - .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) - .ctx(&ctx)?, - order_id - .ok_or_else(|| EbicsErrKind::Custom("missing order ID".into())) - .ctx(&ctx)?, - ); + ctx = ctx.init(); + let xml = u_init(&self.cfg, bank, client, order, &payload); + let U { tx_id, order_id } = self.post_bts(xml, &ctx, parse_u_init).await?; // Transfer phase for segment_nb in 1..=payload.nb_segments() { - ctx = ctx.with_phase(Phase::Transfer(segment_nb)); - let xml = upload_transfer(&self.cfg, client, order, &tx_id, &payload, segment_nb); - self.post_bts(xml, &ctx).await?; + ctx = ctx.transfer(&tx_id, segment_nb); + let xml = u_transfer(&self.cfg, client, order, &tx_id, &payload, segment_nb); + self.post_bts(xml, &ctx, parse_u_transfer).await?; } Ok(order_id) @@ -525,3 +562,159 @@ fn prepare_upload_payload( payload, } } + +#[derive(Debug)] +pub struct TxCheckResult { + pub concurrent_fetch_and_fetch: bool, + pub concurrent_fetch_and_submit: bool, + pub concurrent_submit_and_submit: bool, + pub idempotent_close: bool, +} + +/** + * Test EBICS implementation's transactions semantic: + * - Can two fetch transactions run concurrently ? + * - Can a fetch & submit transactions run concurrently ? + * - Can two submit transactions run concurrently ? + * - Is closing a submit transaction idempotent + */ +pub async fn tx_check( + ebics: &EbicsClient, + db: &PgPool, + client: &ClientPriKeysFile, + bank: &BankPubKeysFile, + fetch: &Order, + submit: &Order, +) -> anyhow::Result<TxCheckResult> { + let mut result = TxCheckResult { + concurrent_fetch_and_fetch: false, + concurrent_fetch_and_submit: false, + concurrent_submit_and_submit: false, + idempotent_close: false, + }; + + let ctx = EbicsCtx::new(fetch).init(); + let DInit { tx_id, .. } = ebics + .post_bts( + d_init(&ebics.cfg, bank, client, fetch, &None), + &ctx, + parse_d_init, + ) + .await?; + ebics_register(db, &tx_id).await?; + { + let ctx = EbicsCtx::new(fetch).init(); + match ebics + .post_bts( + d_init(&ebics.cfg, bank, client, fetch, &None), + &ctx, + parse_d_init, + ) + .await + { + Ok(DInit { tx_id, .. }) => { + ebics_register(db, &tx_id).await?; + result.concurrent_fetch_and_fetch = true; + let ctx = ctx.receipt(&tx_id); + ebics + .post_bts( + receipt(&ebics.cfg, client, fetch, &tx_id, false), + &ctx, + parse_receipt, + ) + .await?; + ebics_remove(db, &tx_id).await?; + } + Err(e) => { + if !matches!(e.kind, EbicsErrKind::Code { .. }) { + return Err(e.into()); + } else { + debug!(target: "testing", "concurrent_fetch_and_fetch {e}") + } + } + } + } + + { + let ctx = EbicsCtx::new(submit).init(); + let random_string: String = rand::rng() + .sample_iter(&Alphanumeric) + .take(2000000) + .map(char::from) + .collect(); + let payload = prepare_upload_payload(&ebics.cfg, client, bank, &random_string); + match ebics + .post_bts( + u_init(&ebics.cfg, bank, client, submit, &payload), + &ctx, + parse_u_init, + ) + .await + { + Ok(U { tx_id, .. }) => { + result.concurrent_fetch_and_submit = true; + let ctx = ctx.transfer(&tx_id, 1); + ebics + .post_bts( + u_transfer(&ebics.cfg, client, fetch, &tx_id, &payload, 1), + &ctx, + parse_u_transfer, + ) + .await?; + let ctx = EbicsCtx::new(submit).init(); + if let Err(e) = ebics + .post_bts( + u_init(&ebics.cfg, bank, client, submit, &payload), + &ctx, + parse_u_init, + ) + .await + { + if !matches!(e.kind, EbicsErrKind::Code { .. }) { + return Err(e.into()); + } else { + debug!(target: "testing", "concurrent_submit_and_submit {e}") + } + } else { + result.concurrent_submit_and_submit = true; + } + } + Err(e) => { + if !matches!(e.kind, EbicsErrKind::Code { .. }) { + return Err(e.into()); + } else { + debug!(target: "testing", "concurrent_fetch_and_submit {e}") + } + } + } + } + + // Close first fetch + let ctx = ctx.receipt(&tx_id); + ebics + .post_bts( + receipt(&ebics.cfg, client, fetch, &tx_id, false), + &ctx, + parse_receipt, + ) + .await?; + + ebics_remove(db, &tx_id).await?; + + // Close first fetch again + let ctx = ctx.interrupt(&tx_id); + if let Err(e) = ebics + .post_bts( + receipt(&ebics.cfg, client, fetch, &tx_id, false), + &ctx, + parse_receipt, + ) + .await + { + debug!(target: "testing", "idempotent_close {e}") + } else { + result.idempotent_close = true + } + + Ok(result) +} diff --git a/src/lib.rs b/src/lib.rs @@ -743,10 +743,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { cmd.run(&pool, &cfg.currency).await?; } Cmd::Config(cmd) => cmd.run(&cfg)?, - Cmd::Testing(cmd) => { - let pool = pool(&cfg).await?; - cmd.run(cfg, &pool).await?; - } + Cmd::Testing(cmd) => cmd.run(cfg).await?, } Ok(()) } diff --git a/src/testing.rs b/src/testing.rs @@ -20,7 +20,6 @@ use anyhow::{anyhow, bail}; use compact_str::CompactString; use jiff::{Timestamp, civil::Date}; -use sqlx::PgPool; use taler_common::{ config::Config, types::{ @@ -34,9 +33,11 @@ use tracing::debug; use crate::{ EbicsClient, EbicsLogs, InTx, config::NexusCfg, + db::pool, ebics::{ EbicsErrKind, - order::{BTF, Order}, + order::{BTF, Order, OrderDoc}, + tx_check, }, keys::expect_full_keys, list::ListCmd, @@ -116,8 +117,11 @@ pub enum TestingCmd { dry_run: bool, }, /// Check transaction semantic - TxCheck, - /// "Listen to EBICS instant notification over websocket + TxCheck { + #[clap(flatten)] + logs: EbicsLogs, + }, + /// Listen to EBICS instant notification over websocket Wss { #[clap(flatten)] logs: EbicsLogs, @@ -125,7 +129,7 @@ pub enum TestingCmd { } impl TestingCmd { - pub async fn run(self, cfg: Config, db: &PgPool) -> anyhow::Result<()> { + pub async fn run(self, cfg: Config) -> anyhow::Result<()> { match self { TestingCmd::Iban(cmd) => cmd.run()?, TestingCmd::FakeIncoming { @@ -134,6 +138,7 @@ impl TestingCmd { subject, payto, } => { + let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; let subject = payto .subject @@ -154,7 +159,7 @@ impl TestingCmd { ); } register_incoming( - db, + &db, &cfg.ingest()?, &InTx { id: InId::new(None, Some(rand_ebics_id()), None), @@ -168,8 +173,9 @@ impl TestingCmd { .await?; } TestingCmd::List(list_cmd) => { + let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; - list_cmd.run(db, &cfg.currency).await?; + list_cmd.run(&db, &cfg.currency).await?; } TestingCmd::EbicsBtd { ty, @@ -184,6 +190,7 @@ impl TestingCmd { peek, dry_run, } => { + let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; let order = Order::from_parts( &ty, @@ -201,7 +208,7 @@ impl TestingCmd { let ebics = EbicsClient::new(&cfg, logs)?; ebics .download( - db, + &db, &client, &bank, &order, @@ -217,8 +224,25 @@ impl TestingCmd { ) .await?; } - TestingCmd::TxCheck => todo!(), + TestingCmd::TxCheck { logs } => { + let db = pool(&cfg).await?; + let cfg = NexusCfg::parse(cfg)?; + let ebics = EbicsClient::new(&cfg, logs)?; + let (client, bank) = expect_full_keys(cfg.keys()?)?; + let dialect = cfg.ebics()?.dialect.standard(); + let res = tx_check( + &ebics, + &db, + &client, + &bank, + &dialect.downloads(&OrderDoc::acknowledgement)[0], + &dialect.direct_debit(), + ) + .await?; + println!("{res:?}") + } TestingCmd::Wss { logs } => { + let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; let ebics = EbicsClient::new(&cfg, logs)?; let (client, bank) = expect_full_keys(cfg.keys()?)?; @@ -228,7 +252,7 @@ impl TestingCmd { debug!(target: "testing", "{orders:?}") } }); - listen_for_notification(&ebics, db, &client, &bank, sender).await + listen_for_notification(&ebics, &db, &client, &bank, sender).await } } Ok(()) diff --git a/src/worker.rs b/src/worker.rs @@ -73,9 +73,9 @@ pub async fn register_incoming( }); if res.completed || res.new { - info!("{fmt}") + info!(target: "worker", "{fmt}") } else { - debug!("{fmt}") + debug!(target: "worker", "{fmt}") } }; let bounce = async |cause: &str| { @@ -121,7 +121,7 @@ pub async fn register_incoming( .await?; match res { IncomingBounceRegistrationResult::Talerable => { - warn!("{payment} tried to bounce a talerable transaction"); + warn!(target: "worker", "{payment} tried to bounce a talerable transaction"); } IncomingBounceRegistrationResult::Success(res) => { log_res(res, "", &format!(": {cause}")); @@ -199,12 +199,12 @@ pub async fn register_outgoing( let res = register_out_tx(db, payment, metadata.as_ref()).await?; if res.new { if res.initiated { - info!("{payment}"); + info!(target: "worker", "{payment}"); } else { - warn!("{payment} recovered"); + warn!(target: "worker", "{payment} recovered"); } } else { - debug!("{payment} already seen"); + debug!(target: "worker", "{payment} already seen"); } Ok(res) } @@ -214,7 +214,7 @@ pub async fn register_outgoing_batch( currency: &Currency, batch: &OutBatch, ) -> sqlx::Result<()> { - info!("{batch}"); + info!(target: "worker", "{batch}"); let txs = unsettled_tx_in_batch(db, currency, &batch.msg_id, &batch.execution_time).await?; for tx in txs { register_outgoing(db, &tx).await?; @@ -224,7 +224,7 @@ pub async fn register_outgoing_batch( pub async fn register_tx(db: &PgPool, cfg: &NexusIngestCfg, tx: &Tx) -> sqlx::Result<()> { if tx.execution_time() < &cfg.ignore_txs_before { - debug!("IGNORE {tx}"); + debug!(target: "worker", "IGNORE {tx}"); } else { match tx { Tx::In(payment) => {