diff --git a/Cargo.toml b/Cargo.toml index 896259e..65dc3ef 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,11 @@ members = ["appa-cosmic"] default-members = ["."] resolver = "3" +# Keep application changes quick to compile while rendering and networking +# dependencies run at a usable speed during `cargo run` development sessions. +[profile.dev.package."*"] +opt-level = 2 + [dependencies] anyhow = "1.0.104" arboard = "3.6.1" @@ -31,7 +36,7 @@ serde_json = "1.0.151" tempfile = "3.27.0" time = { version = "0.3.55", features = ["formatting", "macros", "parsing", "serde"] } toml = "1.1.4" -tokio = { version = "1.53.1", features = ["fs", "macros", "rt-multi-thread", "signal", "sync", "time"] } +tokio = { version = "1.53.1", features = ["fs", "io-util", "macros", "net", "rt-multi-thread", "signal", "sync", "time"] } tracing = "0.1.44" tracing-subscriber = { version = "0.3.23", features = ["env-filter", "fmt"] } uuid = { version = "1.24.0", features = ["serde", "v4"] } diff --git a/README.md b/README.md index 2620082..699e291 100644 --- a/README.md +++ b/README.md @@ -195,6 +195,13 @@ With Nix: nix develop -c cargo run -p appa-cosmic ``` +The GUI talks to the long-lived local Appa daemon. Start it directly while +developing, or install the user service from the GUI: + +```bash +nix develop -c cargo run -- daemon serve +``` + Build the packaged application with `nix build --impure .#appa-cosmic`. ## Troubleshoot diff --git a/appa-cosmic/src/app.rs b/appa-cosmic/src/app.rs index 742d3ce..4ba6b09 100644 --- a/appa-cosmic/src/app.rs +++ b/appa-cosmic/src/app.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; use appa::core::{Dashboard, FolderDetails, FolderSummary}; use cosmic::{ - iced::{Alignment, Length, Subscription}, + iced::{Alignment, Length}, prelude::*, widget::{self, nav_bar}, }; @@ -194,10 +194,6 @@ impl cosmic::Application for AppModel { self.nav.activate(id); Task::none() } - - fn subscription(&self) -> Subscription { - cosmic::iced::time::every(std::time::Duration::from_secs(5)).map(|_| Message::Refresh) - } } impl AppModel { @@ -235,7 +231,12 @@ impl AppModel { content = content.push(self.folder_row(folder)); } } else { - content = content.push(widget::text::body("Loading folder status…")); + content = content + .push(widget::text::body("Appa daemon is unavailable.")) + .push( + widget::button::text("Install background sync service") + .on_press(Message::InstallService), + ); } content.into() } @@ -413,21 +414,23 @@ fn folder_health(folder: &FolderSummary) -> String { } fn load_dashboard() -> Task> { - Task::perform(async { load_dashboard_result() }, Message::DashboardLoaded) - .map(cosmic::Action::App) -} - -fn load_dashboard_result() -> Result { - appa::Appa::open() - .and_then(|appa| appa.dashboard()) - .map_err(|error| error.to_string()) + Task::perform( + async { + let client = appa::ipc::DaemonClient::connect().map_err(|error| error.to_string())?; + client.dashboard().await.map_err(|error| error.to_string()) + }, + Message::DashboardLoaded, + ) + .map(cosmic::Action::App) } fn load_folder_details(folder_path: PathBuf) -> Task> { Task::perform( async move { - appa::Appa::open() - .and_then(|appa| appa.folder_details(&folder_path)) + let client = appa::ipc::DaemonClient::connect().map_err(|error| error.to_string())?; + client + .folder_details(folder_path) + .await .map_err(|error| error.to_string()) }, Message::FolderDetailsLoaded, @@ -436,14 +439,18 @@ fn load_folder_details(folder_path: PathBuf) -> Task> { } fn register_folder(folder_path: String) -> Task> { - run_folder_action(folder_path, |appa, path| appa.register_folder(path)) + run_folder_action(folder_path, |client, path| { + Box::pin(async move { client.register_folder(path).await }) + }) } fn join_folder(folder_path: String, invitation: String) -> Task> { Task::perform( async move { - let appa = appa::Appa::open().map_err(|error| error.to_string())?; - appa.join_folder(PathBuf::from(folder_path).as_path(), &invitation) + let client = appa::ipc::DaemonClient::connect().map_err(|error| error.to_string())?; + client + .join_folder(PathBuf::from(folder_path), invitation) + .await .map_err(|error| error.to_string()) }, folder_action_completed, @@ -453,18 +460,26 @@ fn join_folder(folder_path: String, invitation: String) -> Task Task> { run_folder_action(folder_path.display().to_string(), move |appa, path| { - appa.revoke_member(path, &device_id) + Box::pin(async move { appa.revoke_member(path, device_id).await }) }) } fn run_folder_action( folder_path: String, - action: impl FnOnce(&appa::Appa, &std::path::Path) -> appa::core::Result<()> + Send + 'static, + action: impl FnOnce( + appa::ipc::DaemonClient, + PathBuf, + ) -> std::pin::Pin< + Box> + Send>, + > + Send + + 'static, ) -> Task> { Task::perform( async move { - let appa = appa::Appa::open().map_err(|error| error.to_string())?; - action(&appa, PathBuf::from(folder_path).as_path()).map_err(|error| error.to_string()) + let client = appa::ipc::DaemonClient::connect().map_err(|error| error.to_string())?; + action(client, PathBuf::from(folder_path)) + .await + .map_err(|error| error.to_string()) }, folder_action_completed, ) @@ -481,17 +496,11 @@ fn folder_action_completed(result: Result<(), String>) -> Message { fn create_invite(folder_path: PathBuf) -> Task> { Task::perform( async move { - tokio::task::spawn_blocking(move || { - let runtime = tokio::runtime::Runtime::new().map_err(|error| error.to_string())?; - runtime.block_on(async { - let appa = appa::Appa::open().map_err(|error| error.to_string())?; - appa.create_invite(&folder_path) - .await - .map_err(|error| error.to_string()) - }) - }) - .await - .map_err(|error| error.to_string())? + let client = appa::ipc::DaemonClient::connect().map_err(|error| error.to_string())?; + client + .create_invite(folder_path) + .await + .map_err(|error| error.to_string()) }, Message::InvitationCreated, ) @@ -501,8 +510,8 @@ fn create_invite(folder_path: PathBuf) -> Task> { fn install_service() -> Task> { Task::perform( async { - appa::Appa::open() - .and_then(|appa| appa.install_service()) + appa::install_service() + .map(|_| ()) .map_err(|error| error.to_string()) }, Message::ServiceInstalled, diff --git a/src/app.rs b/src/app.rs index 58aa1d1..d7d6b11 100644 --- a/src/app.rs +++ b/src/app.rs @@ -13,7 +13,7 @@ pub type AppResult = anyhow::Result; mod manifest; mod materialize; -mod run; +pub(crate) mod run; mod sync; pub(super) const MAX_CONCURRENT_BLOB_TRANSFERS: usize = 4; diff --git a/src/app/run.rs b/src/app/run.rs index f1a9f14..d9cb2f6 100644 --- a/src/app/run.rs +++ b/src/app/run.rs @@ -1,9 +1,13 @@ -use std::{collections::BTreeMap, path::Path, time::Duration}; +use std::{ + collections::BTreeMap, + path::{Path, PathBuf}, + time::Duration, +}; use crate::{ app::{AppResult, AppaService}, domain::FolderId, - filesystem::{FolderChangeWatcher, lock_app_process, watch_folder}, + filesystem::{FolderChangeWatcher, FolderLock, lock_app_process, watch_folder}, iroh::NodeHost, storage::FolderConfig, }; @@ -29,36 +33,74 @@ struct FolderWatcher { next_safety_scan: tokio::time::Instant, } -impl AppaService { - pub async fn run_forever(&self, folder_path: Option<&Path>) -> AppResult<()> { - let folders = self.folders_to_run(folder_path)?; - let _lock = lock_app_process(&self.paths.data_directory.join("locks"))?; - let node = self.load_node().await?; +/// Stateful synchronization runtime owned by the daemon process. +pub(crate) struct SyncRunner { + _lock: FolderLock, + node: NodeHost, + watchers: Vec, + selected_folder_path: Option, + next_folder_refresh: tokio::time::Instant, + retries: BTreeMap, +} + +impl SyncRunner { + pub(crate) async fn start( + service: &AppaService, + folder_path: Option<&Path>, + ) -> AppResult { + let folders = service.folders_to_run(folder_path)?; + let lock = lock_app_process(&service.paths.data_directory.join("locks"))?; + let node = service.load_node().await?; node.wait_until_online().await?; - let mut watchers = folders + let watchers = folders .into_iter() .map(FolderWatcher::new) .collect::>>()?; - let mut interval = tokio::time::interval(WATCHER_POLL_INTERVAL); - let mut next_folder_refresh = tokio::time::Instant::now() + FOLDER_REFRESH_INTERVAL; - let mut retries = BTreeMap::new(); tracing::info!(folder_count = watchers.len(), endpoint = %node.endpoint_address().id, "Appa is watching for changes"); + Ok(Self { + _lock: lock, + node, + watchers, + selected_folder_path: folder_path.map(Path::to_path_buf), + next_folder_refresh: tokio::time::Instant::now() + FOLDER_REFRESH_INTERVAL, + retries: BTreeMap::new(), + }) + } + + pub(crate) async fn tick(&mut self, service: &AppaService) -> AppResult<()> { + let now = tokio::time::Instant::now(); + if now >= self.next_folder_refresh { + service.refresh_watched_folders( + &mut self.watchers, + self.selected_folder_path.as_deref(), + &mut self.retries, + )?; + self.next_folder_refresh = now + FOLDER_REFRESH_INTERVAL; + } + service + .sync_watched_folders(&mut self.watchers, &self.node, &mut self.retries) + .await + } + + pub(crate) async fn shutdown(self) -> AppResult<()> { + self.node.shutdown().await + } +} + +impl AppaService { + pub async fn run_forever(&self, folder_path: Option<&Path>) -> AppResult<()> { + let mut runner = SyncRunner::start(self, folder_path).await?; + let mut interval = tokio::time::interval(WATCHER_POLL_INTERVAL); loop { tokio::select! { _ = tokio::signal::ctrl_c() => break, _ = interval.tick() => { - let now = tokio::time::Instant::now(); - if now >= next_folder_refresh { - self.refresh_watched_folders(&mut watchers, folder_path, &mut retries)?; - next_folder_refresh = now + FOLDER_REFRESH_INTERVAL; - } - self.sync_watched_folders(&mut watchers, &node, &mut retries).await?; + runner.tick(self).await?; } } } - node.shutdown().await?; - Ok(()) + runner.shutdown().await } fn refresh_watched_folders( diff --git a/src/cli.rs b/src/cli.rs index 7684e43..9c84434 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -56,6 +56,10 @@ struct CommandLine { #[derive(Debug, Subcommand)] enum Command { + Daemon { + #[command(subcommand)] + command: DaemonCommand, + }, Init { folder: String, }, @@ -181,12 +185,20 @@ enum ServiceCommand { Uninstall, } +#[derive(Debug, Subcommand)] +enum DaemonCommand { + Serve, +} + pub async fn run() -> anyhow::Result<()> { initialize_logging(); let command_line = CommandLine::parse(); tracing::debug!(command = ?command_line.command, "dispatched Appa command"); match command_line.command { + Command::Daemon { + command: DaemonCommand::Serve, + } => crate::daemon::serve().await?, Command::Service { command } => tooling::run_service_command(command)?, Command::Completions { shell } => tooling::print_completions(shell), command => run_app_command(command).await?, @@ -252,6 +264,7 @@ async fn run_app_command(command: Command) -> anyhow::Result<()> { Command::Run { folder } => run_folder(&appa, folder.as_deref().map(Path::new)).await?, Command::Doctor { json } => print_doctor(&appa, json)?, Command::Service { .. } => unreachable!("service commands are handled before opening Appa"), + Command::Daemon { .. } => unreachable!("daemon commands are handled before opening Appa"), } Ok(()) } diff --git a/src/core.rs b/src/core.rs index cabb496..8d9158c 100644 --- a/src/core.rs +++ b/src/core.rs @@ -2,6 +2,7 @@ use std::path::{Path, PathBuf}; +use serde::{Deserialize, Serialize}; use time::OffsetDateTime; use crate::{ @@ -11,7 +12,7 @@ use crate::{ pub type Result = anyhow::Result; -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct Dashboard { pub folders: Vec, pub diagnostics: Diagnostics, @@ -19,7 +20,7 @@ pub struct Dashboard { pub service_is_running: bool, } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct FolderSummary { pub name: String, pub path: PathBuf, @@ -31,7 +32,7 @@ pub struct FolderSummary { pub last_sync_error: Option, } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct FolderDetails { pub summary: FolderSummary, pub conflicts: Vec, @@ -40,20 +41,20 @@ pub struct FolderDetails { pub can_revoke_members: bool, } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct Conflict { pub path: String, pub author_device_id: String, pub modified_at: OffsetDateTime, } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct Member { pub device_id: String, pub label: Option, } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct Diagnostics { pub folder_count: u64, pub issues: Vec, diff --git a/src/daemon.rs b/src/daemon.rs new file mode 100644 index 0000000..5e86cab --- /dev/null +++ b/src/daemon.rs @@ -0,0 +1,136 @@ +//! Long-lived owner of Appa state, synchronization, and local IPC. + +use std::{fs, path::Path}; + +use anyhow::Context; +use tokio::{ + net::{UnixListener, UnixStream}, + time::{self, Duration}, +}; + +use crate::{ + app::{AppaService, run::SyncRunner}, + ipc::{self, Command, PROTOCOL_VERSION, Request, Response}, + storage::AppPaths, +}; + +const SOCKET_POLL_INTERVAL: Duration = Duration::from_millis(250); + +pub async fn serve() -> anyhow::Result<()> { + let paths = AppPaths::discover()?; + let socket_path = paths.daemon_socket_path(); + let service = AppaService::open_at(paths)?; + let mut sync_runner = SyncRunner::start(&service, None).await?; + remove_stale_socket(&socket_path)?; + let listener = UnixListener::bind(&socket_path).with_context(|| { + format!( + "could not bind Appa daemon socket at {}", + socket_path.display() + ) + })?; + let mut interval = time::interval(SOCKET_POLL_INTERVAL); + + tracing::info!("Appa daemon is ready"); + loop { + tokio::select! { + _ = tokio::signal::ctrl_c() => break, + _ = interval.tick() => sync_runner.tick(&service).await?, + accepted = listener.accept() => { + let (mut stream, _) = accepted?; + handle_connection(&service, &mut stream).await; + } + } + } + sync_runner.shutdown().await?; + remove_stale_socket(&socket_path) +} + +async fn handle_connection(service: &AppaService, stream: &mut UnixStream) { + let response = match ipc::read_message::(stream).await { + Ok(request) if request.protocol_version == PROTOCOL_VERSION => { + handle_command(service, request.command).await + } + Ok(request) => Response::Error(format!( + "unsupported Appa daemon protocol version {}; expected {PROTOCOL_VERSION}", + request.protocol_version + )), + Err(error) => Response::Error(error.to_string()), + }; + if let Err(error) = ipc::write_message(stream, &response).await { + tracing::debug!(%error, "Could not send Appa daemon response"); + } +} + +async fn handle_command(service: &AppaService, command: Command) -> Response { + let result = match command { + Command::Dashboard => service_dashboard(service).map(Response::Dashboard), + Command::FolderDetails { folder_path } => { + service_folder_details(service, &folder_path).map(Response::FolderDetails) + } + Command::RegisterFolder { folder_path } => service + .register_folder(&folder_path) + .map(|_| Response::Success), + Command::JoinFolder { + folder_path, + invitation, + } => service + .join_folder(&folder_path, &invitation) + .map(|_| Response::Success), + Command::CreateInvite { folder_path } => service + .create_invite(&folder_path) + .await + .map(Response::Invitation), + Command::RevokeMember { + folder_path, + device_id, + } => service + .revoke_member(&folder_path, &device_id) + .map(|_| Response::Success), + }; + result.unwrap_or_else(|error| Response::Error(error.to_string())) +} + +fn service_dashboard(service: &AppaService) -> anyhow::Result { + Ok(crate::core::Dashboard { + folders: service + .statuses()? + .iter() + .map(crate::core::FolderSummary::from) + .collect(), + diagnostics: crate::core::Diagnostics::from(service.doctor()?), + service_is_installed: crate::service::is_installed()?, + service_is_running: crate::service::is_running()?, + }) +} + +fn service_folder_details( + service: &AppaService, + folder_path: &Path, +) -> anyhow::Result { + let status = service.status(folder_path)?; + Ok(crate::core::FolderDetails { + summary: (&status).into(), + conflicts: service + .conflicts(folder_path)? + .into_iter() + .map(Into::into) + .collect(), + members: service + .members(folder_path)? + .into_iter() + .map(Into::into) + .collect(), + history_revision_count: status.history_revision_count, + can_revoke_members: service.can_revoke_members(folder_path)?, + }) +} + +fn remove_stale_socket(socket_path: &Path) -> anyhow::Result<()> { + match fs::remove_file(socket_path) { + Ok(()) => Ok(()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(error) => { + Err(error).with_context(|| format!("could not remove {}", socket_path.display())) + } + } +} diff --git a/src/ipc.rs b/src/ipc.rs new file mode 100644 index 0000000..32de706 --- /dev/null +++ b/src/ipc.rs @@ -0,0 +1,180 @@ +//! Local IPC protocol and client for the Appa daemon. + +use std::{io, path::PathBuf}; + +use anyhow::Context; +use serde::{Deserialize, Serialize}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::UnixStream, +}; + +use crate::{ + core::{Dashboard, FolderDetails}, + storage::AppPaths, +}; + +pub const PROTOCOL_VERSION: u16 = 1; +const MAX_MESSAGE_BYTES: usize = 1024 * 1024; + +#[derive(Debug, Deserialize, Serialize)] +pub(crate) struct Request { + pub protocol_version: u16, + pub command: Command, +} + +#[derive(Debug, Deserialize, Serialize)] +pub(crate) enum Command { + Dashboard, + FolderDetails { + folder_path: PathBuf, + }, + RegisterFolder { + folder_path: PathBuf, + }, + JoinFolder { + folder_path: PathBuf, + invitation: String, + }, + CreateInvite { + folder_path: PathBuf, + }, + RevokeMember { + folder_path: PathBuf, + device_id: String, + }, +} + +#[derive(Debug, Deserialize, Serialize)] +pub(crate) enum Response { + Dashboard(Dashboard), + FolderDetails(FolderDetails), + Invitation(String), + Success, + Error(String), +} + +/// A short-lived connection to the user-owned Appa daemon. +#[derive(Clone, Debug)] +pub struct DaemonClient { + socket_path: PathBuf, +} + +impl DaemonClient { + pub fn connect() -> anyhow::Result { + Ok(Self { + socket_path: AppPaths::discover()?.daemon_socket_path(), + }) + } + + pub async fn dashboard(&self) -> anyhow::Result { + match self.request(Command::Dashboard).await? { + Response::Dashboard(dashboard) => Ok(dashboard), + response => unexpected_response(response), + } + } + + pub async fn folder_details(&self, folder_path: PathBuf) -> anyhow::Result { + match self.request(Command::FolderDetails { folder_path }).await? { + Response::FolderDetails(details) => Ok(details), + response => unexpected_response(response), + } + } + + pub async fn register_folder(&self, folder_path: PathBuf) -> anyhow::Result<()> { + self.expect_success(Command::RegisterFolder { folder_path }) + .await + } + + pub async fn join_folder( + &self, + folder_path: PathBuf, + invitation: String, + ) -> anyhow::Result<()> { + self.expect_success(Command::JoinFolder { + folder_path, + invitation, + }) + .await + } + + pub async fn create_invite(&self, folder_path: PathBuf) -> anyhow::Result { + match self.request(Command::CreateInvite { folder_path }).await? { + Response::Invitation(invitation) => Ok(invitation), + response => unexpected_response(response), + } + } + + pub async fn revoke_member( + &self, + folder_path: PathBuf, + device_id: String, + ) -> anyhow::Result<()> { + self.expect_success(Command::RevokeMember { + folder_path, + device_id, + }) + .await + } + + async fn expect_success(&self, command: Command) -> anyhow::Result<()> { + match self.request(command).await? { + Response::Success => Ok(()), + response => unexpected_response(response), + } + } + + async fn request(&self, command: Command) -> anyhow::Result { + let mut stream = UnixStream::connect(&self.socket_path) + .await + .with_context(|| { + format!( + "could not connect to Appa daemon at {}", + self.socket_path.display() + ) + })?; + write_message( + &mut stream, + &Request { + protocol_version: PROTOCOL_VERSION, + command, + }, + ) + .await?; + read_message(&mut stream).await + } +} + +fn unexpected_response(response: Response) -> anyhow::Result { + match response { + Response::Error(error) => Err(anyhow::anyhow!(error)), + response => Err(anyhow::anyhow!( + "Appa daemon returned an unexpected response: {response:?}" + )), + } +} + +pub(crate) async fn read_message Deserialize<'de>>( + stream: &mut UnixStream, +) -> anyhow::Result { + let length = stream.read_u32().await? as usize; + if length > MAX_MESSAGE_BYTES { + anyhow::bail!("Appa daemon message exceeds {MAX_MESSAGE_BYTES} bytes"); + } + let mut bytes = vec![0; length]; + stream.read_exact(&mut bytes).await?; + Ok(serde_json::from_slice(&bytes)?) +} + +pub(crate) async fn write_message( + stream: &mut UnixStream, + message: &T, +) -> anyhow::Result<()> { + let bytes = serde_json::to_vec(message)?; + let length = u32::try_from(bytes.len()) + .map_err(|_| io::Error::other("Appa daemon message is too large"))?; + stream.write_u32(length).await?; + stream.write_all(&bytes).await?; + stream.flush().await?; + Ok(()) +} diff --git a/src/lib.rs b/src/lib.rs index 1b01562..efcdb0a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -3,12 +3,14 @@ //! [`Appa`] is the stable Rust interface for companion applications. The CLI //! remains available through [`run`]. -mod app; +pub(crate) mod app; mod cli; mod config; pub mod core; +mod daemon; mod domain; mod filesystem; +pub mod ipc; mod iroh; mod protocol; mod service; @@ -17,3 +19,4 @@ mod storage; /// Runs Appa's command-line interface. pub use cli::run; pub use core::Appa; +pub use service::install as install_service; diff --git a/src/service.rs b/src/service.rs index 2e503e4..2755570 100644 --- a/src/service.rs +++ b/src/service.rs @@ -175,7 +175,7 @@ fn unit_contents(executable: &Path, data_directory: &Path) -> anyhow::Result PathBuf { + self.data_directory.join("appad.sock") + } + pub fn load_identity(&self) -> anyhow::Result { if self.identity_path.exists() { return read_identity(&self.identity_path);