1
use std::{any::Any, collections::HashMap, fmt::Debug, sync::Arc};
2

            
3
use flume::{Receiver, Sender};
4
use tokio::sync::oneshot;
5

            
6
use crate::jobs::{
7
    manager::{ManagedJob, Manager},
8
    task::{Handle, Id},
9
    traits::Executable,
10
    Job, Keyed,
11
};
12

            
13
pub struct Jobs<Key> {
14
    last_task_id: u64,
15
    result_senders: HashMap<Id, Vec<Box<dyn AnySender>>>,
16
    keyed_jobs: HashMap<Key, Id>,
17
    queuer: Sender<Box<dyn Executable>>,
18
    queue: Receiver<Box<dyn Executable>>,
19
}
20

            
21
impl<Key> Debug for Jobs<Key>
22
where
23
    Key: Debug,
24
{
25
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
26
        f.debug_struct("Jobs")
27
            .field("last_task_id", &self.last_task_id)
28
            .field("result_senders", &self.result_senders.len())
29
            .field("keyed_jobs", &self.keyed_jobs)
30
            .field("queuer", &self.queuer)
31
            .field("queue", &self.queue)
32
            .finish()
33
    }
34
}
35

            
36
impl<Key> Default for Jobs<Key> {
37
239
    fn default() -> Self {
38
239
        let (queuer, queue) = flume::unbounded();
39
239

            
40
239
        Self {
41
239
            last_task_id: 0,
42
239
            result_senders: HashMap::new(),
43
239
            keyed_jobs: HashMap::new(),
44
239
            queuer,
45
239
            queue,
46
239
        }
47
239
    }
48
}
49

            
50
impl<Key> Jobs<Key>
51
where
52
    Key: Clone + std::hash::Hash + Eq + Send + Sync + Debug + 'static,
53
{
54
1861
    pub fn queue(&self) -> Receiver<Box<dyn Executable>> {
55
1861
        self.queue.clone()
56
1861
    }
57

            
58
173955
    pub fn enqueue<J: Job + 'static>(
59
173955
        &mut self,
60
173955
        job: J,
61
173955
        key: Option<Key>,
62
173955
        manager: Manager<Key>,
63
173955
    ) -> Handle<J::Output, J::Error, Key> {
64
173955
        self.last_task_id = self.last_task_id.wrapping_add(1);
65
173955
        let id = Id(self.last_task_id);
66
173955
        self.queuer
67
173955
            .send(Box::new(ManagedJob {
68
173955
                id,
69
173955
                job,
70
173955
                key,
71
173955
                manager: manager.clone(),
72
173955
            }))
73
173955
            .unwrap();
74
173955

            
75
173955
        self.create_new_task_handle(id, manager)
76
173955
    }
77

            
78
183307
    pub fn create_new_task_handle<T: Send + Sync + 'static, E: Send + Sync + 'static>(
79
183307
        &mut self,
80
183307
        id: Id,
81
183307
        manager: Manager<Key>,
82
183307
    ) -> Handle<T, E, Key> {
83
183307
        let (sender, receiver) = oneshot::channel();
84
183307
        let senders = self.result_senders.entry(id).or_insert_with(Vec::default);
85
183307
        senders.push(Box::new(Some(sender)));
86
183307

            
87
183307
        Handle {
88
183307
            id,
89
183307
            manager,
90
183307
            receiver,
91
183307
        }
92
183307
    }
93

            
94
130720
    pub fn lookup_or_enqueue<J: Keyed<Key>>(
95
130720
        &mut self,
96
130720
        job: J,
97
130720
        manager: Manager<Key>,
98
130720
    ) -> Handle<<J as Job>::Output, <J as Job>::Error, Key> {
99
130720
        let key = job.key();
100
130720
        if let Some(&id) = self.keyed_jobs.get(&key) {
101
9351
            self.create_new_task_handle(id, manager)
102
        } else {
103
121369
            let handle = self.enqueue(job, Some(key.clone()), manager);
104
121369
            self.keyed_jobs.insert(key, handle.id);
105
121369
            handle
106
        }
107
130720
    }
108

            
109
    pub fn job_completed<T: Clone + Send + Sync + 'static, E: Send + Sync + 'static>(
110
        &mut self,
111
        id: Id,
112
        key: Option<&Key>,
113
        result: Result<T, E>,
114
    ) {
115
173721
        if let Some(key) = key {
116
121135
            self.keyed_jobs.remove(key);
117
158474
        }
118

            
119
173721
        if let Some(senders) = self.result_senders.remove(&id) {
120
173721
            let result = result.map_err(Arc::new);
121
356794
            for mut sender_handle in senders {
122
183073
                let sender = sender_handle
123
183073
                    .as_any_mut()
124
183073
                    .downcast_mut::<Option<oneshot::Sender<Result<T, Arc<E>>>>>()
125
183073
                    .unwrap();
126
183073
                if let Some(sender) = sender.take() {
127
183073
                    drop(sender.send(result.clone()));
128
183073
                }
129
            }
130
        }
131
173721
    }
132
}
133

            
134
pub trait AnySender: Any + Send + Sync {
135
    fn as_any_mut(&mut self) -> &'_ mut dyn Any;
136
}
137

            
138
impl<T> AnySender for Option<oneshot::Sender<T>>
139
where
140
    T: Send + Sync + 'static,
141
{
142
183073
    fn as_any_mut(&mut self) -> &'_ mut dyn Any {
143
183073
        self
144
183073
    }
145
}