[Add] Mesage util
This commit is contained in:
parent
577660884c
commit
2b578693aa
5 changed files with 671 additions and 293 deletions
|
|
@ -14,6 +14,12 @@ pub enum MessageState {
|
|||
Sending,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct MessageStateUpdate {
|
||||
pub send_time: i64,
|
||||
pub state: MessageState,
|
||||
}
|
||||
|
||||
impl MessageState {
|
||||
pub fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
|
|
@ -979,23 +985,29 @@ pub fn change_message_state_by_relay_id(
|
|||
relay_signer_principal: iota_identity::PrincipalHandle,
|
||||
relay_message_id: &str,
|
||||
new_state: MessageState,
|
||||
) -> Result<(), StorageError> {
|
||||
) -> Result<Option<MessageStateUpdate>, StorageError> {
|
||||
db::with_db(|conn| {
|
||||
let tx = conn.unchecked_transaction()?;
|
||||
let Some((msg_id, current)) = tx
|
||||
let Some((msg_id, current, message_time)) = tx
|
||||
.query_row(
|
||||
"SELECT id, message_state FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3",
|
||||
"SELECT id, message_state, message_time FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3",
|
||||
params![storage_owner, relay_signer_principal.0, relay_message_id],
|
||||
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
|
||||
|row| {
|
||||
Ok((
|
||||
row.get::<_, i64>(0)?,
|
||||
row.get::<_, String>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
.optional()?
|
||||
else {
|
||||
return Ok(());
|
||||
return Ok(None);
|
||||
};
|
||||
let state = MessageState::from_str(¤t).upgrade(new_state).as_str();
|
||||
let final_state = MessageState::from_str(¤t).upgrade(new_state);
|
||||
tx.execute(
|
||||
"UPDATE messages SET message_state = ?1 WHERE id = ?2",
|
||||
params![state, msg_id],
|
||||
params![final_state.as_str(), msg_id],
|
||||
)?;
|
||||
sync::record_event(
|
||||
&tx,
|
||||
|
|
@ -1005,7 +1017,10 @@ pub fn change_message_state_by_relay_id(
|
|||
Operation::Upsert,
|
||||
)?;
|
||||
tx.commit()?;
|
||||
Ok(())
|
||||
Ok(Some(MessageStateUpdate {
|
||||
send_time: message_time,
|
||||
state: final_state,
|
||||
}))
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -1020,7 +1035,7 @@ pub fn record_message_receipt(
|
|||
receipt_type: MessageState,
|
||||
event_at: i64,
|
||||
recorded_at: i64,
|
||||
) -> Result<(), StorageError> {
|
||||
) -> Result<Option<MessageStateUpdate>, StorageError> {
|
||||
let receipt_type = match receipt_type {
|
||||
MessageState::Received => "received",
|
||||
MessageState::Read => "read",
|
||||
|
|
@ -1028,11 +1043,19 @@ pub fn record_message_receipt(
|
|||
};
|
||||
db::with_db(|conn| {
|
||||
let tx = conn.unchecked_transaction()?;
|
||||
let Some((message_id, external_principal, authored_at, history_deleted)) = tx
|
||||
let Some((message_id, message_time, external_principal, authored_at, history_deleted)) = tx
|
||||
.query_row(
|
||||
"SELECT id, external_principal, authored_at, history_deleted FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3",
|
||||
"SELECT id, message_time, external_principal, authored_at, history_deleted FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3",
|
||||
params![storage_owner, target_signer_principal.0, target_message_id],
|
||||
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?, row.get::<_, Option<i64>>(2)?, row.get::<_, i64>(3)?)),
|
||||
|row| {
|
||||
Ok((
|
||||
row.get::<_, i64>(0)?,
|
||||
row.get::<_, i64>(1)?,
|
||||
row.get::<_, i64>(2)?,
|
||||
row.get::<_, Option<i64>>(3)?,
|
||||
row.get::<_, i64>(4)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
.optional()?
|
||||
else {
|
||||
|
|
@ -1051,7 +1074,7 @@ pub fn record_message_receipt(
|
|||
));
|
||||
}
|
||||
if history_deleted != 0 {
|
||||
return Ok(());
|
||||
return Ok(None);
|
||||
}
|
||||
tx.execute(
|
||||
"INSERT OR IGNORE INTO message_receipts (storage_owner, target_signer_id, target_signer_principal, target_message_id, receipt_signer_id, receipt_signer_principal, receipt_message_id, receipt_type, event_at, recorded_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
||||
|
|
@ -1062,17 +1085,15 @@ pub fn record_message_receipt(
|
|||
} else {
|
||||
("client_received_at", "client_received_recorded_at")
|
||||
};
|
||||
let state = MessageState::from_str(&tx.query_row(
|
||||
let final_state = MessageState::from_str(&tx.query_row(
|
||||
"SELECT message_state FROM messages WHERE id = ?1",
|
||||
[message_id],
|
||||
|row| row.get::<_, String>(0),
|
||||
)?)
|
||||
.upgrade(MessageState::from_str(receipt_type))
|
||||
.as_str()
|
||||
.to_string();
|
||||
.upgrade(MessageState::from_str(receipt_type));
|
||||
tx.execute(
|
||||
&format!("UPDATE messages SET {state_column} = COALESCE({state_column}, ?1), {recorded_column} = COALESCE({recorded_column}, ?2), message_state = ?3 WHERE id = ?4"),
|
||||
params![event_at, recorded_at, state, message_id],
|
||||
params![event_at, recorded_at, final_state.as_str(), message_id],
|
||||
)?;
|
||||
sync::record_event(
|
||||
&tx,
|
||||
|
|
@ -1082,7 +1103,10 @@ pub fn record_message_receipt(
|
|||
Operation::Upsert,
|
||||
)?;
|
||||
tx.commit()?;
|
||||
Ok(())
|
||||
Ok(Some(MessageStateUpdate {
|
||||
send_time: message_time,
|
||||
state: final_state,
|
||||
}))
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -1579,6 +1603,10 @@ mod tests {
|
|||
MessageState::Received.upgrade(MessageState::Read),
|
||||
MessageState::Read
|
||||
);
|
||||
assert_eq!(
|
||||
MessageState::Read.upgrade(MessageState::Received),
|
||||
MessageState::Read
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
Loading…
Reference in a new issue