//! 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 = std::result::Result; /// 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 { 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 { 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 { 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 { 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 { 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) -> Result { 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> { 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> { 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) } }