server.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346
  1. use std::net::SocketAddr;
  2. use std::path::PathBuf;
  3. use std::pin::Pin;
  4. use std::str::FromStr;
  5. use std::sync::Arc;
  6. use std::time::Duration;
  7. use cdk_common::payment::MintPayment;
  8. use futures::{Stream, StreamExt};
  9. use serde_json::Value;
  10. use tokio::sync::{mpsc, Notify};
  11. use tokio::task::JoinHandle;
  12. use tokio::time::{sleep, Instant};
  13. use tokio_stream::wrappers::ReceiverStream;
  14. use tonic::transport::{Certificate, Identity, Server, ServerTlsConfig};
  15. use tonic::{async_trait, Request, Response, Status};
  16. use tracing::instrument;
  17. use super::cdk_payment_processor_server::{CdkPaymentProcessor, CdkPaymentProcessorServer};
  18. use crate::proto::*;
  19. type ResponseStream =
  20. Pin<Box<dyn Stream<Item = Result<WaitIncomingPaymentResponse, Status>> + Send>>;
  21. /// Payment Processor
  22. #[derive(Clone)]
  23. pub struct PaymentProcessorServer {
  24. inner: Arc<dyn MintPayment<Err = cdk_common::payment::Error> + Send + Sync>,
  25. socket_addr: SocketAddr,
  26. shutdown: Arc<Notify>,
  27. handle: Option<Arc<JoinHandle<anyhow::Result<()>>>>,
  28. }
  29. impl PaymentProcessorServer {
  30. /// Create new [`PaymentProcessorServer`]
  31. pub fn new(
  32. payment_processor: Arc<dyn MintPayment<Err = cdk_common::payment::Error> + Send + Sync>,
  33. addr: &str,
  34. port: u16,
  35. ) -> anyhow::Result<Self> {
  36. let socket_addr = SocketAddr::new(addr.parse()?, port);
  37. Ok(Self {
  38. inner: payment_processor,
  39. socket_addr,
  40. shutdown: Arc::new(Notify::new()),
  41. handle: None,
  42. })
  43. }
  44. /// Start fake wallet grpc server
  45. pub async fn start(&mut self, tls_dir: Option<PathBuf>) -> anyhow::Result<()> {
  46. tracing::info!("Starting RPC server {}", self.socket_addr);
  47. let server = match tls_dir {
  48. Some(tls_dir) => {
  49. tracing::info!("TLS configuration found, starting secure server");
  50. // Check for server.pem
  51. let server_pem_path = tls_dir.join("server.pem");
  52. if !server_pem_path.exists() {
  53. let err_msg = format!(
  54. "TLS certificate file not found: {}",
  55. server_pem_path.display()
  56. );
  57. tracing::error!("{}", err_msg);
  58. return Err(anyhow::anyhow!(err_msg));
  59. }
  60. // Check for server.key
  61. let server_key_path = tls_dir.join("server.key");
  62. if !server_key_path.exists() {
  63. let err_msg = format!("TLS key file not found: {}", server_key_path.display());
  64. tracing::error!("{}", err_msg);
  65. return Err(anyhow::anyhow!(err_msg));
  66. }
  67. // Check for ca.pem
  68. let ca_pem_path = tls_dir.join("ca.pem");
  69. if !ca_pem_path.exists() {
  70. let err_msg =
  71. format!("CA certificate file not found: {}", ca_pem_path.display());
  72. tracing::error!("{}", err_msg);
  73. return Err(anyhow::anyhow!(err_msg));
  74. }
  75. let cert = std::fs::read_to_string(&server_pem_path)?;
  76. let key = std::fs::read_to_string(&server_key_path)?;
  77. let client_ca_cert = std::fs::read_to_string(&ca_pem_path)?;
  78. let client_ca_cert = Certificate::from_pem(client_ca_cert);
  79. let server_identity = Identity::from_pem(cert, key);
  80. let tls_config = ServerTlsConfig::new()
  81. .identity(server_identity)
  82. .client_ca_root(client_ca_cert);
  83. Server::builder()
  84. .tls_config(tls_config)?
  85. .add_service(CdkPaymentProcessorServer::new(self.clone()))
  86. }
  87. None => {
  88. tracing::warn!("No valid TLS configuration found, starting insecure server");
  89. Server::builder().add_service(CdkPaymentProcessorServer::new(self.clone()))
  90. }
  91. };
  92. let shutdown = self.shutdown.clone();
  93. let addr = self.socket_addr;
  94. self.handle = Some(Arc::new(tokio::spawn(async move {
  95. let server = server.serve_with_shutdown(addr, async {
  96. shutdown.notified().await;
  97. });
  98. server.await?;
  99. Ok(())
  100. })));
  101. Ok(())
  102. }
  103. /// Stop fake wallet grpc server
  104. pub async fn stop(&self) -> anyhow::Result<()> {
  105. const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
  106. if let Some(handle) = &self.handle {
  107. tracing::info!("Initiating server shutdown");
  108. self.shutdown.notify_waiters();
  109. let start = Instant::now();
  110. while !handle.is_finished() {
  111. if start.elapsed() >= SHUTDOWN_TIMEOUT {
  112. tracing::error!(
  113. "Server shutdown timed out after {} seconds, aborting handle",
  114. SHUTDOWN_TIMEOUT.as_secs()
  115. );
  116. handle.abort();
  117. break;
  118. }
  119. sleep(Duration::from_millis(100)).await;
  120. }
  121. if handle.is_finished() {
  122. tracing::info!("Server shutdown completed successfully");
  123. }
  124. } else {
  125. tracing::info!("No server handle found, nothing to stop");
  126. }
  127. Ok(())
  128. }
  129. }
  130. impl Drop for PaymentProcessorServer {
  131. fn drop(&mut self) {
  132. tracing::debug!("Dropping payment process server");
  133. self.shutdown.notify_one();
  134. }
  135. }
  136. #[async_trait]
  137. impl CdkPaymentProcessor for PaymentProcessorServer {
  138. async fn get_settings(
  139. &self,
  140. _request: Request<SettingsRequest>,
  141. ) -> Result<Response<SettingsResponse>, Status> {
  142. let settings: Value = self
  143. .inner
  144. .get_settings()
  145. .await
  146. .map_err(|_| Status::internal("Could not get settings"))?;
  147. Ok(Response::new(SettingsResponse {
  148. inner: settings.to_string(),
  149. }))
  150. }
  151. async fn create_payment(
  152. &self,
  153. request: Request<CreatePaymentRequest>,
  154. ) -> Result<Response<CreatePaymentResponse>, Status> {
  155. let CreatePaymentRequest {
  156. amount,
  157. unit,
  158. description,
  159. unix_expiry,
  160. } = request.into_inner();
  161. let unit =
  162. CurrencyUnit::from_str(&unit).map_err(|_| Status::invalid_argument("Invalid unit"))?;
  163. let invoice_response = self
  164. .inner
  165. .create_incoming_payment_request(amount.into(), &unit, description, unix_expiry)
  166. .await
  167. .map_err(|_| Status::internal("Could not create invoice"))?;
  168. Ok(Response::new(invoice_response.into()))
  169. }
  170. async fn get_payment_quote(
  171. &self,
  172. request: Request<PaymentQuoteRequest>,
  173. ) -> Result<Response<PaymentQuoteResponse>, Status> {
  174. let request = request.into_inner();
  175. let options: Option<cdk_common::MeltOptions> =
  176. request.options.as_ref().map(|options| (*options).into());
  177. let payment_quote = self
  178. .inner
  179. .get_payment_quote(
  180. &request.request,
  181. &CurrencyUnit::from_str(&request.unit)
  182. .map_err(|_| Status::invalid_argument("Invalid currency unit"))?,
  183. options,
  184. )
  185. .await
  186. .map_err(|err| {
  187. tracing::error!("Could not get bolt11 melt quote: {}", err);
  188. Status::internal("Could not get melt quote")
  189. })?;
  190. Ok(Response::new(payment_quote.into()))
  191. }
  192. async fn make_payment(
  193. &self,
  194. request: Request<MakePaymentRequest>,
  195. ) -> Result<Response<MakePaymentResponse>, Status> {
  196. let request = request.into_inner();
  197. let pay_invoice = self
  198. .inner
  199. .make_payment(
  200. request
  201. .melt_quote
  202. .ok_or(Status::invalid_argument("Meltquote is required"))?
  203. .try_into()
  204. .map_err(|_err| Status::invalid_argument("Invalid melt quote"))?,
  205. request.partial_amount.map(|a| a.into()),
  206. request.max_fee_amount.map(|a| a.into()),
  207. )
  208. .await
  209. .map_err(|err| {
  210. tracing::error!("Could not make payment: {}", err);
  211. match err {
  212. cdk_common::payment::Error::InvoiceAlreadyPaid => {
  213. Status::already_exists("Payment request already paid")
  214. }
  215. cdk_common::payment::Error::InvoicePaymentPending => {
  216. Status::already_exists("Payment request pending")
  217. }
  218. _ => Status::internal("Could not pay invoice"),
  219. }
  220. })?;
  221. Ok(Response::new(pay_invoice.into()))
  222. }
  223. async fn check_incoming_payment(
  224. &self,
  225. request: Request<CheckIncomingPaymentRequest>,
  226. ) -> Result<Response<CheckIncomingPaymentResponse>, Status> {
  227. let request = request.into_inner();
  228. let check_response = self
  229. .inner
  230. .check_incoming_payment_status(&request.request_lookup_id)
  231. .await
  232. .map_err(|_| Status::internal("Could not check incoming payment status"))?;
  233. Ok(Response::new(CheckIncomingPaymentResponse {
  234. status: QuoteState::from(check_response).into(),
  235. }))
  236. }
  237. async fn check_outgoing_payment(
  238. &self,
  239. request: Request<CheckOutgoingPaymentRequest>,
  240. ) -> Result<Response<MakePaymentResponse>, Status> {
  241. let request = request.into_inner();
  242. let check_response = self
  243. .inner
  244. .check_outgoing_payment(&request.request_lookup_id)
  245. .await
  246. .map_err(|_| Status::internal("Could not check incoming payment status"))?;
  247. Ok(Response::new(check_response.into()))
  248. }
  249. type WaitIncomingPaymentStream = ResponseStream;
  250. // Clippy thinks select is not stable but it compiles fine on MSRV (1.63.0)
  251. #[allow(clippy::incompatible_msrv)]
  252. #[instrument(skip_all)]
  253. async fn wait_incoming_payment(
  254. &self,
  255. _request: Request<WaitIncomingPaymentRequest>,
  256. ) -> Result<Response<Self::WaitIncomingPaymentStream>, Status> {
  257. tracing::debug!("Server waiting for payment stream");
  258. let (tx, rx) = mpsc::channel(128);
  259. let shutdown_clone = self.shutdown.clone();
  260. let ln = self.inner.clone();
  261. tokio::spawn(async move {
  262. loop {
  263. tokio::select! {
  264. _ = shutdown_clone.notified() => {
  265. tracing::info!("Shutdown signal received, stopping task for ");
  266. ln.cancel_wait_invoice();
  267. break;
  268. }
  269. result = ln.wait_any_incoming_payment() => {
  270. match result {
  271. Ok(mut stream) => {
  272. while let Some(request_lookup_id) = stream.next().await {
  273. match tx.send(Result::<_, Status>::Ok(WaitIncomingPaymentResponse{lookup_id: request_lookup_id} )).await {
  274. Ok(_) => {
  275. // item (server response) was queued to be send to client
  276. }
  277. Err(item) => {
  278. tracing::error!("Error adding incoming payment to stream: {}", item);
  279. break;
  280. }
  281. }
  282. }
  283. }
  284. Err(err) => {
  285. tracing::warn!("Could not get invoice stream for {}", err);
  286. tokio::time::sleep(std::time::Duration::from_secs(5)).await;
  287. }
  288. }
  289. }
  290. }
  291. }
  292. });
  293. let output_stream = ReceiverStream::new(rx);
  294. Ok(Response::new(
  295. Box::pin(output_stream) as Self::WaitIncomingPaymentStream
  296. ))
  297. }
  298. }