From 01f1834e25a710b6bb66cd898e83c7107d41d51c Mon Sep 17 00:00:00 2001 From: Alois Date: Sat, 19 Sep 2026 23:04:22 +0200 Subject: [PATCH] feat(settings): add synced user assets --- Cargo.lock | 26 +- client/Cargo.toml | 2 +- flake.lock | 17 + flake.nix | 5 + iota-auth/Cargo.toml | 2 +- iota-connection/Cargo.toml | 2 +- iota-connection/src/message_handlers.rs | 283 +++++++++++ iota-connection/src/relay.rs | 1 + iota-daemon-lib/Cargo.toml | 2 +- iota-identity/Cargo.toml | 2 +- iota-logger/Cargo.toml | 2 +- iota-storage/Cargo.toml | 2 +- iota-storage/src/users/user_manager.rs | 1 + iota-storage/src/util/db.rs | 50 +- iota-storage/src/util/mod.rs | 1 + iota-storage/src/util/sync.rs | 2 + iota-storage/src/util/user_assets.rs | 516 ++++++++++++++++++++ iota-storage/src/util/user_blobs.rs | 23 +- iota-util/Cargo.toml | 2 +- mtp-type-maps | 2 +- omikron-connector/Cargo.toml | 2 +- omikron-connector/src/omikron_connection.rs | 60 +++ other-iota/Cargo.toml | 2 +- web-server/Cargo.toml | 2 +- 24 files changed, 959 insertions(+), 50 deletions(-) create mode 100644 iota-storage/src/util/user_assets.rs diff --git a/Cargo.lock b/Cargo.lock index 1627bab..f4f97ff 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2833,7 +2833,7 @@ dependencies = [ [[package]] name = "mtp" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "mtp-client", "mtp-codec", @@ -2850,7 +2850,7 @@ dependencies = [ [[package]] name = "mtp-client" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "mtp-codec", "mtp-common", @@ -2863,7 +2863,7 @@ dependencies = [ [[package]] name = "mtp-codec" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "base64 0.23.1", "byteorder", @@ -2877,7 +2877,7 @@ dependencies = [ [[package]] name = "mtp-common" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "rand 0.10.2", "thiserror 2.0.20", @@ -2886,7 +2886,7 @@ dependencies = [ [[package]] name = "mtp-core" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "async-trait", "mtp-codec", @@ -2898,7 +2898,7 @@ dependencies = [ [[package]] name = "mtp-crypto" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "argon2", "base64 0.22.1", @@ -2921,7 +2921,7 @@ dependencies = [ [[package]] name = "mtp-files" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "mtp-crypto", "rand 0.10.2", @@ -2932,7 +2932,7 @@ dependencies = [ [[package]] name = "mtp-h3" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "async-trait", "bytes", @@ -2949,7 +2949,7 @@ dependencies = [ [[package]] name = "mtp-host" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "mtp-codec", "mtp-common", @@ -2966,7 +2966,7 @@ dependencies = [ [[package]] name = "mtp-native" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "async-trait", "mtp-codec", @@ -2988,7 +2988,7 @@ dependencies = [ [[package]] name = "mtp-type-map" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "serde", "serde_yaml", @@ -2997,7 +2997,7 @@ dependencies = [ [[package]] name = "mtp-webserver" version = "0.3.0" -source = "git+https://git.methanium.net/Methanium/mtp.git?rev=1f19a0d897c265d1e3f590a876f95e766ff99318#1f19a0d897c265d1e3f590a876f95e766ff99318" +source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "bytes", "h3", @@ -4672,7 +4672,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.4.3", + "getrandom 0.3.4", "once_cell", "rustix", "windows-sys 0.61.2", diff --git a/client/Cargo.toml b/client/Cargo.toml index 3995595..d1d8ac1 100644 --- a/client/Cargo.toml +++ b/client/Cargo.toml @@ -5,7 +5,7 @@ edition = "2024" [dependencies] async-trait = "0.1.89" -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["client", "crypto", "pipes", "web-server"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["client", "crypto", "pipes", "web-server"] } iota-auth = { path = "../iota-auth" } iota-identity = { path = "../iota-identity" } iota-connection = { path = "../iota-connection" } diff --git a/flake.lock b/flake.lock index 3042264..e9fd017 100644 --- a/flake.lock +++ b/flake.lock @@ -18,6 +18,22 @@ "type": "github" } }, + "mtp-type-maps": { + "flake": false, + "locked": { + "lastModified": 1789843665, + "narHash": "sha256-dU7z909w1R6OAVHLf6DsNWDQxx2SDiHj3w7574a7MsA=", + "ref": "refs/heads/main", + "rev": "c1937a37393907bf610ac4f6ba303f19d3071f9b", + "revCount": 24, + "type": "git", + "url": "https://git.methanium.net/tensamin/mtp-type-maps" + }, + "original": { + "type": "git", + "url": "https://git.methanium.net/tensamin/mtp-type-maps" + } + }, "nixpkgs": { "locked": { "lastModified": 1784497964, @@ -52,6 +68,7 @@ "root": { "inputs": { "flake-parts": "flake-parts", + "mtp-type-maps": "mtp-type-maps", "nixpkgs": "nixpkgs", "rust-overlay": "rust-overlay" } diff --git a/flake.nix b/flake.nix index 3f79069..6dd584b 100644 --- a/flake.nix +++ b/flake.nix @@ -8,6 +8,10 @@ url = "github:oxalica/rust-overlay"; inputs.nixpkgs.follows = "nixpkgs"; }; + mtp-type-maps = { + url = "git+https://git.methanium.net/tensamin/mtp-type-maps"; + flake = false; + }; }; outputs = inputs @ { @@ -63,6 +67,7 @@ nativeBuildInputs = commonNativeBuildInputs; buildInputs = commonBuildInputs; dontUseCmakeConfigure = true; + MTP_TYPE_MAPS = "${inputs.mtp-type-maps}/type-maps.yaml"; passthru.dataDir = "/var/lib/iota"; }; diff --git a/iota-auth/Cargo.toml b/iota-auth/Cargo.toml index 719be00..9a135da 100644 --- a/iota-auth/Cargo.toml +++ b/iota-auth/Cargo.toml @@ -7,7 +7,7 @@ edition = "2024" async-trait = "0.1.89" dashmap = "6.2.1" iota-identity = { path = "../iota-identity" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["crypto"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto"] } rand_core = { version = "0.6", features = ["getrandom", "std"] } uuid = { version = "*", features = ["v4"] } diff --git a/iota-connection/Cargo.toml b/iota-connection/Cargo.toml index 8e03415..bdcd27c 100644 --- a/iota-connection/Cargo.toml +++ b/iota-connection/Cargo.toml @@ -11,7 +11,7 @@ iota-identity = { path = "../iota-identity" } iota-logger = { path = "../iota-logger" } iota-storage = { path = "../iota-storage" } iota-util = { path = "../iota-util" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["crypto"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto"] } serde = { version = "1", features = ["derive"] } serde_json = "1" tokio = { version = "1.50.0", features = ["rt"] } diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index fe1880e..ce2b6d0 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -5,6 +5,7 @@ use iota_storage::util::communities_util::CommunitiesUtil; use iota_storage::util::e2ee_storage::{self, ChatSecretQuery}; use iota_storage::util::settings; use iota_storage::util::synced_settings::{self, SettingScope, SyncedSetting}; +use iota_storage::util::user_assets::{self, UploadStatus, UserAssetMetadata}; use iota_storage::util::user_blobs::{self, UserBlob, UserBlobMetadata}; use iota_storage::util::{blocked_users, message_storage_policy, receipt_policy}; use mtp::codec::{ @@ -42,6 +43,8 @@ pub struct BlobMutation { pub changed: Option, } +pub type AssetMutation = BlobMutation; + #[derive(Debug)] pub struct PolicyMutation { pub response: CommunicationValue, @@ -1937,6 +1940,109 @@ fn blob_request(cv: &CommunicationValue) -> Result<(i64, String), CommunicationV Ok((user_id, id.to_owned())) } +fn asset_metadata_value(asset: &UserAssetMetadata) -> DataValue { + typed_container(vec![ + (DataType::AssetId, DataValue::Str(asset.asset_id.clone())), + (DataType::Size, DataValue::SignedNumber(asset.size.into())), + (DataType::Sha256, DataValue::Bytes(asset.sha256.clone())), + (DataType::MimeType, DataValue::Str(asset.mime_type.clone())), + ( + DataType::VersionNumber, + DataValue::SignedNumber(asset.revision.into()), + ), + ( + DataType::CreatedAt, + DataValue::SignedNumber(asset.created_at.into()), + ), + ( + DataType::UpdatedAt, + DataValue::SignedNumber(asset.updated_at.into()), + ), + ]) +} + +fn asset_response( + cv: &CommunicationValue, + ty: CommunicationType, + asset: &UserAssetMetadata, +) -> CommunicationValue { + CommunicationValue::new(ty) + .with_request_id(cv) + .with_receiver(sender_wire_id(asset.user_id)) + .add_typed_default(DataType::AssetId, DataValue::Str(asset.asset_id.clone())) + .add_typed_default(DataType::Size, DataValue::SignedNumber(asset.size.into())) + .add_typed_default(DataType::Sha256, DataValue::Bytes(asset.sha256.clone())) + .add_typed_default(DataType::MimeType, DataValue::Str(asset.mime_type.clone())) + .add_typed_default( + DataType::VersionNumber, + DataValue::SignedNumber(asset.revision.into()), + ) + .add_typed_default( + DataType::CreatedAt, + DataValue::SignedNumber(asset.created_at.into()), + ) + .add_typed_default( + DataType::UpdatedAt, + DataValue::SignedNumber(asset.updated_at.into()), + ) +} + +fn upload_response( + cv: &CommunicationValue, + ty: CommunicationType, + status: &UploadStatus, +) -> CommunicationValue { + let mut response = CommunicationValue::new(ty) + .with_request_id(cv) + .with_receiver(cv.sender().unwrap_or_default()) + .add_typed_default(DataType::UploadId, DataValue::Str(status.upload_id.clone())) + .add_typed_default(DataType::AssetId, DataValue::Str(status.asset_id.clone())) + .add_typed_default(DataType::Size, DataValue::SignedNumber(status.size.into())) + .add_typed_default(DataType::Sha256, DataValue::Bytes(status.sha256.clone())) + .add_typed_default(DataType::MimeType, DataValue::Str(status.mime_type.clone())) + .add_typed_default( + DataType::Offset, + DataValue::SignedNumber(status.offset.into()), + ) + .add_typed_default(DataType::State, DataValue::Str(status.state.clone())) + .add_typed_default( + DataType::CreatedAt, + DataValue::SignedNumber(status.created_at.into()), + ) + .add_typed_default( + DataType::UpdatedAt, + DataValue::SignedNumber(status.updated_at.into()), + ); + if let Some(revision) = status.revision { + response = response.add_typed_default( + DataType::VersionNumber, + DataValue::SignedNumber(revision.into()), + ); + } + if let Some(committed_at) = status.committed_at { + response = response.add_typed_default( + DataType::CommittedAt, + DataValue::SignedNumber(committed_at.into()), + ); + } + response +} + +fn asset_error(cv: &CommunicationValue, error: StorageError) -> CommunicationValue { + error_response(cv, CommunicationType::ErrorInternal) + .add_typed_default(DataType::ErrorType, DataValue::Str(error.to_string())) +} + +fn upload_request(cv: &CommunicationValue) -> Result<(i64, String), CommunicationValue> { + let user_id = required_sender_id(cv)?; + let upload_id = cv + .get_data(DataType::UploadId) + .as_str() + .filter(|id| !id.is_empty()) + .ok_or_else(|| error_response(cv, CommunicationType::ErrorInvalidData))?; + Ok((user_id, upload_id.to_owned())) +} + fn setting_changed(setting: &SyncedSetting) -> CommunicationValue { CommunicationValue::new(CommunicationType::SyncedSettingChanged) .with_receiver(sender_wire_id(setting.user_id)) @@ -2255,6 +2361,183 @@ pub fn handle_user_blob_list(cv: &CommunicationValue) -> CommunicationValue { } } +pub fn handle_user_asset_upload_start(cv: &CommunicationValue) -> CommunicationValue { + let user_id = match required_sender_id(cv) { + Ok(id) if id > 0 => id, + _ => return error_response(cv, CommunicationType::ErrorInvalidData), + }; + let Some(upload_id) = cv.get_data(DataType::UploadId).as_str() else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(asset_id) = cv.get_data(DataType::AssetId).as_str() else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(size) = data_i64(cv, DataType::Size) else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(sha256) = cv.get_data(DataType::Sha256).and_then(DataValue::as_bytes) else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(mime_type) = cv.get_data(DataType::MimeType).as_str() else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + match user_assets::start(user_id, upload_id, asset_id, size, &sha256, mime_type) { + Ok(status) => upload_response(cv, CommunicationType::UserAssetUploadStart, &status), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_upload_chunk(cv: &CommunicationValue) -> CommunicationValue { + let (user_id, upload_id) = match upload_request(cv) { + Ok(value) => value, + Err(response) => return response, + }; + let Some(offset) = data_i64(cv, DataType::Offset) else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(chunk) = cv.get_data(DataType::Blob).and_then(DataValue::as_bytes) else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + match user_assets::append(user_id, &upload_id, offset, &chunk) { + Ok(status) => upload_response(cv, CommunicationType::UserAssetUploadChunk, &status), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_upload_status(cv: &CommunicationValue) -> CommunicationValue { + let (user_id, upload_id) = match upload_request(cv) { + Ok(value) => value, + Err(response) => return response, + }; + match user_assets::status(user_id, &upload_id) { + Ok(Some(status)) => upload_response(cv, CommunicationType::UserAssetUploadStatus, &status), + Ok(None) => error_response(cv, CommunicationType::ErrorNotFound), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_upload_commit(cv: &CommunicationValue) -> AssetMutation { + let (user_id, upload_id) = match upload_request(cv) { + Ok(value) => value, + Err(response) => { + return AssetMutation { + response, + changed: None, + }; + } + }; + match user_assets::commit(user_id, &upload_id) { + Ok((asset, changed)) => AssetMutation { + response: asset_response(cv, CommunicationType::UserAssetUploadCommit, &asset), + changed: changed.then(|| { + asset_response(cv, CommunicationType::UserAssetChanged, &asset) + .with_id(next_notification_id()) + .add_typed_default(DataType::Deleted, DataValue::Bool(false)) + }), + }, + Err(error) => AssetMutation { + response: asset_error(cv, error), + changed: None, + }, + } +} + +pub fn handle_user_asset_upload_abort(cv: &CommunicationValue) -> CommunicationValue { + let (user_id, upload_id) = match upload_request(cv) { + Ok(value) => value, + Err(response) => return response, + }; + match user_assets::abort(user_id, &upload_id) { + Ok(_) => CommunicationValue::new(CommunicationType::UserAssetUploadAbort) + .with_request_id(cv) + .with_receiver(sender_wire_id(user_id)) + .add_typed_default(DataType::UploadId, DataValue::Str(upload_id)), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_get_chunk(cv: &CommunicationValue) -> CommunicationValue { + let user_id = match required_sender_id(cv) { + Ok(id) if id > 0 => id, + _ => return error_response(cv, CommunicationType::ErrorInvalidData), + }; + let Some(asset_id) = cv.get_data(DataType::AssetId).as_str() else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(offset) = data_i64(cv, DataType::Offset) else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + let Some(amount) = data_i64(cv, DataType::Amount).and_then(|value| usize::try_from(value).ok()) + else { + return error_response(cv, CommunicationType::ErrorInvalidData); + }; + match user_assets::read_chunk(user_id, asset_id, offset, amount) { + Ok(Some((asset, bytes))) => { + asset_response(cv, CommunicationType::UserAssetGetChunk, &asset) + .add_typed_default(DataType::Offset, DataValue::SignedNumber(offset.into())) + .add_typed_default(DataType::Blob, DataValue::Bytes(bytes)) + } + Ok(None) => error_response(cv, CommunicationType::ErrorNotFound), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_list(cv: &CommunicationValue) -> CommunicationValue { + let user_id = match required_sender_id(cv) { + Ok(id) if id > 0 => id, + _ => return error_response(cv, CommunicationType::ErrorInvalidData), + }; + match user_assets::list(user_id) { + Ok(assets) => CommunicationValue::new(CommunicationType::UserAssetList) + .with_request_id(cv) + .with_receiver(sender_wire_id(user_id)) + .add_typed_default( + DataType::Assets, + DataValue::Array(assets.iter().map(asset_metadata_value).collect()), + ), + Err(error) => asset_error(cv, error), + } +} + +pub fn handle_user_asset_delete(cv: &CommunicationValue) -> AssetMutation { + let user_id = match required_sender_id(cv) { + Ok(id) if id > 0 => id, + _ => { + return AssetMutation { + response: error_response(cv, CommunicationType::ErrorInvalidData), + changed: None, + }; + } + }; + let Some(asset_id) = cv.get_data(DataType::AssetId).as_str() else { + return AssetMutation { + response: error_response(cv, CommunicationType::ErrorInvalidData), + changed: None, + }; + }; + match user_assets::delete(user_id, asset_id) { + Ok(Some(asset)) => AssetMutation { + response: asset_response(cv, CommunicationType::UserAssetDelete, &asset), + changed: Some( + asset_response(cv, CommunicationType::UserAssetChanged, &asset) + .with_id(next_notification_id()) + .add_typed_default(DataType::Deleted, DataValue::Bool(true)), + ), + }, + Ok(None) => AssetMutation { + response: CommunicationValue::new(CommunicationType::UserAssetDelete) + .with_request_id(cv) + .with_receiver(sender_wire_id(user_id)) + .add_typed_default(DataType::AssetId, DataValue::Str(asset_id.to_owned())), + changed: None, + }, + Err(error) => AssetMutation { + response: asset_error(cv, error), + changed: None, + }, + } +} + fn authenticated_user(cv: &CommunicationValue) -> Result { let user_id = required_sender_id(cv)?; if user_id <= 0 { diff --git a/iota-connection/src/relay.rs b/iota-connection/src/relay.rs index 9a2745c..5ffcb24 100644 --- a/iota-connection/src/relay.rs +++ b/iota-connection/src/relay.rs @@ -66,6 +66,7 @@ mod security_tests { CommunicationType::SyncedSettingDelete, CommunicationType::SyncedSettingsList, CommunicationType::SyncedSettingChanged, + CommunicationType::UserAssetChanged, ] { assert_eq!( message_security_class(&CommunicationValue::new(setting_type)), diff --git a/iota-daemon-lib/Cargo.toml b/iota-daemon-lib/Cargo.toml index dc665ab..8f5dc37 100644 --- a/iota-daemon-lib/Cargo.toml +++ b/iota-daemon-lib/Cargo.toml @@ -18,7 +18,7 @@ iota-updater = { path = "../iota-updater" } iota-util = { path = "../iota-util" } omikron-connector = { path = "../omikron-connector" } other-iota = { path = "../other-iota" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318" } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828" } libc = "0.2" sysinfo = "0.38.0" serde_yaml = "0.9" diff --git a/iota-identity/Cargo.toml b/iota-identity/Cargo.toml index 45b3298..7508211 100644 --- a/iota-identity/Cargo.toml +++ b/iota-identity/Cargo.toml @@ -6,7 +6,7 @@ edition = "2024" [dependencies] async-trait = "0.1.89" base64 = "0.22.1" -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["crypto", "files", "raw"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto", "files", "raw"] } serde = { version = "1", features = ["derive"] } [dev-dependencies] diff --git a/iota-logger/Cargo.toml b/iota-logger/Cargo.toml index 6efc703..fb7537f 100644 --- a/iota-logger/Cargo.toml +++ b/iota-logger/Cargo.toml @@ -8,7 +8,7 @@ iota-paths = { path = "../iota-paths" } iota-state = { path = "../iota-state" } iota-util = { path = "../iota-util" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318" } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828" } ratatui = "0.30.0" json = "0.12.4" diff --git a/iota-storage/Cargo.toml b/iota-storage/Cargo.toml index 4a77d10..6723914 100644 --- a/iota-storage/Cargo.toml +++ b/iota-storage/Cargo.toml @@ -24,5 +24,5 @@ rusqlite = "0.40.0" tokio = { version = "1.50.0", features = ["full"] } [dev-dependencies] -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["crypto"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["crypto"] } tempfile = "3" diff --git a/iota-storage/src/users/user_manager.rs b/iota-storage/src/users/user_manager.rs index 81690e3..19e4ecd 100644 --- a/iota-storage/src/users/user_manager.rs +++ b/iota-storage/src/users/user_manager.rs @@ -393,6 +393,7 @@ fn run_purge_stages( } fn purge_database_rows(user_id: i64) -> Result<(), crate::storage_error::StorageError> { + crate::util::user_assets::purge_user(user_id)?; db::with_db(|conn| { let tx = conn.unchecked_transaction()?; tx.execute( diff --git a/iota-storage/src/util/db.rs b/iota-storage/src/util/db.rs index e0cab2b..b0f1404 100644 --- a/iota-storage/src/util/db.rs +++ b/iota-storage/src/util/db.rs @@ -1959,6 +1959,46 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { )?; } + if current_version < 45 { + conn.execute_batch( + r#" + CREATE TABLE user_assets ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER NOT NULL, + asset_id TEXT NOT NULL, + file_name TEXT NOT NULL, + size INTEGER NOT NULL, + sha256 BLOB NOT NULL, + mime_type TEXT NOT NULL, + revision INTEGER NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL, + published_upload_id TEXT NOT NULL, + deleted INTEGER NOT NULL DEFAULT 0, + UNIQUE(user_id, asset_id) + ); + CREATE INDEX user_assets_owner ON user_assets (user_id, deleted); + CREATE TABLE user_asset_uploads ( + upload_id TEXT PRIMARY KEY, + user_id INTEGER NOT NULL, + asset_id TEXT NOT NULL, + file_name TEXT NOT NULL, + size INTEGER NOT NULL, + sha256 BLOB NOT NULL, + mime_type TEXT NOT NULL, + received INTEGER NOT NULL DEFAULT 0, + state TEXT NOT NULL CHECK (state IN ('uploading', 'committed')), + revision INTEGER, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL, + committed_at INTEGER + ); + CREATE INDEX user_asset_uploads_owner ON user_asset_uploads (user_id, state, updated_at); + PRAGMA user_version = 45; + "#, + )?; + } + Ok(()) } @@ -2033,7 +2073,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 44); + assert_eq!(version, 45); for column in ["height", "reply_to", "edited_count", "deleted_by_external"] { let mut statement = conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?; @@ -2052,7 +2092,7 @@ mod tests { run_migrations_on_connection(&conn)?; run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 44); + assert_eq!(version, 45); for table in [ "sync_heads", "sync_events", @@ -2133,7 +2173,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 44); + assert_eq!(version, 45); for column in [ "id", "user_id", @@ -2228,7 +2268,7 @@ mod tests { )?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; assert_eq!(preserved, "remote_committed"); - assert_eq!(version, 44); + assert_eq!(version, 45); Ok(()) } @@ -2266,7 +2306,7 @@ mod tests { })?; assert_eq!(count, 0); let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 44); + assert_eq!(version, 45); Ok(()) } diff --git a/iota-storage/src/util/mod.rs b/iota-storage/src/util/mod.rs index bcf195e..c9e15b5 100644 --- a/iota-storage/src/util/mod.rs +++ b/iota-storage/src/util/mod.rs @@ -17,4 +17,5 @@ pub mod relay_replay; pub mod settings; pub mod sync; pub mod synced_settings; +pub mod user_assets; pub mod user_blobs; diff --git a/iota-storage/src/util/sync.rs b/iota-storage/src/util/sync.rs index 9d31474..bee377d 100644 --- a/iota-storage/src/util/sync.rs +++ b/iota-storage/src/util/sync.rs @@ -13,6 +13,7 @@ pub enum EntityType { Contact, Setting, UserBlob, + UserAsset, BlockedUser, ReceiptPolicy, MessageStoragePolicy, @@ -24,6 +25,7 @@ impl EntityType { Self::Contact => "contact", Self::Setting => "setting", Self::UserBlob => "user_blob", + Self::UserAsset => "user_asset", Self::BlockedUser => "blocked_user", Self::ReceiptPolicy => "receipt_policy", Self::MessageStoragePolicy => "message_storage_policy", diff --git a/iota-storage/src/util/user_assets.rs b/iota-storage/src/util/user_assets.rs new file mode 100644 index 0000000..0a9d3ea --- /dev/null +++ b/iota-storage/src/util/user_assets.rs @@ -0,0 +1,516 @@ +use crate::storage_error::StorageError; +use crate::util::{db, sync}; +use rand::RngCore; +use rusqlite::{OptionalExtension, Row, params}; +use sha2::{Digest, Sha256}; +use std::collections::HashSet; +use std::fs::{self, File, OpenOptions}; +use std::io::{Read, Seek, SeekFrom, Write}; +use std::path::PathBuf; +use std::sync::atomic::{AtomicI64, Ordering}; + +pub const MAX_CHUNK_BYTES: usize = 8 * 1024 * 1024; +const ABANDONED_UPLOAD_MS: i64 = 24 * 60 * 60 * 1_000; +const CLEANUP_INTERVAL_MS: i64 = 15 * 60 * 1_000; +static LAST_CLEANUP: AtomicI64 = AtomicI64::new(0); + +#[derive(Debug, Clone)] +pub struct UserAssetMetadata { + pub id: i64, + pub user_id: i64, + pub asset_id: String, + pub size: i64, + pub sha256: Vec, + pub mime_type: String, + pub revision: i64, + pub created_at: i64, + pub updated_at: i64, +} + +#[derive(Debug, Clone)] +pub struct UploadStatus { + pub upload_id: String, + pub asset_id: String, + pub size: i64, + pub sha256: Vec, + pub mime_type: String, + pub offset: i64, + pub state: String, + pub revision: Option, + pub created_at: i64, + pub updated_at: i64, + pub committed_at: Option, +} + +fn asset_dir() -> Result { + let path = iota_util::file_util::storage_directory().join("assets"); + fs::create_dir_all(&path)?; + Ok(path) +} + +fn random_file_name() -> String { + let mut bytes = [0_u8; 24]; + rand::thread_rng().fill_bytes(&mut bytes); + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +fn validate(user_id: i64, asset_id: &str, upload_id: Option<&str>) -> Result<(), StorageError> { + if user_id <= 0 || asset_id.is_empty() || asset_id.len() > 256 { + return Err(StorageError::Other("invalid asset owner or id".into())); + } + if upload_id.is_some_and(|id| id.is_empty() || id.len() > 128) { + return Err(StorageError::Other("invalid upload id".into())); + } + Ok(()) +} + +fn upload_from_row(row: &Row<'_>) -> rusqlite::Result<(UploadStatus, String)> { + Ok(( + UploadStatus { + upload_id: row.get(0)?, + asset_id: row.get(2)?, + size: row.get(4)?, + sha256: row.get(5)?, + mime_type: row.get(6)?, + offset: row.get(7)?, + state: row.get(8)?, + revision: row.get(9)?, + created_at: row.get(10)?, + updated_at: row.get(11)?, + committed_at: row.get(12)?, + }, + row.get(3)?, + )) +} + +fn metadata_from_row(row: &Row<'_>) -> rusqlite::Result { + Ok(UserAssetMetadata { + id: row.get(0)?, + user_id: row.get(1)?, + asset_id: row.get(2)?, + size: row.get(4)?, + sha256: row.get(5)?, + mime_type: row.get(6)?, + revision: row.get(7)?, + created_at: row.get(8)?, + updated_at: row.get(9)?, + }) +} + +fn load_upload( + user_id: i64, + upload_id: &str, +) -> Result, StorageError> { + db::with_db(|conn| { + Ok(conn.query_row( + "SELECT upload_id, user_id, asset_id, file_name, size, sha256, mime_type, received, state, revision, created_at, updated_at, committed_at FROM user_asset_uploads WHERE user_id = ?1 AND upload_id = ?2", + params![user_id, upload_id], + upload_from_row, + ).optional()?) + }) +} + +fn reconcile_offset( + user_id: i64, + mut status: UploadStatus, + file_name: &str, +) -> Result { + if status.state != "uploading" { + return Ok(status); + } + let actual = fs::metadata(asset_dir()?.join(file_name))?.len() as i64; + if actual > status.size { + return Err(StorageError::Other( + "upload file exceeds declared size".into(), + )); + } + if actual != status.offset { + let now = sync::now_millis(); + db::with_db(|conn| { + conn.execute( + "UPDATE user_asset_uploads SET received = ?3, updated_at = ?4 WHERE user_id = ?1 AND upload_id = ?2 AND state = 'uploading'", + params![user_id, status.upload_id, actual, now], + )?; + Ok(()) + })?; + status.offset = actual; + status.updated_at = now; + } + Ok(status) +} + +pub fn start( + user_id: i64, + upload_id: &str, + asset_id: &str, + size: i64, + sha256: &[u8], + mime_type: &str, +) -> Result { + validate(user_id, asset_id, Some(upload_id))?; + if size < 0 || sha256.len() != 32 || mime_type.is_empty() || mime_type.len() > 255 { + return Err(StorageError::Other("invalid asset metadata".into())); + } + cleanup_if_due()?; + if let Some((status, file_name)) = load_upload(user_id, upload_id)? { + if status.asset_id != asset_id + || status.size != size + || status.sha256 != sha256 + || status.mime_type != mime_type + { + return Err(StorageError::Other( + "upload id already has different metadata".into(), + )); + } + return reconcile_offset(user_id, status, &file_name); + } + + let file_name = random_file_name(); + let path = asset_dir()?.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( + "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], + )?; + Ok(()) + }); + if let Err(error) = inserted { + let _ = fs::remove_file(path); + return Err(error); + } + Ok(UploadStatus { + upload_id: upload_id.into(), + asset_id: asset_id.into(), + size, + sha256: sha256.into(), + mime_type: mime_type.into(), + offset: 0, + state: "uploading".into(), + revision: None, + created_at: now, + updated_at: now, + committed_at: None, + }) +} + +pub fn status(user_id: i64, upload_id: &str) -> Result, StorageError> { + let Some((status, file_name)) = load_upload(user_id, upload_id)? else { + return Ok(None); + }; + reconcile_offset(user_id, status, &file_name).map(Some) +} + +pub fn append( + user_id: i64, + upload_id: &str, + offset: i64, + chunk: &[u8], +) -> Result { + if chunk.is_empty() || chunk.len() > MAX_CHUNK_BYTES { + return Err(StorageError::Other("invalid asset chunk size".into())); + } + let Some((status, file_name)) = load_upload(user_id, upload_id)? else { + 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 + { + return Err(StorageError::Other( + "asset chunk offset does not match acknowledged offset".into(), + )); + } + let path = asset_dir()?.join(file_name); + let mut file = OpenOptions::new().append(true).open(path)?; + file.write_all(chunk)?; + file.sync_data()?; + let received = offset + chunk.len() as i64; + let now = sync::now_millis(); + db::with_db(|conn| { + conn.execute( + "UPDATE user_asset_uploads SET received = ?3, updated_at = ?4 WHERE user_id = ?1 AND upload_id = ?2 AND state = 'uploading'", + params![user_id, upload_id, received, now], + )?; + Ok(()) + })?; + Ok(UploadStatus { + offset: received, + updated_at: now, + ..status + }) +} + +pub fn commit(user_id: i64, upload_id: &str) -> Result<(UserAssetMetadata, bool), StorageError> { + let Some((status, file_name)) = load_upload(user_id, upload_id)? else { + return Err(StorageError::Other("upload not found".into())); + }; + if status.state == "committed" { + return Ok(( + UserAssetMetadata { + id: 0, + user_id, + asset_id: status.asset_id, + size: status.size, + sha256: status.sha256, + mime_type: status.mime_type, + revision: status.revision.ok_or_else(|| { + StorageError::Other("committed upload has no revision".into()) + })?, + created_at: status.created_at, + updated_at: status.committed_at.unwrap_or(status.updated_at), + }, + false, + )); + } + let status = reconcile_offset(user_id, status, &file_name)?; + if status.offset != status.size { + return Err(StorageError::Other("upload is incomplete".into())); + } + let path = asset_dir()?.join(&file_name); + let mut file = File::open(&path)?; + let mut hash = Sha256::new(); + let mut buffer = [0_u8; 64 * 1024]; + loop { + let read = file.read(&mut buffer)?; + if read == 0 { + break; + } + hash.update(&buffer[..read]); + } + if hash.finalize().as_slice() != status.sha256 { + return Err(StorageError::Other("asset SHA-256 mismatch".into())); + } + + let now = sync::now_millis(); + let (metadata, replaced_file) = db::with_immediate_transaction(|tx| { + let existing = tx.query_row( + "SELECT id, file_name, created_at FROM user_assets WHERE user_id = ?1 AND asset_id = ?2", + params![user_id, status.asset_id], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?, row.get::<_, i64>(2)?)), + ).optional()?; + let id = if let Some(existing) = &existing { + existing.0 + } else { + tx.execute( + "INSERT INTO user_assets (user_id, asset_id, file_name, size, sha256, mime_type, revision, created_at, updated_at, published_upload_id, deleted) VALUES (?1, ?2, ?3, ?4, ?5, ?6, 0, ?7, ?7, ?8, 0)", + params![user_id, status.asset_id, file_name, status.size, status.sha256, status.mime_type, now, upload_id], + )?; + tx.last_insert_rowid() + }; + let revision = sync::record_event( + tx, + user_id, + sync::EntityType::UserAsset, + id, + sync::Operation::Upsert, + )?; + let created_at = existing.as_ref().map(|value| value.2).unwrap_or(now); + tx.execute( + "UPDATE user_assets SET file_name = ?2, size = ?3, sha256 = ?4, mime_type = ?5, revision = ?6, updated_at = ?7, published_upload_id = ?8, deleted = 0 WHERE id = ?1", + params![id, file_name, status.size, status.sha256, status.mime_type, revision, now, upload_id], + )?; + tx.execute( + "UPDATE user_asset_uploads SET state = 'committed', revision = ?3, updated_at = ?4, committed_at = ?4 WHERE user_id = ?1 AND upload_id = ?2", + params![user_id, upload_id, revision, now], + )?; + Ok(( + UserAssetMetadata { + id, + user_id, + asset_id: status.asset_id.clone(), + size: status.size, + sha256: status.sha256.clone(), + mime_type: status.mime_type.clone(), + revision, + created_at, + updated_at: now, + }, + existing + .map(|value| value.1) + .filter(|old| old != &file_name), + )) + })?; + if let Some(old) = replaced_file { + let _ = fs::remove_file(asset_dir()?.join(old)); + } + Ok((metadata, true)) +} + +pub fn abort(user_id: i64, upload_id: &str) -> Result { + let upload = load_upload(user_id, upload_id)?; + let Some((status, file_name)) = upload else { + return Ok(false); + }; + if status.state == "committed" { + return Ok(false); + } + db::with_db(|conn| { + conn.execute("DELETE FROM user_asset_uploads WHERE user_id = ?1 AND upload_id = ?2 AND state = 'uploading'", params![user_id, upload_id])?; + Ok(()) + })?; + match fs::remove_file(asset_dir()?.join(file_name)) { + Ok(()) => Ok(true), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(true), + Err(error) => Err(error.into()), + } +} + +pub fn list(user_id: i64) -> Result, StorageError> { + 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 ORDER BY updated_at DESC")?; + statement + .query_map([user_id], metadata_from_row)? + .collect::, _>>() + .map_err(StorageError::from) + }) +} + +pub fn read_chunk( + user_id: i64, + asset_id: &str, + offset: i64, + amount: usize, +) -> Result)>, StorageError> { + validate(user_id, asset_id, None)?; + if offset < 0 || amount == 0 || amount > MAX_CHUNK_BYTES { + return Err(StorageError::Other("invalid asset chunk range".into())); + } + let value = db::with_db(|conn| { + Ok(conn.query_row( + "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 asset_id = ?2 AND deleted = 0", + params![user_id, asset_id], + |row| Ok((metadata_from_row(row)?, row.get::<_, String>(3)?)), + ).optional()?) + })?; + let Some((metadata, file_name)) = value else { + return Ok(None); + }; + if offset > metadata.size { + return Err(StorageError::Other("asset offset exceeds size".into())); + } + let mut file = File::open(asset_dir()?.join(file_name))?; + file.seek(SeekFrom::Start(offset as u64))?; + let mut bytes = vec![0; amount.min((metadata.size - offset) as usize)]; + file.read_exact(&mut bytes)?; + Ok(Some((metadata, bytes))) +} + +pub fn delete(user_id: i64, asset_id: &str) -> Result, StorageError> { + validate(user_id, asset_id, None)?; + let deleted = db::with_immediate_transaction(|tx| { + let value = tx.query_row( + "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 asset_id = ?2 AND deleted = 0", + params![user_id, asset_id], + |row| Ok((metadata_from_row(row)?, row.get::<_, String>(3)?)), + ).optional()?; + let Some((mut metadata, file_name)) = value else { + return Ok(None); + }; + metadata.revision = sync::record_event( + tx, + user_id, + sync::EntityType::UserAsset, + metadata.id, + sync::Operation::Delete, + )?; + metadata.updated_at = sync::now_millis(); + tx.execute( + "UPDATE user_assets SET deleted = 1, revision = ?2, updated_at = ?3 WHERE id = ?1", + params![metadata.id, metadata.revision, metadata.updated_at], + )?; + Ok(Some((metadata, file_name))) + })?; + let Some((metadata, file_name)) = deleted else { + return Ok(None); + }; + let _ = fs::remove_file(asset_dir()?.join(file_name)); + Ok(Some(metadata)) +} + +pub fn purge_user(user_id: i64) -> Result<(), StorageError> { + let files = db::with_immediate_transaction(|tx| { + let mut files = HashSet::new(); + for query in [ + "SELECT file_name FROM user_assets WHERE user_id = ?1", + "SELECT file_name FROM user_asset_uploads WHERE user_id = ?1", + ] { + let mut statement = tx.prepare(query)?; + for file in statement.query_map([user_id], |row| row.get::<_, String>(0))? { + files.insert(file?); + } + } + tx.execute("DELETE FROM user_assets WHERE user_id = ?1", [user_id])?; + tx.execute( + "DELETE FROM user_asset_uploads WHERE user_id = ?1", + [user_id], + )?; + Ok(files) + })?; + let directory = asset_dir()?; + for file in files { + match fs::remove_file(directory.join(file)) { + Ok(()) => {} + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(error.into()), + } + } + Ok(()) +} + +fn cleanup_if_due() -> Result<(), StorageError> { + let now = sync::now_millis(); + let previous = LAST_CLEANUP.load(Ordering::Relaxed); + if now - previous < CLEANUP_INTERVAL_MS + || LAST_CLEANUP + .compare_exchange(previous, now, Ordering::Relaxed, Ordering::Relaxed) + .is_err() + { + return Ok(()); + } + let cutoff = now - ABANDONED_UPLOAD_MS; + let abandoned = db::with_db(|conn| { + let mut statement = conn.prepare("SELECT file_name FROM user_asset_uploads WHERE state = 'uploading' AND updated_at < ?1")?; + let files = statement + .query_map([cutoff], |row| row.get::<_, String>(0))? + .collect::, _>>()?; + conn.execute( + "DELETE FROM user_asset_uploads WHERE state = 'uploading' AND updated_at < ?1", + [cutoff], + )?; + Ok(files) + })?; + let dir = asset_dir()?; + for file in abandoned { + let _ = fs::remove_file(dir.join(file)); + } + let referenced = db::with_db(|conn| { + let mut files = HashSet::new(); + for query in [ + "SELECT file_name FROM user_assets WHERE deleted = 0", + "SELECT file_name FROM user_asset_uploads WHERE state = 'uploading'", + ] { + let mut statement = conn.prepare(query)?; + for file in statement.query_map([], |row| row.get::<_, String>(0))? { + files.insert(file?); + } + } + Ok(files) + })?; + for entry in fs::read_dir(dir)? { + let entry = entry?; + if entry.file_type()?.is_file() + && !referenced.contains(&entry.file_name().to_string_lossy().into_owned()) + { + let _ = fs::remove_file(entry.path()); + } + } + Ok(()) +} diff --git a/iota-storage/src/util/user_blobs.rs b/iota-storage/src/util/user_blobs.rs index 1cb7531..3ae9e0a 100644 --- a/iota-storage/src/util/user_blobs.rs +++ b/iota-storage/src/util/user_blobs.rs @@ -74,7 +74,7 @@ pub fn put( user_id: i64, blob_id: &str, blob: &[u8], - expected_revision: Option, + _expected_revision: Option, ) -> Result { validate(user_id, blob_id, Some(blob))?; db::with_immediate_transaction(|tx| { @@ -92,21 +92,7 @@ pub fn put( }, ) .optional()?; - let existing_id = if let Some((id, revision, deleted, _)) = existing { - if deleted { - if !matches!(expected_revision, Some(0)) { - return Err(StorageError::RevisionConflict); - } - } else if expected_revision != Some(revision) { - return Err(StorageError::RevisionConflict); - } - Some(id) - } else { - if !matches!(expected_revision, None | Some(0)) { - return Err(StorageError::RevisionConflict); - } - None - }; + let existing_id = existing.map(|(id, _, _, _)| id); let current_size = existing .filter(|(_, _, deleted, _)| !*deleted) .map(|(_, _, _, size)| size) @@ -233,7 +219,7 @@ fn list_by_query(user_id: i64, query: &str, _: &[i64]) -> Result, pub fn delete( user_id: i64, blob_id: &str, - expected_revision: Option, + _expected_revision: Option, ) -> Result, StorageError> { validate(user_id, blob_id, None)?; db::with_immediate_transaction(|tx| { @@ -261,9 +247,6 @@ pub fn delete( changed: false, })); } - if expected_revision != Some(current) { - return Err(StorageError::RevisionConflict); - } let revision = sync::record_event( tx, user_id, diff --git a/iota-util/Cargo.toml b/iota-util/Cargo.toml index d35e9de..ea23d0b 100644 --- a/iota-util/Cargo.toml +++ b/iota-util/Cargo.toml @@ -6,7 +6,7 @@ edition = "2024" [dependencies] iota-paths = { path = "../iota-paths" } iota-identity = { path = "../iota-identity" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = [ +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = [ "crypto" ] } diff --git a/mtp-type-maps b/mtp-type-maps index 10e418c..c1937a3 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit 10e418ca69c27db0e89ac7dc6a2a348ad9b70017 +Subproject commit c1937a37393907bf610ac4f6ba303f19d3071f9b diff --git a/omikron-connector/Cargo.toml b/omikron-connector/Cargo.toml index fd69dff..e09a5de 100644 --- a/omikron-connector/Cargo.toml +++ b/omikron-connector/Cargo.toml @@ -12,7 +12,7 @@ iota-logger = { path = "../iota-logger" } iota-state = { path = "../iota-state" } iota-storage = { path = "../iota-storage" } iota-util = { path = "../iota-util" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = [ +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = [ "client", "crypto", "files", diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index f5aacd4..d8960b3 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -1730,6 +1730,14 @@ impl OmikronConnection { dispatch!(UserBlobGet, handle_user_blob_get); dispatch!(UserBlobDelete, handle_user_blob_delete); dispatch!(UserBlobList, handle_user_blob_list); + dispatch!(UserAssetUploadStart, handle_user_asset_upload_start); + dispatch!(UserAssetUploadChunk, handle_user_asset_upload_chunk); + dispatch!(UserAssetUploadStatus, handle_user_asset_upload_status); + dispatch!(UserAssetUploadCommit, handle_user_asset_upload_commit); + dispatch!(UserAssetUploadAbort, handle_user_asset_upload_abort); + dispatch!(UserAssetGetChunk, handle_user_asset_get_chunk); + dispatch!(UserAssetList, handle_user_asset_list); + dispatch!(UserAssetDelete, handle_user_asset_delete); dispatch!(UserBlock, handle_user_block); dispatch!(UserUnblock, handle_user_unblock); dispatch!(BlockedUsersGet, handle_blocked_users_get); @@ -2905,6 +2913,58 @@ impl OmikronConnection { .await; } + async fn handle_user_asset_upload_start(self: Arc, cv: &CommunicationValue) { + let _ = self + .send_message(&message_handlers::handle_user_asset_upload_start(cv)) + .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; + } + + async fn handle_user_asset_upload_status(self: Arc, cv: &CommunicationValue) { + let _ = self + .send_message(&message_handlers::handle_user_asset_upload_status(cv)) + .await; + } + + async fn handle_user_asset_upload_commit(self: Arc, cv: &CommunicationValue) { + let mutation = message_handlers::handle_user_asset_upload_commit(cv); + let _ = self.send_message(&mutation.response).await; + if let Some(changed) = mutation.changed { + let _ = self.send_message(&changed).await; + } + } + + async fn handle_user_asset_upload_abort(self: Arc, cv: &CommunicationValue) { + let _ = self + .send_message(&message_handlers::handle_user_asset_upload_abort(cv)) + .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; + } + + async fn handle_user_asset_list(self: Arc, cv: &CommunicationValue) { + let _ = self + .send_message(&message_handlers::handle_user_asset_list(cv)) + .await; + } + + async fn handle_user_asset_delete(self: Arc, cv: &CommunicationValue) { + let mutation = message_handlers::handle_user_asset_delete(cv); + let _ = self.send_message(&mutation.response).await; + if let Some(changed) = mutation.changed { + let _ = self.send_message(&changed).await; + } + } + async fn handle_user_block(self: Arc, cv: &CommunicationValue) { let mutation = message_handlers::handle_user_block(cv); let _ = self.send_message(&mutation.response).await; diff --git a/other-iota/Cargo.toml b/other-iota/Cargo.toml index 863317d..49f8081 100644 --- a/other-iota/Cargo.toml +++ b/other-iota/Cargo.toml @@ -8,7 +8,7 @@ async-trait = "0.1.89" iota-connection = { path = "../iota-connection" } iota-identity = { path = "../iota-identity" } iota-util = { path = "../iota-util" } -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["client", "crypto", "pipes", "web-server"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["client", "crypto", "pipes", "web-server"] } reqwest = { version = "0.13", features = ["json"] } serde = "1" tokio = { version = "1.50.0", features = ["full"] } diff --git a/web-server/Cargo.toml b/web-server/Cargo.toml index 706fa90..73bf2c3 100644 --- a/web-server/Cargo.toml +++ b/web-server/Cargo.toml @@ -6,7 +6,7 @@ edition = "2024" [dependencies] async-trait = "0.1.89" -mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "1f19a0d897c265d1e3f590a876f95e766ff99318", features = ["web-server", "crypto", "pipes"] } +mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735de5ebb36bb41dc699260578341828", features = ["web-server", "crypto", "pipes"] } bytes = "1" http = "1" iota-logger = { path = "../iota-logger" }