use super::*;
use crate::peer::{self, Content, Identity, PairState, Peer, Waiting};
use std::collections::HashMap;
#[derive(Default)]
pub(crate) struct Peers {
pub identity: std::sync::Mutex<Option<Identity>>,
pub pairings: std::sync::Mutex<HashMap<String, Pairing>>,
pub wake: Wake,
pub flushing: std::sync::atomic::AtomicBool,
pub again: std::sync::atomic::AtomicBool,
pub held_read: std::sync::atomic::AtomicBool,
pub lines: std::sync::Mutex<HashMap<i64, LineHandle>>,
pub quiet_until: std::sync::atomic::AtomicI64,
}
pub(crate) struct Wake(tokio::sync::watch::Sender<u64>);
impl Default for Wake {
fn default() -> Self {
Wake(tokio::sync::watch::channel(0).0)
}
}
impl Wake {
pub fn notify_one(&self) {
self.0.send_modify(|n| *n += 1);
}
pub fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
self.0.subscribe()
}
}
#[derive(Clone)]
pub(crate) struct LineHandle {
pub tx: tokio::sync::mpsc::Sender<Outbound>,
}
pub(crate) struct Outbound {
pub id: String,
pub chunks: Vec<Vec<u8>>,
}
#[derive(Clone, Debug)]
pub(crate) struct Pairing {
pub state: PairState,
pub until: i64,
}
const PAIRINGS_MAX: usize = 5;
const PAIRING_KEPT: i64 = 600;
const TRIES_MAX: i64 = 20;
const SENT_SHOWN: i64 = 5;
const TOO_LARGE: &str = "too large to send to a friend (over 8 MB)";
pub(crate) fn identity_blocking(app: &App) -> anyhow::Result<Identity> {
if let Some(id) = app.peers.identity.lock().unwrap().clone() {
return Ok(id);
}
let id = Identity::load_or_mint(&app.secrets)?;
*app.peers.identity.lock().unwrap() = Some(id.clone());
Ok(id)
}
pub(crate) fn identity_if_any(app: &App) -> Option<Identity> {
if let Some(id) = app.peers.identity.lock().unwrap().clone() {
return Some(id);
}
let id = Identity::load(&app.secrets)?;
*app.peers.identity.lock().unwrap() = Some(id.clone());
Some(id)
}
pub(crate) fn my_name(paths: &Paths) -> String {
std::fs::read_to_string(paths.config_dir.join("peer-name"))
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.unwrap_or_else(|| {
std::env::var("USER")
.or_else(|_| std::env::var("USERNAME"))
.ok()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| "a friend".to_string())
})
}
fn set_my_name(paths: &Paths, name: &str) {
let name: String = name.trim().chars().take(60).collect();
if name.is_empty() {
return;
}
let _ = std::fs::create_dir_all(&paths.config_dir);
let _ = std::fs::write(paths.config_dir.join("peer-name"), name);
}
pub(crate) fn peers_moved(app: &App) {
emit(app, "peers", json!({}));
}
pub(crate) async fn peers_list(State(app): S) -> Response {
let friends = app.store.peers().unwrap_or_default();
let notes = app.store.peer_notes().unwrap_or_default();
let offers = app.store.peer_offers().unwrap_or_default();
let outbox: Vec<_> = app
.store
.peer_outgoing()
.unwrap_or_default()
.into_iter()
.map(|o| {
let stopped = o.tries >= TRIES_MAX || o.later >= peer::LATER_MAX;
let mut v = json!(o);
v["stopped"] = json!(stopped);
v
})
.collect();
let sent = app.store.peer_sent_recent(SENT_SHOWN).unwrap_or_default();
Json(json!({
"friends": friends,
"notes": notes,
"offers": offers,
"outbox": outbox,
"sent": sent,
"me": { "name": my_name(&app.paths) },
"relay": peer::relay(),
}))
.into_response()
}
#[derive(Deserialize, Default)]
pub(crate) struct PairBody {
#[serde(default)]
pub(crate) name: String,
#[serde(default)]
pub(crate) code: String,
}
fn too_many(app: &App) -> bool {
let now = crate::store::now();
let mut p = app.peers.pairings.lock().unwrap();
p.retain(|_, x| x.until + PAIRING_KEPT > now);
p.values().filter(|x| x.state == PairState::Waiting).count() >= PAIRINGS_MAX
}
pub(crate) async fn pair_start(
State(app): S,
headers: HeaderMap,
Json(b): Json<PairBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
if too_many(&app) {
return (
StatusCode::TOO_MANY_REQUESTS,
Json(json!({ "error": "a few codes are already waiting; let one run out" })),
)
.into_response();
}
let code = match peer::mint_code() {
Ok(c) => c,
Err(e) => return err(e),
};
set_my_name(&app.paths, &b.name);
let until = start_pairing(&app, code.clone(), my_name(&app.paths));
Json(json!({ "code": code, "until": until })).into_response()
}
pub(crate) async fn pair_join(
State(app): S,
headers: HeaderMap,
Json(b): Json<PairBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let code = match peer::normalize_code(&b.code) {
Ok(c) => c,
Err(why) => {
return (
StatusCode::UNPROCESSABLE_ENTITY,
Json(json!({ "error": why })),
)
.into_response()
}
};
if too_many(&app) {
return (
StatusCode::TOO_MANY_REQUESTS,
Json(json!({ "error": "a few codes are already waiting; let one run out" })),
)
.into_response();
}
if app.peers.pairings.lock().unwrap().contains_key(&code) {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "that is this snyvi's own code" })),
)
.into_response();
}
set_my_name(&app.paths, &b.name);
let until = start_pairing(&app, code.clone(), my_name(&app.paths));
Json(json!({ "code": code, "until": until })).into_response()
}
pub(crate) async fn pair_state(State(app): S, Path(code): Path<String>) -> Response {
let p = app.peers.pairings.lock().unwrap().get(&code).cloned();
match p {
Some(p) => Json(json!({ "pairing": p.state, "until": p.until })).into_response(),
None => StatusCode::NOT_FOUND.into_response(),
}
}
fn start_pairing(app: &Arc<App>, code: String, name: String) -> i64 {
let until = crate::store::now() + peer::CODE_TTL;
app.peers.pairings.lock().unwrap().insert(
code.clone(),
Pairing {
state: PairState::Waiting,
until,
},
);
let app = app.clone();
tokio::spawn(async move {
let app2 = app.clone();
let c = code.clone();
let r = tokio::task::spawn_blocking(move || -> anyhow::Result<(Peer, String)> {
let me = identity_blocking(&app2)?;
peer::pair(&me, &c, &name, until)
})
.await;
let state = match r {
Ok(Ok((p, emoji))) => match app.store.pin_peer(&p) {
Ok(pinned) => {
eprintln!("snyvi: paired with {}", pinned.name);
peers_moved(&app);
app.peers.wake.notify_one();
PairState::Done {
emoji,
name: pinned.name,
peer: pinned.id,
}
}
Err(e) => PairState::Failed { why: e.to_string() },
},
Ok(Err(e)) => PairState::Failed {
why: format!("{e:#}"),
},
Err(e) => PairState::Failed { why: e.to_string() },
};
if let Some(p) = app.peers.pairings.lock().unwrap().get_mut(&code) {
p.state = state.clone();
}
emit(&app, "pairing", json!({ "code": code, "pairing": state }));
});
until
}
#[derive(Deserialize, Default)]
pub(crate) struct NameBody {
#[serde(default)]
pub(crate) name: String,
}
pub(crate) async fn peer_rename(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Json(b): Json<NameBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.rename_peer(id, &b.name) {
Ok(true) => {
if let Ok(Some(p)) = app.store.peer(id) {
let _ = app
.store
.rename_project_by_root(&p.project_root(), &p.project_name());
}
peers_moved(&app);
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
#[derive(Deserialize, Default)]
pub(crate) struct MuteBody {
#[serde(default)]
pub(crate) muted: bool,
}
pub(crate) async fn peer_mute(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Json(b): Json<MuteBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.mute_peer(id, b.muted) {
Ok(true) => {
peers_moved(&app);
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
pub(crate) async fn peer_remove(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.remove_peer(id) {
Ok(true) => {
peers_moved(&app);
app.peers.wake.notify_one();
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
pub(crate) async fn peer_restore(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.restore_peer(id) {
Ok(true) => {
peers_moved(&app);
app.peers.wake.notify_one();
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
#[derive(Deserialize, Default)]
pub(crate) struct TextBody {
#[serde(default)]
pub(crate) text: String,
}
pub(crate) async fn peer_note(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Json(b): Json<TextBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let text: String = b.text.trim().chars().take(peer::NOTE_CHARS).collect();
if text.is_empty() {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": "a line of text" })),
)
.into_response();
}
let Ok(Some(p)) = app.store.peer(id) else {
return StatusCode::NOT_FOUND.into_response();
};
if p.removed_at != 0 {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this friend was removed; Restore them first" })),
)
.into_response();
}
let frame = match app.store.peer_queue_note(&p, &text) {
Ok(f) => f,
Err(e) => return err(e),
};
let app2 = app.clone();
let sent = tokio::task::spawn_blocking(move || flush_outbox(&app2, Some(&frame)))
.await
.unwrap_or(false);
peers_moved(&app);
Json(json!({ "sent": sent, "to": p.name })).into_response()
}
#[derive(Deserialize, Default)]
pub(crate) struct SettleBody {
#[serde(default)]
pub(crate) what: String,
#[serde(default)]
pub(crate) desk: i64,
}
pub(crate) async fn peer_note_settle(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Query(q): Query<std::collections::HashMap<String, String>>,
Json(b): Json<SettleBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
if b.what == "keep" {
if let Some(no) = refuse_desk(&app, &headers, &q) {
return no;
}
return keep_line(&app, id, b.desk);
}
match app.store.settle_peer_note(id, &b.what) {
Ok(true) => {
emit(&app, "peernotes", json!({}));
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
fn keep_line(app: &Arc<App>, id: i64, desk: i64) -> Response {
let Ok(Some(n)) = app
.store
.peer_notes()
.map(|ns| ns.into_iter().find(|n| n.id == id))
else {
return StatusCode::NOT_FOUND.into_response();
};
let Ok(Some(d)) = app.store.desk(desk) else {
return (
StatusCode::NOT_FOUND,
Json(json!({ "error": "no such desk" })),
)
.into_response();
};
match app.store.add_desk_note_from(desk, &n.text, &n.from) {
Ok(Some(note)) => {
let _ = app.store.settle_peer_note(id, "taken");
if !n.frame.is_empty() {
let _ = app.store.link_note_frame(note.id, n.peer_id, &n.frame);
}
emit(app, "desknotes", json!({ "desk": desk }));
emit(app, "peernotes", json!({}));
Json(json!({ "note": note, "desk": desk, "name": d.name })).into_response()
}
Ok(None) => (
StatusCode::CONFLICT,
Json(json!({ "error": format!("{} keeps {} notes; take one off first", d.name, crate::desk::NOTES_PER_DESK) })),
)
.into_response(),
Err(e) => err(e),
}
}
#[derive(Deserialize, Default)]
pub(crate) struct SendBody {
#[serde(default)]
pub(crate) peer: i64,
}
pub(crate) async fn doc_send(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
Json(b): Json<SendBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let Ok(Some(doc)) = app.store.get(&id) else {
return StatusCode::NOT_FOUND.into_response();
};
let Ok(Some(p)) = app.store.peer(b.peer) else {
return (
StatusCode::NOT_FOUND,
Json(json!({ "error": "no such friend" })),
)
.into_response();
};
if p.removed_at != 0 {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this friend was removed; Restore them first" })),
)
.into_response();
}
if doc.size as usize > peer::SEND_MAX {
return (
StatusCode::PAYLOAD_TOO_LARGE,
Json(json!({ "error": TOO_LARGE })),
)
.into_response();
}
let frame = match app.store.peer_queue(&p, &id) {
Ok(f) => f,
Err(e) => return err(e),
};
let app2 = app.clone();
let sent = tokio::task::spawn_blocking(move || flush_outbox(&app2, Some(&frame)))
.await
.unwrap_or(false);
peers_moved(&app);
Json(json!({ "sent": sent, "to": p.name })).into_response()
}
#[derive(Deserialize, Default)]
pub(crate) struct AnswerBody {
#[serde(default)]
pub(crate) send: bool,
#[serde(default)]
pub(crate) undo: bool,
}
pub(crate) async fn offer_answer(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Json(b): Json<AnswerBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
if b.undo {
return match app.store.reopen_peer_offer(id) {
Ok(true) => {
emit(&app, "peeroffers", json!({}));
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
};
}
let Ok(Some(o)) = app.store.peer_offer_get(id) else {
return StatusCode::NOT_FOUND.into_response();
};
if b.send {
if let Ok(Some(doc)) = app.store.get(&o.doc_id) {
if doc.size as usize > peer::SEND_MAX {
return (
StatusCode::PAYLOAD_TOO_LARGE,
Json(json!({ "error": TOO_LARGE })),
)
.into_response();
}
}
}
if let Err(e) = app.store.answer_peer_offer(id, b.send) {
return err(e);
}
emit(&app, "peeroffers", json!({}));
if !b.send {
return Json(json!({ "sent": false })).into_response();
}
let Ok(Some(p)) = app.store.peer(o.peer_id) else {
return StatusCode::NOT_FOUND.into_response();
};
let frame = match if o.doc_id.is_empty() {
app.store.peer_queue_note(&p, &o.text)
} else {
app.store.peer_queue(&p, &o.doc_id)
} {
Ok(f) => f,
Err(e) => return err(e),
};
let app2 = app.clone();
let sent = tokio::task::spawn_blocking(move || flush_outbox(&app2, Some(&frame)))
.await
.unwrap_or(false);
peers_moved(&app);
Json(json!({ "sent": sent, "to": p.name })).into_response()
}
#[derive(Deserialize, Default)]
pub(crate) struct OfferBody {
#[serde(default)]
pub(crate) to: String,
#[serde(default)]
pub(crate) doc: String,
#[serde(default)]
pub(crate) text: String,
#[serde(default)]
pub(crate) by: String,
}
pub(crate) async fn pane_offer(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
Json(b): Json<OfferBody>,
) -> Response {
let placed = match agent_pane(&app, &headers, &id) {
Ok(p) => p,
Err(no) => return *no,
};
let to = b.to.trim();
if to.is_empty() {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": "offer_document needs the friend's name" })),
)
.into_response();
}
let Ok(Some(p)) = app.store.peer_by_name(to) else {
return (StatusCode::NOT_FOUND, Json(json!({ "error": format!("the user has no friend called {to}; they pair from Home") }))).into_response();
};
let line: String = b.text.trim().chars().take(peer::NOTE_CHARS).collect();
if b.doc.trim().is_empty() && !line.is_empty() {
return match app.store.peer_offer(p.id, "", &line, &id, &b.by) {
Ok(offer_id) => {
emit(
&app,
"peeroffers",
json!({ "offer": offer_id, "to": p.name, "text": line, "by": b.by, "desk": placed.desk_name }),
);
(
StatusCode::CREATED,
Json(json!({ "to": p.name, "text": line })),
)
.into_response()
}
Err(e) => err(e),
};
}
let Ok(Some(doc)) = app.store.get(b.doc.trim()) else {
return (
StatusCode::NOT_FOUND,
Json(json!({ "error": "no document by that id; send_document answers with one" })),
)
.into_response();
};
match app.store.peer_offer(p.id, &doc.id, "", &id, &b.by) {
Ok(offer_id) => {
emit(
&app,
"peeroffers",
json!({ "offer": offer_id, "to": p.name, "title": doc.title, "by": b.by, "desk": placed.desk_name }),
);
(
StatusCode::CREATED,
Json(json!({ "to": p.name, "title": doc.title })),
)
.into_response()
}
Err(e) => err(e),
}
}
pub(crate) fn flush_all(app: &Arc<App>) {
use std::sync::atomic::Ordering::SeqCst;
app.peers.again.store(true, SeqCst);
if app.peers.flushing.swap(true, SeqCst) {
return;
}
let app = app.clone();
tokio::task::spawn_blocking(move || {
while app.peers.again.swap(false, SeqCst) {
flush_outbox(&app, None);
}
app.peers.flushing.store(false, SeqCst);
if app.peers.again.load(SeqCst) && !app.peers.flushing.swap(true, SeqCst) {
while app.peers.again.swap(false, SeqCst) {
flush_outbox(&app, None);
}
app.peers.flushing.store(false, SeqCst);
}
});
}
pub(crate) fn flush_outbox(app: &App, only: Option<&str>) -> bool {
let now = crate::store::now();
let rows = match only {
Some(id) => app
.store
.peer_unsent_one(id)
.ok()
.flatten()
.into_iter()
.collect(),
None => app.store.peer_unsent(now).unwrap_or_default(),
};
let Ok(me) = identity_blocking(app) else {
return false;
};
let name = my_name(&app.paths);
let mut all = true;
let mut moved = false;
for row in rows {
let (frame_id, peer_id) = (row.id.as_str(), row.peer_id);
if stopped(&row) {
all = false;
continue;
}
let went = match send_one(app, &me, &name, &row) {
Ok(peer::Deposit::Sent) => {
let _ = app.store.peer_sent(frame_id);
if row.kind != "receipt" {
let _ = app.store.touch_peer(peer_id, false);
}
moved = true;
true
}
Ok(peer::Deposit::Pending) => {
let _ = app.store.peer_pending(frame_id);
only.is_some()
}
Ok(peer::Deposit::Later { why, until }) => {
let until = until.or((why == peer::BUSY).then(|| now + 60));
let _ = app.store.peer_waiting(frame_id, why, until);
false
}
Err(e) => {
let _ = app.store.peer_failed(frame_id, &format!("{e:#}"));
false
}
};
all &= went;
}
if moved && only.is_none() {
peers_moved(app);
}
all
}
fn stopped(row: &peer::Unsent) -> bool {
row.tries >= TRIES_MAX || row.later >= peer::LATER_MAX
}
fn send_one(
app: &App,
me: &Identity,
name: &str,
row: &peer::Unsent,
) -> anyhow::Result<peer::Deposit> {
let (frame_id, doc_id) = (row.id.as_str(), row.doc_id.as_str());
let p = app
.store
.peer(row.peer_id)?
.ok_or_else(|| anyhow::anyhow!("no such friend"))?;
if !row.kind.is_empty() {
let unread = row.kind == "receipt" && row.text == "read" && !p.read_receipts;
if row.kind == "receipt" && (p.v < peer::REPLIES_V || unread) {
return Ok(peer::Deposit::Sent);
}
if p.v < peer::REPLIES_V {
return Ok(peer::Deposit::Later {
why: OLDER,
until: None,
});
}
let quiet = app
.peers
.quiet_until
.load(std::sync::atomic::Ordering::Relaxed);
if row.kind == "receipt" && row.text == "read" && quiet > crate::store::now() {
return Ok(peer::Deposit::Later {
why: peer::NEARLY_OUT,
until: Some(quiet),
});
}
let content = match row.kind.as_str() {
"receipt" => Content::Receipt {
of: row.re.clone(),
state: row.text.clone(),
v: peer::CONTENT_V,
},
"reply" => Content::Reply {
re: row.re.clone(),
text: row.text.clone(),
name: name.to_string(),
v: peer::CONTENT_V,
},
"done" => Content::Done {
of: row.re.clone(),
text: row.text.clone(),
commit: Some(row.extra.clone()).filter(|c| !c.is_empty()),
name: name.to_string(),
v: peer::CONTENT_V,
},
other => anyhow::bail!("a frame of a kind this snyvi does not send: {other}"),
};
let frame = peer::seal(me, &p, &content, b"")?;
return leave(app, me, &p, frame_id, &frame);
}
if doc_id.is_empty() {
let content = Content::Note {
text: row.text.clone(),
name: name.to_string(),
at: peer::Folder {
v: peer::CONTENT_V,
..Default::default()
},
};
let frame = peer::seal(me, &p, &content, b"")?;
return leave(app, me, &p, frame_id, &frame);
}
let Some(doc) = app.store.get(doc_id)? else {
app.store.peer_sent(frame_id)?;
return Ok(peer::Deposit::Sent);
};
let mut bytes = std::fs::read(app.store.src_path(doc_id))?;
if bytes.len() > peer::SEND_MAX {
anyhow::bail!("{TOO_LARGE}");
}
let room = peer::SEND_MAX.saturating_sub(bytes.len());
let with = match (
doc.kind,
doc.source_path.as_deref(),
std::str::from_utf8(&bytes),
) {
(crate::render::Kind::Markdown, Some(file), Ok(text)) => {
Some(inline_pictures(text, std::path::Path::new(file), room))
}
_ => None,
};
if let Some(with) = with {
bytes = with.into_bytes();
}
let content = Content::Document {
title: doc.title.clone(),
lang: doc.lang.clone(),
file: doc
.source_path
.as_deref()
.and_then(|sp| std::path::Path::new(sp).file_name())
.map(|f| f.to_string_lossy().to_string()),
name: name.to_string(),
id: doc.id.clone(),
at: folder_of(app, &doc),
};
let frame = peer::seal(me, &p, &content, &bytes)?;
leave(app, me, &p, frame_id, &frame)
}
fn leave(
app: &App,
me: &Identity,
p: &Peer,
frame_id: &str,
frame: &[u8],
) -> anyhow::Result<peer::Deposit> {
if p.v < peer::LINE_V {
return peer::deposit(me, &p.sign_key, frame_id, frame);
}
let line = app.peers.lines.lock().unwrap().get(&p.id).cloned();
let Some(line) = line else {
return Ok(peer::Deposit::Later {
why: LINE_DOWN,
until: None,
});
};
let out = Outbound {
id: frame_id.to_string(),
chunks: peer::chunks(frame_id, frame),
};
match line.tx.try_send(out) {
Ok(()) => Ok(peer::Deposit::Pending),
Err(_) => Ok(peer::Deposit::Later {
why: LINE_DOWN,
until: None,
}),
}
}
const LINE_DOWN: &str = "the line to them is not open";
const OLDER: &str = "their snyvi is older; this goes once they update";
pub(crate) const PRINT_KEPT: i64 = 86_400;
pub(crate) fn folder_of(app: &App, doc: &crate::store::Doc) -> peer::Folder {
let mut at = peer::Folder {
v: peer::CONTENT_V,
..Default::default()
};
let Some(root) = app
.store
.project_root(doc.project_id)
.filter(|r| !r.starts_with("peer:"))
else {
return at;
};
let print = match app.store.print_of(&root) {
Ok(Some((p, when))) if crate::store::now() - when < PRINT_KEPT => Some(p),
_ => {
let p = crate::git::print(std::path::Path::new(&root));
let _ = app.store.set_print(&root, p.as_ref());
p
}
};
let Some(print) = print.filter(|p| p.repo.is_some() || p.remote.is_some()) else {
return at;
};
at.repo = print.repo;
at.remote = print.remote;
at.branch = doc
.branch
.clone()
.or_else(|| crate::project::head_of(std::path::Path::new(&root)));
(at.path, at.key) = place_in(&root, doc.source_path.as_deref());
at
}
pub(crate) fn place_in(root: &str, source: Option<&str>) -> (Option<String>, Option<String>) {
let Some(sp) = source.filter(|s| !s.trim().is_empty()) else {
return (None, None);
};
let file = std::path::Path::new(sp);
if !file.is_absolute() {
return (peer::safe_path(sp), None);
}
match file.strip_prefix(root) {
Ok(rel) => (
peer::safe_path(&rel.to_string_lossy().replace('\\', "/")),
None,
),
Err(_) => {
let h = blake3::hash(format!("snyvi key v1 {sp}").as_bytes());
(None, Some(h.to_hex()[..16].to_string()))
}
}
}
pub(crate) fn inline_pictures(text: &str, file: &std::path::Path, room: usize) -> String {
let Some(dir) = file.parent().filter(|d| d.is_absolute()) else {
return text.to_string();
};
let Ok(root) = crate::project::resolve(dir).root.canonicalize() else {
return text.to_string();
};
let mut left = room;
crate::render::map_md_images(text, |alt, url| {
if !crate::render::relative_url(url) {
return None;
}
let mime = crate::render::picture_mime(&crate::render::ext_of(url))?;
let at = dir.join(url).canonicalize().ok()?;
if !at.starts_with(&root) || std::fs::metadata(&at).ok()?.len() as usize > left {
return None;
}
let bytes = std::fs::read(&at).ok()?;
let md = format!("", crate::render::data_uri(mime, &bytes));
left = left.checked_sub(md.len())?;
Some(md)
})
}
pub(crate) fn take_in(
app: &Arc<App>,
me: &Identity,
w: &Waiting,
bytes: Vec<u8>,
via_line: bool,
) -> anyhow::Result<bool> {
let friend = app
.store
.peer_by_key(&w.sender)?
.filter(|p| p.removed_at == 0);
let Some(p) = friend else {
return Ok(false);
};
if app.store.peer_taken(&w.id)? {
return Ok(false);
}
if w.size as usize > peer::FRAME_MAX || bytes.len() > peer::FRAME_MAX {
return Ok(false);
}
match peer::open(me, &p, &bytes) {
Ok((Content::Other, _)) => {
app.store.peer_hold(&w.id, p.id, &bytes)?;
app.store.peer_take(&w.id)?;
eprintln!("snyvi: {} sent something this snyvi is too old to read; it is kept for after an update", p.name);
Ok(false)
}
Ok((content, body)) => {
if let Err(e) = arrived(app, &p, content, body, &w.id, via_line) {
eprintln!("snyvi: a document from {} could not be kept: {e:#}", p.name);
return Ok(false);
}
app.store.peer_take(&w.id)?;
Ok(true)
}
Err(e) => {
eprintln!("snyvi: a frame from {} was dropped: {e:#}", p.name);
Ok(false)
}
}
}
pub(crate) fn read_held(app: &Arc<App>, me: &Identity) {
let _ = app.store.prune_peer_held();
let Ok(held) = app.store.peer_held() else {
return;
};
for (id, peer_id, bytes) in held {
let friend = app
.store
.peer(peer_id)
.ok()
.flatten()
.filter(|p| p.removed_at == 0);
let Some(p) = friend else {
let _ = app.store.peer_unhold(&id);
continue;
};
match peer::open(me, &p, &bytes) {
Ok((Content::Other, _)) => {}
Ok((content, body)) => {
if let Err(e) = arrived(app, &p, content, body, &id, false) {
eprintln!(
"snyvi: something held from {} could not be kept: {e:#}",
p.name
);
}
let _ = app.store.peer_unhold(&id);
}
Err(_) => {
let _ = app.store.peer_unhold(&id);
}
}
}
}
pub(crate) fn bring_in(app: &Arc<App>, me: &Identity) -> anyhow::Result<usize> {
let waiting = peer::inbox(me)?;
let mut n = 0;
for w in waiting {
let friend = app
.store
.peer_by_key(&w.sender)?
.filter(|p| p.removed_at == 0);
let fetch =
friend.is_some() && !app.store.peer_taken(&w.id)? && w.size as usize <= peer::FRAME_MAX;
if fetch {
let bytes = peer::fetch(me, &w.id)?;
if take_in(app, me, &w, bytes, false)? {
n += 1;
}
}
peer::ack(me, &w.id)?;
}
Ok(n)
}
pub(crate) fn friend_desk(app: &App, p: &Peer) -> Option<crate::desk::Desk> {
if p.desk_id == 0 {
return None;
}
app.store
.desk(p.desk_id)
.ok()
.flatten()
.filter(|d| d.parked.is_none())
}
fn arrived(
app: &Arc<App>,
p: &Peer,
content: Content,
body: Vec<u8>,
frame: &str,
via_line: bool,
) -> anyhow::Result<()> {
if matches!(content, Content::Document { .. } | Content::Note { .. }) {
let _ = app.store.touch_peer(p.id, true);
}
if app.store.peer_due_now(p.id).unwrap_or(0) > 0 {
app.peers.wake.notify_one();
}
let (v, was) = (content.v(), p.v);
let p = &Peer { v, ..p.clone() };
if v != was {
let _ = app.store.peer_set_v(p.id, v);
let _ = app.store.peer_due_now(p.id);
if v > was {
app.peers.wake.notify_one();
}
}
let desk = friend_desk(app, p);
let answer = matches!(content, Content::Document { .. } | Content::Note { .. });
let said = |app: &Arc<App>| {
if answer && v >= peer::REPLIES_V && !via_line {
let _ = app
.store
.peer_queue_kind(p, "receipt", frame, "arrived", "");
app.peers.wake.notify_one();
}
};
match content {
Content::Receipt { of, state, .. } => {
if app.store.peer_receipt(p.id, &of, &state)? {
peers_moved(app);
}
}
Content::Reply { re, text, .. } => {
if app.store.peer_sent_doc(p.id, &re)? {
if let Some(text) = app.store.peer_add_reply(p.id, &re, &text)? {
emit(
app,
"peerreply",
json!({ "doc": re, "from": p.name, "reply": text, "quiet": p.muted }),
);
}
}
}
Content::Done {
of, text, commit, ..
} => {
if let Some(line) = app
.store
.peer_done(p.id, &of, commit.as_deref().unwrap_or(""))?
{
peers_moved(app);
emit(
app,
"peerdone",
json!({ "from": p.name, "done": if line.is_empty() { text.chars().take(peer::NOTE_CHARS).collect() } else { line }, "commit": commit, "quiet": p.muted }),
);
}
}
Content::Document {
title,
lang,
file,
at,
id: sender_id,
..
} => {
let filed = file_into(app, &at);
let desk = match &filed {
Some((_, d)) => d.clone(),
None => desk,
};
let payload = Payload {
title: Some(title),
lang,
origin: Some("peer".into()),
sender: Some(p.name.clone()),
peer: Some(crate::receive::FromPeer {
name: p.name.clone(),
sign_key: p.sign_key.clone(),
bytes: body,
lineage: lineage(&at, file.as_deref()),
file,
desk: desk.as_ref().map(|d| (on_desk(d), d.root.clone())),
root: filed.as_ref().map(|(r, _)| r.clone()),
}),
..Default::default()
};
let received = receive::receive(&app.store, &app.renderer, payload)?;
let _ = app.store.set_peer_key(&received.doc.id, &p.sign_key);
let _ = app
.store
.set_peer_frame(&received.doc.id, frame, &sender_id);
said(app);
if p.muted {
let _ = app.store.mark_read(&received.doc.id);
}
let mut ev = doc_event(app, &received);
ev["from"] = json!(p.name);
ev["quiet"] = json!(p.muted);
emit_doc(app, ev);
if desk.is_some() {
emit(app, "deskdocs", json!({}));
}
if !p.muted && !received.existing {
eprintln!("snyvi: {} sent \"{}\"", p.name, received.doc.title);
}
}
Content::Other => {}
Content::Note { text, .. } => {
if let Some(d) = &desk {
if let Ok(crate::desk::Suggested::Note(n)) =
app.store.suggest_desk_note_from(d.id, &text, &p.name)
{
let _ = app.store.link_note_frame(n.id, p.id, frame);
said(app);
emit(app, "desknotes", json!({ "desk": d.id }));
emit(
app,
"peernotes",
json!({ "from": p.name, "quiet": p.muted, "desk": d.name }),
);
return Ok(());
}
}
app.store.peer_note_arrived(p.id, &text, frame)?;
said(app);
emit(
app,
"peernotes",
json!({ "from": p.name, "quiet": p.muted }),
);
}
}
Ok(())
}
fn file_into(app: &App, at: &peer::Folder) -> Option<(String, Option<crate::desk::Desk>)> {
if !at.named() {
return None;
}
let roots = app.store.roots_by_print(at).ok()?;
if roots.is_empty() {
return None;
}
let desks: Vec<_> = app
.store
.desks()
.unwrap_or_default()
.into_iter()
.filter(|d| d.parked.is_none())
.collect();
roots
.into_iter()
.filter(|(root, _)| std::path::Path::new(root).is_dir())
.map(|(root, newest)| {
let desk = desks
.iter()
.filter(|d| receive::desk_project(&d.root).0 == root)
.max_by_key(|d| d.visited_at)
.cloned();
let branch = at.branch.is_some()
&& crate::project::head_of(std::path::Path::new(&root)) == at.branch;
((desk.is_some(), branch, newest), root, desk)
})
.max_by_key(|x| x.0)
.map(|(_, root, desk)| (root, desk))
}
fn lineage(at: &peer::Folder, file: Option<&str>) -> Option<String> {
let name = file
.and_then(|f| std::path::Path::new(f).file_name())
.map(|f| f.to_string_lossy().to_string())
.filter(|f| !f.trim().is_empty());
if at.named() {
if let Some(path) = at.path.as_deref().and_then(peer::safe_path) {
return Some(path);
}
}
let key = at
.key
.as_deref()
.filter(|k| k.len() <= 64 && k.chars().all(|c| c.is_ascii_alphanumeric()));
match (key, name) {
(Some(k), Some(n)) => Some(format!("{k}/{n}")),
(Some(k), None) => Some(k.to_string()),
(None, n) => n,
}
}
fn on_desk(d: &crate::desk::Desk) -> crate::desk::Origin {
crate::desk::Origin {
id: d.id,
name: d.name.clone(),
slot: 0,
}
}
#[derive(Deserialize, Default)]
pub(crate) struct DeskBody {
#[serde(default)]
pub(crate) desk: i64,
}
pub(crate) async fn peer_desk(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Query(q): Query<std::collections::HashMap<String, String>>,
Json(b): Json<DeskBody>,
) -> Response {
if let Some(no) = refuse_desk(&app, &headers, &q) {
return no;
}
if b.desk != 0 && !matches!(app.store.desk(b.desk), Ok(Some(_))) {
return (
StatusCode::NOT_FOUND,
Json(json!({ "error": "no such desk" })),
)
.into_response();
}
match app.store.set_peer_desk(id, b.desk) {
Ok(true) => {
peers_moved(&app);
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
fn theirs(app: &App, id: &str) -> Result<crate::store::Doc, Box<Response>> {
match app.store.get(id) {
Ok(Some(doc)) if doc.origin == "peer" => Ok(doc),
Ok(Some(_)) => Err(Box::new(
(
StatusCode::CONFLICT,
Json(json!({ "error": "only a friend's document is kept this way" })),
)
.into_response(),
)),
Ok(None) => Err(Box::new(StatusCode::NOT_FOUND.into_response())),
Err(e) => Err(Box::new(err(e))),
}
}
pub(crate) async fn doc_keep(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
Query(q): Query<std::collections::HashMap<String, String>>,
Json(b): Json<DeskBody>,
) -> Response {
if let Some(no) = refuse_desk(&app, &headers, &q) {
return no;
}
if let Err(no) = theirs(&app, &id) {
return *no;
}
let Ok(Some(d)) = app.store.desk(b.desk) else {
return (
StatusCode::NOT_FOUND,
Json(json!({ "error": "no such desk" })),
)
.into_response();
};
let (root, name, _) = receive::desk_project(&d.root);
match app.store.move_lineage(&id, &root, &name, &on_desk(&d)) {
Ok(Some(doc)) => {
emit(&app, "pinned", json!({ "id": id }));
emit(&app, "rendered", json!({ "id": id }));
Json(json!({ "doc": doc, "desk": d.name })).into_response()
}
Ok(None) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
pub(crate) async fn doc_save(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
Query(q): Query<std::collections::HashMap<String, String>>,
) -> Response {
if let Some(no) = refuse_desk(&app, &headers, &q) {
return no;
}
let doc = match theirs(&app, &id) {
Ok(d) => d,
Err(no) => return *no,
};
let Some(d) = doc
.desk
.as_ref()
.and_then(|o| app.store.desk(o.id).ok().flatten())
else {
return (
StatusCode::CONFLICT,
Json(
json!({ "error": "keep it on a desk first; it is saved into that desk's folder" }),
),
)
.into_response();
};
let src = app.store.src_path(&id);
let dir = std::path::Path::new(&d.root).join(format!("from-{}", slug(&doc.sender)));
let name = file_name(&doc);
match tokio::task::spawn_blocking(move || save_new(&src, &dir, &name)).await {
Ok(Ok(at)) => {
let rel = at
.strip_prefix(&d.root)
.map(|r| r.to_string_lossy().replace('\\', "/"))
.unwrap_or_else(|_| at.to_string_lossy().to_string());
if let Err(e) = app.store.set_saved_path(&id, &at.to_string_lossy()) {
eprintln!(
"snyvi: {id} was saved to {} but that was not recorded: {e:#}",
at.display()
);
}
emit(&app, "deskdocs", json!({}));
Json(json!({ "path": at, "rel": rel, "desk": d.name })).into_response()
}
Ok(Err(e)) => err(e),
Err(e) => err(anyhow::anyhow!(e)),
}
}
pub(crate) async fn doc_unfile(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
if let Err(no) = theirs(&app, &id) {
return *no;
}
match app.store.unfile(&id) {
Ok(Some(doc)) => {
emit(&app, "pinned", json!({ "id": id }));
emit(&app, "rendered", json!({ "id": id }));
emit(&app, "deskdocs", json!({}));
Json(json!({ "doc": doc })).into_response()
}
Ok(None) => (
StatusCode::CONFLICT,
Json(json!({ "error": "it does not say which friend sent it" })),
)
.into_response(),
Err(e) => err(e),
}
}
pub(crate) async fn outbox_retry(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.peer_retry(&id) {
Ok(true) => {
app.peers.wake.notify_one();
peers_moved(&app);
StatusCode::NO_CONTENT.into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
pub(crate) async fn doc_reply(
State(app): S,
headers: HeaderMap,
Path(id): Path<String>,
Json(b): Json<TextBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let text: String = b.text.trim().chars().take(peer::REPLY_CHARS).collect();
if text.is_empty() {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": "a line of text" })),
)
.into_response();
}
let (key, _, re) = match app.store.peer_frame(&id) {
Ok(f) => f,
Err(e) => return err(e),
};
let friend = app
.store
.peer_by_key(&key)
.ok()
.flatten()
.filter(|p| p.removed_at == 0);
let (Some(p), false) = (friend, re.is_empty()) else {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this one came before replies could be sent; a line from Home reaches them" })),
)
.into_response();
};
let frame = match app.store.peer_queue_kind(&p, "reply", &re, &text, "") {
Ok(f) => f,
Err(e) => return err(e),
};
let app2 = app.clone();
let sent = tokio::task::spawn_blocking(move || flush_outbox(&app2, Some(&frame)))
.await
.unwrap_or(false);
peers_moved(&app);
Json(json!({ "sent": sent, "to": p.name, "waits": !sent && p.v < peer::REPLIES_V }))
.into_response()
}
#[derive(Deserialize, Default)]
pub(crate) struct OnBody {
#[serde(default)]
pub(crate) on: bool,
}
pub(crate) async fn peer_receipts(
State(app): S,
headers: HeaderMap,
Path(id): Path<i64>,
Json(b): Json<OnBody>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
match app.store.peer_read_receipts(id, b.on) {
Ok(true) => {
peers_moved(&app);
Json(json!({ "on": b.on })).into_response()
}
Ok(false) => StatusCode::NOT_FOUND.into_response(),
Err(e) => err(e),
}
}
pub(crate) fn opened(app: &App, id: &str) {
let Ok((key, frame, _)) = app.store.peer_frame(id) else {
return;
};
if frame.is_empty() {
return;
}
let Some(p) = app
.store
.peer_by_key(&key)
.ok()
.flatten()
.filter(|p| p.removed_at == 0 && p.read_receipts && p.v >= peer::REPLIES_V)
else {
return;
};
if app
.store
.peer_queue_kind(&p, "receipt", &frame, "read", "")
.is_ok()
{
app.peers.wake.notify_one();
}
}
pub(crate) async fn note_tell(
State(app): S,
headers: HeaderMap,
Path((desk, note)): Path<(i64, i64)>,
) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let Ok(Some((peer_id, frame, text, commit))) = app.store.note_to_tell(desk, note) else {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "only a ticked line a friend sent, once" })),
)
.into_response();
};
let Ok(Some(p)) = app.store.peer(peer_id) else {
return StatusCode::NOT_FOUND.into_response();
};
if p.removed_at != 0 {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this friend was removed; Restore them first" })),
)
.into_response();
}
let queued = app
.store
.peer_queue_kind(&p, "done", &frame, &text, &commit)
.and_then(|f| app.store.note_told(desk, note).map(|_| f));
let frame = match queued {
Ok(f) => f,
Err(e) => return err(e),
};
emit(&app, "desknotes", json!({ "desk": desk }));
let app2 = app.clone();
let sent = tokio::task::spawn_blocking(move || flush_outbox(&app2, Some(&frame)))
.await
.unwrap_or(false);
peers_moved(&app);
Json(json!({ "sent": sent, "to": p.name, "waits": !sent && p.v < peer::REPLIES_V }))
.into_response()
}
pub(crate) fn slug(name: &str) -> String {
let mut out = String::new();
for c in name.trim().chars().flat_map(char::to_lowercase) {
if c.is_alphanumeric() {
out.push(c);
} else if !out.ends_with('-') {
out.push('-');
}
}
let out = out.trim_matches('-').chars().take(40).collect::<String>();
if out.is_empty() {
"friend".to_string()
} else {
out
}
}
pub(crate) fn file_name(doc: &crate::store::Doc) -> String {
let sent = doc
.source_path
.as_deref()
.and_then(|p| std::path::Path::new(p).file_name())
.map(|f| f.to_string_lossy().to_string())
.filter(|f| !f.starts_with('.') && !f.trim().is_empty());
sent.unwrap_or_else(|| {
let ext = match doc.kind {
crate::render::Kind::Markdown => "md",
_ => doc
.lang
.as_deref()
.filter(|l| l.len() <= 8)
.unwrap_or("txt"),
};
format!("{}.{ext}", slug(&doc.title))
})
}
pub(crate) fn save_new(
src: &std::path::Path,
dir: &std::path::Path,
name: &str,
) -> anyhow::Result<std::path::PathBuf> {
use std::io::Write;
let bytes = std::fs::read(src)?;
std::fs::create_dir_all(dir)?;
let p = std::path::Path::new(name);
let stem = p
.file_stem()
.map(|s| s.to_string_lossy().to_string())
.unwrap_or_else(|| "document".into());
let ext = p
.extension()
.map(|e| format!(".{}", e.to_string_lossy()))
.unwrap_or_default();
for n in 1..1000 {
let at = dir.join(if n == 1 {
format!("{stem}{ext}")
} else {
format!("{stem}-{n}{ext}")
});
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&at)
{
Ok(mut f) => {
f.write_all(&bytes)?;
return Ok(at);
}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(e) => return Err(e.into()),
}
}
anyhow::bail!("a thousand files by that name are already there")
}
pub(crate) fn sweep_offers(app: &App) {
let Ok(offers) = app.store.peer_offers() else {
return;
};
let mut dropped = 0;
for o in offers {
if !o.pane.is_empty() && !app.panes.is_running(&o.pane) {
dropped += app.store.drop_peer_offers_of(&o.pane).unwrap_or(0);
}
}
if dropped > 0 {
emit(app, "peeroffers", json!({}));
}
}