Skip to main content

asched_core/
registry.rs

1//! Generic project registry.
2//! ref: README.md#architecture
3
4use crate::routine::store::{atomic_toml, read_text_limited};
5use serde::{Deserialize, Serialize};
6use std::collections::HashSet;
7use std::fs::{self, File, OpenOptions};
8use std::os::fd::AsRawFd;
9use std::os::unix::fs::PermissionsExt;
10use std::path::{Path, PathBuf};
11use std::sync::{Mutex, MutexGuard, OnceLock};
12
13const REGISTRY_VERSION: u32 = 1;
14static REGISTRY_WRITE_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
15
16#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
17#[serde(deny_unknown_fields)]
18pub struct Project {
19    pub name: String,
20    pub working_dir: PathBuf,
21}
22
23impl Project {
24    pub fn validated(mut self) -> Result<Self, RegistryError> {
25        self.name = self.name.trim().to_string();
26        validate_name(&self.name)?;
27        self.working_dir = fs::canonicalize(&self.working_dir).map_err(|error| {
28            RegistryError::Validation(format!(
29                "cannot resolve working directory {}: {error}",
30                self.working_dir.display()
31            ))
32        })?;
33        if !self.working_dir.is_dir() {
34            return Err(RegistryError::Validation(format!(
35                "{} is not a directory",
36                self.working_dir.display()
37            )));
38        }
39        Ok(self)
40    }
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(deny_unknown_fields)]
45// ^ [[Project Registry and Scheduling Allowlist]]
46pub struct ProjectRegistry {
47    pub version: u32,
48    pub revision: u64,
49    #[serde(default)]
50    pub projects: Vec<Project>,
51}
52
53impl Default for ProjectRegistry {
54    fn default() -> Self {
55        Self {
56            version: REGISTRY_VERSION,
57            revision: 0,
58            projects: Vec::new(),
59        }
60    }
61}
62
63#[derive(Debug, Clone)]
64pub struct RegistryStore {
65    root: PathBuf,
66}
67
68impl RegistryStore {
69    pub fn new(root: PathBuf) -> Self {
70        Self { root }
71    }
72
73    pub fn default_root() -> Result<PathBuf, RegistryError> {
74        if let Some(path) = std::env::var_os("ASCHED_ROOT") {
75            return Ok(PathBuf::from(path));
76        }
77        dirs::config_dir()
78            .map(|path| path.join("asched"))
79            .ok_or_else(|| RegistryError::Unavailable("no config directory".into()))
80    }
81
82    pub fn root(&self) -> &Path {
83        &self.root
84    }
85
86    pub fn path(&self) -> PathBuf {
87        self.root.join("projects.toml")
88    }
89
90    pub fn load(&self) -> Result<ProjectRegistry, RegistryError> {
91        let path = self.path();
92        if !path.exists() {
93            return Ok(ProjectRegistry::default());
94        }
95        let text = read_text_limited(&path)?;
96        let registry: ProjectRegistry = toml::from_str(&text)
97            .map_err(|error| RegistryError::Corrupt(format!("{}: {error}", path.display())))?;
98        validate_registry(registry)
99    }
100
101    pub fn add(
102        &self,
103        expected_revision: u64,
104        project: Project,
105    ) -> Result<ProjectRegistry, RegistryError> {
106        self.modify(expected_revision, |registry| {
107            let project = project.validated()?;
108            if registry
109                .projects
110                .iter()
111                .any(|item| item.name == project.name)
112            {
113                return Err(RegistryError::Duplicate(project.name));
114            }
115            if let Some(existing) = registry
116                .projects
117                .iter()
118                .find(|item| item.working_dir == project.working_dir)
119            {
120                return Err(RegistryError::Validation(format!(
121                    "{} is already registered as project '{}'",
122                    project.working_dir.display(),
123                    existing.name
124                )));
125            }
126            registry.projects.push(project);
127            registry
128                .projects
129                .sort_by(|left, right| left.name.cmp(&right.name));
130            Ok(())
131        })
132    }
133
134    pub fn remove(
135        &self,
136        expected_revision: u64,
137        name: &str,
138    ) -> Result<ProjectRegistry, RegistryError> {
139        self.modify(expected_revision, |registry| {
140            let before = registry.projects.len();
141            registry.projects.retain(|project| project.name != name);
142            if registry.projects.len() == before {
143                return Err(RegistryError::NotFound(name.to_string()));
144            }
145            Ok(())
146        })
147    }
148
149    pub fn merge(
150        &self,
151        expected_revision: u64,
152        projects: Vec<Project>,
153    ) -> Result<ProjectRegistry, RegistryError> {
154        self.modify(expected_revision, |registry| merge_into(registry, projects))
155    }
156
157    pub fn select(
158        &self,
159        names: &[String],
160        filter: Option<&str>,
161    ) -> Result<Vec<Project>, RegistryError> {
162        let registry = self.load()?;
163        let requested = names.iter().collect::<HashSet<_>>();
164        let missing = names
165            .iter()
166            .filter(|name| {
167                !registry
168                    .projects
169                    .iter()
170                    .any(|project| &project.name == *name)
171            })
172            .cloned()
173            .collect::<Vec<_>>();
174        if !missing.is_empty() {
175            return Err(RegistryError::NotFound(missing.join(", ")));
176        }
177        let filter = filter
178            .map(str::trim)
179            .filter(|value| !value.is_empty())
180            .map(str::to_lowercase);
181        Ok(registry
182            .projects
183            .into_iter()
184            .filter(|project| requested.is_empty() || requested.contains(&project.name))
185            .filter(|project| {
186                filter.as_deref().is_none_or(|value| {
187                    project.name.to_lowercase().contains(value)
188                        || project
189                            .working_dir
190                            .to_string_lossy()
191                            .to_lowercase()
192                            .contains(value)
193                })
194            })
195            .collect())
196    }
197
198    fn modify(
199        &self,
200        expected_revision: u64,
201        update: impl FnOnce(&mut ProjectRegistry) -> Result<(), RegistryError>,
202    ) -> Result<ProjectRegistry, RegistryError> {
203        // ^ Registry writes bypass the daemon; keep file locking, revision checks, and atomic_toml together.
204        let _guard = self.exclusive_lock()?;
205        self.modify_locked(expected_revision, update)
206    }
207
208    pub(crate) fn merge_locked(
209        &self,
210        expected_revision: u64,
211        projects: Vec<Project>,
212    ) -> Result<ProjectRegistry, RegistryError> {
213        self.modify_locked(expected_revision, |registry| merge_into(registry, projects))
214    }
215
216    pub(crate) fn exclusive_lock(&self) -> Result<RegistryWriteGuard, RegistryError> {
217        let lock = REGISTRY_WRITE_LOCK.get_or_init(|| Mutex::new(()));
218        let process = lock
219            .lock()
220            .map_err(|_| RegistryError::Io("project registry lock poisoned".into()))?;
221        let file = RegistryFileLock::acquire(&self.root)?;
222        Ok(RegistryWriteGuard {
223            _process: process,
224            _file: file,
225        })
226    }
227
228    fn modify_locked(
229        &self,
230        expected_revision: u64,
231        update: impl FnOnce(&mut ProjectRegistry) -> Result<(), RegistryError>,
232    ) -> Result<ProjectRegistry, RegistryError> {
233        let mut registry = self.load()?;
234        if registry.revision != expected_revision {
235            return Err(RegistryError::Conflict {
236                expected: expected_revision,
237                actual: registry.revision,
238            });
239        }
240        update(&mut registry)?;
241        registry.revision = registry
242            .revision
243            .checked_add(1)
244            .ok_or_else(|| RegistryError::Corrupt("registry revision overflow".into()))?;
245        atomic_toml(&self.path(), &registry)
246            .map_err(|error| RegistryError::Io(error.to_string()))?;
247        Ok(registry)
248    }
249}
250
251fn merge_into(registry: &mut ProjectRegistry, projects: Vec<Project>) -> Result<(), RegistryError> {
252    for project in projects {
253        let project = project.validated()?;
254        if let Some(existing) = registry
255            .projects
256            .iter()
257            .find(|item| item.name == project.name || item.working_dir == project.working_dir)
258        {
259            if existing != &project {
260                return Err(RegistryError::Validation(format!(
261                    "project '{}' conflicts with existing project '{}' ({})",
262                    project.name,
263                    existing.name,
264                    existing.working_dir.display()
265                )));
266            }
267            continue;
268        }
269        registry.projects.push(project);
270    }
271    registry
272        .projects
273        .sort_by(|left, right| left.name.cmp(&right.name));
274    Ok(())
275}
276
277fn validate_registry(registry: ProjectRegistry) -> Result<ProjectRegistry, RegistryError> {
278    if registry.version != REGISTRY_VERSION {
279        return Err(RegistryError::Corrupt(format!(
280            "unsupported project registry schema {}",
281            registry.version
282        )));
283    }
284    let mut names = HashSet::new();
285    let mut paths = HashSet::new();
286    for project in &registry.projects {
287        validate_name(&project.name).map_err(|error| RegistryError::Corrupt(error.to_string()))?;
288        if project.name.trim() != project.name {
289            return Err(RegistryError::Corrupt(
290                "stored project name must be normalized".into(),
291            ));
292        }
293        if !project.working_dir.is_absolute() {
294            return Err(RegistryError::Corrupt(format!(
295                "stored working directory must be absolute: {}",
296                project.working_dir.display()
297            )));
298        }
299        if project.working_dir.exists() {
300            if !project.working_dir.is_dir() {
301                return Err(RegistryError::Corrupt(format!(
302                    "stored working directory is not a directory: {}",
303                    project.working_dir.display()
304                )));
305            }
306            let canonical = fs::canonicalize(&project.working_dir)?;
307            if canonical != project.working_dir {
308                return Err(RegistryError::Corrupt(format!(
309                    "stored working directory must be canonical: {}",
310                    project.working_dir.display()
311                )));
312            }
313        }
314        if !names.insert(&project.name) {
315            return Err(RegistryError::Duplicate(project.name.clone()));
316        }
317        if !paths.insert(&project.working_dir) {
318            return Err(RegistryError::Corrupt(format!(
319                "working directory registered more than once: {}",
320                project.working_dir.display()
321            )));
322        }
323    }
324    Ok(registry)
325}
326
327fn validate_name(name: &str) -> Result<(), RegistryError> {
328    if name.is_empty() {
329        return Err(RegistryError::Validation(
330            "project name must not be empty".into(),
331        ));
332    }
333    if name.contains('/') || name.chars().any(char::is_control) {
334        return Err(RegistryError::Validation(
335            "project name must not contain '/' or control characters".into(),
336        ));
337    }
338    Ok(())
339}
340
341struct RegistryFileLock(File);
342
343pub(crate) struct RegistryWriteGuard {
344    _process: MutexGuard<'static, ()>,
345    _file: RegistryFileLock,
346}
347
348impl RegistryFileLock {
349    fn acquire(root: &Path) -> Result<Self, RegistryError> {
350        fs::create_dir_all(root)?;
351        fs::set_permissions(root, fs::Permissions::from_mode(0o700))?;
352        let path = root.join("projects.lock");
353        let file = OpenOptions::new()
354            .read(true)
355            .write(true)
356            .create(true)
357            .truncate(false)
358            .open(&path)
359            .map_err(|error| RegistryError::Io(format!("opening {}: {error}", path.display())))?;
360        file.set_permissions(fs::Permissions::from_mode(0o600))?;
361        let result = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) };
362        if result != 0 {
363            return Err(RegistryError::Io(format!(
364                "locking {}: {}",
365                path.display(),
366                std::io::Error::last_os_error()
367            )));
368        }
369        Ok(Self(file))
370    }
371}
372
373impl Drop for RegistryFileLock {
374    fn drop(&mut self) {
375        unsafe {
376            libc::flock(self.0.as_raw_fd(), libc::LOCK_UN);
377        }
378    }
379}
380
381#[derive(Debug, thiserror::Error)]
382pub enum RegistryError {
383    #[error("invalid project: {0}")]
384    Validation(String),
385    #[error("project '{0}' already exists")]
386    Duplicate(String),
387    #[error("project '{0}' not found")]
388    NotFound(String),
389    #[error("stale registry revision: expected {expected}, actual {actual}")]
390    Conflict { expected: u64, actual: u64 },
391    #[error("project registry unavailable: {0}")]
392    Unavailable(String),
393    #[error("project registry I/O error: {0}")]
394    Io(String),
395    #[error("invalid project registry: {0}")]
396    Corrupt(String),
397}
398
399impl From<std::io::Error> for RegistryError {
400    fn from(value: std::io::Error) -> Self {
401        Self::Io(value.to_string())
402    }
403}
404
405#[cfg(test)]
406#[path = "registry_contract_tests.rs"]
407mod registry_contract_tests;