use crate::dal::unified::DAL;
use crate::error::ValidationError;
#[cfg(feature = "postgres")]
use diesel::prelude::*;
#[cfg(feature = "postgres")]
#[derive(Queryable)]
#[diesel(table_name = crate::database::schema::postgres::agent_desired_counts)]
struct AgentDesiredRow {
#[allow(dead_code)]
pub tenant_id: String,
pub desired_count: i32,
#[allow(dead_code)]
pub updated_at: chrono::NaiveDateTime,
pub last_autoscaled_at: Option<chrono::NaiveDateTime>,
}
pub struct AgentDesiredDAL<'a> {
dal: &'a DAL,
}
impl<'a> AgentDesiredDAL<'a> {
pub fn new(dal: &'a DAL) -> Self {
Self { dal }
}
#[cfg(feature = "postgres")]
pub async fn get_desired(&self, tenant_id: &str) -> Result<u32, ValidationError> {
use crate::database::schema::postgres::agent_desired_counts as t;
let conn = self
.dal
.database
.get_postgres_connection()
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))?;
let tenant = tenant_id.to_string();
let row: Option<AgentDesiredRow> = conn
.interact(move |conn| t::table.find(tenant).first(conn).optional())
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))??;
Ok(row.map(|r| r.desired_count.max(0) as u32).unwrap_or(0))
}
#[cfg(feature = "postgres")]
pub async fn list_all(&self) -> Result<Vec<(String, u32)>, ValidationError> {
use crate::database::schema::postgres::agent_desired_counts as t;
let conn = self
.dal
.database
.get_postgres_connection()
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))?;
let rows: Vec<AgentDesiredRow> = conn
.interact(move |conn| t::table.load(conn))
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))??;
Ok(rows
.into_iter()
.map(|r| (r.tenant_id, r.desired_count.max(0) as u32))
.collect())
}
#[cfg(feature = "postgres")]
pub async fn list_all_with_last(
&self,
) -> Result<Vec<(String, u32, Option<chrono::NaiveDateTime>)>, ValidationError> {
use crate::database::schema::postgres::agent_desired_counts as t;
let conn = self
.dal
.database
.get_postgres_connection()
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))?;
let rows: Vec<AgentDesiredRow> = conn
.interact(move |conn| t::table.load(conn))
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))??;
Ok(rows
.into_iter()
.map(|r| {
(
r.tenant_id,
r.desired_count.max(0) as u32,
r.last_autoscaled_at,
)
})
.collect())
}
#[cfg(feature = "postgres")]
pub async fn set_desired(
&self,
tenant_id: &str,
desired_count: u32,
) -> Result<(), ValidationError> {
use crate::database::schema::postgres::agent_desired_counts as t;
let tenant = tenant_id.to_string();
let desired = desired_count as i32;
let conn = self
.dal
.database
.get_postgres_connection()
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))?;
conn.interact(move |conn| {
diesel::insert_into(t::table)
.values((t::tenant_id.eq(&tenant), t::desired_count.eq(desired)))
.on_conflict(t::tenant_id)
.do_update()
.set(t::desired_count.eq(desired))
.execute(conn)
})
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))??;
Ok(())
}
#[cfg(feature = "postgres")]
pub async fn set_desired_autoscaled(
&self,
tenant_id: &str,
desired_count: u32,
) -> Result<(), ValidationError> {
use crate::database::schema::postgres::agent_desired_counts as t;
let tenant = tenant_id.to_string();
let desired = desired_count as i32;
let conn = self
.dal
.database
.get_postgres_connection()
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))?;
conn.interact(move |conn| {
diesel::insert_into(t::table)
.values((
t::tenant_id.eq(&tenant),
t::desired_count.eq(desired),
t::last_autoscaled_at.eq(diesel::dsl::now),
))
.on_conflict(t::tenant_id)
.do_update()
.set((
t::desired_count.eq(desired),
t::last_autoscaled_at.eq(diesel::dsl::now),
))
.execute(conn)
})
.await
.map_err(|e| ValidationError::ConnectionPool(e.to_string()))??;
Ok(())
}
}