Something went wrong. Try again.
A user-space CRDT filesystem
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240//! Thin client for the FSKitBridge wire protocol.//!//! Wraps a `TcpStream` and provides typed helpers for the RPCs that//! `vendor/fskit-rs/src/handler.rs` knows how to answer. Reuses the same//! prost-generated codegen as the daemon (`fskit_rs::pb`) so request ///! response types match exactly.
use std::ffi::OsStr;use std::os::unix::ffi::OsStrExt;
use bytes::{Buf, BytesMut};use fskit_rs::pb::directory_entries::Entry;use fskit_rs::pb::{ Activate, CreateItem, EnumerateDirectory, LookupItem, Read as ReadReq, RemoveItem, Request, Response, TaskOptions, Write as WriteReq, request, response,};use fskit_rs::{Item, ItemAttributes, ItemType};use prost::Message;use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::TcpStream;
/// FSKit reserves item id 2 for the root directory.pub const ROOT_ID: u64 = 2;
/// Errors returned from the wire layer.#[derive(Debug, thiserror::Error)]pub enum Error { #[error("io: {0}")] Io(#[from] std::io::Error),
#[error("protocol: connection closed before response")] ConnectionClosed,
#[error("protocol: response missing content")] EmptyResponse,
#[error("protocol: response id mismatch (expected {expected}, got {got})")] IdMismatch { expected: u64, got: u64 },
#[error("protocol: unexpected response variant: {0}")] UnexpectedResponse(String),
#[error("posix: {0}")] Posix(i32),}
pub type Result<T> = std::result::Result<T, Error>;
/// A connected RPC channel. Each [`Connection`] owns its own request id/// counter.pub struct Connection { stream: TcpStream, buf: BytesMut, next_id: u64,}
impl Connection { /// Connect to a running daemon at `addr` (e.g. `"127.0.0.1:8080"`). pub async fn connect(addr: &str) -> Result<Self> { let stream = TcpStream::connect(addr).await?; Ok(Self { stream, buf: BytesMut::with_capacity(4096), next_id: 1, }) }
/// Wrap an already-connected stream. pub fn from_stream(stream: TcpStream) -> Self { Self { stream, buf: BytesMut::with_capacity(4096), next_id: 1, } }
fn next_id(&mut self) -> u64 { let id = self.next_id; self.next_id += 1; id }
/// Send a request and read the response payload. pub async fn rpc(&mut self, content: request::Content) -> Result<response::Content> { let id = self.next_id(); let req = Request { id, content: Some(content), }; let mut out = Vec::with_capacity(req.encoded_len() + 10); req.encode_length_delimited(&mut out).expect("encode"); self.stream.write_all(&out).await?;
loop { let n = self.stream.read_buf(&mut self.buf).await?; if n == 0 { return Err(Error::ConnectionClosed); } let mut frozen = self.buf.clone().freeze(); let before = frozen.remaining(); match Response::decode_length_delimited(&mut frozen) { Ok(resp) => { let consumed = before - frozen.remaining(); self.buf.advance(consumed); if resp.request_id != id { return Err(Error::IdMismatch { expected: id, got: resp.request_id, }); } let content = resp.content.ok_or(Error::EmptyResponse)?; if let response::Content::PosixError(code) = content { return Err(Error::Posix(code)); } return Ok(content); } Err(_) => continue, } } }
pub async fn activate(&mut self) -> Result<Item> { match self .rpc(request::Content::Activate(Activate { options: Some(TaskOptions::default()), })) .await? { response::Content::Item(item) => Ok(item), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
pub async fn lookup(&mut self, name: &OsStr, directory_id: u64) -> Result<Item> { match self .rpc(request::Content::LookupItem(LookupItem { name: name.as_bytes().to_vec(), directory_id, })) .await? { response::Content::Item(item) => Ok(item), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
pub async fn create( &mut self, name: &OsStr, kind: ItemType, directory_id: u64, ) -> Result<Item> { match self .rpc(request::Content::CreateItem(CreateItem { name: name.as_bytes().to_vec(), r#type: kind as i32, directory_id, attributes: Some(ItemAttributes::default()), })) .await? { response::Content::Item(item) => Ok(item), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
pub async fn write(&mut self, item_id: u64, offset: i64, contents: Vec<u8>) -> Result<i64> { match self .rpc(request::Content::Write(WriteReq { contents, item_id, offset, })) .await? { response::Content::ByteCount(n) => Ok(n), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
pub async fn read(&mut self, item_id: u64, offset: i64, length: i64) -> Result<Vec<u8>> { match self .rpc(request::Content::Read(ReadReq { item_id, offset, length, })) .await? { response::Content::Data(data) => Ok(data), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
pub async fn remove(&mut self, item_id: u64, name: &OsStr, directory_id: u64) -> Result<()> { match self .rpc(request::Content::RemoveItem(RemoveItem { item_id, name: name.as_bytes().to_vec(), directory_id, })) .await? { response::Content::Success(_) => Ok(()), other => Err(Error::UnexpectedResponse(format!("{other:?}"))), } }
/// Enumerate every entry in `directory_id`. Walks pagination cookies /// until the daemon returns an empty page or the cookie stops advancing. pub async fn enumerate(&mut self, directory_id: u64) -> Result<Vec<Entry>> { let mut entries = Vec::new(); let mut cookie = 0u64; let mut verifier = 0u64; loop { let resp = self .rpc(request::Content::EnumerateDirectory(EnumerateDirectory { directory_id, cookie, verifier, })) .await?; let response::Content::DirectoryEntries(page) = resp else { return Err(Error::UnexpectedResponse(format!("{resp:?}"))); }; verifier = page.verifier; if page.entries.is_empty() { break; } let next = page.entries.last().map(|e| e.next_cookie); entries.extend(page.entries); match next { Some(c) if c != cookie => cookie = c, _ => break, } } Ok(entries) }}