[Add] Structured Iota event logging

This commit is contained in:
Alex-Emmet 2026-09-25 21:27:05 +02:00
commit 8d576df557
No known key found for this signature in database
8 changed files with 920 additions and 278 deletions

View file

@ -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<DataValue> = 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<SettingLocator, Comm
}
fn validate_setting_target(
cv: &CommunicationValue,
user_id: i64,
locator: &SettingLocator,
) -> 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),
}
}

View file

@ -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::<std::collections::HashMap<_, _>>(),
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::<Result<Vec<_>, iota_storage::storage_error::StorageError>>();
let Ok(users) = users else {
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(),

View file

@ -1,17 +1,28 @@
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<DaemonMessage>, buffer: Arc<Mutex<LogBuffer>>) {
let Some(mut logs) = subscribe() else {
let Some(logs) = subscribe() else {
return;
};
tokio::spawn(async move {
while let Ok(entry) = logs.recv().await {
tokio::spawn(forward_logs(logs, message_tx, buffer));
}
pub(crate) async fn forward_logs(
mut logs: broadcast::Receiver<UiLogEntry>,
message_tx: broadcast::Sender<DaemonMessage>,
buffer: Arc<Mutex<LogBuffer>>,
) {
loop {
match logs.recv().await {
Ok(entry) => {
let entry = LogEntry {
timestamp_ms: entry.timestamp_ms,
sender: entry.sender,
@ -23,5 +34,62 @@ pub fn spawn(message_tx: broadcast::Sender<DaemonMessage>, buffer: Arc<Mutex<Log
}
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")
);
}
}

View file

@ -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("<redacted>"));
}
#[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 {

View file

@ -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<mpsc::Sender<LogMessage>> = OnceLock::new();
static LOGGER: OnceLock<mpsc::SyncSender<LogMessage>> = OnceLock::new();
static LOG_BROADCASTER: OnceLock<broadcast::Sender<UiLogEntry>> = 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<String>,
format_args: Vec<String>,
message: Option<String>,
@ -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(
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<std::path::PathBuf>) {
let (tx, rx) = mpsc::channel::<LogMessage>();
startup_with_log_dir_and_capacity(log_dir, DEFAULT_LOGGER_QUEUE_CAPACITY);
}
pub fn startup_with_log_dir_and_capacity(log_dir: Option<PathBuf>, queue_capacity: usize) {
let (tx, rx) = mpsc::sync_channel::<LogMessage>(queue_capacity.max(1));
if LOGGER.set(tx).is_err() {
return;
}
@ -72,17 +120,45 @@ pub fn startup_with_log_dir(log_dir: Option<std::path::PathBuf>) {
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<std::path::PathBuf>) {
};
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<std::fs::File> {
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::<Vec<_>>();
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<std::fs::File>,
broadcaster: &broadcast::Sender<UiLogEntry>,
) {
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,9 +348,7 @@ pub fn log_internal_translated(
key: &str,
args: Vec<String>,
) {
if let Some(tx) = LOGGER.get() {
UNIQUE.store(true, Ordering::Relaxed);
let _ = tx.send(LogMessage {
enqueue(LogMessage {
timestamp_ms: SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
@ -163,29 +356,60 @@ pub fn log_internal_translated(
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 {
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,
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(
|| "<opaque payload>".to_string(),
|data| format_data_container(data.to_vec(), version),
);
parts.push(format!("{}", formated_data));
if cv.is_type(CommunicationType::Relay) {
parts.push("<opaque relay payload>".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<String> = 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!("{}=<redacted>", 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}=<array:{}>", values.len()),
DataValue::Str(_) => format!("{name}=<string>"),
DataValue::Bytes(_) => format!("{name}=<bytes>"),
DataValue::Bool(_) | DataValue::BoolTrue | DataValue::BoolFalse => {
format!("{name}=<bool>")
}
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}=<number>")
}
_ => format!("{name}=<value>"),
}
})
.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<DataValue>, version: Version) -> String {
let parts: Vec<String> = 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<char> = 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::<Vec<_>>()
.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("<redacted>").count(), 2);
for secret in ["session-secret-value", "livekit-secret-token", "alice"] {
assert!(!formatted.contains(secret));
}
assert!(formatted.contains("SessionToken=<string>"));
assert!(formatted.contains("UserId=172"));
assert!(
format_cv(&CommunicationValue::new(CommunicationType::Relay))
.contains("<opaque relay payload>")
);
}
}

View file

@ -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<bool> {
}
pub async fn apply_update() -> Result<bool> {
apply_update_observed(&NoopUpdateObserver).await
}
pub async fn apply_update_observed(observer: &dyn UpdateObserver) -> Result<bool> {
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<bool> {
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")

View file

@ -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)

View file

@ -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())
}
}