use super::*;
impl Client {
fn is_own_jid(&self, jid: &Jid) -> bool {
let snapshot = self.persistence_manager.get_device_snapshot();
snapshot
.pn
.as_ref()
.is_some_and(|pn| pn.is_same_user_as(jid))
|| snapshot
.lid
.as_ref()
.is_some_and(|lid| lid.is_same_user_as(jid))
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.send.maybe_tc_token", level = "debug", skip_all, fields(to = %to.observe())))]
pub(super) async fn maybe_include_tc_token(
&self,
to: &Jid,
extra_nodes: &mut Vec<Node>,
sent_at: SendInstant,
) -> bool {
use wacore::iq::abprops::web;
use wacore::iq::tctoken::{
PrivacyTokenChoice, build_cs_token_node, build_tc_token_node, choose_privacy_token,
compute_cs_token, is_tc_token_expired_with_at, should_send_new_tc_token_with_at,
};
let now = sent_at.unix_secs();
if self.is_own_jid(to) {
return false;
}
if to.is_bot() || to.is_status_broadcast() {
return false;
}
let resolved_lid: Option<wacore_binary::CompactString> = if to.is_lid() {
Some(to.user.clone())
} else {
self.lid_pn_cache.get_current_lid(&to.user).await
};
let token_key: &str = resolved_lid.as_deref().unwrap_or(&to.user);
let backend = self.persistence_manager.backend();
let tc_config = self.tc_token_config().await;
let existing = match backend.get_tc_token(token_key).await {
Ok(entry) => entry,
Err(e) => {
log::warn!(target: "Client/TcToken", "Failed to get tc_token for {}: {e}", to.observe());
None
}
};
let should_issue_after_send = should_send_new_tc_token_with_at(
existing.as_ref().and_then(|entry| entry.sender_timestamp),
&tc_config,
now,
);
let valid_tc_token: Option<&[u8]> = existing.as_ref().and_then(|entry| {
(!entry.token.is_empty()
&& !is_tc_token_expired_with_at(entry.token_timestamp, &tc_config, now))
.then_some(entry.token.as_slice())
});
let snapshot = self.persistence_manager.get_device_snapshot();
let cs_token_inputs: Option<(&[u8], &wacore_binary::CompactString)> =
match (&snapshot.nct_salt, &resolved_lid) {
(Some(salt), Some(lid)) => Some((salt.as_slice(), lid)),
_ => None,
};
let tc_send_enabled = self
.ab_props
.is_enabled(web::PRIVACY_TOKEN_SENDING_ON_ALL_1_ON_1_MESSAGES)
.await;
let nct_send_enabled = self
.ab_props
.is_enabled(web::WA_NCT_TOKEN_SEND_ENABLED)
.await;
let choice = choose_privacy_token(
tc_send_enabled,
nct_send_enabled,
valid_tc_token.is_some(),
cs_token_inputs.is_some(),
);
match choice {
PrivacyTokenChoice::TcToken => {
extra_nodes.extend(valid_tc_token.map(build_tc_token_node))
}
PrivacyTokenChoice::CsToken => {
extra_nodes.extend(cs_token_inputs.map(|(salt, lid_user)| {
let recipient_lid = Jid::new(lid_user.as_str(), Server::Lid).to_string();
build_cs_token_node(&compute_cs_token(salt, &recipient_lid))
}));
}
PrivacyTokenChoice::None => {}
}
log::debug!(target: "Client/TcToken", "privacy token for {}: {choice:?}", to.observe());
should_issue_after_send
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.send.issue_tc_token", level = "debug", skip_all, fields(to = %to.observe())))]
pub(crate) async fn issue_tc_token_after_send(&self, to: &Jid) {
use wacore::iq::tctoken::IssuePrivacyTokensSpec;
if to.is_bot() || to.is_status_broadcast() {
return;
}
let issuance_jid = self.resolve_issuance_jid(to).await;
if let Err(e) = self
.execute(IssuePrivacyTokensSpec::new(std::slice::from_ref(
&issuance_jid,
)))
.await
{
log::debug!(target: "Client/TcToken", "Failed to issue tc_token for {}: {e}", issuance_jid.observe());
return;
}
self.record_tc_token_sender_timestamp(to).await;
}
#[cfg(feature = "voip-runtime")]
pub(crate) async fn should_issue_tc_token(&self, to: &Jid) -> bool {
use wacore::iq::tctoken::should_send_new_tc_token_with;
if self.is_own_jid(to) {
return false;
}
if to.is_bot() || to.is_status_broadcast() {
return false;
}
let key = self.resolve_tc_token_key(to).await;
let sender_ts = match self.persistence_manager.backend().get_tc_token(&key).await {
Ok(entry) => entry.and_then(|e| e.sender_timestamp),
Err(e) => {
log::warn!(target: "Client/TcToken", "Failed to read tc_token for {}: {e}", to.observe());
None
}
};
should_send_new_tc_token_with(sender_ts, &self.tc_token_config().await)
}
pub(crate) async fn store_issued_tc_tokens(
&self,
tokens: &[wacore::iq::tctoken::ReceivedTcToken],
) {
use wacore::store::traits::TcTokenEntry;
let backend = self.persistence_manager.backend();
let now = wacore::time::now_secs();
for received in tokens {
if received.token.is_empty() {
log::warn!(target: "Client/TcToken", "Server returned empty tc_token for {}, skipping", received.jid.observe());
continue;
}
let entry = TcTokenEntry {
token: received.token.clone(),
token_timestamp: received.timestamp,
sender_timestamp: Some(now),
};
if let Err(e) = backend.put_tc_token(&received.jid.user, &entry).await {
log::warn!(target: "Client/TcToken", "Failed to store issued tc_token: {e}");
}
}
}
async fn store_issued_tc_tokens_with_sender_ts(
&self,
tokens: &[wacore::iq::tctoken::ReceivedTcToken],
sender_ts: i64,
) {
use wacore::store::traits::TcTokenEntry;
let backend = self.persistence_manager.backend();
for received in tokens {
if received.token.is_empty() {
continue;
}
let entry = TcTokenEntry {
token: received.token.clone(),
token_timestamp: received.timestamp,
sender_timestamp: Some(sender_ts),
};
if let Err(e) = backend.put_tc_token(&received.jid.user, &entry).await {
log::warn!(target: "Client/TcToken", "Failed to store re-issued tc_token: {e}");
}
}
}
async fn record_tc_token_sender_timestamp(&self, to: &Jid) {
let key = self.resolve_tc_token_key(to).await;
let now = wacore::time::now_secs();
if let Err(e) = self
.persistence_manager
.backend()
.touch_tc_token_sender_timestamp(&key, now)
.await
{
log::warn!(target: "Client/TcToken", "Failed to record tc_token sender_timestamp for {}: {e}", to.observe());
}
}
#[cfg_attr(feature = "tracing", tracing::instrument(name = "wa.send.reissue_tc_token", level = "debug", skip_all, fields(sender = %sender.observe())))]
pub(crate) async fn reissue_tc_token_after_identity_change(&self, sender: &Jid) {
use wacore::iq::tctoken::{IssuePrivacyTokensSpec, is_sender_tc_token_expired};
let bare = sender.to_non_ad_string();
let mutex = self.session_lock_for(&bare).await;
let Some(_guard) = mutex.try_lock() else {
return;
};
let token_jid = self.resolve_tc_token_key(sender).await;
let backend = self.persistence_manager.backend();
let entry = match backend.get_tc_token(&token_jid).await {
Ok(Some(e)) => e,
_ => return,
};
let Some(sender_ts) = entry.sender_timestamp else {
return;
};
let tc_config = self.tc_token_config().await;
if is_sender_tc_token_expired(sender_ts, &tc_config) {
return;
}
let issuance_jid = self.resolve_issuance_jid(sender).await;
match self
.execute(IssuePrivacyTokensSpec::with_timestamp(
std::slice::from_ref(&issuance_jid),
sender_ts,
))
.await
{
Ok(response) => {
self.store_issued_tc_tokens_with_sender_ts(&response.tokens, sender_ts)
.await;
log::debug!(
target: "Client/TcToken",
"Re-issued tctoken after identity change for {}",
sender.observe()
);
}
Err(e) => {
log::debug!(
target: "Client/TcToken",
"Failed to re-issue tctoken after identity change for {}: {e}",
sender.observe()
);
}
}
}
pub(crate) async fn lookup_tc_token_for_jid(&self, jid: &Jid) -> Option<Vec<u8>> {
use wacore::iq::tctoken::is_tc_token_expired_with;
let token_key = self.resolve_tc_token_key(jid).await;
let tc_config = self.tc_token_config().await;
let backend = self.persistence_manager.backend();
match backend.get_tc_token(&token_key).await {
Ok(Some(entry))
if !entry.token.is_empty()
&& !is_tc_token_expired_with(entry.token_timestamp, &tc_config) =>
{
Some(entry.token)
}
Ok(_) => None,
Err(e) => {
log::warn!(target: "Client/TcToken", "Failed to get tc_token for {}: {e}", jid.observe());
None
}
}
}
pub(crate) async fn tc_token_config(&self) -> wacore::iq::tctoken::TcTokenConfig {
use wacore::iq::abprops::web;
use wacore::iq::tctoken::TcTokenConfig;
TcTokenConfig {
bucket_duration: self.ab_props.get_int(web::TCTOKEN_DURATION).await,
num_buckets: self.ab_props.get_int(web::TCTOKEN_NUM_BUCKETS).await,
sender_bucket_duration: self.ab_props.get_int(web::TCTOKEN_DURATION_SENDER).await,
sender_num_buckets: self.ab_props.get_int(web::TCTOKEN_NUM_BUCKETS_SENDER).await,
}
.clamped()
}
pub(crate) async fn resolve_tc_token_key(&self, jid: &Jid) -> String {
if jid.is_lid() {
jid.user.to_string()
} else if let Some(lid_user) = self.lid_pn_cache.get_current_lid(&jid.user).await {
lid_user.to_string()
} else {
jid.user.to_string()
}
}
async fn resolve_to_lid_jid(&self, jid: &Jid) -> Jid {
if jid.is_lid() {
return jid.to_non_ad();
}
if let Some(lid_user) = self.lid_pn_cache.get_current_lid(&jid.user).await {
Jid::new(lid_user, Server::Lid)
} else {
jid.to_non_ad()
}
}
async fn resolve_issuance_jid(&self, jid: &Jid) -> Jid {
use wacore::iq::abprops::web;
let issue_to_lid = self
.ab_props
.is_enabled(web::LID_TRUSTED_TOKEN_ISSUE_TO_LID)
.await;
let resolved = if issue_to_lid {
self.resolve_to_lid_jid(jid).await
} else if jid.is_lid() {
if let Some(pn) = self.lid_pn_cache.get_phone_number(&jid.user).await {
Jid::new(&pn, Server::Pn)
} else {
jid.to_non_ad()
}
} else {
jid.to_non_ad()
};
resolved.into_non_ad()
}
}
#[cfg(test)]
mod tests {
use crate::test_utils::create_test_client;
use wacore::store::traits::TcTokenEntry;
use wacore_binary::{Jid, Server};
#[tokio::test]
async fn record_sender_timestamp_creates_byteless_placeholder() {
let client = create_test_client().await;
let jid = Jid::new("770000001", Server::Lid);
client.record_tc_token_sender_timestamp(&jid).await;
let entry = client
.persistence_manager
.backend()
.get_tc_token("770000001")
.await
.unwrap()
.expect("placeholder entry should be created");
assert!(entry.token.is_empty(), "placeholder carries no token bytes");
assert!(
entry.sender_timestamp.is_some(),
"placeholder records the issuance timestamp"
);
}
#[tokio::test]
async fn record_sender_timestamp_preserves_existing_token() {
let client = create_test_client().await;
let backend = client.persistence_manager.backend();
backend
.put_tc_token(
"770000002",
&TcTokenEntry {
token: vec![1, 2, 3],
token_timestamp: 1_700_000_000,
sender_timestamp: None,
},
)
.await
.unwrap();
let jid = Jid::new("770000002", Server::Lid);
client.record_tc_token_sender_timestamp(&jid).await;
let entry = backend.get_tc_token("770000002").await.unwrap().unwrap();
assert_eq!(entry.token, vec![1, 2, 3], "received token is preserved");
assert_eq!(entry.token_timestamp, 1_700_000_000);
assert!(
entry.sender_timestamp.is_some(),
"sender_timestamp is advanced on issuance"
);
}
#[cfg(feature = "voip-runtime")]
#[tokio::test]
async fn should_issue_tc_token_true_for_unknown_contact() {
let client = create_test_client().await;
let jid = Jid::new("770000003", Server::Lid);
assert!(
client.should_issue_tc_token(&jid).await,
"a contact with no recorded issuance should get a token"
);
}
#[cfg(feature = "voip-runtime")]
#[tokio::test]
async fn should_issue_tc_token_false_within_sender_bucket() {
let client = create_test_client().await;
client
.persistence_manager
.backend()
.put_tc_token(
"770000004",
&TcTokenEntry {
token: Vec::new(),
token_timestamp: 0,
sender_timestamp: Some(wacore::time::now_secs()),
},
)
.await
.unwrap();
let jid = Jid::new("770000004", Server::Lid);
assert!(
!client.should_issue_tc_token(&jid).await,
"a fresh issuance in the current bucket must not re-issue"
);
}
#[cfg(feature = "voip-runtime")]
#[tokio::test]
async fn should_issue_tc_token_false_for_self() {
let client = create_test_client().await;
let own = Jid::new("999000111", Server::Lid);
client
.persistence_manager
.process_command(crate::store::commands::DeviceCommand::SetLid(Some(
own.clone(),
)))
.await;
assert!(
!client.should_issue_tc_token(&own).await,
"a self-call must never issue a tc token for our own account"
);
}
}