Add 9P2000 server implementation

- p9.rs: Full 9P2000 protocol handler
  - Tversion/Tattach/Twalk/Topen/Tread/Twrite/Tclunk/Tstat
  - Directory listing with stat encoding
  - File read/write with buffering
  - Rdwr support for command files

- ork-server binary:
  - TCP or Unix socket listener
  - Threaded connection handling
  - Usage: ork-server -a unix:///tmp/ork.sock ~/org

Tested with plan9port 9p client - read/write/ls all working.
This commit is contained in:
Levi Neely 2026-09-30 15:09:59 +02:00
parent 96f89acd65
commit 9e2fcb9010
5 changed files with 810 additions and 0 deletions

1
Cargo.lock generated
View File

@ -563,6 +563,7 @@ name = "ork-server"
version = "0.1.0"
dependencies = [
"anyhow",
"bitflags 2.13.2",
"nine",
"notify",
"org-ast",

View File

@ -5,12 +5,17 @@ edition = "2021"
license = "GPL-3.0-or-later"
description = "9P server for orkmode"
[[bin]]
name = "ork-server"
path = "src/bin/ork-server.rs"
[dependencies]
org-ast = { path = "../org-ast" }
org-parser = { path = "../org-parser" }
# 9P protocol
nine = "0.2"
bitflags = "2"
# Async runtime
tokio = { version = "1", features = ["full"] }

View File

@ -0,0 +1,124 @@
//! ork-server - 9P server for org-mode files.
use ork_server::{namespace, OrkState, Server};
use std::env;
use std::io;
use std::net::TcpListener;
use std::os::unix::net::UnixListener;
use std::path::PathBuf;
use std::sync::Arc;
use std::thread;
fn main() -> io::Result<()> {
let args: Vec<String> = env::args().collect();
// Parse arguments
let mut org_dir = env::current_dir()?;
let mut addr = "tcp://localhost:5640".to_string();
let mut i = 1;
while i < args.len() {
match args[i].as_str() {
"-d" | "--dir" => {
i += 1;
if i < args.len() {
org_dir = PathBuf::from(&args[i]);
}
}
"-a" | "--addr" => {
i += 1;
if i < args.len() {
addr = args[i].clone();
}
}
"-h" | "--help" => {
print_help();
return Ok(());
}
_ => {
// Assume it's the org directory
org_dir = PathBuf::from(&args[i]);
}
}
i += 1;
}
// Initialize state
eprintln!("Loading org files from: {}", org_dir.display());
let state = Arc::new(OrkState::new(&org_dir)?);
let root = namespace::build_namespace(state.clone());
let server = Arc::new(Server::new(root));
// Parse address and start server
if addr.starts_with("unix://") {
let path = &addr[7..];
// Remove existing socket
let _ = std::fs::remove_file(path);
let listener = UnixListener::bind(path)?;
eprintln!("Listening on unix://{}", path);
for stream in listener.incoming() {
match stream {
Ok(stream) => {
let server = server.clone();
thread::spawn(move || {
let reader = stream.try_clone().unwrap();
let writer = stream;
if let Err(e) = server.serve(reader, writer) {
eprintln!("Connection error: {}", e);
}
});
}
Err(e) => eprintln!("Accept error: {}", e),
}
}
} else {
// TCP
let tcp_addr = if addr.starts_with("tcp://") {
&addr[6..]
} else {
&addr
};
let listener = TcpListener::bind(tcp_addr)?;
eprintln!("Listening on tcp://{}", tcp_addr);
for stream in listener.incoming() {
match stream {
Ok(stream) => {
let server = server.clone();
thread::spawn(move || {
let reader = stream.try_clone().unwrap();
let writer = stream;
if let Err(e) = server.serve(reader, writer) {
eprintln!("Connection error: {}", e);
}
});
}
Err(e) => eprintln!("Accept error: {}", e),
}
}
}
Ok(())
}
fn print_help() {
eprintln!("ork-server - 9P server for org-mode files
USAGE:
ork-server [OPTIONS] [ORG_DIR]
OPTIONS:
-d, --dir <DIR> Directory containing org files (default: current dir)
-a, --addr <ADDR> Listen address (default: tcp://localhost:5640)
Examples: tcp://0.0.0.0:5640, unix:///tmp/ork.sock
-h, --help Print help
EXAMPLES:
ork-server ~/org
ork-server -a unix:///tmp/ork.sock ~/org
ork-server -a tcp://0.0.0.0:5640 -d ~/org
");
}

View File

@ -5,5 +5,7 @@
pub mod virtfs;
pub mod state;
pub mod namespace;
pub mod p9;
pub use state::OrkState;
pub use p9::Server;

678
crates/ork-server/src/p9.rs Normal file
View File

@ -0,0 +1,678 @@
//! 9P2000 server implementation.
//!
//! Handles the 9P protocol and maps operations to the virtual filesystem.
use crate::virtfs::FsNode;
use nine::p2000::*;
use std::collections::HashMap;
use std::io::{self, Read, Write};
use std::sync::RwLock;
use std::hash::{Hash, Hasher};
use std::collections::hash_map::DefaultHasher;
/// Maximum message size.
const MSIZE: u32 = 8192;
/// 9P protocol version.
const VERSION: &str = "9P2000";
/// Fid state: tracks open files.
struct Fid {
path: String,
node: FsNode,
open_mode: Option<OpenMode>,
dir_offset: usize, // For directory reads
read_cache: Vec<u8>, // Cached content for reads
write_buf: Vec<u8>, // Buffer for writes
}
/// 9P server state.
pub struct Server {
root: FsNode,
fids: RwLock<HashMap<u32, Fid>>,
msize: u32,
}
impl Server {
/// Create a new 9P server with the given root filesystem.
pub fn new(root: FsNode) -> Self {
Self {
root,
fids: RwLock::new(HashMap::new()),
msize: MSIZE,
}
}
/// Serve a single connection.
pub fn serve<R: Read, W: Write>(&self, mut reader: R, mut writer: W) -> io::Result<()> {
let mut buf = vec![0u8; MSIZE as usize];
loop {
// Read message size (4 bytes)
let mut size_buf = [0u8; 4];
match reader.read_exact(&mut size_buf) {
Ok(_) => {}
Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(()),
Err(e) => return Err(e),
}
let size = u32::from_le_bytes(size_buf) as usize;
if size < 7 || size > self.msize as usize {
return Err(io::Error::new(io::ErrorKind::InvalidData, "invalid message size"));
}
// Read rest of message
buf.resize(size, 0);
buf[0..4].copy_from_slice(&size_buf);
reader.read_exact(&mut buf[4..size])?;
// Parse and handle message
let response = self.handle_message(&buf[..size])?;
// Write response
writer.write_all(&response)?;
writer.flush()?;
}
}
/// Handle a single 9P message.
fn handle_message(&self, data: &[u8]) -> io::Result<Vec<u8>> {
if data.len() < 7 {
return self.error_response(0, "message too short");
}
let msg_type = data[4];
let tag = u16::from_le_bytes([data[5], data[6]]);
let result = match msg_type {
100 => self.handle_tversion(data, tag),
102 => self.handle_tauth(tag),
104 => self.handle_tattach(data, tag),
110 => self.handle_twalk(data, tag),
112 => self.handle_topen(data, tag),
114 => self.handle_tcreate(data, tag),
116 => self.handle_tread(data, tag),
118 => self.handle_twrite(data, tag),
120 => self.handle_tclunk(data, tag),
122 => self.handle_tremove(data, tag),
124 => self.handle_tstat(data, tag),
126 => self.handle_twstat(data, tag),
108 => self.handle_tflush(tag),
_ => self.error_response(tag, &format!("unknown message type: {}", msg_type)),
};
result
}
/// Create an error response.
fn error_response(&self, tag: u16, msg: &str) -> io::Result<Vec<u8>> {
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]); // size placeholder
resp.push(107); // Rerror
resp.extend_from_slice(&tag.to_le_bytes());
// String: 2-byte length + data
let msg_bytes = msg.as_bytes();
resp.extend_from_slice(&(msg_bytes.len() as u16).to_le_bytes());
resp.extend_from_slice(msg_bytes);
// Fill in size
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tversion handler.
fn handle_tversion(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
// Parse client msize and version
if data.len() < 13 {
return self.error_response(tag, "tversion too short");
}
let client_msize = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let msize = client_msize.min(self.msize);
// Build Rversion
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]); // size placeholder
resp.push(101); // Rversion
resp.extend_from_slice(&tag.to_le_bytes());
resp.extend_from_slice(&msize.to_le_bytes());
let version = VERSION.as_bytes();
resp.extend_from_slice(&(version.len() as u16).to_le_bytes());
resp.extend_from_slice(version);
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tauth handler (no auth required).
fn handle_tauth(&self, tag: u16) -> io::Result<Vec<u8>> {
self.error_response(tag, "no authentication required")
}
/// Tflush handler.
fn handle_tflush(&self, tag: u16) -> io::Result<Vec<u8>> {
// Just acknowledge the flush
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(109); // Rflush
resp.extend_from_slice(&tag.to_le_bytes());
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tattach handler.
fn handle_tattach(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 19 {
return self.error_response(tag, "tattach too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
// afid at 11-14, uname and aname follow
// Create root fid
let mut fids = self.fids.write().unwrap();
fids.insert(fid, Fid {
path: "/".to_string(),
node: self.root.clone(),
open_mode: None,
dir_offset: 0,
read_cache: Vec::new(),
write_buf: Vec::new(),
});
// Build Rattach
let qid = self.node_qid(&self.root, "/");
self.qid_response(105, tag, &qid)
}
/// Twalk handler.
fn handle_twalk(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 17 {
return self.error_response(tag, "twalk too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let newfid = u32::from_le_bytes([data[11], data[12], data[13], data[14]]);
let nwname = u16::from_le_bytes([data[15], data[16]]) as usize;
// Parse path elements
let mut offset = 17;
let mut wnames = Vec::new();
for _ in 0..nwname {
if offset + 2 > data.len() {
return self.error_response(tag, "twalk truncated");
}
let len = u16::from_le_bytes([data[offset], data[offset + 1]]) as usize;
offset += 2;
if offset + len > data.len() {
return self.error_response(tag, "twalk truncated");
}
let name = String::from_utf8_lossy(&data[offset..offset + len]).to_string();
offset += len;
wnames.push(name);
}
// Get source fid
let fids = self.fids.read().unwrap();
let source = match fids.get(&fid) {
Some(f) => f,
None => return self.error_response(tag, "unknown fid"),
};
// Walk the path
let mut current_path = source.path.clone();
let mut current_node = source.node.clone();
let mut qids = Vec::new();
for name in &wnames {
if name == ".." {
// Go up one level
if let Some(pos) = current_path.rfind('/') {
if pos > 0 {
current_path = current_path[..pos].to_string();
} else {
current_path = "/".to_string();
}
}
// Re-walk from root
current_node = match self.root.walk(&current_path[1..]) {
Some(n) => n,
None => return self.error_response(tag, "walk failed"),
};
} else {
// Walk to child
let new_node = match current_node.walk(name) {
Some(n) => n,
None => {
// Partial walk - return qids so far
break;
}
};
if current_path == "/" {
current_path = format!("/{}", name);
} else {
current_path = format!("{}/{}", current_path, name);
}
current_node = new_node;
}
qids.push(self.node_qid(&current_node, &current_path));
}
drop(fids);
// If we walked all names, create/update newfid
if qids.len() == wnames.len() {
let mut fids = self.fids.write().unwrap();
fids.insert(newfid, Fid {
path: current_path,
node: current_node,
open_mode: None,
dir_offset: 0,
read_cache: Vec::new(),
write_buf: Vec::new(),
});
}
// Build Rwalk
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(111); // Rwalk
resp.extend_from_slice(&tag.to_le_bytes());
resp.extend_from_slice(&(qids.len() as u16).to_le_bytes());
for qid in &qids {
self.encode_qid(&mut resp, qid);
}
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Topen handler.
fn handle_topen(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 12 {
return self.error_response(tag, "topen too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let mode = OpenMode::from_bits_truncate(data[11]);
let mut fids = self.fids.write().unwrap();
let f = match fids.get_mut(&fid) {
Some(f) => f,
None => return self.error_response(tag, "unknown fid"),
};
// Check permissions
let node_mode = f.node.mode;
if mode.is_readable() && (node_mode & 0o444) == 0 && !f.node.is_dir {
return self.error_response(tag, "permission denied");
}
if mode.is_writable() && (node_mode & 0o222) == 0 {
return self.error_response(tag, "permission denied");
}
f.open_mode = Some(mode);
f.dir_offset = 0;
f.read_cache.clear();
f.write_buf.clear();
let qid = self.node_qid(&f.node, &f.path);
let iounit = self.msize - 24; // Leave room for headers
// Build Ropen
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(113); // Ropen
resp.extend_from_slice(&tag.to_le_bytes());
self.encode_qid(&mut resp, &qid);
resp.extend_from_slice(&iounit.to_le_bytes());
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tcreate handler.
fn handle_tcreate(&self, _data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
self.error_response(tag, "create not supported")
}
/// Tread handler.
fn handle_tread(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 23 {
return self.error_response(tag, "tread too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let offset = u64::from_le_bytes([
data[11], data[12], data[13], data[14],
data[15], data[16], data[17], data[18],
]) as usize;
let count = u32::from_le_bytes([data[19], data[20], data[21], data[22]]) as usize;
let mut fids = self.fids.write().unwrap();
let f = match fids.get_mut(&fid) {
Some(f) => f,
None => return self.error_response(tag, "unknown fid"),
};
if f.open_mode.is_none() {
return self.error_response(tag, "not open");
}
let read_data = if f.node.is_dir {
// Directory read
self.read_dir(f, offset, count)?
} else {
// File read
if f.read_cache.is_empty() {
f.read_cache = f.node.read().unwrap_or_default();
}
if offset >= f.read_cache.len() {
Vec::new()
} else {
let end = (offset + count).min(f.read_cache.len());
f.read_cache[offset..end].to_vec()
}
};
// Build Rread
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(117); // Rread
resp.extend_from_slice(&tag.to_le_bytes());
resp.extend_from_slice(&(read_data.len() as u32).to_le_bytes());
resp.extend_from_slice(&read_data);
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Read directory entries.
fn read_dir(&self, fid: &mut Fid, offset: usize, count: usize) -> io::Result<Vec<u8>> {
// Get all entries and encode them
if fid.read_cache.is_empty() || offset == 0 {
fid.read_cache.clear();
let entries = fid.node.list().unwrap_or_default();
for name in entries {
if let Some(child) = fid.node.walk(&name) {
let child_path = if fid.path == "/" {
format!("/{}", name)
} else {
format!("{}/{}", fid.path, name)
};
let stat = self.node_stat(&child, &child_path, &name);
self.encode_stat(&mut fid.read_cache, &stat);
}
}
}
if offset >= fid.read_cache.len() {
return Ok(Vec::new());
}
let end = (offset + count).min(fid.read_cache.len());
Ok(fid.read_cache[offset..end].to_vec())
}
/// Twrite handler.
fn handle_twrite(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 23 {
return self.error_response(tag, "twrite too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let _offset = u64::from_le_bytes([
data[11], data[12], data[13], data[14],
data[15], data[16], data[17], data[18],
]);
let count = u32::from_le_bytes([data[19], data[20], data[21], data[22]]) as usize;
if data.len() < 23 + count {
return self.error_response(tag, "twrite truncated");
}
let write_data = &data[23..23 + count];
let mut fids = self.fids.write().unwrap();
let f = match fids.get_mut(&fid) {
Some(f) => f,
None => return self.error_response(tag, "unknown fid"),
};
let mode = match f.open_mode {
Some(m) => m,
None => return self.error_response(tag, "not open"),
};
if !mode.is_writable() {
return self.error_response(tag, "not open for writing");
}
// Check if this is rdwr mode (has rdwr handler)
if f.node.rdwr_fn.is_some() {
// For rdwr files, write triggers immediate response
match f.node.rdwr(write_data) {
Ok(response) => {
f.read_cache = response;
}
Err(e) => return self.error_response(tag, &e.to_string()),
}
} else if f.node.write_fn.is_some() {
// Buffer writes for regular write files
f.write_buf.extend_from_slice(write_data);
} else {
return self.error_response(tag, "not writable");
}
// Build Rwrite
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(119); // Rwrite
resp.extend_from_slice(&tag.to_le_bytes());
resp.extend_from_slice(&(count as u32).to_le_bytes());
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tclunk handler.
fn handle_tclunk(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 11 {
return self.error_response(tag, "tclunk too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let mut fids = self.fids.write().unwrap();
if let Some(f) = fids.remove(&fid) {
// Flush write buffer if not rdwr mode
if !f.write_buf.is_empty() && f.node.rdwr_fn.is_none() {
if let Err(e) = f.node.write(&f.write_buf) {
return self.error_response(tag, &e.to_string());
}
}
}
// Build Rclunk
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(121); // Rclunk
resp.extend_from_slice(&tag.to_le_bytes());
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Tremove handler.
fn handle_tremove(&self, _data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
self.error_response(tag, "remove not supported")
}
/// Tstat handler.
fn handle_tstat(&self, data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
if data.len() < 11 {
return self.error_response(tag, "tstat too short");
}
let fid = u32::from_le_bytes([data[7], data[8], data[9], data[10]]);
let fids = self.fids.read().unwrap();
let f = match fids.get(&fid) {
Some(f) => f,
None => return self.error_response(tag, "unknown fid"),
};
let name = f.path.rsplit('/').next().unwrap_or("");
let name = if name.is_empty() { "/" } else { name };
let stat = self.node_stat(&f.node, &f.path, name);
// Build Rstat
let mut stat_buf = Vec::new();
self.encode_stat(&mut stat_buf, &stat);
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(125); // Rstat
resp.extend_from_slice(&tag.to_le_bytes());
resp.extend_from_slice(&(stat_buf.len() as u16).to_le_bytes());
resp.extend_from_slice(&stat_buf);
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Twstat handler.
fn handle_twstat(&self, _data: &[u8], tag: u16) -> io::Result<Vec<u8>> {
self.error_response(tag, "wstat not supported")
}
// ═══════════════════════════════════════════════════════════════════
// Helpers
// ═══════════════════════════════════════════════════════════════════
/// Generate a Qid for a node.
fn node_qid(&self, node: &FsNode, path: &str) -> Qid {
let file_type = if node.is_dir {
FileType::DIR
} else {
FileType::FILE
};
let mut hasher = DefaultHasher::new();
path.hash(&mut hasher);
let path_hash = hasher.finish();
Qid {
file_type,
version: 0,
path: path_hash,
}
}
/// Generate a Stat for a node.
fn node_stat(&self, node: &FsNode, path: &str, name: &str) -> Stat {
let qid = self.node_qid(node, path);
let mut mode = FileMode::from_bits_truncate(node.mode);
if node.is_dir {
mode |= FileMode::DIR;
}
// Estimate length
let length = if node.is_dir {
0
} else {
node.read().map(|d| d.len() as u64).unwrap_or(0)
};
Stat {
type_: 0,
dev: 0,
qid,
mode,
atime: 0,
mtime: 0,
length,
name: name.to_string().into(),
uid: "ork".to_string().into(),
gid: "ork".to_string().into(),
muid: "ork".to_string().into(),
}
}
/// Build a response with just a qid.
fn qid_response(&self, msg_type: u8, tag: u16, qid: &Qid) -> io::Result<Vec<u8>> {
let mut resp = Vec::new();
resp.extend_from_slice(&[0u8; 4]);
resp.push(msg_type);
resp.extend_from_slice(&tag.to_le_bytes());
self.encode_qid(&mut resp, qid);
let size = resp.len() as u32;
resp[0..4].copy_from_slice(&size.to_le_bytes());
Ok(resp)
}
/// Encode a Qid into a buffer.
fn encode_qid(&self, buf: &mut Vec<u8>, qid: &Qid) {
buf.push(qid.file_type.bits());
buf.extend_from_slice(&qid.version.to_le_bytes());
buf.extend_from_slice(&qid.path.to_le_bytes());
}
/// Encode a Stat into a buffer.
fn encode_stat(&self, buf: &mut Vec<u8>, stat: &Stat) {
// Stat has a leading 2-byte size field
let start = buf.len();
buf.extend_from_slice(&[0u8; 2]); // size placeholder
buf.extend_from_slice(&stat.type_.to_le_bytes());
buf.extend_from_slice(&stat.dev.to_le_bytes());
self.encode_qid(buf, &stat.qid);
buf.extend_from_slice(&stat.mode.bits().to_le_bytes());
buf.extend_from_slice(&stat.atime.to_le_bytes());
buf.extend_from_slice(&stat.mtime.to_le_bytes());
buf.extend_from_slice(&stat.length.to_le_bytes());
// Strings
self.encode_string(buf, &stat.name);
self.encode_string(buf, &stat.uid);
self.encode_string(buf, &stat.gid);
self.encode_string(buf, &stat.muid);
// Fill in size
let stat_size = (buf.len() - start - 2) as u16;
buf[start..start + 2].copy_from_slice(&stat_size.to_le_bytes());
}
/// Encode a string into a buffer (2-byte length prefix).
fn encode_string(&self, buf: &mut Vec<u8>, s: &str) {
let bytes = s.as_bytes();
buf.extend_from_slice(&(bytes.len() as u16).to_le_bytes());
buf.extend_from_slice(bytes);
}
}