Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE chats DROP COLUMN archive_appstate_seen;
ALTER TABLE chats DROP COLUMN mute_appstate_seen;
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
-- A live message may create a chat before HistorySync supplies its preference
-- snapshot. Keep that snapshot renewable until an explicit app-state action
-- has spoken for each preference; its false/NULL answer is authoritative too.
ALTER TABLE chats ADD COLUMN mute_appstate_seen BOOLEAN NOT NULL DEFAULT FALSE;
ALTER TABLE chats ADD COLUMN archive_appstate_seen BOOLEAN NOT NULL DEFAULT FALSE;
4 changes: 2 additions & 2 deletions crates/chat-store/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,6 @@ pub use materialize::is_control_only;
pub use store::ChatStore;
pub use types::{
ArrivalCursor, AvatarDescriptor, ChatCursor, ChatEntry, ChatNameExpected, ChatNameWrite,
ContactEntry, MediaRef, MessageCoverage, MessageCursor, MessageKind, MessageStatus,
ReactionEntry, ReceiptEntry, StoreChange, StoredMessage,
ChatNotificationMetadata, ContactEntry, MediaRef, MessageCoverage, MessageCursor, MessageKind,
MessageStatus, ReactionEntry, ReceiptEntry, StoreChange, StoredMessage,
};
64 changes: 60 additions & 4 deletions crates/chat-store/src/queries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,15 +8,15 @@ use std::str::FromStr;
use chrono::{DateTime, Utc};
use diesel::prelude::*;
use log::warn;
use wacore_binary::Jid;
use wacore_binary::jid::{Jid, JidExt};

use crate::error::{ChatStoreError, Result, db_err};
use crate::schema;
use crate::store::ChatStore;
use crate::types::{
ArrivalCursor, AvatarDescriptor, ChatCursor, ChatEntry, ContactEntry, MediaRef,
MessageCoverage, MessageCursor, MessageKind, MessageStatus, ReactionEntry, ReceiptEntry,
StoredMessage,
ArrivalCursor, AvatarDescriptor, ChatCursor, ChatEntry, ChatNotificationMetadata, ContactEntry,
MediaRef, MessageCoverage, MessageCursor, MessageKind, MessageStatus, ReactionEntry,
ReceiptEntry, StoredMessage,
};

/// How many keys one batched lookup may bind at a time.
Expand Down Expand Up @@ -165,6 +165,10 @@ struct ChatRow {
read_boundary_ms: i64,
#[allow(dead_code)]
read_boundary_ids: Option<String>,
#[allow(dead_code)]
mute_appstate_seen: bool,
#[allow(dead_code)]
archive_appstate_seen: bool,
}

impl From<ChatRow> for ChatEntry {
Expand Down Expand Up @@ -442,6 +446,58 @@ impl ChatStore {
Ok(row.map(Into::into))
}

/// Alert policy and title from the durable conversation rows.
///
/// A live message may precede GUI hydration, so the front end's `Chat`
/// cannot answer whether the phone muted or archived the conversation.
/// Unlike `chat()`, this reads *both* sides of a split PN/LID pair: a
/// stale alias with an active mute must not be bypassed just because the
/// other side has the newer message. No row means unknown, not allowed.
pub async fn notification_metadata(
&self,
jid: &Jid,
) -> Result<Option<ChatNotificationMetadata>> {
use schema::chats::dsl;
let device_id = self.device_id();
let is_group = jid.is_group();
let jid = jid.to_string();
let now_ms = wacore::time::now_utc().timestamp_millis();
let rows: Vec<(Option<i64>, bool, Option<String>)> = self
.db()
.read(move |conn| {
let keys =
crate::lid::chat_key_candidates(conn, device_id, &jid).map_err(db_err)?;
dsl::chats
.filter(dsl::device_id.eq(device_id).and(dsl::jid.eq_any(keys)))
.order((dsl::last_message_ts.desc(), dsl::jid.desc()))
.select((dsl::muted_until, dsl::archived, dsl::name))
.load(conn)
.map_err(db_err)
})
.await?;
if rows.is_empty() {
return Ok(None);
}
let muted = rows
.iter()
.any(|(mute, _, _)| mute.is_some_and(|until| until > now_ms));
let archived = rows.iter().any(|(_, archived, _)| *archived);
let name = rows
.into_iter()
.filter_map(|(_, _, name)| name)
.find(|name| {
!name.trim().is_empty()
&& !(is_group
&& matches!(name.trim(), "Unnamed group" | "Group name unavailable"))
});
Ok(Some(ChatNotificationMetadata {
muted,
archived,
allowed: !muted && !archived,
name,
}))
}

/// Every special chat's identity columns, in one read.
///
/// The chat-name resolver's full pass works from this, not from the
Expand Down
2 changes: 2 additions & 0 deletions crates/chat-store/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ diesel::table! {
ephemeral_expiration -> Nullable<Integer>,
read_boundary_ms -> BigInt,
read_boundary_ids -> Nullable<Text>,
mute_appstate_seen -> Bool,
archive_appstate_seen -> Bool,
}
}

Expand Down
38 changes: 33 additions & 5 deletions crates/chat-store/src/store/chat_rows.rs
Original file line number Diff line number Diff line change
Expand Up @@ -220,8 +220,10 @@ pub(super) fn ensure_chat(
/// rows have already moved). Activity and preview re-derive from the merged
/// messages; the self-read state is the union of both sides so neither
/// side's covered messages re-badge; sticky user prefs (pin/mute/archive,
/// name, ephemeral) keep dest's value and fall back to src's. A manual-unread
/// marker on either side survives; otherwise the badge is recounted.
/// name, ephemeral) keep dest's value and fall back to src's. For mute and
/// archive, an explicit app-state answer (including unmute/unarchive) beats a
/// history-only value from the other alias. A manual-unread marker on either
/// side survives; otherwise the badge is recounted.
pub(crate) fn merge_chat_metadata(
conn: &mut SqliteConnection,
device_id: i32,
Expand All @@ -237,6 +239,8 @@ pub(crate) fn merge_chat_metadata(
bool,
Option<i32>,
Option<String>,
bool,
bool,
);
let prefs = |conn: &mut SqliteConnection, key: &str| -> QueryResult<Option<PrefRow>> {
chat_row(device_id, key)
Expand All @@ -248,6 +252,8 @@ pub(crate) fn merge_chat_metadata(
dsl::archived,
dsl::ephemeral_expiration,
dsl::name,
dsl::mute_appstate_seen,
dsl::archive_appstate_seen,
))
.first(conn)
.optional()
Expand All @@ -257,7 +263,8 @@ pub(crate) fn merge_chat_metadata(
};
let src_state = read_state(conn, device_id, src)?;
ensure_chat(conn, device_id, dest)?;
let dest_row = prefs(conn, dest)?.unwrap_or((0, 0, None, None, false, None, None));
let dest_row =
prefs(conn, dest)?.unwrap_or((0, 0, None, None, false, None, None, false, false));
let dest_state = read_state(conn, device_id, dest)?;

let mut merged = ReadState {
Expand All @@ -277,15 +284,36 @@ pub(crate) fn merge_chat_metadata(
} else {
count_unread(conn, device_id, dest, &merged)?
};
// A split pair may have received an app-state action on only one side.
// Preserve that explicit answer even when it is false/NULL; otherwise
// keep the established sticky fallback between history-only rows. If
// both sides have explicit answers, the surviving (more active) row wins
// as it did before this provenance was tracked.
let muted_until = if dest_row.7 {
dest_row.3
} else if src_row.7 {
src_row.3
} else {
dest_row.3.or(src_row.3)
};
let archived = if dest_row.8 {
dest_row.4
} else if src_row.8 {
src_row.4
} else {
dest_row.4 || src_row.4
};
diesel::update(chat_row(device_id, dest))
.set((
dsl::last_message_ts.eq(src_row.0.max(dest_row.0)),
dsl::unread_count.eq(unread),
dsl::pinned_at.eq(dest_row.2.or(src_row.2)),
dsl::muted_until.eq(dest_row.3.or(src_row.3)),
dsl::archived.eq(dest_row.4 || src_row.4),
dsl::muted_until.eq(muted_until),
dsl::archived.eq(archived),
dsl::ephemeral_expiration.eq(dest_row.5.or(src_row.5)),
dsl::name.eq(dest_row.6.or(src_row.6)),
dsl::mute_appstate_seen.eq(dest_row.7 || src_row.7),
dsl::archive_appstate_seen.eq(dest_row.8 || src_row.8),
dsl::read_boundary_ms.eq(merged.watermark_ms),
dsl::read_boundary_ids.eq(ids_json),
))
Expand Down
39 changes: 29 additions & 10 deletions crates/chat-store/src/store/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,11 @@ pub(super) fn apply_event(
Ok(())
}
Event::MuteUpdate(update) => {
let muted_until = if update.action.muted.unwrap_or(false) {
let Some(muted) = update.action.muted else {
// Missing is not an explicit unmute or app-state authority.
return Ok(());
};
let muted_until = if muted {
// Absent or non-positive (WA Web sends -1 for indefinite,
// this crate's own mute_chat() included) = muted forever.
Some(
Expand All @@ -183,30 +187,45 @@ pub(super) fn apply_event(
};
let chat = crate::lid::route_chat_key(conn, device_id, &update.jid.to_string(), cs)?;
ensure_chat(conn, device_id, &chat)?;
let stored: Option<i64> = chat_row(device_id, &chat)
.select(schema::chats::muted_until)
let (stored, seen): (Option<i64>, bool) = chat_row(device_id, &chat)
.select((
schema::chats::muted_until,
schema::chats::mute_appstate_seen,
))
.first(conn)?;
if stored == muted_until {
if stored == muted_until && seen {
return Ok(());
}
diesel::update(chat_row(device_id, &chat))
.set(schema::chats::muted_until.eq(muted_until))
.set((
schema::chats::muted_until.eq(muted_until),
schema::chats::mute_appstate_seen.eq(true),
))
.execute(conn)?;
cs.chats = true;
Ok(())
}
Event::ArchiveUpdate(update) => {
let Some(archived) = update.action.archived else {
// Missing is not an explicit unarchive.
return Ok(());
};
let chat = crate::lid::route_chat_key(conn, device_id, &update.jid.to_string(), cs)?;
ensure_chat(conn, device_id, &chat)?;
let archived = update.action.archived.unwrap_or(false);
let stored: bool = chat_row(device_id, &chat)
.select(schema::chats::archived)
let (stored, seen): (bool, bool) = chat_row(device_id, &chat)
.select((
schema::chats::archived,
schema::chats::archive_appstate_seen,
))
.first(conn)?;
if stored == archived {
if stored == archived && seen {
return Ok(());
}
diesel::update(chat_row(device_id, &chat))
.set(schema::chats::archived.eq(archived))
.set((
schema::chats::archived.eq(archived),
schema::chats::archive_appstate_seen.eq(true),
))
.execute(conn)?;
cs.chats = true;
Ok(())
Expand Down
24 changes: 19 additions & 5 deletions crates/chat-store/src/store/history_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,18 +110,32 @@ fn apply_history_conversation(
))
.on_conflict((dsl::device_id, dsl::jid))
.do_update()
// Live rows already track unread/mute/pin; history only refreshes
// identity + activity floor. A nameless chunk preserves an
// existing name rather than clobbering it with NULL; a named one
// still updates. Two statements rather than one, because the SET
// clause is static and the two cases write different columns.
// A live message can create the row before the phone's history
// snapshot has supplied its mute/archive preferences. History
// keeps refreshing those until an explicit app-state update has
// spoken for each one; then even false/NULL is newer authority
// than a delayed snapshot. Unread/pin remain owned by live state.
// A nameless chunk preserves an existing name rather than
// clobbering it with NULL; a named one still updates.
.set((
dsl::name.eq(diesel::dsl::sql::<
diesel::sql_types::Nullable<diesel::sql_types::Text>,
>("COALESCE(excluded.name, name)")),
dsl::last_message_ts.eq(diesel::dsl::sql::<diesel::sql_types::BigInt>(
"MAX(last_message_ts, excluded.last_message_ts)",
)),
dsl::muted_until.eq(diesel::dsl::sql::<
diesel::sql_types::Nullable<diesel::sql_types::BigInt>,
>(
"CASE WHEN mute_appstate_seen THEN muted_until ELSE excluded.muted_until END",
)),
dsl::archived.eq(diesel::dsl::sql::<diesel::sql_types::Bool>(
if conv.archived.is_some() {
"CASE WHEN archive_appstate_seen THEN archived ELSE excluded.archived END"
} else {
"archived"
},
)),
))
.execute(conn)?;
}
Expand Down
47 changes: 47 additions & 0 deletions crates/chat-store/src/store/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -564,6 +564,38 @@ impl ChatStore {
.map_err(|_| ChatStoreError::Store(StoreError::Validation("writer stopped".into())))
}

/// Apply a locally requested delete-for-me through the same materializer
/// as the app-state event that will later arrive from linked devices.
pub fn record_delete_for_me(
&self,
chat: &Jid,
target_id: &str,
from_me: bool,
participant: Option<Jid>,
message_timestamp_ms: i64,
timestamp: DateTime<Utc>,
) -> Result<()> {
use wacore::types::events::DeleteMessageForMeUpdate;

let update = DeleteMessageForMeUpdate::builder()
.chat_jid(chat.clone())
.maybe_participant_jid(participant)
.message_id(target_id.to_owned())
.from_me(from_me)
.timestamp(timestamp)
.action(Box::new(wa::sync_action_value::DeleteMessageForMeAction {
delete_media: Some(false),
message_timestamp: Some(message_timestamp_ms),
}))
.from_full_sync(false)
.build();
self.tx
.send(WriterMsg::Event(Arc::new(Event::DeleteMessageForMeUpdate(
update,
))))
.map_err(|_| ChatStoreError::Store(StoreError::Validation("writer stopped".into())))
}

/// Record a reaction this client just sent. An empty `emoji` removes this
/// client's existing reaction, matching the inbound event semantics.
///
Expand Down Expand Up @@ -769,6 +801,21 @@ mod migration_tests {
.expect("create store");
ChatStore::new(&store).await.expect("run migrations");

// The current top migration only tracks which source last supplied
// mute/archive preferences. Revert it first so the historical
// downgrade assertions below still start at account-cascade.
store
.shared()
.run(|conn| {
conn.revert_last_migration(MIGRATIONS)
.map(|_| ())
.map_err(StoreError::Migration)
})
.await
.expect("preference-provenance downgrade is reversible");
assert!(!has_column(&store, "chats", "mute_appstate_seen").await);
assert!(!has_column(&store, "chats", "archive_appstate_seen").await);

// Reverted in reverse application order. On top is the account-cascade
// follow-up: it only adds the `device` foreign key to the descriptors
// and the labels table, so reverting it leaves both in place without
Expand Down
Loading