diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index 32fa3c0..8b5ff37 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -1,5 +1,5 @@ use crate::message_common::*; -use iota_logger::log; +use iota_logger::{LogLevel, PrintType, log_event}; use iota_storage::util::chat_files::{self, MessageState}; use iota_storage::util::chats_util::{self, get_user, has_user, mod_user}; use iota_storage::util::communities_util::CommunitiesUtil; @@ -19,6 +19,24 @@ use std::sync::atomic::{AtomicU32, Ordering}; static NEXT_NOTIFICATION_ID: AtomicU32 = AtomicU32::new(1); +fn internal_handler_error( + cv: &CommunicationValue, + operation: &'static str, + error: impl std::fmt::Display, +) -> CommunicationValue { + log_event!( + PrintType::Client, + LogLevel::Error, + "request.internal_error", + "operation={} message_id={:?} sender={:?} error={}", + operation, + cv.id(), + cv.sender(), + error + ); + error_response(cv, CommunicationType::ErrorInternal) +} + fn next_notification_id() -> u32 { NEXT_NOTIFICATION_ID.fetch_add(1, Ordering::Relaxed).max(1) } @@ -915,14 +933,14 @@ pub fn handle_get_chat_secret(cv: &CommunicationValue) -> CommunicationValue { None => return error_response(cv, CommunicationType::ErrorNotFound), }, Ok(None) => return error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "chat_secret.contact_get", error), }; let owner_principal = match iota_storage::identity::SqlitePrincipalStore .principal_for_local_user(iota_identity::LocalUserId(sender_id)) { Ok(Some(principal)) => principal, Ok(None) => return error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "chat_secret.principal_get", error), }; let Some(principal_chat_id) = e2ee_storage::principal_chat_id(owner_principal, partner_principal) @@ -1131,15 +1149,15 @@ pub fn handle_account_state_request(cv: &CommunicationValue) -> CommunicationVal }; let contacts = match chats_util::get_users(user_id) { Ok(contacts) => contacts, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let messages = match chat_files::get_all_messages(user_id) { Ok(messages) => messages, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let settings = match synced_settings::list(user_id) { Ok(settings) => settings, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let message_values = messages .iter() @@ -1147,23 +1165,23 @@ pub fn handle_account_state_request(cv: &CommunicationValue) -> CommunicationVal .collect(); let blocked_users = match blocked_users::list(user_id) { Ok(users) => users, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let receipt_policy = match receipt_policy::get(user_id) { Ok(policy) => policy, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let message_storage_policy = match message_storage_policy::get(user_id) { Ok(policy) => policy, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let contact_ids = match current_contact_ids(user_id) { Ok(contact_ids) => contact_ids, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let communities = match community_values(user_id) { Ok(communities) => communities, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "account_state.snapshot", error), }; let retention_duration = match message_storage_policy.retention { message_storage_policy::MessageRetention::Forever => None, @@ -1284,7 +1302,7 @@ pub fn handle_messages_get(cv: &CommunicationValue) -> CommunicationValue { }; let messages = match chat_files::get_messages(my_id_i64, partner_id, offset, amount) { Ok(messages) => messages, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "messages.get", error), }; let mut msg_array: Vec = Vec::new(); for m in &messages { @@ -1312,12 +1330,12 @@ pub fn handle_message_get(cv: &CommunicationValue) -> CommunicationValue { { Ok(Some((message, offset))) => (message, Some(offset)), Ok(None) => return error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "message.get", error), }, None => match chat_files::get_message(owner, send_time, None, None) { Ok(Some(message)) => (message, None), Ok(None) => return error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "message.get", error), }, }; @@ -1344,7 +1362,7 @@ pub fn handle_get_chats(cv: &CommunicationValue) -> CommunicationValue { }; let users = match chats_util::get_users(user_id_i64) { Ok(users) => users, - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "chats.get", error), }; let mut user_array = Vec::new(); for user in users { @@ -1554,8 +1572,8 @@ pub fn handle_global_settings_save(cv: &CommunicationValue) -> CommunicationValu ); }; - if settings::save_global(my_id_i64, settings_value).is_err() { - return error_response(cv, CommunicationType::ErrorInternal); + if let Err(error) = settings::save_global(my_id_i64, settings_value) { + return internal_handler_error(cv, "settings.global_save", error); } let mut response = CommunicationValue::new(CommunicationType::GlobalSettingsSave) @@ -1702,8 +1720,8 @@ pub fn handle_settings_save( return error_response(cv, CommunicationType::ErrorInvalidData); }; - if settings::save(my_id_i64, session_id_i64, settings_name, settings_value).is_err() { - return error_response(cv, CommunicationType::ErrorInternal); + if let Err(error) = settings::save(my_id_i64, session_id_i64, settings_name, settings_value) { + return internal_handler_error(cv, "settings.save", error); } CommunicationValue::new(CommunicationType::SettingsSave) @@ -2085,17 +2103,23 @@ fn upload_response( } fn asset_error(cv: &CommunicationValue, error: StorageError) -> CommunicationValue { - log!( - "Asset request rejected sender={:?} request_id={:?}: {error}", - cv.sender(), - cv.id() - ); - let (kind, category) = match error { + let (kind, category) = match &error { StorageError::AssetResourceLimit(_) => { (CommunicationType::ErrorInvalidData, "asset_resource_limit") } _ => (CommunicationType::ErrorInternal, "asset_request_failed"), }; + if kind == CommunicationType::ErrorInternal { + log_event!( + PrintType::Client, + LogLevel::Error, + "request.internal_error", + "operation=user_assets message_id={:?} sender={:?} error={}", + cv.id(), + cv.sender(), + error + ); + } error_response(cv, kind).add_typed_default(DataType::ErrorType, DataValue::Str(category.into())) } @@ -2212,6 +2236,7 @@ fn parse_setting_locator(cv: &CommunicationValue) -> Result Result<(), CommunicationType> { @@ -2225,14 +2250,20 @@ fn validate_setting_target( match has_user(user_id, contact_id) { Ok(true) => Ok(()), Ok(false) => Err(CommunicationType::ErrorInvalidData), - Err(_) => Err(CommunicationType::ErrorInternal), + Err(error) => { + internal_handler_error(cv, "settings.validate_contact", error); + Err(CommunicationType::ErrorInternal) + } } } SettingScope::Community => { match CommunitiesUtil::has_community(user_id, &locator.scope_key) { Ok(true) => Ok(()), Ok(false) => Err(CommunicationType::ErrorInvalidData), - Err(_) => Err(CommunicationType::ErrorInternal), + Err(error) => { + internal_handler_error(cv, "settings.validate_community", error); + Err(CommunicationType::ErrorInternal) + } } } } @@ -2262,7 +2293,7 @@ pub fn handle_synced_setting_set(cv: &CommunicationValue) -> SettingMutation { }; } }; - if let Err(error_type) = validate_setting_target(user_id, &locator) { + if let Err(error_type) = validate_setting_target(cv, user_id, &locator) { return setting_mutation_error(cv, error_type); } let Some(payload) = cv.get_data(DataType::Payload).as_str() else { @@ -2279,7 +2310,10 @@ pub fn handle_synced_setting_set(cv: &CommunicationValue) -> SettingMutation { response: setting_response(cv, CommunicationType::SyncedSettingSet, &setting), changed: Some(setting_changed(&setting)), }, - Err(_) => setting_mutation_error(cv, CommunicationType::ErrorInternal), + Err(error) => SettingMutation { + response: internal_handler_error(cv, "settings.set", error), + changed: None, + }, } } @@ -2292,13 +2326,13 @@ pub fn handle_synced_setting_get(cv: &CommunicationValue) -> CommunicationValue Ok(locator) => locator, Err(response) => return response, }; - if let Err(error_type) = validate_setting_target(user_id, &locator) { + if let Err(error_type) = validate_setting_target(cv, user_id, &locator) { return error_response(cv, error_type); } match synced_settings::get(user_id, locator.scope, &locator.scope_key, &locator.name) { Ok(Some(setting)) => setting_response(cv, CommunicationType::SyncedSettingGet, &setting), Ok(None) => error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "settings.get", error), } } @@ -2316,7 +2350,7 @@ pub fn handle_synced_setting_delete(cv: &CommunicationValue) -> SettingMutation }; } }; - if let Err(error_type) = validate_setting_target(user_id, &locator) { + if let Err(error_type) = validate_setting_target(cv, user_id, &locator) { return setting_mutation_error(cv, error_type); } match synced_settings::delete(user_id, locator.scope, &locator.scope_key, &locator.name) { @@ -2340,7 +2374,10 @@ pub fn handle_synced_setting_delete(cv: &CommunicationValue) -> SettingMutation .with_receiver(sender_wire_id(user_id)), changed: None, }, - Err(_) => setting_mutation_error(cv, CommunicationType::ErrorInternal), + Err(error) => SettingMutation { + response: internal_handler_error(cv, "settings.delete", error), + changed: None, + }, } } @@ -2357,7 +2394,7 @@ pub fn handle_synced_settings_list(cv: &CommunicationValue) -> CommunicationValu DataType::Settings, DataValue::Array(settings.iter().map(synced_setting_value).collect()), ), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "settings.list", error), } } @@ -2384,8 +2421,8 @@ pub fn handle_user_blob_put(cv: &CommunicationValue) -> BlobMutation { response: error_response(cv, CommunicationType::ErrorInvalidData), changed: None, }, - Err(_) => BlobMutation { - response: error_response(cv, CommunicationType::ErrorInternal), + Err(error) => BlobMutation { + response: internal_handler_error(cv, "user_blobs.put", error), changed: None, }, } @@ -2398,7 +2435,7 @@ pub fn handle_user_blob_get(cv: &CommunicationValue) -> CommunicationValue { match user_blobs::get(user_id, &blob_id) { Ok(Some(blob)) => blob_response(cv, CommunicationType::UserBlobGet, &blob), Ok(None) => error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "user_blobs.get", error), } } @@ -2433,8 +2470,8 @@ pub fn handle_user_blob_delete(cv: &CommunicationValue) -> BlobMutation { response: error_response(cv, CommunicationType::ErrorInvalidData), changed: None, }, - Err(_) => BlobMutation { - response: error_response(cv, CommunicationType::ErrorInternal), + Err(error) => BlobMutation { + response: internal_handler_error(cv, "user_blobs.delete", error), changed: None, }, } @@ -2466,7 +2503,7 @@ pub fn handle_user_blob_list(cv: &CommunicationValue) -> CommunicationValue { DataType::Blobs, DataValue::Array(blobs.iter().map(blob_metadata_value).collect()), ), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "user_blobs.list", error), } } @@ -2719,8 +2756,8 @@ pub fn handle_user_block(cv: &CommunicationValue) -> PolicyMutation { .add_typed_default(DataType::Deleted, DataValue::Bool(record.deleted)), ), }, - Err(_) => PolicyMutation { - response: error_response(cv, CommunicationType::ErrorInternal), + Err(error) => PolicyMutation { + response: internal_handler_error(cv, "blocked_users.block", error), changed: None, }, } @@ -2775,8 +2812,8 @@ pub fn handle_user_unblock(cv: &CommunicationValue) -> PolicyMutation { .add_typed_default(DataType::Deleted, DataValue::Bool(mutation.deleted)) }), }, - Err(_) => PolicyMutation { - response: error_response(cv, CommunicationType::ErrorInternal), + Err(error) => PolicyMutation { + response: internal_handler_error(cv, "blocked_users.unblock", error), changed: None, }, } @@ -2799,7 +2836,7 @@ pub fn handle_blocked_users_get(cv: &CommunicationValue) -> CommunicationValue { .collect(), ), ), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "blocked_users.get", error), } } @@ -2819,7 +2856,7 @@ pub fn handle_receipt_policy_get(cv: &CommunicationValue) -> CommunicationValue }; match receipt_policy::get(user_id) { Ok(policy) => receipt_response(cv, CommunicationType::ReceiptPolicyGet, policy), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "receipt_policy.get", error), } } pub fn handle_receipt_policy_set(cv: &CommunicationValue) -> PolicyMutation { @@ -2850,8 +2887,8 @@ pub fn handle_receipt_policy_set(cv: &CommunicationValue) -> PolicyMutation { changed: Some(changed), } } - Err(_) => PolicyMutation { - response: error_response(cv, CommunicationType::ErrorInternal), + Err(error) => PolicyMutation { + response: internal_handler_error(cv, "receipt_policy.set", error), changed: None, }, } @@ -2900,7 +2937,7 @@ pub fn handle_message_storage_policy_get(cv: &CommunicationValue) -> Communicati Ok(policy) => { storage_policy_response(cv, CommunicationType::MessageStoragePolicyGet, policy) } - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "message_storage_policy.get", error), } } @@ -2948,8 +2985,8 @@ pub fn handle_message_storage_policy_set(cv: &CommunicationValue) -> PolicyMutat changed: Some(changed), } } - Err(_) => PolicyMutation { - response: error_response(cv, CommunicationType::ErrorInvalidData), + Err(error) => PolicyMutation { + response: internal_handler_error(cv, "message_storage_policy.set", error), changed: None, }, } @@ -2970,13 +3007,13 @@ pub fn handle_user_block_check(cv: &CommunicationValue) -> CommunicationValue { None => return error_response(cv, CommunicationType::ErrorNotFound), }, Ok(None) => return error_response(cv, CommunicationType::ErrorNotFound), - Err(_) => return error_response(cv, CommunicationType::ErrorInternal), + Err(error) => return internal_handler_error(cv, "blocked_users.check_contact", error), }; match blocked_users::is_principal_blocked(receiver_id, sender_principal) { Ok(blocked) => CommunicationValue::new(CommunicationType::UserBlockCheck) .with_request_id(cv) .add_typed_default(DataType::IsBlocked, DataValue::Bool(blocked)), - Err(_) => error_response(cv, CommunicationType::ErrorInternal), + Err(error) => internal_handler_error(cv, "blocked_users.check", error), } } diff --git a/iota-daemon-lib/src/command_router.rs b/iota-daemon-lib/src/command_router.rs index 9fb5b86..3f9fd23 100644 --- a/iota-daemon-lib/src/command_router.rs +++ b/iota-daemon-lib/src/command_router.rs @@ -8,7 +8,7 @@ use iota_ipc::{ TaskSummary, UpdateStatusResponse, UserDetailResponse, UserDiagnostics, UserOperationKind, UserOperationSummary, UserReconcileResult, UserSummary, }; -use iota_logger::{log, log_command}; +use iota_logger::{LogLevel, PrintType, log, log_command, log_event}; use iota_storage::users::pending_operations::{self, PendingUserOperationKind}; use iota_storage::users::user_manager; use iota_storage::util::config_util::{self}; @@ -64,6 +64,7 @@ fn now_millis() -> i64 { } fn invitation_created_after_mirror( + request_id: u64, authority: InvitationAuthority, invitation_id: i64, raw_token: String, @@ -74,8 +75,12 @@ fn invitation_created_after_mirror( let mirror_synced = match mirror_result { Ok(()) => true, Err(error) => { - log!( - "Invitation {} was created by Omega, but its local mirror could not be stored: {}", + log_event!( + PrintType::Command, + LogLevel::Error, + "invitation.mirror_store_failed", + "request_id={} invitation_id={} error={}", + request_id, invitation_id, error ); @@ -139,12 +144,16 @@ impl CommandRouter { request: LocalRequest, ) -> ResponseEnvelope { if !peer.role.allows(request.required_role()) { - log!( - "IPC authorization denied: pid={}, uid={}, role={:?}, request={:?}", + log_event!( + PrintType::Command, + LogLevel::Warn, + "ipc.authorization_denied", + "request_id={} pid={} uid={} role={:?} request={}", + request_id, peer.pid, peer.uid, peer.role, - request + request.log_name() ); return ResponseEnvelope { request_id, @@ -153,17 +162,18 @@ impl CommandRouter { } log_command!( - "pid={} uid={} role={:?} request={:?}", + "request_id={} pid={} uid={} role={:?} request={}", + request_id, peer.pid, peer.uid, peer.role, - request + request.log_name() ); - let result = self.execute(request).await; + let result = self.execute(request_id, request).await; ResponseEnvelope { request_id, result } } - async fn execute(&self, request: LocalRequest) -> ResponseResult { + async fn execute(&self, request_id: u64, request: LocalRequest) -> ResponseResult { if !self.services.active && !matches!( request, @@ -248,11 +258,31 @@ impl CommandRouter { ) }) .collect::>(), - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.pending_operations_load_failed", + "request_id={} error={}", + request_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; let users = match user_manager::get_residency() { Ok(users) => users, - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.residency_load_failed", + "request_id={} error={}", + request_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } } .into_iter() .map(|user| { @@ -275,8 +305,19 @@ impl CommandRouter { }) }) .collect::, iota_storage::storage_error::StorageError>>(); - let Ok(users) = users else { - return ResponseResult::Error(IpcErrorCode::StorageFailure); + let users = match users { + Ok(users) => users, + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.profile_load_failed", + "request_id={} error={}", + request_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; ResponseResult::Ok(ResponsePayload::Users(users)) } @@ -386,7 +427,17 @@ impl CommandRouter { Err(omikron_connector::OmikronError::Timeout(_)) => { return ResponseResult::Error(IpcErrorCode::Timeout); } - Err(_) => return ResponseResult::Error(IpcErrorCode::OmikronUnavailable), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Warn, + "invitation.create_failed", + "request_id={} error={}", + request_id, + error + ); + return ResponseResult::Error(IpcErrorCode::OmikronUnavailable); + } }; let invitation_id = response .get_data(DataType::InvitationId) @@ -419,6 +470,13 @@ impl CommandRouter { .zip(expires_at) .map(|(((id, token), created), expires)| (id, token, created, expires)) else { + log_event!( + PrintType::Command, + LogLevel::Error, + "invitation.invalid_response", + "request_id={} reason=missing_required_fields", + request_id + ); return ResponseResult::Error(IpcErrorCode::InternalFailure); }; let summary = iota_storage::users::invitations::InvitationSummary { @@ -443,6 +501,7 @@ impl CommandRouter { iota_storage::users::invitations::insert(&summary, None, now_millis()); ResponseResult::Ok(ResponsePayload::InvitationCreated( invitation_created_after_mirror( + request_id, authority, invitation_id, raw_token, @@ -471,12 +530,30 @@ impl CommandRouter { Err(_) => return ResponseResult::Error(IpcErrorCode::InternalFailure), } } - if iota_storage::users::invitations::expire_pending(now_millis()).is_err() { + if let Err(error) = iota_storage::users::invitations::expire_pending(now_millis()) { + log_event!( + PrintType::Command, + LogLevel::Error, + "invitation.expire_failed", + "request_id={} error={}", + request_id, + error + ); return ResponseResult::Error(IpcErrorCode::StorageFailure); } let invitations = match iota_storage::users::invitations::list() { Ok(invitations) => invitations, - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "invitation.list_failed", + "request_id={} error={}", + request_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; let mut summaries = Vec::new(); for invitation in invitations { @@ -875,10 +952,30 @@ impl CommandRouter { Err(AccountError::Unavailable(_)) => { ResponseResult::Error(IpcErrorCode::OmikronUnavailable) } - Err(AccountError::Storage(_)) => { + Err(AccountError::Storage(error)) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.release_storage_failed", + "request_id={} user_id={} error={:?}", + request_id, + user_id, + error + ); ResponseResult::Error(IpcErrorCode::StorageFailure) } - Err(_) => ResponseResult::Error(IpcErrorCode::InternalFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.release_failed", + "request_id={} user_id={} error={:?}", + request_id, + user_id, + error + ); + ResponseResult::Error(IpcErrorCode::InternalFailure) + } } } LocalRequest::ReconnectOmikron => match self.services.omikron() { @@ -886,7 +983,17 @@ impl CommandRouter { Ok(()) => ResponseResult::Ok(ResponsePayload::Acknowledged { message: "Reconnected to Omikron server".into(), }), - Err(_) => ResponseResult::Error(IpcErrorCode::OmikronUnavailable), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Warn, + "ipc.omikron_reconnect_failed", + "request_id={} error={}", + request_id, + error + ); + ResponseResult::Error(IpcErrorCode::OmikronUnavailable) + } }, None => ResponseResult::Error(IpcErrorCode::OmikronUnavailable), }, @@ -943,17 +1050,48 @@ impl CommandRouter { } LocalRequest::SetConfig { key, value } => { match config_util::modify_config_value(&key, &value) { - Ok(()) => ResponseResult::Ok(ResponsePayload::Acknowledged { - message: format!("Set {key} = {value}"), - }), - Err(_e) => ResponseResult::Error(IpcErrorCode::InvalidRequest), + Ok(()) => { + log_event!( + PrintType::Command, + LogLevel::Info, + "config.changed", + "request_id={} key={:?}", + request_id, + key + ); + ResponseResult::Ok(ResponsePayload::Acknowledged { + message: format!("Set {key}"), + }) + } + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Warn, + "config.change_failed", + "request_id={} key={:?} error={}", + request_id, + key, + error + ); + ResponseResult::Error(IpcErrorCode::InvalidRequest) + } } } LocalRequest::ReloadConfig => match config_util::load_config() { Ok(()) => ResponseResult::Ok(ResponsePayload::Acknowledged { message: "Configuration reloaded".into(), }), - Err(_) => ResponseResult::Error(IpcErrorCode::InvalidRequest), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Warn, + "config.reload_failed", + "request_id={} error={}", + request_id, + error + ); + ResponseResult::Error(IpcErrorCode::InvalidRequest) + } }, LocalRequest::GetOmikronStatus => { let connected = match self.services.omikron() { @@ -983,7 +1121,18 @@ impl CommandRouter { let residency = match user_manager::get_residency_by_id(user_id) { Ok(Some(residency)) => residency, Ok(None) => return ResponseResult::Error(IpcErrorCode::NotFound), - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.residency_load_failed", + "request_id={} user_id={} error={}", + request_id, + user_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; match user_manager::get_user(user_id) { Ok(Some(user)) => { @@ -1016,14 +1165,36 @@ impl CommandRouter { })) } Ok(None) => ResponseResult::Error(IpcErrorCode::StorageFailure), - Err(_) => ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.profile_load_failed", + "request_id={} user_id={} error={}", + request_id, + user_id, + error + ); + ResponseResult::Error(IpcErrorCode::StorageFailure) + } } } LocalRequest::ExportUserCredential { user_id } => { let residency = match user_manager::get_residency_by_id(user_id) { Ok(Some(residency)) => residency, Ok(None) => return ResponseResult::Error(IpcErrorCode::NotFound), - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.residency_load_failed", + "request_id={} user_id={} error={}", + request_id, + user_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; if residency.state != user_manager::LocalUserState::Managed || residency.credential_origin != user_manager::CredentialOrigin::Local @@ -1036,7 +1207,18 @@ impl CommandRouter { ) { Ok(Some(credential)) => credential, Ok(None) => return ResponseResult::Error(IpcErrorCode::NotFound), - Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "user.credential_read_failed", + "request_id={} user_id={} error={}", + request_id, + user_id, + error + ); + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } }; let parsed = match iota_util::tu::TuCredential::parse(&credential) { Ok(parsed) if parsed.user_id == user_id => parsed, @@ -1063,7 +1245,17 @@ impl CommandRouter { available, })) } - Err(_e) => ResponseResult::Error(IpcErrorCode::InternalFailure), + Err(error) => { + log_event!( + PrintType::Command, + LogLevel::Error, + "update.check_failed", + "request_id={} error={}", + request_id, + error + ); + ResponseResult::Error(IpcErrorCode::InternalFailure) + } }, LocalRequest::ListCommunities => { let iota_id = config_util::CONFIG.load().iota_id; @@ -1182,6 +1374,7 @@ mod tests { #[test] fn omega_creation_credentials_survive_a_local_mirror_failure() { let created = invitation_created_after_mirror( + 1, InvitationAuthority::Omega, 7, "raw-token".into(), diff --git a/iota-daemon-lib/src/log_broadcaster.rs b/iota-daemon-lib/src/log_broadcaster.rs index 9bb1d05..51f1feb 100644 --- a/iota-daemon-lib/src/log_broadcaster.rs +++ b/iota-daemon-lib/src/log_broadcaster.rs @@ -1,27 +1,95 @@ use crate::log_buffer::LogBuffer; use iota_ipc::{DaemonMessage, LogEntry}; use iota_logger::subscribe; +use iota_state::UiLogEntry; use std::sync::{Arc, Mutex}; use tokio::sync::broadcast; +use tokio::sync::broadcast::error::RecvError; /* The daemon adapts logger output to the wire protocol so the logger stays * independent from both the socket implementation and TUI state. */ pub fn spawn(message_tx: broadcast::Sender, buffer: Arc>) { - let Some(mut logs) = subscribe() else { + let Some(logs) = subscribe() else { return; }; - tokio::spawn(async move { - while let Ok(entry) = logs.recv().await { - let entry = LogEntry { - timestamp_ms: entry.timestamp_ms, - sender: entry.sender, - message: entry.message, - is_error: entry.is_error, - }; - if let Ok(mut buf) = buffer.lock() { - buf.push(entry.clone()); + tokio::spawn(forward_logs(logs, message_tx, buffer)); +} + +pub(crate) async fn forward_logs( + mut logs: broadcast::Receiver, + message_tx: broadcast::Sender, + buffer: Arc>, +) { + loop { + match logs.recv().await { + Ok(entry) => { + let entry = LogEntry { + timestamp_ms: entry.timestamp_ms, + sender: entry.sender, + message: entry.message, + is_error: entry.is_error, + }; + if let Ok(mut buf) = buffer.lock() { + buf.push(entry.clone()); + } + let _ = message_tx.send(DaemonMessage::LogEntry(entry)); } - let _ = message_tx.send(DaemonMessage::LogEntry(entry)); + Err(RecvError::Lagged(skipped)) => { + let _ = message_tx.send(DaemonMessage::Gap { skipped }); + } + Err(RecvError::Closed) => break, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn lag_reports_gap_and_continues_forwarding() { + let (logs_tx, logs_rx) = broadcast::channel(2); + let (message_tx, mut messages) = broadcast::channel(16); + let buffer = Arc::new(Mutex::new(LogBuffer::new(8))); + for index in 0..4 { + let _ = logs_tx.send(UiLogEntry { + timestamp_ms: index, + sender: "Client".into(), + message: index.to_string(), + is_error: false, + }); + } + let task = tokio::spawn(forward_logs(logs_rx, message_tx, buffer.clone())); + assert!(matches!( + messages.recv().await.unwrap(), + DaemonMessage::Gap { skipped: 2 } + )); + assert!(matches!( + messages.recv().await.unwrap(), + DaemonMessage::LogEntry(_) + )); + let _ = logs_tx.send(UiLogEntry { + timestamp_ms: 4, + sender: "Client".into(), + message: "after lag".into(), + is_error: false, + }); + let mut found = false; + for _ in 0..2 { + if let DaemonMessage::LogEntry(entry) = messages.recv().await.unwrap() { + found |= entry.message == "after lag"; + } } - }); + assert!(found); + drop(logs_tx); + task.await.unwrap(); + assert!( + buffer + .lock() + .unwrap() + .recent(4) + .iter() + .any(|entry| entry.message == "after lag") + ); + } } diff --git a/iota-ipc/src/protocol.rs b/iota-ipc/src/protocol.rs index 09e881b..9ca95e3 100644 --- a/iota-ipc/src/protocol.rs +++ b/iota-ipc/src/protocol.rs @@ -171,6 +171,52 @@ impl IpcRole { } impl LocalRequest { + pub const fn log_name(&self) -> &'static str { + match self { + Self::GetStatus => "get_status", + Self::ListTasks => "list_tasks", + Self::ListUsers => "list_users", + Self::ListTAuthApps => "list_tauth_apps", + Self::GetTAuthApp { .. } => "get_tauth_app", + Self::CreateTAuthApp { .. } => "create_tauth_app", + Self::ExportTAuthApp { .. } => "export_tauth_app", + Self::DeleteTAuthApp { .. } => "delete_tauth_app", + Self::CreateInvitation { .. } => "create_invitation", + Self::ListInvitations { .. } => "list_invitations", + Self::RevokeInvitation { .. } => "revoke_invitation", + Self::CreateUser { .. } => "create_user", + Self::InspectTuCredential { .. } => "inspect_tu_credential", + Self::AttachUserFromTu { .. } => "attach_user_from_tu", + Self::ReconcileUser { .. } => "reconcile_user", + Self::ForceDetachUser { .. } => "force_detach_user", + Self::ForgetReleasedUser { .. } => "forget_released_user", + Self::GetUserDiagnostics { .. } => "get_user_diagnostics", + Self::RevokeTAuthGrant { .. } => "revoke_tauth_grant", + Self::RevokeAllTAuthGrants { .. } => "revoke_all_tauth_grants", + Self::ExportUserCredential { .. } => "export_user_credential", + Self::PurgeUserData { .. } => "purge_user_data", + Self::ReleaseUser { .. } => "release_user", + Self::CompleteDeleteUser { .. } => "complete_delete_user", + Self::RemoveUser { .. } => "remove_user", + Self::ReconnectOmikron => "reconnect_omikron", + Self::RotateIotaIdentity => "rotate_iota_identity", + Self::RequestProcessExit { .. } => "request_process_exit", + Self::GetDaemonStatus => "get_daemon_status", + Self::RestartDaemon => "restart_daemon", + Self::StopDaemon => "stop_daemon", + Self::GetConfig => "get_config", + Self::SetConfig { .. } => "set_config", + Self::ReloadConfig => "reload_config", + Self::GetOmikronStatus => "get_omikron_status", + Self::ListComponents => "list_components", + Self::GetUser { .. } => "get_user", + Self::ImportUser { .. } => "import_user", + Self::GetLogs { .. } => "get_logs", + Self::CheckUpdate => "check_update", + Self::ListCommunities => "list_communities", + } + } + /// Return the minimum authenticated local role required to execute a /// request. New request variants must be assigned explicitly here. pub fn required_role(&self) -> IpcRole { @@ -620,6 +666,16 @@ mod error_tests { assert!(output.contains("")); } + #[test] + fn request_log_name_does_not_include_config_value() { + let request = super::LocalRequest::SetConfig { + key: "omikron_host".into(), + value: "private-config-value".into(), + }; + assert_eq!(request.log_name(), "set_config"); + assert!(!request.log_name().contains("private-config-value")); + } + #[test] fn user_summary_serializes_pending_operation_without_lifecycle_secrets() { let summary = UserSummary { diff --git a/iota-logger/src/lib.rs b/iota-logger/src/lib.rs index 9e5ea98..58f20ad 100644 --- a/iota-logger/src/lib.rs +++ b/iota-logger/src/lib.rs @@ -1,12 +1,17 @@ use std::{ fs::{self, OpenOptions}, io::Write, - sync::{OnceLock, atomic::Ordering, mpsc}, + path::{Path, PathBuf}, + sync::{ + OnceLock, + atomic::{AtomicU64, Ordering}, + mpsc::{self, RecvTimeoutError, TrySendError}, + }, thread, - time::{SystemTime, UNIX_EPOCH}, + time::{Duration, Instant, SystemTime, UNIX_EPOCH}, }; -use mtp::codec::{CommunicationValue, DataTypeId, DataValue, TypeMap, Version}; +use mtp::codec::{CommunicationType, CommunicationValue, DataTypeId, DataValue, TypeMap}; use ratatui::style::Color; use iota_state::{UNIQUE, UiLogEntry}; @@ -14,8 +19,31 @@ use tokio::sync::broadcast; pub mod language_creator; pub mod language_manager; -static LOGGER: OnceLock> = OnceLock::new(); +static LOGGER: OnceLock> = OnceLock::new(); static LOG_BROADCASTER: OnceLock> = OnceLock::new(); +static DROPPED_LOGS: AtomicU64 = AtomicU64::new(0); +const DEFAULT_LOGGER_QUEUE_CAPACITY: usize = 1024; +const DROPPED_LOG_REPORT_INTERVAL: Duration = Duration::from_secs(5); +const MAX_LOG_FILE_BYTES: u64 = 16 * 1024 * 1024; +const RETAINED_LOG_FILES: usize = 8; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum LogLevel { + Debug, + Info, + Warn, + Error, +} +impl LogLevel { + pub const fn as_str(self) -> &'static str { + match self { + Self::Debug => "DEBUG", + Self::Info => "INFO", + Self::Warn => "WARN", + Self::Error => "ERROR", + } + } +} #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] #[allow(unused)] @@ -29,6 +57,17 @@ pub enum PrintType { Command, } impl PrintType { + pub const fn as_str(self) -> &'static str { + match self { + Self::Call => "call", + Self::Client => "client", + Self::Iota => "iota", + Self::Omikron => "omikron", + Self::Omega => "omega", + Self::General => "general", + Self::Command => "command", + } + } pub fn prefix_color(self) -> Color { match self { PrintType::Call => Color::Magenta, @@ -47,6 +86,8 @@ struct LogMessage { prefix: String, kind: PrintType, is_error: bool, + level: LogLevel, + event: &'static str, translation_key: Option, format_args: Vec, message: Option, @@ -55,16 +96,23 @@ struct LogMessage { /* The logger owns file persistence while consumers receive rendered entries * through a process-local broadcast subscription. */ pub fn startup() { - startup_with_log_dir(Some( - iota_paths::IotaPaths::resolve(iota_paths::Scope::User) - .expect("resolve Iota user paths") - .log_dir, - )); + startup_with_log_dir_and_capacity( + Some( + iota_paths::IotaPaths::resolve(iota_paths::Scope::User) + .expect("resolve Iota user paths") + .log_dir, + ), + DEFAULT_LOGGER_QUEUE_CAPACITY, + ); } /// `None` keeps logging on stderr only (the systemd default). pub fn startup_with_log_dir(log_dir: Option) { - let (tx, rx) = mpsc::channel::(); + startup_with_log_dir_and_capacity(log_dir, DEFAULT_LOGGER_QUEUE_CAPACITY); +} + +pub fn startup_with_log_dir_and_capacity(log_dir: Option, queue_capacity: usize) { + let (tx, rx) = mpsc::sync_channel::(queue_capacity.max(1)); if LOGGER.set(tx).is_err() { return; } @@ -72,17 +120,45 @@ pub fn startup_with_log_dir(log_dir: Option) { let _ = LOG_BROADCASTER.set(broadcast_tx.clone()); thread::spawn(move || { - let mut file = log_dir.and_then(|log_dir| { - fs::create_dir_all(&log_dir).ok()?; - let start_ts = SystemTime::now().duration_since(UNIX_EPOCH).ok()?.as_secs(); - OpenOptions::new() - .create(true) - .append(true) - .open(log_dir.join(format!("log_{start_ts}.txt"))) - .ok() + let start_ts = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs(); + let mut sequence = 0; + let mut file = log_dir.as_ref().and_then(|log_dir| { + if let Err(error) = fs::create_dir_all(log_dir) { + eprintln!( + "Unable to create Iota log directory {}: {error}", + log_dir.display() + ); + return None; + } + prune_logs(log_dir); + match open_log(log_dir, start_ts, sequence) { + Ok(file) => Some(file), + Err(error) => { + eprintln!( + "Unable to open Iota log file in {}: {error}", + log_dir.display() + ); + None + } + } }); - - for msg in rx { + let mut last_drop_report = Instant::now(); + loop { + let msg = match rx.recv_timeout(DROPPED_LOG_REPORT_INTERVAL) { + Ok(msg) => msg, + Err(RecvTimeoutError::Timeout) => { + report_dropped_logs(&mut file, &broadcast_tx); + last_drop_report = Instant::now(); + continue; + } + Err(RecvTimeoutError::Disconnected) => { + report_dropped_logs(&mut file, &broadcast_tx); + break; + } + }; let resolved_message = if let Some(key) = msg.translation_key { let args: Vec<&str> = msg.format_args.iter().map(|s| s.as_str()).collect(); language_manager::format(&key, &args) @@ -99,27 +175,49 @@ pub fn startup_with_log_dir(log_dir: Option) { }; let line = format!( - "{} {} {}{}", - fixed_box(&msg.timestamp_ms.to_string(), 13), - timestamp, - prefix, - resolved_message + "{} level={} component={} direction={} sender=- event={} message={:?}", + msg.timestamp_ms, + msg.level.as_str(), + msg.kind.as_str(), + match msg.prefix.trim() { + ">" => "in", + "<" => "out", + ">>" => "internal", + _ => "local", + }, + msg.event, + format!("{prefix}{resolved_message}") ); if let Some(file) = file.as_mut() { - let _ = writeln!(file, "{}", line); + if let Err(error) = writeln!(file, "{line}") { + eprintln!("Unable to write Iota log file: {error}"); + } } - let _ = writeln!(std::io::stderr(), "{}", line); + let _ = writeln!(std::io::stderr(), "{timestamp} {line}"); + + let ui_message = if msg.event == "message" { + resolved_message.clone() + } else { + format!("event={} {}", msg.event, resolved_message) + }; let entry = UiLogEntry { timestamp_ms: msg.timestamp_ms, sender: format!("{:?}", msg.kind), - message: resolved_message, + message: ui_message, is_error: msg.is_error, }; let _ = broadcast_tx.send(entry); + if last_drop_report.elapsed() >= DROPPED_LOG_REPORT_INTERVAL { + report_dropped_logs(&mut file, &broadcast_tx); + last_drop_report = Instant::now(); + } + if let (Some(file), Some(log_dir)) = (file.as_mut(), log_dir.as_ref()) { + rotate_log_if_needed(file, log_dir, start_ts, &mut sequence); + } } }); } @@ -136,13 +234,110 @@ fn format_timestamp_inline(timestamp_ms: u128) -> String { format!("[{:02}:{:02}:{:02}]", hours, minutes, seconds) } -fn fixed_box(content: &str, width: usize) -> String { - let s: String = content.chars().take(width).collect(); - let len = s.chars().count(); - if len < width { - format!("[{}{}]", " ".repeat(width - len), s) - } else { - s +fn open_log(dir: &Path, start_ts: u64, sequence: u32) -> std::io::Result { + let mut options = OpenOptions::new(); + options.create(true).append(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + options.open(dir.join(format!("log_{start_ts}_{sequence}.txt"))) +} + +fn rotate_log_if_needed(file: &mut std::fs::File, dir: &Path, start_ts: u64, sequence: &mut u32) { + if !file + .metadata() + .is_ok_and(|meta| meta.len() >= MAX_LOG_FILE_BYTES) + { + return; + } + let next = sequence.saturating_add(1); + match open_log(dir, start_ts, next) { + Ok(new_file) => { + *file = new_file; + *sequence = next; + prune_logs(dir); + } + Err(error) => eprintln!("Unable to rotate Iota log file: {error}"), + } +} + +fn prune_logs(dir: &Path) { + let entries = match fs::read_dir(dir) { + Ok(entries) => entries, + Err(error) => { + eprintln!( + "Unable to enumerate Iota log directory {}: {error}", + dir.display() + ); + return; + } + }; + let mut logs = entries + .filter_map(Result::ok) + .filter(|entry| { + entry + .file_name() + .to_str() + .is_some_and(|name| name.starts_with("log_") && name.ends_with(".txt")) + }) + .collect::>(); + logs.sort_by_key(|entry| { + entry + .metadata() + .and_then(|meta| meta.modified()) + .unwrap_or(UNIX_EPOCH) + }); + let remove_count = logs.len().saturating_sub(RETAINED_LOG_FILES); + for entry in logs.into_iter().take(remove_count) { + if let Err(error) = fs::remove_file(entry.path()) { + eprintln!("Unable to remove old Iota log file: {error}"); + } + } +} + +fn report_dropped_logs( + file: &mut Option, + broadcaster: &broadcast::Sender, +) { + let dropped = DROPPED_LOGS.swap(0, Ordering::Relaxed); + if dropped == 0 { + return; + } + let timestamp_ms = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis(); + let message = format!("count={dropped}"); + let line = format!( + "{timestamp_ms} level=ERROR component=general direction=internal sender=- event=logger.dropped message={message:?}" + ); + if let Some(file) = file { + if let Err(error) = writeln!(file, "{line}") { + eprintln!("Unable to write Iota log file: {error}"); + } + } + eprintln!("{line}"); + let _ = broadcaster.send(UiLogEntry { + timestamp_ms, + sender: "General".into(), + message: format!("event=logger.dropped {message}"), + is_error: true, + }); +} + +fn enqueue(message: LogMessage) { + let Some(tx) = LOGGER.get() else { + return; + }; + UNIQUE.store(true, Ordering::Relaxed); + match tx.try_send(message) { + Ok(()) => {} + Err(TrySendError::Full(_)) => { + DROPPED_LOGS.fetch_add(1, Ordering::Relaxed); + } + Err(TrySendError::Disconnected(_)) => eprintln!("Iota logger thread has stopped"), } } @@ -153,39 +348,68 @@ pub fn log_internal_translated( key: &str, args: Vec, ) { - if let Some(tx) = LOGGER.get() { - UNIQUE.store(true, Ordering::Relaxed); - let _ = tx.send(LogMessage { - timestamp_ms: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis(), - prefix, - kind, - is_error, - translation_key: Some(key.to_string()), - format_args: args, - message: None, - }); - } + enqueue(LogMessage { + timestamp_ms: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis(), + prefix, + kind, + is_error, + level: if is_error { + LogLevel::Error + } else { + LogLevel::Info + }, + event: "message", + translation_key: Some(key.to_string()), + format_args: args, + message: None, + }); } pub fn log_internal(kind: PrintType, prefix: String, is_error: bool, message: String) { - if let Some(tx) = LOGGER.get() { - UNIQUE.store(true, Ordering::Relaxed); - let _ = tx.send(LogMessage { - timestamp_ms: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis(), - prefix, - kind, - is_error, - translation_key: None, - format_args: Vec::new(), - message: Some(message), - }); - } + log_event_internal( + kind, + if is_error { + LogLevel::Error + } else { + LogLevel::Info + }, + "message", + prefix, + message, + ); +} + +pub fn log_event_internal( + kind: PrintType, + level: LogLevel, + event: &'static str, + prefix: String, + message: String, +) { + enqueue(LogMessage { + timestamp_ms: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis(), + prefix, + kind, + is_error: level == LogLevel::Error, + level, + event, + translation_key: None, + format_args: Vec::new(), + message: Some(message), + }); +} + +#[macro_export] +macro_rules! log_event { + ($kind:expr, $level:expr, $event:expr, $($arg:tt)*) => { + $crate::log_event_internal($kind, $level, $event, String::new(), format!($($arg)*)) + }; } #[macro_export] @@ -307,10 +531,15 @@ pub fn log_cv_internal( ) { let formatted = format_cv(cv); - log_internal( + log_event_internal( print_type.unwrap_or(PrintType::General), + LogLevel::Debug, + if prefix.trim() == "<" { + "protocol.sent" + } else { + "protocol.received" + }, prefix.to_string(), - false, formatted, ); } @@ -334,136 +563,103 @@ pub fn format_cv(cv: &CommunicationValue) -> String { .map_or_else(|| "none".to_string(), |value| value.to_string()); parts.push(format!("{} (id={})", comm_type, id)); - let version = cv - .type_map() - .map(|type_map| type_map.version.clone()) - .unwrap_or_else(|| Version(3, 0)); - let formated_data = cv.data().map_or_else( - || "".to_string(), - |data| format_data_container(data.to_vec(), version), - ); - - parts.push(format!("{}", formated_data)); + if cv.is_type(CommunicationType::Relay) { + parts.push("".into()); + return parts.join(": "); + } + if let Some(data) = cv.data() { + let type_map = cv.type_map().cloned().unwrap_or_else(TypeMap::latest); + parts.push(format_data_container(data, &type_map)); + } parts.join(": ") } -fn format_data_container(data: Vec<(DataTypeId, DataValue)>, version: Version) -> String { - let parts: Vec = data - .into_iter() +fn format_data_container(data: &[(DataTypeId, DataValue)], type_map: &TypeMap) -> String { + data.iter() .map(|(key, value)| { - let key_str = key.to_string(); - - if is_secret_data_type(key) { - return format!("{}=", key_str); - } - + let name = type_map + .data_type_name(key.0) + .map(str::to_owned) + .unwrap_or_else(|| key.to_string()); match value { - DataValue::Str(s) => format!("{}=\"{}\"", key_str, abbreviate_string(&s)), - + DataValue::SignedNumber(value) + if matches!( + name.as_str(), + "UserId" + | "IotaId" + | "OmikronId" + | "InvitationId" + | "RelayMessageId" + | "VersionNumber" + | "Offset" + | "Amount" + ) => + { + format!("{name}={value}") + } + DataValue::UnsignedNumber(value) + if matches!( + name.as_str(), + "UserId" + | "IotaId" + | "OmikronId" + | "InvitationId" + | "RelayMessageId" + | "VersionNumber" + | "Offset" + | "Amount" + ) => + { + format!("{name}={value}") + } DataValue::Container(inner) => { - let inner_formatted = format_data_container(inner, version.clone()); - format!("{}={{ {} }}", key_str, inner_formatted) + format!("{name}={{ {} }}", format_data_container(inner, type_map)) } - - DataValue::Array(arr) => { - let arr_formatted = format_array(arr, version.clone()); - format!("{}=[{}]", key_str, arr_formatted) + DataValue::Array(values) => format!("{name}=", values.len()), + DataValue::Str(_) => format!("{name}="), + DataValue::Bytes(_) => format!("{name}="), + DataValue::Bool(_) | DataValue::BoolTrue | DataValue::BoolFalse => { + format!("{name}=") } - - DataValue::Bool(b) => format!("{}={}", key_str, b), - - DataValue::BoolTrue => format!("{}=true", key_str), - DataValue::BoolFalse => format!("{}=false", key_str), - - DataValue::SignedNumber(num) => format!("{}={}", key_str, num), - - _ => "".to_string(), + DataValue::SignedNumber(_) | DataValue::UnsignedNumber(_) => { + format!("{name}=") + } + _ => format!("{name}="), } }) - .collect(); - - parts.join(", ") -} - -fn is_secret_data_type(key: DataTypeId) -> bool { - matches!( - TypeMap::latest().data_type_name(key.0), - Some("InvitationToken" | "InvitationPassword" | "ResetToken" | "RegisterId" | "NewToken") - ) -} - -fn format_array(arr: Vec, version: Version) -> String { - let parts: Vec = arr - .into_iter() - .map(|value| match value { - DataValue::Str(s) => format!("\"{}\"", abbreviate_string(&s)), - - DataValue::Container(inner) => { - let inner_formatted = format_data_container(inner, version.clone()); - format!("{{ {} }}", inner_formatted) - } - - DataValue::Array(inner_arr) => { - let formatted = format_array(inner_arr, version.clone()); - format!("[{}]", formatted) - } - - DataValue::Bool(b) => b.to_string(), - - DataValue::BoolTrue => "true".to_string(), - DataValue::BoolFalse => "false".to_string(), - - DataValue::SignedNumber(num) => num.to_string(), - - _ => String::new(), - }) - .collect(); - - parts.join(", ") -} - -fn abbreviate_string(value: &str) -> String { - const EDGE_LENGTH: usize = 4; - - let chars: Vec = value.chars().collect(); - if chars.len() <= EDGE_LENGTH * 2 { - return value.to_string(); - } - - let prefix: String = chars.iter().take(EDGE_LENGTH).collect(); - let suffix: String = chars.iter().rev().take(EDGE_LENGTH).rev().collect(); - format!("{prefix}...{suffix}") + .collect::>() + .join(", ") } #[cfg(test)] mod tests { - use super::{abbreviate_string, format_cv}; - use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue, TypeMap}; + use super::format_cv; + use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; #[test] - fn abbreviates_only_strings_longer_than_eight_characters() { - assert_eq!(abbreviate_string("12345678"), "12345678"); - assert_eq!(abbreviate_string("123456789"), "1234...6789"); - assert_eq!(abbreviate_string("YWJjZGVmZ2hpag=="), "YWJj...ag=="); - } - - #[test] - fn redacts_invitation_and_account_credentials() { - let reset_token = DataType::ResetToken.try_to_id(&TypeMap::latest()).unwrap(); - let value = CommunicationValue::new(CommunicationType::RedeemUserInvitation) - .add_typed_default(DataType::InvitationToken, DataValue::Str("short123".into())) + fn protocol_values_are_metadata_only() { + let value = CommunicationValue::new(CommunicationType::Success) .add_typed_default( - DataType::Invitations, - DataValue::Array(vec![DataValue::Container(vec![( - reset_token, - DataValue::Str("reset-secret".into()), - )])]), - ); + DataType::SessionToken, + DataValue::Str("session-secret-value".into()), + ) + .add_typed_default( + DataType::CallToken, + DataValue::Str("livekit-secret-token".into()), + ) + .add_typed_default(DataType::Username, DataValue::Str("alice".into())) + .add_typed_default(DataType::UserId, DataValue::SignedNumber(172)); let formatted = format_cv(&value); - assert!(!formatted.contains("short123")); - assert!(!formatted.contains("reset-secret")); - assert_eq!(formatted.matches("").count(), 2); + for secret in ["session-secret-value", "livekit-secret-token", "alice"] { + assert!(!formatted.contains(secret)); + } + assert!(formatted.contains("SessionToken=")); + assert!(formatted.contains("UserId=172")); + assert!( + format_cv(&CommunicationValue::new(CommunicationType::Relay)) + .contains("") + ); } } diff --git a/iota-updater/src/lib.rs b/iota-updater/src/lib.rs index 5f649b8..d245a2b 100644 --- a/iota-updater/src/lib.rs +++ b/iota-updater/src/lib.rs @@ -30,6 +30,26 @@ const DEFAULT_ACTIVATION_TIMEOUT_SECONDS: u64 = 60; const DEFAULT_ACTIVATION_RETRY_MILLISECONDS: u64 = 250; const REQUIRED_HOST_ARTIFACTS: &str = include_str!("../artifacts.tsv"); +#[derive(Clone, Copy, Debug)] +pub enum UpdatePhase { + CandidateEvaluated, + ArtifactsStaged, + ActivationStarted, + DaemonRestarted, + ActivationCommitted, + RollbackStarted, + RollbackCompleted, +} + +pub trait UpdateObserver: Send + Sync { + fn event(&self, phase: UpdatePhase, product_version: Option<&str>); +} + +pub struct NoopUpdateObserver; +impl UpdateObserver for NoopUpdateObserver { + fn event(&self, _phase: UpdatePhase, _product_version: Option<&str>) {} +} + #[derive(Clone, Debug)] pub struct ActivationPolicy { pub timeout: Duration, @@ -176,16 +196,21 @@ pub async fn check_update() -> Result { } pub async fn apply_update() -> Result { + apply_update_observed(&NoopUpdateObserver).await +} + +pub async fn apply_update_observed(observer: &dyn UpdateObserver) -> Result { let paths = iota_paths::IotaPaths::resolve(iota_paths::Scope::System) .map_err(|error| anyhow::anyhow!(error))?; let policy = UpdatePolicy::from_environment()?; - apply_update_with(&paths, &SystemdSupervisor, &policy).await + apply_update_with_observer(&paths, &SystemdSupervisor, &policy, observer).await } -async fn apply_update_with( +async fn apply_update_with_observer( paths: &iota_paths::IotaPaths, supervisor: &impl DaemonSupervisor, policy: &UpdatePolicy, + observer: &dyn UpdateObserver, ) -> Result { let transaction = UpdateTransaction::from_paths(paths)?; let _lock = transaction.acquire()?; @@ -209,6 +234,10 @@ async fn apply_update_with( &policy.signing_key_id, Utc::now(), )?; + observer.event( + UpdatePhase::CandidateEvaluated, + Some(&release.manifest.product_version), + ); if decision == CandidateDecision::NoUpdate { return Ok(false); } @@ -236,6 +265,10 @@ async fn apply_update_with( serde_json::to_vec_pretty(&release.manifest)?, ) .context("write staged release manifest")?; + observer.event( + UpdatePhase::ArtifactsStaged, + Some(&release.manifest.product_version), + ); let daemon_was_active = supervisor.is_active()?; begin_activation_state( paths, @@ -244,6 +277,10 @@ async fn apply_update_with( &release.manifest, daemon_was_active, )?; + observer.event( + UpdatePhase::ActivationStarted, + Some(&release.manifest.product_version), + ); let activation = match transaction.activate(&release.manifest.product_version) { Ok(activation) => activation, Err(error) => { @@ -263,13 +300,25 @@ async fn apply_update_with( &policy.activation, ) .await?; + observer.event( + UpdatePhase::DaemonRestarted, + Some(&release.manifest.product_version), + ); } commit_update_state(paths, &mut state, &release.manifest)?; + observer.event( + UpdatePhase::ActivationCommitted, + Some(&release.manifest.product_version), + ); Ok::<_, anyhow::Error>(()) } .await; if let Err(update_error) = activation_result { + observer.event( + UpdatePhase::RollbackStarted, + Some(&release.manifest.product_version), + ); restore_previous_release( paths, &transaction, @@ -279,6 +328,10 @@ async fn apply_update_with( &policy.activation, ) .await?; + observer.event( + UpdatePhase::RollbackCompleted, + Some(&release.manifest.product_version), + ); record_failed_release(paths, &mut state, &release.manifest)?; return Err( update_error.context("candidate release failed runtime activation and was rolled back") diff --git a/iota-updater/src/main.rs b/iota-updater/src/main.rs index eaf02ba..958cc5e 100644 --- a/iota-updater/src/main.rs +++ b/iota-updater/src/main.rs @@ -6,7 +6,18 @@ async fn main() -> Result<()> { match command.as_str() { "check" => println!("{}", iota_updater::check_update().await?), "status" => println!("updater ready"), - "apply" => println!("{}", iota_updater::apply_update().await?), + "apply" => { + struct StderrObserver; + impl iota_updater::UpdateObserver for StderrObserver { + fn event(&self, phase: iota_updater::UpdatePhase, product_version: Option<&str>) { + eprintln!("update.phase phase={phase:?} product_version={product_version:?}"); + } + } + println!( + "{}", + iota_updater::apply_update_observed(&StderrObserver).await? + ); + } "rollback" => { let version = std::env::args() .nth(2) diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index 5b9a9d9..6a577f8 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -1,6 +1,6 @@ use base64::{Engine as _, engine::general_purpose::STANDARD}; use dashmap::{DashMap, DashSet}; -use iota_logger::{log, log_cv_in, log_cv_out, log_t}; +use iota_logger::{LogLevel, PrintType, log, log_cv_in, log_cv_out, log_event, log_t}; use iota_state::AppState; use iota_storage::util::config_util::{CONFIG, modify_config}; use iota_storage::util::relay_replay; @@ -3103,6 +3103,15 @@ impl OmikronConnection { let sender_guard = self.sender.read().await; if let Some(sender) = sender_guard.as_ref() { if !sender.is_open() { + log_event!( + PrintType::Omikron, + LogLevel::Warn, + "omikron.send_failed", + "connection_id={} message_id={:?} message_type={} reason=connection_closed", + self.connection_id, + cv.id(), + cv.get_type_name().unwrap_or("unknown") + ); drop(sender_guard); if let Some(sender) = self.sender.write().await.take() { sender.close().await; @@ -3121,6 +3130,16 @@ impl OmikronConnection { log_cv_out!(cv); if let Err(e) = sender_clone.send(cv).await { + log_event!( + PrintType::Omikron, + LogLevel::Error, + "omikron.send_failed", + "connection_id={} message_id={:?} message_type={} error={}", + self.connection_id, + cv.id(), + cv.get_type_name().unwrap_or("unknown"), + e + ); self.fail_all_waiting_tasks(format!( "Send failed: {} (connection_id={})", e, self.connection_id @@ -3131,6 +3150,15 @@ impl OmikronConnection { Ok(()) } else { + log_event!( + PrintType::Omikron, + LogLevel::Debug, + "omikron.send_skipped", + "connection_id={} message_id={:?} message_type={} reason=not_connected", + self.connection_id, + cv.id(), + cv.get_type_name().unwrap_or("unknown") + ); Err("not connected".to_string()) } }