mas_tasks/
lib.rs

1// Copyright 2024, 2025 New Vector Ltd.
2// Copyright 2021-2024 The Matrix.org Foundation C.I.C.
3//
4// SPDX-License-Identifier: AGPL-3.0-only
5// Please see LICENSE in the repository root for full details.
6
7use std::sync::{Arc, LazyLock};
8
9use mas_data_model::SiteConfig;
10use mas_email::Mailer;
11use mas_matrix::HomeserverConnection;
12use mas_router::UrlBuilder;
13use mas_storage::{BoxClock, BoxRepository, RepositoryError, SystemClock};
14use mas_storage_pg::PgRepository;
15use new_queue::QueueRunnerError;
16use opentelemetry::metrics::Meter;
17use rand::SeedableRng;
18use sqlx::{Pool, Postgres};
19use tokio_util::{sync::CancellationToken, task::TaskTracker};
20
21mod database;
22mod email;
23mod matrix;
24mod new_queue;
25mod recovery;
26mod sessions;
27mod user;
28
29static METER: LazyLock<Meter> = LazyLock::new(|| {
30    let scope = opentelemetry::InstrumentationScope::builder(env!("CARGO_PKG_NAME"))
31        .with_version(env!("CARGO_PKG_VERSION"))
32        .with_schema_url(opentelemetry_semantic_conventions::SCHEMA_URL)
33        .build();
34
35    opentelemetry::global::meter_with_scope(scope)
36});
37
38#[derive(Clone)]
39struct State {
40    pool: Pool<Postgres>,
41    mailer: Mailer,
42    clock: SystemClock,
43    homeserver: Arc<dyn HomeserverConnection>,
44    url_builder: UrlBuilder,
45    site_config: SiteConfig,
46}
47
48impl State {
49    pub fn new(
50        pool: Pool<Postgres>,
51        clock: SystemClock,
52        mailer: Mailer,
53        homeserver: impl HomeserverConnection + 'static,
54        url_builder: UrlBuilder,
55        site_config: SiteConfig,
56    ) -> Self {
57        Self {
58            pool,
59            mailer,
60            clock,
61            homeserver: Arc::new(homeserver),
62            url_builder,
63            site_config,
64        }
65    }
66
67    pub fn pool(&self) -> &Pool<Postgres> {
68        &self.pool
69    }
70
71    pub fn clock(&self) -> BoxClock {
72        Box::new(self.clock.clone())
73    }
74
75    pub fn mailer(&self) -> &Mailer {
76        &self.mailer
77    }
78
79    // This is fine for now, we may move that to a trait at some point.
80    #[allow(clippy::unused_self, clippy::disallowed_methods)]
81    pub fn rng(&self) -> rand_chacha::ChaChaRng {
82        rand_chacha::ChaChaRng::from_rng(rand::thread_rng()).expect("failed to seed rng")
83    }
84
85    pub async fn repository(&self) -> Result<BoxRepository, RepositoryError> {
86        let repo = PgRepository::from_pool(self.pool())
87            .await
88            .map_err(RepositoryError::from_error)?
89            .boxed();
90
91        Ok(repo)
92    }
93
94    pub fn matrix_connection(&self) -> &dyn HomeserverConnection {
95        self.homeserver.as_ref()
96    }
97
98    pub fn url_builder(&self) -> &UrlBuilder {
99        &self.url_builder
100    }
101
102    pub fn site_config(&self) -> &SiteConfig {
103        &self.site_config
104    }
105}
106
107/// Initialise the workers.
108///
109/// # Errors
110///
111/// This function can fail if the database connection fails.
112pub async fn init(
113    pool: &Pool<Postgres>,
114    mailer: &Mailer,
115    homeserver: impl HomeserverConnection + 'static,
116    url_builder: UrlBuilder,
117    site_config: &SiteConfig,
118    cancellation_token: CancellationToken,
119    task_tracker: &TaskTracker,
120) -> Result<(), QueueRunnerError> {
121    let state = State::new(
122        pool.clone(),
123        SystemClock::default(),
124        mailer.clone(),
125        homeserver,
126        url_builder,
127        site_config.clone(),
128    );
129    let mut worker = self::new_queue::QueueWorker::new(state, cancellation_token).await?;
130
131    worker
132        .register_handler::<mas_storage::queue::CleanupExpiredTokensJob>()
133        .register_handler::<mas_storage::queue::DeactivateUserJob>()
134        .register_handler::<mas_storage::queue::DeleteDeviceJob>()
135        .register_handler::<mas_storage::queue::ProvisionDeviceJob>()
136        .register_handler::<mas_storage::queue::ProvisionUserJob>()
137        .register_handler::<mas_storage::queue::ReactivateUserJob>()
138        .register_handler::<mas_storage::queue::SendAccountRecoveryEmailsJob>()
139        .register_handler::<mas_storage::queue::SendEmailAuthenticationCodeJob>()
140        .register_handler::<mas_storage::queue::SyncDevicesJob>()
141        .register_handler::<mas_storage::queue::VerifyEmailJob>()
142        .register_handler::<mas_storage::queue::ExpireInactiveSessionsJob>()
143        .register_handler::<mas_storage::queue::ExpireInactiveCompatSessionsJob>()
144        .register_handler::<mas_storage::queue::ExpireInactiveOAuthSessionsJob>()
145        .register_handler::<mas_storage::queue::ExpireInactiveUserSessionsJob>()
146        .register_handler::<mas_storage::queue::PruneStalePolicyDataJob>()
147        .add_schedule(
148            "cleanup-expired-tokens",
149            "0 0 * * * *".parse()?,
150            mas_storage::queue::CleanupExpiredTokensJob,
151        )
152        .add_schedule(
153            "expire-inactive-sessions",
154            // Run this job every 15 minutes
155            "30 */15 * * * *".parse()?,
156            mas_storage::queue::ExpireInactiveSessionsJob,
157        )
158        .add_schedule(
159            "prune-stale-policy-data",
160            // Run once a day
161            "0 0 2 * * *".parse()?,
162            mas_storage::queue::PruneStalePolicyDataJob,
163        );
164
165    task_tracker.spawn(worker.run());
166
167    Ok(())
168}