[Fix] Bound Iota storage, relay and transport resources
This commit is contained in:
parent
46078cbc4a
commit
e19c3c3d12
19 changed files with 609 additions and 48 deletions
2
Cargo.lock
generated
2
Cargo.lock
generated
|
|
@ -2255,6 +2255,7 @@ version = "0.1.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"base64 0.22.1",
|
"base64 0.22.1",
|
||||||
|
"libc",
|
||||||
"mtp",
|
"mtp",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
|
|
@ -2326,6 +2327,7 @@ dependencies = [
|
||||||
"arc-swap",
|
"arc-swap",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"base64 0.22.1",
|
"base64 0.22.1",
|
||||||
|
"fs2",
|
||||||
"iota-identity",
|
"iota-identity",
|
||||||
"iota-logger",
|
"iota-logger",
|
||||||
"iota-paths",
|
"iota-paths",
|
||||||
|
|
|
||||||
22
README.md
22
README.md
|
|
@ -66,3 +66,25 @@ usermod -aG iota-operators USER
|
||||||
The user must start a new login session before supplementary group membership
|
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
|
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`.
|
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.
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
use crate::message_common::*;
|
use crate::message_common::*;
|
||||||
|
use iota_logger::log;
|
||||||
use iota_storage::util::chat_files::{self, MessageState};
|
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::chats_util::{self, get_user, has_user, mod_user};
|
||||||
use iota_storage::util::communities_util::CommunitiesUtil;
|
use iota_storage::util::communities_util::CommunitiesUtil;
|
||||||
|
|
@ -2084,8 +2085,50 @@ fn upload_response(
|
||||||
}
|
}
|
||||||
|
|
||||||
fn asset_error(cv: &CommunicationValue, error: StorageError) -> CommunicationValue {
|
fn asset_error(cv: &CommunicationValue, error: StorageError) -> CommunicationValue {
|
||||||
error_response(cv, CommunicationType::ErrorInternal)
|
log!(
|
||||||
.add_typed_default(DataType::ErrorType, DataValue::Str(error.to_string()))
|
"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> {
|
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 {
|
if user_id <= 0 {
|
||||||
return error_response(cv, CommunicationType::ErrorInvalidData);
|
return error_response(cv, CommunicationType::ErrorInvalidData);
|
||||||
}
|
}
|
||||||
match user_blobs::list_metadata(user_id) {
|
let after_id = match cv.get_data(DataType::Offset).and_then(DataValue::as_number) {
|
||||||
Ok(blobs) => CommunicationValue::new(CommunicationType::UserBlobList)
|
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_request_id(cv)
|
||||||
.with_receiver(sender_wire_id(user_id))
|
.with_receiver(sender_wire_id(user_id))
|
||||||
|
.add_typed_default(
|
||||||
|
DataType::Offset,
|
||||||
|
DataValue::SignedNumber(next.unwrap_or(0).into()),
|
||||||
|
)
|
||||||
.add_typed_default(
|
.add_typed_default(
|
||||||
DataType::Blobs,
|
DataType::Blobs,
|
||||||
DataValue::Array(blobs.iter().map(blob_metadata_value).collect()),
|
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,
|
Ok(id) if id > 0 => id,
|
||||||
_ => return error_response(cv, CommunicationType::ErrorInvalidData),
|
_ => return error_response(cv, CommunicationType::ErrorInvalidData),
|
||||||
};
|
};
|
||||||
match user_assets::list(user_id) {
|
let after_id = match cv.get_data(DataType::Offset).and_then(DataValue::as_number) {
|
||||||
Ok(assets) => CommunicationValue::new(CommunicationType::UserAssetList)
|
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_request_id(cv)
|
||||||
.with_receiver(sender_wire_id(user_id))
|
.with_receiver(sender_wire_id(user_id))
|
||||||
|
.add_typed_default(
|
||||||
|
DataType::Offset,
|
||||||
|
DataValue::SignedNumber(next.unwrap_or(0).into()),
|
||||||
|
)
|
||||||
.add_typed_default(
|
.add_typed_default(
|
||||||
DataType::Assets,
|
DataType::Assets,
|
||||||
DataValue::Array(assets.iter().map(asset_metadata_value).collect()),
|
DataValue::Array(assets.iter().map(asset_metadata_value).collect()),
|
||||||
|
|
|
||||||
|
|
@ -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 iota_util::route_target::RouteTarget;
|
||||||
use mtp::codec::{
|
use mtp::codec::{
|
||||||
CommunicationValue, ProtectionPolicy, RelayError, RelayOpenOptions, SignaturePolicy, TypeMap,
|
CommunicationValue, ProtectionPolicy, RelayError, RelayOpenOptions, SignaturePolicy, TypeMap,
|
||||||
|
|
@ -107,6 +111,9 @@ pub enum RelayValidationError {
|
||||||
MissingTypeMap,
|
MissingTypeMap,
|
||||||
InvalidRouteTarget(u64),
|
InvalidRouteTarget(u64),
|
||||||
KeyLookup(String),
|
KeyLookup(String),
|
||||||
|
Clock(String),
|
||||||
|
StaleRelay { created_at: u64, now: u64 },
|
||||||
|
RelayFromFuture { created_at: u64, now: u64 },
|
||||||
Relay(RelayError),
|
Relay(RelayError),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -129,6 +136,14 @@ impl fmt::Display for RelayValidationError {
|
||||||
write!(formatter, "relay has invalid route target {target}")
|
write!(formatter, "relay has invalid route target {target}")
|
||||||
}
|
}
|
||||||
Self::KeyLookup(error) => write!(formatter, "trusted signer lookup failed: {error}"),
|
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),
|
Self::Relay(error) => error.fmt(formatter),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -142,6 +157,23 @@ impl From<RelayError> 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
|
* Relay metadata is opened only after the claimed signer selects trusted key
|
||||||
* history. Replay reservation happens after verification and durable
|
* history. Replay reservation happens after verification and durable
|
||||||
|
|
@ -195,6 +227,14 @@ where
|
||||||
RelayOpenOptions::new(RELAY_PROTECTION_POLICY),
|
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 {
|
let context = VerifiedRelayContext {
|
||||||
signer_id: metadata.signer_id(),
|
signer_id: metadata.signer_id(),
|
||||||
final_recipient_id: metadata.final_recipient_id(),
|
final_recipient_id: metadata.final_recipient_id(),
|
||||||
|
|
@ -244,6 +284,23 @@ mod tests {
|
||||||
use mtp::codec::SealedRelayBuilder;
|
use mtp::codec::SealedRelayBuilder;
|
||||||
use mtp::crypto::{DualSigner, Ed25519Signer, Keyring};
|
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> {
|
fn relay(message_id: u128) -> Result<(Keyring, Keyring, CommunicationValue), String> {
|
||||||
let signer_keyring = Keyring::generate();
|
let signer_keyring = Keyring::generate();
|
||||||
let recipient_keyring = Keyring::generate();
|
let recipient_keyring = Keyring::generate();
|
||||||
|
|
@ -264,7 +321,7 @@ mod tests {
|
||||||
&signer,
|
&signer,
|
||||||
)
|
)
|
||||||
.message_id(message_id)
|
.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()])
|
.metadata_recipients(vec![recipient_keyring.public_key_bundle()])
|
||||||
.content_recipients(vec![recipient_keyring.public_key_bundle()])
|
.content_recipients(vec![recipient_keyring.public_key_bundle()])
|
||||||
.build()
|
.build()
|
||||||
|
|
@ -339,7 +396,7 @@ mod tests {
|
||||||
&signer,
|
&signer,
|
||||||
)
|
)
|
||||||
.message_id(7_u128)
|
.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()])
|
.metadata_recipients(vec![recipient_keyring.public_key_bundle()])
|
||||||
.content_recipients(vec![recipient_keyring.public_key_bundle()])
|
.content_recipients(vec![recipient_keyring.public_key_bundle()])
|
||||||
.build()
|
.build()
|
||||||
|
|
|
||||||
|
|
@ -327,11 +327,11 @@ fn peer_credentials(stream: &UnixStream) -> Result<PeerIdentity> {
|
||||||
}
|
}
|
||||||
#[cfg(not(target_os = "linux"))]
|
#[cfg(not(target_os = "linux"))]
|
||||||
{
|
{
|
||||||
Ok(PeerIdentity {
|
let _ = stream;
|
||||||
pid: 0,
|
Err(std::io::Error::new(
|
||||||
uid: 0,
|
std::io::ErrorKind::Unsupported,
|
||||||
_gid: 0,
|
"IPC peer credentials are unsupported on this platform",
|
||||||
})
|
))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -576,6 +576,7 @@ async fn main() -> ExitCode {
|
||||||
key: resolve_config_path(&paths.config_file, &key, &paths.config_dir),
|
key: resolve_config_path(&paths.config_file, &key, &paths.config_dir),
|
||||||
}),
|
}),
|
||||||
required: web.required,
|
required: web.required,
|
||||||
|
max_mtp_sessions: web.max_mtp_sessions,
|
||||||
authority_discovery: iota_identity::AuthorityDiscoveryDocument {
|
authority_discovery: iota_identity::AuthorityDiscoveryDocument {
|
||||||
version: 1,
|
version: 1,
|
||||||
service: iota_identity::AuthorityKind::Iota,
|
service: iota_identity::AuthorityKind::Iota,
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ version = "0.1.0"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
libc = "0.2"
|
||||||
async-trait = "0.1.89"
|
async-trait = "0.1.89"
|
||||||
base64 = "0.22.1"
|
base64 = "0.22.1"
|
||||||
mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto", "files", "raw"] }
|
mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto", "files", "raw"] }
|
||||||
|
|
|
||||||
|
|
@ -179,6 +179,34 @@ impl fmt::Display for LocalNodeIdentityError {
|
||||||
|
|
||||||
impl std::error::Error 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)]
|
#[derive(Clone)]
|
||||||
pub struct LocalNodeIdentity {
|
pub struct LocalNodeIdentity {
|
||||||
keyring: Arc<mtp::crypto::Keyring>,
|
keyring: Arc<mtp::crypto::Keyring>,
|
||||||
|
|
@ -210,6 +238,10 @@ impl LocalNodeIdentity {
|
||||||
source,
|
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) {
|
let keyring = match mtp::files::load_keyring_raw(path) {
|
||||||
Ok(keyring) => keyring,
|
Ok(keyring) => keyring,
|
||||||
Err(mtp::files::FileError::Io(error))
|
Err(mtp::files::FileError::Io(error))
|
||||||
|
|
@ -250,6 +282,10 @@ impl LocalNodeIdentity {
|
||||||
source: error.to_string(),
|
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| {
|
let persisted = mtp::files::load_keyring_raw(path).map_err(|error| {
|
||||||
LocalNodeIdentityError::Storage {
|
LocalNodeIdentityError::Storage {
|
||||||
path: path.display().to_string(),
|
path: path.display().to_string(),
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ sha2 = "0.10"
|
||||||
thiserror = "2"
|
thiserror = "2"
|
||||||
rand = "0.8"
|
rand = "0.8"
|
||||||
rusqlite = "0.40.0"
|
rusqlite = "0.40.0"
|
||||||
|
fs2 = "0.4.3"
|
||||||
tokio = { version = "1.50.0", features = ["full"] }
|
tokio = { version = "1.50.0", features = ["full"] }
|
||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
|
|
|
||||||
|
|
@ -14,6 +14,8 @@ pub enum StorageError {
|
||||||
RevisionConflict,
|
RevisionConflict,
|
||||||
#[error("pending relay ownership is unknown")]
|
#[error("pending relay ownership is unknown")]
|
||||||
PendingRelayOwnershipUnknown,
|
PendingRelayOwnershipUnknown,
|
||||||
|
#[error("asset resource limit: {0}")]
|
||||||
|
AssetResourceLimit(&'static str),
|
||||||
#[error("{0}")]
|
#[error("{0}")]
|
||||||
Other(String),
|
Other(String),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -40,8 +40,14 @@ pub enum ConfigError {
|
||||||
InvalidRelayRouterKey(String),
|
InvalidRelayRouterKey(String),
|
||||||
#[error("relay router certificate path must not be empty")]
|
#[error("relay router certificate path must not be empty")]
|
||||||
MissingRelayRouterCertificate,
|
MissingRelayRouterCertificate,
|
||||||
|
#[error("web.max_mtp_sessions must be greater than zero")]
|
||||||
|
InvalidMaxMtpSessions,
|
||||||
#[error("max_ipc_clients must be greater than zero")]
|
#[error("max_ipc_clients must be greater than zero")]
|
||||||
InvalidMaxIpcClients,
|
InvalidMaxIpcClients,
|
||||||
|
#[error("invalid storage limit: {0}")]
|
||||||
|
InvalidStorageLimit(&'static str),
|
||||||
|
#[error("invalid max_relay_future_skew_millis")]
|
||||||
|
InvalidRelayFutureSkew,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
|
@ -71,6 +77,44 @@ pub struct IotaConfig {
|
||||||
pub read_receipts_enabled: bool,
|
pub read_receipts_enabled: bool,
|
||||||
#[serde(default = "default_max_ipc_clients")]
|
#[serde(default = "default_max_ipc_clients")]
|
||||||
pub max_ipc_clients: usize,
|
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)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
|
@ -109,11 +153,16 @@ pub struct WebSettings {
|
||||||
pub key: Option<String>,
|
pub key: Option<String>,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
pub required: bool,
|
pub required: bool,
|
||||||
|
#[serde(default = "default_max_mtp_sessions")]
|
||||||
|
pub max_mtp_sessions: usize,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
pub direct_endpoints: Vec<String>,
|
pub direct_endpoints: Vec<String>,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
pub relay_hints: Vec<String>,
|
pub relay_hints: Vec<String>,
|
||||||
}
|
}
|
||||||
|
const fn default_max_mtp_sessions() -> usize {
|
||||||
|
256
|
||||||
|
}
|
||||||
fn default_web_bind() -> String {
|
fn default_web_bind() -> String {
|
||||||
"127.0.0.1".into()
|
"127.0.0.1".into()
|
||||||
}
|
}
|
||||||
|
|
@ -130,6 +179,7 @@ impl Default for WebSettings {
|
||||||
certificate: None,
|
certificate: None,
|
||||||
key: None,
|
key: None,
|
||||||
required: false,
|
required: false,
|
||||||
|
max_mtp_sessions: default_max_mtp_sessions(),
|
||||||
direct_endpoints: Vec::new(),
|
direct_endpoints: Vec::new(),
|
||||||
relay_hints: Vec::new(),
|
relay_hints: Vec::new(),
|
||||||
}
|
}
|
||||||
|
|
@ -148,6 +198,10 @@ const fn default_max_ipc_clients() -> usize {
|
||||||
64
|
64
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const fn default_max_relay_future_skew_millis() -> u64 {
|
||||||
|
5 * 60 * 1_000
|
||||||
|
}
|
||||||
|
|
||||||
impl Default for IotaConfig {
|
impl Default for IotaConfig {
|
||||||
fn default() -> Self {
|
fn default() -> Self {
|
||||||
Self {
|
Self {
|
||||||
|
|
@ -163,6 +217,8 @@ impl Default for IotaConfig {
|
||||||
private_key: None,
|
private_key: None,
|
||||||
read_receipts_enabled: default_read_receipts_enabled(),
|
read_receipts_enabled: default_read_receipts_enabled(),
|
||||||
max_ipc_clients: default_max_ipc_clients(),
|
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(),
|
bind: config.web.bind.clone(),
|
||||||
source,
|
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 {
|
if config.max_ipc_clients == 0 {
|
||||||
return Err(ConfigError::InvalidMaxIpcClients);
|
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
|
for endpoint in config
|
||||||
.web
|
.web
|
||||||
.direct_endpoints
|
.direct_endpoints
|
||||||
|
|
@ -357,6 +437,30 @@ mod tests {
|
||||||
use super::{ConfigError, IotaConfig, RelayRouterSettings, parse_config, validate_config};
|
use super::{ConfigError, IotaConfig, RelayRouterSettings, parse_config, validate_config};
|
||||||
use std::path::Path;
|
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]
|
#[test]
|
||||||
fn malformed_yaml_is_rejected() {
|
fn malformed_yaml_is_rejected() {
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,9 @@ use crate::storage_error::StorageError;
|
||||||
use crate::util::db;
|
use crate::util::db;
|
||||||
use rusqlite::{OptionalExtension, params};
|
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)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct DeliveredRelay {
|
pub struct DeliveredRelay {
|
||||||
pub signer_id: i64,
|
pub signer_id: i64,
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
use crate::storage_error::StorageError;
|
use crate::storage_error::StorageError;
|
||||||
use crate::util::{db, sync};
|
use crate::util::{config_util::CONFIG, db, sync};
|
||||||
use rand::RngCore;
|
use rand::RngCore;
|
||||||
use rusqlite::{OptionalExtension, Row, params};
|
use rusqlite::{OptionalExtension, Row, params};
|
||||||
use sha2::{Digest, Sha256};
|
use sha2::{Digest, Sha256};
|
||||||
|
|
@ -118,7 +118,8 @@ fn reconcile_offset(
|
||||||
if status.state != "uploading" {
|
if status.state != "uploading" {
|
||||||
return Ok(status);
|
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 {
|
if actual > status.size {
|
||||||
return Err(StorageError::Other(
|
return Err(StorageError::Other(
|
||||||
"upload file exceeds declared size".into(),
|
"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 {
|
if size < 0 || sha256.len() != 32 || mime_type.is_empty() || mime_type.len() > 255 {
|
||||||
return Err(StorageError::Other("invalid asset metadata".into()));
|
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()?;
|
cleanup_if_due()?;
|
||||||
if let Some((status, file_name)) = load_upload(user_id, upload_id)? {
|
if let Some((status, file_name)) = load_upload(user_id, upload_id)? {
|
||||||
if status.asset_id != asset_id
|
if status.asset_id != asset_id
|
||||||
|
|
@ -166,15 +172,52 @@ pub fn start(
|
||||||
}
|
}
|
||||||
|
|
||||||
let file_name = random_file_name();
|
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()
|
OpenOptions::new()
|
||||||
.write(true)
|
.write(true)
|
||||||
.create_new(true)
|
.create_new(true)
|
||||||
.open(&path)?
|
.open(&path)?
|
||||||
.sync_all()?;
|
.sync_all()?;
|
||||||
let now = sync::now_millis();
|
let now = sync::now_millis();
|
||||||
let inserted = db::with_db(|conn| {
|
let inserted = db::with_immediate_transaction(|tx| {
|
||||||
conn.execute(
|
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)",
|
"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],
|
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()));
|
return Err(StorageError::Other("upload not found".into()));
|
||||||
};
|
};
|
||||||
let status = reconcile_offset(user_id, status, &file_name)?;
|
let status = reconcile_offset(user_id, status, &file_name)?;
|
||||||
if status.state != "uploading"
|
let chunk_len = i64::try_from(chunk.len())
|
||||||
|| offset != status.offset
|
.map_err(|_| StorageError::Other("asset chunk length overflow".into()))?;
|
||||||
|| offset + chunk.len() as i64 > status.size
|
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(
|
return Err(StorageError::Other(
|
||||||
"asset chunk offset does not match acknowledged offset".into(),
|
"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)?;
|
let mut file = OpenOptions::new().append(true).open(path)?;
|
||||||
file.write_all(chunk)?;
|
file.write_all(chunk)?;
|
||||||
file.sync_data()?;
|
file.sync_data()?;
|
||||||
let received = offset + chunk.len() as i64;
|
let received = end;
|
||||||
let now = sync::now_millis();
|
let now = sync::now_millis();
|
||||||
db::with_db(|conn| {
|
db::with_db(|conn| {
|
||||||
conn.execute(
|
conn.execute(
|
||||||
|
|
@ -362,6 +407,34 @@ pub fn abort(user_id: i64, upload_id: &str) -> Result<bool, StorageError> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub const ASSET_PAGE_SIZE: i64 = 128;
|
||||||
|
|
||||||
|
pub fn list_page(
|
||||||
|
user_id: i64,
|
||||||
|
after_id: Option<i64>,
|
||||||
|
) -> Result<(Vec<UserAssetMetadata>, Option<i64>), 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::<Result<Vec<_>, _>>()?;
|
||||||
|
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<Vec<UserAssetMetadata>, StorageError> {
|
pub fn list(user_id: i64) -> Result<Vec<UserAssetMetadata>, StorageError> {
|
||||||
cleanup_if_due()?;
|
cleanup_if_due()?;
|
||||||
db::with_db(|conn| {
|
db::with_db(|conn| {
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
/* Opaque client-owned bytes. Confidentiality is provided by the client-side blob format. */
|
/* Opaque client-owned bytes. Confidentiality is provided by the client-side blob format. */
|
||||||
use crate::storage_error::StorageError;
|
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};
|
use rusqlite::{Connection, OptionalExtension, Row, Transaction, params};
|
||||||
|
|
||||||
pub const MAX_BLOB_ID_BYTES: usize = 256;
|
pub const MAX_BLOB_ID_BYTES: usize = 256;
|
||||||
|
|
@ -97,6 +97,16 @@ pub fn put(
|
||||||
.filter(|(_, _, deleted, _)| !*deleted)
|
.filter(|(_, _, deleted, _)| !*deleted)
|
||||||
.map(|(_, _, _, size)| size)
|
.map(|(_, _, _, size)| size)
|
||||||
.unwrap_or(0);
|
.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 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
|
let proposed_size = aggregate
|
||||||
.checked_sub(current_size)
|
.checked_sub(current_size)
|
||||||
|
|
@ -145,6 +155,46 @@ pub fn list(user_id: i64) -> Result<Vec<UserBlob>, StorageError> {
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub const BLOB_PAGE_SIZE: i64 = 128;
|
||||||
|
|
||||||
|
pub fn list_metadata_page(
|
||||||
|
user_id: i64,
|
||||||
|
after_id: Option<i64>,
|
||||||
|
) -> Result<(Vec<UserBlobMetadata>, Option<i64>), 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<Vec<UserBlobMetadata>, StorageError> {
|
pub fn list_metadata(user_id: i64) -> Result<Vec<UserBlobMetadata>, StorageError> {
|
||||||
if user_id <= 0 {
|
if user_id <= 0 {
|
||||||
return Err(StorageError::Other("invalid blob owner".into()));
|
return Err(StorageError::Other("invalid blob owner".into()));
|
||||||
|
|
|
||||||
50
iota-storage/tests/asset_quota.rs
Normal file
50
iota-storage/tests/asset_quota.rs
Normal file
|
|
@ -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<dyn std::error::Error>> {
|
||||||
|
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(())
|
||||||
|
}
|
||||||
|
|
@ -38,6 +38,31 @@ use iota_identity::AuthorityLocator;
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|
||||||
const IOTA_KEYRING_PATH: &str = "iota.mk";
|
const IOTA_KEYRING_PATH: &str = "iota.mk";
|
||||||
|
static ASSET_IO_PERMITS: LazyLock<Arc<Semaphore>> = LazyLock::new(|| {
|
||||||
|
Arc::new(Semaphore::new(
|
||||||
|
CONFIG.load().storage_limits.max_asset_io_workers,
|
||||||
|
))
|
||||||
|
});
|
||||||
|
|
||||||
|
async fn run_asset_handler<T: Send + 'static>(
|
||||||
|
cv: &CommunicationValue,
|
||||||
|
handler: fn(&CommunicationValue) -> T,
|
||||||
|
) -> Result<T, CommunicationValue> {
|
||||||
|
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<PathBuf> = std::sync::OnceLock::new();
|
static IDENTITY_PATH: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
|
||||||
static OMIKRON_TRUST_DIRECTORY: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
|
static OMIKRON_TRUST_DIRECTORY: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
|
@ -231,7 +256,6 @@ fn parse_omega_invitation(
|
||||||
const TASK_CLEANUP_INTERVAL: Duration = Duration::from_secs(60);
|
const TASK_CLEANUP_INTERVAL: Duration = Duration::from_secs(60);
|
||||||
const TASK_MAX_AGE: Duration = Duration::from_secs(60);
|
const TASK_MAX_AGE: Duration = Duration::from_secs(60);
|
||||||
const MAX_CONCURRENT_HANDLERS: usize = 20;
|
const MAX_CONCURRENT_HANDLERS: usize = 20;
|
||||||
const RELAY_RETENTION_MILLIS: i64 = 30 * 24 * 60 * 60 * 1000;
|
|
||||||
|
|
||||||
struct ResolvedOmikronEndpoint {
|
struct ResolvedOmikronEndpoint {
|
||||||
id: Option<i64>,
|
id: Option<i64>,
|
||||||
|
|
@ -1126,7 +1150,7 @@ impl OmikronConnection {
|
||||||
invitation_self.flush_pending_invitation_actions().await;
|
invitation_self.flush_pending_invitation_actions().await;
|
||||||
});
|
});
|
||||||
if let Err(error) = relay_replay::prune_completed(
|
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);
|
log!("Relay replay cleanup failed: {}", error);
|
||||||
}
|
}
|
||||||
|
|
@ -2914,25 +2938,51 @@ impl OmikronConnection {
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_user_asset_upload_start(self: Arc<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_upload_start(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result =
|
||||||
.send_message(&message_handlers::handle_user_asset_upload_start(cv))
|
match run_asset_handler(cv, message_handlers::handle_user_asset_upload_start).await {
|
||||||
.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<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_upload_chunk(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result =
|
||||||
.send_message(&message_handlers::handle_user_asset_upload_chunk(cv))
|
match run_asset_handler(cv, message_handlers::handle_user_asset_upload_chunk).await {
|
||||||
.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<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_upload_status(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result =
|
||||||
.send_message(&message_handlers::handle_user_asset_upload_status(cv))
|
match run_asset_handler(cv, message_handlers::handle_user_asset_upload_status).await {
|
||||||
.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<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_upload_commit(self: Arc<Self>, 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;
|
let _ = self.send_message(&mutation.response).await;
|
||||||
if let Some(changed) = mutation.changed {
|
if let Some(changed) = mutation.changed {
|
||||||
let _ = self.send_message(&changed).await;
|
let _ = self.send_message(&changed).await;
|
||||||
|
|
@ -2940,25 +2990,49 @@ impl OmikronConnection {
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_user_asset_upload_abort(self: Arc<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_upload_abort(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result =
|
||||||
.send_message(&message_handlers::handle_user_asset_upload_abort(cv))
|
match run_asset_handler(cv, message_handlers::handle_user_asset_upload_abort).await {
|
||||||
.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<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_get_chunk(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result =
|
||||||
.send_message(&message_handlers::handle_user_asset_get_chunk(cv))
|
match run_asset_handler(cv, message_handlers::handle_user_asset_get_chunk).await {
|
||||||
.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<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_list(self: Arc<Self>, cv: &CommunicationValue) {
|
||||||
let _ = self
|
let result = match run_asset_handler(cv, message_handlers::handle_user_asset_list).await {
|
||||||
.send_message(&message_handlers::handle_user_asset_list(cv))
|
Ok(result) => result,
|
||||||
.await;
|
Err(response) => {
|
||||||
|
let _ = self.send_message(&response).await;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let _ = self.send_message(&result).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_user_asset_delete(self: Arc<Self>, cv: &CommunicationValue) {
|
async fn handle_user_asset_delete(self: Arc<Self>, 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;
|
let _ = self.send_message(&mutation.response).await;
|
||||||
if let Some(changed) = mutation.changed {
|
if let Some(changed) = mutation.changed {
|
||||||
let _ = self.send_message(&changed).await;
|
let _ = self.send_message(&changed).await;
|
||||||
|
|
@ -3539,6 +3613,12 @@ mod tests {
|
||||||
raw.push(1);
|
raw.push(1);
|
||||||
raw.extend_from_slice(&keyring.try_to_bytes().expect("keyring serializes"));
|
raw.extend_from_slice(&keyring.try_to_bytes().expect("keyring serializes"));
|
||||||
fs::write(&path, raw).expect("legacy fixture writes");
|
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");
|
let migrated = load_or_migrate_keyring_at(&path, None).expect("legacy identity loads");
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|
|
||||||
|
|
@ -40,6 +40,7 @@ pub struct WebConfig {
|
||||||
pub asset_dir: PathBuf,
|
pub asset_dir: PathBuf,
|
||||||
pub tls: Option<TlsConfig>,
|
pub tls: Option<TlsConfig>,
|
||||||
pub required: bool,
|
pub required: bool,
|
||||||
|
pub max_mtp_sessions: usize,
|
||||||
pub authority_discovery: iota_identity::AuthorityDiscoveryDocument,
|
pub authority_discovery: iota_identity::AuthorityDiscoveryDocument,
|
||||||
pub local_users: Arc<dyn LocalUserStore>,
|
pub local_users: Arc<dyn LocalUserStore>,
|
||||||
pub descriptor_publisher: Arc<dyn LocalDescriptorPublisher>,
|
pub descriptor_publisher: Arc<dyn LocalDescriptorPublisher>,
|
||||||
|
|
@ -187,7 +188,7 @@ pub async fn start(
|
||||||
}),
|
}),
|
||||||
Box::new(|_, _| Box::pin(async { 0 })),
|
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 assets = config.asset_dir.clone();
|
||||||
let discovery = serde_json::to_vec(&config.authority_discovery)
|
let discovery = serde_json::to_vec(&config.authority_discovery)
|
||||||
.map(Bytes::from)
|
.map(Bytes::from)
|
||||||
|
|
@ -290,13 +291,21 @@ pub async fn start(
|
||||||
.map_err(|e| WebServerError::Startup(e.to_string()))?;
|
.map_err(|e| WebServerError::Startup(e.to_string()))?;
|
||||||
let cancellation = parent.child_token();
|
let cancellation = parent.child_token();
|
||||||
let task_cancellation = cancellation.clone();
|
let task_cancellation = cancellation.clone();
|
||||||
|
let session_limit = Arc::new(tokio::sync::Semaphore::new(config.max_mtp_sessions));
|
||||||
let join = tokio::spawn(async move {
|
let join = tokio::spawn(async move {
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
result = server.accept() => match result {
|
result = server.accept() => match result {
|
||||||
Ok(Some(connection)) => {
|
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();
|
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,
|
Ok(None) => break,
|
||||||
Err(error) => log!("MTP webserver connection failed: {}", error),
|
Err(error) => log!("MTP webserver connection failed: {}", error),
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,9 @@ name = "web-ui"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
|
|
||||||
|
[features]
|
||||||
|
legacy-web-admin = []
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
iota-storage = { path = "../iota-storage" }
|
iota-storage = { path = "../iota-storage" }
|
||||||
iota-state = { path = "../iota-state" }
|
iota-state = { path = "../iota-state" }
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,5 @@
|
||||||
|
#[cfg(feature = "legacy-web-admin")]
|
||||||
pub mod api;
|
pub mod api;
|
||||||
|
#[cfg(feature = "legacy-web-admin")]
|
||||||
pub mod server;
|
pub mod server;
|
||||||
pub mod web_path_parser;
|
pub mod web_path_parser;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue