Make IPC concurrent so relay checks and feeds never block the UI

This commit is contained in:
Avi 2026-08-23 19:03:16 -05:00
commit 4d70b25628
4 changed files with 136 additions and 89 deletions

View file

@ -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,
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?;
.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<W>(writer: &mut W, envelope: ReplyEnvelope) -> Result<(), AppError>
where
W: tokio::io::AsyncWriteExt + Unpin,
{
async fn write_line(
writer: &tokio::sync::Mutex<tokio::io::Stdout>,
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<Mutex<App>>,
signer: &Signer,
signer: Arc<Signer>,
request: Request,
) -> Reply<serde_json::Value> {
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<Mutex<App>>,
signer: &Signer,
request: Request,
) -> Result<serde_json::Value, AppError> {
// 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<serde_json::Value, AppError> {
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<serde_json::Val
Ok(json!(report))
}
Request::FeedGet {
limit,
contacts_only,
} => {
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<serde_json::Val
Ok(json!(app.settings))
}
Request::RelayTest { url } => {
let result = relays::test_connection(&url, RELAY_TEST_TIMEOUT).await?;
Ok(json!(result))
}
Request::SettingsGet => Ok(json!(app.settings)),
Request::BackupNow => {

View file

@ -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<String, AppError> {
let uri = args
.get(3)
.ok_or_else(|| AppError::config("Usage: signer connect <uri>"))?;
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)");

View file

@ -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 {
// 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<Option<Result<String, RelayFailure>>> =
(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(_)) => succeeded.push(url.clone()),
Ok(Err(err)) => failed.push(failure_for(url, &err)),
Err(_) => failed.push(RelayFailure {
url: url.clone(),
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) => failed.push(failure_for(url, &err)),
Err(err) => (index, Err(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 }
}

View file

@ -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<Mutex<App>>, uri: &str) -> Result<(), AppError> {
pub fn connect(&self, app: Arc<tokio::sync::Mutex<App>>, 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<String, String> {
/// 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<Mutex<App>>, uri: ConnectUri) {
async fn run_sign_task(signer: Signer, app: Arc<tokio::sync::Mutex<App>>, 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());
}