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 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 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 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}