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
88
89
90
91
92
93
94
use crossbeam_channel::{unbounded, Receiver, Sender};
use std::collections::BTreeMap;
use Uuid;

pub struct RxHandle<T: Clone> {
    uuid: Uuid,
    receiver: Receiver<T>,
    register_sender: Sender<SenderRegistration<T>>,
}

impl<T: Clone> RxHandle<T> {
    pub fn receiver(&self) -> &Receiver<T> {
        &self.receiver
    }

    pub fn recv(&self) -> Option<T> {
        self.receiver.recv()
    }
}

impl<T: Clone> RxHandle<T> {
    pub fn new(register_sender: &Sender<SenderRegistration<T>>) -> Self {
        let uuid = Uuid::new_v4();
        let (sender, receiver) = unbounded();
        register_sender.send(SenderRegistration::AttachSender {
            uuid: uuid.clone(),
            sender,
        });
        RxHandle {
            uuid,
            receiver,
            register_sender: register_sender.clone(),
        }
    }
}

impl<T: Clone> Drop for RxHandle<T> {
    fn drop(&mut self) {
        self.register_sender
            .send(SenderRegistration::DetachSender {
                uuid: self.uuid.clone(),
            });
    }
}

#[derive(Clone)]
pub enum SenderRegistration<T: Clone> {
    AttachSender {
        uuid: Uuid,
        sender: Sender<T>,
    },
    DetachSender {
        uuid: Uuid,
    },
}

pub struct RegisteredSenders<T: Clone> {
    senders: BTreeMap<Uuid, Sender<T>>,
    registration: Receiver<SenderRegistration<T>>,
}

impl<T: Clone> RegisteredSenders<T> {
    pub fn new(registration: Receiver<SenderRegistration<T>>) -> Self {
        RegisteredSenders {
            senders: BTreeMap::new(),
            registration,
        }
    }

    pub fn process_pending(&mut self) {
        self.process_pending_with(|_|{});
    }

    pub fn process_pending_with<F: FnMut(&SenderRegistration<T>)>(&mut self, mut pre_process: F) {
        use self::SenderRegistration::*;
        while let Some(s) = self.registration.try_recv() {
            pre_process(&s);
            match s {
                AttachSender { uuid, sender } => {
                    self.senders.insert(uuid, sender);
                }
                DetachSender { uuid } => {
                    self.senders.remove(&uuid);
                }
            }
        }
    }

    pub fn send(&mut self, item: T) {
        for sender in self.senders.values() {
            sender.send(item.clone());
        }
    }
}