Commit

Author:

Hash:

Timestamp:

+0 -170 +/-2 browse

Kevin Schoon [me@kevinschoon.com]

4f3056d6df5bf5abb5fade5ea458058c6d660a6c

Wed, 24 Jun 2026 13:37:30 +0000 (2 months ago)

drop very old scheduler crate
1diff --git a/crates/scheduler/Cargo.toml b/crates/scheduler/Cargo.toml
2deleted file mode 100644
3index 0ee043b..0000000
4--- a/crates/scheduler/Cargo.toml
5+++ /dev/null
6 @@ -1,9 +0,0 @@
7- [package]
8- name = "ayllu_scheduler"
9- version = "0.2.1"
10- edition = "2024"
11-
12- # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
13-
14- [dependencies]
15- tracing = { workspace = true }
16 diff --git a/crates/scheduler/src/lib.rs b/crates/scheduler/src/lib.rs
17deleted file mode 100644
18index 399e376..0000000
19--- a/crates/scheduler/src/lib.rs
20+++ /dev/null
21 @@ -1,161 +0,0 @@
22- use std::fmt::Debug;
23- use std::sync::mpsc::{channel, sync_channel, Receiver, Sender, SyncSender};
24- use std::thread::spawn;
25- use std::time::Duration;
26-
27- use tracing::log::{info, warn};
28-
29- const TIMEOUT_MS: u64 = 500;
30-
31- // blocking worker bound to a single thread
32- pub struct Worker<T, F>
33- where
34- T: Debug + Clone + Send + 'static,
35- F: Fn(T),
36- {
37- name: String,
38- callback: F,
39- rx: Receiver<T>,
40- }
41-
42- impl<T, F> Worker<T, F>
43- where
44- T: Debug + Clone + Send + 'static,
45- F: Fn(T),
46- {
47- fn new(name: &str, rx: Receiver<T>, callback: F) -> Self {
48- Self {
49- name: name.to_string(),
50- rx,
51- callback,
52- }
53- }
54-
55- fn process(&self) {
56- info!("{} is starting up", self.name);
57- loop {
58- match self.rx.recv() {
59- Ok(msg) => {
60- info!("processing job: {}: {:?}", self.name, msg);
61- let cb = &self.callback;
62- cb(msg);
63- }
64- Err(_) => {
65- warn!("worker {} shutting down", self.name);
66- break;
67- }
68- }
69- }
70- }
71- }
72-
73- // simple thread-based job scheduler
74- pub struct Scheduler<T>
75- where
76- T: Debug + Clone,
77- {
78- size: usize,
79- sender_ch: Option<SyncSender<T>>,
80- shutdown_ch: Option<Sender<bool>>,
81- }
82-
83- impl<T> Scheduler<T>
84- where
85- T: Debug + Clone + Send + 'static,
86- {
87- pub fn new(size: usize) -> Self {
88- Scheduler {
89- size,
90- sender_ch: None,
91- shutdown_ch: None,
92- }
93- }
94-
95- // process all jobs
96- pub fn initialize<F>(&mut self, callback: F)
97- where
98- F: Fn(T) + Copy + Send + 'static,
99- {
100- let mut workers: Vec<SyncSender<T>> = Vec::new();
101- let size = self.size;
102- for i in 0..self.size {
103- let (tx, rx) = sync_channel::<T>(size);
104- spawn(move || {
105- Worker::<T, F>::new(format!("Worker{}", i).as_str(), rx, callback).process();
106- });
107- workers.push(tx);
108- }
109- let (tx_shutdown, rx_shutdown) = channel::<bool>();
110- let (tx, rx) = sync_channel::<T>(size);
111- let n_workers = size;
112- spawn(move || {
113- let mut position = 0;
114- loop {
115- if let Ok(next_job) = rx.recv_timeout(Duration::from_millis(TIMEOUT_MS)) {
116- info!("sending job to worker: {}", position);
117- let worker = workers.get(position).unwrap();
118- worker.send(next_job).unwrap();
119- if position + 1 == n_workers {
120- position = 0;
121- } else {
122- position += 1;
123- }
124- }
125- if rx_shutdown
126- .recv_timeout(Duration::from_millis(TIMEOUT_MS))
127- .is_ok()
128- {
129- info!("scheduler is shutting down")
130- }
131- }
132- });
133- self.sender_ch = Some(tx);
134- self.shutdown_ch = Some(tx_shutdown);
135- }
136-
137- pub fn submit(&self, job: T) {
138- self.sender_ch.as_ref().unwrap().send(job).unwrap()
139- }
140-
141- pub fn shutdown(&self) {
142- self.shutdown_ch.as_ref().unwrap().send(true).unwrap();
143- }
144- }
145-
146- #[cfg(test)]
147- mod tests {
148- use super::*;
149-
150- use std::fmt::Debug;
151-
152- #[derive(Debug, Clone)]
153- struct Job {}
154-
155- #[test]
156- fn test_scheduler() {
157- let mut scheduler = Scheduler::new(5);
158- scheduler.initialize(|job| {
159- println!("its a job: {:?}", job);
160- });
161- scheduler.submit(Job {});
162- scheduler.submit(Job {});
163- scheduler.submit(Job {});
164- scheduler.submit(Job {});
165- scheduler.submit(Job {});
166- scheduler.submit(Job {});
167- scheduler.submit(Job {});
168- scheduler.submit(Job {});
169- scheduler.submit(Job {});
170- scheduler.shutdown();
171- }
172-
173- #[test]
174- fn test_scheduler_shutdown_send() {
175- let mut scheduler = Scheduler::new(5);
176- scheduler.initialize(|job| {
177- println!("its a job: {:?}", job);
178- });
179- scheduler.shutdown();
180- scheduler.submit(Job {});
181- }
182- }