Files
aho_corasick
ansi_term
atty
backtrace
backtrace_sys
bitflags
blindbid
block_buffer
block_padding
bulletproofs
byte_tools
byteorder
cfg_if
chrono
clap
clear_on_drop
curve25519_dalek
digest
dusk_blindbidproof
dusk_tlv
dusk_uds
env_logger
failure
failure_derive
fake_simd
futures
futures_channel
futures_core
futures_executor
futures_io
futures_macro
futures_sink
futures_task
futures_util
async_await
future
io
lock
sink
stream
task
generic_array
humantime
keccak
lazy_static
libc
log
memchr
merlin
num_cpus
num_integer
num_traits
opaque_debug
packed_simd
pin_utils
proc_macro2
proc_macro_hack
proc_macro_nested
quick_error
quote
rand
rand_chacha
rand_core
rand_hc
rand_isaac
rand_jitter
rand_os
rand_pcg
rand_xorshift
regex
regex_syntax
rustc_demangle
serde
serde_derive
sha2
sha3
slab
strsim
subtle
syn
synstructure
termcolor
textwrap
thread_local
time
typenum
unicode_width
unicode_xid
vec_map
 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
use crate::{Message, Task, TaskProvider};

use std::sync::{mpsc, Arc, Mutex};

use futures::executor::block_on;

/// This function will panic if the task channel is broken.
///
/// There is no point in preserving the event loop in case there is no available channel.
///
/// Therefore, this function should be called from a thread.
pub fn worker<T: TaskProvider>(
    tx: mpsc::Sender<Task>,
    rx: Arc<Mutex<mpsc::Receiver<Task>>>,
    provider: T,
) {
    loop {
        let task = rx
            .lock()
            .map_err(|e| {
                error!("Error trying to lock the task channel: {}", e);
                e
            })
            .unwrap()
            .recv()
            .map_err(|e| {
                error!(
                    "Error trying to receive a task from the respective channel: {}",
                    e
                );
                e
            })
            .unwrap();

        match task {
            Task::Socket(stream) => {
                let mut p = provider.clone();

                p.set_socket(stream);

                // TODO - Naive implementation, will not reschedule if the poll returns pending
                let message = block_on(p);

                if Message::ShouldQuit == message {
                    tx.send(Task::Message(Message::ShouldQuit))
                        .map_err(|e| {
                            error!(
                                "Error trying to send a ShouldQuit message to the task channel: {}",
                                e
                            );
                            e
                        })
                        .unwrap();
                }
            }

            Task::Message(m) if m == Message::ShouldQuit => {
                tx.send(Task::Message(Message::ShouldQuit))
                    .map_err(|e| {
                        error!(
                            "Error trying to send a ShouldQuit message to the task channel: {}",
                            e
                        );
                        e
                    })
                    .unwrap();
                break;
            }

            _ => (),
        }
    }
}