diff --git a/docs/architecture.md b/docs/architecture.md index 2420cef..044d046 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -33,6 +33,28 @@ place and keeps ordinary commands from starting another synchronizer. a pending materialization record first, so an interrupted write can be recovered on the next sync. +## Audit log + +Appa keeps a separate signed audit log for each folder. The compact manifest +history supports local restore and is allowed to prune old revisions. The audit +log is append-only and retained independently. + +Each device signs its own sequence of events. Every event names its parent +hash, the manifest root, and the roster hash that was current when it was +created. Devices exchange missing events through the authenticated control +protocol in bounded batches. A device that has retained a newer signed head can +detect a later rollback, omission, altered event, or conflicting event at the +same author sequence. + +This is not a global blockchain or a replacement for Appa's vector clocks. +Vector clocks continue to resolve file causality. Merkle trees continue to make +manifest comparison efficient. The audit log records who committed a state and +which prior event they extended. + +Audit evidence begins when an honest member has observed and retained a head. +It cannot prove that a first-seen state was not fabricated, protect against a +stolen signing key, or prove facts that no folder member retained. + All control streams and blob transfers are authenticated and encrypted by Iroh. The daemon uses mDNS on a LAN to find direct routes. Iroh uses relay routes when no direct route works. Discovery only finds routes. It does not @@ -55,7 +77,7 @@ Appa has four deliberately independent version signals. They protect different persisted or networked formats and should only change with their corresponding format. -- The Iroh control-stream ALPN is `appa/sync/3`. It selects a wire-compatible +- The Iroh control-stream ALPN is `appa/sync/4`. It selects a wire-compatible control protocol before either peer processes a request. - The invitation protocol version is `5`. It covers the signed invitation payload and its validation rules. diff --git a/src/app/sync.rs b/src/app/sync.rs index 06a3462..b2503ab 100644 --- a/src/app/sync.rs +++ b/src/app/sync.rs @@ -136,6 +136,7 @@ impl AppaService { active_member_ids: roster.active_device_ids(), roster: roster.clone(), peer_endpoints: self.state_store.peers(folder.id)?, + audit_events: self.state_store.audit_events(folder.id)?, }, manifest, ) @@ -150,6 +151,11 @@ impl AppaService { let audit_event = self.manifest_audit_event(manifest)?; self.state_store .save_manifest_with_audit(manifest, &audit_event)?; + node.publish_audit_events( + manifest.folder_id, + self.state_store.audit_events(manifest.folder_id)?, + ) + .await?; node.publish_manifest(manifest.clone()).await } @@ -281,6 +287,8 @@ impl AppaService { return Ok(None); } self.save_peer_response(folder, &peer, &roster, peer_endpoints, &capability)?; + self.receive_peer_audit_events(context, &peer, &roster) + .await?; if summary.root_hash == manifest_root_hash(local_manifest)? { tracing::debug!(folder = %folder.name, peer = %peer.id, "Folder roots match; skipping manifest transfer"); return Ok(None); @@ -318,6 +326,62 @@ impl AppaService { } } } + + async fn receive_peer_audit_events( + &self, + context: &SyncContext<'_>, + peer: &iroh::EndpointAddr, + roster: &FolderRoster, + ) -> AppResult<()> { + let events = match context + .node + .request_audit_events( + peer.clone(), + context.folder.id, + context.folder.capability.clone(), + self.state_store.join_invite(context.folder.id)?, + self.state_store + .audit_heads(context.folder.id)? + .into_iter() + .map(|head| (head.author_device_id, head.sequence)) + .collect(), + ) + .await + { + Ok(events) => events, + Err(error) => { + tracing::debug!(folder = %context.folder.name, peer = %peer.id, %error, "peer did not provide audit events"); + return Ok(()); + } + }; + for event in events { + if event.folder_id != context.folder.id { + anyhow::bail!("peer returned an audit event for another folder"); + } + if !roster.contains_member(&event.author_device_id) { + anyhow::bail!("peer returned an audit event from a non-member"); + } + if let Err(error) = self.state_store.append_audit_event(&event) { + let detail = format!( + "rejected audit event from {}: {error}", + event.author_device_id + ); + self.state_store.record_audit_fault( + context.folder.id, + &event.author_device_id, + &detail, + )?; + anyhow::bail!(detail); + } + } + context + .node + .publish_audit_events( + context.folder.id, + self.state_store.audit_events(context.folder.id)?, + ) + .await + } } async fn request_peer_summaries( diff --git a/src/app/tests.rs b/src/app/tests.rs index 97a7c3b..54599d9 100644 --- a/src/app/tests.rs +++ b/src/app/tests.rs @@ -183,6 +183,7 @@ async fn synchronizes_and_deletes_a_file_between_two_devices() -> anyhow::Result .load_roster(source_config.id)? .expect("roster"), peer_endpoints: vec![], + audit_events: vec![], }, source_manifest.clone(), ) @@ -292,6 +293,7 @@ async fn preserves_both_versions_after_offline_edits() -> anyhow::Result<()> { .load_roster(source_config.id)? .expect("roster"), peer_endpoints: vec![], + audit_events: vec![], }, source_manifest.clone(), ) @@ -339,6 +341,7 @@ async fn preserves_both_versions_after_offline_edits() -> anyhow::Result<()> { .load_roster(source_config.id)? .expect("roster"), peer_endpoints: vec![], + audit_events: vec![], }, target_manifest.clone(), ) diff --git a/src/iroh.rs b/src/iroh.rs index a737b8c..5b51ccd 100644 --- a/src/iroh.rs +++ b/src/iroh.rs @@ -25,7 +25,9 @@ use iroh_network::{ use tokio::sync::RwLock; use crate::{ - domain::{FolderId, FolderRoster, Manifest, ManifestMerkleIndex, ManifestMerkleNode}, + domain::{ + AuditEvent, FolderId, FolderRoster, Manifest, ManifestMerkleIndex, ManifestMerkleNode, + }, protocol::{ControlMessage, Invite, ManifestRequestKind, ManifestSummary}, }; @@ -33,8 +35,9 @@ mod handler; use handler::AppaProtocol; /// ALPN for Appa's Iroh control streams. A new value is wire-incompatible. -pub const APPA_ALPN: &[u8] = b"appa/sync/3"; +pub const APPA_ALPN: &[u8] = b"appa/sync/4"; const MAX_CONTROL_MESSAGE_BYTES: usize = 16 * 1024 * 1024; +pub(super) const MAX_AUDIT_EVENTS_PER_RESPONSE: usize = 256; const MAX_CONCURRENT_ANNOUNCEMENTS: usize = 4; const ANNOUNCEMENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); pub(super) const CONTROL_REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); @@ -49,6 +52,7 @@ pub struct FolderSession { pub active_member_ids: BTreeSet, pub roster: FolderRoster, pub peer_endpoints: Vec, + pub audit_events: Vec, } #[derive(Clone, Debug)] @@ -79,6 +83,7 @@ pub(super) struct HostedFolder { pub(super) peer_endpoints: Vec, pub(super) manifest: Manifest, pub(super) merkle_index: ManifestMerkleIndex, + pub(super) audit_events: Vec, } pub struct NodeHost { @@ -192,6 +197,7 @@ impl NodeHost { peer_endpoints: session.peer_endpoints, manifest, merkle_index, + audit_events: session.audit_events, }, ); Ok(()) @@ -239,6 +245,19 @@ impl NodeHost { Ok(()) } + pub async fn publish_audit_events( + &self, + folder_id: FolderId, + audit_events: Vec, + ) -> anyhow::Result<()> { + let mut folders = self.folders.write().await; + let Some(folder) = folders.get_mut(&folder_id) else { + anyhow::bail!("folder is not registered with this Appa node"); + }; + folder.audit_events = audit_events; + Ok(()) + } + pub async fn announce_manifest( &self, peers: Vec, @@ -383,6 +402,32 @@ impl NodeHost { Ok(remote_manifest) } + pub async fn request_audit_events( + &self, + peer_address: EndpointAddr, + folder_id: FolderId, + capability: String, + invite: Option, + after_sequences: BTreeMap, + ) -> anyhow::Result> { + let connection = self.connect_to_peer(peer_address).await?; + let response = request_control_message( + &connection, + &ControlMessage::Request { + folder_id, + capability, + invite: invite.map(Box::new), + requester_endpoint: self.endpoint_address(), + request: ManifestRequestKind::AuditEvents { after_sequences }, + }, + ) + .await?; + let ControlMessage::AuditEvents { events } = response else { + anyhow::bail!("peer returned an unexpected audit response"); + }; + Ok(events) + } + async fn request_manifest_merkle_node( &self, connection: &Connection, diff --git a/src/iroh/handler.rs b/src/iroh/handler.rs index 5b8f505..a70efd8 100644 --- a/src/iroh/handler.rs +++ b/src/iroh/handler.rs @@ -13,8 +13,8 @@ use tokio::sync::RwLock; use crate::{ domain::FolderId, iroh::{ - CONTROL_REQUEST_TIMEOUT, DiscoveredPeer, HostedFolder, MAX_CONTROL_MESSAGE_BYTES, - recover_mutex, + CONTROL_REQUEST_TIMEOUT, DiscoveredPeer, HostedFolder, MAX_AUDIT_EVENTS_PER_RESPONSE, + MAX_CONTROL_MESSAGE_BYTES, recover_mutex, }, protocol::{ControlMessage, ManifestRequestKind, ManifestSummary, validate_invite}, }; @@ -204,6 +204,23 @@ impl AppaProtocol { .map_err(protocol_error)? .ok_or_else(|| protocol_error("unknown Merkle prefix"))?, }), + ManifestRequestKind::AuditEvents { after_sequences } => { + Ok(ControlMessage::AuditEvents { + events: folder + .audit_events + .iter() + .filter(|event| { + event.sequence + > after_sequences + .get(&event.author_device_id) + .copied() + .unwrap_or_default() + }) + .take(MAX_AUDIT_EVENTS_PER_RESPONSE) + .cloned() + .collect(), + }) + } } } } diff --git a/src/iroh/tests.rs b/src/iroh/tests.rs index bc32405..5ba271e 100644 --- a/src/iroh/tests.rs +++ b/src/iroh/tests.rs @@ -120,6 +120,7 @@ async fn announces_folder_changes_to_an_authorized_peer() -> anyhow::Result<()> active_member_ids: BTreeSet::from([source.endpoint_address().id.to_string()]), roster: test_roster(folder_id, [source.endpoint_address().id.to_string()]), peer_endpoints: vec![source.endpoint_address()], + audit_events: vec![], }, crate::domain::Manifest::empty(folder_id), ) @@ -173,6 +174,7 @@ async fn requests_remote_merkle_children_missing_from_a_small_local_manifest() - active_member_ids: BTreeSet::from([target.endpoint_address().id.to_string()]), roster: test_roster(folder_id, [target.endpoint_address().id.to_string()]), peer_endpoints: vec![target.endpoint_address()], + audit_events: vec![], }, remote_manifest.clone(), ) @@ -249,6 +251,7 @@ async fn serves_multiple_isolated_folders_from_one_endpoint() -> anyhow::Result< ], ), peer_endpoints: vec![target.endpoint_address(), third.endpoint_address()], + audit_events: vec![], }, crate::domain::Manifest::empty(first_folder), ) @@ -261,6 +264,7 @@ async fn serves_multiple_isolated_folders_from_one_endpoint() -> anyhow::Result< active_member_ids: BTreeSet::from([target.endpoint_address().id.to_string()]), roster: test_roster(second_folder, [target.endpoint_address().id.to_string()]), peer_endpoints: vec![], + audit_events: vec![], }, crate::domain::Manifest::empty(second_folder), ) @@ -329,6 +333,7 @@ async fn serves_multiple_isolated_folders_from_one_endpoint() -> anyhow::Result< active_member_ids: BTreeSet::from([target.endpoint_address().id.to_string()]), roster: test_roster(rotated_folder, [target.endpoint_address().id.to_string()]), peer_endpoints: vec![], + audit_events: vec![], }, crate::domain::Manifest::empty(rotated_folder), ) diff --git a/src/protocol.rs b/src/protocol.rs index 34c54df..5a431ae 100644 --- a/src/protocol.rs +++ b/src/protocol.rs @@ -1,6 +1,10 @@ +use std::collections::BTreeMap; + use anyhow::Context; -use crate::domain::{DeviceId, FolderId, FolderRoster, ManifestMerkleNode, canonical_json_bytes}; +use crate::domain::{ + AuditEvent, DeviceId, FolderId, FolderRoster, ManifestMerkleNode, canonical_json_bytes, +}; #[cfg(test)] use crate::domain::{Manifest, manifest_root_hash}; use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; @@ -32,7 +36,12 @@ impl ManifestSummary { #[serde(rename_all = "snake_case")] pub enum ManifestRequestKind { Summary, - MerkleNode { prefix: String }, + MerkleNode { + prefix: String, + }, + AuditEvents { + after_sequences: BTreeMap, + }, } #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] @@ -119,6 +128,9 @@ pub enum ControlMessage { MerkleNode { node: ManifestMerkleNode, }, + AuditEvents { + events: Vec, + }, Announcement { folder_id: FolderId, capability: String, diff --git a/src/storage/audit.rs b/src/storage/audit.rs index f247bc5..a3acc99 100644 --- a/src/storage/audit.rs +++ b/src/storage/audit.rs @@ -44,7 +44,6 @@ impl StateStore { }) } - #[cfg(test)] pub fn append_audit_event(&self, event: &AuditEvent) -> anyhow::Result<()> { event.validate()?; let transaction = self.connection.unchecked_transaction()?; @@ -101,6 +100,19 @@ impl StateStore { }) } + pub fn record_audit_fault( + &self, + folder_id: FolderId, + author_device_id: &str, + detail: &str, + ) -> anyhow::Result<()> { + self.connection.execute( + "INSERT INTO audit_integrity_faults (folder_id, author_device_id, detected_at, detail) VALUES (?1, ?2, datetime('now'), ?3)", + params![folder_id.to_string(), author_device_id, detail], + )?; + Ok(()) + } + fn audit_head( &self, folder_id: FolderId, @@ -144,6 +156,9 @@ pub(super) fn append_audit_event( event: &AuditEvent, ) -> anyhow::Result<()> { event.validate()?; + if event_already_exists(transaction, event)? { + return Ok(()); + } let current_head = transaction .query_row( "SELECT sequence, event_hash FROM audit_heads WHERE folder_id = ?1 AND author_device_id = ?2", @@ -164,6 +179,24 @@ pub(super) fn append_audit_event( Ok(()) } +fn event_already_exists( + transaction: &rusqlite::Transaction<'_>, + event: &AuditEvent, +) -> anyhow::Result { + let stored_hash = transaction + .query_row( + "SELECT event_hash FROM audit_events WHERE folder_id = ?1 AND author_device_id = ?2 AND sequence = ?3", + params![event.folder_id.to_string(), event.author_device_id, i64::try_from(event.sequence)?], + |row| row.get::<_, String>(0), + ) + .optional()?; + match stored_hash { + None => Ok(false), + Some(stored_hash) if stored_hash == event.hash()? => Ok(true), + Some(_) => anyhow::bail!("audit event conflicts with an existing author sequence"), + } +} + fn validate_new_event( event: &AuditEvent, current_head: Option<(i64, String)>,