[Fix] Propper message storage and loading
[Todo] Respond to Message Send
This commit is contained in:
parent
c6a59abd25
commit
19c8c14fb8
2 changed files with 95 additions and 91 deletions
|
|
@ -518,12 +518,28 @@ impl OmikronConnection {
|
|||
height,
|
||||
);
|
||||
|
||||
// persist message for the sender (storage_owner = sender_id)
|
||||
chat_files::add_message(
|
||||
timestamp_u128,
|
||||
true,
|
||||
sender_id as i64,
|
||||
receiver_id as i64,
|
||||
&content,
|
||||
height,
|
||||
);
|
||||
|
||||
// send confirmation back to sender
|
||||
let conf_msg = CommunicationValue::new(CommunicationType::message_send)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(sender_id as u64);
|
||||
self.send_message(&conf_msg).await;
|
||||
|
||||
// Build a live-delivery message for the local client (recipient)
|
||||
let user_forward = CommunicationValue::new(CommunicationType::message_live)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(receiver_id as u64)
|
||||
.add_data(DataTypes::send_time, DataValue::Number(timestamp_i64))
|
||||
.add_data(DataTypes::message, DataValue::Str(content.clone()))
|
||||
.add_data(DataTypes::content, DataValue::Str(content.clone()))
|
||||
.add_data(DataTypes::sender_id, DataValue::Number(sender_id as i64))
|
||||
.add_data(DataTypes::height, DataValue::Number(height));
|
||||
|
||||
|
|
@ -540,7 +556,7 @@ impl OmikronConnection {
|
|||
.unwrap_or_else(|| "".to_string());
|
||||
let ms = MessageState::from_str(&ms_raw).upgrade(MessageState::Received);
|
||||
|
||||
// update stored message state
|
||||
// update stored message state for receiver
|
||||
let _ = chat_files::change_message_state(
|
||||
timestamp_i64,
|
||||
receiver_id as i64,
|
||||
|
|
@ -548,6 +564,14 @@ impl OmikronConnection {
|
|||
ms.clone(),
|
||||
);
|
||||
|
||||
// update stored message state for sender
|
||||
let _ = chat_files::change_message_state(
|
||||
timestamp_i64,
|
||||
sender_id as i64,
|
||||
receiver_id as i64,
|
||||
ms.clone(),
|
||||
);
|
||||
|
||||
// notify original sender about the delivered/read state
|
||||
self.send_message(
|
||||
&CommunicationValue::new(CommunicationType::message_state)
|
||||
|
|
@ -562,7 +586,7 @@ impl OmikronConnection {
|
|||
)
|
||||
.await;
|
||||
} else {
|
||||
// Delivery failed or timed out; mark as Sent and notify sender
|
||||
// Delivery failed or timed out; mark as Sent
|
||||
let _ = chat_files::change_message_state(
|
||||
timestamp_i64,
|
||||
receiver_id as i64,
|
||||
|
|
@ -570,6 +594,14 @@ impl OmikronConnection {
|
|||
MessageState::Sent,
|
||||
);
|
||||
|
||||
let _ = chat_files::change_message_state(
|
||||
timestamp_i64,
|
||||
sender_id as i64,
|
||||
receiver_id as i64,
|
||||
MessageState::Sent,
|
||||
);
|
||||
|
||||
// notify sender
|
||||
self.send_message(
|
||||
&CommunicationValue::new(CommunicationType::message_state)
|
||||
.with_id(cv.get_id())
|
||||
|
|
@ -627,7 +659,7 @@ impl OmikronConnection {
|
|||
.with_id(cv.get_id())
|
||||
.with_receiver(*receiver_id)
|
||||
.add_data(DataTypes::send_time, DataValue::Number(timestamp))
|
||||
.add_data(DataTypes::message, DataValue::Str(content.clone()))
|
||||
.add_data(DataTypes::content, DataValue::Str(content.clone()))
|
||||
.add_data(DataTypes::sender_id, DataValue::Number(*sender_id as i64))
|
||||
.add_data(DataTypes::height, DataValue::Number(height));
|
||||
|
||||
|
|
@ -687,37 +719,21 @@ impl OmikronConnection {
|
|||
return;
|
||||
}
|
||||
|
||||
// Duplicate handling for CommunicationType::message_send removed.
|
||||
// Rationale: This branch duplicated logic present earlier that handles incoming
|
||||
// stored messages and live delivery to local clients. Keeping a single,
|
||||
// well-defined code path for `message_send` reduces ambiguity and avoids
|
||||
// accidental early returns that block other handlers. If the protocol needs
|
||||
// distinct handling for client-originated sends vs stored deliveries, prefer
|
||||
// using distinct CommunicationType variants or an explicit field/flag.
|
||||
|
||||
if cv.is_type(CommunicationType::messages_get) {
|
||||
let my_id = cv.get_sender();
|
||||
let partner_id = cv.get_data(DataTypes::user_id).as_number().unwrap_or(0);
|
||||
let offset = cv.get_data(DataTypes::offset).as_number().unwrap_or(0);
|
||||
let amount = cv.get_data(DataTypes::amount).as_number().unwrap_or(0);
|
||||
// retrieve raw JSON messages
|
||||
let messages = chat_files::get_messages(my_id as i64, partner_id, offset, amount);
|
||||
// convert JSON array -> protocol Array of Containers (send_time, content, sender_id, message_state, height)
|
||||
let mut msg_array: Vec<DataValue> = Vec::new();
|
||||
for m in messages.members() {
|
||||
// extract fields defensively
|
||||
let message_time: i64 = m["message_time"].as_i64().unwrap_or(0);
|
||||
let content: String = m["content"].as_str().unwrap_or("").to_string();
|
||||
let sent_by_self: bool = m["sent_by_self"].as_bool().unwrap_or(false);
|
||||
let height: i64 = m["height"].as_i64().unwrap_or(0);
|
||||
// determine sender id:
|
||||
// - if sent_by_self => sender is the requester (my_id)
|
||||
// - otherwise prefer an explicit chat_partner_id if present on the request,
|
||||
// fallback to the partner_id parameter
|
||||
let sender_id: i64 = if sent_by_self {
|
||||
my_id as i64
|
||||
} else {
|
||||
// check for chat_partner_id in the incoming request (accept number or string)
|
||||
if let Some(n) = cv.get_data(DataTypes::chat_partner_id).as_number() {
|
||||
n as i64
|
||||
} else if let Some(s) = cv.get_data(DataTypes::chat_partner_id).as_str() {
|
||||
|
|
@ -730,14 +746,11 @@ impl OmikronConnection {
|
|||
|
||||
let mut container = Vec::new();
|
||||
container.push((DataTypes::send_time, DataValue::Number(message_time)));
|
||||
container.push((DataTypes::message, DataValue::Str(content)));
|
||||
container.push((DataTypes::content, DataValue::Str(content)));
|
||||
container.push((DataTypes::sender_id, DataValue::Number(sender_id)));
|
||||
container.push((DataTypes::message_state, DataValue::Str(message_state)));
|
||||
container.push((DataTypes::height, DataValue::Number(height)));
|
||||
container.push((
|
||||
DataTypes::parse("sent_by_self".to_string()),
|
||||
DataValue::Bool(sent_by_self),
|
||||
));
|
||||
container.push((DataTypes::sent_by_self, DataValue::Bool(sent_by_self)));
|
||||
msg_array.push(DataValue::Container(container));
|
||||
}
|
||||
|
||||
|
|
@ -1022,9 +1035,6 @@ impl OmikronConnection {
|
|||
let response = CommunicationValue::new(CommunicationType::error)
|
||||
.with_id(key)
|
||||
.add_data(DataTypes::message, DataValue::Str(reason.clone()));
|
||||
// Historically this used the global `OMIKRON_CONNECTION`. Using the global here
|
||||
// preserves the original behavior and avoids ownership/borrow issues when
|
||||
// invoking the waiting-task closures from a &self context.
|
||||
let _ = (waiting_task.task)(OMIKRON_CONNECTION.clone(), response);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue