From e19c3c3d121c514dfc39cb73cfe01d56c820ff0b Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Thu, 24 Sep 2026 18:37:28 +0200 Subject: [PATCH] [Fix] Bound Iota storage, relay and transport resources --- Cargo.lock | 2 + README.md | 22 ++++ iota-connection/src/message_handlers.rs | 77 +++++++++++- iota-connection/src/relay.rs | 61 +++++++++- iota-daemon-lib/src/ipc_server.rs | 10 +- iota-daemon/src/main.rs | 1 + iota-identity/Cargo.toml | 1 + iota-identity/src/lib.rs | 36 ++++++ iota-storage/Cargo.toml | 1 + iota-storage/src/storage_error.rs | 2 + iota-storage/src/util/config_util.rs | 104 ++++++++++++++++ iota-storage/src/util/relay_replay.rs | 3 + iota-storage/src/util/user_assets.rs | 93 +++++++++++++-- iota-storage/src/util/user_blobs.rs | 52 +++++++- iota-storage/tests/asset_quota.rs | 50 ++++++++ omikron-connector/src/omikron_connection.rs | 124 ++++++++++++++++---- web-server/src/lib.rs | 13 +- web-ui/Cargo.toml | 3 + web-ui/src/lib.rs | 2 + 19 files changed, 609 insertions(+), 48 deletions(-) create mode 100644 iota-storage/tests/asset_quota.rs diff --git a/Cargo.lock b/Cargo.lock index a15b555..ef01e01 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2255,6 +2255,7 @@ version = "0.1.0" dependencies = [ "async-trait", "base64 0.22.1", + "libc", "mtp", "serde", "serde_json", @@ -2326,6 +2327,7 @@ dependencies = [ "arc-swap", "async-trait", "base64 0.22.1", + "fs2", "iota-identity", "iota-logger", "iota-paths", diff --git a/README.md b/README.md index 5a8c16a..0e619af 100644 --- a/README.md +++ b/README.md @@ -66,3 +66,25 @@ usermod -aG iota-operators USER The user must start a new login session before supplementary group membership is visible. Unix per-user deployments must set `IOTA_SOCKET` to an absolute path; Iota does not derive its IPC socket from `XDG_RUNTIME_DIR`. + +## Asset storage limits + +`storage_limits` in Iota's configuration controls upload admission. Defaults are 256 MiB per asset, 2 GiB of committed assets plus active upload reservations per user, four active uploads per user, 512 MiB of free filesystem space reserved, four concurrent asset I/O workers, and 4096 blobs per user. Set these fields to match deployment capacity: + +```yaml +storage_limits: + max_asset_bytes: 268435456 + max_user_asset_bytes: 2147483648 + max_active_asset_uploads_per_user: 4 + min_free_asset_storage_bytes: 536870912 + max_asset_io_workers: 4 + max_user_blobs: 4096 +``` + +Replacement uploads temporarily count both old and new assets until commit. + +## Relay freshness + +Signed relays expire after 30 days less five minutes; replay rows remain for 30 days. `max_relay_future_skew_millis` defaults to 300000, or five minutes, and can be reduced for deployment clock tolerance. Fixed five-minute margin keeps freshness within replay retention even if configuration changes. Timestamps use Unix milliseconds. + +Asset and blob list responses contain at most 128 entries. Send the returned positive `Offset` as the next request cursor; `Offset: 0` ends pagination. `web.max_mtp_sessions` defaults to 256 and limits concurrent authenticated MTP sessions. Legacy web administration requires the non-default `legacy-web-admin` feature. diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index 07983aa..32fa3c0 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -1,4 +1,5 @@ use crate::message_common::*; +use iota_logger::log; use iota_storage::util::chat_files::{self, MessageState}; use iota_storage::util::chats_util::{self, get_user, has_user, mod_user}; use iota_storage::util::communities_util::CommunitiesUtil; @@ -2084,8 +2085,50 @@ fn upload_response( } fn asset_error(cv: &CommunicationValue, error: StorageError) -> CommunicationValue { - error_response(cv, CommunicationType::ErrorInternal) - .add_typed_default(DataType::ErrorType, DataValue::Str(error.to_string())) + log!( + "Asset request rejected sender={:?} request_id={:?}: {error}", + cv.sender(), + cv.id() + ); + let (kind, category) = match error { + StorageError::AssetResourceLimit(_) => { + (CommunicationType::ErrorInvalidData, "asset_resource_limit") + } + _ => (CommunicationType::ErrorInternal, "asset_request_failed"), + }; + error_response(cv, kind).add_typed_default(DataType::ErrorType, DataValue::Str(category.into())) +} + +#[cfg(test)] +mod asset_error_tests { + use super::*; + + #[test] + fn storage_errors_do_not_expose_internal_paths() { + let request = CommunicationValue::new(CommunicationType::UserAssetUploadStart).with_id(7); + let internal = asset_error( + &request, + StorageError::Other("private path /srv/assets/key".into()), + ); + assert!(internal.is_type(CommunicationType::ErrorInternal)); + assert_eq!( + internal + .get_data(DataType::ErrorType) + .and_then(DataValue::as_str), + Some("asset_request_failed") + ); + let limit = asset_error( + &request, + StorageError::AssetResourceLimit("max_asset_bytes"), + ); + assert!(limit.is_type(CommunicationType::ErrorInvalidData)); + assert_eq!( + limit + .get_data(DataType::ErrorType) + .and_then(DataValue::as_str), + Some("asset_resource_limit") + ); + } } fn upload_request(cv: &CommunicationValue) -> Result<(i64, String), CommunicationValue> { @@ -2404,10 +2447,21 @@ pub fn handle_user_blob_list(cv: &CommunicationValue) -> CommunicationValue { if user_id <= 0 { return error_response(cv, CommunicationType::ErrorInvalidData); } - match user_blobs::list_metadata(user_id) { - Ok(blobs) => CommunicationValue::new(CommunicationType::UserBlobList) + let after_id = match cv.get_data(DataType::Offset).and_then(DataValue::as_number) { + Some(raw) => match i64::try_from(raw) { + Ok(id) if id > 0 => Some(id), + _ => return error_response(cv, CommunicationType::ErrorInvalidData), + }, + None => None, + }; + match user_blobs::list_metadata_page(user_id, after_id) { + Ok((blobs, next)) => CommunicationValue::new(CommunicationType::UserBlobList) .with_request_id(cv) .with_receiver(sender_wire_id(user_id)) + .add_typed_default( + DataType::Offset, + DataValue::SignedNumber(next.unwrap_or(0).into()), + ) .add_typed_default( DataType::Blobs, DataValue::Array(blobs.iter().map(blob_metadata_value).collect()), @@ -2542,10 +2596,21 @@ pub fn handle_user_asset_list(cv: &CommunicationValue) -> CommunicationValue { Ok(id) if id > 0 => id, _ => return error_response(cv, CommunicationType::ErrorInvalidData), }; - match user_assets::list(user_id) { - Ok(assets) => CommunicationValue::new(CommunicationType::UserAssetList) + let after_id = match cv.get_data(DataType::Offset).and_then(DataValue::as_number) { + Some(raw) => match i64::try_from(raw) { + Ok(id) if id > 0 => Some(id), + _ => return error_response(cv, CommunicationType::ErrorInvalidData), + }, + None => None, + }; + match user_assets::list_page(user_id, after_id) { + Ok((assets, next)) => CommunicationValue::new(CommunicationType::UserAssetList) .with_request_id(cv) .with_receiver(sender_wire_id(user_id)) + .add_typed_default( + DataType::Offset, + DataValue::SignedNumber(next.unwrap_or(0).into()), + ) .add_typed_default( DataType::Assets, DataValue::Array(assets.iter().map(asset_metadata_value).collect()), diff --git a/iota-connection/src/relay.rs b/iota-connection/src/relay.rs index 5ffcb24..6d05d5a 100644 --- a/iota-connection/src/relay.rs +++ b/iota-connection/src/relay.rs @@ -1,3 +1,7 @@ +use iota_storage::util::{ + config_util::CONFIG, + relay_replay::{MAX_RELAY_FUTURE_SKEW_MILLIS, RELAY_RETENTION_MILLIS}, +}; use iota_util::route_target::RouteTarget; use mtp::codec::{ CommunicationValue, ProtectionPolicy, RelayError, RelayOpenOptions, SignaturePolicy, TypeMap, @@ -107,6 +111,9 @@ pub enum RelayValidationError { MissingTypeMap, InvalidRouteTarget(u64), KeyLookup(String), + Clock(String), + StaleRelay { created_at: u64, now: u64 }, + RelayFromFuture { created_at: u64, now: u64 }, Relay(RelayError), } @@ -129,6 +136,14 @@ impl fmt::Display for RelayValidationError { write!(formatter, "relay has invalid route target {target}") } Self::KeyLookup(error) => write!(formatter, "trusted signer lookup failed: {error}"), + Self::Clock(error) => write!(formatter, "relay clock failed: {error}"), + Self::StaleRelay { created_at, now } => { + write!(formatter, "relay timestamp {created_at} is stale at {now}") + } + Self::RelayFromFuture { created_at, now } => write!( + formatter, + "relay timestamp {created_at} is in the future at {now}" + ), Self::Relay(error) => error.fmt(formatter), } } @@ -142,6 +157,23 @@ impl From for RelayValidationError { } } +fn validate_created_at( + created_at: u64, + now: u64, + max_future_skew_millis: u64, +) -> Result<(), RelayValidationError> { + if created_at > now.saturating_add(max_future_skew_millis) { + return Err(RelayValidationError::RelayFromFuture { created_at, now }); + } + let max_age = u64::try_from(RELAY_RETENTION_MILLIS) + .unwrap_or(0) + .saturating_sub(MAX_RELAY_FUTURE_SKEW_MILLIS); + if now.saturating_sub(created_at) > max_age { + return Err(RelayValidationError::StaleRelay { created_at, now }); + } + Ok(()) +} + /* * Relay metadata is opened only after the claimed signer selects trusted key * history. Replay reservation happens after verification and durable @@ -195,6 +227,14 @@ where RelayOpenOptions::new(RELAY_PROTECTION_POLICY), )?; + let now = mtp::common::unix_time_millis() + .map_err(|error| RelayValidationError::Clock(error.to_string()))?; + validate_created_at( + metadata.created_at(), + now, + CONFIG.load().max_relay_future_skew_millis, + )?; + let context = VerifiedRelayContext { signer_id: metadata.signer_id(), final_recipient_id: metadata.final_recipient_id(), @@ -244,6 +284,23 @@ mod tests { use mtp::codec::SealedRelayBuilder; use mtp::crypto::{DualSigner, Ed25519Signer, Keyring}; + #[test] + fn relay_freshness_has_bounded_future_and_retention() { + let now = 1_000_000_000_000_u64; + let max_age = u64::try_from(RELAY_RETENTION_MILLIS).unwrap_or(0) - 300_000; + assert!(validate_created_at(now, now, 300_000).is_ok()); + assert!(validate_created_at(now - max_age, now, 300_000).is_ok()); + assert!(matches!( + validate_created_at(now - max_age - 1, now, 300_000), + Err(RelayValidationError::StaleRelay { .. }) + )); + assert!(validate_created_at(now + 300_000, now, 300_000).is_ok()); + assert!(matches!( + validate_created_at(now + 300_001, now, 300_000), + Err(RelayValidationError::RelayFromFuture { .. }) + )); + } + fn relay(message_id: u128) -> Result<(Keyring, Keyring, CommunicationValue), String> { let signer_keyring = Keyring::generate(); let recipient_keyring = Keyring::generate(); @@ -264,7 +321,7 @@ mod tests { &signer, ) .message_id(message_id) - .created_at(123) + .created_at(mtp::common::unix_time_millis().map_err(|error| error.to_string())?) .metadata_recipients(vec![recipient_keyring.public_key_bundle()]) .content_recipients(vec![recipient_keyring.public_key_bundle()]) .build() @@ -339,7 +396,7 @@ mod tests { &signer, ) .message_id(7_u128) - .created_at(123) + .created_at(mtp::common::unix_time_millis().map_err(|error| error.to_string())?) .metadata_recipients(vec![recipient_keyring.public_key_bundle()]) .content_recipients(vec![recipient_keyring.public_key_bundle()]) .build() diff --git a/iota-daemon-lib/src/ipc_server.rs b/iota-daemon-lib/src/ipc_server.rs index 1e597c6..b90e71f 100644 --- a/iota-daemon-lib/src/ipc_server.rs +++ b/iota-daemon-lib/src/ipc_server.rs @@ -327,11 +327,11 @@ fn peer_credentials(stream: &UnixStream) -> Result { } #[cfg(not(target_os = "linux"))] { - Ok(PeerIdentity { - pid: 0, - uid: 0, - _gid: 0, - }) + let _ = stream; + Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "IPC peer credentials are unsupported on this platform", + )) } } diff --git a/iota-daemon/src/main.rs b/iota-daemon/src/main.rs index 910a983..d746017 100644 --- a/iota-daemon/src/main.rs +++ b/iota-daemon/src/main.rs @@ -576,6 +576,7 @@ async fn main() -> ExitCode { key: resolve_config_path(&paths.config_file, &key, &paths.config_dir), }), required: web.required, + max_mtp_sessions: web.max_mtp_sessions, authority_discovery: iota_identity::AuthorityDiscoveryDocument { version: 1, service: iota_identity::AuthorityKind::Iota, diff --git a/iota-identity/Cargo.toml b/iota-identity/Cargo.toml index 7508211..d7e56eb 100644 --- a/iota-identity/Cargo.toml +++ b/iota-identity/Cargo.toml @@ -4,6 +4,7 @@ version = "0.1.0" edition = "2024" [dependencies] +libc = "0.2" async-trait = "0.1.89" base64 = "0.22.1" mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto", "files", "raw"] } diff --git a/iota-identity/src/lib.rs b/iota-identity/src/lib.rs index 468581a..02aeab2 100644 --- a/iota-identity/src/lib.rs +++ b/iota-identity/src/lib.rs @@ -179,6 +179,34 @@ impl fmt::Display for LocalNodeIdentityError { impl std::error::Error for LocalNodeIdentityError {} +#[cfg(unix)] +fn verify_private_keyring(path: &Path) -> std::io::Result<()> { + use std::os::unix::fs::MetadataExt; + let metadata = match fs::symlink_metadata(path) { + Ok(metadata) => metadata, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()), + Err(error) => return Err(error), + }; + if !metadata.file_type().is_file() + || metadata.mode() & 0o077 != 0 + || metadata.uid() != unsafe { libc::geteuid() } + { + return Err(std::io::Error::new( + std::io::ErrorKind::PermissionDenied, + "keyring must be a regular owner-only file owned by the service user", + )); + } + Ok(()) +} + +#[cfg(not(unix))] +fn verify_private_keyring(_: &Path) -> std::io::Result<()> { + Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "keyring permissions cannot be verified on this platform", + )) +} + #[derive(Clone)] pub struct LocalNodeIdentity { keyring: Arc, @@ -210,6 +238,10 @@ impl LocalNodeIdentity { source, })?; } + verify_private_keyring(path).map_err(|error| LocalNodeIdentityError::Storage { + path: path.display().to_string(), + source: error.to_string(), + })?; let keyring = match mtp::files::load_keyring_raw(path) { Ok(keyring) => keyring, Err(mtp::files::FileError::Io(error)) @@ -250,6 +282,10 @@ impl LocalNodeIdentity { source: error.to_string(), } })?; + verify_private_keyring(path).map_err(|error| LocalNodeIdentityError::Storage { + path: path.display().to_string(), + source: error.to_string(), + })?; let persisted = mtp::files::load_keyring_raw(path).map_err(|error| { LocalNodeIdentityError::Storage { path: path.display().to_string(), diff --git a/iota-storage/Cargo.toml b/iota-storage/Cargo.toml index 6723914..ff6ce06 100644 --- a/iota-storage/Cargo.toml +++ b/iota-storage/Cargo.toml @@ -21,6 +21,7 @@ sha2 = "0.10" thiserror = "2" rand = "0.8" rusqlite = "0.40.0" +fs2 = "0.4.3" tokio = { version = "1.50.0", features = ["full"] } [dev-dependencies] diff --git a/iota-storage/src/storage_error.rs b/iota-storage/src/storage_error.rs index f407b7e..2a2f053 100644 --- a/iota-storage/src/storage_error.rs +++ b/iota-storage/src/storage_error.rs @@ -14,6 +14,8 @@ pub enum StorageError { RevisionConflict, #[error("pending relay ownership is unknown")] PendingRelayOwnershipUnknown, + #[error("asset resource limit: {0}")] + AssetResourceLimit(&'static str), #[error("{0}")] Other(String), } diff --git a/iota-storage/src/util/config_util.rs b/iota-storage/src/util/config_util.rs index 1e01050..b7040dd 100644 --- a/iota-storage/src/util/config_util.rs +++ b/iota-storage/src/util/config_util.rs @@ -40,8 +40,14 @@ pub enum ConfigError { InvalidRelayRouterKey(String), #[error("relay router certificate path must not be empty")] MissingRelayRouterCertificate, + #[error("web.max_mtp_sessions must be greater than zero")] + InvalidMaxMtpSessions, #[error("max_ipc_clients must be greater than zero")] InvalidMaxIpcClients, + #[error("invalid storage limit: {0}")] + InvalidStorageLimit(&'static str), + #[error("invalid max_relay_future_skew_millis")] + InvalidRelayFutureSkew, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -71,6 +77,44 @@ pub struct IotaConfig { pub read_receipts_enabled: bool, #[serde(default = "default_max_ipc_clients")] pub max_ipc_clients: usize, + #[serde(default)] + pub storage_limits: StorageLimits, + #[serde(default = "default_max_relay_future_skew_millis")] + pub max_relay_future_skew_millis: u64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct StorageLimits { + pub max_asset_bytes: i64, + pub max_user_asset_bytes: i64, + #[serde(default = "default_max_user_blobs")] + pub max_user_blobs: i64, + pub max_active_asset_uploads_per_user: usize, + pub min_free_asset_storage_bytes: u64, + #[serde(default = "default_max_asset_io_workers")] + pub max_asset_io_workers: usize, +} + +const fn default_max_user_blobs() -> i64 { + 4096 +} + +const fn default_max_asset_io_workers() -> usize { + 4 +} + +impl Default for StorageLimits { + fn default() -> Self { + Self { + max_asset_bytes: 256 * 1024 * 1024, + max_user_asset_bytes: 2 * 1024 * 1024 * 1024, + max_user_blobs: default_max_user_blobs(), + max_active_asset_uploads_per_user: 4, + min_free_asset_storage_bytes: 512 * 1024 * 1024, + max_asset_io_workers: default_max_asset_io_workers(), + } + } } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -109,11 +153,16 @@ pub struct WebSettings { pub key: Option, #[serde(default)] pub required: bool, + #[serde(default = "default_max_mtp_sessions")] + pub max_mtp_sessions: usize, #[serde(default)] pub direct_endpoints: Vec, #[serde(default)] pub relay_hints: Vec, } +const fn default_max_mtp_sessions() -> usize { + 256 +} fn default_web_bind() -> String { "127.0.0.1".into() } @@ -130,6 +179,7 @@ impl Default for WebSettings { certificate: None, key: None, required: false, + max_mtp_sessions: default_max_mtp_sessions(), direct_endpoints: Vec::new(), relay_hints: Vec::new(), } @@ -148,6 +198,10 @@ const fn default_max_ipc_clients() -> usize { 64 } +const fn default_max_relay_future_skew_millis() -> u64 { + 5 * 60 * 1_000 +} + impl Default for IotaConfig { fn default() -> Self { Self { @@ -163,6 +217,8 @@ impl Default for IotaConfig { private_key: None, read_receipts_enabled: default_read_receipts_enabled(), max_ipc_clients: default_max_ipc_clients(), + storage_limits: StorageLimits::default(), + max_relay_future_skew_millis: default_max_relay_future_skew_millis(), } } } @@ -211,9 +267,33 @@ pub fn validate_config(config: &IotaConfig) -> Result<(), ConfigError> { bind: config.web.bind.clone(), source, })?; + if config.storage_limits.max_user_blobs <= 0 { + return Err(ConfigError::InvalidStorageLimit("max_user_blobs")); + } + if config.web.max_mtp_sessions == 0 { + return Err(ConfigError::InvalidMaxMtpSessions); + } if config.max_ipc_clients == 0 { return Err(ConfigError::InvalidMaxIpcClients); } + if config.max_relay_future_skew_millis > super::relay_replay::MAX_RELAY_FUTURE_SKEW_MILLIS { + return Err(ConfigError::InvalidRelayFutureSkew); + } + let limits = &config.storage_limits; + if limits.max_asset_bytes <= 0 { + return Err(ConfigError::InvalidStorageLimit("max_asset_bytes")); + } + if limits.max_user_asset_bytes < limits.max_asset_bytes { + return Err(ConfigError::InvalidStorageLimit("max_user_asset_bytes")); + } + if limits.max_active_asset_uploads_per_user == 0 { + return Err(ConfigError::InvalidStorageLimit( + "max_active_asset_uploads_per_user", + )); + } + if limits.max_asset_io_workers == 0 { + return Err(ConfigError::InvalidStorageLimit("max_asset_io_workers")); + } for endpoint in config .web .direct_endpoints @@ -357,6 +437,30 @@ mod tests { use super::{ConfigError, IotaConfig, RelayRouterSettings, parse_config, validate_config}; use std::path::Path; + #[test] + fn asset_limits_reject_invalid_quota_and_upload_count() { + let mut config = IotaConfig::default(); + config.storage_limits.max_asset_bytes = 0; + assert!(matches!( + validate_config(&config), + Err(ConfigError::InvalidStorageLimit("max_asset_bytes")) + )); + config.storage_limits.max_asset_bytes = 10; + config.storage_limits.max_user_asset_bytes = 9; + assert!(matches!( + validate_config(&config), + Err(ConfigError::InvalidStorageLimit("max_user_asset_bytes")) + )); + config.storage_limits.max_user_asset_bytes = 10; + config.storage_limits.max_active_asset_uploads_per_user = 0; + assert!(matches!( + validate_config(&config), + Err(ConfigError::InvalidStorageLimit( + "max_active_asset_uploads_per_user" + )) + )); + } + #[test] fn malformed_yaml_is_rejected() { assert!(matches!( diff --git a/iota-storage/src/util/relay_replay.rs b/iota-storage/src/util/relay_replay.rs index ee4e476..65c5a95 100644 --- a/iota-storage/src/util/relay_replay.rs +++ b/iota-storage/src/util/relay_replay.rs @@ -2,6 +2,9 @@ use crate::storage_error::StorageError; use crate::util::db; use rusqlite::{OptionalExtension, params}; +pub const RELAY_RETENTION_MILLIS: i64 = 30 * 24 * 60 * 60 * 1_000; +pub const MAX_RELAY_FUTURE_SKEW_MILLIS: u64 = 5 * 60 * 1_000; + #[derive(Debug, Clone, PartialEq, Eq)] pub struct DeliveredRelay { pub signer_id: i64, diff --git a/iota-storage/src/util/user_assets.rs b/iota-storage/src/util/user_assets.rs index 0a9d3ea..93752b1 100644 --- a/iota-storage/src/util/user_assets.rs +++ b/iota-storage/src/util/user_assets.rs @@ -1,5 +1,5 @@ use crate::storage_error::StorageError; -use crate::util::{db, sync}; +use crate::util::{config_util::CONFIG, db, sync}; use rand::RngCore; use rusqlite::{OptionalExtension, Row, params}; use sha2::{Digest, Sha256}; @@ -118,7 +118,8 @@ fn reconcile_offset( if status.state != "uploading" { return Ok(status); } - let actual = fs::metadata(asset_dir()?.join(file_name))?.len() as i64; + let actual = i64::try_from(fs::metadata(asset_dir()?.join(file_name))?.len()) + .map_err(|_| StorageError::Other("upload file exceeds supported size".into()))?; if actual > status.size { return Err(StorageError::Other( "upload file exceeds declared size".into(), @@ -151,6 +152,11 @@ pub fn start( if size < 0 || sha256.len() != 32 || mime_type.is_empty() || mime_type.len() > 255 { return Err(StorageError::Other("invalid asset metadata".into())); } + let config = CONFIG.load(); + let limits = &config.storage_limits; + if size > limits.max_asset_bytes { + return Err(StorageError::AssetResourceLimit("max_asset_bytes")); + } cleanup_if_due()?; if let Some((status, file_name)) = load_upload(user_id, upload_id)? { if status.asset_id != asset_id @@ -166,15 +172,52 @@ pub fn start( } let file_name = random_file_name(); - let path = asset_dir()?.join(&file_name); + let directory = asset_dir()?; + let required_free = limits + .min_free_asset_storage_bytes + .checked_add( + u64::try_from(size).map_err(|_| StorageError::Other("invalid asset size".into()))?, + ) + .ok_or_else(|| StorageError::Other("filesystem reservation overflow".into()))?; + if fs2::available_space(&directory)? < required_free { + return Err(StorageError::AssetResourceLimit( + "min_free_asset_storage_bytes", + )); + } + let path = directory.join(&file_name); OpenOptions::new() .write(true) .create_new(true) .open(&path)? .sync_all()?; let now = sync::now_millis(); - let inserted = db::with_db(|conn| { - conn.execute( + let inserted = db::with_immediate_transaction(|tx| { + let committed: i64 = tx.query_row( + "SELECT COALESCE(SUM(size), 0) FROM user_assets WHERE user_id = ?1 AND deleted = 0", + [user_id], + |row| row.get(0), + )?; + let (reserved, active): (i64, i64) = tx.query_row( + "SELECT COALESCE(SUM(size), 0), COUNT(*) FROM user_asset_uploads WHERE user_id = ?1 AND state = 'uploading'", + [user_id], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + if usize::try_from(active) + .ok() + .is_none_or(|count| count >= limits.max_active_asset_uploads_per_user) + { + return Err(StorageError::AssetResourceLimit( + "max_active_asset_uploads_per_user", + )); + } + let proposed = committed + .checked_add(reserved) + .and_then(|value| value.checked_add(size)) + .ok_or_else(|| StorageError::Other("asset quota accounting overflow".into()))?; + if proposed > limits.max_user_asset_bytes { + return Err(StorageError::AssetResourceLimit("max_user_asset_bytes")); + } + tx.execute( "INSERT INTO user_asset_uploads (upload_id, user_id, asset_id, file_name, size, sha256, mime_type, received, state, created_at, updated_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 0, 'uploading', ?8, ?8)", params![upload_id, user_id, asset_id, file_name, size, sha256, mime_type, now], )?; @@ -219,10 +262,12 @@ pub fn append( return Err(StorageError::Other("upload not found".into())); }; let status = reconcile_offset(user_id, status, &file_name)?; - if status.state != "uploading" - || offset != status.offset - || offset + chunk.len() as i64 > status.size - { + let chunk_len = i64::try_from(chunk.len()) + .map_err(|_| StorageError::Other("asset chunk length overflow".into()))?; + let end = offset + .checked_add(chunk_len) + .ok_or_else(|| StorageError::Other("asset chunk offset overflow".into()))?; + if status.state != "uploading" || offset != status.offset || end > status.size { return Err(StorageError::Other( "asset chunk offset does not match acknowledged offset".into(), )); @@ -231,7 +276,7 @@ pub fn append( let mut file = OpenOptions::new().append(true).open(path)?; file.write_all(chunk)?; file.sync_data()?; - let received = offset + chunk.len() as i64; + let received = end; let now = sync::now_millis(); db::with_db(|conn| { conn.execute( @@ -362,6 +407,34 @@ pub fn abort(user_id: i64, upload_id: &str) -> Result { } } +pub const ASSET_PAGE_SIZE: i64 = 128; + +pub fn list_page( + user_id: i64, + after_id: Option, +) -> Result<(Vec, Option), StorageError> { + if user_id <= 0 || after_id.is_some_and(|id| id <= 0) { + return Err(StorageError::Other("invalid asset list cursor".into())); + } + cleanup_if_due()?; + db::with_db(|conn| { + let mut statement = conn.prepare("SELECT id, user_id, asset_id, file_name, size, sha256, mime_type, revision, created_at, updated_at FROM user_assets WHERE user_id = ?1 AND deleted = 0 AND id < ?2 ORDER BY id DESC LIMIT ?3")?; + let mut items = statement + .query_map( + params![user_id, after_id.unwrap_or(i64::MAX), ASSET_PAGE_SIZE + 1], + metadata_from_row, + )? + .collect::, _>>()?; + let next = if items.len() > ASSET_PAGE_SIZE as usize { + items.pop(); + items.last().map(|asset| asset.id) + } else { + None + }; + Ok((items, next)) + }) +} + pub fn list(user_id: i64) -> Result, StorageError> { cleanup_if_due()?; db::with_db(|conn| { diff --git a/iota-storage/src/util/user_blobs.rs b/iota-storage/src/util/user_blobs.rs index 3ae9e0a..0f0017d 100644 --- a/iota-storage/src/util/user_blobs.rs +++ b/iota-storage/src/util/user_blobs.rs @@ -1,6 +1,6 @@ /* Opaque client-owned bytes. Confidentiality is provided by the client-side blob format. */ use crate::storage_error::StorageError; -use crate::util::{db, sync}; +use crate::util::{config_util::CONFIG, db, sync}; use rusqlite::{Connection, OptionalExtension, Row, Transaction, params}; pub const MAX_BLOB_ID_BYTES: usize = 256; @@ -97,6 +97,16 @@ pub fn put( .filter(|(_, _, deleted, _)| !*deleted) .map(|(_, _, _, size)| size) .unwrap_or(0); + if existing.is_none_or(|(_, _, deleted, _)| deleted) { + let count: i64 = tx.query_row( + "SELECT COUNT(*) FROM user_blobs WHERE user_id = ?1 AND deleted = 0", + [user_id], + |row| row.get(0), + )?; + if count >= CONFIG.load().storage_limits.max_user_blobs { + return Err(StorageError::Other("user blob count quota exceeded".into())); + } + } let aggregate: i64 = tx.query_row("SELECT COALESCE(SUM(length(blob)), 0) FROM user_blobs WHERE user_id = ?1 AND deleted = 0", [user_id], |r| r.get(0))?; let proposed_size = aggregate .checked_sub(current_size) @@ -145,6 +155,46 @@ pub fn list(user_id: i64) -> Result, StorageError> { ) } +pub const BLOB_PAGE_SIZE: i64 = 128; + +pub fn list_metadata_page( + user_id: i64, + after_id: Option, +) -> Result<(Vec, Option), StorageError> { + if user_id <= 0 || after_id.is_some_and(|id| id <= 0) { + return Err(StorageError::Other("invalid blob list cursor".into())); + } + db::with_db(|conn| { + let mut statement = conn.prepare("SELECT id, blob_id, revision, updated_at FROM user_blobs WHERE user_id = ?1 AND deleted = 0 AND id < ?2 ORDER BY id DESC LIMIT ?3")?; + let mut rows = statement.query(params![ + user_id, + after_id.unwrap_or(i64::MAX), + BLOB_PAGE_SIZE + 1 + ])?; + let mut items = Vec::new(); + while let Some(row) = rows.next()? { + items.push(( + row.get::<_, i64>(0)?, + UserBlobMetadata { + blob_id: row.get(1)?, + revision: row.get(2)?, + updated_at: row.get(3)?, + }, + )); + } + let next = if items.len() > BLOB_PAGE_SIZE as usize { + items.pop(); + items.last().map(|(id, _)| *id) + } else { + None + }; + Ok(( + items.into_iter().map(|(_, metadata)| metadata).collect(), + next, + )) + }) +} + pub fn list_metadata(user_id: i64) -> Result, StorageError> { if user_id <= 0 { return Err(StorageError::Other("invalid blob owner".into())); diff --git a/iota-storage/tests/asset_quota.rs b/iota-storage/tests/asset_quota.rs new file mode 100644 index 0000000..1c645e8 --- /dev/null +++ b/iota-storage/tests/asset_quota.rs @@ -0,0 +1,50 @@ +use iota_storage::util::{config_util, db, user_assets}; +use sha2::{Digest, Sha256}; +use std::sync::Arc; + +#[test] +fn upload_quota_reserves_and_releases_capacity() -> Result<(), Box> { + let directory = tempfile::tempdir()?; + iota_util::file_util::configure_storage_directory(directory.path().into()); + db::initialize_database()?; + let mut config = config_util::IotaConfig::default(); + config.storage_limits.max_asset_bytes = 5; + config.storage_limits.max_user_asset_bytes = 8; + config.storage_limits.max_active_asset_uploads_per_user = 2; + config.storage_limits.min_free_asset_storage_bytes = 0; + config_util::CONFIG.store(Arc::new(config)); + + let hash = Sha256::digest(b"data"); + assert!(user_assets::start(42, "oversize", "a", 6, &hash, "text/plain").is_err()); + user_assets::start(42, "first", "a", 4, &hash, "text/plain")?; + user_assets::start(42, "second", "b", 4, &hash, "text/plain")?; + assert!(user_assets::start(42, "third", "c", 1, &hash, "text/plain").is_err()); + assert!(user_assets::abort(42, "second")?); + user_assets::start(42, "third", "c", 1, &hash, "text/plain")?; + user_assets::append(42, "first", 0, b"data")?; + user_assets::commit(42, "first")?; + assert!(user_assets::start(42, "fourth", "d", 4, &hash, "text/plain").is_err()); + assert!(user_assets::abort(42, "third")?); + user_assets::start(42, "fourth", "d", 4, &hash, "text/plain")?; + + let barrier = Arc::new(std::sync::Barrier::new(2)); + let attempts: Vec<_> = ["concurrent-a", "concurrent-b"] + .into_iter() + .map(|upload_id| { + let barrier = barrier.clone(); + std::thread::spawn(move || { + barrier.wait(); + user_assets::start(77, upload_id, upload_id, 5, &hash, "text/plain").is_ok() + }) + }) + .collect(); + assert_eq!( + attempts + .into_iter() + .filter_map(|handle| handle.join().ok()) + .filter(|ok| *ok) + .count(), + 1 + ); + Ok(()) +} diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index d8960b3..5b9a9d9 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -38,6 +38,31 @@ use iota_identity::AuthorityLocator; // ============================================================================ const IOTA_KEYRING_PATH: &str = "iota.mk"; +static ASSET_IO_PERMITS: LazyLock> = LazyLock::new(|| { + Arc::new(Semaphore::new( + CONFIG.load().storage_limits.max_asset_io_workers, + )) +}); + +async fn run_asset_handler( + cv: &CommunicationValue, + handler: fn(&CommunicationValue) -> T, +) -> Result { + let permit = ASSET_IO_PERMITS + .clone() + .try_acquire_owned() + .map_err(|_| error_response(cv, CommunicationType::ErrorInternal))?; + let frame = cv.clone(); + tokio::task::spawn_blocking(move || { + let _permit = permit; + handler(&frame) + }) + .await + .map_err(|error| { + log!("Asset I/O worker failed: {error}"); + error_response(cv, CommunicationType::ErrorInternal) + }) +} static IDENTITY_PATH: std::sync::OnceLock = std::sync::OnceLock::new(); static OMIKRON_TRUST_DIRECTORY: std::sync::OnceLock = std::sync::OnceLock::new(); @@ -231,7 +256,6 @@ fn parse_omega_invitation( const TASK_CLEANUP_INTERVAL: Duration = Duration::from_secs(60); const TASK_MAX_AGE: Duration = Duration::from_secs(60); const MAX_CONCURRENT_HANDLERS: usize = 20; -const RELAY_RETENTION_MILLIS: i64 = 30 * 24 * 60 * 60 * 1000; struct ResolvedOmikronEndpoint { id: Option, @@ -1126,7 +1150,7 @@ impl OmikronConnection { invitation_self.flush_pending_invitation_actions().await; }); if let Err(error) = relay_replay::prune_completed( - now_millis_i64().saturating_sub(RELAY_RETENTION_MILLIS), + now_millis_i64().saturating_sub(relay_replay::RELAY_RETENTION_MILLIS), ) { log!("Relay replay cleanup failed: {}", error); } @@ -2914,25 +2938,51 @@ impl OmikronConnection { } async fn handle_user_asset_upload_start(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_upload_start(cv)) - .await; + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_upload_start).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_upload_chunk(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_upload_chunk(cv)) - .await; + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_upload_chunk).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_upload_status(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_upload_status(cv)) - .await; + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_upload_status).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_upload_commit(self: Arc, cv: &CommunicationValue) { - let mutation = message_handlers::handle_user_asset_upload_commit(cv); + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_upload_commit).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let mutation = result; let _ = self.send_message(&mutation.response).await; if let Some(changed) = mutation.changed { let _ = self.send_message(&changed).await; @@ -2940,25 +2990,49 @@ impl OmikronConnection { } async fn handle_user_asset_upload_abort(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_upload_abort(cv)) - .await; + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_upload_abort).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_get_chunk(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_get_chunk(cv)) - .await; + let result = + match run_asset_handler(cv, message_handlers::handle_user_asset_get_chunk).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_list(self: Arc, cv: &CommunicationValue) { - let _ = self - .send_message(&message_handlers::handle_user_asset_list(cv)) - .await; + let result = match run_asset_handler(cv, message_handlers::handle_user_asset_list).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let _ = self.send_message(&result).await; } async fn handle_user_asset_delete(self: Arc, cv: &CommunicationValue) { - let mutation = message_handlers::handle_user_asset_delete(cv); + let result = match run_asset_handler(cv, message_handlers::handle_user_asset_delete).await { + Ok(result) => result, + Err(response) => { + let _ = self.send_message(&response).await; + return; + } + }; + let mutation = result; let _ = self.send_message(&mutation.response).await; if let Some(changed) = mutation.changed { let _ = self.send_message(&changed).await; @@ -3539,6 +3613,12 @@ mod tests { raw.push(1); raw.extend_from_slice(&keyring.try_to_bytes().expect("keyring serializes")); fs::write(&path, raw).expect("legacy fixture writes"); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + fs::set_permissions(&path, fs::Permissions::from_mode(0o600)) + .expect("legacy fixture is owner-only"); + } let migrated = load_or_migrate_keyring_at(&path, None).expect("legacy identity loads"); assert_eq!( diff --git a/web-server/src/lib.rs b/web-server/src/lib.rs index bcd0765..d690df0 100644 --- a/web-server/src/lib.rs +++ b/web-server/src/lib.rs @@ -40,6 +40,7 @@ pub struct WebConfig { pub asset_dir: PathBuf, pub tls: Option, pub required: bool, + pub max_mtp_sessions: usize, pub authority_discovery: iota_identity::AuthorityDiscoveryDocument, pub local_users: Arc, pub descriptor_publisher: Arc, @@ -187,7 +188,7 @@ pub async fn start( }), Box::new(|_, _| Box::pin(async { 0 })), ) - .with_authentication_policy(mtp::host::AuthenticationPolicy::AllowAuthentication); + .with_authentication_policy(mtp::host::AuthenticationPolicy::ForceAuthentication); let assets = config.asset_dir.clone(); let discovery = serde_json::to_vec(&config.authority_discovery) .map(Bytes::from) @@ -290,13 +291,21 @@ pub async fn start( .map_err(|e| WebServerError::Startup(e.to_string()))?; let cancellation = parent.child_token(); let task_cancellation = cancellation.clone(); + let session_limit = Arc::new(tokio::sync::Semaphore::new(config.max_mtp_sessions)); let join = tokio::spawn(async move { loop { tokio::select! { result = server.accept() => match result { Ok(Some(connection)) => { + let Ok(permit) = session_limit.clone().try_acquire_owned() else { + log!("Rejected MTP connection: session limit reached"); + continue; + }; let handler = mtp_handler.clone(); - tokio::spawn(async move { handler.accept(connection).await; }); + tokio::spawn(async move { + let _permit = permit; + handler.accept(connection).await; + }); } Ok(None) => break, Err(error) => log!("MTP webserver connection failed: {}", error), diff --git a/web-ui/Cargo.toml b/web-ui/Cargo.toml index 679a1f2..d1f3516 100644 --- a/web-ui/Cargo.toml +++ b/web-ui/Cargo.toml @@ -3,6 +3,9 @@ name = "web-ui" version = "0.1.0" edition = "2024" +[features] +legacy-web-admin = [] + [dependencies] iota-storage = { path = "../iota-storage" } iota-state = { path = "../iota-state" } diff --git a/web-ui/src/lib.rs b/web-ui/src/lib.rs index de2a79a..f2f958a 100644 --- a/web-ui/src/lib.rs +++ b/web-ui/src/lib.rs @@ -1,3 +1,5 @@ +#[cfg(feature = "legacy-web-admin")] pub mod api; +#[cfg(feature = "legacy-web-admin")] pub mod server; pub mod web_path_parser;