dexter-daemon 0.1.0

Worker repurposing plans for staged Dexter daemon pipelines
Documentation
/// Worker allocation plan for staged daemon pipelines where fetch and persistence
/// run concurrently.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WorkerRepurposePlan {
    /// Number of fetch workers actively scraping upstream sources.
    pub fetch_workers: usize,
    /// Number of dedicated persistence workers (CSV + inserts + post-fetch side effects).
    pub persist_workers: usize,
    /// Capacity of the fetch queue feeding fetch workers.
    pub fetch_queue_capacity: usize,
    /// Capacity of the persistence queue feeding persist workers.
    pub persist_queue_capacity: usize,
    /// Number of randomized fetch jobs to keep in flight.
    pub dispatch_target: usize,
}

impl WorkerRepurposePlan {
    /// Build a plan for randomized registry ingestion.
    ///
    /// `dispatch_target` should already include provider-specific prefetch shaping.
    pub fn for_randomized_registry(
        fetch_workers: usize,
        fetch_queue_capacity: usize,
        dispatch_target: usize,
    ) -> Self {
        let fetch_workers = fetch_workers.max(1);
        let fetch_queue_capacity = fetch_queue_capacity.max(fetch_workers * 2).max(4);
        let dispatch_target = dispatch_target.max(1);
        let persist_workers = 1;
        let persist_queue_capacity = fetch_queue_capacity.max(dispatch_target).max(4);
        Self {
            fetch_workers,
            persist_workers,
            fetch_queue_capacity,
            persist_queue_capacity,
            dispatch_target,
        }
    }
}

#[cfg(test)]
mod tests {
    use super::WorkerRepurposePlan;

    #[test]
    fn randomized_registry_plan_keeps_minimum_bounds() {
        let plan = WorkerRepurposePlan::for_randomized_registry(0, 0, 0);
        assert_eq!(plan.fetch_workers, 1);
        assert_eq!(plan.persist_workers, 1);
        assert_eq!(plan.fetch_queue_capacity, 4);
        assert_eq!(plan.persist_queue_capacity, 4);
        assert_eq!(plan.dispatch_target, 1);
    }

    #[test]
    fn randomized_registry_plan_scales_persist_queue_to_dispatch_target() {
        let plan = WorkerRepurposePlan::for_randomized_registry(4, 8, 64);
        assert_eq!(plan.fetch_workers, 4);
        assert_eq!(plan.fetch_queue_capacity, 8);
        assert_eq!(plan.persist_workers, 1);
        assert_eq!(plan.persist_queue_capacity, 64);
        assert_eq!(plan.dispatch_target, 64);
    }
}