use map_core::client::MapClient;
use map_core::folders::Folder;
use map_core::MessageStatus;
use store::{Direction, NewMessage, OutgoingStatus, PhoneField, Store};
use tokio::io::{AsyncRead, AsyncWrite};
use tokio::sync::{mpsc, watch};
use crate::outbox::drain_outbox;
use crate::sync::backfill_catch_up;
use crate::util::now_ms;
use crate::{EventType, MnsEvent};
fn parse_folder(s: &str) -> Option<Folder> {
let leaf = s.rsplit('/').next().unwrap_or(s);
match leaf.to_ascii_lowercase().as_str() {
"inbox" => Some(Folder::Inbox),
"sent" => Some(Folder::Sent),
"outbox" => Some(Folder::Outbox),
"deleted" => Some(Folder::Deleted),
_ => None,
}
}
pub async fn handle_mns_event<T: AsyncRead + AsyncWrite + Unpin>(
event: &MnsEvent,
client: &mut MapClient<T>,
store: &Store,
now: i64,
) -> anyhow::Result<()> {
match event.event_type() {
EventType::NewMessage => handle_new_message(event, client, store, now).await?,
EventType::MessageDeleted => {
if let Some(handle) = event.handle() {
store.delete_by_handle(handle).await?;
}
}
EventType::MessageShift => {
if let (Some(handle), Some(folder)) = (event.handle(), event.folder()) {
store.update_folder(handle, folder).await?;
}
}
EventType::ReadStatusChanged => handle_read_status_changed(event, client, store).await?,
EventType::DeliverySuccess | EventType::SendingSuccess => {
mark_outgoing(event, store, OutgoingStatus::SentConfirmed).await?;
}
EventType::DeliveryFailure | EventType::SendingFailure => {
mark_outgoing(event, store, OutgoingStatus::FailedPermanent).await?;
}
EventType::MemoryFull => {
tracing::warn!(
"device message store is full — new messages may be rejected until freed"
);
}
EventType::MemoryAvailable => {
tracing::info!("device message store has space available again");
}
}
Ok(())
}
async fn handle_new_message<T: AsyncRead + AsyncWrite + Unpin>(
event: &MnsEvent,
client: &mut MapClient<T>,
store: &Store,
now: i64,
) -> anyhow::Result<()> {
let (Some(handle), Some(folder_raw)) = (event.handle(), event.folder()) else {
tracing::warn!("NewMessage event missing handle or folder — skipped");
return Ok(());
};
let Some(folder) = parse_folder(folder_raw) else {
tracing::warn!("NewMessage event unknown folder {folder_raw} — skipped");
return Ok(());
};
client.set_folder(folder).await?;
let bmsg = client.get_message(handle).await?;
let address = bmsg.originator().map(|o| o.tel.clone()).unwrap_or_default();
let status = i32::from(matches!(bmsg.status(), MessageStatus::Read));
let msg = NewMessage {
map_handle: handle.to_owned(),
timestamp_ms: now,
folder: folder_raw.to_owned(),
direction: Direction::Received,
address: PhoneField::new(&address, None),
status,
synced_at: now,
text: bmsg.envelope().body.text.clone(),
outgoing_status: None,
};
store.upsert(msg).await?;
Ok(())
}
async fn handle_read_status_changed<T: AsyncRead + AsyncWrite + Unpin>(
event: &MnsEvent,
client: &mut MapClient<T>,
store: &Store,
) -> anyhow::Result<()> {
let (Some(handle), Some(folder_raw)) = (event.handle(), event.folder()) else {
tracing::warn!("ReadStatusChanged event missing handle or folder — skipped");
return Ok(());
};
let Some(folder) = parse_folder(folder_raw) else {
tracing::warn!("ReadStatusChanged event unknown folder {folder_raw} — skipped");
return Ok(());
};
client.set_folder(folder).await?;
let bmsg = client.get_message(handle).await?;
let status = i32::from(matches!(bmsg.status(), MessageStatus::Read));
store.update_status(handle, status).await?;
Ok(())
}
async fn mark_outgoing(
event: &MnsEvent,
store: &Store,
status: OutgoingStatus,
) -> anyhow::Result<()> {
if let Some(handle) = event.handle() {
store.update_outgoing_status(handle, status).await?;
}
Ok(())
}
pub async fn run_watch<T: AsyncRead + AsyncWrite + Unpin>(
event_rx: &mut mpsc::Receiver<MnsEvent>,
client: &mut MapClient<T>,
store: &Store,
mut cancel_rx: watch::Receiver<bool>,
) -> anyhow::Result<()> {
backfill_catch_up(client, store).await?;
drain_outbox(client, store, now_ms()).await?;
loop {
tokio::select! {
biased;
_ = cancel_rx.changed() => {
if *cancel_rx.borrow() { break; }
}
event = event_rx.recv() => {
let Some(ev) = event else { break; };
handle_mns_event(&ev, client, store, now_ms()).await?;
}
}
}
Ok(())
}
#[cfg(test)]
mod tests;