1pub mod store;
12pub mod writer;
13
14use std::path::Path;
15
16use anyhow::{Context, Result};
17use axum::{
18 Router,
19 extract::{
20 State, WebSocketUpgrade,
21 ws::{Message, WebSocket},
22 },
23 response::Response,
24 routing::any,
25};
26use blockworx_doc::{encode, protocol::ClientMsg, session::Host};
27use futures_util::{SinkExt, StreamExt};
28use tokio::{net::TcpListener, sync::mpsc};
29
30use crate::{
31 store::{Store, replay_into},
32 writer::{Request, Writer},
33};
34
35type Writers = mpsc::UnboundedSender<Request>;
36
37pub fn start(database: &Path) -> Result<Writers> {
44 let store = Store::open(database)?;
45 let mut host = Host::default();
46 replay_into(&store, &mut host).context("replaying the log")?;
47 tracing::info!("replayed to rev {}", host.rev().get());
48
49 let (requests, receiver) = mpsc::unbounded_channel();
50 tokio::spawn(Writer::new(host, store).run(receiver));
51 Ok(requests)
52}
53
54pub async fn serve(listener: TcpListener, writer: Writers) -> Result<()> {
57 let app = Router::new().route("/ws", any(upgrade)).with_state(writer);
58 axum::serve(listener, app).await.context("serving")
59}
60
61async fn upgrade(State(writer): State<Writers>, ws: WebSocketUpgrade) -> Response {
62 ws.on_upgrade(move |socket| connection(socket, writer))
63}
64
65async fn connection(socket: WebSocket, writer: Writers) {
68 let (mut sink, mut stream) = socket.split();
69 let (outbox, mut inbox) = mpsc::unbounded_channel();
70 let (reply, id) = tokio::sync::oneshot::channel();
71
72 if writer.send(Request::Connect { outbox, reply }).is_err() {
73 return;
74 }
75 let Ok(id) = id.await else { return };
76
77 let pump = tokio::spawn(async move {
78 while let Some(message) = inbox.recv().await {
79 if sink
80 .send(Message::Binary(encode::to_bytes(&message).into()))
81 .await
82 .is_err()
83 {
84 break;
85 }
86 }
87 });
88
89 while let Some(Ok(message)) = stream.next().await {
90 let Message::Binary(bytes) = message else {
91 continue;
92 };
93 match encode::from_bytes::<ClientMsg>(&bytes) {
94 Ok(ClientMsg::Submit { nonce, commit }) => {
95 if writer
96 .send(Request::Submit {
97 from: id,
98 nonce,
99 commit,
100 })
101 .is_err()
102 {
103 break;
104 }
105 }
106 Err(refusal) => {
110 tracing::warn!("closing a connection that sent an unreadable frame: {refusal}");
111 break;
112 }
113 }
114 }
115
116 let _ = writer.send(Request::Disconnect(id));
117 pump.abort();
118}