1use 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 #[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
107pub 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 "30 */15 * * * *".parse()?,
156 mas_storage::queue::ExpireInactiveSessionsJob,
157 )
158 .add_schedule(
159 "prune-stale-policy-data",
160 "0 0 2 * * *".parse()?,
162 mas_storage::queue::PruneStalePolicyDataJob,
163 );
164
165 task_tracker.spawn(worker.run());
166
167 Ok(())
168}