mirror of
https://github.com/4yn/slidershim.git
synced 2026-09-07 08:15:33 -05:00
brokenithm working
This commit is contained in:
@@ -5,41 +5,41 @@ use std::{future::Future, io, time::Duration};
|
||||
|
||||
use tokio::{select, time::sleep};
|
||||
|
||||
use slidershim::slider_io::worker::{AsyncJob, AsyncWorker};
|
||||
// use slidershim::slider_io::worker::{AsyncJob, AsyncWorker};
|
||||
|
||||
struct CounterJob;
|
||||
// struct CounterJob;
|
||||
|
||||
#[async_trait]
|
||||
impl AsyncJob for CounterJob {
|
||||
async fn do_work<F: Future<Output = ()> + Send>(self, stop_signal: F) {
|
||||
let job_a = async {
|
||||
println!("Start job A");
|
||||
let mut x = 0;
|
||||
loop {
|
||||
x += 1;
|
||||
println!("{}", x);
|
||||
sleep(Duration::from_millis(100)).await;
|
||||
}
|
||||
};
|
||||
let job_b = async move {
|
||||
println!("Start job B");
|
||||
stop_signal.await;
|
||||
println!("Stop signal hit at job B");
|
||||
};
|
||||
// #[async_trait]
|
||||
// impl AsyncJob for CounterJob {
|
||||
// async fn run<F: Future<Output = ()> + Send>(self, stop_signal: F) {
|
||||
// let job_a = async {
|
||||
// println!("Start job A");
|
||||
// let mut x = 0;
|
||||
// loop {
|
||||
// x += 1;
|
||||
// println!("{}", x);
|
||||
// sleep(Duration::from_millis(100)).await;
|
||||
// }
|
||||
// };
|
||||
// let job_b = async move {
|
||||
// println!("Start job B");
|
||||
// stop_signal.await;
|
||||
// println!("Stop signal hit at job B");
|
||||
// };
|
||||
|
||||
select! {
|
||||
_ = job_a => {},
|
||||
_ = job_b => {},
|
||||
}
|
||||
}
|
||||
}
|
||||
// select! {
|
||||
// _ = job_a => {},
|
||||
// _ = job_b => {},
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
|
||||
fn main() {
|
||||
env_logger::Builder::new()
|
||||
.filter_level(log::LevelFilter::Debug)
|
||||
.init();
|
||||
|
||||
let worker = AsyncWorker::new("counter", CounterJob);
|
||||
// let worker = AsyncWorker::new("counter", CounterJob);
|
||||
let mut input = String::new();
|
||||
let string = io::stdin().read_line(&mut input).unwrap();
|
||||
}
|
||||
|
||||
@@ -4,14 +4,16 @@ use std::{io, time::Duration};
|
||||
|
||||
use tokio::time::sleep;
|
||||
|
||||
// use slidershim::slider_io::{brokenithm::BrokenithmJob, worker::AsyncWorker};
|
||||
use slidershim::slider_io::{
|
||||
brokenithm::BrokenithmJob, controller_state::FullState, worker::AsyncWorker,
|
||||
};
|
||||
|
||||
fn main() {
|
||||
env_logger::Builder::new()
|
||||
.filter_level(log::LevelFilter::Debug)
|
||||
.init();
|
||||
|
||||
// let worker = AsyncWorker::new("brokenithm", BrokenithmJob);
|
||||
let worker = AsyncWorker::new("brokenithm", BrokenithmJob::new(FullState::new()));
|
||||
let mut input = String::new();
|
||||
let string = io::stdin().read_line(&mut input).unwrap();
|
||||
}
|
||||
|
||||
@@ -1,3 +1,44 @@
|
||||
use std::{
|
||||
env, fs,
|
||||
path::{Path, PathBuf},
|
||||
};
|
||||
|
||||
const COPY_DIR: &'static str = "res";
|
||||
|
||||
fn copy_dir<P, Q>(from: P, to: Q)
|
||||
where
|
||||
P: AsRef<Path>,
|
||||
Q: AsRef<Path>,
|
||||
{
|
||||
// https://stackoverflow.com/a/68950006
|
||||
let to = to.as_ref().to_path_buf();
|
||||
|
||||
for path in fs::read_dir(from).unwrap() {
|
||||
let path = path.unwrap().path();
|
||||
let to = to.clone().join(path.file_name().unwrap());
|
||||
|
||||
if path.is_file() {
|
||||
fs::copy(&path, to).unwrap();
|
||||
} else if path.is_dir() {
|
||||
if !to.exists() {
|
||||
fs::create_dir(&to).unwrap();
|
||||
}
|
||||
|
||||
copy_dir(&path, to);
|
||||
} else { /* Skip other content */
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn main() {
|
||||
let out = env::var("PROFILE").unwrap();
|
||||
let out = PathBuf::from(format!("target/{}/{}", out, COPY_DIR));
|
||||
|
||||
if out.exists() {
|
||||
fs::remove_dir_all(&out).unwrap()
|
||||
};
|
||||
fs::create_dir(&out).unwrap();
|
||||
copy_dir(COPY_DIR, &out);
|
||||
|
||||
tauri_build::build();
|
||||
}
|
||||
|
||||
@@ -32,8 +32,12 @@ fn main() {
|
||||
env_logger::init();
|
||||
|
||||
let config = Arc::new(Mutex::new(Some(slider_io::Config::default())));
|
||||
let manager: slider_io::Manager;
|
||||
{
|
||||
config.lock().unwrap().as_ref().unwrap().save();
|
||||
let c = config.lock().unwrap();
|
||||
let cr = c.as_ref().unwrap();
|
||||
cr.save();
|
||||
manager = slider_io::Manager::new(cr.clone());
|
||||
}
|
||||
|
||||
tauri::Builder::default()
|
||||
|
||||
@@ -1,50 +1,213 @@
|
||||
use std::{convert::Infallible, net::SocketAddr};
|
||||
|
||||
use log::info;
|
||||
use tokio::time::sleep;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use futures::{SinkExt, StreamExt};
|
||||
use hyper::{
|
||||
header,
|
||||
server::conn::AddrStream,
|
||||
service::{make_service_fn, service_fn},
|
||||
Body, Request, Response, Server,
|
||||
upgrade::{self, Upgraded},
|
||||
Body, Method, Request, Response, Server, StatusCode,
|
||||
};
|
||||
use log::{error, info};
|
||||
use path_clean::PathClean;
|
||||
use std::{convert::Infallible, future::Future, net::SocketAddr, path::PathBuf};
|
||||
use tokio::{fs::File, select};
|
||||
use tokio_tungstenite::WebSocketStream;
|
||||
use tokio_util::codec::{BytesCodec, FramedRead};
|
||||
use tungstenite::{handshake, Error, Message};
|
||||
|
||||
// use crate::slider_io::worker::{AsyncJob, AsyncJobFut, AsyncJobRecvStop};
|
||||
use crate::slider_io::{controller_state::FullState, worker::AsyncJob};
|
||||
|
||||
// https://levelup.gitconnected.com/handling-websocket-and-http-on-the-same-port-with-rust-f65b770722c9
|
||||
|
||||
async fn error_response() -> Result<Response<Body>, Infallible> {
|
||||
Ok(
|
||||
Response::builder()
|
||||
.status(StatusCode::NOT_FOUND)
|
||||
.body(Body::from(format!("Not found")))
|
||||
.unwrap(),
|
||||
)
|
||||
}
|
||||
|
||||
async fn serve_file(path: &str) -> Result<Response<Body>, Infallible> {
|
||||
let mut pb = PathBuf::from("res/www/");
|
||||
pb.push(path);
|
||||
pb.clean();
|
||||
|
||||
// println!("CWD {:?}", std::env::current_dir());
|
||||
// println!("Serving file {:?}", pb);
|
||||
|
||||
match File::open(pb).await {
|
||||
Ok(f) => {
|
||||
let stream = FramedRead::new(f, BytesCodec::new());
|
||||
let body = Body::wrap_stream(stream);
|
||||
Ok(Response::new(body))
|
||||
}
|
||||
Err(_) => error_response().await,
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_brokenithm(ws_stream: WebSocketStream<Upgraded>, state: FullState) {
|
||||
let (mut ws_write, mut ws_read) = ws_stream.split();
|
||||
|
||||
loop {
|
||||
match ws_read.next().await {
|
||||
Some(msg) => match msg {
|
||||
Ok(msg) => match msg {
|
||||
Message::Text(msg) => {
|
||||
let mut chars = msg.chars();
|
||||
let head = chars.next().unwrap();
|
||||
match head {
|
||||
'a' => {
|
||||
ws_write.send(Message::Text("alive".to_string())).await;
|
||||
}
|
||||
'b' => {
|
||||
let flat_state: Vec<bool> = chars
|
||||
.map(|x| match x {
|
||||
'0' => false,
|
||||
'1' => true,
|
||||
_ => unreachable!(),
|
||||
})
|
||||
.collect();
|
||||
let mut controller_state_handle = state.controller_state.lock().unwrap();
|
||||
for (idx, c) in flat_state[0..32].iter().enumerate() {
|
||||
controller_state_handle.ground_state[idx] = match c {
|
||||
false => 0,
|
||||
true => 255,
|
||||
}
|
||||
}
|
||||
for (idx, c) in flat_state[32..38].iter().enumerate() {
|
||||
controller_state_handle.air_state[idx] = match c {
|
||||
false => 0,
|
||||
true => 1,
|
||||
}
|
||||
}
|
||||
// println!(
|
||||
// "{:?} {:?}",
|
||||
// controller_state_handle.ground_state, controller_state_handle.air_state
|
||||
// );
|
||||
}
|
||||
_ => {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
Message::Close(_) => {
|
||||
info!("Websocket connection closed");
|
||||
break;
|
||||
}
|
||||
_ => {}
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Websocket connection error: {}", e);
|
||||
break;
|
||||
}
|
||||
},
|
||||
None => {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_websocket(
|
||||
mut request: Request<Body>,
|
||||
state: FullState,
|
||||
) -> Result<Response<Body>, Infallible> {
|
||||
let res = match handshake::server::create_response_with_body(&request, || Body::empty()) {
|
||||
Ok(res) => {
|
||||
tokio::spawn(async move {
|
||||
match upgrade::on(&mut request).await {
|
||||
Ok(upgraded) => {
|
||||
let ws_stream = WebSocketStream::from_raw_socket(
|
||||
upgraded,
|
||||
tokio_tungstenite::tungstenite::protocol::Role::Server,
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
|
||||
handle_brokenithm(ws_stream, state).await;
|
||||
}
|
||||
|
||||
Err(e) => {
|
||||
error!("Websocket upgrade error: {}", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
res
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Websocket creation error: {}", e);
|
||||
Response::builder()
|
||||
.status(StatusCode::BAD_REQUEST)
|
||||
.body(Body::from(format!("Failed to create websocket: {}", e)))
|
||||
.unwrap()
|
||||
}
|
||||
};
|
||||
Ok(res)
|
||||
}
|
||||
|
||||
async fn handle_request(
|
||||
request: Request<Body>,
|
||||
remote_addr: SocketAddr,
|
||||
state: FullState,
|
||||
) -> Result<Response<Body>, Infallible> {
|
||||
Ok(Response::new(Body::from(format!(
|
||||
"Hello there connection {}\n",
|
||||
remote_addr
|
||||
))))
|
||||
}
|
||||
let method = request.method();
|
||||
let path = request.uri().path();
|
||||
if method != Method::GET {
|
||||
error!("Server unknown method {} {}", method, path);
|
||||
return error_response().await;
|
||||
}
|
||||
info!("Server {} {}", method, path);
|
||||
|
||||
async fn brokenithm_server() {
|
||||
let addr = SocketAddr::from(([0, 0, 0, 0], 1666));
|
||||
|
||||
info!("Brokenithm opening on {:?}", addr);
|
||||
|
||||
let make_svc = make_service_fn(|conn: &AddrStream| {
|
||||
let remote_addr = conn.remote_addr();
|
||||
async move {
|
||||
Ok::<_, Infallible>(service_fn(move |request: Request<Body>| {
|
||||
handle_request(request, remote_addr)
|
||||
}))
|
||||
}
|
||||
});
|
||||
|
||||
let server = Server::bind(&addr).serve(make_svc);
|
||||
if let Err(e) = server.await {
|
||||
eprintln!("Server error: {}", e);
|
||||
match (
|
||||
request.uri().path(),
|
||||
request.headers().contains_key(header::UPGRADE),
|
||||
) {
|
||||
("/", false) | ("/index.html", false) => serve_file("index.html").await,
|
||||
(filename, false) => serve_file(&filename[1..]).await,
|
||||
("/ws", true) => handle_websocket(request, state).await,
|
||||
_ => error_response().await,
|
||||
}
|
||||
}
|
||||
|
||||
// struct BrokenithmJob;
|
||||
pub struct BrokenithmJob {
|
||||
state: FullState,
|
||||
}
|
||||
|
||||
// impl AsyncJob {
|
||||
// fn job(self, mut recv_stop: AsyncJobRecvStop) -> AsyncJobFut {
|
||||
// return Box::pin()
|
||||
// }
|
||||
// }
|
||||
impl BrokenithmJob {
|
||||
pub fn new(state: FullState) -> Self {
|
||||
Self { state }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl AsyncJob for BrokenithmJob {
|
||||
async fn run<F: Future<Output = ()> + Send>(self, stop_signal: F) {
|
||||
let state = self.state.clone();
|
||||
let make_svc = make_service_fn(|conn: &AddrStream| {
|
||||
let remote_addr = conn.remote_addr();
|
||||
let make_svc_state = state.clone();
|
||||
async move {
|
||||
Ok::<_, Infallible>(service_fn(move |request: Request<Body>| {
|
||||
let svc_state = make_svc_state.clone();
|
||||
handle_request(request, remote_addr, svc_state)
|
||||
}))
|
||||
}
|
||||
});
|
||||
|
||||
let addr = SocketAddr::from(([0, 0, 0, 0], 1606));
|
||||
info!("Brokenithm server listening on {}", addr);
|
||||
|
||||
let server = Server::bind(&addr)
|
||||
// .http1_keepalive(false)
|
||||
// .http2_keep_alive_interval(None)
|
||||
// .tcp_keepalive(None)
|
||||
.serve(make_svc)
|
||||
.with_graceful_shutdown(stop_signal);
|
||||
|
||||
if let Err(e) = server.await {
|
||||
info!("Brokenithm server stopped: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,10 +2,10 @@ mod config;
|
||||
mod utils;
|
||||
pub mod worker;
|
||||
|
||||
mod controller_state;
|
||||
pub mod controller_state;
|
||||
mod voltex;
|
||||
|
||||
mod brokenithm;
|
||||
pub mod brokenithm;
|
||||
mod gamepad;
|
||||
mod keyboard;
|
||||
|
||||
|
||||
@@ -68,7 +68,7 @@ impl Drop for ThreadWorker {
|
||||
|
||||
#[async_trait]
|
||||
pub trait AsyncJob: Send + 'static {
|
||||
async fn do_work<F: Future<Output = ()> + Send>(self, stop_signal: F);
|
||||
async fn run<F: Future<Output = ()> + Send>(self, stop_signal: F);
|
||||
}
|
||||
|
||||
pub struct AsyncWorker {
|
||||
@@ -94,7 +94,7 @@ impl AsyncWorker {
|
||||
|
||||
let task = runtime.spawn(async move {
|
||||
job
|
||||
.do_work(async move {
|
||||
.run(async move {
|
||||
recv_stop.await;
|
||||
})
|
||||
.await;
|
||||
|
||||
Reference in New Issue
Block a user