-
Notifications
You must be signed in to change notification settings - Fork 7
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Gateway: implement backend-doubler (#206)
- Loading branch information
Showing
26 changed files
with
768 additions
and
261 deletions.
There are no files selected for viewing
Large diffs are not rendered by default.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,7 +1,7 @@ | ||
[package] | ||
name = "gateway" | ||
version = "0.12.2" | ||
authors = ["Terkwood <metaterkhorn@gmail.com>"] | ||
version = "0.12.3" | ||
authors = ["terkwood <[email protected].com>"] | ||
edition = "2018" | ||
|
||
[dependencies] | ||
|
@@ -10,8 +10,8 @@ crossbeam-channel = "0.4.0" | |
mio-extras = "2.0.6" | ||
rand = "0.7.2" | ||
rdkafka = "0.23.0" | ||
serde = "1.0.103" | ||
serde_derive = "1.0.103" | ||
serde = "1.0.106" | ||
serde_derive = "1.0.106" | ||
serde_json = "1.0.44" | ||
time = "0.1.42" | ||
uuid = { version = "0.8.1", features = ["v4", "serde"] } | ||
|
@@ -23,6 +23,10 @@ dotenv = "0.15.0" | |
futures = "0.3.0" | ||
chrono = { version = "0.4.10", features = ["serde"] } | ||
r2d2_redis = "0.12.0" | ||
log = "0.4.8" | ||
env_logger = "0.7.1" | ||
micro_model_moves = { git = "https://github.com/Terkwood/BUGOUT", branch = "unstable" } | ||
micro_model_bot = { git = "https://github.com/Terkwood/BUGOUT", branch = "unstable" } | ||
|
||
[dev-dependencies] | ||
rand = "0.7.2" |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -24,4 +24,6 @@ COPY . /var/BUGOUT/gateway/. | |
|
||
RUN cargo install --path . | ||
|
||
ENV RUST_LOG info | ||
|
||
CMD ["gateway"] |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,121 @@ | ||
use crate::backend_commands::BackendCommands; | ||
|
||
use crossbeam_channel::{select, Receiver, Sender}; | ||
use log::error; | ||
pub fn double_commands(opts: DoublerOpts) { | ||
loop { | ||
select! { | ||
recv(opts.session_commands_out) -> msg => match msg { | ||
Ok( backend_command ) => { | ||
if let Err(e) = opts.redis_commands_in.send(backend_command.clone()) { | ||
error!("err doubler 0 {:?}",e) | ||
} | ||
|
||
if let Err(e) = opts.kafka_commands_in.send(backend_command) { | ||
error!("FAILED doubler TO BACKEND {:?}", e) | ||
} | ||
} | ||
Err(e) => error!("session command out: {:?}",e) | ||
} | ||
} | ||
} | ||
} | ||
|
||
pub struct DoublerOpts { | ||
pub session_commands_out: Receiver<BackendCommands>, | ||
pub kafka_commands_in: Sender<BackendCommands>, | ||
pub redis_commands_in: Sender<BackendCommands>, | ||
} | ||
#[cfg(test)] | ||
mod tests { | ||
use super::*; | ||
use crate::backend_commands::*; | ||
use crate::model::*; | ||
|
||
use crossbeam_channel::{select, unbounded}; | ||
use std::thread; | ||
use uuid::Uuid; | ||
|
||
#[test] | ||
fn test_double_commands() { | ||
let (session_commands_in, session_commands_out): ( | ||
Sender<BackendCommands>, | ||
Receiver<BackendCommands>, | ||
) = unbounded(); | ||
|
||
let (kafka_commands_in, kafka_commands_out): ( | ||
Sender<BackendCommands>, | ||
Receiver<BackendCommands>, | ||
) = unbounded(); | ||
|
||
let (redis_commands_in, redis_commands_out): ( | ||
Sender<BackendCommands>, | ||
Receiver<BackendCommands>, | ||
) = unbounded(); | ||
thread::spawn(move || { | ||
let opts = DoublerOpts { | ||
kafka_commands_in, | ||
redis_commands_in, | ||
session_commands_out, | ||
}; | ||
double_commands(opts) | ||
}); | ||
|
||
{ | ||
let session_id = Uuid::new_v4(); | ||
let client_id = Uuid::new_v4(); | ||
session_commands_in | ||
.send(BackendCommands::FindPublicGame( | ||
FindPublicGameBackendCommand { | ||
session_id, | ||
client_id, | ||
}, | ||
)) | ||
.expect("send0") | ||
} | ||
|
||
{ | ||
let game_id = micro_model_moves::GameId(Uuid::new_v4()); | ||
let player = Player::WHITE; | ||
session_commands_in | ||
.send(BackendCommands::AttachBot( | ||
micro_model_bot::gateway::AttachBot { | ||
game_id, | ||
player: match player { | ||
Player::WHITE => micro_model_moves::Player::WHITE, | ||
_ => micro_model_moves::Player::BLACK, | ||
}, | ||
}, | ||
)) | ||
.expect("send1") | ||
} | ||
|
||
select! { recv(kafka_commands_out) -> co => | ||
match co.expect("kafka co 0 select") { | ||
BackendCommands::FindPublicGame(_) => assert!(true), | ||
_ => assert!(false) | ||
} | ||
} | ||
|
||
select! { recv(redis_commands_out) -> co => | ||
match co.expect("redis co 0 select") { | ||
BackendCommands::FindPublicGame(_) => assert!(true), | ||
_ => assert!(false) | ||
} | ||
} | ||
|
||
select! { recv(kafka_commands_out) -> co => | ||
match co.expect("kafka co 1 select") { | ||
BackendCommands::AttachBot(_) => assert!(true), | ||
_ => assert!(false) | ||
} | ||
} | ||
|
||
select! { recv(redis_commands_out) -> co => | ||
match co.expect("redis co 1 select") { | ||
BackendCommands::AttachBot(_) => assert!(true), | ||
_ => assert!(false) | ||
} | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,7 @@ | ||
mod doubler; | ||
mod start; | ||
|
||
use crate::backend_events::{BackendEvents, KafkaShutdownEvent}; | ||
use crate::idle_status::KafkaActivityObserved; | ||
pub use doubler::double_commands; | ||
pub use start::{start_all, BackendInitOptions}; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,49 @@ | ||
use crate::backend_commands::BackendCommands; | ||
use crate::kafka_io; | ||
use crate::redis_io; | ||
|
||
use crossbeam_channel::{unbounded, Receiver, Sender}; | ||
use futures::executor::block_on; | ||
use std::thread; | ||
|
||
use super::*; | ||
|
||
pub fn start_all(opts: BackendInitOptions) { | ||
let (kafka_commands_in, kafka_commands_out): ( | ||
Sender<BackendCommands>, | ||
Receiver<BackendCommands>, | ||
) = unbounded(); | ||
|
||
let (redis_commands_in, redis_commands_out): ( | ||
Sender<BackendCommands>, | ||
Receiver<BackendCommands>, | ||
) = unbounded(); | ||
|
||
thread::spawn(move || redis_io::xadd_commands(redis_commands_out, &redis_io::create_pool())); | ||
|
||
let bei = opts.backend_events_in.clone(); | ||
thread::spawn(move || redis_io::stream::process(bei)); | ||
|
||
let soc = opts.session_commands_out; | ||
thread::spawn(move || { | ||
double_commands(super::doubler::DoublerOpts { | ||
session_commands_out: soc, | ||
kafka_commands_in, | ||
redis_commands_in, | ||
}) | ||
}); | ||
|
||
block_on(kafka_io::start( | ||
opts.backend_events_in.clone(), | ||
opts.shutdown_in.clone(), | ||
opts.kafka_activity_in.clone(), | ||
kafka_commands_out, | ||
)) | ||
} | ||
|
||
pub struct BackendInitOptions { | ||
pub backend_events_in: Sender<BackendEvents>, | ||
pub shutdown_in: Sender<KafkaShutdownEvent>, | ||
pub kafka_activity_in: Sender<KafkaActivityObserved>, | ||
pub session_commands_out: Receiver<BackendCommands>, | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.