diff --git a/.forgejo/workflows/release.yml b/.forgejo/workflows/release.yml index 8b7600c..b0500ec 100644 --- a/.forgejo/workflows/release.yml +++ b/.forgejo/workflows/release.yml @@ -18,6 +18,18 @@ on: description: "Release description" required: true type: string + release_sequence: + description: "Monotonic sequence allocated for this update channel" + required: true + type: string + expires_at: + description: "Signed manifest expiry in RFC 3339 format" + required: true + type: string + +concurrency: + group: iota-release-${{ inputs.release_type }} + cancel-in-progress: false jobs: build: @@ -88,8 +100,11 @@ jobs: CHANNEL_TAG: ${{ steps.version.outputs.channel_tag }} CHANNEL_MANIFEST_NAME: ${{ steps.version.outputs.channel_manifest_name }} RELEASE_TYPE: ${{ inputs.release_type }} + RELEASE_SEQUENCE: ${{ inputs.release_sequence }} + EXPIRES_AT: ${{ inputs.expires_at }} SERVER_URL: ${{ forgejo.server_url }} REPO: ${{ forgejo.repository }} + TOKEN: ${{ forgejo.token }} IOTA_RELEASE_SIGNING_KEY: ${{ secrets.IOTA_RELEASE_SIGNING_KEY }} run: | set -eu @@ -103,17 +118,45 @@ jobs: UPDATE_BASE_URL="${SERVER_URL%/}/${REPO}/releases/download/${TAG}" UPDATE_CHANNEL_BASE_URL="${SERVER_URL%/}/${REPO}/releases/download/${CHANNEL_TAG}" + export UPDATE_CHANNEL_BASE_URL + nix-shell -p curl jq --run ' + set -eu + CHANNEL_STATUS="$(curl -L -sS -w "%{http_code}" -o previous-channel-manifest.json \ + -H "Authorization: token $TOKEN" \ + "$UPDATE_CHANNEL_BASE_URL/$CHANNEL_MANIFEST_NAME")" + case "$CHANNEL_STATUS" in + 200) + jq -e \ + --argjson proposed "$RELEASE_SEQUENCE" \ + --arg version "$TAG" \ + "(.release_sequence < \$proposed) or (.release_sequence == \$proposed and .product_version == \$version)" \ + previous-channel-manifest.json >/dev/null || { + echo "release sequence must increase, or identify the same release during a refresh" + exit 1 + } + ;; + 404) ;; + *) + echo "could not read current channel manifest: HTTP $CHANNEL_STATUS" + exit 1 + ;; + esac + ' PUBLISHED_AT="$(date -u +%Y-%m-%dT%H:%M:%SZ)" - export ARCH PUBLISHED_AT RELEASE_TYPE TAG UPDATE_BASE_URL UPDATE_MANIFEST_NAME - nix-shell -p coreutils jq --run 'bash scripts/build-update-manifest.sh dist/bin "$TAG" "$RELEASE_TYPE" "$PUBLISHED_AT" linux "$ARCH" "$UPDATE_BASE_URL" "dist/$UPDATE_MANIFEST_NAME"' + export ARCH EXPIRES_AT PUBLISHED_AT RELEASE_SEQUENCE RELEASE_TYPE TAG UPDATE_BASE_URL UPDATE_MANIFEST_NAME + nix-shell -p coreutils jq --run 'bash scripts/build-update-manifest.sh dist/bin "$TAG" "$RELEASE_TYPE" "$RELEASE_SEQUENCE" "$PUBLISHED_AT" "$EXPIRES_AT" linux "$ARCH" "$UPDATE_BASE_URL" "dist/$UPDATE_MANIFEST_NAME"' UPDATE_PUBLIC_KEY="$(result/bin/iota-release sign "dist/$UPDATE_MANIFEST_NAME" "dist/$UPDATE_MANIFEST_NAME.sig")" - bash scripts/build-release-bundle.sh \ + export UPDATE_PUBLIC_KEY + nix-shell -p jq zip --run 'bash scripts/build-release-bundle.sh \ dist/bin \ "$TAG" \ + "dist/$UPDATE_MANIFEST_NAME" \ "$UPDATE_CHANNEL_BASE_URL/$CHANNEL_MANIFEST_NAME" \ "$UPDATE_PUBLIC_KEY" \ "$UPDATE_CHANNEL_BASE_URL/$CHANNEL_MANIFEST_NAME.sig" \ - "dist/$ASSET_NAME" + "$RELEASE_TYPE" \ + primary \ + "dist/$ASSET_NAME"' result/bin/iota-bundle "dist/$ASSET_NAME" - name: Create release and upload release assets diff --git a/Cargo.lock b/Cargo.lock index 788506a..71512c0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2290,6 +2290,7 @@ name = "iota-updater" version = "0.1.0" dependencies = [ "anyhow", + "chrono", "ed25519-dalek 2.2.0", "fs2", "hex", diff --git a/iota-cli/src/ipc_client.rs b/iota-cli/src/ipc_client.rs index 22b81e0..99e8653 100644 --- a/iota-cli/src/ipc_client.rs +++ b/iota-cli/src/ipc_client.rs @@ -484,6 +484,38 @@ impl IpcClient { .join("\n") } } + ResponsePayload::InvitationCreated(invitation) => format!( + "Created {:?} invitation {}\nToken: {}\nExpires: {}{}", + invitation.authority, + invitation.invitation_id, + invitation.raw_token.0, + invitation.expires_at, + invitation + .short_url + .as_ref() + .map(|url| format!("\nLink: {url}")) + .unwrap_or_default(), + ), + ResponsePayload::Invitations(invitations) => { + if invitations.is_empty() { + "No invitations.".into() + } else { + invitations + .iter() + .map(|invitation| { + format!( + "{} ({:?}, {:?})", + invitation.invitation_id, invitation.authority, invitation.state + ) + }) + .collect::>() + .join("\n") + } + } + ResponsePayload::InvitationUpdated(invitation) => format!( + "Invitation {} is {:?}", + invitation.invitation_id, invitation.state + ), ResponsePayload::UserCreated { user_id, username } => { format!("Created user {} ({})", username, user_id) } @@ -588,6 +620,7 @@ impl IpcClient { iota_ipc::IpcErrorCode::OmikronUnavailable => { "Omikron is unavailable; try reconnecting." } + iota_ipc::IpcErrorCode::Unsupported => "The daemon does not implement this operation.", iota_ipc::IpcErrorCode::UnsupportedVersion => { "CLI and daemon versions are incompatible." } diff --git a/iota-cli/src/screens/users/add_flow.rs b/iota-cli/src/screens/users/add_flow.rs index 27d1f1f..1b697b0 100644 --- a/iota-cli/src/screens/users/add_flow.rs +++ b/iota-cli/src/screens/users/add_flow.rs @@ -25,6 +25,7 @@ pub enum AddUserPhase { pub struct AddUserFlow { pub phase: AddUserPhase, pub methods: MenuState, + pub invitation_authorities: MenuState, pub username: TextInput, pub import_path: TextInput, pub credential: Option, @@ -39,7 +40,7 @@ impl AddUserFlow { methods: MenuState::new(vec![ MenuItem { label: "Share an invitation".into(), - description: Some("Not supported by the connected daemon".into()), + description: Some("Create a single-use onboarding credential".into()), value: AddUserMethod::Invitation, enabled: invitation_supported, disabled_reason: Some("Not supported by the connected daemon".into()), @@ -60,6 +61,22 @@ impl AddUserFlow { .then(|| "TU inspection is not supported by the connected daemon".into()), }, ]), + invitation_authorities: MenuState::new(vec![ + MenuItem { + label: "Omega".into(), + description: Some("Central invitation authority".into()), + value: iota_ipc::InvitationAuthority::Omega, + enabled: true, + disabled_reason: None, + }, + MenuItem { + label: "This Iota".into(), + description: Some("Local invitation authority".into()), + value: iota_ipc::InvitationAuthority::Iota, + enabled: false, + disabled_reason: Some("Local invitations are not implemented yet".into()), + }, + ]), username: TextInput::new("Username"), import_path: TextInput::new("TU path"), credential: None, diff --git a/iota-cli/src/screens/users/mod.rs b/iota-cli/src/screens/users/mod.rs index e8e29f9..62afad6 100644 --- a/iota-cli/src/screens/users/mod.rs +++ b/iota-cli/src/screens/users/mod.rs @@ -399,6 +399,42 @@ impl UsersScreen { }), } } + fn start_invitation(&mut self, authority: iota_ipc::InvitationAuthority) -> InteractionResult { + self.pending = true; + self.message = Some("Creating invitation…".into()); + self.overlay = None; + let ipc = self.ipc.clone(); + InteractionResult::AppTask { + task: Box::pin(async move { + let result = match ipc + .send_request(iota_ipc::LocalRequest::CreateInvitation { + authority, + lifetime_seconds: 7 * 24 * 60 * 60, + password: None, + label: None, + }) + .await + { + Ok(iota_ipc::ResponseResult::Ok( + iota_ipc::ResponsePayload::InvitationCreated(invitation), + )) => Ok(format!( + "Invitation {} created. Token: {} Link: {}", + invitation.invitation_id, + invitation.raw_token.0, + invitation.short_url.as_deref().unwrap_or("unavailable") + )), + Ok(iota_ipc::ResponseResult::Error(error)) => { + Err(format!("Create invitation failed: {error}")) + } + Ok(_) => Err( + "Create invitation failed: daemon returned an unexpected response.".into(), + ), + Err(error) => Err(format!("Create invitation failed: {error}")), + }; + UiEvent::App(AppEvent::UserOperationFinished(result)) + }), + } + } fn refresh(&self) -> InteractionResult { let ipc = self.ipc.clone(); InteractionResult::AppTask { @@ -649,6 +685,33 @@ impl UsersScreen { inner, ); } + AddUserPhase::ConfigureInvitation => { + let lines = flow + .invitation_authorities + .items() + .iter() + .enumerate() + .map(|(index, item)| { + Line::from(format!( + "{} {}{}", + if flow.invitation_authorities.selected_index() == Some(index) { + ">" + } else { + " " + }, + item.label, + item.disabled_reason + .as_ref() + .map(|reason| format!(" ({reason})")) + .unwrap_or_default() + )) + }) + .chain(std::iter::once(Line::from( + "\nEnter creates a single-use invitation valid for one week.", + ))) + .collect::>(); + frame.render_widget(Paragraph::new(lines), inner); + } AddUserPhase::InspectingTu => { frame.render_widget(Paragraph::new("Inspecting credential…"), inner); } @@ -991,6 +1054,29 @@ impl Screen for UsersScreen { InteractionResult::Handled } }, + AddUserPhase::ConfigureInvitation => { + if flow.invitation_authorities.handle_key(key.code) { + return InteractionResult::Handled; + } + match key.code { + KeyCode::Esc => { + flow.phase = AddUserPhase::ChooseMethod; + InteractionResult::Handled + } + KeyCode::Enter => { + let authority = flow + .invitation_authorities + .selected_item() + .map(|item| item.value); + if let Some(authority) = authority { + self.start_invitation(authority) + } else { + InteractionResult::Handled + } + } + _ => InteractionResult::Handled, + } + } AddUserPhase::ReviewImport => match key.code { KeyCode::Esc => { flow.phase = AddUserPhase::ConfigureImport; diff --git a/iota-daemon-lib/src/command_router.rs b/iota-daemon-lib/src/command_router.rs index 6d0871c..f28d5d7 100644 --- a/iota-daemon-lib/src/command_router.rs +++ b/iota-daemon-lib/src/command_router.rs @@ -2,10 +2,11 @@ use crate::log_buffer::LogBuffer; use crate::{DaemonRuntime, DaemonServices}; use iota_ipc::{ CommunitySummary, ComponentStatusResponse, ConfigResponse, DaemonMessage, ExitIntent, - IpcErrorCode, LocalRequest, LogEntriesResponse, LogEntry, MAX_MESSAGE_SIZE, - OmikronStatusResponse, ReconcileAction, ResponseEnvelope, ResponsePayload, ResponseResult, - StatusResponse, TaskSummary, UpdateStatusResponse, UserDetailResponse, UserDiagnostics, - UserOperationKind, UserOperationSummary, UserReconcileResult, UserSummary, + InvitationAuthority, InvitationCreated, InvitationState, InvitationSummary, IpcErrorCode, + LocalRequest, LogEntriesResponse, LogEntry, MAX_MESSAGE_SIZE, OmikronStatusResponse, + ReconcileAction, ResponseEnvelope, ResponsePayload, ResponseResult, StatusResponse, + TaskSummary, UpdateStatusResponse, UserDetailResponse, UserDiagnostics, UserOperationKind, + UserOperationSummary, UserReconcileResult, UserSummary, }; use iota_logger::{log, log_command}; use iota_storage::users::pending_operations::{ @@ -13,6 +14,7 @@ use iota_storage::users::pending_operations::{ }; use iota_storage::users::user_manager; use iota_storage::util::config_util::{self}; +use iota_util::mtp_compat::OptionalDataValueExt; use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; use std::sync::{Arc, Mutex}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; @@ -149,6 +151,10 @@ impl CommandRouter { | LocalRequest::ReconcileUser { .. } | LocalRequest::ReleaseUser { .. } | LocalRequest::CompleteDeleteUser { .. } + | LocalRequest::CreateInvitation { + authority: InvitationAuthority::Omega, + .. + } ); if needs_omikron && !self.services.omikron.is_connected().await { return ResponseResult::Error( @@ -240,6 +246,269 @@ impl CommandRouter { }; ResponseResult::Ok(ResponsePayload::Users(users)) } + LocalRequest::CreateInvitation { + authority, + lifetime_seconds, + password, + label, + } => { + if authority == InvitationAuthority::Iota { + return ResponseResult::Error(IpcErrorCode::Unsupported); + } + if lifetime_seconds == 0 || lifetime_seconds > 7 * 24 * 60 * 60 { + return ResponseResult::Error(IpcErrorCode::InvalidRequest); + } + let password_protected = password.is_some(); + let mut request = CommunicationValue::new(CommunicationType::CreateUserInvitation) + .add_typed_default( + DataType::InvitationAuthority, + DataValue::Str("omega".into()), + ) + .add_typed_default( + DataType::InvitationLifetimeSeconds, + DataValue::SignedNumber(lifetime_seconds.into()), + ); + if let Some(password) = password { + request = request.add_typed_default( + DataType::InvitationPassword, + DataValue::Str(password.0), + ); + } + if let Some(label) = label.as_ref() { + request = request.add_typed_default( + DataType::InvitationLabel, + DataValue::Str(label.clone()), + ); + } + let response = match self + .services + .omikron + .await_response(&request, Duration::from_secs(20)) + .await + { + Ok(response) => response, + Err(omikron_connector::OmikronError::Timeout(_)) => { + return ResponseResult::Error(IpcErrorCode::Timeout); + } + Err(_) => return ResponseResult::Error(IpcErrorCode::OmikronUnavailable), + }; + let invitation_id = response + .get_data(DataType::InvitationId) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()); + let raw_token = response + .get_data(DataType::InvitationToken) + .as_str() + .map(str::to_owned); + let expires_at = response + .get_data(DataType::InvitationExpiresAt) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()); + let created_at = response + .get_data(DataType::InvitationCreatedAt) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()); + let remote_revision = response + .get_data(DataType::InvitationRevision) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()) + .unwrap_or(1); + let short_url = response + .get_data(DataType::Link) + .as_str() + .map(str::to_owned); + let Some((invitation_id, raw_token, created_at, expires_at)) = invitation_id + .zip(raw_token) + .zip(created_at) + .zip(expires_at) + .map(|(((id, token), created), expires)| (id, token, created, expires)) + else { + return ResponseResult::Error(IpcErrorCode::InternalFailure); + }; + let summary = iota_storage::users::invitations::InvitationSummary { + invitation_id, + authority: iota_storage::users::invitations::InvitationAuthority::Omega, + label, + password_protected, + created_at, + expires_at: Some(expires_at), + state: iota_storage::users::invitations::InvitationState::Pending, + remote_revision, + redeemed_user_id: None, + redeemed_at: None, + revoked_at: None, + pending_action: None, + pending_action_at: None, + last_synced_at: Some(now_millis()), + local_provisioned_user_id: None, + local_provisioned_at: None, + }; + if iota_storage::users::invitations::insert(&summary, None, now_millis()).is_err() { + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } + ResponseResult::Ok(ResponsePayload::InvitationCreated(InvitationCreated { + authority, + invitation_id, + raw_token: iota_ipc::SecretString(raw_token), + expires_at, + short_url, + })) + } + LocalRequest::ListInvitations { authority } => { + if authority != Some(InvitationAuthority::Iota) + && self.services.omikron.is_connected().await + { + match self.services.omikron.sync_omega_invitations().await { + Ok(()) => {} + Err(omikron_connector::OmikronError::Storage(_)) => { + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } + Err(omikron_connector::OmikronError::Timeout(_)) => { + return ResponseResult::Error(IpcErrorCode::Timeout); + } + Err(omikron_connector::OmikronError::Disconnected(_)) => { + return ResponseResult::Error(IpcErrorCode::OmikronUnavailable); + } + Err(_) => return ResponseResult::Error(IpcErrorCode::InternalFailure), + } + } + if iota_storage::users::invitations::expire_pending(now_millis()).is_err() { + return ResponseResult::Error(IpcErrorCode::StorageFailure); + } + let invitations = match iota_storage::users::invitations::list() { + Ok(invitations) => invitations, + Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + }; + let mut summaries = Vec::new(); + for invitation in invitations { + let invitation_authority = match invitation.authority { + iota_storage::users::invitations::InvitationAuthority::Omega => { + InvitationAuthority::Omega + } + iota_storage::users::invitations::InvitationAuthority::Iota => { + InvitationAuthority::Iota + } + }; + if authority.is_some_and(|filter| filter != invitation_authority) { + continue; + } + let display_state = if invitation.pending_action + == Some(iota_storage::users::invitations::PendingAction::Revoke) + { + InvitationState::RevocationPending + } else if invitation.authority + == iota_storage::users::invitations::InvitationAuthority::Omega + && invitation.state + == iota_storage::users::invitations::InvitationState::Pending + && invitation + .expires_at + .is_some_and(|expires_at| expires_at <= now_millis()) + { + InvitationState::Expired + } else { + match invitation.state { + iota_storage::users::invitations::InvitationState::Pending => { + InvitationState::Pending + } + iota_storage::users::invitations::InvitationState::Provisioning => { + InvitationState::Provisioning + } + iota_storage::users::invitations::InvitationState::Redeemed => { + InvitationState::Redeemed + } + iota_storage::users::invitations::InvitationState::Revoked => { + InvitationState::Revoked + } + iota_storage::users::invitations::InvitationState::Expired => { + InvitationState::Expired + } + } + }; + summaries.push(InvitationSummary { + authority: invitation_authority, + invitation_id: invitation.invitation_id, + label: invitation.label, + password_protected: invitation.password_protected, + created_at: invitation.created_at, + expires_at: invitation.expires_at, + state: display_state, + redeemed_user_id: invitation.redeemed_user_id, + local_provisioned_user_id: invitation.local_provisioned_user_id, + local_provisioned_at: invitation.local_provisioned_at, + }); + } + ResponseResult::Ok(ResponsePayload::Invitations(summaries)) + } + LocalRequest::RevokeInvitation { + authority, + invitation_id, + } => { + if authority == InvitationAuthority::Iota { + return ResponseResult::Error(IpcErrorCode::Unsupported); + } + let changed = match iota_storage::users::invitations::mark_revoke_pending( + invitation_id, + now_millis(), + ) { + Ok(changed) => changed, + Err(_) => return ResponseResult::Error(IpcErrorCode::StorageFailure), + }; + if !changed { + return ResponseResult::Error(IpcErrorCode::Conflict); + } + if self.services.omikron.is_connected().await { + self.services + .omikron + .flush_pending_invitation_actions() + .await; + } + let invitation = + match iota_storage::users::invitations::list() + .ok() + .and_then(|items| { + items + .into_iter() + .find(|item| item.invitation_id == invitation_id) + }) { + Some(invitation) => invitation, + None => return ResponseResult::Error(IpcErrorCode::StorageFailure), + }; + let state = if invitation.pending_action + == Some(iota_storage::users::invitations::PendingAction::Revoke) + { + InvitationState::RevocationPending + } else { + match invitation.state { + iota_storage::users::invitations::InvitationState::Pending => { + InvitationState::Pending + } + iota_storage::users::invitations::InvitationState::Provisioning => { + InvitationState::Provisioning + } + iota_storage::users::invitations::InvitationState::Redeemed => { + InvitationState::Redeemed + } + iota_storage::users::invitations::InvitationState::Revoked => { + InvitationState::Revoked + } + iota_storage::users::invitations::InvitationState::Expired => { + InvitationState::Expired + } + } + }; + ResponseResult::Ok(ResponsePayload::InvitationUpdated(InvitationSummary { + authority: InvitationAuthority::Omega, + invitation_id, + label: invitation.label, + password_protected: invitation.password_protected, + created_at: invitation.created_at, + expires_at: invitation.expires_at, + state, + redeemed_user_id: invitation.redeemed_user_id, + local_provisioned_user_id: invitation.local_provisioned_user_id, + local_provisioned_at: invitation.local_provisioned_at, + })) + } LocalRequest::CreateUser { username } => { match omikron_connector::user_ops::create_user( self.services.omikron.as_ref(), diff --git a/iota-daemon-lib/src/ipc_server.rs b/iota-daemon-lib/src/ipc_server.rs index 980c739..1e597c6 100644 --- a/iota-daemon-lib/src/ipc_server.rs +++ b/iota-daemon-lib/src/ipc_server.rs @@ -22,8 +22,6 @@ use uuid::Uuid; const CLIENT_CHANNEL_SIZE: usize = 256; const MAX_CONFIGURED_IPC_CLIENTS: usize = 4096; -/// Maximum handshake retries before giving up. -const MAX_HANDSHAKE_RETRIES: u32 = 1; const CLIENT_IO_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15); /// Minimum metric subscription interval to prevent excessive update rates. @@ -368,73 +366,67 @@ async fn handle_client( eprintln!("IPC handshake started (pid={}, uid={})", peer.pid, peer.uid); // --- Handshake --- - let mut negotiated_version: Option = None; - for _ in 0..MAX_HANDSHAKE_RETRIES { - match timeout( - std::time::Duration::from_secs(15), - read_msg::<_, ClientMessage>(&mut reader), - ) - .await - { - Err(_) => { + let negotiated_version = match timeout( + std::time::Duration::from_secs(15), + read_msg::<_, ClientMessage>(&mut reader), + ) + .await + { + Err(_) => { + return Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + "IPC Hello timed out", + )); + } + Ok(result) => match result { + Ok(ClientMessage::Hello { supported_versions }) => { + let version = supported_versions + .iter() + .copied() + .filter(|v| *v >= MIN_PROTOCOL_VERSION && *v <= PROTOCOL_VERSION) + .max() + .ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::Unsupported, + "No compatible IPC protocol version", + ) + })?; + let ack = DaemonMessage::HelloAck(HelloAck { + protocol_version: version, + daemon_version: env!("CARGO_PKG_VERSION").to_string(), + instance_id: instance_id.clone(), + startup_phase: runtime.current_startup_phase().into(), + capabilities: vec![ + "commands".into(), + "metrics".into(), + "logs".into(), + "user_management_v2".into(), + "tu_inspection_v1".into(), + "credential_export_v1".into(), + "user_invitations_v1".into(), + ], + lifecycle: *runtime.lifecycle.borrow(), + health: runtime.overall_health(), + deployment_mode: from_environment().mode, + supervisor: from_environment().supervisor, + }); + write_client_message(&mut writer, &ack).await?; + eprintln!( + "IPC handshake acknowledged (pid={}, uid={})", + peer.pid, peer.uid + ); + version + } + Ok(_) => { + // Unexpected first message, send an error and close. return Err(std::io::Error::new( - std::io::ErrorKind::TimedOut, - "IPC Hello timed out", + std::io::ErrorKind::InvalidData, + "Expected Hello as first message", )); } - Ok(result) => match result { - Ok(ClientMessage::Hello { supported_versions }) => { - let version = supported_versions - .iter() - .copied() - .filter(|v| *v >= MIN_PROTOCOL_VERSION && *v <= PROTOCOL_VERSION) - .max() - .ok_or_else(|| { - std::io::Error::new( - std::io::ErrorKind::Unsupported, - "No compatible IPC protocol version", - ) - })?; - negotiated_version = Some(version); - let ack = DaemonMessage::HelloAck(HelloAck { - protocol_version: version, - daemon_version: env!("CARGO_PKG_VERSION").to_string(), - instance_id: instance_id.clone(), - startup_phase: runtime.current_startup_phase().into(), - capabilities: vec![ - "commands".into(), - "metrics".into(), - "logs".into(), - "user_management_v2".into(), - "tu_inspection_v1".into(), - "credential_export_v1".into(), - ], - lifecycle: *runtime.lifecycle.borrow(), - health: runtime.overall_health(), - deployment_mode: from_environment().mode, - supervisor: from_environment().supervisor, - }); - write_client_message(&mut writer, &ack).await?; - eprintln!( - "IPC handshake acknowledged (pid={}, uid={})", - peer.pid, peer.uid - ); - break; - } - Ok(_) => { - // Unexpected first message, send an error and close. - return Err(std::io::Error::new( - std::io::ErrorKind::InvalidData, - "Expected Hello as first message", - )); - } - Err(e) => return Err(e), - }, - } - } - let negotiated_version = negotiated_version.ok_or_else(|| { - std::io::Error::new(std::io::ErrorKind::Other, "Handshake failed after retries") - })?; + Err(e) => return Err(e), + }, + }; log!("IPC client connected (pid={}, uid={})", peer.pid, peer.uid); diff --git a/iota-installer/src/lib.rs b/iota-installer/src/lib.rs index c9fa167..282953b 100644 --- a/iota-installer/src/lib.rs +++ b/iota-installer/src/lib.rs @@ -183,6 +183,41 @@ fn product_version(staging: &Path) -> Result { .with_context(|| format!("read bundle manifest {}", manifest_path.display()))?; let manifest: serde_json::Value = serde_json::from_str(&contents) .with_context(|| format!("parse bundle manifest {}", manifest_path.display()))?; + for name in [ + "product_version", + "channel", + "published_at", + "expires_at", + "release_signing_key_id", + ] { + manifest + .get(name) + .and_then(|value| value.as_str()) + .filter(|value| !value.is_empty()) + .with_context(|| format!("bundle manifest {name} must be a non-empty string"))?; + } + for name in [ + "release_sequence", + "minimum_data_schema", + "supported_ipc_min", + "supported_ipc_max", + ] { + let value = manifest + .get(name) + .and_then(|value| value.as_u64()) + .with_context(|| format!("bundle manifest {name} must be an unsigned integer"))?; + if name == "release_sequence" && value == 0 { + bail!("bundle manifest release_sequence must be greater than zero"); + } + } + manifest + .get("artifacts") + .and_then(|value| value.as_array()) + .context("bundle manifest artifacts must be an array")?; + manifest + .get("rollback_compatible") + .and_then(|value| value.as_bool()) + .context("bundle manifest rollback_compatible must be a boolean")?; let version = manifest .get("product_version") .and_then(|value| value.as_str()) @@ -218,6 +253,8 @@ fn validate_update_environment(staging: &Path) -> Result<()> { "IOTA_UPDATE_MANIFEST", "IOTA_UPDATE_PUBLIC_KEY", "IOTA_UPDATE_SIGNATURE", + "IOTA_UPDATE_CHANNEL", + "IOTA_UPDATE_SIGNING_KEY_ID", ] { let value = entries .get(name) @@ -227,7 +264,7 @@ fn validate_update_environment(staging: &Path) -> Result<()> { bail!("updater environment variable {name} contains whitespace"); } } - if entries.len() != 3 { + if entries.len() != 5 { bail!("updater environment contains unexpected variables"); } let public_key = entries @@ -240,6 +277,27 @@ fn validate_update_environment(staging: &Path) -> Result<()> { { bail!("IOTA_UPDATE_PUBLIC_KEY must be 32-byte hex"); } + let manifest: serde_json::Value = serde_json::from_slice( + &fs::read(staging.join("manifest.json")).context("read bundle manifest")?, + ) + .context("parse bundle manifest")?; + for (environment_name, manifest_name) in [ + ("IOTA_UPDATE_CHANNEL", "channel"), + ("IOTA_UPDATE_SIGNING_KEY_ID", "release_signing_key_id"), + ] { + let environment_value = entries + .get(environment_name) + .with_context(|| format!("updater environment is missing {environment_name}"))?; + let manifest_value = manifest + .get(manifest_name) + .and_then(|value| value.as_str()) + .with_context(|| format!("bundle manifest is missing {manifest_name}"))?; + if *environment_value != manifest_value { + bail!( + "updater environment {environment_name} does not match bundle manifest {manifest_name}" + ); + } + } Ok(()) } @@ -279,11 +337,26 @@ mod tests { .start_file(name, SimpleFileOptions::default()) .unwrap(); if name == "manifest.json" { - write!(archive, "{{\"product_version\":\"{product_version}\"}}").unwrap(); + let manifest = serde_json::json!({ + "product_version": product_version, + "channel": "stable", + "release_sequence": 1, + "published_at": "2026-09-10T00:00:00Z", + "expires_at": "2026-10-10T00:00:00Z", + "minimum_data_schema": 1, + "supported_ipc_min": 2, + "supported_ipc_max": 4, + "artifacts": [], + "release_signing_key_id": "primary", + "rollback_compatible": true + }); + archive + .write_all(&serde_json::to_vec(&manifest).unwrap()) + .unwrap(); } else if name == "systemd/update.env" { archive .write_all( - b"IOTA_UPDATE_MANIFEST=https://example.invalid/manifest.json\nIOTA_UPDATE_PUBLIC_KEY=0707070707070707070707070707070707070707070707070707070707070707\nIOTA_UPDATE_SIGNATURE=https://example.invalid/manifest.json.sig\n", + b"IOTA_UPDATE_MANIFEST=https://example.invalid/manifest.json\nIOTA_UPDATE_PUBLIC_KEY=0707070707070707070707070707070707070707070707070707070707070707\nIOTA_UPDATE_SIGNATURE=https://example.invalid/manifest.json.sig\nIOTA_UPDATE_CHANNEL=stable\nIOTA_UPDATE_SIGNING_KEY_ID=primary\n", ) .unwrap(); } else { @@ -354,4 +427,18 @@ mod tests { assert!(error.to_string().contains("IOTA_UPDATE_PUBLIC_KEY")); } + + #[test] + fn rejects_legacy_bundle_manifest_without_sequence_baseline() { + let directory = tempfile::tempdir().unwrap(); + fs::write( + directory.path().join("manifest.json"), + br#"{"product_version":"1.0.0"}"#, + ) + .unwrap(); + + let error = product_version(directory.path()).unwrap_err(); + + assert!(error.to_string().contains("channel")); + } } diff --git a/iota-ipc/src/lib.rs b/iota-ipc/src/lib.rs index 64603e6..f8c251f 100644 --- a/iota-ipc/src/lib.rs +++ b/iota-ipc/src/lib.rs @@ -5,8 +5,9 @@ pub mod transport; pub use protocol::{ ClientMessage, CommunitySummary, ComponentHealth, ComponentId, ComponentStatusResponse, ConfigResponse, ConnectionStatus, CredentialStatus, DaemonMessage, DaemonStatusResponse, - DeploymentMode, ExitIntent, HealthStatus, HelloAck, IpcErrorCode, IpcRole, LifecycleEvent, - LifecyclePhase, LocalRequest, LocalUserState, LogEntriesResponse, LogEntry, MetricSample, + DeploymentMode, ExitIntent, HealthStatus, HelloAck, InvitationAuthority, InvitationCreated, + InvitationState, InvitationSummary, IpcErrorCode, IpcRole, LifecycleEvent, LifecyclePhase, + LocalRequest, LocalUserState, LogEntriesResponse, LogEntry, MetricSample, OmikronStatusResponse, ReconcileAction, RequestEnvelope, ResponseEnvelope, ResponsePayload, ResponseResult, SecretString, StartupPhase, StateSnapshot, StatusResponse, SupervisorKind, TaskSummary, TuCredentialPreview, UpdateStatusResponse, UserDetailResponse, UserDiagnostics, @@ -15,6 +16,6 @@ pub use protocol::{ pub use transport::{MAX_MESSAGE_SIZE, read_msg, write_msg}; /// Current IPC protocol version. -pub const PROTOCOL_VERSION: u16 = 4; +pub const PROTOCOL_VERSION: u16 = 5; /// Minimum protocol version this daemon understands. pub const MIN_PROTOCOL_VERSION: u16 = 2; diff --git a/iota-ipc/src/protocol.rs b/iota-ipc/src/protocol.rs index f5051c6..8d3b40e 100644 --- a/iota-ipc/src/protocol.rs +++ b/iota-ipc/src/protocol.rs @@ -45,6 +45,19 @@ pub enum LocalRequest { GetStatus, ListTasks, ListUsers, + CreateInvitation { + authority: InvitationAuthority, + lifetime_seconds: u64, + password: Option, + label: Option, + }, + ListInvitations { + authority: Option, + }, + RevokeInvitation { + authority: InvitationAuthority, + invitation_id: i64, + }, CreateUser { username: String, }, @@ -148,6 +161,7 @@ impl LocalRequest { Self::GetStatus | Self::ListTasks | Self::ListUsers + | Self::ListInvitations { .. } | Self::GetDaemonStatus | Self::GetOmikronStatus | Self::ListComponents @@ -160,6 +174,8 @@ impl LocalRequest { Self::ReconnectOmikron | Self::ReloadConfig => IpcRole::Operate, Self::CreateUser { .. } + | Self::CreateInvitation { .. } + | Self::RevokeInvitation { .. } | Self::InspectTuCredential { .. } | Self::AttachUserFromTu { .. } | Self::ReconcileUser { .. } @@ -249,6 +265,9 @@ pub enum ResponsePayload { Status(StatusResponse), Tasks(Vec), Users(Vec), + InvitationCreated(InvitationCreated), + Invitations(Vec), + InvitationUpdated(InvitationSummary), UserCreated { user_id: i64, username: String, @@ -327,6 +346,49 @@ pub struct CommunitySummary { pub title: String, } +#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum InvitationAuthority { + Omega, + Iota, +} + +#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum InvitationState { + Pending, + RevocationPending, + Provisioning, + Redeemed, + Revoked, + Expired, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct InvitationCreated { + pub authority: InvitationAuthority, + pub invitation_id: i64, + pub raw_token: SecretString, + pub expires_at: i64, + pub short_url: Option, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct InvitationSummary { + pub authority: InvitationAuthority, + pub invitation_id: i64, + pub label: Option, + pub password_protected: bool, + pub created_at: i64, + pub expires_at: Option, + pub state: InvitationState, + pub redeemed_user_id: Option, + #[serde(default)] + pub local_provisioned_user_id: Option, + #[serde(default)] + pub local_provisioned_at: Option, +} + #[derive(Clone, Debug, Deserialize, Serialize)] pub struct StatusResponse { pub phase: String, @@ -427,6 +489,7 @@ pub enum IpcErrorCode { Conflict, StorageFailure, OmikronUnavailable, + Unsupported, UnsupportedVersion, NotReady, Disconnected, @@ -444,6 +507,7 @@ impl std::fmt::Display for IpcErrorCode { Self::Conflict => "the request conflicts with current daemon state", Self::StorageFailure => "the daemon could not access local storage", Self::OmikronUnavailable => "Omikron is unavailable", + Self::Unsupported => "the requested operation is not implemented", Self::UnsupportedVersion => "the client and daemon protocol versions are incompatible", Self::NotReady => "the daemon is not ready yet", Self::Disconnected => "the daemon connection was lost", diff --git a/iota-logger/src/lib.rs b/iota-logger/src/lib.rs index d4b3575..9e5ea98 100644 --- a/iota-logger/src/lib.rs +++ b/iota-logger/src/lib.rs @@ -6,7 +6,7 @@ use std::{ time::{SystemTime, UNIX_EPOCH}, }; -use mtp::codec::{CommunicationValue, DataTypeId, DataValue, Version}; +use mtp::codec::{CommunicationValue, DataTypeId, DataValue, TypeMap, Version}; use ratatui::style::Color; use iota_state::{UNIQUE, UiLogEntry}; @@ -354,6 +354,10 @@ fn format_data_container(data: Vec<(DataTypeId, DataValue)>, version: Version) - .map(|(key, value)| { let key_str = key.to_string(); + if is_secret_data_type(key) { + return format!("{}=", key_str); + } + match value { DataValue::Str(s) => format!("{}=\"{}\"", key_str, abbreviate_string(&s)), @@ -382,6 +386,13 @@ fn format_data_container(data: Vec<(DataTypeId, DataValue)>, version: Version) - 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() @@ -427,7 +438,8 @@ fn abbreviate_string(value: &str) -> String { #[cfg(test)] mod tests { - use super::abbreviate_string; + use super::{abbreviate_string, format_cv}; + use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue, TypeMap}; #[test] fn abbreviates_only_strings_longer_than_eight_characters() { @@ -435,6 +447,24 @@ mod tests { 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())) + .add_typed_default( + DataType::Invitations, + DataValue::Array(vec![DataValue::Container(vec![( + reset_token, + DataValue::Str("reset-secret".into()), + )])]), + ); + let formatted = format_cv(&value); + assert!(!formatted.contains("short123")); + assert!(!formatted.contains("reset-secret")); + assert_eq!(formatted.matches("").count(), 2); + } } #[macro_export] diff --git a/iota-storage/src/users/invitations.rs b/iota-storage/src/users/invitations.rs index 48f8b1d..b4b80d6 100644 --- a/iota-storage/src/users/invitations.rs +++ b/iota-storage/src/users/invitations.rs @@ -1,9 +1,58 @@ -use crate::{storage_error::StorageError, util::db}; -use rusqlite::params; +use crate::{storage_error::StorageError, users::user_profile::UserProfile, util::db}; +use rusqlite::{OptionalExtension, params}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum InvitationAuthority { + Omega, + Iota, +} + +impl InvitationAuthority { + pub fn as_str(self) -> &'static str { + match self { + Self::Omega => "omega", + Self::Iota => "iota", + } + } + + pub fn parse(value: &str) -> Result { + match value { + "omega" => Ok(Self::Omega), + "iota" => Ok(Self::Iota), + _ => Err(StorageError::Other("unknown invitation authority".into())), + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum PendingAction { + Revoke, +} + +impl PendingAction { + fn parse(value: &str) -> Result { + match value { + "revoke" => Ok(Self::Revoke), + _ => Err(StorageError::Other( + "unknown pending invitation action".into(), + )), + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ProvisioningResult { + Created, + AlreadyApplied, + RevocationPending, + MissingInvitation, + Conflict, +} #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum InvitationState { Pending, + Provisioning, Redeemed, Revoked, Expired, @@ -13,6 +62,7 @@ impl InvitationState { fn as_str(self) -> &'static str { match self { Self::Pending => "pending", + Self::Provisioning => "provisioning", Self::Redeemed => "redeemed", Self::Revoked => "revoked", Self::Expired => "expired", @@ -22,6 +72,7 @@ impl InvitationState { fn parse(value: &str) -> Result { match value { "pending" => Ok(Self::Pending), + "provisioning" => Ok(Self::Provisioning), "redeemed" => Ok(Self::Redeemed), "revoked" => Ok(Self::Revoked), "expired" => Ok(Self::Expired), @@ -32,52 +83,361 @@ impl InvitationState { #[derive(Clone, Debug, PartialEq, Eq)] pub struct InvitationSummary { - pub invitation_id: String, + pub invitation_id: i64, + pub authority: InvitationAuthority, + pub label: Option, + pub password_protected: bool, pub created_at: i64, pub expires_at: Option, pub state: InvitationState, + pub remote_revision: i64, pub redeemed_user_id: Option, pub redeemed_at: Option, pub revoked_at: Option, + pub pending_action: Option, + pub pending_action_at: Option, + pub last_synced_at: Option, + pub local_provisioned_user_id: Option, + pub local_provisioned_at: Option, } -pub fn insert(summary: &InvitationSummary, token_hash: &[u8]) -> Result<(), StorageError> { +pub fn insert( + summary: &InvitationSummary, + token_hash: Option<&[u8]>, + synced_at: i64, +) -> Result<(), StorageError> { + db::with_immediate_transaction(|tx| insert_in_tx(tx, summary, token_hash, synced_at)) +} + +pub fn merge_omega_snapshot( + invitations: &[InvitationSummary], + synced_at: i64, +) -> Result<(), StorageError> { db::with_immediate_transaction(|tx| { - tx.execute( - "INSERT INTO user_invitations (invitation_id, token_hash, created_at, expires_at, state, redeemed_user_id, redeemed_at, revoked_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)", - params![summary.invitation_id, token_hash, summary.created_at, summary.expires_at, summary.state.as_str(), summary.redeemed_user_id, summary.redeemed_at, summary.revoked_at], - )?; + for invitation in invitations { + insert_in_tx(tx, invitation, None, synced_at)?; + } Ok(()) }) } +fn insert_in_tx( + tx: &rusqlite::Transaction<'_>, + summary: &InvitationSummary, + token_hash: Option<&[u8]>, + synced_at: i64, +) -> Result<(), StorageError> { + tx.execute( + "INSERT INTO user_invitations (invitation_id, authority, token_hash, label, password_protected, created_at, expires_at, state, redeemed_user_id, redeemed_at, revoked_at, pending_action, pending_action_at, remote_revision, last_synced_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, NULL, ?12, ?13) ON CONFLICT(invitation_id) DO UPDATE SET authority = excluded.authority, token_hash = excluded.token_hash, label = excluded.label, password_protected = excluded.password_protected, expires_at = excluded.expires_at, state = excluded.state, redeemed_user_id = excluded.redeemed_user_id, redeemed_at = excluded.redeemed_at, revoked_at = excluded.revoked_at, remote_revision = excluded.remote_revision, last_synced_at = excluded.last_synced_at WHERE excluded.remote_revision >= user_invitations.remote_revision", + params![summary.invitation_id, summary.authority.as_str(), token_hash, summary.label, summary.password_protected, summary.created_at, summary.expires_at, summary.state.as_str(), summary.redeemed_user_id, summary.redeemed_at, summary.revoked_at, summary.remote_revision, synced_at], + )?; + Ok(()) +} + pub fn list() -> Result, StorageError> { db::with_db(|conn| { - let mut statement = conn.prepare("SELECT invitation_id, created_at, expires_at, state, redeemed_user_id, redeemed_at, revoked_at FROM user_invitations ORDER BY created_at DESC")?; + let mut statement = conn.prepare("SELECT invitation_id, authority, label, password_protected, created_at, expires_at, state, remote_revision, redeemed_user_id, redeemed_at, revoked_at, pending_action, pending_action_at, last_synced_at, local_provisioned_user_id, local_provisioned_at FROM user_invitations ORDER BY created_at DESC")?; Ok(statement .query_map([], |row| { Ok(InvitationSummary { invitation_id: row.get(0)?, - created_at: row.get(1)?, - expires_at: row.get(2)?, - state: InvitationState::parse(&row.get::<_, String>(3)?).map_err(|error| { + authority: InvitationAuthority::parse(&row.get::<_, String>(1)?).map_err( + |error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)), + )?, + label: row.get(2)?, + password_protected: row.get(3)?, + created_at: row.get(4)?, + expires_at: row.get(5)?, + state: InvitationState::parse(&row.get::<_, String>(6)?).map_err(|error| { rusqlite::Error::ToSqlConversionFailure(Box::new(error)) })?, - redeemed_user_id: row.get(4)?, - redeemed_at: row.get(5)?, - revoked_at: row.get(6)?, + remote_revision: row.get(7)?, + redeemed_user_id: row.get(8)?, + redeemed_at: row.get(9)?, + revoked_at: row.get(10)?, + pending_action: row + .get::<_, Option>(11)? + .map(|value| PendingAction::parse(&value)) + .transpose() + .map_err(|error| { + rusqlite::Error::ToSqlConversionFailure(Box::new(error)) + })?, + pending_action_at: row.get(12)?, + last_synced_at: row.get(13)?, + local_provisioned_user_id: row.get(14)?, + local_provisioned_at: row.get(15)?, }) })? .collect::, _>>()?) }) } -pub fn revoke(invitation_id: &str, revoked_at: i64) -> Result { +pub fn expire_pending(now: i64) -> Result<(), StorageError> { + db::with_immediate_transaction(|tx| { + tx.execute( + "UPDATE user_invitations SET state = 'expired' WHERE authority = 'iota' AND state = 'pending' AND expires_at IS NOT NULL AND expires_at <= ?1", + params![now], + )?; + Ok(()) + }) +} + +pub fn mark_revoke_pending(invitation_id: i64, requested_at: i64) -> Result { + db::with_immediate_transaction(|tx| mark_revoke_pending_in_tx(tx, invitation_id, requested_at)) +} + +fn mark_revoke_pending_in_tx( + tx: &rusqlite::Transaction<'_>, + invitation_id: i64, + requested_at: i64, +) -> Result { + let changed = tx.execute( + "UPDATE user_invitations SET pending_action = 'revoke', pending_action_at = COALESCE(pending_action_at, ?2) WHERE invitation_id = ?1 AND authority = 'omega' AND state IN ('pending', 'provisioning', 'revoked')", + params![invitation_id, requested_at], + )?; + Ok(changed == 1) +} + +pub fn apply_revoke_result( + invitation_id: i64, + state: InvitationState, + remote_revision: i64, + revoked_at: Option, + synced_at: i64, +) -> Result { db::with_immediate_transaction(|tx| { let changed = tx.execute( - "UPDATE user_invitations SET state = 'revoked', revoked_at = ?2 WHERE invitation_id = ?1 AND state = 'pending'", - params![invitation_id, revoked_at], + "UPDATE user_invitations SET state = ?2, remote_revision = ?3, revoked_at = ?4, pending_action = NULL, pending_action_at = NULL, last_synced_at = ?5 WHERE invitation_id = ?1 AND ?3 >= remote_revision", + params![invitation_id, state.as_str(), remote_revision, revoked_at, synced_at], )?; Ok(changed == 1) }) } + +pub fn apply_external_invitation_provisioning( + invitation_id: i64, + invitation_revision: i64, + user: &UserProfile, + changed_at: i64, +) -> Result { + db::with_immediate_transaction(|tx| { + apply_external_invitation_provisioning_in_tx( + tx, + invitation_id, + invitation_revision, + user, + changed_at, + ) + }) +} + +fn apply_external_invitation_provisioning_in_tx( + tx: &rusqlite::Transaction<'_>, + invitation_id: i64, + invitation_revision: i64, + user: &UserProfile, + changed_at: i64, +) -> Result { + tx.execute( + "UPDATE user_invitations SET state = 'provisioning', remote_revision = ?2, last_synced_at = ?3 WHERE invitation_id = ?1 AND authority = 'omega' AND ?2 >= remote_revision AND state NOT IN ('revoked', 'expired')", + params![invitation_id, invitation_revision, changed_at], + )?; + let invitation = tx + .query_row( + "SELECT state, pending_action, redeemed_user_id FROM user_invitations WHERE invitation_id = ?1 AND authority = 'omega'", + params![invitation_id], + |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option>(1)?, row.get::<_, Option>(2)?)), + ) + .optional()?; + let Some((state, pending_action, redeemed_user_id)) = invitation else { + return Ok(ProvisioningResult::MissingInvitation); + }; + if pending_action.as_deref() == Some("revoke") { + return Ok(ProvisioningResult::RevocationPending); + } + if !matches!(state.as_str(), "pending" | "provisioning" | "redeemed") + || redeemed_user_id.is_some_and(|id| id != user.user_id) + { + return Ok(ProvisioningResult::Conflict); + } + let existing = tx + .query_row( + "SELECT username, public_key FROM users WHERE user_id = ?1", + params![user.user_id], + |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)), + ) + .optional()?; + if let Some((username, public_key)) = existing { + if username != user.username || public_key != user.public_key { + return Ok(ProvisioningResult::Conflict); + } + tx.execute( + "UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = COALESCE(local_provisioned_at, ?3) WHERE invitation_id = ?1", + params![invitation_id, user.user_id, changed_at], + )?; + return Ok(ProvisioningResult::AlreadyApplied); + } + tx.execute( + "INSERT INTO users (user_id, username, public_key, private_key_hash, reset_token, created_at, display_name) VALUES (?1, ?2, ?3, NULL, NULL, ?4, NULL)", + params![user.user_id, user.username, user.public_key, user.created_at], + )?; + tx.execute( + "INSERT INTO user_residency (user_id, username, lifecycle_state, data_state, credential_origin, updated_at) VALUES (?1, ?2, 'managed', 'present', 'external', ?3)", + params![user.user_id, user.username, changed_at], + )?; + tx.execute( + "UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = ?3 WHERE invitation_id = ?1", + params![invitation_id, user.user_id, changed_at], + )?; + Ok(ProvisioningResult::Created) +} + +#[cfg(test)] +mod tests { + use super::*; + use rusqlite::Connection; + + fn database() -> Result { + let connection = Connection::open_in_memory()?; + connection.execute_batch( + "CREATE TABLE user_invitations (invitation_id INTEGER PRIMARY KEY, authority TEXT NOT NULL, token_hash BLOB, label TEXT, password_protected INTEGER NOT NULL, created_at INTEGER NOT NULL, expires_at INTEGER, state TEXT NOT NULL, redeemed_user_id INTEGER, redeemed_at INTEGER, revoked_at INTEGER, pending_action TEXT, pending_action_at INTEGER, remote_revision INTEGER NOT NULL DEFAULT 0, last_synced_at INTEGER, local_provisioned_user_id INTEGER, local_provisioned_at INTEGER); CREATE TABLE users (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, public_key TEXT NOT NULL, private_key_hash TEXT, reset_token TEXT, created_at INTEGER NOT NULL, display_name TEXT); CREATE TABLE user_residency (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, lifecycle_state TEXT NOT NULL, data_state TEXT NOT NULL, credential_origin TEXT NOT NULL, updated_at INTEGER NOT NULL);", + )?; + Ok(connection) + } + + fn summary(revision: i64, state: InvitationState) -> InvitationSummary { + InvitationSummary { + invitation_id: 7, + authority: InvitationAuthority::Omega, + label: None, + password_protected: false, + created_at: 10, + expires_at: Some(20), + state, + remote_revision: revision, + redeemed_user_id: None, + redeemed_at: None, + revoked_at: None, + pending_action: None, + pending_action_at: None, + last_synced_at: None, + local_provisioned_user_id: None, + local_provisioned_at: None, + } + } + + fn profile(public_key: &str) -> UserProfile { + UserProfile::new_with_created_at(9, "alice".into(), None, public_key.into(), None, None, 30) + } + + #[test] + fn offline_revoke_preserves_authoritative_state() -> Result<(), StorageError> { + let mut connection = database()?; + let tx = connection.transaction()?; + insert_in_tx(&tx, &summary(3, InvitationState::Pending), None, 30)?; + assert!(mark_revoke_pending_in_tx(&tx, 7, 31)?); + let row: (String, Option) = tx.query_row( + "SELECT state, pending_action FROM user_invitations WHERE invitation_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + assert_eq!(row, ("pending".into(), Some("revoke".into()))); + Ok(()) + } + + #[test] + fn stale_snapshot_is_ignored_and_pending_action_is_preserved() -> Result<(), StorageError> { + let mut connection = database()?; + let tx = connection.transaction()?; + insert_in_tx(&tx, &summary(5, InvitationState::Pending), None, 30)?; + mark_revoke_pending_in_tx(&tx, 7, 31)?; + insert_in_tx(&tx, &summary(4, InvitationState::Expired), None, 32)?; + let stale: (String, i64, Option) = tx.query_row( + "SELECT state, remote_revision, pending_action FROM user_invitations WHERE invitation_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(stale, ("pending".into(), 5, Some("revoke".into()))); + insert_in_tx(&tx, &summary(6, InvitationState::Expired), None, 33)?; + let fresh: (String, i64, Option) = tx.query_row( + "SELECT state, remote_revision, pending_action FROM user_invitations WHERE invitation_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(fresh, ("expired".into(), 6, Some("revoke".into()))); + Ok(()) + } + + #[test] + fn provisioning_replay_is_idempotent_and_conflicts_do_not_overwrite() -> Result<(), StorageError> + { + let mut connection = database()?; + let tx = connection.transaction()?; + insert_in_tx(&tx, &summary(1, InvitationState::Pending), None, 30)?; + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-a"), 31)?, + ProvisioningResult::Created + ); + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-a"), 32)?, + ProvisioningResult::AlreadyApplied + ); + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-b"), 33)?, + ProvisioningResult::Conflict + ); + let public_key: String = tx.query_row( + "SELECT public_key FROM users WHERE user_id = 9", + [], + |row| row.get(0), + )?; + assert_eq!(public_key, "key-a"); + let invitation: (String, i64, Option) = tx.query_row( + "SELECT state, remote_revision, local_provisioned_user_id FROM user_invitations WHERE invitation_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(invitation, ("provisioning".into(), 2, Some(9))); + Ok(()) + } + + #[test] + fn stale_provisioning_cannot_override_authoritative_revocation() -> Result<(), StorageError> { + let mut connection = database()?; + let tx = connection.transaction()?; + insert_in_tx(&tx, &summary(3, InvitationState::Revoked), None, 30)?; + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-a"), 31)?, + ProvisioningResult::Conflict + ); + let users: i64 = tx.query_row("SELECT COUNT(*) FROM users", [], |row| row.get(0))?; + assert_eq!(users, 0); + Ok(()) + } + + #[test] + fn pending_revoke_blocks_provisioning() -> Result<(), StorageError> { + let mut connection = database()?; + let tx = connection.transaction()?; + insert_in_tx(&tx, &summary(1, InvitationState::Pending), None, 30)?; + mark_revoke_pending_in_tx(&tx, 7, 31)?; + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-a"), 32)?, + ProvisioningResult::RevocationPending + ); + let users: i64 = tx.query_row("SELECT COUNT(*) FROM users", [], |row| row.get(0))?; + assert_eq!(users, 0); + Ok(()) + } + + #[test] + fn missing_invitation_is_reported_separately() -> Result<(), StorageError> { + let mut connection = database()?; + let tx = connection.transaction()?; + assert_eq!( + apply_external_invitation_provisioning_in_tx(&tx, 7, 2, &profile("key-a"), 32)?, + ProvisioningResult::MissingInvitation + ); + Ok(()) + } +} diff --git a/iota-storage/src/util/db.rs b/iota-storage/src/util/db.rs index 14bb08f..85e1c16 100644 --- a/iota-storage/src/util/db.rs +++ b/iota-storage/src/util/db.rs @@ -934,6 +934,47 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { conn.pragma_update(None, "user_version", 25)?; } + if current_version < 27 { + let invitation_exists: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'user_invitations')", + [], + |row| row.get(0), + )?; + if invitation_exists { + conn.execute_batch( + r#" + ALTER TABLE user_invitations RENAME TO user_invitations_before_authority; + CREATE TABLE user_invitations ( + invitation_id INTEGER PRIMARY KEY, + authority TEXT NOT NULL CHECK (authority IN ('omega', 'iota')), + token_hash BLOB NOT NULL, + label TEXT, + password_protected INTEGER NOT NULL DEFAULT 0, + created_at INTEGER NOT NULL, + expires_at INTEGER, + state TEXT NOT NULL CHECK (state IN ('pending', 'provisioning', 'redeemed', 'revoked', 'expired')), + redeemed_user_id INTEGER, + redeemed_at INTEGER, + revoked_at INTEGER, + local_revocation_pending INTEGER NOT NULL DEFAULT 0 + ); + INSERT INTO user_invitations ( + invitation_id, authority, token_hash, created_at, expires_at, + state, redeemed_user_id, redeemed_at, revoked_at + ) + SELECT CAST(invitation_id AS INTEGER), 'omega', token_hash, created_at, expires_at, + state, redeemed_user_id, redeemed_at, revoked_at + FROM user_invitations_before_authority + WHERE CAST(invitation_id AS INTEGER) > 0; + DROP TABLE user_invitations_before_authority; + CREATE INDEX idx_user_invitations_state_created + ON user_invitations (state, created_at); + "#, + )?; + } + conn.pragma_update(None, "user_version", 27)?; + } + if current_version < 26 { conn.execute_batch( r#" @@ -947,7 +988,63 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { ); "#, )?; - conn.pragma_update(None, "user_version", 26)?; + conn.pragma_update(None, "user_version", 27)?; + } + + if current_version < 28 { + let invitation_exists: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'user_invitations')", + [], + |row| row.get(0), + )?; + if invitation_exists { + conn.execute_batch( + r#" + ALTER TABLE user_invitations RENAME TO user_invitations_before_revision; + CREATE TABLE user_invitations ( + invitation_id INTEGER PRIMARY KEY, + authority TEXT NOT NULL CHECK (authority IN ('omega', 'iota')), + token_hash BLOB, + label TEXT, + password_protected INTEGER NOT NULL DEFAULT 0, + created_at INTEGER NOT NULL, + expires_at INTEGER, + state TEXT NOT NULL CHECK (state IN ('pending', 'provisioning', 'redeemed', 'revoked', 'expired')), + redeemed_user_id INTEGER, + redeemed_at INTEGER, + revoked_at INTEGER, + pending_action TEXT CHECK (pending_action IS NULL OR pending_action = 'revoke'), + pending_action_at INTEGER, + remote_revision INTEGER NOT NULL DEFAULT 0, + last_synced_at INTEGER + ); + INSERT INTO user_invitations ( + invitation_id, authority, token_hash, label, password_protected, + created_at, expires_at, state, redeemed_user_id, redeemed_at, + revoked_at, pending_action, pending_action_at + ) + SELECT invitation_id, authority, + CASE WHEN authority = 'iota' THEN token_hash ELSE NULL END, + label, password_protected, created_at, expires_at, state, + redeemed_user_id, redeemed_at, revoked_at, + CASE WHEN local_revocation_pending = 1 THEN 'revoke' ELSE NULL END, + CASE WHEN local_revocation_pending = 1 THEN revoked_at ELSE NULL END + FROM user_invitations_before_revision; + DROP TABLE user_invitations_before_revision; + CREATE INDEX idx_user_invitations_state_created + ON user_invitations (state, created_at); + "#, + )?; + } + conn.pragma_update(None, "user_version", 28)?; + } + + if current_version < 29 { + conn.execute_batch( + "ALTER TABLE user_invitations ADD COLUMN local_provisioned_user_id INTEGER; \ + ALTER TABLE user_invitations ADD COLUMN local_provisioned_at INTEGER;", + )?; + conn.pragma_update(None, "user_version", 29)?; } Ok(()) @@ -1024,7 +1121,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 26); + assert_eq!(version, 29); for column in ["height", "reply_to", "edited_count", "deleted_by_external"] { let mut statement = conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?; @@ -1043,7 +1140,7 @@ mod tests { run_migrations_on_connection(&conn)?; run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 26); + assert_eq!(version, 29); for table in [ "sync_heads", "sync_events", @@ -1090,7 +1187,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 26); + assert_eq!(version, 29); for column in [ "id", "user_id", @@ -1185,7 +1282,40 @@ mod tests { )?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; assert_eq!(preserved, "remote_committed"); - assert_eq!(version, 26); + assert_eq!(version, 29); + Ok(()) + } + + #[test] + fn invitation_authority_migration_preserves_numeric_invitations() -> Result<(), StorageError> { + let conn = Connection::open_in_memory()?; + conn.execute_batch( + r#" + CREATE TABLE user_invitations ( + invitation_id TEXT PRIMARY KEY, + token_hash BLOB NOT NULL, + created_at INTEGER NOT NULL, + expires_at INTEGER, + state TEXT NOT NULL, + redeemed_user_id INTEGER, + redeemed_at INTEGER, + revoked_at INTEGER + ); + INSERT INTO user_invitations ( + invitation_id, token_hash, created_at, expires_at, state + ) VALUES ('42', X'0102', 10, 20, 'pending'); + PRAGMA user_version = 26; + "#, + )?; + + run_migrations_on_connection(&conn)?; + + let row: (i64, String, String, Option) = conn.query_row( + "SELECT invitation_id, authority, state, pending_action FROM user_invitations", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)), + )?; + assert_eq!(row, (42, "omega".into(), "pending".into(), None)); Ok(()) } } diff --git a/iota-updater/Cargo.toml b/iota-updater/Cargo.toml index 8811f02..0d0dd9b 100644 --- a/iota-updater/Cargo.toml +++ b/iota-updater/Cargo.toml @@ -17,3 +17,4 @@ serde_json = "1.0" reqwest = "0.13.2" fs2 = "0.4.3" rusqlite = "0.40.0" +chrono = "0.4.43" diff --git a/iota-updater/src/lib.rs b/iota-updater/src/lib.rs index 60cd293..5f649b8 100644 --- a/iota-updater/src/lib.rs +++ b/iota-updater/src/lib.rs @@ -2,37 +2,214 @@ pub mod manifest; pub mod transaction; use anyhow::{Context, Result, bail}; +use chrono::{DateTime, Utc}; +use iota_ipc::{ + ClientMessage, DaemonMessage, HealthStatus, MIN_PROTOCOL_VERSION, PROTOCOL_VERSION, + StartupPhase, read_msg, write_msg, +}; use manifest::{Artifact, ReleaseManifest, verify_signature}; -use std::{fs, path::Path}; -use transaction::UpdateTransaction; +use serde::{Deserialize, Serialize}; +use std::{ + collections::BTreeMap, + fs, + path::{Path, PathBuf}, + process::Command, + time::Duration, +}; +use tokio::net::UnixStream; +use transaction::{Activation, UpdateTransaction}; const MANIFEST_ENV: &str = "IOTA_UPDATE_MANIFEST"; const SIGNATURE_ENV: &str = "IOTA_UPDATE_SIGNATURE"; const PUBLIC_KEY_ENV: &str = "IOTA_UPDATE_PUBLIC_KEY"; +const CHANNEL_ENV: &str = "IOTA_UPDATE_CHANNEL"; +const SIGNING_KEY_ID_ENV: &str = "IOTA_UPDATE_SIGNING_KEY_ID"; +const ACTIVATION_TIMEOUT_ENV: &str = "IOTA_UPDATE_ACTIVATION_TIMEOUT_SECONDS"; +const ACTIVATION_RETRY_ENV: &str = "IOTA_UPDATE_ACTIVATION_RETRY_MILLISECONDS"; +const DEFAULT_ACTIVATION_TIMEOUT_SECONDS: u64 = 60; +const DEFAULT_ACTIVATION_RETRY_MILLISECONDS: u64 = 250; const REQUIRED_HOST_ARTIFACTS: &str = include_str!("../artifacts.tsv"); -/// Reads the configured signed release manifest and reports whether it differs -/// from the release currently selected by the installation's `current` link. -/// A configured manifest always requires `IOTA_UPDATE_PUBLIC_KEY` (32-byte hex) -/// and a hex signature, supplied by `IOTA_UPDATE_SIGNATURE` or `.sig`. +#[derive(Clone, Debug)] +pub struct ActivationPolicy { + pub timeout: Duration, + pub retry_interval: Duration, +} + +impl ActivationPolicy { + fn from_environment() -> Result { + Ok(Self { + timeout: Duration::from_secs(environment_u64( + ACTIVATION_TIMEOUT_ENV, + DEFAULT_ACTIVATION_TIMEOUT_SECONDS, + )?), + retry_interval: Duration::from_millis(environment_u64( + ACTIVATION_RETRY_ENV, + DEFAULT_ACTIVATION_RETRY_MILLISECONDS, + )?), + }) + } +} + +#[derive(Clone, Debug)] +struct UpdatePolicy { + channel: String, + signing_key_id: String, + activation: ActivationPolicy, +} + +impl UpdatePolicy { + fn from_environment() -> Result { + Ok(Self { + channel: required_environment(CHANNEL_ENV)?, + signing_key_id: required_environment(SIGNING_KEY_ID_ENV)?, + activation: ActivationPolicy::from_environment()?, + }) + } +} + +#[derive(Clone, Debug, Default, Deserialize, Serialize)] +struct UpdateState { + #[serde(default)] + channels: BTreeMap, + #[serde(default)] + pending_activation: Option, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +struct PendingActivation { + previous_target: PathBuf, + new_target: PathBuf, + channel: String, + release_sequence: u64, + product_version: String, + daemon_was_active: bool, + selected: bool, +} + +#[derive(Clone, Debug, Deserialize, Serialize)] +struct ChannelState { + highest_accepted_sequence: u64, + highest_accepted_version: String, + active_version: String, + #[serde(default)] + last_failed_sequence: Option, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum CandidateDecision { + NoUpdate, + Install, +} + +pub trait DaemonSupervisor { + fn is_active(&self) -> Result; + fn restart(&self) -> Result<()>; + fn main_pid(&self) -> Result; +} + +struct SystemdSupervisor; + +impl DaemonSupervisor for SystemdSupervisor { + fn is_active(&self) -> Result { + let status = Command::new("systemctl") + .args(["is-active", "--quiet", "iota-daemon.service"]) + .status() + .context("query iota-daemon.service state")?; + Ok(status.success()) + } + + fn restart(&self) -> Result<()> { + let status = Command::new("systemctl") + .args(["restart", "iota-daemon.service"]) + .status() + .context("restart iota-daemon.service")?; + if !status.success() { + bail!("systemctl restart iota-daemon.service failed with {status}"); + } + Ok(()) + } + + fn main_pid(&self) -> Result { + let output = Command::new("systemctl") + .args([ + "show", + "--property=MainPID", + "--value", + "iota-daemon.service", + ]) + .output() + .context("query iota-daemon.service MainPID")?; + if !output.status.success() { + bail!("could not query iota-daemon.service MainPID"); + } + let text = + String::from_utf8(output.stdout).context("systemctl returned non-UTF-8 MainPID")?; + text.trim() + .parse() + .context("parse iota-daemon.service MainPID") + } +} + pub async fn check_update() -> Result { let Some(release) = configured_release().await? else { return Ok(false); }; let paths = iota_paths::IotaPaths::resolve(iota_paths::Scope::System) .map_err(|error| anyhow::anyhow!(error))?; - Ok(current_version(&paths.install_root)?.as_deref() != Some(&release.manifest.product_version)) + let policy = UpdatePolicy::from_environment()?; + let transaction = UpdateTransaction::from_paths(&paths)?; + let _lock = transaction.acquire()?; + let state = load_or_initialize_update_state(&paths)?; + if state.pending_activation.is_some() { + bail!("an interrupted update requires iota-updater apply recovery"); + } + let decision = evaluate_candidate( + &release.manifest, + current_version(&paths.install_root)?.as_deref(), + state.channels.get(&policy.channel), + &policy.channel, + &policy.signing_key_id, + Utc::now(), + )?; + Ok(decision == CandidateDecision::Install) } -/// Downloads, verifies, stages, and atomically activates the configured release. -/// Returns `false` when no manifest is configured or it already is current. pub async fn apply_update() -> 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 +} + +async fn apply_update_with( + paths: &iota_paths::IotaPaths, + supervisor: &impl DaemonSupervisor, + policy: &UpdatePolicy, +) -> Result { + let transaction = UpdateTransaction::from_paths(paths)?; + let _lock = transaction.acquire()?; + let mut state = load_or_initialize_update_state(paths)?; + recover_interrupted_activation( + paths, + &transaction, + &mut state, + supervisor, + &policy.activation, + ) + .await?; let Some(release) = configured_release().await? else { return Ok(false); }; - let paths = iota_paths::IotaPaths::resolve(iota_paths::Scope::System) - .map_err(|error| anyhow::anyhow!(error))?; - if current_version(&paths.install_root)?.as_deref() == Some(&release.manifest.product_version) { + let decision = evaluate_candidate( + &release.manifest, + current_version(&paths.install_root)?.as_deref(), + state.channels.get(&policy.channel), + &policy.channel, + &policy.signing_key_id, + Utc::now(), + )?; + if decision == CandidateDecision::NoUpdate { return Ok(false); } validate_manifest_compatibility( @@ -42,8 +219,6 @@ pub async fn apply_update() -> Result { iota_ipc::PROTOCOL_VERSION, )?; - let transaction = UpdateTransaction::from_paths(&paths)?; - let _lock = transaction.acquire()?; if transaction.staging.exists() { fs::remove_dir_all(&transaction.staging).with_context(|| { format!( @@ -61,10 +236,493 @@ pub async fn apply_update() -> Result { serde_json::to_vec_pretty(&release.manifest)?, ) .context("write staged release manifest")?; - transaction.activate(&release.manifest.product_version)?; + let daemon_was_active = supervisor.is_active()?; + begin_activation_state( + paths, + &mut state, + &transaction, + &release.manifest, + daemon_was_active, + )?; + let activation = match transaction.activate(&release.manifest.product_version) { + Ok(activation) => activation, + Err(error) => { + remove_unselected_candidate(&state, &transaction.root)?; + clear_pending_activation(paths, &mut state)?; + return Err(error); + } + }; + let activation_result = async { + validate_activation_receipt(&state, &activation)?; + if daemon_was_active { + supervisor.restart()?; + wait_for_ready_daemon( + &iota_paths::socket_path(iota_paths::Scope::System), + &activation.new_target, + supervisor, + &policy.activation, + ) + .await?; + } + commit_update_state(paths, &mut state, &release.manifest)?; + Ok::<_, anyhow::Error>(()) + } + .await; + + if let Err(update_error) = activation_result { + restore_previous_release( + paths, + &transaction, + &activation, + daemon_was_active, + supervisor, + &policy.activation, + ) + .await?; + record_failed_release(paths, &mut state, &release.manifest)?; + return Err( + update_error.context("candidate release failed runtime activation and was rolled back") + ); + } Ok(true) } +async fn recover_interrupted_activation( + paths: &iota_paths::IotaPaths, + transaction: &UpdateTransaction, + state: &mut UpdateState, + supervisor: &impl DaemonSupervisor, + policy: &ActivationPolicy, +) -> Result<()> { + let Some(pending) = state.pending_activation.clone() else { + return Ok(()); + }; + validate_pending_candidate(&pending, &transaction.root)?; + let current_target = transaction.current_target()?; + let candidate_was_selected = pending.selected || current_target == pending.new_target; + if current_target == pending.new_target { + transaction + .restore_activation(&Activation { + previous_target: pending.previous_target.clone(), + new_target: pending.new_target.clone(), + }) + .context("restore previous release after interrupted activation")?; + } else if current_target != pending.previous_target { + bail!( + "cannot recover interrupted activation: current target {} matches neither {} nor {}", + current_target.display(), + pending.previous_target.display(), + pending.new_target.display() + ); + } + if candidate_was_selected && pending.daemon_was_active { + supervisor + .restart() + .context("restart previous daemon after interrupted activation")?; + wait_for_ready_daemon( + &iota_paths::socket_path(iota_paths::Scope::System), + &activation_target_path(&paths.install_root, &pending.previous_target), + supervisor, + policy, + ) + .await + .context("previous daemon failed during interrupted update recovery")?; + } + if candidate_was_selected { + record_failed_sequence(paths, state, &pending.channel, pending.release_sequence)?; + } else { + remove_unselected_candidate(state, &transaction.root)?; + clear_pending_activation(paths, state)?; + } + Ok(()) +} + +fn remove_unselected_candidate(state: &UpdateState, install_root: &Path) -> Result<()> { + let Some(pending) = state.pending_activation.as_ref() else { + return Ok(()); + }; + validate_pending_candidate(pending, install_root)?; + if pending.new_target.exists() { + fs::remove_dir_all(&pending.new_target).with_context(|| { + format!( + "remove unselected candidate release {}", + pending.new_target.display() + ) + })?; + } + Ok(()) +} + +fn validate_pending_candidate(pending: &PendingActivation, install_root: &Path) -> Result<()> { + validate_version(&pending.product_version)?; + let expected = install_root.join("versions").join(&pending.product_version); + if pending.new_target != expected { + bail!( + "pending candidate path {} does not match expected path {}", + pending.new_target.display(), + expected.display() + ); + } + Ok(()) +} + +fn begin_activation_state( + paths: &iota_paths::IotaPaths, + state: &mut UpdateState, + transaction: &UpdateTransaction, + manifest: &ReleaseManifest, + daemon_was_active: bool, +) -> Result<()> { + let new_target = transaction + .root + .join("versions") + .join(&manifest.product_version); + if new_target.exists() { + bail!("release version already exists: {}", new_target.display()); + } + let mut updated = state.clone(); + updated.pending_activation = Some(PendingActivation { + previous_target: transaction.current_target()?, + new_target, + channel: manifest.channel.clone(), + release_sequence: manifest.release_sequence, + product_version: manifest.product_version.clone(), + daemon_was_active, + selected: true, + }); + save_update_state(&paths.update_status_file(), &updated)?; + *state = updated; + Ok(()) +} + +fn validate_activation_receipt(state: &UpdateState, activation: &Activation) -> Result<()> { + let pending = state + .pending_activation + .as_ref() + .context("activation has no persisted rollback receipt")?; + if pending.previous_target != activation.previous_target + || pending.new_target != activation.new_target + { + bail!("activation receipt does not match persisted rollback state"); + } + Ok(()) +} + +fn clear_pending_activation(paths: &iota_paths::IotaPaths, state: &mut UpdateState) -> Result<()> { + let mut updated = state.clone(); + updated.pending_activation = None; + save_update_state(&paths.update_status_file(), &updated)?; + *state = updated; + Ok(()) +} + +async fn restore_previous_release( + paths: &iota_paths::IotaPaths, + transaction: &UpdateTransaction, + activation: &Activation, + daemon_was_active: bool, + supervisor: &impl DaemonSupervisor, + policy: &ActivationPolicy, +) -> Result<()> { + transaction + .restore_activation(activation) + .context("restore previous release after failed activation")?; + if daemon_was_active { + supervisor + .restart() + .context("restart previous daemon after failed activation")?; + wait_for_ready_daemon( + &iota_paths::socket_path(iota_paths::Scope::System), + &activation_target_path(&paths.install_root, &activation.previous_target), + supervisor, + policy, + ) + .await + .context("previous daemon failed after update rollback")?; + } + Ok(()) +} + +fn activation_target_path(install_root: &Path, target: &Path) -> PathBuf { + if target.is_absolute() { + target.to_owned() + } else { + install_root.join(target) + } +} + +async fn probe_daemon(socket: &Path) -> Result { + let mut stream = UnixStream::connect(socket) + .await + .with_context(|| format!("connect to {}", socket.display()))?; + write_msg( + &mut stream, + &ClientMessage::Hello { + supported_versions: (MIN_PROTOCOL_VERSION..=PROTOCOL_VERSION).collect(), + }, + ) + .await + .context("send updater IPC hello")?; + match read_msg::<_, DaemonMessage>(&mut stream).await? { + DaemonMessage::HelloAck(ack) => Ok(ack), + other => bail!("expected daemon HelloAck, received {other:?}"), + } +} + +async fn wait_for_ready_daemon( + socket: &Path, + expected_release: &Path, + supervisor: &impl DaemonSupervisor, + policy: &ActivationPolicy, +) -> Result<()> { + let deadline = tokio::time::Instant::now() + policy.timeout; + loop { + if let Ok(ack) = probe_daemon(socket).await + && ack.startup_phase == StartupPhase::Ready + && ack.health != HealthStatus::Failed + { + verify_running_executable(supervisor, expected_release)?; + return Ok(()); + } + if tokio::time::Instant::now() >= deadline { + bail!("updated daemon did not reach the required ready state"); + } + tokio::time::sleep(policy.retry_interval).await; + } +} + +fn verify_running_executable( + supervisor: &impl DaemonSupervisor, + expected_release: &Path, +) -> Result<()> { + let pid = supervisor.main_pid()?; + if pid == 0 { + bail!("iota-daemon.service has no MainPID"); + } + let actual = fs::canonicalize(format!("/proc/{pid}/exe")) + .context("resolve running daemon executable")?; + let expected = fs::canonicalize(expected_release.join("bin/iota-daemon")) + .context("resolve expected daemon executable")?; + if actual != expected { + bail!( + "running daemon executable mismatch: expected {}, got {}", + expected.display(), + actual.display() + ); + } + Ok(()) +} + +fn evaluate_candidate( + manifest: &ReleaseManifest, + current_version: Option<&str>, + channel_state: Option<&ChannelState>, + expected_channel: &str, + expected_key_id: &str, + now: DateTime, +) -> Result { + validate_manifest_identity(manifest, expected_channel, expected_key_id)?; + validate_manifest_freshness(manifest, now)?; + if manifest.release_sequence == 0 { + bail!("release sequence must be greater than zero"); + } + if let Some(state) = channel_state { + if manifest.release_sequence < state.highest_accepted_sequence { + bail!( + "release sequence {} is older than trusted sequence {}", + manifest.release_sequence, + state.highest_accepted_sequence + ); + } + if manifest.release_sequence == state.highest_accepted_sequence + && manifest.product_version != state.highest_accepted_version + { + bail!( + "release sequence {} was already associated with version {}", + manifest.release_sequence, + state.highest_accepted_version + ); + } + if state.last_failed_sequence == Some(manifest.release_sequence) { + bail!( + "release sequence {} previously failed activation", + manifest.release_sequence + ); + } + if manifest.release_sequence == state.highest_accepted_sequence { + return Ok(CandidateDecision::NoUpdate); + } + } + if current_version == Some(manifest.product_version.as_str()) { + return Ok(CandidateDecision::NoUpdate); + } + Ok(CandidateDecision::Install) +} + +fn validate_manifest_identity( + manifest: &ReleaseManifest, + expected_channel: &str, + expected_key_id: &str, +) -> Result<()> { + if manifest.channel != expected_channel { + bail!( + "release channel mismatch: expected {}, got {}", + expected_channel, + manifest.channel + ); + } + if manifest.release_signing_key_id != expected_key_id { + bail!( + "release signing key id mismatch: expected {}, got {}", + expected_key_id, + manifest.release_signing_key_id + ); + } + Ok(()) +} + +fn validate_manifest_freshness(manifest: &ReleaseManifest, now: DateTime) -> Result<()> { + let published_at = DateTime::parse_from_rfc3339(&manifest.published_at) + .context("parse manifest published_at")? + .with_timezone(&Utc); + let expires_at = DateTime::parse_from_rfc3339(&manifest.expires_at) + .context("parse manifest expires_at")? + .with_timezone(&Utc); + if expires_at <= published_at { + bail!("release manifest expires_at must be after published_at"); + } + if now > expires_at { + bail!("release manifest expired at {}", manifest.expires_at); + } + Ok(()) +} + +fn load_or_initialize_update_state(paths: &iota_paths::IotaPaths) -> Result { + let path = paths.update_status_file(); + if path.exists() { + let bytes = + fs::read(&path).with_context(|| format!("read update status {}", path.display()))?; + return serde_json::from_slice(&bytes) + .with_context(|| format!("parse update status {}", path.display())); + } + let installed = read_installed_manifest(&paths.install_root.join("current")).context( + "initialize anti-rollback state from installed manifest; legacy installations require an installer migration", + )?; + if installed.release_sequence == 0 { + bail!("installed release sequence must be greater than zero"); + } + let mut state = UpdateState::default(); + state.channels.insert( + installed.channel.clone(), + ChannelState { + highest_accepted_sequence: installed.release_sequence, + highest_accepted_version: installed.product_version.clone(), + active_version: installed.product_version, + last_failed_sequence: None, + }, + ); + save_update_state(&path, &state)?; + Ok(state) +} + +fn save_update_state(path: &Path, state: &UpdateState) -> Result<()> { + let parent = path.parent().context("update status path has no parent")?; + fs::create_dir_all(parent)?; + let mut temporary = + tempfile::NamedTempFile::new_in(parent).context("create temporary update status file")?; + serde_json::to_writer_pretty(temporary.as_file_mut(), state) + .context("serialize update status")?; + temporary + .as_file_mut() + .sync_all() + .context("flush update status")?; + temporary + .persist(path) + .map_err(|error| error.error) + .context("replace update status")?; + fs::File::open(parent) + .context("open update status directory")? + .sync_all() + .context("flush update status directory")?; + Ok(()) +} + +fn commit_update_state( + paths: &iota_paths::IotaPaths, + state: &mut UpdateState, + manifest: &ReleaseManifest, +) -> Result<()> { + let mut updated = state.clone(); + updated.channels.insert( + manifest.channel.clone(), + ChannelState { + highest_accepted_sequence: manifest.release_sequence, + highest_accepted_version: manifest.product_version.clone(), + active_version: manifest.product_version.clone(), + last_failed_sequence: None, + }, + ); + updated.pending_activation = None; + save_update_state(&paths.update_status_file(), &updated)?; + *state = updated; + Ok(()) +} + +fn record_failed_release( + paths: &iota_paths::IotaPaths, + state: &mut UpdateState, + manifest: &ReleaseManifest, +) -> Result<()> { + record_failed_sequence(paths, state, &manifest.channel, manifest.release_sequence) +} + +fn record_failed_sequence( + paths: &iota_paths::IotaPaths, + state: &mut UpdateState, + channel_name: &str, + release_sequence: u64, +) -> Result<()> { + let mut updated = state.clone(); + let active_version = current_version(&paths.install_root)? + .context("active release has no product version while recording failed update")?; + let channel = updated + .channels + .entry(channel_name.to_owned()) + .or_insert_with(|| ChannelState { + highest_accepted_sequence: 0, + highest_accepted_version: String::new(), + active_version, + last_failed_sequence: None, + }); + channel.last_failed_sequence = Some(release_sequence); + updated.pending_activation = None; + save_update_state(&paths.update_status_file(), &updated)?; + *state = updated; + Ok(()) +} + +fn required_environment(name: &'static str) -> Result { + let value = std::env::var(name).with_context(|| format!("{name} is required"))?; + if value.is_empty() { + bail!("{name} must not be empty"); + } + Ok(value) +} + +fn environment_u64(name: &'static str, default: u64) -> Result { + let value = match std::env::var(name) { + Ok(value) => value + .parse::() + .with_context(|| format!("{name} must be an unsigned integer"))?, + Err(std::env::VarError::NotPresent) => default, + Err(std::env::VarError::NotUnicode(_)) => bail!("{name} must be valid UTF-8"), + }; + if value == 0 { + bail!("{name} must be greater than zero"); + } + Ok(value) +} + /// Switches `current` to an already-installed release version. pub fn rollback(version: &str) -> Result<()> { validate_version(version)?; @@ -72,6 +730,10 @@ pub fn rollback(version: &str) -> Result<()> { .map_err(|error| anyhow::anyhow!(error))?; let transaction = UpdateTransaction::from_paths(&paths)?; let _lock = transaction.acquire()?; + let mut state = load_or_initialize_update_state(&paths)?; + if state.pending_activation.is_some() { + bail!("an interrupted update requires iota-updater apply recovery"); + } let release_dir = transaction.root.join("versions").join(version); if !release_dir.is_dir() { bail!( @@ -88,7 +750,43 @@ pub fn rollback(version: &str) -> Result<()> { iota_ipc::MIN_PROTOCOL_VERSION, iota_ipc::PROTOCOL_VERSION, )?; - transaction.rollback(version) + if active_manifest.channel != target_manifest.channel { + bail!( + "rollback target channel {} does not match active channel {}", + target_manifest.channel, + active_manifest.channel + ); + } + let previous_target = fs::read_link(transaction.root.join("current")) + .context("read current release link before rollback")?; + transaction.rollback(version)?; + let activation = Activation { + previous_target, + new_target: release_dir, + }; + if let Err(error) = set_active_version(&paths, &mut state, &target_manifest) { + transaction + .restore_activation(&activation) + .context("restore active release after rollback state commit failed")?; + return Err(error.context("persist rollback update state")); + } + Ok(()) +} + +fn set_active_version( + paths: &iota_paths::IotaPaths, + state: &mut UpdateState, + manifest: &ReleaseManifest, +) -> Result<()> { + let mut updated = state.clone(); + let channel = updated + .channels + .get_mut(&manifest.channel) + .context("rollback channel has no trusted baseline")?; + channel.active_version = manifest.product_version.clone(); + save_update_state(&paths.update_status_file(), &updated)?; + *state = updated; + Ok(()) } struct ConfiguredRelease { @@ -355,6 +1053,40 @@ mod tests { use super::*; use ed25519_dalek::{Signer, SigningKey}; + struct CurrentProcessSupervisor; + + impl DaemonSupervisor for CurrentProcessSupervisor { + fn is_active(&self) -> Result { + Ok(true) + } + + fn restart(&self) -> Result<()> { + Ok(()) + } + + fn main_pid(&self) -> Result { + Ok(std::process::id()) + } + } + + fn test_paths(root: &Path) -> iota_paths::IotaPaths { + let state_dir = root.join("state"); + iota_paths::IotaPaths { + scope: iota_paths::Scope::User, + config_dir: root.join("config"), + config_file: root.join("config/config.yaml"), + state_dir: state_dir.clone(), + storage_dir: state_dir.join("storage"), + identity_dir: state_dir.join("identity"), + cache_dir: root.join("cache"), + runtime_dir: Some(root.join("run")), + log_dir: root.join("logs"), + asset_dir: root.join("web"), + install_root: root.join("install"), + ipc_endpoint: iota_paths::IpcEndpoint::UnixSocket(root.join("run/iota.sock")), + } + } + fn artifact(role: &str, path: &str, os: &str, architecture: &str) -> Artifact { Artifact { role: role.into(), @@ -371,7 +1103,9 @@ mod tests { ReleaseManifest { product_version: "1.2.3".into(), channel: "stable".into(), + release_sequence: 12, published_at: "2026-09-10T00:00:00Z".into(), + expires_at: "2030-09-10T00:00:00Z".into(), minimum_data_schema: 1, supported_ipc_min: 1, supported_ipc_max: 1, @@ -517,4 +1251,224 @@ mod tests { assert!(error.to_string().contains("does not permit rollback")); } + + #[test] + fn candidate_policy_rejects_replay_and_sequence_reuse() { + let state = ChannelState { + highest_accepted_sequence: 12, + highest_accepted_version: "1.2.3".into(), + active_version: "1.1.0".into(), + last_failed_sequence: None, + }; + let mut replay = release_manifest(host_artifacts()); + replay.release_sequence = 11; + assert!( + evaluate_candidate( + &replay, + Some("1.1.0"), + Some(&state), + "stable", + "test", + Utc::now(), + ) + .unwrap_err() + .to_string() + .contains("older than trusted sequence") + ); + + let mut reused = release_manifest(host_artifacts()); + reused.product_version = "1.2.4".into(); + assert!( + evaluate_candidate( + &reused, + Some("1.1.0"), + Some(&state), + "stable", + "test", + Utc::now(), + ) + .unwrap_err() + .to_string() + .contains("already associated") + ); + } + + #[test] + fn accepted_sequence_is_not_reinstalled_after_explicit_rollback() { + let state = ChannelState { + highest_accepted_sequence: 12, + highest_accepted_version: "1.2.3".into(), + active_version: "1.1.0".into(), + last_failed_sequence: None, + }; + let decision = evaluate_candidate( + &release_manifest(host_artifacts()), + Some("1.1.0"), + Some(&state), + "stable", + "test", + Utc::now(), + ) + .unwrap(); + assert_eq!(decision, CandidateDecision::NoUpdate); + } + + #[test] + fn candidate_policy_rejects_expired_wrong_channel_and_failed_release() { + let mut expired = release_manifest(host_artifacts()); + expired.published_at = "2020-01-01T00:00:00Z".into(); + expired.expires_at = "2020-01-02T00:00:00Z".into(); + assert!( + evaluate_candidate(&expired, None, None, "stable", "test", Utc::now(),) + .unwrap_err() + .to_string() + .contains("expired") + ); + + let manifest = release_manifest(host_artifacts()); + assert!( + evaluate_candidate(&manifest, None, None, "dev", "test", Utc::now()) + .unwrap_err() + .to_string() + .contains("channel mismatch") + ); + + let state = ChannelState { + highest_accepted_sequence: 11, + highest_accepted_version: "1.2.2".into(), + active_version: "1.2.2".into(), + last_failed_sequence: Some(12), + }; + assert!( + evaluate_candidate( + &manifest, + Some("1.2.2"), + Some(&state), + "stable", + "test", + Utc::now(), + ) + .unwrap_err() + .to_string() + .contains("previously failed activation") + ); + } + + #[test] + fn update_state_is_replaced_atomically() { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("update-status.json"); + let mut state = UpdateState::default(); + state.channels.insert( + "stable".into(), + ChannelState { + highest_accepted_sequence: 12, + highest_accepted_version: "1.2.3".into(), + active_version: "1.2.3".into(), + last_failed_sequence: None, + }, + ); + save_update_state(&path, &state).unwrap(); + let loaded: UpdateState = serde_json::from_slice(&fs::read(path).unwrap()).unwrap(); + assert_eq!(loaded.channels["stable"].highest_accepted_sequence, 12); + } + + #[tokio::test] + async fn readiness_probe_checks_ipc_state_and_running_executable() { + let directory = tempfile::tempdir().unwrap(); + let socket = directory.path().join("iota.sock"); + let release = directory.path().join("versions/1.2.3"); + fs::create_dir_all(release.join("bin")).unwrap(); + std::os::unix::fs::symlink( + fs::canonicalize("/proc/self/exe").unwrap(), + release.join("bin/iota-daemon"), + ) + .unwrap(); + let listener = tokio::net::UnixListener::bind(&socket).unwrap(); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.unwrap(); + let hello = read_msg::<_, ClientMessage>(&mut stream).await.unwrap(); + assert!(matches!(hello, ClientMessage::Hello { .. })); + write_msg( + &mut stream, + &DaemonMessage::HelloAck(iota_ipc::HelloAck { + protocol_version: PROTOCOL_VERSION, + daemon_version: "test".into(), + instance_id: "test".into(), + startup_phase: StartupPhase::Ready, + capabilities: Vec::new(), + lifecycle: iota_ipc::LifecyclePhase::Ready, + health: HealthStatus::Degraded, + deployment_mode: iota_ipc::DeploymentMode::default(), + supervisor: iota_ipc::SupervisorKind::default(), + }), + ) + .await + .unwrap(); + }); + wait_for_ready_daemon( + &socket, + &release, + &CurrentProcessSupervisor, + &ActivationPolicy { + timeout: Duration::from_secs(1), + retry_interval: Duration::from_millis(10), + }, + ) + .await + .unwrap(); + server.await.unwrap(); + } + + #[tokio::test] + async fn interrupted_selected_candidate_is_restored_and_marked_failed() { + let directory = tempfile::tempdir().unwrap(); + let paths = test_paths(directory.path()); + let transaction = UpdateTransaction::new(&paths.install_root); + let old_release = paths.install_root.join("versions/1.0.0"); + let new_release = paths.install_root.join("versions/1.1.0"); + fs::create_dir_all(&old_release).unwrap(); + fs::create_dir_all(&new_release).unwrap(); + fs::write( + old_release.join("manifest.json"), + br#"{"product_version":"1.0.0"}"#, + ) + .unwrap(); + fs::create_dir_all(&paths.install_root).unwrap(); + std::os::unix::fs::symlink(&new_release, paths.install_root.join("current")).unwrap(); + let mut state = UpdateState::default(); + state.channels.insert( + "stable".into(), + ChannelState { + highest_accepted_sequence: 1, + highest_accepted_version: "1.0.0".into(), + active_version: "1.0.0".into(), + last_failed_sequence: None, + }, + ); + state.pending_activation = Some(PendingActivation { + previous_target: old_release.clone(), + new_target: new_release, + channel: "stable".into(), + release_sequence: 2, + product_version: "1.1.0".into(), + daemon_was_active: false, + selected: true, + }); + recover_interrupted_activation( + &paths, + &transaction, + &mut state, + &CurrentProcessSupervisor, + &ActivationPolicy { + timeout: Duration::from_secs(1), + retry_interval: Duration::from_millis(10), + }, + ) + .await + .unwrap(); + assert_eq!(transaction.current_target().unwrap(), old_release); + assert_eq!(state.channels["stable"].last_failed_sequence, Some(2)); + assert!(state.pending_activation.is_none()); + } } diff --git a/iota-updater/src/manifest.rs b/iota-updater/src/manifest.rs index f2bae0a..610449e 100644 --- a/iota-updater/src/manifest.rs +++ b/iota-updater/src/manifest.rs @@ -7,7 +7,9 @@ use sha2::{Digest, Sha256}; pub struct ReleaseManifest { pub product_version: String, pub channel: String, + pub release_sequence: u64, pub published_at: String, + pub expires_at: String, pub minimum_data_schema: u64, pub supported_ipc_min: u16, pub supported_ipc_max: u16, diff --git a/iota-updater/src/transaction.rs b/iota-updater/src/transaction.rs index bb1693e..2f64980 100644 --- a/iota-updater/src/transaction.rs +++ b/iota-updater/src/transaction.rs @@ -6,6 +6,12 @@ use std::{ path::{Component, Path, PathBuf}, }; +#[derive(Clone, Debug)] +pub struct Activation { + pub previous_target: PathBuf, + pub new_target: PathBuf, +} + #[derive(Clone, Debug)] pub struct UpdateTransaction { pub root: PathBuf, @@ -43,6 +49,11 @@ impl UpdateTransaction { .context("update already in progress")?; Ok(UpdateLock { _file: file }) } + pub fn current_target(&self) -> Result { + let current = self.root.join("current"); + fs::read_link(¤t) + .with_context(|| format!("read current release link {}", current.display())) + } pub fn stage_artifact(&self, source: &Path, artifact: &Artifact) -> Result { let artifact_path = Path::new(&artifact.path); if artifact_path.is_absolute() @@ -65,9 +76,15 @@ impl UpdateTransaction { set_executable_if_binary(&target, artifact)?; Ok(target) } - pub fn activate(&self, version: &str) -> Result<()> { + pub fn activate(&self, version: &str) -> Result { + let current = self.root.join("current"); + let previous_target = self.current_target()?; let version_dir = self.root.join("versions").join(version); - fs::create_dir_all(version_dir.parent().unwrap())?; + fs::create_dir_all( + version_dir + .parent() + .context("release version directory has no parent")?, + )?; if version_dir.exists() { anyhow::bail!("release version already exists: {}", version_dir.display()); } @@ -80,22 +97,46 @@ impl UpdateTransaction { } Err(error) => return Err(error).context("activate staged release"), } - let current_tmp = self.root.join("current.new"); - let _ = fs::remove_file(¤t_tmp); - std::os::unix::fs::symlink(&version_dir, ¤t_tmp)?; - fs::rename(current_tmp, self.root.join("current"))?; - Ok(()) + replace_symlink(¤t, &version_dir, &self.root.join("current.new"))?; + Ok(Activation { + previous_target, + new_target: version_dir, + }) + } + pub fn restore_activation(&self, activation: &Activation) -> Result<()> { + replace_symlink( + &self.root.join("current"), + &activation.previous_target, + &self.root.join("current.rollback"), + ) } pub fn rollback(&self, previous: &str) -> Result<()> { - let current = self.root.join("current"); - let tmp = self.root.join("current.rollback"); - let _ = fs::remove_file(&tmp); - std::os::unix::fs::symlink(self.root.join("versions").join(previous), &tmp)?; - fs::rename(tmp, current)?; - Ok(()) + replace_symlink( + &self.root.join("current"), + &self.root.join("versions").join(previous), + &self.root.join("current.rollback"), + ) } } +#[cfg(unix)] +fn replace_symlink(link: &Path, target: &Path, temporary: &Path) -> Result<()> { + match fs::remove_file(temporary) { + Ok(()) => {} + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(error).context("remove stale temporary release link"), + } + std::os::unix::fs::symlink(target, temporary).with_context(|| { + format!( + "create temporary symlink {} -> {}", + temporary.display(), + target.display() + ) + })?; + fs::rename(temporary, link).with_context(|| format!("replace symlink {}", link.display()))?; + Ok(()) +} + #[cfg(unix)] fn set_executable_if_binary(path: &Path, artifact: &Artifact) -> Result<()> { use std::os::unix::fs::PermissionsExt; @@ -169,12 +210,21 @@ mod tests { drop(lock); assert!(tx.acquire().is_ok()); assert!(tx.lock_file.exists()); + let old_release = dir.path().join("versions/0.9.0"); + std::fs::create_dir_all(&old_release).unwrap(); + std::os::unix::fs::symlink(&old_release, dir.path().join("current")).unwrap(); std::fs::create_dir_all(&tx.staging).unwrap(); std::fs::write(tx.staging.join("manifest.json"), b"ok").unwrap(); - tx.activate("1.0.0").unwrap(); + let activation = tx.activate("1.0.0").unwrap(); assert_eq!( std::fs::read_to_string(dir.path().join("current/manifest.json")).unwrap(), "ok" ); + assert_eq!(activation.previous_target, old_release); + tx.restore_activation(&activation).unwrap(); + assert_eq!( + std::fs::read_link(dir.path().join("current")).unwrap(), + old_release + ); } } diff --git a/iota/src/main.rs b/iota/src/main.rs index e211ce9..c0e593c 100644 --- a/iota/src/main.rs +++ b/iota/src/main.rs @@ -843,6 +843,32 @@ async fn run_command( } } } + ResponsePayload::InvitationCreated(invitation) => { + println!("Invitation: {}", invitation.invitation_id); + println!("Token: {}", invitation.raw_token.0); + if let Some(url) = invitation.short_url { + println!("Link: {url}"); + } + println!("Expires: {}", invitation.expires_at); + } + ResponsePayload::Invitations(invitations) => { + if invitations.is_empty() { + println!("No invitations."); + } else { + for invitation in invitations { + println!( + "{} ({:?}, {:?})", + invitation.invitation_id, invitation.authority, invitation.state + ); + } + } + } + ResponsePayload::InvitationUpdated(invitation) => { + println!( + "Invitation {} is {:?}", + invitation.invitation_id, invitation.state + ); + } ResponsePayload::UserCreated { user_id, username } => { println!( "{} {} ({})", @@ -1023,6 +1049,35 @@ fn render_structured(payload: &ResponsePayload, output: OutputFormat) -> Result< fn render_table(payload: &ResponsePayload) { match payload { + ResponsePayload::InvitationCreated(invitation) => { + println!("{:<15} {}", "Invitation ID", invitation.invitation_id); + println!( + "{:<15} {}", + "Authority", + format!("{:?}", invitation.authority) + ); + println!("{:<15} {}", "Token", invitation.raw_token.0); + println!("{:<15} {}", "Expires", invitation.expires_at); + if let Some(url) = &invitation.short_url { + println!("{:<15} {}", "Link", url); + } + } + ResponsePayload::Invitations(invitations) => { + println!("{:<16} {:<10} {:<14} LABEL", "ID", "AUTHORITY", "STATE"); + for invitation in invitations { + println!( + "{:<16} {:<10} {:<14} {}", + invitation.invitation_id, + format!("{:?}", invitation.authority), + format!("{:?}", invitation.state), + invitation.label.as_deref().unwrap_or("-") + ); + } + } + ResponsePayload::InvitationUpdated(invitation) => { + println!("{:<15} {}", "Invitation ID", invitation.invitation_id); + println!("{:<15} {:?}", "State", invitation.state); + } ResponsePayload::Users(users) => { if users.is_empty() { println!("No users."); diff --git a/mtp-type-maps b/mtp-type-maps index 2388c22..91d0397 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit 2388c225b50db4d566653fb220786931eb3989c7 +Subproject commit 91d03974632a2897de5457d52170dda9d9e36c3d diff --git a/omikron-connector/src/client.rs b/omikron-connector/src/client.rs index 03cb832..49f734a 100644 --- a/omikron-connector/src/client.rs +++ b/omikron-connector/src/client.rs @@ -8,6 +8,7 @@ pub enum OmikronError { Timeout(String), Authentication(String), Rejected(CommunicationType, String), + Storage(String), Internal(String), } @@ -18,6 +19,7 @@ impl std::fmt::Display for OmikronError { | Self::Timeout(v) | Self::Authentication(v) | Self::Rejected(_, v) + | Self::Storage(v) | Self::Internal(v) => f.write_str(v), } } @@ -42,6 +44,12 @@ pub trait OmikronClient: Send + Sync { value: &CommunicationValue, timeout: Duration, ) -> Result; + async fn sync_omega_invitations(&self) -> Result<(), OmikronError> { + Err(OmikronError::Internal( + "invitation synchronization is unavailable".into(), + )) + } + async fn flush_pending_invitation_actions(&self) {} async fn reconnect(&self) -> Result<(), OmikronError>; /// Replace the local Iota identity and wait for the new identity to /// register/authenticate. This is deliberately available while offline: diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index 9867fa9..e9e0bc3 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -9,13 +9,13 @@ use iota_storage::util::{ use iota_util::crypto_helper::{self, keyring_from_base64}; use iota_util::crypto_util::{self}; use mtp::client::{Client, ClientConfig, MTPConnection, Policy, SendMode, Sender}; -use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; +use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataTypeId, DataValue, TypeMap}; use mtp::crypto::{Keyring, PublicKeyBundle}; use rand_core::RngCore; use std::env; use std::fs; use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::sync::{Arc, LazyLock}; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use tokio::sync::{Mutex, RwLock, Semaphore, oneshot, watch}; @@ -128,6 +128,119 @@ const RECONNECT_DELAY: Duration = Duration::from_secs(5); const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(300); const CONNECTION_TIMEOUT: Duration = Duration::from_secs(45); const MAINTENANCE_INTERVAL: Duration = Duration::from_secs(5); + +fn container_value<'a>( + fields: &'a [(DataTypeId, DataValue)], + kind: DataType, + type_map: &TypeMap, +) -> Option<&'a DataValue> { + let id = kind.try_to_id(type_map)?; + fields + .iter() + .find_map(|(field_id, value)| (*field_id == id).then_some(value)) +} + +fn container_signed_i64( + fields: &[(DataTypeId, DataValue)], + kind: DataType, + type_map: &TypeMap, +) -> Option { + container_value(fields, kind, type_map) + .and_then(DataValue::as_signed_number) + .and_then(|value| i64::try_from(value).ok()) +} + +fn parse_omega_invitation( + fields: &[(DataTypeId, DataValue)], + type_map: &TypeMap, + expected_iota_id: i64, + synced_at: i64, +) -> Result { + use iota_storage::users::invitations::{InvitationState, InvitationSummary}; + + let invitation_id = container_signed_i64(fields, DataType::InvitationId, type_map) + .filter(|value| *value > 0) + .ok_or_else(|| OmikronError::Internal("invalid Omega invitation ID".into()))?; + if container_value(fields, DataType::InvitationAuthority, type_map).and_then(DataValue::as_str) + != Some("omega") + { + return Err(OmikronError::Internal( + "invalid Omega invitation authority".into(), + )); + } + if container_signed_i64(fields, DataType::IotaId, type_map) != Some(expected_iota_id) { + return Err(OmikronError::Internal( + "Omega invitation snapshot contained the wrong Iota ID".into(), + )); + } + let remote_revision = container_signed_i64(fields, DataType::InvitationRevision, type_map) + .filter(|value| *value > 0) + .ok_or_else(|| OmikronError::Internal("invalid Omega invitation revision".into()))?; + let state = match container_value(fields, DataType::InvitationState, type_map) + .and_then(DataValue::as_str) + { + Some("pending") => InvitationState::Pending, + Some("provisioning") => InvitationState::Provisioning, + Some("redeemed") => InvitationState::Redeemed, + Some("revoked") => InvitationState::Revoked, + Some("expired") => InvitationState::Expired, + _ => { + return Err(OmikronError::Internal( + "invalid Omega invitation state".into(), + )); + } + }; + let label = match container_value(fields, DataType::InvitationLabel, type_map) { + Some(value) => Some( + value + .as_str() + .ok_or_else(|| OmikronError::Internal("invalid Omega invitation label".into()))?, + ), + None => None, + }; + let password_protected = + container_value(fields, DataType::InvitationPasswordProtected, type_map) + .and_then(DataValue::as_bool) + .ok_or_else(|| { + OmikronError::Internal("invalid Omega invitation password protection flag".into()) + })?; + let created_at = container_signed_i64(fields, DataType::InvitationCreatedAt, type_map) + .filter(|value| *value > 0) + .ok_or_else(|| OmikronError::Internal("invalid Omega invitation creation time".into()))?; + let expires_at = container_signed_i64(fields, DataType::InvitationExpiresAt, type_map) + .filter(|value| *value > 0) + .ok_or_else(|| OmikronError::Internal("invalid Omega invitation expiry time".into()))?; + let redeemed_user_id = match container_value(fields, DataType::UserId, type_map) { + Some(value) => Some( + value + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0) + .ok_or_else(|| { + OmikronError::Internal("invalid Omega invitation redeemed user ID".into()) + })?, + ), + None => None, + }; + Ok(InvitationSummary { + invitation_id, + authority: iota_storage::users::invitations::InvitationAuthority::Omega, + label: label.map(str::to_owned), + password_protected, + created_at, + expires_at: Some(expires_at), + state, + remote_revision, + redeemed_user_id, + redeemed_at: None, + revoked_at: None, + pending_action: None, + pending_action_at: None, + last_synced_at: Some(synced_at), + local_provisioned_user_id: None, + local_provisioned_at: None, + }) +} const TASK_CLEANUP_INTERVAL: Duration = Duration::from_secs(60); const TASK_MAX_AGE: Duration = Duration::from_secs(60); const MAX_CONCURRENT_HANDLERS: usize = 20; @@ -140,6 +253,26 @@ struct ResolvedOmikronEndpoint { public_key: PublicKeyBundle, } +struct InvitationSyncRuntime { + flush_in_progress: AtomicBool, + next_retry_at: Mutex>, + retry_delay: Mutex, +} + +impl InvitationSyncRuntime { + fn new() -> Self { + Self { + flush_in_progress: AtomicBool::new(false), + next_retry_at: Mutex::new(None), + retry_delay: Mutex::new(Duration::from_secs(5)), + } + } +} + +fn next_invitation_retry_delay(current: Duration) -> Duration { + (current * 2).min(Duration::from_secs(300)) +} + #[derive(Debug)] pub enum IdentityError { Storage(mtp::files::FileError), @@ -320,6 +453,7 @@ pub struct OmikronConnection { cancellation: CancellationToken, pub(crate) active_tasks: Arc>, pub(crate) app: Arc>, + invitation_sync: Arc, } impl OmikronConnection { @@ -353,6 +487,7 @@ impl OmikronConnection { cancellation, active_tasks, app, + invitation_sync: Arc::new(InvitationSyncRuntime::new()), } } @@ -546,6 +681,13 @@ impl OmikronConnection { log_t!("omikron_authenticated"); self.classify_legacy_pending_relays().await; + if let Err(error) = self.sync_omega_invitations().await { + log!( + "Omega invitation snapshot synchronization failed: {}", + error + ); + } + self.flush_pending_invitation_actions().await; let maintenance_self = self.clone(); let maintenance_handle = tokio::spawn(async move { @@ -585,6 +727,155 @@ impl OmikronConnection { load_or_migrate_keyring_at(identity_path(), CONFIG.load().keyring.clone()) } + async fn sync_pending_invitation_revocations(&self) -> Result<(), String> { + let invitations = iota_storage::users::invitations::list() + .map_err(|error| format!("failed to load pending invitation actions: {error}"))?; + for invitation in invitations.into_iter().filter(|invitation| { + invitation.authority == iota_storage::users::invitations::InvitationAuthority::Omega + && invitation.pending_action + == Some(iota_storage::users::invitations::PendingAction::Revoke) + }) { + let request = CommunicationValue::new(CommunicationType::RevokeUserInvitation) + .add_typed_default( + DataType::InvitationAuthority, + DataValue::Str("omega".into()), + ) + .add_typed_default( + DataType::InvitationId, + DataValue::SignedNumber(invitation.invitation_id.into()), + ); + let response = self + .await_response(&request, Some(Duration::from_secs(20))) + .await?; + if response.is_type(CommunicationType::Success) { + let Some(revision) = response + .get_data(DataType::InvitationRevision) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()) + else { + return Err("Omega revoke response omitted invitation revision".into()); + }; + let state = match response.get_data(DataType::InvitationState).as_str() { + Some("pending") => iota_storage::users::invitations::InvitationState::Pending, + Some("provisioning") => { + iota_storage::users::invitations::InvitationState::Provisioning + } + Some("redeemed") => iota_storage::users::invitations::InvitationState::Redeemed, + Some("revoked") => iota_storage::users::invitations::InvitationState::Revoked, + Some("expired") => iota_storage::users::invitations::InvitationState::Expired, + _ => return Err("Omega revoke response contained invalid state".into()), + }; + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as i64; + iota_storage::users::invitations::apply_revoke_result( + invitation.invitation_id, + state, + revision, + (state == iota_storage::users::invitations::InvitationState::Revoked) + .then_some(now), + now, + ) + .map_err(|error| format!("failed to apply invitation revoke result: {error}"))?; + } else { + return Err("Omega returned the wrong invitation revoke response type".into()); + } + } + Ok(()) + } + + pub async fn flush_pending_invitation_actions(&self) { + if self + .invitation_sync + .flush_in_progress + .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire) + .is_err() + { + return; + } + let now = Instant::now(); + if self + .invitation_sync + .next_retry_at + .lock() + .await + .is_some_and(|retry_at| retry_at > now) + { + self.invitation_sync + .flush_in_progress + .store(false, Ordering::Release); + return; + } + match self.sync_pending_invitation_revocations().await { + Ok(()) => { + *self.invitation_sync.next_retry_at.lock().await = None; + *self.invitation_sync.retry_delay.lock().await = Duration::from_secs(5); + } + Err(error) => { + log!("Pending invitation action flush failed: {}", error); + let mut delay = self.invitation_sync.retry_delay.lock().await; + *self.invitation_sync.next_retry_at.lock().await = Some(now + *delay); + *delay = next_invitation_retry_delay(*delay); + } + } + self.invitation_sync + .flush_in_progress + .store(false, Ordering::Release); + } + + pub async fn sync_omega_invitations(&self) -> Result<(), OmikronError> { + let request = CommunicationValue::new(CommunicationType::ListUserInvitations) + .add_typed_default( + DataType::InvitationAuthority, + DataValue::Str("omega".into()), + ); + let response = self + .await_response(&request, Some(Duration::from_secs(20))) + .await + .map_err(map_await_response_error)?; + if !response.is_type(CommunicationType::ListUserInvitations) { + return Err(OmikronError::Internal( + "Omega returned the wrong invitation snapshot response type".into(), + )); + } + let Some(DataValue::Array(remote)) = response.get_data(DataType::Invitations) else { + return Err(OmikronError::Internal( + "Omega invitation snapshot omitted Invitations".into(), + )); + }; + let type_map = TypeMap::latest(); + let synced_at = now_millis_i64(); + let expected_iota_id = CONFIG + .load() + .iota_id + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0) + .ok_or_else(|| OmikronError::Internal("Iota identity is not configured".into()))?; + let mut invitations = Vec::with_capacity(remote.len()); + for invitation in remote { + let DataValue::Container(fields) = invitation else { + return Err(OmikronError::Internal( + "Omega invitation snapshot contained a non-container record".into(), + )); + }; + invitations.push(parse_omega_invitation( + fields, + &type_map, + expected_iota_id, + synced_at, + )?); + } + iota_storage::users::invitations::merge_omega_snapshot(&invitations, synced_at).map_err( + |error| { + OmikronError::Storage(format!( + "failed to persist Omega invitation snapshot: {error}" + )) + }, + )?; + Ok(()) + } + // ------------------------------------------------------------------------- // Omikron discovery (via Omega's HTTP API, replacing the static // host/port/public-key-file model) @@ -827,6 +1118,10 @@ impl OmikronConnection { self.classify_legacy_pending_relays().await; self.flush_pending_relays().await; + let invitation_self = self.clone(); + tokio::spawn(async move { + invitation_self.flush_pending_invitation_actions().await; + }); if let Err(error) = relay_replay::prune_completed( now_millis_i64().saturating_sub(RELAY_RETENTION_MILLIS), ) { @@ -2407,14 +2702,22 @@ impl OmikronConnection { let public_key = cv.get_data(DataType::PublicKey).as_str().map(str::to_owned); let invitation_id = cv .get_data(DataType::InvitationId) - .as_str() - .filter(|value| uuid::Uuid::parse_str(value).is_ok()) - .map(str::to_owned); - let Some((user_id, username, public_key, invitation_id)) = user_id + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0); + let invitation_revision = cv + .get_data(DataType::InvitationRevision) + .as_signed_number() + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0); + let Some((user_id, username, public_key, invitation_id, invitation_revision)) = user_id .zip(username) .zip(public_key) .zip(invitation_id) - .map(|(((id, username), key), invitation_id)| (id, username, key, invitation_id)) + .zip(invitation_revision) + .map(|((((id, username), key), invitation_id), revision)| { + (id, username, key, invitation_id, revision) + }) else { let _ = self .send_message(&error_response(cv, CommunicationType::ErrorInvalidData)) @@ -2424,22 +2727,65 @@ impl OmikronConnection { let profile = iota_storage::users::user_profile::UserProfile::new( user_id, username, None, public_key, None, None, ); - if iota_storage::users::user_manager::try_add_user_with_credential_origin( - profile, - iota_storage::users::user_manager::CredentialOrigin::External, - ) - .is_err() - { - let _ = self - .send_message(&error_response(cv, CommunicationType::ErrorInternal)) - .await; - return; + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as i64; + let mut result = iota_storage::users::invitations::apply_external_invitation_provisioning( + invitation_id, + invitation_revision, + &profile, + now, + ); + if matches!( + result, + Ok(iota_storage::users::invitations::ProvisioningResult::MissingInvitation) + ) { + if self.sync_omega_invitations().await.is_ok() { + result = iota_storage::users::invitations::apply_external_invitation_provisioning( + invitation_id, + invitation_revision, + &profile, + now, + ); + } + } + match result { + Ok( + iota_storage::users::invitations::ProvisioningResult::Created + | iota_storage::users::invitations::ProvisioningResult::AlreadyApplied, + ) => {} + Ok(iota_storage::users::invitations::ProvisioningResult::RevocationPending) => { + self.flush_pending_invitation_actions().await; + return; + } + Ok(iota_storage::users::invitations::ProvisioningResult::Conflict) => { + let _ = self + .send_message(&error_response(cv, CommunicationType::ErrorInvalidData)) + .await; + return; + } + Ok(iota_storage::users::invitations::ProvisioningResult::MissingInvitation) => { + let _ = self + .send_message(&error_response(cv, CommunicationType::ErrorInvalidData)) + .await; + return; + } + Err(_) => { + let _ = self + .send_message(&error_response(cv, CommunicationType::ErrorInternal)) + .await; + return; + } } let acknowledgement = CommunicationValue::new(CommunicationType::AcknowledgeIotaUserProvision) .with_request_id(cv) .add_typed_default(DataType::UserId, DataValue::SignedNumber(user_id.into())) - .add_typed_default(DataType::InvitationId, DataValue::Str(invitation_id)); + .add_typed_default( + DataType::InvitationId, + DataValue::SignedNumber(invitation_id.into()), + ); let _ = self.send_message(&acknowledgement).await; } @@ -3555,6 +3901,14 @@ impl OmikronClient for OmikronConnection { .map_err(map_await_response_error) } + async fn sync_omega_invitations(&self) -> Result<(), OmikronError> { + Self::sync_omega_invitations(self).await + } + + async fn flush_pending_invitation_actions(&self) { + Self::flush_pending_invitation_actions(self).await; + } + async fn reconnect(&self) -> Result<(), OmikronError> { let this = Arc::new(Self { state: self.state.clone(), @@ -3574,6 +3928,7 @@ impl OmikronClient for OmikronConnection { cancellation: self.cancellation.clone(), active_tasks: self.active_tasks.clone(), app: self.app.clone(), + invitation_sync: self.invitation_sync.clone(), }); Self::reconnect(&this).await; Ok(()) @@ -3598,6 +3953,7 @@ impl OmikronClient for OmikronConnection { cancellation: self.cancellation.clone(), active_tasks: self.active_tasks.clone(), app: self.app.clone(), + invitation_sync: self.invitation_sync.clone(), }); Self::rotate_identity(&this).await } @@ -3704,4 +4060,87 @@ mod tests { assert!(recipient_block_policy_applies(true, 42, 43)); assert!(!recipient_block_policy_applies(false, 42, 43)); } + + #[test] + fn omega_invitation_snapshot_parser_requires_authoritative_fields() { + let type_map = TypeMap::latest(); + let mut fields = Vec::new(); + for (kind, value) in [ + ( + DataType::InvitationAuthority, + DataValue::Str("omega".into()), + ), + (DataType::InvitationId, DataValue::SignedNumber(7)), + (DataType::IotaId, DataValue::SignedNumber(42)), + (DataType::InvitationRevision, DataValue::SignedNumber(3)), + ( + DataType::InvitationState, + DataValue::Str("provisioning".into()), + ), + (DataType::InvitationCreatedAt, DataValue::SignedNumber(10)), + (DataType::InvitationExpiresAt, DataValue::SignedNumber(20)), + ( + DataType::InvitationPasswordProtected, + DataValue::Bool(false), + ), + ] { + fields.push((kind.try_to_id(&type_map).expect("type is mapped"), value)); + } + let parsed = parse_omega_invitation(&fields, &type_map, 42, 20).expect("snapshot parses"); + assert_eq!(parsed.invitation_id, 7); + assert_eq!(parsed.remote_revision, 3); + assert_eq!( + parsed.authority, + iota_storage::users::invitations::InvitationAuthority::Omega + ); + + fields.retain(|(id, _)| Some(*id) != DataType::InvitationRevision.try_to_id(&type_map)); + assert!(parse_omega_invitation(&fields, &type_map, 42, 20).is_err()); + } + + #[test] + fn omega_invitation_snapshot_parser_rejects_wrong_authority_or_iota() { + let type_map = TypeMap::latest(); + let fields = [ + (DataType::InvitationAuthority, DataValue::Str("iota".into())), + (DataType::InvitationId, DataValue::SignedNumber(7)), + (DataType::IotaId, DataValue::SignedNumber(42)), + (DataType::InvitationRevision, DataValue::SignedNumber(3)), + (DataType::InvitationState, DataValue::Str("pending".into())), + (DataType::InvitationCreatedAt, DataValue::SignedNumber(10)), + (DataType::InvitationExpiresAt, DataValue::SignedNumber(20)), + ( + DataType::InvitationPasswordProtected, + DataValue::Bool(false), + ), + ] + .into_iter() + .map(|(kind, value)| (kind.try_to_id(&type_map).expect("type is mapped"), value)) + .collect::>(); + + assert!(parse_omega_invitation(&fields, &type_map, 42, 20).is_err()); + + let mut fields = fields; + let authority_id = DataType::InvitationAuthority + .try_to_id(&type_map) + .expect("type is mapped"); + let authority = fields + .iter_mut() + .find(|(id, _)| *id == authority_id) + .expect("authority field exists"); + authority.1 = DataValue::Str("omega".into()); + assert!(parse_omega_invitation(&fields, &type_map, 43, 20).is_err()); + } + + #[test] + fn invitation_action_retry_delay_backs_off_and_caps() { + assert_eq!( + next_invitation_retry_delay(Duration::from_secs(5)), + Duration::from_secs(10) + ); + assert_eq!( + next_invitation_retry_delay(Duration::from_secs(300)), + Duration::from_secs(300) + ); + } } diff --git a/scripts/build-release-bundle.sh b/scripts/build-release-bundle.sh index 779744d..3ea6f05 100644 --- a/scripts/build-release-bundle.sh +++ b/scripts/build-release-bundle.sh @@ -1,17 +1,20 @@ #!/usr/bin/env bash set -euo pipefail -if [[ "$#" -ne 6 ]]; then - echo "usage: $0 BINARY_DIRECTORY PRODUCT_VERSION UPDATE_MANIFEST_URL UPDATE_PUBLIC_KEY UPDATE_SIGNATURE_URL OUTPUT.zip" >&2 +if [[ "$#" -ne 9 ]]; then + echo "usage: $0 BINARY_DIRECTORY PRODUCT_VERSION RELEASE_MANIFEST UPDATE_MANIFEST_URL UPDATE_PUBLIC_KEY UPDATE_SIGNATURE_URL UPDATE_CHANNEL UPDATE_SIGNING_KEY_ID OUTPUT.zip" >&2 exit 2 fi binary_directory="$1" product_version="$2" -update_manifest_url="$3" -update_public_key="$4" -update_signature_url="$5" -output="$6" +release_manifest="$3" +update_manifest_url="$4" +update_public_key="$5" +update_signature_url="$6" +update_channel="$7" +update_signing_key_id="$8" +output="$9" repository_directory="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" contract="$repository_directory/iota-installer/bundle-files.txt" staging_directory="$(mktemp -d)" @@ -21,6 +24,20 @@ if [[ ! "$product_version" =~ ^[0-9A-Za-z][0-9A-Za-z.+_-]*$ ]]; then echo "product version contains unsupported characters: $product_version" >&2 exit 2 fi +if ! jq -e \ + --arg product_version "$product_version" \ + --arg channel "$update_channel" \ + --arg key_id "$update_signing_key_id" \ + '.product_version == $product_version + and .channel == $channel + and .release_signing_key_id == $key_id + and (.release_sequence | type == "number" and . > 0) + and (.published_at | type == "string" and length > 0) + and (.expires_at | type == "string" and length > 0)' \ + "$release_manifest" >/dev/null; then + echo "release manifest identity or anti-rollback metadata is invalid" >&2 + exit 2 +fi mkdir -p "$(dirname "$output")" output="$(cd "$(dirname "$output")" && pwd)/$(basename "$output")" @@ -37,13 +54,15 @@ while IFS= read -r bundle_path; do install -m 0755 "$binary_directory/${bundle_path#bin/}" "$destination" ;; manifest.json) - printf '{"product_version":"%s"}\n' "$product_version" > "$destination" + install -m 0644 "$release_manifest" "$destination" ;; systemd/update.env) - printf 'IOTA_UPDATE_MANIFEST=%s\nIOTA_UPDATE_PUBLIC_KEY=%s\nIOTA_UPDATE_SIGNATURE=%s\n' \ + printf 'IOTA_UPDATE_MANIFEST=%s\nIOTA_UPDATE_PUBLIC_KEY=%s\nIOTA_UPDATE_SIGNATURE=%s\nIOTA_UPDATE_CHANNEL=%s\nIOTA_UPDATE_SIGNING_KEY_ID=%s\n' \ "$update_manifest_url" \ "$update_public_key" \ "$update_signature_url" \ + "$update_channel" \ + "$update_signing_key_id" \ > "$destination" ;; *) diff --git a/scripts/build-update-manifest.sh b/scripts/build-update-manifest.sh index 456b3e7..40740de 100644 --- a/scripts/build-update-manifest.sh +++ b/scripts/build-update-manifest.sh @@ -1,25 +1,38 @@ #!/usr/bin/env bash set -euo pipefail -if [[ "$#" -ne 8 ]]; then - echo "usage: $0 BINARY_DIRECTORY PRODUCT_VERSION CHANNEL PUBLISHED_AT OS ARCHITECTURE BASE_URL OUTPUT.json" >&2 +if [[ "$#" -ne 10 ]]; then + echo "usage: $0 BINARY_DIRECTORY PRODUCT_VERSION CHANNEL RELEASE_SEQUENCE PUBLISHED_AT EXPIRES_AT OS ARCHITECTURE BASE_URL OUTPUT.json" >&2 exit 2 fi binary_directory="$1" product_version="$2" channel="$3" -published_at="$4" -operating_system="$5" -architecture="$6" -base_url="$(printf '%s' "$7" | sed 's#/$##')" -output="$8" +release_sequence="$4" +published_at="$5" +expires_at="$6" +operating_system="$7" +architecture="$8" +base_url="$(printf '%s' "$9" | sed 's#/$##')" +output="${10}" repository_directory="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" contract="$repository_directory/iota-updater/artifacts.tsv" staging_directory="$(mktemp -d)" artifacts="$staging_directory/artifacts.jsonl" trap 'rm -rf "$staging_directory"' EXIT +if [[ ! "$release_sequence" =~ ^[1-9][0-9]*$ ]]; then + echo "release sequence must be a positive integer" >&2 + exit 2 +fi +published_epoch="$(date -d "$published_at" +%s)" +expires_epoch="$(date -d "$expires_at" +%s)" +if (( expires_epoch <= published_epoch )); then + echo "expires_at must be after published_at" >&2 + exit 2 +fi + while IFS=$'\t' read -r role artifact_path; do [[ -n "$role" && -n "$artifact_path" ]] || continue asset_name="$(basename "$artifact_path")" @@ -40,12 +53,16 @@ mkdir -p "$(dirname "$output")" jq -s \ --arg product_version "$product_version" \ --arg channel "$channel" \ + --argjson release_sequence "$release_sequence" \ --arg published_at "$published_at" \ + --arg expires_at "$expires_at" \ --arg release_signing_key_id "primary" \ '{ product_version: $product_version, channel: $channel, + release_sequence: $release_sequence, published_at: $published_at, + expires_at: $expires_at, minimum_data_schema: 1, supported_ipc_min: 2, supported_ipc_max: 4, diff --git a/web-ui/src/api.rs b/web-ui/src/api.rs index c2c52de..ce95302 100755 --- a/web-ui/src/api.rs +++ b/web-ui/src/api.rs @@ -226,6 +226,7 @@ fn ipc_error_response(code: IpcErrorCode) -> HttpResponse { let status = match code { IpcErrorCode::InvalidRequest => actix_web::http::StatusCode::BAD_REQUEST, IpcErrorCode::Conflict => actix_web::http::StatusCode::CONFLICT, + IpcErrorCode::Unsupported => actix_web::http::StatusCode::NOT_IMPLEMENTED, IpcErrorCode::NotReady | IpcErrorCode::OmikronUnavailable => { actix_web::http::StatusCode::SERVICE_UNAVAILABLE }