mirror of
https://github.com/chatmail/core.git
synced 2026-10-04 04:00:27 +03:00
Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
720b5d3c5c | ||
|
|
73b45064f3 | ||
|
|
d1b3b9fbb7 | ||
|
|
89cafbd8b0 | ||
|
|
58061e0689 | ||
|
|
4937e0cd4f |
@@ -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(())
|
||||
}
|
||||
|
||||
@@ -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
@@ -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
@@ -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),
|
||||
|
||||
@@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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()
|
||||
);
|
||||
|
||||
@@ -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?
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user