use crate::error::Error;
use crate::node::Address;
use crate::transaction::Transaction;
use crate::util::Channel;
pub struct Topic {
pub address: Address,
pub channel: Channel<Command>,
pub subscribers: SubscriberBucket,
pub cache: Vec<Transaction>,
public: Address,
}
#[derive(Debug)]
pub enum Command {
Subscriber(Address),
Unsubscriber(Address),
Broadcast(Address, Vec<u8>),
Message(Transaction),
Drop(Address),
}
#[derive(Debug, Clone)]
pub struct SubscriberBucket {
subscribers: Vec<Address>,
}
pub struct TopicBucket {
pub topics: Vec<Simple>,
}
pub struct Simple {
pub address: Address,
pub channel: Channel<Command>,
}
impl Topic {
pub fn new(
address: Address,
channel: Channel<Command>,
subscribers: Vec<Address>,
public: Address,
) -> Self {
Self {
address,
channel,
subscribers: SubscriberBucket::new(subscribers),
cache: Vec::new(),
public,
}
}
pub fn recv(&mut self) -> Option<Transaction> {
if self.cache.len() != 0 {
return self.cache.pop();
}
loop {
match self.channel.recv() {
Some(m) => match m {
Command::Message(t) => {
if t.source() != self.public {
return Some(t);
}
}
Command::Subscriber(addr) => {
if addr != self.address && addr != self.public {
self.subscribers.add(addr);
}
}
Command::Unsubscriber(addr) => {
self.subscribers.remove(&addr);
}
_ => {
continue;
}
},
None => {
return None;
}
}
}
}
pub fn try_recv(&mut self) -> Option<Transaction> {
if self.cache.len() != 0 {
return self.cache.pop();
}
loop {
match self.channel.try_recv() {
Some(m) => match m {
Command::Message(t) => {
if t.source() != self.public {
return Some(t);
}
}
Command::Subscriber(addr) => {
if addr != self.address && addr != self.public {
self.subscribers.add(addr);
}
}
Command::Unsubscriber(addr) => {
self.subscribers.remove(&addr);
}
_ => {
continue;
}
},
None => {
return None;
}
}
}
}
pub fn broadcast(&mut self, body: Vec<u8>) -> Result<(), Error> {
loop {
match self.channel.try_recv() {
Some(m) => match m {
Command::Message(t) => {
if t.source() != self.public {
self.cache.push(t);
}
}
Command::Subscriber(addr) => {
if addr != self.address {
self.subscribers.add(addr);
}
}
Command::Unsubscriber(addr) => {
self.subscribers.remove(&addr);
}
_ => {
continue;
}
},
None => {
break;
}
}
}
for sub in &self.subscribers.subscribers {
let action = Command::Broadcast(sub.clone(), body.clone());
let e = self.channel.send(action);
if e.is_err() {
log::error!("channel is unavailable, it is possible the thread crashed.")
}
}
return Ok(());
}
pub fn unsubscribe(&mut self) {
for sub in &self.subscribers.subscribers {
let action = Command::Drop(sub.clone());
let e = self.channel.send(action);
if e.is_err() {
log::error!("channel is unavailable, it is possible the thread crashed.")
}
}
}
pub fn address(&self) -> Address {
self.address.clone()
}
}
impl Drop for Topic {
fn drop(&mut self) {
for sub in self.subscribers.clone().into_iter() {
let command = Command::Drop(sub);
let _ = self.channel.send(command);
}
}
}
impl Simple {
pub fn new(address: Address, channel: Channel<Command>) -> Self {
Self { address, channel }
}
}
impl SubscriberBucket {
pub fn new(subscribers: Vec<Address>) -> Self {
Self { subscribers }
}
pub fn add(&mut self, address: Address) {
match self.get(&address) {
Some(_) => {}
None => self.subscribers.push(address),
}
}
pub fn get(&self, search: &Address) -> Option<&Address> {
let index = self.subscribers.iter().position(|e| e == search);
match index {
Some(i) => self.subscribers.get(i),
None => None,
}
}
pub fn remove(&mut self, target: &Address) {
let index = self.subscribers.iter().position(|e| e == target);
match index {
Some(i) => {
self.subscribers.remove(i);
}
None => {}
}
}
pub fn add_bulk(&mut self, data: Vec<Address>) {
data.iter().for_each(|x| self.add(x.clone()));
}
pub fn len(&self) -> usize {
self.subscribers.len()
}
}
impl Iterator for SubscriberBucket {
type Item = Address;
fn next(&mut self) -> Option<Self::Item> {
self.subscribers.pop()
}
}
impl TopicBucket {
pub fn new() -> Self {
Self { topics: Vec::new() }
}
pub fn add(&mut self, simple: Simple) {
if TopicBucket::find(&self, &simple.address).is_none() {
self.topics.push(simple);
}
}
pub fn find(&self, search: &Address) -> Option<&Simple> {
let index = self.topics.iter().position(|e| &e.address == search);
match index {
Some(i) => self.topics.get(i),
None => None,
}
}
pub fn find_mut(&mut self, search: &Address) -> Option<&mut Simple> {
let index = self.topics.iter().position(|e| &e.address == search);
match index {
Some(i) => self.topics.get_mut(i),
None => None,
}
}
pub fn remove(&mut self, target: &Address) {
let index = self.topics.iter().position(|e| &e.address == target);
match index {
Some(i) => {
self.topics.remove(i);
}
None => {}
}
}
pub fn is_local(&self, query: &Address) -> bool {
match self.find(query) {
Some(_) => true,
None => false,
}
}
pub fn len(&self) -> usize {
self.topics.len()
}
}
impl Iterator for TopicBucket {
type Item = Simple;
fn next(&mut self) -> Option<Self::Item> {
self.topics.pop()
}
}