mod database;
pub mod filesystem;
mod package;
pub use database::{build_queue_stats, reconciler_stats, BuildQueueStats, ReconcilerStats};
use async_trait::async_trait;
use std::collections::HashMap;
use std::path::Path;
use uuid::Uuid;
use crate::database::Database;
use crate::registry::error::RegistryError;
use crate::registry::loader::{PackageLoader, TaskRegistrar};
use crate::registry::traits::{RegistryStorage, WorkflowRegistry};
use crate::registry::types::{
LoadedWorkflow, WorkflowMetadata, WorkflowPackageId, WorkflowSourceFile,
};
use crate::task::TaskNamespace;
const MAX_SOURCE_FILE_BYTES: u64 = 1024 * 1024;
pub struct WorkflowRegistryImpl<S: RegistryStorage> {
pub(super) storage: S,
pub(super) database: Database,
#[allow(dead_code)]
loader: PackageLoader,
registrar: TaskRegistrar,
pub(super) loaded_packages: HashMap<Uuid, Vec<TaskNamespace>>,
}
impl<S: RegistryStorage> WorkflowRegistryImpl<S> {
pub fn new(storage: S, database: Database) -> Result<Self, RegistryError> {
let loader = PackageLoader::new().map_err(RegistryError::Loader)?;
let registrar = TaskRegistrar::new().map_err(RegistryError::Loader)?;
Ok(Self {
storage,
database,
loader,
registrar,
loaded_packages: HashMap::new(),
})
}
pub fn loaded_package_count(&self) -> usize {
self.loaded_packages.len()
}
pub fn total_registered_tasks(&self) -> usize {
self.loaded_packages.values().map(|tasks| tasks.len()).sum()
}
pub async fn register_workflow_package(
&mut self,
package_data: Vec<u8>,
) -> Result<Uuid, RegistryError> {
WorkflowRegistry::register_workflow(self, package_data).await
}
pub async fn get_source_for_build(
&self,
package_id: Uuid,
) -> Result<Option<(WorkflowMetadata, Vec<u8>)>, RegistryError> {
let ins = match self.inspect_package_by_id(package_id).await? {
Some(ins) => ins,
None => return Ok(None),
};
let registry_id = ins.metadata.registry_id.to_string();
let package_data = match self.storage.retrieve_binary(®istry_id).await? {
Some(data) => data,
None => {
return Err(RegistryError::Internal(
"Package metadata exists but binary data is missing".to_string(),
));
}
};
Ok(Some((ins.metadata, package_data)))
}
pub async fn get_workflow_source(
&self,
package_id: Uuid,
) -> Result<Option<(WorkflowMetadata, Vec<WorkflowSourceFile>)>, RegistryError> {
let (metadata, archive_bytes) = match self.get_source_for_build(package_id).await? {
Some(pair) => pair,
None => return Ok(None),
};
let files = tokio::task::spawn_blocking(move || extract_source_files(&archive_bytes))
.await
.map_err(|e| {
RegistryError::Internal(format!("source extraction task panicked: {}", e))
})??;
Ok(Some((metadata, files)))
}
pub async fn is_workflow_paused(&self, name: &str) -> Result<bool, RegistryError> {
let workflows = self.list_workflows().await?;
Ok(workflows
.into_iter()
.find(|w| w.workflow_name == name || w.package_name == name)
.map(|w| w.paused)
.unwrap_or(false))
}
pub async fn get_workflow_declared_params(
&self,
name: &str,
) -> Result<Vec<cloacina_api_types::InputSlot>, RegistryError> {
let workflows = self.list_workflows().await?;
Ok(workflows
.into_iter()
.find(|w| w.workflow_name == name || w.package_name == name)
.map(|w| w.declared_params)
.unwrap_or_default())
}
pub async fn find_trigger_subscribers(
&self,
trigger_name: &str,
) -> Result<Vec<String>, RegistryError> {
let workflows = self.list_workflows().await?;
let mut seen = std::collections::HashSet::new();
Ok(workflows
.into_iter()
.filter(|w| w.workflow_triggers.iter().any(|t| t == trigger_name))
.map(|w| w.workflow_name)
.filter(|name| seen.insert(name.clone()))
.collect())
}
pub async fn find_surface_input_slots(
&self,
kind: &str,
name: &str,
) -> Result<Vec<cloacina_api_types::InputSlot>, RegistryError> {
let workflows = self.list_workflows().await?;
for w in workflows {
for surface in w.declared_surfaces {
if surface.kind == kind && surface.name == name {
return Ok(surface.slots);
}
}
}
Ok(Vec::new())
}
pub async fn find_package_for_surface(
&self,
kind: &str,
name: &str,
) -> Result<Option<String>, RegistryError> {
let workflows = self.list_workflows().await?;
for w in workflows {
for surface in &w.declared_surfaces {
if surface.kind == kind && surface.name == name {
return Ok(Some(w.package_name.clone()));
}
}
}
Ok(None)
}
pub async fn find_accumulator_input_slot(
&self,
name: &str,
) -> Result<Option<cloacina_api_types::InputSlot>, RegistryError> {
let workflows = self.list_workflows().await?;
for w in workflows {
for surface in w.declared_surfaces {
if let Some(slot) = surface.slots.into_iter().find(|s| s.name == name) {
return Ok(Some(slot));
}
}
}
Ok(None)
}
pub async fn set_workflow_paused(
&self,
name: &str,
paused: bool,
) -> Result<Option<Uuid>, RegistryError> {
let workflows = self.list_workflows().await?;
let Some(target) = workflows
.into_iter()
.find(|w| w.workflow_name == name || w.package_name == name)
else {
return Ok(None);
};
self.set_package_paused(target.id, paused).await?;
Ok(Some(target.id))
}
pub async fn get_workflow_package_by_id(
&self,
package_id: Uuid,
) -> Result<Option<(WorkflowMetadata, Vec<u8>)>, RegistryError> {
let (registry_id, metadata, _compiled) =
match self.get_package_metadata_by_id(package_id).await? {
Some(data) => data,
None => return Ok(None),
};
let package_data = match self.storage.retrieve_binary(®istry_id).await? {
Some(data) => data,
None => {
return Err(RegistryError::Internal(
"Package metadata exists but binary data is missing".to_string(),
));
}
};
Ok(Some((metadata, package_data)))
}
pub async fn get_workflow_package_by_name(
&self,
package_name: &str,
version: &str,
) -> Result<Option<(WorkflowMetadata, Vec<u8>)>, RegistryError> {
match self.get_workflow(package_name, version).await? {
Some(loaded) => Ok(Some((loaded.metadata, loaded.package_data))),
None => Ok(None),
}
}
pub async fn exists_by_id(&self, package_id: Uuid) -> Result<bool, RegistryError> {
Ok(self.get_package_metadata_by_id(package_id).await?.is_some())
}
pub async fn exists_by_name(
&self,
package_name: &str,
version: &str,
) -> Result<bool, RegistryError> {
Ok(self
.get_package_metadata(package_name, version)
.await?
.is_some())
}
pub async fn list_packages(&self) -> Result<Vec<WorkflowMetadata>, RegistryError> {
self.list_all_packages().await
}
pub async fn unregister_workflow_package_by_id(
&mut self,
package_id: Uuid,
) -> Result<(), RegistryError> {
let (registry_id, _metadata, _compiled) =
match self.get_package_metadata_by_id(package_id).await? {
Some(data) => data,
None => return Ok(()), };
if let Some(_namespaces) = self.loaded_packages.remove(&package_id) {
self.registrar
.unregister_package_tasks(&package_id.to_string())
.map_err(RegistryError::Loader)?;
}
self.delete_package_metadata_by_id(package_id).await?;
self.storage.delete_binary(®istry_id).await?;
Ok(())
}
pub async fn unregister_workflow_package_by_name(
&mut self,
package_name: &str,
version: &str,
) -> Result<(), RegistryError> {
if self
.get_package_metadata(package_name, version)
.await?
.is_none()
{
return Ok(()); }
self.unregister_workflow(package_name, version).await
}
}
#[async_trait]
impl<S: RegistryStorage + Send + Sync> WorkflowRegistry for WorkflowRegistryImpl<S> {
async fn register_workflow(
&mut self,
package_data: Vec<u8>,
) -> Result<WorkflowPackageId, RegistryError> {
if !Self::is_cloacina_package(&package_data) {
return Err(RegistryError::ValidationError {
reason: "Package data is not a valid .cloacina bzip2 source archive. \
Raw library registration is not supported."
.to_string(),
});
}
let work_dir = tempfile::TempDir::new()
.map_err(|e| RegistryError::Internal(format!("Failed to create temp dir: {}", e)))?;
let archive_path = work_dir.path().join("pkg.cloacina");
std::fs::write(&archive_path, &package_data)
.map_err(|e| RegistryError::Internal(format!("Failed to write archive: {}", e)))?;
let extract_dir = work_dir.path().join("source");
std::fs::create_dir_all(&extract_dir)
.map_err(|e| RegistryError::Internal(format!("Failed to create extract dir: {}", e)))?;
let source_dir = fidius_core::package::unpack_package(&archive_path, &extract_dir)
.map_err(|e| RegistryError::ValidationError {
reason: format!("Failed to unpack source archive: {}", e),
})?;
let manifest = fidius_core::package::load_manifest::<
cloacina_workflow_plugin::CloacinaMetadata,
>(&source_dir)
.map_err(|e| RegistryError::ValidationError {
reason: format!("Failed to load package.toml: {}", e),
})?;
let pkg_name = manifest.package.name.clone();
let pkg_version = manifest.package.version.clone();
let content_hash = {
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(&package_data);
format!("{:x}", hasher.finalize())
};
let active = self.get_active_package_by_name(&pkg_name).await?;
if let Some((existing_id, _, ref existing_hash)) = active {
if existing_hash == &content_hash {
return Ok(existing_id);
}
}
let package_metadata = crate::registry::loader::package_loader::PackageMetadata {
package_name: pkg_name,
workflow_name: manifest.metadata.workflow_name.clone().unwrap_or_default(),
version: pkg_version,
description: manifest.metadata.description.clone(),
author: manifest.metadata.author.clone(),
tasks: vec![],
graph_data: None,
architecture: std::env::consts::ARCH.to_string(),
symbols: vec![],
workflow_triggers: vec![],
declared_params: vec![],
declared_surfaces: vec![],
};
let registry_id = self.storage.store_binary(package_data).await?;
let old_id = active.map(|(id, _, _)| id);
let prebuilt = self.find_success_by_hash(&content_hash).await?;
let prebuilt_bytes = prebuilt.map(|(_, bytes)| bytes);
let package_id = self
.supersede_and_insert_with_prebuilt(
old_id,
®istry_id,
&package_metadata,
&content_hash,
prebuilt_bytes,
)
.await?;
Ok(package_id)
}
async fn get_workflow(
&self,
package_name: &str,
version: &str,
) -> Result<Option<LoadedWorkflow>, RegistryError> {
let (registry_id, package_metadata, compiled_data) =
match self.get_package_metadata(package_name, version).await? {
Some(data) => data,
None => return Ok(None),
};
let package_data = match self.storage.retrieve_binary(®istry_id).await? {
Some(data) => data,
None => {
return Err(RegistryError::Internal(
"Package metadata exists but binary data is missing".to_string(),
));
}
};
let workflow_metadata = WorkflowMetadata {
id: Uuid::new_v4(), registry_id: Uuid::parse_str(®istry_id).map_err(RegistryError::InvalidUuid)?,
workflow_name: if package_metadata.workflow_name.is_empty() {
package_metadata.package_name.clone()
} else {
package_metadata.workflow_name.clone()
},
package_name: package_metadata.package_name.clone(),
version: package_metadata.version.clone(),
description: package_metadata.description.clone(),
author: package_metadata.author.clone(),
tasks: package_metadata
.tasks
.iter()
.map(|t| t.local_id.clone())
.collect(),
task_graph: database::build_task_graph(&package_metadata),
schedules: Vec::new(),
created_at: chrono::Utc::now(),
updated_at: chrono::Utc::now(),
paused: false,
declared_params: package_metadata.declared_params.clone(),
declared_surfaces: package_metadata.declared_surfaces.clone(),
workflow_triggers: package_metadata.workflow_triggers.clone(),
};
Ok(Some(LoadedWorkflow {
metadata: workflow_metadata,
package_data,
compiled_data,
}))
}
async fn list_workflows(&self) -> Result<Vec<WorkflowMetadata>, RegistryError> {
self.list_all_packages().await
}
async fn persist_task_graph(
&self,
package_id: crate::registry::types::WorkflowPackageId,
tasks: Vec<(String, Vec<String>)>,
) -> Result<(), RegistryError> {
self.persist_task_graph_db(package_id, tasks).await
}
async fn unregister_workflow(
&mut self,
package_name: &str,
version: &str,
) -> Result<(), RegistryError> {
let (registry_id, _, _) = self
.get_package_metadata(package_name, version)
.await?
.ok_or_else(|| RegistryError::PackageNotFound {
package_name: package_name.to_string(),
version: version.to_string(),
})?;
let package_uuid = Uuid::parse_str(®istry_id).map_err(RegistryError::InvalidUuid)?;
if let Some(_namespaces) = self.loaded_packages.remove(&package_uuid) {
self.registrar
.unregister_package_tasks(&package_uuid.to_string())
.map_err(RegistryError::Loader)?;
}
self.delete_package_metadata(package_name, version).await?;
self.storage.delete_binary(®istry_id).await?;
Ok(())
}
async fn find_signature(&self, package_hash: &str) -> Result<bool, RegistryError> {
use crate::security::{DbPackageSigner, PackageSigner};
let signer = DbPackageSigner::new(crate::dal::DAL::new(self.database.clone()));
match signer.find_signature(package_hash).await {
Ok(opt) => Ok(opt.is_some()),
Err(e) => Err(RegistryError::Internal(format!(
"find_signature DAL query failed for hash {}: {}",
package_hash, e
))),
}
}
}
fn extract_source_files(archive_bytes: &[u8]) -> Result<Vec<WorkflowSourceFile>, RegistryError> {
let work_dir = tempfile::TempDir::new()
.map_err(|e| RegistryError::Internal(format!("Failed to create temp dir: {}", e)))?;
let archive_path = work_dir.path().join("pkg.cloacina");
std::fs::write(&archive_path, archive_bytes)
.map_err(|e| RegistryError::Internal(format!("Failed to write archive: {}", e)))?;
let extract_dir = work_dir.path().join("source");
std::fs::create_dir_all(&extract_dir)
.map_err(|e| RegistryError::Internal(format!("Failed to create extract dir: {}", e)))?;
let source_dir =
fidius_core::package::unpack_package(&archive_path, &extract_dir).map_err(|e| {
RegistryError::ValidationError {
reason: format!("Failed to unpack source archive: {}", e),
}
})?;
let mut files = Vec::new();
collect_source_files(&source_dir, &source_dir, &mut files)?;
files.sort_by(|a, b| a.path.cmp(&b.path));
Ok(files)
}
fn collect_source_files(
root: &Path,
dir: &Path,
out: &mut Vec<WorkflowSourceFile>,
) -> Result<(), RegistryError> {
let entries = std::fs::read_dir(dir)
.map_err(|e| RegistryError::Internal(format!("Failed to read source dir: {}", e)))?;
for entry in entries {
let entry = entry
.map_err(|e| RegistryError::Internal(format!("Failed to read dir entry: {}", e)))?;
let path = entry.path();
let file_type = entry
.file_type()
.map_err(|e| RegistryError::Internal(format!("Failed to stat dir entry: {}", e)))?;
if file_type.is_dir() {
collect_source_files(root, &path, out)?;
} else if file_type.is_file() {
if let Ok(meta) = entry.metadata() {
if meta.len() > MAX_SOURCE_FILE_BYTES {
continue;
}
}
let Ok(bytes) = std::fs::read(&path) else {
continue;
};
let Ok(contents) = String::from_utf8(bytes) else {
continue;
};
let rel = path.strip_prefix(root).unwrap_or(&path);
out.push(WorkflowSourceFile {
path: rel.to_string_lossy().replace('\\', "/"),
contents,
});
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use crate::registry::storage::FilesystemRegistryStorage;
use tempfile::TempDir;
#[tokio::test]
async fn test_registry_creation() {
let temp_dir = TempDir::new().unwrap();
let _storage = FilesystemRegistryStorage::new(temp_dir.path()).unwrap();
assert!(temp_dir.path().exists());
}
#[test]
fn test_registry_metrics() {
let temp_dir = TempDir::new().unwrap();
let _storage = FilesystemRegistryStorage::new(temp_dir.path()).unwrap();
assert!(temp_dir.path().exists());
}
}