pub mod models;
mod operations;
pub(in crate) mod schema;
use diesel::r2d2::{ConnectionManager, Pool};
use super::diesel::models::{AgentModel, NewAgentModel, NewRoleModel, RoleModel};
use super::{Agent, AgentStore, AgentStoreError, Role};
use crate::commits::MAX_COMMIT_NUM;
use crate::error::{
ConstraintViolationError, ConstraintViolationType, InternalError,
ResourceTemporarilyUnavailableError,
};
use operations::add_agent::AgentStoreAddAgentOperation as _;
use operations::fetch_agent::AgentStoreFetchAgentOperation as _;
use operations::list_agents::AgentStoreListAgentsOperation as _;
use operations::update_agent::AgentStoreUpdateAgentOperation as _;
use operations::AgentStoreOperations;
#[derive(Clone)]
pub struct DieselAgentStore<C: diesel::Connection + 'static> {
connection_pool: Pool<ConnectionManager<C>>,
}
impl<C: diesel::Connection> DieselAgentStore<C> {
#[allow(dead_code)]
pub fn new(connection_pool: Pool<ConnectionManager<C>>) -> Self {
DieselAgentStore { connection_pool }
}
}
#[cfg(feature = "postgres")]
impl AgentStore for DieselAgentStore<diesel::pg::PgConnection> {
fn add_agent(&self, agent: Agent) -> Result<(), AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.add_agent(agent.clone().into(), make_role_models(&agent))
}
fn list_agents(&self, service_id: Option<&str>) -> Result<Vec<Agent>, AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.list_agents(service_id)
}
fn fetch_agent(
&self,
pub_key: &str,
service_id: Option<&str>,
) -> Result<Option<Agent>, AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.fetch_agent(pub_key, service_id)
}
fn update_agent(&self, agent: Agent) -> Result<(), AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.update_agent(agent.clone().into(), make_role_models(&agent))
}
}
#[cfg(feature = "sqlite")]
impl AgentStore for DieselAgentStore<diesel::sqlite::SqliteConnection> {
fn add_agent(&self, agent: Agent) -> Result<(), AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.add_agent(agent.clone().into(), make_role_models(&agent))
}
fn list_agents(&self, service_id: Option<&str>) -> Result<Vec<Agent>, AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.list_agents(service_id)
}
fn fetch_agent(
&self,
pub_key: &str,
service_id: Option<&str>,
) -> Result<Option<Agent>, AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.fetch_agent(pub_key, service_id)
}
fn update_agent(&self, agent: Agent) -> Result<(), AgentStoreError> {
AgentStoreOperations::new(&*self.connection_pool.get().map_err(|err| {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
})?)
.update_agent(agent.clone().into(), make_role_models(&agent))
}
}
impl From<RoleModel> for Role {
fn from(role: RoleModel) -> Self {
Self {
public_key: role.public_key,
role_name: role.role_name,
start_commit_num: role.start_commit_num,
end_commit_num: role.end_commit_num,
service_id: role.service_id,
}
}
}
impl From<(AgentModel, Vec<RoleModel>)> for Agent {
fn from((agent_model, role_models): (AgentModel, Vec<RoleModel>)) -> Self {
Self {
public_key: agent_model.public_key,
org_id: agent_model.org_id,
active: agent_model.active,
metadata: agent_model.metadata,
roles: role_models
.iter()
.map(|role| role.role_name.to_string())
.collect(),
start_commit_num: agent_model.start_commit_num,
end_commit_num: agent_model.end_commit_num,
service_id: agent_model.service_id,
}
}
}
impl From<Agent> for NewAgentModel {
fn from(agent: Agent) -> NewAgentModel {
Self {
public_key: agent.public_key,
org_id: agent.org_id,
active: agent.active,
metadata: agent.metadata,
start_commit_num: agent.start_commit_num,
end_commit_num: MAX_COMMIT_NUM,
service_id: agent.service_id,
}
}
}
pub fn make_role_models(agent: &Agent) -> Vec<NewRoleModel> {
let mut roles = Vec::new();
for role in &agent.roles {
roles.push(NewRoleModel {
public_key: agent.public_key.to_string(),
role_name: role.to_string(),
start_commit_num: agent.start_commit_num,
end_commit_num: agent.end_commit_num,
service_id: agent.service_id.clone(),
})
}
roles
}
impl From<diesel::result::Error> for AgentStoreError {
fn from(err: diesel::result::Error) -> AgentStoreError {
match err {
diesel::result::Error::DatabaseError(
diesel::result::DatabaseErrorKind::UniqueViolation,
_,
) => AgentStoreError::ConstraintViolationError(
ConstraintViolationError::from_source_with_violation_type(
ConstraintViolationType::Unique,
Box::new(err),
),
),
diesel::result::Error::DatabaseError(
diesel::result::DatabaseErrorKind::ForeignKeyViolation,
_,
) => AgentStoreError::ConstraintViolationError(
ConstraintViolationError::from_source_with_violation_type(
ConstraintViolationType::ForeignKey,
Box::new(err),
),
),
_ => AgentStoreError::InternalError(InternalError::from_source(Box::new(err))),
}
}
}
impl From<diesel::r2d2::PoolError> for AgentStoreError {
fn from(err: diesel::r2d2::PoolError) -> AgentStoreError {
AgentStoreError::ResourceTemporarilyUnavailableError(
ResourceTemporarilyUnavailableError::from_source(Box::new(err)),
)
}
}