Skip to main content

shared/models/
database_agent_host.rs

1use crate::{
2    models::{InsertQueryBuilder, UpdateQueryBuilder},
3    prelude::*,
4};
5use garde::Validate;
6use rand::distr::SampleString;
7use serde::{Deserialize, Serialize};
8use sqlx::{Row, postgres::PgRow};
9use std::{
10    collections::BTreeMap,
11    sync::{Arc, LazyLock},
12};
13use utoipa::ToSchema;
14
15#[inline]
16fn default_true() -> bool {
17    true
18}
19
20#[derive(ToSchema, Serialize, Deserialize, Validate, Clone)]
21pub struct DatabaseAgentHostTypeSettings {
22    #[garde(skip)]
23    #[serde(default = "default_true")]
24    pub enabled: bool,
25
26    #[garde(
27        length(chars, min = 3, max = 255),
28        inner(custom(crate::utils::validate_host))
29    )]
30    #[schema(min_length = 3, max_length = 255)]
31    #[serde(default)]
32    pub public_host: Option<compact_str::CompactString>,
33    #[garde(range(min = 1))]
34    #[schema(minimum = 1)]
35    #[serde(default)]
36    pub public_port: Option<u16>,
37}
38
39impl Default for DatabaseAgentHostTypeSettings {
40    fn default() -> Self {
41        Self {
42            enabled: true,
43            public_host: None,
44            public_port: None,
45        }
46    }
47}
48
49#[derive(ToSchema, Serialize, Deserialize, Validate, Clone, Default)]
50pub struct DatabaseAgentHostTypes {
51    #[garde(dive)]
52    #[serde(default)]
53    pub postgres: DatabaseAgentHostTypeSettings,
54    #[garde(dive)]
55    #[serde(default)]
56    pub mariadb: DatabaseAgentHostTypeSettings,
57    #[garde(dive)]
58    #[serde(default)]
59    pub mongodb: DatabaseAgentHostTypeSettings,
60    #[garde(dive)]
61    #[serde(default)]
62    pub redis: DatabaseAgentHostTypeSettings,
63}
64
65impl DatabaseAgentHostTypes {
66    #[inline]
67    pub fn get(&self, r#type: db_agent_api::DatabaseAgentType) -> &DatabaseAgentHostTypeSettings {
68        match r#type {
69            db_agent_api::DatabaseAgentType::Postgres => &self.postgres,
70            db_agent_api::DatabaseAgentType::Mariadb => &self.mariadb,
71            db_agent_api::DatabaseAgentType::Mongodb => &self.mongodb,
72            db_agent_api::DatabaseAgentType::Redis => &self.redis,
73        }
74    }
75}
76
77#[derive(Serialize, Deserialize, Clone)]
78pub struct DatabaseAgentHost {
79    pub uuid: uuid::Uuid,
80
81    pub name: compact_str::CompactString,
82    pub description: Option<compact_str::CompactString>,
83
84    pub deployment_enabled: bool,
85    pub maintenance_enabled: bool,
86
87    pub url: reqwest::Url,
88
89    pub memory: i64,
90    pub disk: i64,
91
92    pub types: DatabaseAgentHostTypes,
93
94    pub token: Vec<u8>,
95
96    pub created: chrono::NaiveDateTime,
97
98    extension_data: super::ModelExtensionData,
99}
100
101impl BaseModel for DatabaseAgentHost {
102    const NAME: &'static str = "database_agent_host";
103
104    fn get_extension_list() -> &'static super::ModelExtensionList {
105        static EXTENSIONS: LazyLock<super::ModelExtensionList> =
106            LazyLock::new(|| parking_lot::RwLock::new(Vec::new()));
107
108        &EXTENSIONS
109    }
110
111    fn get_extension_data(&self) -> &super::ModelExtensionData {
112        &self.extension_data
113    }
114
115    #[inline]
116    fn base_columns(prefix: Option<&str>) -> BTreeMap<&'static str, compact_str::CompactString> {
117        let prefix = prefix.unwrap_or_default();
118
119        BTreeMap::from([
120            (
121                "database_agent_hosts.uuid",
122                compact_str::format_compact!("{prefix}uuid"),
123            ),
124            (
125                "database_agent_hosts.name",
126                compact_str::format_compact!("{prefix}name"),
127            ),
128            (
129                "database_agent_hosts.description",
130                compact_str::format_compact!("{prefix}description"),
131            ),
132            (
133                "database_agent_hosts.deployment_enabled",
134                compact_str::format_compact!("{prefix}deployment_enabled"),
135            ),
136            (
137                "database_agent_hosts.maintenance_enabled",
138                compact_str::format_compact!("{prefix}maintenance_enabled"),
139            ),
140            (
141                "database_agent_hosts.url",
142                compact_str::format_compact!("{prefix}url"),
143            ),
144            (
145                "database_agent_hosts.memory",
146                compact_str::format_compact!("{prefix}memory"),
147            ),
148            (
149                "database_agent_hosts.disk",
150                compact_str::format_compact!("{prefix}disk"),
151            ),
152            (
153                "database_agent_hosts.types",
154                compact_str::format_compact!("{prefix}types"),
155            ),
156            (
157                "database_agent_hosts.token",
158                compact_str::format_compact!("{prefix}token"),
159            ),
160            (
161                "database_agent_hosts.created",
162                compact_str::format_compact!("{prefix}created"),
163            ),
164        ])
165    }
166
167    #[inline]
168    fn map(prefix: Option<&str>, row: &PgRow) -> Result<Self, crate::database::DatabaseError> {
169        let prefix = prefix.unwrap_or_default();
170
171        Ok(Self {
172            uuid: row.try_get(compact_str::format_compact!("{prefix}uuid").as_str())?,
173            name: row.try_get(compact_str::format_compact!("{prefix}name").as_str())?,
174            description: row
175                .try_get(compact_str::format_compact!("{prefix}description").as_str())?,
176            deployment_enabled: row
177                .try_get(compact_str::format_compact!("{prefix}deployment_enabled").as_str())?,
178            maintenance_enabled: row
179                .try_get(compact_str::format_compact!("{prefix}maintenance_enabled").as_str())?,
180            url: row
181                .try_get::<String, _>(compact_str::format_compact!("{prefix}url").as_str())?
182                .parse()
183                .map_err(anyhow::Error::new)?,
184            memory: row.try_get(compact_str::format_compact!("{prefix}memory").as_str())?,
185            disk: row.try_get(compact_str::format_compact!("{prefix}disk").as_str())?,
186            types: serde_json::from_value(
187                row.try_get(compact_str::format_compact!("{prefix}types").as_str())?,
188            )?,
189            token: row.try_get(compact_str::format_compact!("{prefix}token").as_str())?,
190            created: row.try_get(compact_str::format_compact!("{prefix}created").as_str())?,
191            extension_data: Self::map_extensions(prefix, row)?,
192        })
193    }
194}
195
196impl DatabaseAgentHost {
197    pub async fn all_with_pagination(
198        database: &crate::database::Database,
199        page: i64,
200        per_page: i64,
201        search: Option<&str>,
202    ) -> Result<super::Pagination<Self>, crate::database::DatabaseError> {
203        let offset = (page - 1) * per_page;
204
205        let rows = sqlx::query(sqlx::AssertSqlSafe(format!(
206            r#"
207            SELECT {}, COUNT(*) OVER() AS total_count
208            FROM database_agent_hosts
209            WHERE ($1 IS NULL OR database_agent_hosts.name ILIKE '%' || $1 || '%')
210            ORDER BY database_agent_hosts.created
211            LIMIT $2 OFFSET $3
212            "#,
213            Self::columns_sql(None)
214        )))
215        .bind(search)
216        .bind(per_page)
217        .bind(offset)
218        .fetch_all(database.read())
219        .await?;
220
221        Ok(super::Pagination {
222            total: rows
223                .first()
224                .map_or(Ok(0), |row| row.try_get("total_count"))?,
225            per_page,
226            page,
227            data: rows
228                .into_iter()
229                .map(|row| Self::map(None, &row))
230                .try_collect_vec()?,
231        })
232    }
233
234    pub async fn by_node_most_eligible(
235        database: &crate::database::Database,
236        node: &super::node::Node,
237        r#type: db_agent_api::DatabaseAgentType,
238        memory: i64,
239        disk: i64,
240    ) -> Result<Vec<Self>, crate::database::DatabaseError> {
241        let rows = sqlx::query(sqlx::AssertSqlSafe(format!(
242            r#"
243            WITH database_usage AS (
244                SELECT
245                    server_database_instances.database_agent_host_uuid,
246                    COALESCE(SUM(COALESCE(server_database_instances.memory, t.memory)), 0)::BIGINT AS used_memory,
247                    COALESCE(SUM(COALESCE(server_database_instances.disk, t.disk)), 0)::BIGINT AS used_disk
248                FROM server_database_instances
249                LEFT JOIN database_agent_templates t ON t.uuid = server_database_instances.database_agent_template_uuid
250                GROUP BY server_database_instances.database_agent_host_uuid
251            )
252            SELECT {}
253            FROM database_agent_hosts
254            LEFT JOIN database_usage u ON database_agent_hosts.uuid = u.database_agent_host_uuid
255            WHERE (
256                EXISTS (
257                    SELECT 1 FROM node_database_agent_hosts
258                    WHERE node_database_agent_hosts.database_agent_host_uuid = database_agent_hosts.uuid AND node_database_agent_hosts.node_uuid = $1
259                )
260                OR EXISTS (
261                    SELECT 1 FROM location_database_agent_hosts
262                    WHERE location_database_agent_hosts.database_agent_host_uuid = database_agent_hosts.uuid AND location_database_agent_hosts.location_uuid = $2
263                )
264            )
265            AND database_agent_hosts.deployment_enabled
266            AND NOT database_agent_hosts.maintenance_enabled
267            AND COALESCE((database_agent_hosts.types -> $3 ->> 'enabled')::BOOL, TRUE)
268            AND COALESCE(u.used_memory, 0) + $4 <= database_agent_hosts.memory
269            AND COALESCE(u.used_disk, 0) + $5 <= database_agent_hosts.disk
270            ORDER BY
271                GREATEST(
272                    (COALESCE(u.used_memory, 0) + $4)::FLOAT / NULLIF(database_agent_hosts.memory, 0),
273                    (COALESCE(u.used_disk, 0) + $5)::FLOAT / NULLIF(database_agent_hosts.disk, 0)
274                )
275            "#,
276            Self::columns_sql(None)
277        )))
278        .bind(node.uuid)
279        .bind(node.location.uuid)
280        .bind(r#type.as_str())
281        .bind(memory)
282        .bind(disk)
283        .fetch_all(database.read())
284        .await?;
285
286        rows.into_iter()
287            .map(|row| Self::map(None, &row))
288            .try_collect_vec()
289    }
290
291    #[inline]
292    pub fn generate_token() -> String {
293        rand::distr::Alphanumeric.sample_string(&mut rand::rng(), 64)
294    }
295
296    pub async fn reset_token(&self, state: &crate::State) -> Result<String, anyhow::Error> {
297        let token = Self::generate_token();
298
299        sqlx::query(
300            r#"
301            UPDATE database_agent_hosts
302            SET token = $2
303            WHERE database_agent_hosts.uuid = $1
304            "#,
305        )
306        .bind(self.uuid)
307        .bind(state.database.encrypt(token.clone()).await?)
308        .execute(state.database.write())
309        .await?;
310
311        Ok(token)
312    }
313
314    #[inline]
315    pub async fn api_client(
316        &self,
317        database: &crate::database::Database,
318    ) -> Result<db_agent_api::client::DbAgentClient, anyhow::Error> {
319        Ok(db_agent_api::client::DbAgentClient::new(
320            self.url.to_string(),
321            database.decrypt(self.token.to_vec()).await?.into(),
322        ))
323    }
324
325    /// Fetch the current configuration of this database agent host
326    ///
327    /// Cached for 120 seconds.
328    pub async fn fetch_configuration(
329        &self,
330        database: &crate::database::Database,
331    ) -> Result<db_agent_api::system_config::get::Response, anyhow::Error> {
332        database
333            .cache
334            .cached(
335                &format!("database_agent_host::{}::configuration", self.uuid),
336                120,
337                || async {
338                    Ok::<_, anyhow::Error>(
339                        self.api_client(database).await?.get_system_config().await?,
340                    )
341                },
342            )
343            .await
344    }
345
346    /// Fetch the current resource usages of all databases on this host.
347    ///
348    /// Cached for 15 seconds.
349    pub async fn fetch_database_resources(
350        &self,
351        database: &crate::database::Database,
352    ) -> Result<std::collections::HashMap<uuid::Uuid, db_agent_api::ResourceUsage>, anyhow::Error>
353    {
354        database
355            .cache
356            .cached(
357                &format!("database_agent_host::{}::database_resources", self.uuid),
358                15,
359                || async {
360                    let resources = self
361                        .api_client(database)
362                        .await?
363                        .get_instances_utilization()
364                        .await?;
365
366                    Ok::<_, anyhow::Error>(resources.into_iter().collect())
367                },
368            )
369            .await
370    }
371
372    /// Update the configuration of this database agent host
373    ///
374    /// Invalidates the cached configuration.
375    pub async fn update_configuration(
376        &self,
377        database: &crate::database::Database,
378        config_patch: &serde_json::Value,
379    ) -> Result<bool, anyhow::Error> {
380        let response = self
381            .api_client(database)
382            .await?
383            .patch_system_config(config_patch)
384            .await?;
385        if !response.applied {
386            return Ok(false);
387        }
388
389        database
390            .cache
391            .invalidate(&format!(
392                "database_agent_host::{}::configuration",
393                self.uuid
394            ))
395            .await?;
396
397        Ok(true)
398    }
399}
400
401#[async_trait::async_trait]
402impl IntoAdminApiObject for DatabaseAgentHost {
403    type AdminApiObject = AdminApiDatabaseAgentHost;
404    type ExtraArgs<'a> = ();
405
406    async fn into_admin_api_object<'a>(
407        self,
408        state: &crate::State,
409        _args: Self::ExtraArgs<'a>,
410    ) -> Result<Self::AdminApiObject, crate::database::DatabaseError> {
411        let api_object = AdminApiDatabaseAgentHost::init_hooks(&self, state).await?;
412
413        let api_object = finish_extendible!(
414            AdminApiDatabaseAgentHost {
415                uuid: self.uuid,
416                name: self.name,
417                description: self.description,
418                deployment_enabled: self.deployment_enabled,
419                maintenance_enabled: self.maintenance_enabled,
420                url: self.url.to_string(),
421                memory: self.memory,
422                disk: self.disk,
423                types: self.types,
424                created: self.created.and_utc(),
425            },
426            api_object,
427            state
428        )?;
429
430        Ok(api_object)
431    }
432}
433
434#[async_trait::async_trait]
435impl ByUuid for DatabaseAgentHost {
436    async fn by_uuid(
437        database: &crate::database::Database,
438        uuid: uuid::Uuid,
439    ) -> Result<Self, crate::database::DatabaseError> {
440        let row = sqlx::query(sqlx::AssertSqlSafe(format!(
441            r#"
442            SELECT {}
443            FROM database_agent_hosts
444            WHERE database_agent_hosts.uuid = $1
445            "#,
446            Self::columns_sql(None)
447        )))
448        .bind(uuid)
449        .fetch_one(database.read())
450        .await?;
451
452        Self::map(None, &row)
453    }
454
455    async fn by_uuid_with_transaction(
456        transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
457        uuid: uuid::Uuid,
458    ) -> Result<Self, crate::database::DatabaseError> {
459        let row = sqlx::query(sqlx::AssertSqlSafe(format!(
460            r#"
461            SELECT {}
462            FROM database_agent_hosts
463            WHERE database_agent_hosts.uuid = $1
464            "#,
465            Self::columns_sql(None)
466        )))
467        .bind(uuid)
468        .fetch_one(&mut **transaction)
469        .await?;
470
471        Self::map(None, &row)
472    }
473}
474
475#[derive(ToSchema, Deserialize, Validate)]
476pub struct CreateDatabaseAgentHostOptions {
477    #[garde(length(chars, min = 1, max = 255))]
478    #[schema(min_length = 1, max_length = 255)]
479    pub name: compact_str::CompactString,
480    #[garde(length(chars, min = 1, max = 1024))]
481    #[schema(min_length = 1, max_length = 1024)]
482    pub description: Option<compact_str::CompactString>,
483
484    #[garde(skip)]
485    pub deployment_enabled: bool,
486    #[garde(skip)]
487    pub maintenance_enabled: bool,
488
489    #[garde(length(chars, min = 3, max = 255), url)]
490    #[schema(min_length = 3, max_length = 255, format = "uri")]
491    pub url: compact_str::CompactString,
492
493    #[garde(range(min = 1))]
494    #[schema(minimum = 1)]
495    pub memory: i64,
496    #[garde(range(min = 1))]
497    #[schema(minimum = 1)]
498    pub disk: i64,
499
500    #[garde(dive)]
501    #[serde(default)]
502    pub types: DatabaseAgentHostTypes,
503}
504
505#[async_trait::async_trait]
506impl CreatableModel for DatabaseAgentHost {
507    type CreateOptions<'a> = CreateDatabaseAgentHostOptions;
508    type CreateResult = Self;
509
510    fn get_create_handlers() -> &'static LazyLock<CreateListenerList<Self>> {
511        static CREATE_LISTENERS: LazyLock<CreateListenerList<DatabaseAgentHost>> =
512            LazyLock::new(|| Arc::new(ModelHandlerList::default()));
513
514        &CREATE_LISTENERS
515    }
516
517    async fn create_with_transaction(
518        state: &crate::State,
519        mut options: Self::CreateOptions<'_>,
520        transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
521    ) -> Result<Self, crate::database::DatabaseError> {
522        options.validate()?;
523
524        let mut query_builder = InsertQueryBuilder::new("database_agent_hosts");
525
526        Self::run_create_handlers(&mut options, &mut query_builder, state, transaction).await?;
527
528        let token = Self::generate_token();
529
530        query_builder
531            .set("name", &options.name)
532            .set("description", &options.description)
533            .set("deployment_enabled", options.deployment_enabled)
534            .set("maintenance_enabled", options.maintenance_enabled)
535            .set("url", &options.url)
536            .set("memory", options.memory)
537            .set("disk", options.disk)
538            .set("types", serde_json::to_value(&options.types)?)
539            .set("token", state.database.encrypt(token.clone()).await?);
540
541        let row = query_builder
542            .returning(&Self::columns_sql(None))
543            .fetch_one(&mut **transaction)
544            .await?;
545        let mut database_agent_host = Self::map(None, &row)?;
546
547        Self::run_after_create_handlers(&mut database_agent_host, &options, state, transaction)
548            .await?;
549
550        Ok(database_agent_host)
551    }
552}
553
554#[derive(ToSchema, Serialize, Deserialize, Validate, Clone, Default)]
555pub struct UpdateDatabaseAgentHostOptions {
556    #[garde(length(chars, min = 1, max = 255))]
557    #[schema(min_length = 1, max_length = 255)]
558    pub name: Option<compact_str::CompactString>,
559    #[garde(length(chars, min = 1, max = 1024))]
560    #[schema(min_length = 1, max_length = 1024)]
561    #[serde(
562        default,
563        skip_serializing_if = "Option::is_none",
564        with = "::serde_with::rust::double_option"
565    )]
566    pub description: Option<Option<compact_str::CompactString>>,
567
568    #[garde(skip)]
569    pub deployment_enabled: Option<bool>,
570    #[garde(skip)]
571    pub maintenance_enabled: Option<bool>,
572
573    #[garde(length(chars, min = 3, max = 255), url)]
574    #[schema(min_length = 3, max_length = 255, format = "uri")]
575    pub url: Option<compact_str::CompactString>,
576
577    #[garde(range(min = 1))]
578    #[schema(minimum = 1)]
579    pub memory: Option<i64>,
580    #[garde(range(min = 1))]
581    #[schema(minimum = 1)]
582    pub disk: Option<i64>,
583
584    #[garde(dive)]
585    pub types: Option<DatabaseAgentHostTypes>,
586}
587
588#[async_trait::async_trait]
589impl UpdatableModel for DatabaseAgentHost {
590    type UpdateOptions = UpdateDatabaseAgentHostOptions;
591
592    fn get_update_handlers() -> &'static LazyLock<UpdateHandlerList<Self>> {
593        static UPDATE_LISTENERS: LazyLock<UpdateHandlerList<DatabaseAgentHost>> =
594            LazyLock::new(|| Arc::new(ModelHandlerList::default()));
595
596        &UPDATE_LISTENERS
597    }
598
599    async fn update_with_transaction(
600        &mut self,
601        state: &crate::State,
602        mut options: Self::UpdateOptions,
603        transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
604    ) -> Result<(), crate::database::DatabaseError> {
605        options.validate()?;
606
607        let mut query_builder = UpdateQueryBuilder::new("database_agent_hosts");
608
609        self.run_update_handlers(&mut options, &mut query_builder, state, transaction)
610            .await?;
611
612        query_builder
613            .set("name", options.name.as_ref())
614            .set(
615                "description",
616                options.description.as_ref().map(|d| d.as_ref()),
617            )
618            .set("deployment_enabled", options.deployment_enabled)
619            .set("maintenance_enabled", options.maintenance_enabled)
620            .set("url", options.url.as_ref())
621            .set("memory", options.memory)
622            .set("disk", options.disk)
623            .set(
624                "types",
625                options
626                    .types
627                    .as_ref()
628                    .map(serde_json::to_value)
629                    .transpose()?,
630            )
631            .where_eq("uuid", self.uuid);
632
633        query_builder.execute(&mut **transaction).await?;
634
635        if let Some(name) = options.name {
636            self.name = name;
637        }
638        if let Some(description) = options.description {
639            self.description = description;
640        }
641        if let Some(deployment_enabled) = options.deployment_enabled {
642            self.deployment_enabled = deployment_enabled;
643        }
644        if let Some(maintenance_enabled) = options.maintenance_enabled {
645            self.maintenance_enabled = maintenance_enabled;
646        }
647        if let Some(url) = options.url {
648            self.url = url.parse().map_err(anyhow::Error::new)?;
649        }
650        if let Some(memory) = options.memory {
651            self.memory = memory;
652        }
653        if let Some(disk) = options.disk {
654            self.disk = disk;
655        }
656        if let Some(types) = options.types {
657            self.types = types;
658        }
659
660        self.run_after_update_handlers(state, transaction).await?;
661
662        Ok(())
663    }
664}
665
666#[async_trait::async_trait]
667impl DeletableModel for DatabaseAgentHost {
668    type DeleteOptions = ();
669
670    fn get_delete_handlers() -> &'static LazyLock<DeleteHandlerList<Self>> {
671        static DELETE_LISTENERS: LazyLock<DeleteHandlerList<DatabaseAgentHost>> =
672            LazyLock::new(|| Arc::new(ModelHandlerList::default()));
673
674        &DELETE_LISTENERS
675    }
676
677    async fn delete_with_transaction(
678        &self,
679        state: &crate::State,
680        options: Self::DeleteOptions,
681        transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
682    ) -> Result<(), anyhow::Error> {
683        self.run_delete_handlers(&options, state, transaction)
684            .await?;
685
686        sqlx::query(
687            r#"
688            DELETE FROM database_agent_hosts
689            WHERE database_agent_hosts.uuid = $1
690            "#,
691        )
692        .bind(self.uuid)
693        .execute(&mut **transaction)
694        .await?;
695
696        self.run_after_delete_handlers(&options, state, transaction)
697            .await?;
698
699        Ok(())
700    }
701}
702
703#[schema_extension_derive::extendible]
704#[init_args(DatabaseAgentHost, crate::State)]
705#[hook_args(crate::State)]
706#[derive(ToSchema, Serialize)]
707#[schema(title = "DatabaseAgentHost")]
708pub struct AdminApiDatabaseAgentHost {
709    pub uuid: uuid::Uuid,
710
711    pub name: compact_str::CompactString,
712    pub description: Option<compact_str::CompactString>,
713
714    pub deployment_enabled: bool,
715    pub maintenance_enabled: bool,
716
717    #[schema(format = "uri")]
718    pub url: String,
719
720    pub memory: i64,
721    pub disk: i64,
722
723    pub types: DatabaseAgentHostTypes,
724
725    pub created: chrono::DateTime<chrono::Utc>,
726}