1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
use decoder::MAX_COMPONENTS;
use error::Result;
use std::{mem, sync::mpsc::{self, Sender}};
use std::thread;
use super::{RowData, Worker};
use super::immediate::ImmediateWorker;
enum WorkerMsg {
Start(RowData),
AppendRow(Vec<i16>),
GetResult(Sender<Vec<u8>>),
}
pub struct MultiThreadedWorker {
senders: [Option<Sender<WorkerMsg>>; MAX_COMPONENTS]
}
impl Worker for MultiThreadedWorker {
fn new() -> Result<Self> {
Ok(MultiThreadedWorker {
senders: [None, None, None, None]
})
}
fn start(&mut self, row_data: RowData) -> Result<()> {
let component = row_data.index;
if let None = self.senders[component] {
let sender = spawn_worker_thread(component)?;
self.senders[component] = Some(sender);
}
let sender = mem::replace(&mut self.senders[component], None).unwrap();
sender.send(WorkerMsg::Start(row_data)).expect("jpeg-decoder worker thread error");
self.senders[component] = Some(sender);
Ok(())
}
fn append_row(&mut self, row: (usize, Vec<i16>)) -> Result<()> {
let component = row.0;
let sender = mem::replace(&mut self.senders[component], None).unwrap();
sender.send(WorkerMsg::AppendRow(row.1)).expect("jpeg-decoder worker thread error");
self.senders[component] = Some(sender);
Ok(())
}
fn get_result(&mut self, index: usize) -> Result<Vec<u8>> {
let (tx, rx) = mpsc::channel();
let sender = mem::replace(&mut self.senders[index], None).unwrap();
sender.send(WorkerMsg::GetResult(tx)).expect("jpeg-decoder worker thread error");
Ok(rx.recv().expect("jpeg-decoder worker thread error"))
}
}
fn spawn_worker_thread(component: usize) -> Result<Sender<WorkerMsg>> {
let thread_builder = thread::Builder::new().name(format!("worker thread for component {}", component));
let (tx, rx) = mpsc::channel();
thread_builder.spawn(move || {
let mut worker = ImmediateWorker::new_immediate();
while let Ok(message) = rx.recv() {
match message {
WorkerMsg::Start(mut data) => {
data.index = 0;
worker.start_immediate(data);
},
WorkerMsg::AppendRow(row) => {
worker.append_row_immediate((0, row));
},
WorkerMsg::GetResult(chan) => {
let _ = chan.send(worker.get_result_immediate(0));
break;
},
}
}
})?;
Ok(tx)
}