[Fix] UserBlobs
This commit is contained in:
parent
8c96f22c16
commit
fa359bc2c8
2 changed files with 54 additions and 42 deletions
|
|
@ -5,7 +5,7 @@ use iota_storage::util::communities_util::CommunitiesUtil;
|
||||||
use iota_storage::util::e2ee_storage::{self, ChatSecretQuery};
|
use iota_storage::util::e2ee_storage::{self, ChatSecretQuery};
|
||||||
use iota_storage::util::settings;
|
use iota_storage::util::settings;
|
||||||
use iota_storage::util::synced_settings::{self, SettingScope, SyncedSetting};
|
use iota_storage::util::synced_settings::{self, SettingScope, SyncedSetting};
|
||||||
use iota_storage::util::user_blobs::{self, UserBlob};
|
use iota_storage::util::user_blobs::{self, UserBlob, UserBlobMetadata};
|
||||||
use iota_storage::util::{blocked_users, message_storage_policy, receipt_policy};
|
use iota_storage::util::{blocked_users, message_storage_policy, receipt_policy};
|
||||||
use mtp::codec::{
|
use mtp::codec::{
|
||||||
CommunicationType, CommunicationValue, DataType, DataValue, TypeMap, VerifiedRelayContent,
|
CommunicationType, CommunicationValue, DataType, DataValue, TypeMap, VerifiedRelayContent,
|
||||||
|
|
@ -1053,28 +1053,6 @@ pub fn handle_client_connected(cv: &CommunicationValue) -> CommunicationValue {
|
||||||
{
|
{
|
||||||
return sync_error(cv);
|
return sync_error(cv);
|
||||||
}
|
}
|
||||||
let (blobs, deleted_blobs) = if mode == "delta" {
|
|
||||||
match sync::delta(user_id, reported_version, head) {
|
|
||||||
Ok(delta) => {
|
|
||||||
let blobs = match user_blobs::list_by_ids(user_id, &delta.blob_upserts) {
|
|
||||||
Ok(blobs) => blobs,
|
|
||||||
Err(_) => return sync_error(cv),
|
|
||||||
};
|
|
||||||
let deleted =
|
|
||||||
match user_blobs::list_deleted_by_ids(user_id, &delta.deleted_blob_ids) {
|
|
||||||
Ok(blobs) => blobs,
|
|
||||||
Err(_) => return sync_error(cv),
|
|
||||||
};
|
|
||||||
(blobs, deleted)
|
|
||||||
}
|
|
||||||
Err(_) => return sync_error(cv),
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
match user_blobs::list(user_id) {
|
|
||||||
Ok(blobs) => (blobs, Vec::new()),
|
|
||||||
Err(_) => return sync_error(cv),
|
|
||||||
}
|
|
||||||
};
|
|
||||||
let blocked_users = match blocked_users::list(user_id) {
|
let blocked_users = match blocked_users::list(user_id) {
|
||||||
Ok(users) => users,
|
Ok(users) => users,
|
||||||
Err(_) => return sync_error(cv),
|
Err(_) => return sync_error(cv),
|
||||||
|
|
@ -1125,19 +1103,6 @@ pub fn handle_client_connected(cv: &CommunicationValue) -> CommunicationValue {
|
||||||
DataType::Settings,
|
DataType::Settings,
|
||||||
DataValue::Array(settings.iter().map(synced_setting_value).collect()),
|
DataValue::Array(settings.iter().map(synced_setting_value).collect()),
|
||||||
)
|
)
|
||||||
.add_typed_default(
|
|
||||||
DataType::Blobs,
|
|
||||||
DataValue::Array(blobs.iter().map(blob_value).collect()),
|
|
||||||
)
|
|
||||||
.add_typed_default(
|
|
||||||
DataType::DeletedBlobIds,
|
|
||||||
DataValue::Array(
|
|
||||||
deleted_blobs
|
|
||||||
.iter()
|
|
||||||
.map(|blob| DataValue::Str(blob.blob_id.clone()))
|
|
||||||
.collect(),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
.add_typed_default(
|
.add_typed_default(
|
||||||
DataType::BlockedUserIds,
|
DataType::BlockedUserIds,
|
||||||
DataValue::Array(
|
DataValue::Array(
|
||||||
|
|
@ -1856,10 +1821,9 @@ fn setting_response(
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn blob_value(blob: &UserBlob) -> DataValue {
|
fn blob_metadata_value(blob: &UserBlobMetadata) -> DataValue {
|
||||||
typed_container(vec![
|
typed_container(vec![
|
||||||
(DataType::BlobId, DataValue::Str(blob.blob_id.clone())),
|
(DataType::BlobId, DataValue::Str(blob.blob_id.clone())),
|
||||||
(DataType::Blob, DataValue::Bytes(blob.blob.clone())),
|
|
||||||
(
|
(
|
||||||
DataType::VersionNumber,
|
DataType::VersionNumber,
|
||||||
DataValue::SignedNumber(blob.revision.into()),
|
DataValue::SignedNumber(blob.revision.into()),
|
||||||
|
|
@ -1891,6 +1855,25 @@ fn blob_response(
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn blob_mutation_response(
|
||||||
|
cv: &CommunicationValue,
|
||||||
|
ty: CommunicationType,
|
||||||
|
blob: &UserBlob,
|
||||||
|
) -> CommunicationValue {
|
||||||
|
CommunicationValue::new(ty)
|
||||||
|
.with_request_id(cv)
|
||||||
|
.with_receiver(sender_wire_id(blob.user_id))
|
||||||
|
.add_typed_default(DataType::BlobId, DataValue::Str(blob.blob_id.clone()))
|
||||||
|
.add_typed_default(
|
||||||
|
DataType::VersionNumber,
|
||||||
|
DataValue::SignedNumber(blob.revision.into()),
|
||||||
|
)
|
||||||
|
.add_typed_default(
|
||||||
|
DataType::UpdatedAt,
|
||||||
|
DataValue::SignedNumber(blob.updated_at.into()),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
fn blob_changed(user_id: i64, blob_id: String, revision: i64, deleted: bool) -> CommunicationValue {
|
fn blob_changed(user_id: i64, blob_id: String, revision: i64, deleted: bool) -> CommunicationValue {
|
||||||
CommunicationValue::new(CommunicationType::UserBlobChanged)
|
CommunicationValue::new(CommunicationType::UserBlobChanged)
|
||||||
.with_id(next_notification_id())
|
.with_id(next_notification_id())
|
||||||
|
|
@ -2152,7 +2135,7 @@ pub fn handle_user_blob_put(cv: &CommunicationValue) -> BlobMutation {
|
||||||
let expected_revision = data_i64(cv, DataType::ExpectedRevision);
|
let expected_revision = data_i64(cv, DataType::ExpectedRevision);
|
||||||
match user_blobs::put(user_id, &blob_id, &blob, expected_revision) {
|
match user_blobs::put(user_id, &blob_id, &blob, expected_revision) {
|
||||||
Ok(stored) => BlobMutation {
|
Ok(stored) => BlobMutation {
|
||||||
response: blob_response(cv, CommunicationType::UserBlobPut, &stored),
|
response: blob_mutation_response(cv, CommunicationType::UserBlobPut, &stored),
|
||||||
changed: Some(blob_changed(
|
changed: Some(blob_changed(
|
||||||
user_id,
|
user_id,
|
||||||
stored.blob_id.clone(),
|
stored.blob_id.clone(),
|
||||||
|
|
@ -2227,13 +2210,13 @@ 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(user_id) {
|
match user_blobs::list_metadata(user_id) {
|
||||||
Ok(blobs) => CommunicationValue::new(CommunicationType::UserBlobList)
|
Ok(blobs) => 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(
|
.add_typed_default(
|
||||||
DataType::Blobs,
|
DataType::Blobs,
|
||||||
DataValue::Array(blobs.iter().map(blob_value).collect()),
|
DataValue::Array(blobs.iter().map(blob_metadata_value).collect()),
|
||||||
),
|
),
|
||||||
Err(_) => error_response(cv, CommunicationType::ErrorInternal),
|
Err(_) => error_response(cv, CommunicationType::ErrorInternal),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
/* Opaque client-owned data is stored without interpreting its encrypted bytes. */
|
/* 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::{db, sync};
|
||||||
use rusqlite::{Connection, OptionalExtension, Row, Transaction, params};
|
use rusqlite::{Connection, OptionalExtension, Row, Transaction, params};
|
||||||
|
|
@ -17,6 +17,13 @@ pub struct UserBlob {
|
||||||
pub updated_at: i64,
|
pub updated_at: i64,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct UserBlobMetadata {
|
||||||
|
pub blob_id: String,
|
||||||
|
pub revision: i64,
|
||||||
|
pub updated_at: i64,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct DeletedUserBlob {
|
pub struct DeletedUserBlob {
|
||||||
pub id: i64,
|
pub id: i64,
|
||||||
|
|
@ -51,6 +58,14 @@ fn from_row(row: &Row<'_>) -> rusqlite::Result<UserBlob> {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn metadata_from_row(row: &Row<'_>) -> rusqlite::Result<UserBlobMetadata> {
|
||||||
|
Ok(UserBlobMetadata {
|
||||||
|
blob_id: row.get(0)?,
|
||||||
|
revision: row.get(1)?,
|
||||||
|
updated_at: row.get(2)?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
fn load(tx: &Transaction<'_>, id: i64) -> Result<UserBlob, StorageError> {
|
fn load(tx: &Transaction<'_>, id: i64) -> Result<UserBlob, StorageError> {
|
||||||
Ok(tx.query_row("SELECT id, user_id, blob_id, blob, revision, updated_at FROM user_blobs WHERE id = ?1 AND deleted = 0", [id], from_row)?)
|
Ok(tx.query_row("SELECT id, user_id, blob_id, blob, revision, updated_at FROM user_blobs WHERE id = ?1 AND deleted = 0", [id], from_row)?)
|
||||||
}
|
}
|
||||||
|
|
@ -143,6 +158,20 @@ pub fn list(user_id: i64) -> Result<Vec<UserBlob>, StorageError> {
|
||||||
&[],
|
&[],
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn list_metadata(user_id: i64) -> Result<Vec<UserBlobMetadata>, StorageError> {
|
||||||
|
if user_id <= 0 {
|
||||||
|
return Err(StorageError::Other("invalid blob owner".into()));
|
||||||
|
}
|
||||||
|
db::with_db(|conn| {
|
||||||
|
let mut statement = conn.prepare(
|
||||||
|
"SELECT blob_id, revision, updated_at FROM user_blobs WHERE user_id = ?1 AND deleted = 0",
|
||||||
|
)?;
|
||||||
|
let rows = statement.query_map([user_id], metadata_from_row)?;
|
||||||
|
rows.collect::<Result<Vec<_>, _>>()
|
||||||
|
.map_err(StorageError::from)
|
||||||
|
})
|
||||||
|
}
|
||||||
pub fn list_by_ids(user_id: i64, ids: &[i64]) -> Result<Vec<UserBlob>, StorageError> {
|
pub fn list_by_ids(user_id: i64, ids: &[i64]) -> Result<Vec<UserBlob>, 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()));
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue