Compare commits

...
11 changed files with 325 additions and 68 deletions
+12
View File
@@ -525,6 +525,18 @@ impl CommandApi {
ctx.add_transport_from_qr(&qr).await
}
/// Automatically adds up to three transports.
///
/// If the user just scanned a QR code of type `Account`, `Login`,
/// `AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
/// then UI implementations should pass it as the `qr` parameter.
/// The host(s) from the QR code will then also be considered
/// for creating an account there.
async fn init_transports(&self, account_id: u32, qr: Option<String>) -> Result<()> {
let ctx = self.get_context(account_id).await?;
ctx.init_transports(qr.as_deref()).await
}
/// Returns the list of all email accounts that are used as a transport in the current profile.
/// Use [Self::add_or_update_transport()] to add or change a transport
/// and [Self::delete_transport()] to remove a transport.
@@ -139,6 +139,18 @@ class Account:
"""Add a new transport using a QR code."""
yield self._rpc.add_transport_from_qr.future(self.id, qr)
@futuremethod
def init_transports(self, qr: Optional[str] = None):
"""Automatically adds up to three transports.
If the user just scanned a QR code of type `Account`, `Login`,
`AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
then UI implementations should pass it as the `qr` parameter.
The host(s) from the QR code will then also be considered
for creating an account there.
"""
yield self._rpc.init_transports.future(self.id, qr)
def delete_transport(self, addr: str):
"""Delete a transport."""
self._rpc.delete_transport(self.id, addr)
@@ -92,9 +92,10 @@ class RPCAccountFactory:
"""Create a new configured account."""
account = self.get_unconfigured_account()
qr = self.get_account_qr()
yield account.add_transport_from_qr.future(qr)
yield account.init_transports.future(qr)
assert account.is_configured()
assert len(account.list_transports()) == 1
return account
def new_configured_bot(self) -> Bot:
@@ -32,6 +32,12 @@ def test_add_second_address(acf) -> None:
account.add_transport_from_qr(qr)
assert len(account.list_transports()) == 2
# init_transports() only works on an unconfigured profile:
with pytest.raises(JsonRpcError):
account.init_transports(qr)
with pytest.raises(JsonRpcError):
account.init_transports()
account.add_transport_from_qr(qr)
assert len(account.list_transports()) == 3
+9
View File
@@ -631,6 +631,15 @@ CREATE TABLE broadcast_secrets(
-- Candidate chatmail relays for automatic relay management.
CREATE TABLE relay_candidates(
host TEXT PRIMARY KEY NOT NULL,
last_tried INTEGER NOT NULL DEFAULT 0 -- Deprecated 2026-09, replaced with separate relay_candidates_last_tried table.
) STRICT;
-- This table is used for storing the timestamp of the last
-- connection attempt per chatmail relay candidate.
-- This table can contain relays that were removed from the list of candidates,
-- and it does not contain the default relays.
CREATE TABLE relay_candidates_last_tried(
host TEXT PRIMARY KEY NOT NULL,
last_tried INTEGER NOT NULL DEFAULT 0 -- Timestamp of the last connection attempt.
) STRICT;
+104 -25
View File
@@ -2,8 +2,8 @@
//!
//! Chatmail relays create an account on first login,
//! so a profile can add further transports on its own without user interaction.
//! Candidate hosts come from the `relay_candidates` table,
//! which migrations seed with a list of known chatmail relays.
//! Candidate hosts come from the `relay_candidates` table
//! as well as the [`DEFAULT_RELAY_CANDIDATES`] list.
//!
//! Status of implementation:
//! Additions are attempted right before going into IMAP IDLE,
@@ -13,6 +13,7 @@
//! [`Config::AutorelayFinished`] is set and nothing is ever added again,
//! so deleting a transport later does not pull in a replacement.
use std::collections::BTreeSet;
use std::pin::Pin;
use anyhow::Result;
@@ -23,6 +24,7 @@ use rand::seq::IndexedRandom;
use crate::config::{self, Config};
use crate::log::{LogExt, warn};
use crate::login_param::{EnteredCertificateChecks, EnteredImapLoginParam};
use crate::sql::TransactionExt as _;
use crate::{configure::EnteredLoginParam, context::Context, tools::time};
/// The target number of transports.
@@ -32,6 +34,59 @@ const AUTOMATIC_ADDITION_DEBOUNCE_SECONDS: i64 = 60 * 60; // one hour
/// How long we ignore a relay candidate after failing to connect to it:
const BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY: i64 = 60 * 60 * 24 * 7; // one week
/// The list of relays to which we onboard
/// if no other relays are provided via QR codes.
/// Please keep this list alphabetically sorted.
const DEFAULT_RELAY_CANDIDATES: &[&str] = &[
"chat.adminforge.de",
"chat.feld.me",
"chat.me.ke",
"chat.nuvon.app",
"chat.tinydispatch.org",
"chat.vim.wtf",
"chatmail.uk",
"chtml.ca",
"deltachat.me",
"e2e.sus.fr",
"jp.deltachat.me",
"mailchat.pl",
"nchrcht.la10cy.net",
"nine.testrun.org",
"sweetfern.net",
"tarpit.fun",
];
pub(crate) async fn init_transports_inner(
context: &Context,
addrs_from_qr: Vec<String>,
skip_network: bool,
) -> Result<()> {
context
.sql
.transaction(|transaction| {
let mut stmt = transaction.prepare("INSERT INTO relay_candidates(host) VALUES(?)")?;
for addr in addrs_from_qr {
stmt.execute((addr,))?;
}
Ok(())
})
.await?;
let host = "nine.testrun.org";
let param = login_param_from_host(host);
let res = crate::configure::configure(context, &param, skip_network).await;
if let Err(err) = &res {
warn!(context, "Failed to init transports: {err:#}.");
} else {
info!(context, "Initialized with transport {host}");
}
res?;
context.set_config_bool(Config::Autorelay, true).await?;
Ok(())
}
pub(crate) fn maybe_add_additional_relays(
context: Context,
) -> Pin<Box<dyn Future<Output = ()> + Send>> {
@@ -85,12 +140,16 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
let mut relay_added = false;
// Using `for` instead of `while` to prevent infinite loop
for _ in 0..NUM_TRANSPORTS_TARGET {
if context.count_transports().await? >= NUM_TRANSPORTS_TARGET {
let num_transports = context.count_transports().await?;
if num_transports >= NUM_TRANSPORTS_TARGET {
info!(context, "Transports target reached at {num_transports}");
context
.set_config_internal(Config::AutorelayFinished, config::from_bool(true))
.await?;
return Ok(relay_added);
} else {
info!(context, "There are {num_transports} relays, will add more");
}
// First, query all candidates that were not tried since `BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY` seconds.
@@ -111,13 +170,7 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
candidates.len(),
);
context
.sql
.execute(
"UPDATE relay_candidates SET last_tried=? WHERE host=?",
(now, host),
)
.await?;
set_relay_candidate_last_tried(context, host, now).await?;
let param = login_param_from_host(host);
let res = crate::configure::configure(context, &param, skip_network).await;
if let Err(e) = res {
@@ -134,28 +187,54 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
Ok(relay_added)
}
async fn load_relay_candidates(context: &Context, now: i64) -> Result<Vec<String>, anyhow::Error> {
let cutoff_timestamp = now.saturating_sub(BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY);
let candidates: Vec<String> = context
async fn set_relay_candidate_last_tried(
context: &Context,
host: &str,
now: i64,
) -> Result<(), anyhow::Error> {
context
.sql
.query_map_vec(
// This also selects candidates which have last_tried in the future,
.execute(
"INSERT OR REPLACE INTO relay_candidates_last_tried(host, last_tried) VALUES(?, ?)",
(host, now),
)
.await?;
Ok(())
}
async fn load_relay_candidates(context: &Context, now: i64) -> Result<Vec<String>> {
let res = context
.sql
.transaction(|transaction| {
let mut candidates: BTreeSet<String> =
transaction.query_map_collect("SELECT host FROM relay_candidates", (), |row| {
Ok(row.get(0)?)
})?;
candidates.extend(DEFAULT_RELAY_CANDIDATES.iter().map(|s| s.to_string()));
let cutoff_timestamp = now.saturating_sub(BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY);
// This does not select candidates which have last_tried in the future,
// essentially treating them as never tried,
// so if some timestamp far in the future is accidentally stored,
// we are not stuck never trying the candidate.
// After trying the candidate, last_tried will be corrected to the current time.
"SELECT host FROM relay_candidates WHERE (last_tried<? OR last_tried>?)
AND NOT EXISTS (
SELECT 1
FROM transports
WHERE substr(addr, instr(addr, '@') + 1) = host
)",
(cutoff_timestamp, now),
|row| Ok(row.get::<_, String>(0)?),
)
let exclude: BTreeSet<String> = transaction.query_map_collect(
"SELECT host FROM relay_candidates_last_tried WHERE (last_tried>=? AND last_tried<=?)
UNION
SELECT substr(addr, instr(addr, '@') + 1) FROM transports",
(cutoff_timestamp, now),
|row| Ok(row.get(0)?),
)?;
Ok(candidates
.difference(&exclude)
.map(|s| s.to_string())
.collect::<Vec<String>>())
})
.await?;
Ok(candidates)
Ok(res)
}
pub(crate) fn login_param_from_host(host: &str) -> EnteredLoginParam {
+82 -40
View File
@@ -4,41 +4,78 @@ use super::*;
use crate::test_utils::TestContext;
use crate::tools::SystemTime;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_init_transports_basic() -> Result<()> {
let t = &TestContext::new().await;
assert!(t.list_transports().await?.is_empty());
let skip_network = true;
init_transports_inner(t, vec![], skip_network).await?;
let relays = get_configured_relays(t).await;
assert_eq!(relays.len(), 1);
for relay in &relays {
assert!(DEFAULT_RELAY_CANDIDATES.contains(&relay.as_ref()));
}
assert_eq!(t.get_config_bool(Config::Autorelay).await?, true);
Ok(())
}
async fn get_configured_relays(t: &TestContext) -> Vec<String> {
let transports = t.list_transports().await.unwrap();
let mut relays: Vec<_> = transports
.iter()
.map(|t| t.addr.split_once('@').unwrap().1)
.collect();
// Check that every relay is used only once:
relays.sort();
relays.dedup();
assert_eq!(relays.len(), transports.len());
relays.into_iter().map(|s| s.to_string()).collect()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_load_relay_candidates_single() -> Result<()> {
let t = &TestContext::new_alice().await;
enable_config(t).await;
let now = time();
t.sql.execute("DELETE FROM relay_candidates", ()).await?;
// This host should be returned by load_relay_candidates():
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
("never_tried.example", 0),
"INSERT INTO relay_candidates (host) VALUES (?)",
("never_tried.example",),
)
.await?;
// This host was recently tried and should not be returned:
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
("recent.example", now),
"INSERT INTO relay_candidates (host) VALUES (?)",
("recent.example",),
)
.await?;
set_relay_candidate_last_tried(t, "recent.example", now).await?;
// This host is already in use (alice@example.org) and should not be returned:
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
("example.org", 0),
"INSERT INTO relay_candidates (host) VALUES (?)",
("example.org",),
)
.await?;
let candidates = load_relay_candidates(t, now).await?;
assert_eq!(candidates, vec!["never_tried.example".to_string()]);
assert!(candidates.contains(&"never_tried.example".to_string()));
assert_eq!(candidates.contains(&"recent.example".to_string()), false);
assert_eq!(candidates.contains(&"example.org".to_string()), false);
assert_eq!(candidates.len(), DEFAULT_RELAY_CANDIDATES.len() + 1);
Ok(())
}
@@ -49,27 +86,22 @@ async fn test_load_relay_candidates_multiple() -> Result<()> {
enable_config(t).await;
let now = time();
t.sql.execute("DELETE FROM relay_candidates", ()).await?;
for host in ["a.example", "b.example", "c.example"] {
const EXAMPLE_CANDIDATES: &[&str] = &["a.example", "b.example", "c.example"];
for host in EXAMPLE_CANDIDATES {
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
(host, 0),
)
.execute("INSERT INTO relay_candidates (host) VALUES (?)", (host,))
.await?;
}
let mut candidates = load_relay_candidates(t, now).await?;
candidates.sort();
assert_eq!(
candidates,
vec![
"a.example".to_string(),
"b.example".to_string(),
"c.example".to_string()
]
);
let mut expected = EXAMPLE_CANDIDATES.to_vec();
expected.extend(DEFAULT_RELAY_CANDIDATES);
expected.sort();
assert_eq!(candidates, expected);
Ok(())
}
@@ -160,11 +192,21 @@ async fn test_maybe_add_additional_relays_add_one() -> Result<()> {
enable_config(t).await;
let now = time();
t.sql.execute("DELETE FROM relay_candidates", ()).await?;
// Make sure that default relay candidates
// are not used by setting last_used to now:
for candidate in DEFAULT_RELAY_CANDIDATES {
t.sql
.execute(
"INSERT INTO relay_candidates_last_tried(host, last_tried) VALUES(?,?)",
(candidate, now),
)
.await?;
}
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
("relay.example", 0),
"INSERT INTO relay_candidates (host) VALUES (?)",
("relay.example",),
)
.await?;
@@ -189,16 +231,6 @@ async fn test_maybe_add_additional_relays_add_multiple() -> Result<()> {
enable_config(t).await;
let now = time();
t.sql.execute("DELETE FROM relay_candidates", ()).await?;
for host in ["a.example", "b.example", "c.example", "d.example"] {
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
(host, 0),
)
.await?;
}
let skip_network = true;
let relay_added = maybe_add_additional_relays_inner(t, skip_network).await?;
assert!(relay_added);
@@ -218,12 +250,22 @@ async fn test_maybe_add_additional_relays_failure() -> Result<()> {
enable_config(t).await;
let now = time();
t.sql.execute("DELETE FROM relay_candidates", ()).await?;
// Make sure that default relay candidates
// are not used by setting last_used to now:
for candidate in DEFAULT_RELAY_CANDIDATES {
t.sql
.execute(
"INSERT INTO relay_candidates_last_tried(host, last_tried) VALUES(?,?)",
(candidate, now - 2),
)
.await?;
}
for i in 1..10 {
t.sql
.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)",
(format!("{i}.invalid.example"), 0),
"INSERT INTO relay_candidates (host) VALUES (?)",
(format!("{i}.invalid.example"),),
)
.await?;
}
@@ -246,7 +288,7 @@ async fn test_maybe_add_additional_relays_failure() -> Result<()> {
assert!(
t.sql
.exists(
"SELECT COUNT(*) FROM relay_candidates WHERE last_tried>=?",
"SELECT COUNT(*) FROM relay_candidates_last_tried WHERE last_tried>=?",
(now,)
)
.await?
+50 -1
View File
@@ -40,7 +40,7 @@ use crate::transport::{
ConnectionCandidate, delete_transport_row, maybe_update_sending_transport,
purge_transport_caches, send_sync_transports, transport_addrs,
};
use crate::{EventType, stock_str};
use crate::{EventType, autorelay, stock_str};
/// Maximum number of relays.
///
@@ -190,6 +190,55 @@ impl Context {
Ok(())
}
/// Automatically adds up to three transports.
///
/// If the user just scanned a QR code of type `Account`, `Login`,
/// `AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
/// then UI implementations should pass it as the `qr` parameter.
/// The host(s) from the QR code will then also be considered
/// for creating an account there.
pub async fn init_transports(&self, qr: Option<&str>) -> Result<()> {
if self.is_configured().await? {
bail!("Transports are already initialized");
}
let mut addrs_from_qr = vec![];
if let Some(qr) = qr {
match crate::qr::check_qr(self, qr).await? {
crate::qr::Qr::Account { .. } | crate::qr::Qr::Login { .. } => {
return self.add_transport_from_qr(qr).await;
}
crate::qr::Qr::AskVerifyContact { addrs, .. }
| crate::qr::Qr::AskVerifyGroup { addrs, .. }
| crate::qr::Qr::AskJoinBroadcast { addrs, .. } => addrs_from_qr = addrs,
_ => bail!("This QR code can't be used to initialize transports"),
}
}
self.stop_io().await;
let cancel_channel = self.alloc_ongoing().await?;
let skip_network = false;
let res = autorelay::init_transports_inner(self, addrs_from_qr, skip_network)
.race(cancel_channel.recv().map(|_| Err(format_err!("Canceled"))))
.await;
self.free_ongoing().await;
if let Err(err) = res {
let error_msg = stock_str::configuration_failed(self, &format!("{err:#}"));
progress!(self, 0, Some(error_msg.clone()));
bail!(error_msg);
}
progress!(self, 1000);
self.start_io().await;
Ok(())
}
/// Returns the list of all email accounts that are used as a transport in the current profile.
/// Use [Self::add_or_update_transport()] to add or change a transport
/// and [Self::delete_transport()] to delete a transport.
+1 -1
View File
@@ -231,7 +231,7 @@ pub struct InnerContext {
/// This is a global mutex-like state for operations which should be modal in the
/// clients.
running_state: RwLock<RunningState>,
/// Mutex to prevent running housekeeping or relay management from multiple threads at once.
/// Lock to prevent running housekeeping or relay management from multiple threads at once.
pub(crate) background_task_mutex: Mutex<()>,
/// Mutex to prevent multiple IMAP loops from fetching the messages at once.
+34
View File
@@ -684,6 +684,40 @@ impl Sql {
}
}
pub(crate) trait TransactionExt {
/// Prepares and executes the statement and maps a function over the resulting rows.
///
/// Collects the resulting rows into a collection.
fn query_map_collect<T, C, F>(
&self,
sql: &str,
params: impl rusqlite::Params + Send,
f: F,
) -> Result<C>
where
T: Send + 'static,
C: Send + 'static + std::iter::FromIterator<T>,
F: Send + FnMut(&rusqlite::Row) -> Result<T>;
}
impl TransactionExt for rusqlite::Transaction<'_> {
fn query_map_collect<T, C, F>(
&self,
sql: &str,
params: impl rusqlite::Params + Send,
f: F,
) -> Result<C>
where
T: Send + 'static,
C: Send + 'static + std::iter::FromIterator<T>,
F: Send + FnMut(&rusqlite::Row) -> Result<T>,
{
let mut stmt = self.prepare(sql)?;
let res = stmt.query_and_then(params, f)?;
res.collect()
}
}
/// Creates a new SQLite connection.
///
/// `path` is the database path.
+13
View File
@@ -2673,6 +2673,19 @@ CREATE TABLE smtp2 (
.await?;
}
inc_and_check(&mut migration_version, 167)?;
if dbversion < migration_version {
sql.execute_migration(
"DELETE FROM relay_candidates;
CREATE TABLE relay_candidates_last_tried(
host TEXT PRIMARY KEY NOT NULL,
last_tried INTEGER NOT NULL DEFAULT 0
) STRICT",
migration_version,
)
.await?;
}
let new_version = sql
.get_raw_config_int(VERSION_CFG)
.await?