From 4d70b25628f907448a18e96ba89df0dcdc402ce5 Mon Sep 17 00:00:00 2001 From: Avi Date: Sun, 23 Aug 2026 19:03:16 -0500 Subject: [PATCH] Make IPC concurrent so relay checks and feeds never block the UI --- src/ipc.rs | 138 ++++++++++++++++++++++++++++-------------------- src/main.rs | 4 +- src/profiles.rs | 66 ++++++++++++++++------- src/signer.rs | 17 +++--- 4 files changed, 136 insertions(+), 89 deletions(-) diff --git a/src/ipc.rs b/src/ipc.rs index 89c75ea..c0c1b04 100644 --- a/src/ipc.rs +++ b/src/ipc.rs @@ -1,8 +1,9 @@ -use std::sync::{Arc, Mutex}; +use std::sync::Arc; use std::time::Duration; use serde::{Deserialize, Serialize}; use serde_json::json; +use tokio::sync::Mutex; use crate::app::App; use crate::errors::AppError; @@ -157,20 +158,26 @@ pub struct ReplyEnvelope { /// Run the JSON-lines IPC server on stdin/stdout. /// -/// The Electron main process spawns `keynectr serve` and -/// exchanges one JSON object per line. Requests are processed sequentially so -/// the shared state never sees concurrent mutations. +/// The Electron main process spawns `keynectr serve` and exchanges one JSON +/// object per line. Requests are handled concurrently — replies are +/// id-correlated, so they may arrive out of order — while every handler that +/// touches shared state locks it for the duration, so mutations remain +/// serialized and never interleave. Long network-only requests (relay tests, +/// feed reads) run without holding that lock so they cannot delay interactive +/// ones such as selecting a profile. pub async fn serve() -> Result<(), AppError> { use tokio::io::AsyncBufReadExt; + use tokio::task::JoinSet; // Shared state, so the NIP-46 signer's background task and the request loop // both see the same vault (including its unlock key) without racing writes. let app = Arc::new(Mutex::new(App::load()?)); - let signer = Signer::new(); + let signer = Arc::new(Signer::new()); + let stdout = Arc::new(tokio::sync::Mutex::new(tokio::io::stdout())); let stdin = tokio::io::stdin(); let mut lines = tokio::io::BufReader::new(stdin).lines(); - let mut stdout = tokio::io::stdout(); + let mut tasks = JoinSet::new(); while let Some(line) = lines .next_line() @@ -189,29 +196,40 @@ pub async fn serve() -> Result<(), AppError> { message: "The request could not be understood.".to_string(), details: Some(e.to_string()), }; - write_line(&mut stdout, ReplyEnvelope { id: 0, reply }).await?; + write_line(&stdout, ReplyEnvelope { id: 0, reply }).await?; continue; } }; - let reply = handle(app.clone(), &signer, envelope.request).await; - write_line( - &mut stdout, - ReplyEnvelope { - id: envelope.id, - reply, - }, - ) - .await?; + let task_app = app.clone(); + let task_signer = signer.clone(); + let task_stdout = stdout.clone(); + tasks.spawn(async move { + let reply = handle(task_app, task_signer, envelope.request).await; + let _ = write_line( + &task_stdout, + ReplyEnvelope { + id: envelope.id, + reply, + }, + ) + .await; + }); } + // Finish in-flight requests before returning so the GUI never sees the + // backend disappear mid-request. + while tasks.join_next().await.is_some() {} Ok(()) } -async fn write_line(writer: &mut W, envelope: ReplyEnvelope) -> Result<(), AppError> -where - W: tokio::io::AsyncWriteExt + Unpin, -{ +async fn write_line( + writer: &tokio::sync::Mutex, + envelope: ReplyEnvelope, +) -> Result<(), AppError> { + use tokio::io::AsyncWriteExt; + + let mut writer = writer.lock().await; let mut line = serde_json::to_string(&envelope) .map_err(|e| AppError::json("Could not prepare a response", e))?; line.push('\n'); @@ -228,10 +246,10 @@ where async fn handle( app: Arc>, - signer: &Signer, + signer: Arc, request: Request, ) -> Reply { - let result = run(&app, signer, request).await; + let result = run(&app, &signer, request).await; match result { Ok(value) => Reply::Ok { data: value }, Err(err) => Reply::Error { @@ -255,19 +273,16 @@ fn error_code(err: &AppError) -> String { } /// Signer control commands never touch the vault directly, so they take the -/// shared handle (a clone) rather than locking the state. Everything else -/// locks the vault for the duration of the call, mirroring the old -/// single-threaded model. -#[allow(clippy::await_holding_lock)] +/// shared handle (a clone) rather than locking the state. Read-only network +/// requests (relay tests, feed reads) grab what they need under a short lock +/// and then run without it, so slow relays cannot delay interactive requests. +/// Everything else locks the state for the duration of the call, so mutations +/// remain serialized and never interleave. async fn run( app: &Arc>, signer: &Signer, request: Request, ) -> Result { - // Signer control commands never touch the vault directly, so they take the - // shared handle (a clone) rather than locking the state. Everything else - // locks the vault for the duration of the call, mirroring the old - // single-threaded model. match request { Request::SignerConnect { uri } => { signer.connect(app.clone(), &uri)?; @@ -282,16 +297,49 @@ async fn run( signer.approve(&id, approved)?; Ok(json!(signer.status())) } + Request::RelayTest { url } => { + // Pure network probe against the given URL; no shared state. + let result = relays::test_connection(&url, RELAY_TEST_TIMEOUT).await?; + Ok(json!(result)) + } + Request::FeedGet { + limit, + contacts_only, + } => { + let limit = limit.unwrap_or(feed::DEFAULT_LIMIT); + let contacts_only = contacts_only.unwrap_or(false); + // Copy the inputs out of shared state under a short lock so the + // multi-second relay fetches below never block a Select or save. + let (settings, owner_hex) = { + let guard = app.lock().await; + let owner_hex = if contacts_only { + let npub = guard + .vault + .active_profile + .as_deref() + .ok_or_else(AppError::no_active_profile)?; + Some(feed::owner_pubkey(npub)?.to_hex()) + } else { + None + }; + (guard.settings.clone(), owner_hex) + }; + let items = match owner_hex { + Some(hex) => feed::contact_feed(&settings, limit, &hex).await?, + None => feed::aggregate_feed(&settings, limit).await?, + }; + Ok(json!(items)) + } other => { - let mut guard = app.lock().expect("app mutex poisoned"); + let mut guard = app.lock().await; run_with_app(&mut guard, other).await } } } /// Requests dispatched to the vault state. The shared mutex guard is held across -/// the awaited operation on purpose: requests remain effectively sequential, and -/// a concurrent `await` never yields back into a state the loop expects to own. +/// the awaited operation on purpose: mutations stay serialized against each +/// other even though requests themselves are handled concurrently. async fn run_with_app(app: &mut App, request: Request) -> Result { match request { Request::Init | Request::GetState => Ok(json!(app.state_view())), @@ -351,25 +399,6 @@ async fn run_with_app(app: &mut App, request: Request) -> Result { - let limit = limit.unwrap_or(feed::DEFAULT_LIMIT); - let items = if contacts_only.unwrap_or(false) { - let npub = app - .vault - .active_profile - .as_deref() - .ok_or_else(AppError::no_active_profile)?; - let pubkey = feed::owner_pubkey(npub)?; - feed::contact_feed(&app.settings, limit, &pubkey.to_hex()).await? - } else { - feed::aggregate_feed(&app.settings, limit).await? - }; - Ok(json!(items)) - } - Request::RelayAdd { url } => { relays::add_relay(&mut app.settings, &url)?; app.save_settings()?; @@ -388,11 +417,6 @@ async fn run_with_app(app: &mut App, request: Request) -> Result { - let result = relays::test_connection(&url, RELAY_TEST_TIMEOUT).await?; - Ok(json!(result)) - } - Request::SettingsGet => Ok(json!(app.settings)), Request::BackupNow => { diff --git a/src/main.rs b/src/main.rs index 725c999..56a53e9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,5 +1,5 @@ use std::process::ExitCode; -use std::sync::{Arc, Mutex}; +use std::sync::Arc; use keynectr::app::App; use keynectr::errors::{AppError, ErrorKind}; @@ -509,7 +509,7 @@ async fn cli_signer(args: &[String]) -> Result { let uri = args .get(3) .ok_or_else(|| AppError::config("Usage: signer connect "))?; - let app = Arc::new(Mutex::new(load_app_with_unlock()?)); + let app = Arc::new(tokio::sync::Mutex::new(load_app_with_unlock()?)); let signer = Signer::new(); signer.connect(app, uri)?; println!("Connecting to the NIP-46 app… (interrupt with Ctrl-C to stop)"); diff --git a/src/profiles.rs b/src/profiles.rs index 3d42a3f..5f542d4 100644 --- a/src/profiles.rs +++ b/src/profiles.rs @@ -11,9 +11,10 @@ use crate::settings::Settings; use crate::vault::{unix_timestamp, StoredProfile, Vault}; /// How long to wait for relays to accept a connection attempt. -const METADATA_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); -/// How long to wait for a single relay to accept the metadata event. -const METADATA_SEND_TIMEOUT: Duration = Duration::from_secs(15); +const METADATA_CONNECT_TIMEOUT: Duration = Duration::from_secs(5); +/// How long to wait for a single relay to accept the metadata event. Relays +/// are sent to in parallel, so this caps the whole publish, not each relay. +const METADATA_SEND_TIMEOUT: Duration = Duration::from_secs(8); /// A safe view of a profile that contains no secret key material. #[derive(Debug, Clone, Serialize, PartialEq, Eq)] @@ -343,29 +344,54 @@ async fn publish_metadata_async( client.connect().await; let _ = client.wait_for_connection(METADATA_CONNECT_TIMEOUT).await; - let mut succeeded = Vec::new(); - let mut failed = Vec::new(); - for url in &relay_urls { - match client.relay(url.as_str()).await { - Ok(relay) => { - match tokio::time::timeout(METADATA_SEND_TIMEOUT, relay.send_event(&event)).await { - Ok(Ok(_)) => succeeded.push(url.clone()), - Ok(Err(err)) => failed.push(failure_for(url, &err)), - Err(_) => failed.push(RelayFailure { - url: url.clone(), - error: "The relay did not respond in time.".to_string(), - details: Some( - "Timed out while waiting for the relay to accept the metadata." - .to_string(), + // Send to every relay in parallel so one slow or dead relay cannot drag + // the whole publish out to (relays x timeout). + let mut outcomes: Vec>> = + (0..relay_urls.len()).map(|_| None).collect(); + let mut sends = tokio::task::JoinSet::new(); + for (index, url) in relay_urls.iter().cloned().enumerate() { + let client = client.clone(); + let event = event.clone(); + sends.spawn(async move { + match client.relay(url.as_str()).await { + Ok(relay) => { + match tokio::time::timeout(METADATA_SEND_TIMEOUT, relay.send_event(&event)) + .await + { + Ok(Ok(_)) => (index, Ok(url)), + Ok(Err(err)) => (index, Err(failure_for(&url, &err))), + Err(_) => ( + index, + Err(RelayFailure { + url, + error: "The relay did not respond in time.".to_string(), + details: Some( + "Timed out while waiting for the relay to accept the metadata." + .to_string(), + ), + }), ), - }), + } } + Err(err) => (index, Err(failure_for(&url, &err))), } - Err(err) => failed.push(failure_for(url, &err)), - } + }); + } + while let Some(joined) = sends.join_next().await { + let (index, outcome) = joined.expect("metadata send task panicked"); + outcomes[index] = Some(outcome); } client.disconnect().await; + let mut succeeded = Vec::new(); + let mut failed = Vec::new(); + for outcome in outcomes.into_iter().flatten() { + match outcome { + Ok(url) => succeeded.push(url), + Err(failure) => failed.push(failure), + } + } + MetadataPublishReport { succeeded, failed } } diff --git a/src/signer.rs b/src/signer.rs index 086f0b3..84e8d7a 100644 --- a/src/signer.rs +++ b/src/signer.rs @@ -238,7 +238,7 @@ impl Signer { /// /// `app` is the shared vault state, so the signer reflects the current unlock /// key and active profile. A locked vault cannot sign until it is unlocked. - pub fn connect(&self, app: Arc>, uri: &str) -> Result<(), AppError> { + pub fn connect(&self, app: Arc>, uri: &str) -> Result<(), AppError> { if uri.trim().starts_with("bunker://") { return Err(AppError::config( "That is a bunker:// link, which means routing through *another* signer. \ @@ -596,16 +596,10 @@ fn nip44(keys: &Keys, request: &RawRequest) -> Result { /// The background loop: connect to the client's relays, announce ourselves, /// subscribe to kind 24133 events, and answer requests until stopped. -async fn run_sign_task(signer: Signer, app: Arc>, uri: ConnectUri) { +async fn run_sign_task(signer: Signer, app: Arc>, uri: ConnectUri) { // 1. Resolve the active profile's key under the current vault lock. let keys = { - let guard = match app.lock() { - Ok(guard) => guard, - Err(_) => { - signer.fail("The vault could not be read."); - return; - } - }; + let guard = app.lock().await; let hex = match profiles::resolve_active_secret_key(&guard.vault, guard.vault_key()) { Ok(hex) => hex, Err(err) => { @@ -764,7 +758,10 @@ mod tests { fn bunker_link_is_rejected() { let signer = Signer::new(); assert!(signer - .connect(Arc::new(Mutex::new(App::load().unwrap())), "bunker://abc") + .connect( + Arc::new(tokio::sync::Mutex::new(App::load().unwrap())), + "bunker://abc" + ) .is_err()); }