[Add] timestamping
This commit is contained in:
parent
29d46d1c95
commit
3bfec96848
4 changed files with 133 additions and 46 deletions
|
|
@ -185,6 +185,9 @@ pub fn apply_verified_relay_content(
|
|||
&context.type_map,
|
||||
)
|
||||
.ok_or_else(|| "Relay MessageSend is missing RelayMessageId".to_string())?;
|
||||
if relay_message_id != context.message_id {
|
||||
return Err("Relay MessageSend identity does not match its protected message ID".into());
|
||||
}
|
||||
chat_files::add_message(chat_files::NewMessage {
|
||||
relay_signer_id: sender_id,
|
||||
relay_message_id,
|
||||
|
|
|
|||
|
|
@ -81,13 +81,12 @@ pub fn mark_delivered_for_frame(destination_id: u64, frame_id: u32) -> Result<()
|
|||
})
|
||||
}
|
||||
|
||||
pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), StorageError> {
|
||||
if !matches!(
|
||||
state,
|
||||
"received" | "applied" | "queued" | "delivered" | "rejected"
|
||||
) {
|
||||
return Err(StorageError::Other("invalid relay inbox state".into()));
|
||||
}
|
||||
fn mark_transition(
|
||||
signer_id: u64,
|
||||
message_id: &str,
|
||||
state: &str,
|
||||
column: &str,
|
||||
) -> Result<(), StorageError> {
|
||||
let signer_id = i64::try_from(signer_id)
|
||||
.map_err(|_| StorageError::Other("relay signer ID exceeds SQLite range".into()))?;
|
||||
let timestamp = std::time::SystemTime::now()
|
||||
|
|
@ -95,14 +94,6 @@ pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), S
|
|||
.unwrap_or_default()
|
||||
.as_millis() as i64;
|
||||
db::with_db(|connection| {
|
||||
let column = match state {
|
||||
"applied" => "applied_at",
|
||||
"queued" => "queued_at",
|
||||
"delivered" => "downstream_acked_at",
|
||||
"rejected" => "rejected_at",
|
||||
"received" => "accepted_at",
|
||||
_ => return Err(StorageError::Other("invalid relay inbox state".into())),
|
||||
};
|
||||
connection.execute(
|
||||
&format!("UPDATE relay_inbox SET state = ?3, {column} = COALESCE({column}, ?4) WHERE signer_id = ?1 AND message_id = ?2"),
|
||||
params![signer_id, message_id, state, timestamp],
|
||||
|
|
@ -111,6 +102,22 @@ pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), S
|
|||
})
|
||||
}
|
||||
|
||||
pub fn mark_applied(signer_id: u64, message_id: &str) -> Result<(), StorageError> {
|
||||
mark_transition(signer_id, message_id, "applied", "applied_at")
|
||||
}
|
||||
|
||||
pub fn mark_queued(signer_id: u64, message_id: &str) -> Result<(), StorageError> {
|
||||
mark_transition(signer_id, message_id, "queued", "queued_at")
|
||||
}
|
||||
|
||||
pub fn mark_downstream_acked(signer_id: u64, message_id: &str) -> Result<(), StorageError> {
|
||||
mark_transition(signer_id, message_id, "delivered", "downstream_acked_at")
|
||||
}
|
||||
|
||||
pub fn mark_rejected(signer_id: u64, message_id: &str) -> Result<(), StorageError> {
|
||||
mark_transition(signer_id, message_id, "rejected", "rejected_at")
|
||||
}
|
||||
|
||||
pub fn prune_completed(before_terminal_at: i64) -> Result<(), StorageError> {
|
||||
db::with_db(|connection| {
|
||||
connection.execute(
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ use crate::util::db;
|
|||
use rusqlite::{Transaction, params};
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
pub const CACHE_SCHEMA_VERSION: i64 = 1;
|
||||
pub const CACHE_SCHEMA_VERSION: i64 = 2;
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum EntityType {
|
||||
|
|
|
|||
|
|
@ -893,6 +893,7 @@ impl OmikronConnection {
|
|||
iota_id: u64,
|
||||
relay_message_id: &str,
|
||||
accepted_at: i64,
|
||||
include_origin_timestamp: bool,
|
||||
) {
|
||||
let Some(frame_id) = frame_id else { return };
|
||||
let response = CommunicationValue::new(CommunicationType::Success)
|
||||
|
|
@ -906,6 +907,19 @@ impl OmikronConnection {
|
|||
DataType::RelayAcceptedAt,
|
||||
DataValue::SignedNumber(accepted_at.into()),
|
||||
);
|
||||
let response = if include_origin_timestamp {
|
||||
response
|
||||
.add_typed_default(
|
||||
DataType::OriginIotaReceivedAt,
|
||||
DataValue::SignedNumber(accepted_at.into()),
|
||||
)
|
||||
.add_typed_default(
|
||||
DataType::DestinationIotaReceivedAt,
|
||||
DataValue::SignedNumber(accepted_at.into()),
|
||||
)
|
||||
} else {
|
||||
response
|
||||
};
|
||||
if let Err(error) = self.send_message(&response).await {
|
||||
log!("Relay response could not be sent: {}", error);
|
||||
}
|
||||
|
|
@ -1039,6 +1053,63 @@ impl OmikronConnection {
|
|||
relay_replay::RelayReservation::Existing { .. } => false,
|
||||
};
|
||||
|
||||
/* A shared Iota owns both independent replicas before delivering to its
|
||||
* local recipient. The destination path below writes the recipient copy. */
|
||||
if signer_is_local && recipient_is_local && !already_applied {
|
||||
let content = match open_verified_relay_content(
|
||||
&verified,
|
||||
&[&keyring],
|
||||
verified.context.signer_id,
|
||||
) {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
log!("Relay shared-Iota origin content verification failed: {}", error);
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
let owner = match i64::try_from(verified.context.signer_id) {
|
||||
Ok(value) => value,
|
||||
Err(_) => {
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
if let Err(error) = message_handlers::apply_verified_relay_content(
|
||||
&verified.context,
|
||||
&content,
|
||||
accepted_at,
|
||||
owner,
|
||||
true,
|
||||
) {
|
||||
log!("Relay shared-Iota origin application failed: {}", error);
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
if let Err(error) = chat_files::record_destination_iota_received(
|
||||
owner,
|
||||
owner,
|
||||
&verified.context.message_id,
|
||||
accepted_at,
|
||||
) {
|
||||
log!("Relay shared-Iota destination timestamp storage failed: {}", error);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInternal)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
if signer_is_local && !recipient_is_local {
|
||||
if !already_applied {
|
||||
let content = match open_verified_relay_content(
|
||||
|
|
@ -1049,10 +1120,9 @@ impl OmikronConnection {
|
|||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
log!("Relay origin content verification failed: {}", error);
|
||||
let _ = relay_replay::mark_state(
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"rejected",
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
|
|
@ -1067,10 +1137,9 @@ impl OmikronConnection {
|
|||
true,
|
||||
) {
|
||||
log!("Relay origin application failed: {}", error);
|
||||
let _ = relay_replay::mark_state(
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"rejected",
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
|
|
@ -1119,10 +1188,9 @@ impl OmikronConnection {
|
|||
.await;
|
||||
return;
|
||||
}
|
||||
if let Err(error) = relay_replay::mark_state(
|
||||
if let Err(error) = relay_replay::mark_queued(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"queued",
|
||||
) {
|
||||
log!("Relay origin state update failed: {}", error);
|
||||
}
|
||||
|
|
@ -1134,7 +1202,7 @@ impl OmikronConnection {
|
|||
let returned_id = response
|
||||
.get_data(DataType::RelayMessageId)
|
||||
.as_str();
|
||||
let accepted_at = response
|
||||
let destination_accepted_at = response
|
||||
.get_data(DataType::RelayAcceptedAt)
|
||||
.as_number()
|
||||
.and_then(|value| i64::try_from(value).ok());
|
||||
|
|
@ -1144,7 +1212,12 @@ impl OmikronConnection {
|
|||
.await;
|
||||
return;
|
||||
}
|
||||
if let Some(accepted_at) = accepted_at {
|
||||
let Some(destination_accepted_at) = destination_accepted_at else {
|
||||
log!("Relay acknowledgement is missing RelayAcceptedAt");
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
return;
|
||||
};
|
||||
if let (Ok(owner), Ok(signer)) = (
|
||||
i64::try_from(verified.context.signer_id),
|
||||
i64::try_from(verified.context.signer_id),
|
||||
|
|
@ -1153,26 +1226,33 @@ impl OmikronConnection {
|
|||
owner,
|
||||
signer,
|
||||
&verified.context.message_id,
|
||||
accepted_at,
|
||||
destination_accepted_at,
|
||||
) {
|
||||
log!("Relay destination acknowledgement storage failed: {}", error);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Err(error) = relay_queue::acknowledge_iota(router, frame_id) {
|
||||
log!(
|
||||
"Relay origin acknowledgement could not clear the queue: {}",
|
||||
error
|
||||
);
|
||||
}
|
||||
if let Err(error) = relay_replay::mark_state(
|
||||
if let Err(error) = relay_replay::mark_downstream_acked(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"delivered",
|
||||
) {
|
||||
log!("Relay origin delivery state update failed: {}", error);
|
||||
}
|
||||
let response = response.with_id(frame_id);
|
||||
let response = response
|
||||
.add_typed_default(
|
||||
DataType::OriginIotaReceivedAt,
|
||||
DataValue::SignedNumber(accepted_at.into()),
|
||||
)
|
||||
.add_typed_default(
|
||||
DataType::DestinationIotaReceivedAt,
|
||||
DataValue::SignedNumber(destination_accepted_at.into()),
|
||||
)
|
||||
.with_id(frame_id);
|
||||
if let Err(error) = self.send_message(&response).await {
|
||||
log!("Relay response could not be sent: {}", error);
|
||||
}
|
||||
|
|
@ -1248,10 +1328,9 @@ impl OmikronConnection {
|
|||
queue_error
|
||||
);
|
||||
}
|
||||
let _ = relay_replay::mark_state(
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"rejected",
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
|
|
@ -1284,19 +1363,17 @@ impl OmikronConnection {
|
|||
{
|
||||
log!("Relay application queue cleanup failed: {}", queue_error);
|
||||
}
|
||||
let _ = relay_replay::mark_state(
|
||||
let _ = relay_replay::mark_rejected(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"rejected",
|
||||
);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
if let Err(error) = relay_replay::mark_state(
|
||||
if let Err(error) = relay_replay::mark_applied(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"applied",
|
||||
) {
|
||||
log!("Relay application state update failed: {}", error);
|
||||
self.send_relay_response(frame.id(), CommunicationType::ErrorInternal)
|
||||
|
|
@ -1305,10 +1382,9 @@ impl OmikronConnection {
|
|||
}
|
||||
}
|
||||
|
||||
if let Err(error) = relay_replay::mark_state(
|
||||
if let Err(error) = relay_replay::mark_queued(
|
||||
verified.context.signer_id,
|
||||
&verified.context.message_id,
|
||||
"queued",
|
||||
) {
|
||||
log!("Relay queue state update failed: {}", error);
|
||||
}
|
||||
|
|
@ -1317,6 +1393,7 @@ impl OmikronConnection {
|
|||
local_iota_id,
|
||||
&verified.context.message_id,
|
||||
accepted_at,
|
||||
signer_is_local,
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = self.send_message(&forwarded).await {
|
||||
|
|
|
|||
Loading…
Reference in a new issue