This commit is contained in:
Piyush मिश्रः 2021-05-10 00:48:28 +05:30
parent 1f9d1f7bfd
commit 6d18337f10
3 changed files with 0 additions and 1168 deletions

View File

@ -1,404 +0,0 @@
//! Chat Pinnd(पिण्ड) is Actor to manage Websocket Chat related action
use std::{collections::HashMap, vec};
use actix::prelude::*;
use actix_broker::BrokerSubscribe;
use crate::{ws_sansad, messages as ms};
#[allow(dead_code)]
pub struct ChatPinnd {
kaksh: HashMap<String, Kaksh>, // kunjika, Kaksh
vyaktigat_waitlist: Vec<VyaktiWatchlist>,
}
pub struct Kaksh {
length: Option<usize>,
last_message_id: u128,
loog: Vec<Loog>
}
pub struct Loog {
addr: Addr<ws_sansad::WsSansad>,
kunjika: String,
name: String,
tags: Option<Vec<String>>
}
#[derive(Debug, Clone)]
pub struct Vyakti {
name: String,
tags: Vec<String>
}
pub struct VyaktiWatchlist {
kunjika: String,
name: String,
tags: Vec<String>,
addr: Addr<ws_sansad::WsSansad>
}
impl Actor for ChatPinnd {
type Context = Context<Self>;
fn started(&mut self, ctx: &mut Self::Context) {
// for actix broker
self.subscribe_system_async::<ms::SendText>(ctx);
self.subscribe_system_async::<ms::SendStatus>(ctx);
self.subscribe_system_async::<ms::LeaveUser>(ctx);
}
}
/// Join kaksh
impl Handler<ms::JoinKaksh> for ChatPinnd {
type Result = ms::Resp;
fn handle(&mut self, msg: ms::JoinKaksh, _: &mut Self::Context) -> Self::Result {
// check if user exist
if let Some(_) = self.vyaktigat_waitlist.iter().position(|vk| vk.kunjika == msg.kunjika) {
return ms::Resp::Err("Kunjika already exist".to_owned());
}
if let Some(_) = self.kaksh.iter().position(|(_,g)| {
match g.loog.iter().position(|a| a.kunjika == msg.kunjika) {
Some(_) => true,
None => false
}
}) {
return ms::Resp::Err("Kunjika already exist".to_owned());
}
// check if kaksh exist and add user
match self.kaksh.get_mut(&msg.kaksh_kunjika) {
Some(kaksh) =>{ // exist
// check if kaksh have no space left
if let Some(n) = kaksh.length {
if kaksh.loog.len() >= n {
return ms::Resp::Err("Kaksh have no space".to_owned());
}
}
kaksh.loog.iter().for_each(|a: &Loog| {
a.addr.do_send(ms::WsConnected {
name: msg.name.to_owned(),
kunjika: msg.kunjika.to_owned()
})
});
kaksh.loog.push(Loog::new(msg.addr, msg.kunjika,msg.name, None));
}, None => { // don't exist
// add kaksh and notify
msg.addr.do_send(ms::WsConnected {
name: msg.name.to_owned(),
kunjika: msg.kunjika.to_owned()
});
self.kaksh.insert(msg.kaksh_kunjika, Kaksh {
length: msg.length,
last_message_id: 0,
loog: vec![Loog::new(msg.addr,msg.kunjika,msg.name, None)]
});
}
}
ms::Resp::Ok
}
}
/// Join random vayakti
/// Works as:
/// Check if watchlist is empty, if yes add the kunjika andaddr to watchlist
/// if watchlist have people get 0th person an connect it
impl Handler<ms::JoinRandom> for ChatPinnd {
type Result = ms::Resp;
fn handle(&mut self, msg: ms::JoinRandom, _: &mut Self::Context) -> Self::Result {
// check if user exist
if let Some(_) = self.vyaktigat_waitlist.iter().position(|vk| vk.kunjika == msg.kunjika) {
return ms::Resp::Err("Kunjika already exist".to_owned());
}
if let Some(_) = self.kaksh.iter().position(|(_,g)| {
match g.loog.iter().position(|a| a.kunjika == msg.kunjika) {
Some(_) => true,
None => false
}
}) {
return ms::Resp::Err("Kunjika already exist".to_owned());
}
// Check if watch list is empty
if self.vyaktigat_waitlist.len() == 0 {
self.vyaktigat_waitlist.push(VyaktiWatchlist {
kunjika: msg.kunjika,
addr: msg.addr,
name: msg.name,
tags: msg.tags
});
return ms::Resp::None;
}
// connect person with tag
let pos = if msg.tags.len() > 0 {
match self.vyaktigat_waitlist.iter().position(|vk| {
match vk.tags.iter().position(|t| msg.tags.contains(t)) {
Some(_) => true,
None => false
}
}) {
Some(i) => i,
None => {
self.vyaktigat_waitlist.push(VyaktiWatchlist {
kunjika: msg.kunjika,
addr: msg.addr,
name: msg.name,
tags: msg.tags
});
return ms::Resp::None;
}
}
} else { 0 };
let vayakti_watchlist = self.vyaktigat_waitlist.remove(pos);
let group_kunjika = format!("gupt_{}>{}",msg.kunjika.to_owned(), vayakti_watchlist.kunjika);
self.kaksh.insert(group_kunjika.to_owned(), Kaksh {
length: Some(2),
last_message_id: 0,
loog: vec![Loog::new(msg.addr.clone(), msg.kunjika.to_owned(), msg.name.to_owned(), Some(msg.tags.clone())),
Loog::new(vayakti_watchlist.addr.clone(), vayakti_watchlist.kunjika.to_owned(), vayakti_watchlist.name.to_owned(), Some(vayakti_watchlist.tags.clone()))]
});
// notify about connection
msg.addr.do_send(ms::WsConnectedRandom {
name: vayakti_watchlist.name,
kunjika: vayakti_watchlist.kunjika,
kaksh_kunjika: group_kunjika.to_owned()
});
vayakti_watchlist.addr.do_send(ms::WsConnectedRandom {
name: msg.name,
kunjika: msg.kunjika.to_owned(),
kaksh_kunjika: group_kunjika
});
ms::Resp::Ok
}
}
/// Next Random user
impl Handler<ms::JoinRandomNext> for ChatPinnd {
type Result = ms::Resp;
fn handle(&mut self, msg: ms::JoinRandomNext, _: &mut Self::Context) -> Self::Result {
let kaksh = match self.kaksh.get_mut(&msg.kaksh_kunjika) {
Some(v) => v,
None => return ms::Resp::Err("Failed to join, check entries!".to_owned())
};
let loog_i = match kaksh.loog.iter().position(|a| a.kunjika == msg.kunjika) {
Some(v) => v,
None => return ms::Resp::Err("Failed to join, check entries!".to_owned())
};
let addr;
let name;
let tags;
{
let loog = match kaksh.loog.get(loog_i) {
Some(v) => v,
None => return ms::Resp::Err("Failed to join, check entries!".to_owned())
};
if let None = loog.tags {
return ms::Resp::Err("You are not a randome vyakti!".to_owned());
}
addr = loog.addr.clone();
name = loog.name.to_owned();
tags = match loog.tags.clone() {
Some(v) => v,
None => return ms::Resp::Err("Failed to join, check entries!".to_owned())
};
}
// remove from old kaksh
kaksh.loog.remove(loog_i);
kaksh.loog.iter().for_each(|a| {
a.addr.do_send(ms::WsDisconnected {
kunjika: msg.kunjika.to_owned(),
name: name.to_owned()
})
});
// Check if watch list is empty
if self.vyaktigat_waitlist.len() == 0 {
self.vyaktigat_waitlist.push(VyaktiWatchlist {
kunjika: msg.kunjika,
addr,
name,
tags
});
return ms::Resp::None;
}
// connect person with tag or to zero
let pos = if tags.len() > 0 {
match self.vyaktigat_waitlist.iter().position(|vk| {
match vk.tags.iter().position(|t| tags.contains(t)) {
Some(_) => true,
None => false
}
}) {
Some(i) => i,
None => {
self.vyaktigat_waitlist.push(VyaktiWatchlist {
kunjika: msg.kunjika,
addr,
name,
tags
});
return ms::Resp::None;
}
}
} else { 0 };
let vayakti_watchlist = self.vyaktigat_waitlist.remove(pos);
let group_kunjika = format!("gupt_{}>{}",msg.kunjika.to_owned(), vayakti_watchlist.kunjika);
let log_count = kaksh.loog.len();
drop(kaksh);
if log_count == 0 {
self.kaksh.remove(&msg.kaksh_kunjika);
}
self.kaksh.insert(group_kunjika.to_owned(), Kaksh {
length: Some(2),
last_message_id: 0,
loog: vec![Loog::new(addr.clone(), msg.kunjika.to_owned(), name.to_owned(), Some(tags.clone())),
Loog::new(vayakti_watchlist.addr.clone(), vayakti_watchlist.kunjika.to_owned(), vayakti_watchlist.name.to_owned(), Some(vayakti_watchlist.tags.clone()))]
});
// notify about connection
addr.do_send(ms::WsConnectedRandom {
name: vayakti_watchlist.name,
kunjika: vayakti_watchlist.kunjika,
kaksh_kunjika: group_kunjika.to_owned()
});
vayakti_watchlist.addr.do_send(ms::WsConnectedRandom {
name,
kunjika: msg.kunjika.to_owned(),
kaksh_kunjika: group_kunjika
});
ms::Resp::Ok
}
}
/// send text to everyone
impl Handler<ms::SendText> for ChatPinnd {
type Result = ();
fn handle(&mut self, msg: ms::SendText, _: &mut Self::Context) -> Self::Result {
if let Some(kaksh) = self.kaksh.get_mut(&msg.kaksh_kunjika) {
kaksh.last_message_id += 1;
let msg_id = kaksh.last_message_id;
kaksh.loog.iter().for_each(|c| {
c.addr.do_send(ms::WsText {
sender_kunjika: msg.kunjika.to_owned(),
text: msg.text.to_owned(),
reply: msg.reply.to_owned(),
msg_id
});
});
}
}
}
/// send status to everyone
impl Handler<ms::SendStatus> for ChatPinnd {
type Result = ();
fn handle(&mut self, msg: ms::SendStatus, _: &mut Self::Context) -> Self::Result {
if let Some(kaksh) = self.kaksh.get(&msg.kaksh_kunjika) {
kaksh.loog.iter().for_each(|c| {
if c.kunjika == msg.kunjika {
return;
}
c.addr.do_send(ms::WsStatus {
sender_kunjika: msg.kunjika.to_owned(),
status: msg.status.to_owned(),
});
});
}
}
}
/// send list of users
impl Handler<ms::List> for ChatPinnd {
type Result = String;
fn handle(&mut self, msg: ms::List, _: &mut Self::Context) -> Self::Result {
if let Some(kaksh) = self.kaksh.get(&msg.kaksh_kunjika) {
let mut list = Vec::new();
for x in kaksh.loog.iter() {
list.push((x.kunjika.to_owned(),x.name.to_owned()));
}
serde_json::json!(list).to_string()
} else {
"".to_string()
}
}
}
/// Notifiy a user disconnected and trim kaksh
impl Handler<ms::LeaveUser> for ChatPinnd {
type Result = ();
fn handle(&mut self, msg: ms::LeaveUser, _: &mut Self::Context) -> Self::Result {
if let Some(kaksh_kunjika) = &msg.kaksh_kunjika {
if let Some(kaksh) = self.kaksh.get_mut(kaksh_kunjika) {
let name = if let Some(i) = kaksh.loog.iter().position(|x| x.addr == msg.addr) {
kaksh.loog.remove(i).name
} else { "".to_owned() };
if kaksh.loog.len() == 0 {
self.kaksh.remove(kaksh_kunjika);
} else {
kaksh.loog.iter().for_each(|a| {
a.addr.do_send(ms::WsDisconnected {
kunjika: msg.kunjika.to_owned(),
name: name.to_owned()
})
});
}
}
}
if let Some(i) = self.vyaktigat_waitlist.iter().position(|a| a.kunjika == msg.kunjika) {
self.vyaktigat_waitlist.remove(i);
}
}
}
impl Default for ChatPinnd {
fn default() -> Self {
ChatPinnd {
kaksh: HashMap::new(),
vyaktigat_waitlist: Vec::new()
}
}
}
impl Loog {
fn new(addr: Addr<ws_sansad::WsSansad>,
kunjika: String,
name: String,
tags: Option<Vec<String>>) -> Self {
Loog {
addr,
kunjika,
name,
tags
}
}
}
impl SystemService for ChatPinnd {}
impl Supervised for ChatPinnd {}

View File

@ -1,172 +0,0 @@
//! Messages to be sent between Actors
use actix::prelude::*;
use dev::{MessageResponse, ResponseChannel};
use crate::ws_sansad::WsSansad;
//################################################## For ChatPinnd ##################################################
/// Request to change information of vayakti to list of vayakti im ChatPind
/// Request to Kaksh with its kunjika
#[derive(Clone, Message)]
#[rtype(result = "Resp")]
pub struct JoinKaksh {
pub kaksh_kunjika: String,
pub length: Option<usize>,
pub addr: Addr<WsSansad>,
pub kunjika: String,
pub name: String,
}
/// Request to connect Random vayakti
#[derive(Clone, Message)]
#[rtype(result = "Resp")]
pub struct JoinRandom {
pub addr: Addr<WsSansad>,
pub kunjika: String,
pub name: String,
pub tags: Vec<String>,
}
/// Request to connect Random vayakti
#[derive(Clone, Message)]
#[rtype(result = "Resp")]
pub struct JoinRandomNext {
pub kaksh_kunjika: String,
pub kunjika: String
}
/// Request to send text
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct SendText {
pub kaksh_kunjika: String,
pub kunjika: String,
pub text: String,
pub reply: Option<String>,
}
/// Request to send image t
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct SendImage {
pub kaksh_kunjika: String,
pub kunjika: String,
pub part: String,
pub image_id: i32
}
// Request to send text t
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct SendStatus {
pub kaksh_kunjika: String,
pub kunjika: String,
pub status: String
}
// Request to send text t
#[derive(Clone, Message)]
#[rtype(result = "String")]
pub struct List {
pub kaksh_kunjika: String
}
/// Request to leave kaksh
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct LeaveUser {
pub kaksh_kunjika: Option<String>,
pub kunjika: String,
pub addr: Addr<WsSansad>
}
//################################################## For WsSansad ##################################################
// Request to send own kunjika hash
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsKunjikaHash {
pub kunjika: String
}
// Request to send transfer text
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsText {
pub text: String,
pub reply: Option<String>,
pub sender_kunjika: String,
pub msg_id: u128
}
// Request to send transfer text
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsImage {
pub text: String,
pub reply: Option<String>,
pub sender_kunjika: String
}
// Request to send transfer text
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsStatus {
pub status: String,
pub sender_kunjika: String
}
// Request to send transfer text
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsList {
pub json: String
}
// Notify Someone connected
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsConnected {
pub name: String,
pub kunjika: String
}
// Notify someone disconnected
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsDisconnected {
pub kunjika: String,
pub name: String
}
// Give response message
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsResponse {
pub result: String,
pub message: String
}
// Got connected to random vayakti
#[derive(Clone, Message)]
#[rtype(result = "()")]
pub struct WsConnectedRandom {
pub name: String,
pub kunjika: String,
pub kaksh_kunjika: String
}
//################################################## Helper ##################################################
#[derive(Debug)]
pub enum Resp {
Ok,
Err(String),
None
}
impl<A, M> MessageResponse<A, M> for Resp
where
A: Actor,
M: Message<Result = Resp>,
{
fn handle<R: ResponseChannel<M>>(self, _: &mut A::Context, tx: Option<R>) {
if let Some(tx) = tx {
tx.send(self);
}
}
}

View File

@ -1,592 +0,0 @@
//! Ws Sansad manage websocket of each client
use actix::prelude::*;
use actix_broker::{Broker, SystemBroker};
use actix_web_actors::ws;
use ms::Resp;
use serde_json::{json, Value};
use std::time::{Duration, Instant};
use crate::{chat_pinnd::ChatPinnd, messages as ms, validator::{Validation as vl, validate}};
/// How often heartbeat pings are sent
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(5);
/// How long before lack of client response causes a timeout
const CLIENT_TIMEOUT: Duration = Duration::from_secs(15);
// for phones if browser kept websocket on
const SPECIAL_HEARTBEAT_INTERVAL: Duration = Duration::from_secs(5*60);
const SPECIAL_CLIENT_TIMEOUT: Duration = Duration::from_secs(10*60);
pub struct WsSansad {
kunjika: String,
isthiti: Isthiti,
addr: Option<Addr<Self>>,
hb: Instant,
special_hb: Instant
}
#[derive(Debug)]
enum Isthiti {
None,
Kaksh(String),
VraktigatWaitlist
}
impl Actor for WsSansad {
type Context = ws::WebsocketContext<Self>;
fn started(&mut self, ctx: &mut Self::Context) {
self.addr = Some(ctx.address().clone()); // own addr
self.hb(ctx);
self.special_hb(ctx);
}
fn stopping(&mut self, _: &mut Self::Context) -> Running {
futures::executor::block_on(self.leave_kaksh()); // notify leaving
Running::Stop
}
}
/// manage stream
impl StreamHandler<Result<ws::Message, ws::ProtocolError>> for WsSansad {
fn handle(&mut self, msg: Result<ws::Message, ws::ProtocolError>, ctx: &mut Self::Context) {
match msg {
Ok(ws::Message::Ping(msg)) => {
ctx.ping(&msg);
self.hb = Instant::now();
}, Ok(ws::Message::Pong(_)) => {
self.hb = Instant::now();
}, Ok(ws::Message::Text(msg)) => {
self.special_hb = Instant::now();
futures::executor::block_on(self.parse_text_handle(msg));
}, Ok(ws::Message::Close(msg)) => {
ctx.close(msg);
ctx.stop();
}
_ => ctx.stop()
}
}
}
/// send text message
impl Handler<ms::WsText> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsText, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "text",
"text": msg.text,
"reply": msg.reply,
"kunjika": msg.sender_kunjika, // Sender's kunjuka
"msg_id": msg.msg_id.to_string()
});
ctx.text(json.to_string());
}
}
/// send text status
impl Handler<ms::WsStatus> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsStatus, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "status",
"status": msg.status,
"kunjika": msg.sender_kunjika // Sender's kunjuka
});
ctx.text(json.to_string());
}
}
/// List Vayakti
impl Handler<ms::WsList> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsList, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "list",
"vayakti": msg.json
});
ctx.text(json.to_string());
}
}
/// Own Kunjika hash
impl Handler<ms::WsKunjikaHash> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsKunjikaHash, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "kunjika",
"kunjika": msg.kunjika
});
ctx.text(json.to_string());
}
}
/// send response ok, error
impl Handler<ms::WsResponse> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsResponse, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "resp",
"result": msg.result,
"message": msg.message
});
ctx.text(json.to_string());
}
}
/// notify someone got connected
impl Handler<ms::WsConnected> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsConnected, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "connected",
"name": msg.name,
"kunjika": msg.kunjika
});
ctx.text(json.to_string());
}
}
/// notify someone got disconnected
impl Handler<ms::WsDisconnected> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsDisconnected, ctx: &mut Self::Context) -> Self::Result {
let json = json!({
"cmd": "disconnected",
"name": msg.name,
"kunjika": msg.kunjika
});
ctx.text(json.to_string());
}
}
/// notify got connected to random person
impl Handler<ms::WsConnectedRandom> for WsSansad {
type Result = ();
fn handle(&mut self, msg: ms::WsConnectedRandom, ctx: &mut Self::Context) -> Self::Result {
self.isthiti = Isthiti::Kaksh(msg.kaksh_kunjika);
let json = json!({
"cmd": "random",
"name": msg.name,
"kunjika": msg.kunjika
});
ctx.text(json.to_string());
}
}
impl WsSansad {
pub fn new() -> Self {
WsSansad {
kunjika: String::new(),
isthiti: Isthiti::None,
addr: None,
hb: Instant::now(),
special_hb: Instant::now()
}
}
/// helper method that sends ping to client every second.
///
/// also this method checks heartbeats from client
fn hb(&self, ctx: &mut <Self as Actor>::Context) {
ctx.run_interval(HEARTBEAT_INTERVAL, |act, ctx| {
// check client heartbeats
if Instant::now().duration_since(act.hb) > CLIENT_TIMEOUT {
// heartbeat timed out
// stop actor
futures::executor::block_on(act.leave_kaksh()); // notify leaving
ctx.stop();
// don't try to send a ping
return;
}
ctx.ping(b"");
});
}
/// helper method that sends ping to client every second.
///
/// also this method checks heartbeats from client
fn special_hb(&self, ctx: &mut <Self as Actor>::Context) {
ctx.run_interval(SPECIAL_HEARTBEAT_INTERVAL, |act, ctx| {
// check client heartbeats
if Instant::now().duration_since(act.special_hb) > SPECIAL_CLIENT_TIMEOUT {
// heartbeat timed out
// stop actor
futures::executor::block_on(act.leave_kaksh()); // notify leaving
ctx.stop();
// don't try to send a ping
return;
}
});
}
/// parse the request text from client
async fn parse_text_handle(&mut self, msg: String) {
if let Ok(val) = serde_json::from_str::<Value>(&msg) {
// let cmd = match val.get("cmd") {
// Some(v) => v,
// None => return
// };
// let cmd = match cmd.as_str() {
// Some(v) => v,
// None => return
// };
match val.get("cmd").unwrap().as_str().unwrap() {
"join" => { self.join_kaksh(val).await },
"rand" => { self.join_random(val).await },
"randnext" => { self.join_random_next().await },
"text" => { self.send_text(val).await },
"status" => { self.send_status(val).await },
"list" => { self.list().await },
"leave" => { self.leave_kaksh().await },
_ => ()
}
}
}
/// send ok response
fn send_ok_response(&self, text: &str) {
self.addr.clone().unwrap().do_send(ms::WsResponse {
result: "Ok".to_owned(),
message: text.to_owned()
});
}
/// send error response
fn send_err_response(&self, text: &str) {
self.addr.clone().unwrap().do_send(ms::WsResponse {
result: "Err".to_owned(),
message: text.to_owned()
});
}
/// Request for joining to random person
async fn join_random(&mut self, val: Value) {
// Check is already joined
match self.isthiti {
Isthiti::None => (),
Isthiti::VraktigatWaitlist => {
self.send_ok_response("watchlist");
return;
}, Isthiti::Kaksh(_) => return
}
let kunjika = match val.get("kunjika") {
Some(val ) => val.as_str().unwrap().to_owned(),
None => {
self.send_err_response("Invalid request");
return;
}
};
// kunjika to hash
let mut m = sha1::Sha1::new();
m.update(format!("{}{}",kunjika,
std::env::var("SALT").unwrap_or("".to_owned())).as_bytes());
let kunjika = m.digest().to_string();
let name = match val.get("name") {
Some(val ) => val.as_str().unwrap().to_owned(),
None => {
self.send_err_response("Invalid request");
return;
}
};
let tags = match val.get("tags") {
Some(val ) => {
let mut v = Vec::new();
for x in val.as_str().unwrap().split_ascii_whitespace() {
v.push(x.to_owned());
}
v
},
None => {
Vec::new()
}
};
// Validate
if let Some(val ) = validate(vec![vl::NonEmpty, vl::NoSpace, vl::NoHashtag], &kunjika, "Kunjika") {
self.send_err_response(&val);
return;
} else if let Some(val ) = validate(vec![vl::NonEmpty], &name, "Name") {
self.send_err_response(&val);
return;
}
// request
let result: Resp = ChatPinnd::from_registry().send(ms::JoinRandom{
addr: self.addr.clone().unwrap(),
kunjika: kunjika.to_owned(),
name,
tags
}).await.unwrap();
match result {
Resp::Err(err) => self.send_err_response(&err),
Resp::Ok => {
self.addr.clone().unwrap().do_send(ms::WsKunjikaHash{ kunjika: kunjika.clone() });
self.kunjika = kunjika;
},
Resp::None => {
self.addr.clone().unwrap().do_send(ms::WsResponse{
result: "watch".to_owned() ,
message: "Watchlist".to_owned()
});
self.isthiti = Isthiti::VraktigatWaitlist;
self.addr.clone().unwrap().do_send(ms::WsKunjikaHash{ kunjika: kunjika.clone() });
self.kunjika = kunjika
}
}
}
/// Request for joining to random person
async fn join_random_next(&mut self) {
// Check is already joined
let kaksh_kunjika = match &self.isthiti {
Isthiti::VraktigatWaitlist => {
self.send_ok_response("watchlist");
return;
},
Isthiti::Kaksh(kaksh_kunjika) => kaksh_kunjika,
Isthiti::None => {
self.send_ok_response("Not allowed");
return;
}
};
// request
let result: Resp = ChatPinnd::from_registry().send(ms::JoinRandomNext {
kunjika: self.kunjika.to_owned(),
kaksh_kunjika: kaksh_kunjika.to_owned(),
}).await.unwrap();
match result {
Resp::Err(err) => self.send_err_response(&err),
Resp::None => {
self.addr.clone().unwrap().do_send(ms::WsResponse{
result: "watch".to_owned() ,
message: "Watchlist".to_owned()
});
self.isthiti = Isthiti::VraktigatWaitlist;
self.kunjika = self.kunjika.to_owned()
}
_ => ()
}
}
/// Request to join to kaksh
async fn join_kaksh(&mut self, val: Value) {
// Check is already joined
match self.isthiti {
Isthiti::None => (),
_ => return
}
// is vayakti in watch list
if let Isthiti::VraktigatWaitlist = self.isthiti {
self.send_ok_response("watchlist");
return;
}
let kunjika = match val.get("kunjika") {
Some(val ) => val.as_str().unwrap().to_owned(),
None => {
self.send_err_response("Invalid request");
return;
}
};
// kunjika to hash
let mut m = sha1::Sha1::new();
m.update(format!("{}{}",kunjika,
std::env::var("SALT").unwrap_or("".to_owned())).as_bytes());
let kunjika = m.digest().to_string();
let name = match val.get("name") {
Some(val ) => val.as_str().unwrap().to_owned(),
None => {
self.send_err_response("Invalid request");
return;
}
};
let kaksh_kunjika = match val.get("kaksh_kunjika") {
Some(val ) => val.as_str().unwrap().to_owned(),
None => {
self.send_err_response("Invalid request");
return;
}
};
let length: Option<usize> = match val.get("length") {
Some(val) => match val.as_i64(){
Some(val) => Some(val as usize),
None => None
},
None => None
};
// Validate
if let Some(val ) = validate(vec![vl::NonEmpty, vl::NoGupt, vl::NoSpace], &kaksh_kunjika, "Kaksh Kunjika") {
self.send_err_response(&val);
return;
} else if let Some(val ) = validate(vec![vl::NonEmpty, vl::NoSpace, vl::NoHashtag], &kunjika, "Kunjika") {
self.send_err_response(&val);
return;
} else if let Some(val ) = validate(vec![vl::NonEmpty], &name, "Name") {
self.send_err_response(&val);
return;
}
// request
let result: Resp = ChatPinnd::from_registry().send(ms::JoinKaksh {
kaksh_kunjika: kaksh_kunjika.to_owned(),
length,
addr: self.addr.clone().unwrap(),
kunjika: kunjika.to_owned(),
name
}).await.unwrap();
match result {
Resp::Err(err) => self.send_err_response(&err),
Resp::Ok => {
self.isthiti = Isthiti::Kaksh(kaksh_kunjika);
self.addr.clone().unwrap().do_send(ms::WsKunjikaHash{ kunjika: kunjika.clone() });
self.kunjika = kunjika;
self.send_ok_response("joined")
}
_ => ()
}
}
/// Request to join to kaksh
async fn list(&mut self) {
// check if vayakti exist
if let Isthiti::None = self.isthiti {
self.send_err_response("Not in any Kaksh");
return;
}
// check if connected to any kaksh
match &self.isthiti {
Isthiti::Kaksh(kunjika) => {
let json: String = ChatPinnd::from_registry().send(ms::List {
kaksh_kunjika: kunjika.to_owned()
}).await.unwrap();
self.addr.clone().unwrap().do_send(ms::WsList {
json
})
},
_ => {
self.send_err_response("Kaksh not connected");
return;
}
}
}
/// send text to vayakti in kaksh
async fn send_text(&mut self, val: Value) {
// check if vayakti exist
if let Isthiti::None = self.isthiti {
self.send_err_response("Not in any Kaksh");
return;
}
// check if connected to any kaksh
match self.isthiti {
Isthiti::Kaksh(_) => (),
_ => {
self.send_err_response("Kaksh not connected");
return;
}
}
// sent text
let text = match val.get("text") {
Some(val) => val,
None => {
self.send_err_response("Invalid request");
return;
}
}.as_str().unwrap().to_owned();
let reply: Option<String> = match val.get("reply") {
Some(val) => Some(val.as_str().unwrap().to_owned()),
None => None
};
let kaksh_kunjika = match &self.isthiti {
Isthiti::Kaksh(kaksh_kunjika) => {
kaksh_kunjika.to_owned()
}, _ => {
return;
}
};
Broker::<SystemBroker>::issue_async(ms::SendText {
kaksh_kunjika,
kunjika: self.kunjika.to_owned(),
text,
reply
});
}
/// send status to vayakti in kaksh
async fn send_status(&mut self, val: Value) {
// check if vayakti exist
if let Isthiti::None = self.isthiti {
self.send_err_response("Not in any Kaksh");
return;
}
// check if connected to any kaksh
match self.isthiti {
Isthiti::Kaksh(_) => (),
_ => {
self.send_err_response("Kaksh not connected");
return;
}
}
// sent status
let status = match val.get("status") {
Some(val) => val,
None => {
self.send_err_response("Invalid request");
return;
}
}.as_str().unwrap().to_owned();
let kaksh_kunjika = match &self.isthiti {
Isthiti::Kaksh(kaksh_kunjika) => {
kaksh_kunjika.to_owned()
}, _ => {
return;
}
};
Broker::<SystemBroker>::issue_async(ms::SendStatus {
kaksh_kunjika,
kunjika: self.kunjika.to_owned(),
status
});
}
/// notify leaving
async fn leave_kaksh(&mut self) {
let kaksh_kunjika = match &self.isthiti {
Isthiti::Kaksh(val) => Some(val.to_owned()),
_ => None
};
Broker::<SystemBroker>::issue_async(ms::LeaveUser {
kaksh_kunjika,
kunjika: self.kunjika.to_owned(),
addr: self.addr.clone().unwrap()
});
self.isthiti = Isthiti::None;
self.send_ok_response("left");
}
}