Compare commits

...
9 changed files with 368 additions and 16 deletions
+147
View File
@@ -5,9 +5,12 @@ use deltachat_contact_tools::addr_normalize;
use rand::distr::{Alphanumeric, SampleString};
use rand::seq::IndexedRandom;
use crate::chat::add_device_msg;
use crate::config::{self, Config};
use crate::contact::ContactId;
use crate::log::{LogExt, warn};
use crate::login_param::{EnteredCertificateChecks, EnteredImapLoginParam};
use crate::message::{Message, MsgId};
use crate::{configure::EnteredLoginParam, context::Context, tools::time};
/// The target number of transports.
@@ -174,5 +177,149 @@ pub(crate) fn login_param_from_host(host: &str) -> EnteredLoginParam {
}
}
pub(crate) async fn record_message_sent_via_transport_by_msg_id(
context: &Context,
msg_id: MsgId,
transport_id: u32,
) -> Result<()> {
let Some(contact_id): Option<ContactId> = context
.sql
.query_get_value("SELECT from_id FROM msgs WHERE id=?", (msg_id,))
.await?
else {
warn!(context, "Invalid msg_id");
return Ok(());
};
record_message_sent_via_transport(context, contact_id, transport_id).await?;
Ok(())
}
/// Records that the specified contact sent a message via the specified transport,
/// i.e. that we can be sure that the contact knows about this transport.
/// This info is saved into the `transport_awareness_by_contacts` table.
pub(crate) async fn record_message_sent_via_transport(
context: &Context,
contact_id: ContactId,
transport_id: u32,
) -> Result<()> {
if contact_id.is_special() {
return Ok(());
}
context
.sql
.execute(
"INSERT OR REPLACE INTO transport_awareness_by_contacts(contact_id, transport_id, last_seen)
VALUES (?,?,?)",
(contact_id, transport_id, time()),
)
.await?;
Ok(())
}
pub(crate) async fn output_debug_transport_awareness(context: &Context) -> Result<()> {
let msg = get_debug_transport_awareness(context).await?;
let mut msg = Message::new_text(msg);
add_device_msg(context, None, Some(&mut msg)).await?;
Ok(())
}
async fn get_debug_transport_awareness(context: &Context) -> Result<String, anyhow::Error> {
fn get_contact_name(row: &rusqlite::Row<'_>) -> Result<String> {
let name: String = row.get(0)?;
let authname: String = row.get(1)?;
if name.is_empty() {
Ok(authname)
} else {
Ok(name)
}
}
fn filter_out_empty<F>(rows: rusqlite::AndThenRows<'_, F>) -> Result<Vec<String>>
where
F: FnMut(&rusqlite::Row<'_>) -> Result<String>,
{
rows.filter(|name| name.as_ref().is_ok_and(|n| !n.is_empty()))
.collect()
}
let mut msg = "=== Usage of transports by contacts ===".to_string();
let transports = context
.sql
.query_map_vec("SELECT id, addr FROM transports", (), |row| {
let id: u32 = row.get(0)?;
let addr: String = row.get(1)?;
Ok((id, addr))
})
.await?;
for (transport_id, addr) in transports {
let lost_contacts: Vec<String> = context
.sql
.query_map(
// Select all contacts that can reach us via this transport,
// but can't reach us via the other transports.
// These are the contacts we may lose by removing this transport.
"SELECT name, authname FROM contacts
WHERE id IN (SELECT contact_id FROM transport_awareness_by_contacts WHERE transport_id = ?1)
AND id NOT IN (
SELECT contact_id FROM transport_awareness_by_contacts WHERE transport_id IN (
SELECT id FROM transports WHERE id != ?1
)
)
AND id>9
ORDER BY last_seen DESC",
(transport_id,),
get_contact_name,
filter_out_empty,
)
.await?;
msg += &format!("\nThese contacts fail to reach you if you remove transport {addr}:\n");
msg += &lost_contacts.join("\n");
}
let lost_contacts: Vec<String> = context
.sql
.query_map(
// Select all contacts that cannot reach us via one of the current transports,
// but could reach us via some transport that's not in use anymore:
"SELECT name, authname FROM contacts
WHERE id NOT IN (
SELECT contact_id FROM transport_awareness_by_contacts WHERE transport_id IN (
SELECT id FROM transports
)
)
AND id IN (SELECT contact_id FROM transport_awareness_by_contacts)
AND id>9
ORDER BY last_seen DESC",
(),
get_contact_name,
filter_out_empty,
)
.await?;
msg += "\nThese contacts likely can't reach you anymore:\n";
msg += &lost_contacts.join("\n");
let unknown_contacts: Vec<String> = context
.sql
.query_map(
// Select all contacts that never sent a message
// since we started recording transport usage
"SELECT name, authname FROM contacts
WHERE id NOT IN (SELECT contact_id FROM transport_awareness_by_contacts)
AND id>9
ORDER BY last_seen DESC",
(),
get_contact_name,
filter_out_empty,
)
.await?;
msg += "\nFor these contacts, it's unclear how they can reach you:\n";
msg += &unknown_contacts.join("\n");
Ok(msg)
}
#[cfg(test)]
mod automatic_relay_management_tests;
@@ -1,8 +1,12 @@
use std::time::Duration;
use anyhow::Context as _;
use super::*;
use crate::test_utils::TestContext;
use crate::imap::prefetch_should_download;
use crate::test_utils::{TestContext, TestContextManager};
use crate::tools::SystemTime;
use crate::transport::add_pseudo_transport;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_load_relay_candidates_single() -> Result<()> {
@@ -290,3 +294,156 @@ async fn enable_config(context: &Context) {
.await
.unwrap();
}
/// Tests that `record_message_sent_via_transport()`,
/// `record_message_sent_via_transport_by_msg_id()` and `prefetch_should_download()`
/// correctly populate the `transport_awareness_by_contacts` table,
/// and that `output_debug_transport_awareness()`
/// turns that table into a human-readable device message.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_transport_awareness() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
bob.set_config(Config::Displayname, Some("Bob")).await?;
let charlie = &tcm.charlie().await;
charlie
.set_config(Config::Displayname, Some("Charlie"))
.await?;
let dom = &tcm.dom().await;
dom.set_config(Config::Displayname, Some("Dom")).await?;
let fiona = &tcm.fiona().await;
fiona.set_config(Config::Displayname, Some("Fiona")).await?;
add_pseudo_transport(alice, "transport-a@example.org").await?;
add_pseudo_transport(alice, "transport-b@example.org").await?;
let transport_a: u32 = alice
.sql
.query_get_value(
"SELECT id FROM transports WHERE addr=?",
("transport-a@example.org",),
)
.await?
.context("transport A not found")?;
let transport_b: u32 = alice
.sql
.query_get_value(
"SELECT id FROM transports WHERE addr=?",
("transport-b@example.org",),
)
.await?
.context("transport B not found")?;
// Fiona is only reachable via transport A, recorded directly
// via `record_message_sent_via_transport()`.
let fiona_id = alice.add_or_lookup_contact_id(fiona).await;
record_message_sent_via_transport(alice, fiona_id, transport_a).await?;
// Bob is only reachable via transport B, recorded indirectly
// via `record_message_sent_via_transport_by_msg_id()`.
let bob_id = alice.add_or_lookup_contact_id(bob).await;
alice
.sql
.execute(
"INSERT INTO msgs (rfc724_mid, from_id) VALUES (?, ?)",
("bob-message@localhost", bob_id),
)
.await?;
let bob_msg_id: MsgId = alice
.sql
.query_get_value(
"SELECT id FROM msgs WHERE rfc724_mid=?",
("bob-message@localhost",),
)
.await?
.context("Bob's message not found")?;
record_message_sent_via_transport_by_msg_id(alice, bob_msg_id, transport_b).await?;
// Charlie is reachable via both transports, so she should not show up as
// "at risk" for either of them. His usage of transport A is recorded
// directly, his usage of transport B is recorded by simulating that an
// already-fetched message of him is prefetched again (as happens e.g.
// when the same message is visible on two transports).
let charlie_id = alice.add_or_lookup_contact_id(charlie).await;
record_message_sent_via_transport(alice, charlie_id, transport_a).await?;
alice
.sql
.execute(
"INSERT INTO msgs (rfc724_mid, from_id) VALUES (?, ?)",
("charlie-message@localhost", charlie_id),
)
.await?;
let download = prefetch_should_download(
alice,
&[],
"charlie-message@localhost",
std::iter::empty(),
transport_b,
)
.await?;
// The message was already fetched, so it should not be downloaded again
assert_eq!(download, false);
// Dom never sent us anything on any transport, so it's unclear how they
// can reach us.
let _dom_id = alice.add_or_lookup_contact_id(dom).await;
let mut recorded: Vec<(ContactId, u32)> = alice
.sql
.query_map_vec(
"SELECT contact_id, transport_id FROM transport_awareness_by_contacts",
(),
|row| {
let contact_id: ContactId = row.get(0)?;
let transport_id: u32 = row.get(1)?;
Ok((contact_id, transport_id))
},
)
.await?;
recorded.sort();
let mut expected = vec![
(fiona_id, transport_a),
(bob_id, transport_b),
(charlie_id, transport_a),
(charlie_id, transport_b),
];
expected.sort();
assert_eq!(recorded, expected);
let actual = get_debug_transport_awareness(alice).await?;
assert_eq!(
actual,
"=== Usage of transports by contacts ===
These contacts fail to reach you if you remove transport alice@example.org:
These contacts fail to reach you if you remove transport transport-a@example.org:
Fiona
These contacts fail to reach you if you remove transport transport-b@example.org:
Bob
These contacts likely can't reach you anymore:
For these contacts, it's unclear how they can reach you:
Dom",
"Transport usage output didn't match, actual output was:\n{actual}\n"
);
alice.delete_transport("transport-a@example.org").await?;
let actual = get_debug_transport_awareness(alice).await?;
assert_eq!(
actual,
"=== Usage of transports by contacts ===
These contacts fail to reach you if you remove transport alice@example.org:
These contacts fail to reach you if you remove transport transport-b@example.org:
Bob
Charlie
These contacts likely can't reach you anymore:
Fiona
For these contacts, it's unclear how they can reach you:
Dom",
"Transport usage output didn't match, actual output was:\n{actual}\n"
);
Ok(())
}
+6
View File
@@ -16,6 +16,7 @@ use mail_builder::mime::MimePart;
use serde::{Deserialize, Serialize};
use strum_macros::EnumIter;
use crate::automatic_relay_management::output_debug_transport_awareness;
use crate::blob::BlobObject;
use crate::chatlist::Chatlist;
use crate::chatlist_events;
@@ -2750,6 +2751,11 @@ async fn prepare_send_msg(
}
chat.prepare_msg_raw(context, msg, update_msg_id).await?;
info!(context, "dbg sending msg {msg:?}");
if msg.text == "/transport_awareness" && chat.is_self_talk() {
output_debug_transport_awareness(context).await?;
}
let row_ids = create_send_msg_jobs(context, msg)
.await
.context("Failed to create send jobs")?;
+25 -8
View File
@@ -21,8 +21,6 @@ use futures_lite::FutureExt;
use ratelimit::Ratelimit;
use url::Url;
use crate::chat::{self, add_device_msg};
use crate::config::Config;
use crate::constants::DC_VERSION_STR;
use crate::context::Context;
use crate::ensure_and_debug_assert;
@@ -43,7 +41,11 @@ use crate::transport::{
ConfiguredLoginParam, ConfiguredServerLoginParam, prioritize_server_login_params,
};
use crate::{
automatic_relay_management::record_message_sent_via_transport,
automatic_relay_management::record_message_sent_via_transport_by_msg_id,
calls::{UnresolvedIceServer, create_fallback_ice_servers, create_ice_servers_from_metadata},
chat::{self, add_device_msg},
config::Config,
ephemeral::delete_expired_imap_messages,
};
@@ -651,9 +653,15 @@ impl Imap {
// message, move it to the movebox and then download the second message before
// downloading the first one, if downloading from inbox before moving is allowed.
if folder == target
&& prefetch_should_download(context, &headers, &message_id, fetch_response.flags())
.await
.context("prefetch_should_download")?
&& prefetch_should_download(
context,
&headers,
&message_id,
fetch_response.flags(),
session.transport_id(),
)
.await
.context("prefetch_should_download")?
{
if headers
.get_header_value(HeaderDef::ChatIsPostMessage)
@@ -1238,6 +1246,11 @@ impl Session {
}
Ok(msg) => msg,
};
if let Some(received_msg) = &received_msg {
record_message_sent_via_transport(context, received_msg.from_id, transport_id)
.await
.context("Recording transport")?;
}
received_msgs_channel
.send((request_uid, received_msg))
.await?;
@@ -1620,15 +1633,19 @@ pub(crate) fn create_message_id() -> String {
pub(crate) async fn prefetch_should_download(
context: &Context,
headers: &[mailparse::MailHeader<'_>],
message_id: &str,
rfc724_mid: &str,
mut flags: impl Iterator<Item = Flag<'_>>,
transport_id: u32,
) -> Result<bool> {
if message::rfc724_mid_fetch_tried(context, message_id).await? {
if let Some(msg_id) = message::rfc724_mid_fetch_tried(context, rfc724_mid).await? {
if let Some(from) = mimeparser::get_from(headers)
&& context.is_self_addr(&from.addr).await?
{
markseen_on_imap_table(context, message_id).await?;
markseen_on_imap_table(context, rfc724_mid).await?;
}
record_message_sent_via_transport_by_msg_id(context, msg_id, transport_id)
.await
.context("Recording transport")?;
return Ok(false);
}
+9 -6
View File
@@ -2187,18 +2187,21 @@ pub(crate) async fn rfc724_mid_exists_ex(
Ok(res)
}
/// Returns `true` if the given `rfc724_mid` has nothing left to fetch from a server,
/// Returns `Some(msg_id)` if the given `rfc724_mid` has nothing left to fetch from a server,
/// i.e. it was already fetched or is an outgoing message.
///
/// For post-messages, this returns `true` if an attempt to fetch was made or is ongoing,
/// For post-messages, this returns `Some(msg_id)` if an attempt to fetch was made or is ongoing,
/// even if this was not successful,
/// because we don't want to automatically try fetching these messages over and over again
/// (this function is not called when the user manually clicked "Download").
pub(crate) async fn rfc724_mid_fetch_tried(context: &Context, rfc724_mid: &str) -> Result<bool> {
pub(crate) async fn rfc724_mid_fetch_tried(
context: &Context,
rfc724_mid: &str,
) -> Result<Option<MsgId>> {
let rfc724_mid = rfc724_mid.trim_start_matches('<').trim_end_matches('>');
if rfc724_mid.is_empty() {
warn!(context, "Empty rfc724_mid passed to rfc724_mid_fetch_tried");
return Ok(false);
return Ok(None);
}
// Explanation of the SQL statement:
@@ -2215,8 +2218,8 @@ pub(crate) async fn rfc724_mid_fetch_tried(context: &Context, rfc724_mid: &str)
// so that we do not need to check the download state.
let res = context
.sql
.exists(
"SELECT COUNT(*) FROM msgs
.query_get_value(
"SELECT id FROM msgs
WHERE (rfc724_mid=?1 AND download_state<>?2)
OR pre_rfc724_mid=?1",
(rfc724_mid, DownloadState::Available),
+6
View File
@@ -82,6 +82,9 @@ pub struct ReceivedMsg {
/// Whether IMAP messages should be immediately deleted.
pub needs_delete_job: bool,
/// The database ID of the contact that sent the message.
pub(crate) from_id: ContactId,
}
/// Decision on which kind of chat the message
@@ -492,6 +495,7 @@ pub(crate) async fn receive_imf_inner(
sort_timestamp: 0,
msg_ids,
needs_delete_job: false,
from_id: ContactId::UNDEFINED,
}))
};
@@ -669,6 +673,7 @@ pub(crate) async fn receive_imf_inner(
sort_timestamp: mime_parser.timestamp_sent,
msg_ids: vec![msg_id],
needs_delete_job: res == securejoin::HandshakeMessage::Done,
from_id,
});
}
securejoin::HandshakeMessage::Propagate => {
@@ -2353,6 +2358,7 @@ INSERT INTO msgs
sort_timestamp,
msg_ids: created_db_entries,
needs_delete_job: false,
from_id,
})
}
+1 -1
View File
@@ -648,7 +648,7 @@ async fn test_parse_ndn(
// Check that the ndn would be downloaded:
let headers = mailparse::parse_mail(raw_ndn).unwrap().headers;
assert!(
prefetch_should_download(&t, &headers, "some-other-message-id", std::iter::empty(),)
prefetch_should_download(&t, &headers, "some-other-message-id", std::iter::empty(), 0)
.await
.unwrap()
);
+15
View File
@@ -2610,6 +2610,21 @@ UPDATE msgs SET state=24 WHERE state=18; -- Change OutPreparing to OutFailed.
.await?;
}
inc_and_check(&mut migration_version, 164)?;
if dbversion < migration_version {
sql.execute_migration(
"CREATE TABLE transport_awareness_by_contacts(
contact_id INTEGER NOT NULL,
transport_id INTEGER NOT NULL,
last_seen INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY(contact_id, transport_id),
FOREIGN KEY(contact_id) REFERENCES contacts(id) ON DELETE CASCADE
) STRICT",
migration_version,
)
.await?;
}
let new_version = sql
.get_raw_config_int(VERSION_CFG)
.await?
+1
View File
@@ -173,6 +173,7 @@ async fn test_receive_pre_message_and_dl_post_message() -> Result<()> {
&headers,
&headers.get_header_value(HeaderDef::MessageId).unwrap(),
std::iter::empty(),
0,
)
.await
.unwrap()