commit d0741663e0eb238449527bbee1682d4376fecb4f parent 344a0399178b0efe0bd5f47aa7e2dd5095d455fc Author: Florian Dold <dold@taler.net> Date: Tue, 8 Sep 2026 00:46:10 +0200 services: apply DD102 supervision and retain operation retries Classify configuration and credential errors as status 6 and propagate terminal worker failures. Keep API retries and database reconnection local with exponential backoff. Retain healthy clients and sessions, and keep Cyclos notification recovery independent of polling. Poll Wise balances independently so a failing balance cannot stall the others. Attempt all balances in transient mode and reject an empty worker balance configuration. Preserve APNs bounded API and local database retries. Use ten-second systemd restart delays, bind activation sockets to their services, and request PostgreSQL through product targets. Remove socket service runtime limits for systemd 257, preserve service state on upgrade, and explicitly stop targets and slices on package removal. Diffstat:
29 files changed, 617 insertions(+), 239 deletions(-)
diff --git a/adapters/taler-cyclos/src/worker.rs b/adapters/taler-cyclos/src/worker.rs @@ -26,7 +26,7 @@ use taler_common::{ config::Config, types::amount::{self, Currency}, }; -use tokio::{join, sync::Notify}; +use tokio::sync::Notify; use tracing::{debug, error, info, trace, warn}; use crate::{ @@ -58,6 +58,26 @@ pub enum WorkerError { pub type WorkerResult = Result<(), WorkerError>; +/// Retry an operation, not a terminated task. A closed pool cannot recover, +/// and injected interruptions must remain visible to the caller/supervisor. +fn retry_delay( + err: &WorkerError, + jitter: &mut ExpoBackoffDecorr, + frequency: Duration, +) -> Option<Duration> { + match err { + WorkerError::Db(sqlx::Error::PoolClosed) | WorkerError::Injected(_) => None, + WorkerError::Concurrency => Some(jitter.backoff().max(Duration::from_secs(15))), + WorkerError::Api(e) + if matches!(*e.err, CyclosErr::Input(InputError::Validation { .. })) => + { + // Invalid input must not be resubmitted at notification speed. + Some(jitter.backoff().max(frequency)) + } + _ => Some(jitter.backoff()), + } +} + pub async fn run_worker( cfg: &Config, pool: &PgPool, @@ -95,37 +115,13 @@ pub async fn run_worker( }; let worker = async { let mut jitter = ExpoBackoffDecorr::default(); - let mut skip_notifications: bool = true; loop { - info!(target: "worker", "running at initialisation"); let res: WorkerResult = async { let db = &mut PgListener::connect_with(pool).await?; - - // Listen to all channels db.listen_all(["transfer"]).await?; - + info!(target: "worker", "database listener connected"); loop { - if !skip_notifications { - tokio::select! { - _ = tokio::time::sleep(cfg.frequency) => { - info!(target: "worker", "running at frequency"); - } - res = db.try_recv() => { - let mut ntf = res?; - // Conflate all notifications - while let Some(n) = ntf { - debug!(target: "worker", "notification from {}", n.channel()); - ntf = db.next_buffered(); - } - info!(target: "worker", "running at db trigger"); - } - _ = notification.notified() => { - info!(target: "worker", "running at notification trigger"); - } - }; - } - skip_notifications = false; - Worker { + let result = Worker { client: &client, db: db.acquire().await?, account_type_id: *cfg.account_type_id, @@ -134,35 +130,51 @@ pub async fn run_worker( currency: cfg.currency, } .run() - .await?; - jitter.reset(); + .await; + match result { + Ok(()) => jitter.reset(), + Err(err @ WorkerError::Db(_)) => return Err(err), + Err(err) => { + let Some(delay) = retry_delay(&err, &mut jitter, cfg.frequency) else { + return Err(err); + }; + error!(target: "worker", ?delay, "synchronization pass failed: {err}"); + tokio::time::sleep(delay).await; + continue; + } + } + tokio::select! { + _ = tokio::time::sleep(cfg.frequency) => { + info!(target: "worker", "running at frequency"); + } + res = db.try_recv() => { + let mut ntf = res?; + while let Some(n) = ntf { + debug!(target: "worker", "notification from {}", n.channel()); + ntf = db.next_buffered(); + } + } + _ = notification.notified() => { + info!(target: "worker", "running at notification trigger"); + } + } } } .await; let err = res.unwrap_err(); - error!(target: "worker", "{err}"); - - match err { - WorkerError::Concurrency => { - // This error won't resolve by itself easily and it mean we are actually making progress - // in another worker so we can jitter more aggressively - tokio::time::sleep(Duration::from_secs(15)).await; - skip_notifications = false; - } - WorkerError::Api(ApiErr { ctx: _, err }) - if matches!(*err, CyclosErr::Input(InputError::Validation { .. })) => - { - // In case of validation failure we do not want to retry right away as it can DOS the service - skip_notifications = false; - } - WorkerError::Api(_) | WorkerError::Db(_) | WorkerError::Injected(_) => { - skip_notifications = true; - } - } - tokio::time::sleep(jitter.backoff()).await; + let Some(delay) = retry_delay(&err, &mut jitter, cfg.frequency) else { + break Err(err); + }; + error!(target: "worker", ?delay, "database session failed: {err}"); + tokio::time::sleep(delay).await; } }; - join!(watcher, worker); // TODO try_join + // The optional notification listener reconnects independently. Polling + // continues during its outages; a terminal worker result is not swallowed. + tokio::select! { + result = worker => result?, + _ = watcher => unreachable!("notification listener does not return"), + } Ok(()) } @@ -523,3 +535,39 @@ pub fn extract_tx_info(tx: HistoryItem) -> Tx { }) } } + +#[cfg(test)] +mod retry_tests { + use super::*; + + #[test] + fn database_outages_retry_but_terminal_interruptions_propagate() { + let mut jitter = ExpoBackoffDecorr::default(); + let err = WorkerError::Db(sqlx::Error::Io( + std::io::ErrorKind::ConnectionRefused.into(), + )); + let delay = retry_delay(&err, &mut jitter, Duration::from_secs(60)).unwrap(); + assert!(delay >= Duration::from_millis(400)); + assert!(delay <= Duration::from_secs(30)); + assert!( + retry_delay( + &WorkerError::Db(sqlx::Error::PoolClosed), + &mut jitter, + Duration::ZERO + ) + .is_none() + ); + assert!( + retry_delay( + &WorkerError::Injected(InjectedErr("worker interrupted")), + &mut jitter, + Duration::ZERO + ) + .is_none() + ); + assert!( + retry_delay(&WorkerError::Concurrency, &mut jitter, Duration::ZERO).unwrap() + >= Duration::from_secs(15) + ); + } +} diff --git a/adapters/taler-magnet-bank/src/worker.rs b/adapters/taler-magnet-bank/src/worker.rs @@ -60,6 +60,16 @@ pub enum WorkerError { pub type WorkerResult = Result<(), WorkerError>; +/// Retry an operation, not a terminated task. A closed pool cannot recover, +/// and injected interruptions must remain visible to the caller/supervisor. +fn retry_delay(err: &WorkerError, jitter: &mut ExpoBackoffDecorr) -> Option<Duration> { + match err { + WorkerError::Db(sqlx::Error::PoolClosed) | WorkerError::Injected(_) => None, + WorkerError::Concurrency => Some(jitter.backoff().max(Duration::from_secs(15))), + _ => Some(jitter.backoff()), + } +} + pub async fn run_worker( cfg: &Config, pool: &PgPool, @@ -67,7 +77,12 @@ pub async fn run_worker( transient: bool, ) -> anyhow::Result<()> { let cfg = WorkerCfg::parse(cfg)?; - let keys = setup::load(&cfg)?; + let keys = setup::load(&cfg).map_err(|err| taler_common::config::ValueErr::Invalid { + ty: "keys file".into(), + section: "magnet-bank-worker".into(), + option: "KEYS_FILE".into(), + err: err.to_string(), + })?; let client = AuthClient::new(client, &cfg.api_url, &cfg.consumer).upgrade(&keys.access_token); if transient { @@ -89,20 +104,28 @@ pub async fn run_worker( } let mut jitter = ExpoBackoffDecorr::default(); + // Account discovery is an API operation; retry it without reloading keys + // or rebuilding the HTTP client. + let account = loop { + match client.account(cfg.payto.bban()).await { + Ok(account) => break account, + Err(err) => { + error!(target: "worker", "account lookup failed: {err}"); + tokio::time::sleep(jitter.backoff()).await; + } + } + }; + jitter.reset(); loop { + // Only database/session failures reconstruct the listener. API failures + // retry a pass with the existing client and database session. let res: WorkerResult = async { - let account = client.account(cfg.payto.bban()).await?; let db = &mut PgListener::connect_with(pool).await?; - - // Listen to all channels db.listen_all(["transfer"]).await?; - - info!(target: "worker", "running at initialisation"); - + info!(target: "worker", "database listener connected"); loop { - debug!(target: "worker", "running"); - Worker { + let result = Worker { client: &client, db: db.acquire().await?, account_number: &account.number, @@ -113,36 +136,40 @@ pub async fn run_worker( ignore_bounces_before: cfg.ignore_bounces_before, } .run() - .await?; - jitter.reset(); - - // Wait for notifications or sync timeout - if let Ok(res) = tokio::time::timeout(cfg.frequency, db.try_recv()).await { - let mut ntf = res?; - // Conflate all notifications - while let Some(n) = ntf { - debug!(target: "worker", "notification from {}", n.channel()); - ntf = db.next_buffered(); + .await; + match result { + Ok(()) => jitter.reset(), + Err(err @ WorkerError::Db(_)) => return Err(err), + Err(err) => { + let Some(delay) = retry_delay(&err, &mut jitter) else { + return Err(err); + }; + error!(target: "worker", ?delay, "synchronization pass failed: {err}"); + tokio::time::sleep(delay).await; + continue; } - - if ntf.is_some() { - info!(target: "worker", "running at db trigger"); - } else { - info!(target: "worker", "running at frequency"); + } + // Notifications may accelerate successful polling, but never + // bypass the backoff after a failed operation. + match tokio::time::timeout(cfg.frequency, db.try_recv()).await { + Ok(res) => { + let mut ntf = res?; + while let Some(n) = ntf { + debug!(target: "worker", "notification from {}", n.channel()); + ntf = db.next_buffered(); + } } + Err(_) => info!(target: "worker", "running at frequency"), } } } .await; let err = res.unwrap_err(); - error!(target: "worker", "{err}"); - - if matches!(err, WorkerError::Concurrency) { - // This error won't resolve by itself easily and it mean we are actually making progress - // in another worker so we can jitter more aggressively - tokio::time::sleep(Duration::from_secs(15)).await; - } - tokio::time::sleep(jitter.backoff()).await; + let Some(delay) = retry_delay(&err, &mut jitter) else { + return Err(err.into()); + }; + error!(target: "worker", ?delay, "database session failed: {err}"); + tokio::time::sleep(delay).await; } } @@ -642,3 +669,30 @@ pub fn parse_bounce_outgoing(subject: &str) -> Result<u32, BounceSubjectErr> { let id: u32 = id.parse()?; Ok(id) } + +#[cfg(test)] +mod retry_tests { + use super::*; + + #[test] + fn database_outages_retry_but_terminal_interruptions_propagate() { + let mut jitter = ExpoBackoffDecorr::default(); + let err = WorkerError::Db(sqlx::Error::Io( + std::io::ErrorKind::ConnectionRefused.into(), + )); + let delay = retry_delay(&err, &mut jitter).unwrap(); + assert!(delay >= Duration::from_millis(400)); + assert!(delay <= Duration::from_secs(30)); + assert!(retry_delay(&WorkerError::Db(sqlx::Error::PoolClosed), &mut jitter).is_none()); + assert!( + retry_delay( + &WorkerError::Injected(InjectedErr("worker interrupted")), + &mut jitter + ) + .is_none() + ); + assert!( + retry_delay(&WorkerError::Concurrency, &mut jitter).unwrap() >= Duration::from_secs(15) + ); + } +} diff --git a/adapters/taler-wise/Cargo.toml b/adapters/taler-wise/Cargo.toml @@ -30,3 +30,8 @@ tokio.workspace = true taler-test-utils.workspace = true regex.workspace = true uuid = { version = "1.24.0", features = ["serde"] } + +futures-util = { version = "0.3", default-features = false, features = ["alloc"] } + +[dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } diff --git a/adapters/taler-wise/src/config.rs b/adapters/taler-wise/src/config.rs @@ -134,11 +134,35 @@ pub struct WorkerCfg { impl WorkerCfg { pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { let s = cfg.section("wise-worker"); - Ok(Self { + let cfg = Self { profile_id: s.number("PROFILE_ID").require()?, token: s.str("TOKEN").require()?, frequency: s.duration("FREQUENCY").require()?, balances: balances(cfg)?, - }) + }; + if cfg.balances.is_empty() { + return Err(ValueErr::Missing { + ty: "balance ID".into(), + section: "wise-balance-*".into(), + option: "ID".into(), + }); + } + Ok(cfg) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn worker_requires_a_balance() { + let cfg = + Config::from_mem("[wise-worker]\nPROFILE_ID=1\nTOKEN=test\nFREQUENCY=1 min\n").unwrap(); + assert!(matches!( + WorkerCfg::parse(&cfg), + Err(ValueErr::Missing { section, option, .. }) + if section == "wise-balance-*" && option == "ID" + )); } } diff --git a/adapters/taler-wise/src/worker.rs b/adapters/taler-wise/src/worker.rs @@ -14,7 +14,9 @@ TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> */ -use std::sync::LazyLock; +use std::{sync::LazyLock, time::Duration}; + +use futures_util::future::{join_all, try_join_all}; use http_client::ApiErr; use jiff::Timestamp; @@ -25,7 +27,7 @@ use taler_common::{ExpoBackoffDecorr, config::Config, types::payto::BankID}; use tracing::{error, info, trace, warn}; use crate::{ - config::WorkerCfg, + config::{WiseBalance, WorkerCfg}, db::{AddIncomingResult, TxIn, register_tx_in}, payto::WiseAccount, wise_api::{ @@ -61,6 +63,36 @@ fn parse_account(account_str: &str) -> Option<WiseAccount> { None } +/// Each balance owns its retry schedule. A slow or failing balance must not +/// delay other balances, and success elsewhere must not reset its backoff. +async fn poll_balance( + balance_id: u32, + frequency: Duration, + transient: bool, + mut sync: impl AsyncFnMut() -> WorkerResult, +) -> WorkerResult { + let mut jitter = ExpoBackoffDecorr::default(); + loop { + let result = sync().await; + if transient { + return result; + } + let delay = match result { + Ok(()) => { + jitter.reset(); + frequency + } + Err(err @ WorkerError::Db(sqlx::Error::PoolClosed)) => return Err(err), + Err(err) => { + let delay = jitter.backoff(); + error!(target: "worker", balance_id, ?delay, "balance synchronization failed: {err}"); + delay + } + }; + tokio::time::sleep(delay).await; + } +} + pub async fn run_worker( cfg: &Config, pool: &PgPool, @@ -69,99 +101,186 @@ pub async fn run_worker( ) -> anyhow::Result<()> { let cfg = WorkerCfg::parse(cfg)?; let client = Client::new(client, &cfg.token); - let mut jitter = ExpoBackoffDecorr::default(); - loop { - if !transient { - info!(target: "worker", "running at initialisation"); + let cfg = &cfg; + let client = &client; + let workers = cfg.balances.iter().map(|balance| { + poll_balance(balance.id, cfg.frequency, transient, async move || { + sync_balance(cfg, balance, pool, client).await + }) + }); + if transient { + // Attempt every balance once, even if another balance fails. + for result in join_all(workers).await { + result?; } - let res: WorkerResult = async { - loop { - // Sync - for balance in &cfg.balances { - let stmt = client - .balance_statement(cfg.profile_id, balance.id, &balance.currency, "2026-07-16T00:00:00.000Z".parse().unwrap(), Timestamp::now() ) - .await - .unwrap(); - let now = Timestamp::now(); - for tx in stmt.transactions { - match tx.direction { - Direction::Debit => { - // TODO support outgoing transaction + } else { + // These are operation loops, not detached tasks: panics and terminal + // failures still reach the process, cancelling the remaining loops. + try_join_all(workers).await?; + } + Ok(()) +} + +async fn sync_balance( + cfg: &WorkerCfg, + balance: &WiseBalance, + pool: &PgPool, + client: &Client<'_>, +) -> WorkerResult { + let stmt = client + .balance_statement( + cfg.profile_id, + balance.id, + &balance.currency, + "2026-07-16T00:00:00.000Z".parse().unwrap(), + Timestamp::now(), + ) + .await?; + let now = Timestamp::now(); + for tx in stmt.transactions { + match tx.direction { + Direction::Debit => { + // TODO support outgoing transaction + } + Direction::Credit => { + let subject = parse_incoming_unstructured(&tx.details.payment_reference); + // Parse sender account + let payto = parse_account(&tx.details.sender_account); + // + let t = TxIn { + balance_id: balance.id, + wise_ref: Some(tx.reference_number), + amount: tx.amount.into(), + subject: tx.details.payment_reference, + name: tx.details.sender_name, + debtor: payto, + value_at: tx.date, + }; + let subject = &match subject { + Ok(IncomingSubject::Key(key)) => Some(key), + Ok(IncomingSubject::AdminBalanceAdjust) | Err(_) => None, + }; + let failure = match register_tx_in(pool, &t, subject, &now).await? { + AddIncomingResult::Success { new, .. } => { + if new { + info!(target: "worker", "in {t}"); + if t.debtor.is_none() { + warn!(target: "worker", "Couldn't parse creditor account from '{}'", tx.details.sender_account) } - Direction::Credit => { - let subject = - parse_incoming_unstructured(&tx.details.payment_reference); - // Parse sender account - let payto = parse_account(&tx.details.sender_account); - // - let t = TxIn { - balance_id: balance.id, - wise_ref: Some(tx.reference_number), - amount: tx.amount.into(), - subject: tx.details.payment_reference, - name: tx.details.sender_name, - debtor: payto, - value_at: tx.date, - }; - let subject = &match subject { - Ok(IncomingSubject::Key(key)) => Some(key), - Ok(IncomingSubject::AdminBalanceAdjust) | Err(_) => None, - }; - let failure = - match register_tx_in(pool, &t, subject, &now).await? { - AddIncomingResult::Success { new, .. } => { - if new { - info!(target: "worker", "in {t}"); - if t.debtor.is_none() { - warn!(target: "worker", "Couldn't parse creditor account from '{}'", tx.details.sender_account) - } - } else { - trace!(target: "worker", "in {t} already seen"); - } - continue; - } - AddIncomingResult::ReservePubReuse => "reserve pub reuse", - AddIncomingResult::UnknownMapping => "unknown mapping", - AddIncomingResult::MappingReuse => "mapping reuse", - }; - - match register_tx_in(pool, &t, &None, &now).await? { - AddIncomingResult::Success { new, .. } => { - if new { - info!(target: "worker", "in {t}: {failure}"); - if t.debtor.is_none() { - warn!(target: "worker", "Couldn't parse creditor account from '{}'", tx.details.sender_account) - } - } else { - trace!(target: "worker", "in {t} already seen: {failure}"); - } - continue; - } - AddIncomingResult::ReservePubReuse - | AddIncomingResult::UnknownMapping - | AddIncomingResult::MappingReuse => unreachable!(), - }; + } else { + trace!(target: "worker", "in {t} already seen"); + } + continue; + } + AddIncomingResult::ReservePubReuse => "reserve pub reuse", + AddIncomingResult::UnknownMapping => "unknown mapping", + AddIncomingResult::MappingReuse => "mapping reuse", + }; + + match register_tx_in(pool, &t, &None, &now).await? { + AddIncomingResult::Success { new, .. } => { + if new { + info!(target: "worker", "in {t}: {failure}"); + if t.debtor.is_none() { + warn!(target: "worker", "Couldn't parse creditor account from '{}'", tx.details.sender_account) } + } else { + trace!(target: "worker", "in {t} already seen: {failure}"); } + continue; } - } + AddIncomingResult::ReservePubReuse + | AddIncomingResult::UnknownMapping + | AddIncomingResult::MappingReuse => unreachable!(), + }; + } + } + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use futures_util::FutureExt; + use std::cell::Cell; + use tokio::time::Instant; + + fn unavailable() -> WorkerError { + sqlx::Error::Io(std::io::ErrorKind::ConnectionRefused.into()).into() + } - // then wait - if transient { - break Ok(()); + #[tokio::test(start_paused = true)] + async fn failed_balance_retries_without_stalling_other_balances() { + let attempts = Cell::new(0); + let healthy = Cell::new(0); + let start = Instant::now(); + let bad = poll_balance(1, Duration::from_secs(60), false, async || { + attempts.set(attempts.get() + 1); + if attempts.get() == 3 { + Err(sqlx::Error::PoolClosed.into()) + } else { + Err(unavailable()) + } + }); + let good = poll_balance(2, Duration::from_millis(100), false, async || { + healthy.set(healthy.get() + 1); + Ok(()) + }); + let result = try_join_all([bad.boxed_local(), good.boxed_local()]).await; + assert!(matches!( + result, + Err(WorkerError::Db(sqlx::Error::PoolClosed)) + )); + assert_eq!(attempts.get(), 3); + assert!(healthy.get() > attempts.get()); + assert!(start.elapsed() >= Duration::from_millis(800)); + } + + #[tokio::test(start_paused = true)] + async fn success_resets_only_this_balances_backoff() { + let attempts = Cell::new(0); + let failed_at = Cell::new(Instant::now()); + let result = poll_balance(1, Duration::from_millis(10), false, async || { + attempts.set(attempts.get() + 1); + match attempts.get() { + 11 => Ok(()), + 12 => { + failed_at.set(Instant::now()); + Err(unavailable()) } - jitter.reset(); - tokio::time::sleep(cfg.frequency).await; - info!(target: "worker", "running at frequency"); + 13 => Err(sqlx::Error::PoolClosed.into()), + _ => Err(unavailable()), } - } + }) .await; - if transient { - res?; - return Ok(()); - } - let err = res.unwrap_err(); - error!(target: "worker", "{err}"); - tokio::time::sleep(jitter.backoff()).await; + assert!(result.is_err()); + let delay = failed_at.get().elapsed(); + assert!(delay >= Duration::from_millis(400)); + assert!(delay < Duration::from_secs(1)); + } + + #[tokio::test(start_paused = true)] + async fn transient_attempts_each_balance_once_and_reports_failure() { + let attempts = Cell::new(0); + let results = join_all([ + poll_balance( + 1, + Duration::from_secs(60), + true, + async || Err(unavailable()), + ) + .boxed_local(), + poll_balance(2, Duration::from_secs(60), true, async || { + tokio::time::sleep(Duration::from_millis(100)).await; + attempts.set(attempts.get() + 1); + Ok(()) + }) + .boxed_local(), + ]) + .await; + assert!(results[0].is_err()); + assert!(results[1].is_ok()); + assert_eq!(attempts.get(), 1); } } diff --git a/common/taler-common/src/lib.rs b/common/taler-common/src/lib.rs @@ -64,7 +64,7 @@ pub fn taler_main( Ok(cfg) => cfg, Err(err) => { error!(target: "config", "{}", err); - std::process::exit(1); + std::process::exit(6); } }; @@ -78,7 +78,52 @@ pub fn taler_main( let result = runtime.block_on(app(&cfg)); if let Err(err) = result { error!(target: "cli", "{}", err); - std::process::exit(1); + std::process::exit(error_exit_status(&err)); + } +} + +/// DD102: only diagnosed configuration errors suppress service recovery. +fn error_exit_status(err: &anyhow::Error) -> i32 { + if err.chain().any(|cause| { + cause.is::<config::parser::ParserErr>() + || cause.is::<config::ValueErr>() + || cause.is::<config::PathsubErr>() + }) { + 6 + } else { + 1 + } +} + +#[cfg(test)] +mod exit_status_tests { + use super::*; + + #[test] + fn configuration_errors_keep_their_status_through_context() { + let cfg = Config::from_mem("[test]\nport = invalid\n").unwrap(); + let err = cfg + .section("test") + .number::<u16>("port") + .require() + .unwrap_err(); + assert_eq!( + 6, + error_exit_status(&anyhow::Error::new(err).context("starting server")) + ); + let err = cfg.section("test").str("missing").require().unwrap_err(); + assert_eq!(6, error_exit_status(&anyhow::Error::new(err))); + let err = Config::from_mem("not a configuration entry").unwrap_err(); + assert_eq!(6, error_exit_status(&anyhow::Error::new(err))); + } + + #[test] + fn unavailable_dependencies_remain_restartable() { + let err = std::io::Error::from(std::io::ErrorKind::ConnectionRefused); + assert_eq!( + 1, + error_exit_status(&anyhow::Error::new(err).context("database")) + ); } } diff --git a/debian/rules b/debian/rules @@ -32,4 +32,4 @@ override_dh_installsystemd: dh_installsystemd --no-enable --no-start --no-stop-on-upgrade --name taler-apns-relay-httpd dh_installsystemd --no-enable --no-start --no-stop-on-upgrade --name taler-apns-relay-worker dh_installsystemd --no-enable --no-start --no-stop-on-upgrade --name taler-apns-relay - dh_installsystemd + dh_installsystemd --no-start --no-enable --no-stop-on-upgrade diff --git a/debian/taler-apns-relay.prerm b/debian/taler-apns-relay.prerm @@ -0,0 +1,15 @@ +#!/bin/sh + +set -e + +# --no-start also suppresses debhelper's stop-on-remove snippet on Trixie. +# Stop the product explicitly on removal, while leaving upgrades passive. +if [ -d /run/systemd/system ] && [ "$1" = remove ]; +then + deb-systemd-invoke stop 'taler-apns-relay.target' >/dev/null || true + deb-systemd-invoke stop 'taler-apns-relay.slice' >/dev/null || true +fi + +#DEBHELPER# + +exit 0 diff --git a/debian/taler-apns-relay.taler-apns-relay-httpd.service b/debian/taler-apns-relay.taler-apns-relay-httpd.service @@ -1,20 +1,20 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler APNs relay REST API Requires=taler-apns-relay-httpd.socket -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-apns-relay.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-apns-relay-httpd Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-apns-relay serve -c /etc/taler-apns-relay/taler-apns-relay.conf diff --git a/debian/taler-apns-relay.taler-apns-relay-httpd.socket b/debian/taler-apns-relay.taler-apns-relay-httpd.socket @@ -1,6 +1,6 @@ [Unit] +BindsTo=taler-apns-relay-httpd.service Description=GNU Taler APNs relay socket -PartOf=taler-apns-relay-httpd.service [Socket] ListenStream=/run/taler-apns-relay/httpd/apns-relay-http.sock diff --git a/debian/taler-apns-relay.taler-apns-relay-worker.service b/debian/taler-apns-relay.taler-apns-relay-worker.service @@ -1,19 +1,19 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler APNs relay worker -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-apns-relay.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-apns-relay-worker Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-apns-relay worker -c /etc/taler-apns-relay/taler-apns-relay.conf diff --git a/debian/taler-apns-relay.taler-apns-relay.target b/debian/taler-apns-relay.taler-apns-relay.target @@ -1,6 +1,7 @@ [Unit] +Wants=postgresql.service Description=GNU Taler APNs relay -After=postgres.service network.target +After=postgresql.service network.target Wants=taler-apns-relay-httpd.service Wants=taler-apns-relay-worker.service diff --git a/debian/taler-cyclos.prerm b/debian/taler-cyclos.prerm @@ -0,0 +1,15 @@ +#!/bin/sh + +set -e + +# --no-start also suppresses debhelper's stop-on-remove snippet on Trixie. +# Stop the product explicitly on removal, while leaving upgrades passive. +if [ -d /run/systemd/system ] && [ "$1" = remove ]; +then + deb-systemd-invoke stop 'taler-cyclos.target' >/dev/null || true + deb-systemd-invoke stop 'taler-cyclos.slice' >/dev/null || true +fi + +#DEBHELPER# + +exit 0 diff --git a/debian/taler-cyclos.taler-cyclos-httpd.service b/debian/taler-cyclos.taler-cyclos-httpd.service @@ -1,20 +1,20 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Cyclos adapter REST API Requires=taler-cyclos-httpd.socket -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-cyclos.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-cyclos-httpd Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-cyclos serve -c /etc/taler-cyclos/taler-cyclos.conf ExecCondition=/usr/bin/taler-cyclos serve -c /etc/taler-cyclos/taler-cyclos.conf --check diff --git a/debian/taler-cyclos.taler-cyclos-httpd.socket b/debian/taler-cyclos.taler-cyclos-httpd.socket @@ -1,6 +1,6 @@ [Unit] +BindsTo=taler-cyclos-httpd.service Description=GNU Taler Cyclos adapter socket -PartOf=taler-cyclos-httpd.service [Socket] ListenStream=/run/taler-cyclos/httpd/cyclos-http.sock diff --git a/debian/taler-cyclos.taler-cyclos-worker.service b/debian/taler-cyclos.taler-cyclos-worker.service @@ -1,19 +1,19 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Cyclos adapter worker -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-cyclos.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-cyclos-worker Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-cyclos worker -c /etc/taler-cyclos/taler-cyclos.conf diff --git a/debian/taler-cyclos.taler-cyclos.target b/debian/taler-cyclos.taler-cyclos.target @@ -1,6 +1,7 @@ [Unit] +Wants=postgresql.service Description=GNU Taler Cyclos adapter -After=postgres.service network.target +After=postgresql.service network.target Wants=taler-cyclos-httpd.service Wants=taler-cyclos-worker.service diff --git a/debian/taler-magnet-bank.prerm b/debian/taler-magnet-bank.prerm @@ -0,0 +1,15 @@ +#!/bin/sh + +set -e + +# --no-start also suppresses debhelper's stop-on-remove snippet on Trixie. +# Stop the product explicitly on removal, while leaving upgrades passive. +if [ -d /run/systemd/system ] && [ "$1" = remove ]; +then + deb-systemd-invoke stop 'taler-magnet-bank.target' >/dev/null || true + deb-systemd-invoke stop 'taler-magnet-bank.slice' >/dev/null || true +fi + +#DEBHELPER# + +exit 0 diff --git a/debian/taler-magnet-bank.taler-magnet-bank-httpd.service b/debian/taler-magnet-bank.taler-magnet-bank-httpd.service @@ -1,20 +1,20 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Magnet Bank adapter REST API Requires=taler-magnet-bank-httpd.socket -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-magnet-bank.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-magnet-bank-httpd Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-magnet-bank serve -c /etc/taler-magnet-bank/taler-magnet-bank.conf ExecCondition=/usr/bin/taler-magnet-bank serve -c /etc/taler-magnet-bank/taler-magnet-bank.conf --check diff --git a/debian/taler-magnet-bank.taler-magnet-bank-httpd.socket b/debian/taler-magnet-bank.taler-magnet-bank-httpd.socket @@ -1,6 +1,6 @@ [Unit] +BindsTo=taler-magnet-bank-httpd.service Description=GNU Taler Magnet Bank adapter socket -PartOf=taler-magnet-bank-httpd.service [Socket] ListenStream=/run/taler-magnet-bank/httpd/magnet-bank-http.sock diff --git a/debian/taler-magnet-bank.taler-magnet-bank-worker.service b/debian/taler-magnet-bank.taler-magnet-bank-worker.service @@ -1,19 +1,19 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Magnet Bank adapter worker -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-magnet-bank.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-magnet-bank-worker Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-magnet-bank worker -c /etc/taler-magnet-bank/taler-magnet-bank.conf diff --git a/debian/taler-magnet-bank.taler-magnet-bank.target b/debian/taler-magnet-bank.taler-magnet-bank.target @@ -1,6 +1,7 @@ [Unit] +Wants=postgresql.service Description=GNU Taler Magnet Bank adapter -After=postgres.service network.target +After=postgresql.service network.target Wants=taler-magnet-bank-httpd.service Wants=taler-magnet-bank-worker.service diff --git a/debian/taler-wise.prerm b/debian/taler-wise.prerm @@ -0,0 +1,15 @@ +#!/bin/sh + +set -e + +# --no-start also suppresses debhelper's stop-on-remove snippet on Trixie. +# Stop the product explicitly on removal, while leaving upgrades passive. +if [ -d /run/systemd/system ] && [ "$1" = remove ]; +then + deb-systemd-invoke stop 'taler-wise.target' >/dev/null || true + deb-systemd-invoke stop 'taler-wise.slice' >/dev/null || true +fi + +#DEBHELPER# + +exit 0 diff --git a/debian/taler-wise.taler-wise-httpd.service b/debian/taler-wise.taler-wise-httpd.service @@ -1,20 +1,20 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Wise adapter REST API Requires=taler-wise-httpd.socket -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-wise.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-wise-httpd Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-wise serve -c /etc/taler-wise/taler-wise.conf ExecCondition=/usr/bin/taler-wise serve -c /etc/taler-wise/taler-wise.conf --check diff --git a/debian/taler-wise.taler-wise-httpd.socket b/debian/taler-wise.taler-wise-httpd.socket @@ -1,6 +1,6 @@ [Unit] +BindsTo=taler-wise-httpd.service Description=GNU Taler Wise adapter socket -PartOf=taler-wise-httpd.service [Socket] ListenStream=/run/taler-wise/httpd/wise-http.sock diff --git a/debian/taler-wise.taler-wise-worker.service b/debian/taler-wise.taler-wise-worker.service @@ -1,19 +1,19 @@ [Unit] +StartLimitIntervalSec=0 Description=GNU Taler Wise adapter worker -After=network.target postgres.service +After=network.target postgresql.service PartOf=taler-wise.target [Service] +# DD102: retry transient failures without a start-rate limit. +Restart=always +RestartSec=10s User=taler-wise-worker Type=exec -Restart=always RestartMode=direct -RestartSec=1ms -RestartPreventExitStatus=9 +RestartPreventExitStatus=6 9 -StartLimitBurst=5 -StartLimitInterval=5s ExecStart=/usr/bin/taler-wise worker -c /etc/taler-wise/taler-wise.conf diff --git a/debian/taler-wise.taler-wise.target b/debian/taler-wise.taler-wise.target @@ -1,6 +1,7 @@ [Unit] +Wants=postgresql.service Description=GNU Taler Wise adapter -After=postgres.service network.target +After=postgresql.service network.target Wants=taler-wise-httpd.service Wants=taler-wise-worker.service diff --git a/taler-apns-relay/src/apns.rs b/taler-apns-relay/src/apns.rs @@ -16,7 +16,6 @@ use std::time::Duration; -use anyhow::{anyhow, bail}; use aws_lc_rs::{ rand::SystemRandom, signature::{ECDSA_P256_SHA256_FIXED_SIGNING, EcdsaKeyPair}, @@ -193,17 +192,30 @@ impl Client { .install_default() .expect("failed to install the default TLS provider"); + let invalid_key = |message: String| taler_common::config::ValueErr::Invalid { + ty: "PKCS#8 private key file".into(), + section: "apns-relay-worker".into(), + option: "KEY_FILE".into(), + err: message, + }; // Load the signature key pair let private_key_der = PrivateKeyDer::from_pem_file(key_path) - .map_err(|e| anyhow!("failed to read key file at '{key_path}': {e}"))?; + .map_err(|e| invalid_key(format!("failed to read key file at '{key_path}': {e}")))?; let PrivateKeyDer::Pkcs8(pkcs8_der) = private_key_der else { - bail!("invalid key file at '{key_path}': not a valid PKCS#8 private key"); + return Err(invalid_key(format!( + "invalid key file at '{key_path}': not a valid PKCS#8 private key" + )) + .into()); }; let key_pair = EcdsaKeyPair::from_pkcs8( &ECDSA_P256_SHA256_FIXED_SIGNING, pkcs8_der.secret_pkcs8_der(), ) - .map_err(|_| anyhow!("invalid key file at '{key_path}': not a valid PKCS#8 private key"))?; + .map_err(|_| { + invalid_key(format!( + "invalid key file at '{key_path}': not a valid PKCS#8 private key" + )) + })?; // Make a signature let now = Timestamp::now(); diff --git a/taler-apns-relay/src/worker.rs b/taler-apns-relay/src/worker.rs @@ -88,7 +88,7 @@ async fn wakeup(pool: &PgPool, client: &mut Client) -> Result<(), WorkerError> { | Reason::BadEnvironmentKeyIdInToken | Reason::PayloadTooLarge => { error!(target: "worker", "config error, check the configuration: {e}"); - std::process::exit(9); + std::process::exit(6); } // Unregister Reason::DeviceTokenNotForTopic @@ -132,7 +132,6 @@ async fn wakeup(pool: &PgPool, client: &mut Client) -> Result<(), WorkerError> { pub async fn run(cfg: &Config, pool: &PgPool, transient: bool) -> anyhow::Result<()> { let cfg = WorkerCfg::parse(cfg)?; let mut client = Client::new(&cfg.apns)?; - let mut jitter = ExpoBackoffDecorr::default(); if transient { wakeup(pool, &mut client).await?; @@ -140,10 +139,18 @@ pub async fn run(cfg: &Config, pool: &PgPool, transient: bool) -> anyhow::Result } info!(target: "worker", "running at initialisation"); + let mut jitter = ExpoBackoffDecorr::default(); loop { - while let Err(WorkerError::Db(e)) = wakeup(pool, &mut client).await { - error!(target: "worker", "{e}"); - tokio::time::sleep(jitter.backoff()).await; + match wakeup(pool, &mut client).await { + Ok(()) => jitter.reset(), + Err(WorkerError::Db(err)) if !matches!(err, sqlx::Error::PoolClosed) => { + let delay = jitter.backoff(); + error!(target: "worker", ?delay, "database operation failed: {err}"); + tokio::time::sleep(delay).await; + continue; + } + // APNs errors reach here only after their bounded operation retries. + Err(err) => return Err(err.into()), } // TODO take sending time into account