use crate::chat::edits::PendingEditStore;
use crate::fswalk::path_to_forward_slash;
use crate::markdown::extract_referenced_assets_for_file;
use crate::search::SearchIndex;
use crate::workspace_fs::WorkspaceFs;
use arc_swap::ArcSwapOption;
use notify::{
event::{CreateKind, ModifyKind, RemoveKind},
EventKind, RecursiveMode, Watcher,
};
use serde::{Deserialize, Serialize};
use std::{
collections::{BTreeSet, HashMap, HashSet},
path::{Path, PathBuf},
sync::{
atomic::{AtomicBool, Ordering},
Arc, RwLock,
},
};
use tokio::sync::broadcast;
const LIVE_RELOAD_EXTENSIONS: &[&str] = &[
"md", "markdown", "png", "jpg", "jpeg", "gif", "webp", "avif", "svg", "css", "js",
];
const LIVE_RELOAD_IGNORED_DIRS: &[&str] = &[".git", "node_modules", "target"];
const WATCH_STOP_POLL: std::time::Duration = std::time::Duration::from_millis(500);
const WATCH_DEBOUNCE: std::time::Duration = std::time::Duration::from_millis(250);
const WATCH_MAX_BATCH_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkspaceFlags {
#[serde(default)]
pub enable_search: bool,
#[serde(default)]
pub enable_viewed: bool,
#[serde(default)]
pub enable_edit: bool,
#[serde(default)]
pub enable_live: bool,
#[serde(default)]
pub enable_chat: bool,
#[serde(default)]
pub shared_annotation: bool,
}
#[derive(Clone, Default)]
pub struct WorkspaceConfig {
pub path: PathBuf,
pub flags: WorkspaceFlags,
pub single_file: Option<String>,
pub collaborator_access_code_hash: String,
pub alias: String,
}
pub(crate) struct WorkspaceEntry {
pub id: String,
pub fs: Arc<WorkspaceFs>,
pub enable_search: AtomicBool,
pub enable_viewed: AtomicBool,
pub enable_edit: AtomicBool,
pub enable_live: AtomicBool,
pub enable_chat: AtomicBool,
pub shared_annotation: AtomicBool,
pub config_tx: broadcast::Sender<()>,
pub events_tx: broadcast::Sender<WorkspaceEvent>,
pub search_index: ArcSwapOption<SearchIndex>,
pub single_file: Option<String>,
pub pending_edits: Arc<PendingEditStore>,
pub collaborator_access_code_hash: RwLock<String>,
pub alias: RwLock<String>,
stopped: Arc<AtomicBool>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) enum WorkspaceEvent {
Channel { channel: String, payload: String },
Workspace { payload: String },
}
impl WorkspaceEntry {
pub(crate) fn search_ready(&self) -> bool {
self.enable_search.load(Ordering::Relaxed) && self.search_index.load().is_some()
}
pub(crate) fn flags(&self) -> WorkspaceFlags {
WorkspaceFlags {
enable_search: self.enable_search.load(Ordering::Relaxed),
enable_viewed: self.enable_viewed.load(Ordering::Relaxed),
enable_edit: self.enable_edit.load(Ordering::Relaxed),
enable_live: self.enable_live.load(Ordering::Relaxed),
enable_chat: self.enable_chat.load(Ordering::Relaxed),
shared_annotation: self.shared_annotation.load(Ordering::Relaxed),
}
}
pub(crate) fn is_ephemeral(&self) -> bool {
self.fs.is_single_file()
}
pub(crate) fn collaborator_access_code_hash(&self) -> String {
self.collaborator_access_code_hash.read().unwrap().clone()
}
pub(crate) fn alias(&self) -> String {
self.alias.read().unwrap().clone()
}
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct WorkspaceInfo {
pub id: String,
pub path: String,
#[serde(flatten)]
pub flags: WorkspaceFlags,
pub search_ready: bool,
pub ephemeral: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub single_file: Option<String>,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub collaborator_access_code_hash: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub alias: String,
}
pub type PersistHook = Arc<dyn Fn(&WorkspaceRegistry) + Send + Sync>;
pub struct WorkspaceRegistry {
inner: RwLock<HashMap<String, Arc<WorkspaceEntry>>>,
pub(crate) salt: String,
persist: RwLock<Option<PersistHook>>,
}
pub fn hash_id(path: &Path, salt: &str) -> String {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(salt.as_bytes());
h.update(b"\0");
h.update(path.as_os_str().to_string_lossy().as_bytes());
let digest = h.finalize();
format!(
"{:02x}{:02x}{:02x}{:02x}",
digest[0], digest[1], digest[2], digest[3]
)
}
pub const MIN_ACCESS_CODE_LEN: usize = 8;
pub fn validate_access_code(code: &str) -> Result<(), String> {
let code = code.trim();
let len = code.chars().count();
if !code.is_empty() && len < MIN_ACCESS_CODE_LEN {
return Err(format!(
"access code must be at least {MIN_ACCESS_CODE_LEN} characters (got {len})"
));
}
Ok(())
}
pub fn hash_access_code(salt: &str, code: &str) -> String {
if code.is_empty() {
return String::new();
}
access_code_digest(salt, code)
}
fn access_code_digest(salt: &str, code: &str) -> String {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(salt.as_bytes());
h.update(b"\0mk-access\0");
h.update(code.as_bytes());
h.finalize().iter().map(|b| format!("{b:02x}")).collect()
}
pub fn access_code_matches(salt: &str, code: &str, stored: &str) -> bool {
if stored.is_empty() {
return false;
}
let full = access_code_digest(salt, code);
ct_eq(full.as_bytes(), stored.as_bytes())
}
pub(crate) fn ct_eq(a: &[u8], b: &[u8]) -> bool {
if a.len() != b.len() {
return false;
}
let mut diff = 0u8;
for (x, y) in a.iter().zip(b) {
diff |= x ^ y;
}
diff == 0
}
pub(crate) fn generate_token() -> String {
uuid::Uuid::new_v4().simple().to_string()
}
pub(crate) fn write_file_user_private(path: &Path, contents: &[u8]) -> std::io::Result<()> {
use std::io::Write;
static TMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let dir = path.parent().unwrap_or_else(|| Path::new("."));
let stem = path
.file_name()
.and_then(|s| s.to_str())
.unwrap_or("settings");
let seq = TMP_COUNTER.fetch_add(1, Ordering::Relaxed);
let tmp = dir.join(format!(".{stem}.tmp.{}.{seq}", std::process::id()));
let write_tmp = || -> std::io::Result<()> {
let mut opts = std::fs::OpenOptions::new();
opts.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
opts.mode(0o600);
}
let mut f = opts.open(&tmp)?;
f.write_all(contents)?;
f.sync_all()?;
Ok(())
};
if let Err(e) = write_tmp() {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
if let Err(e) = std::fs::rename(&tmp, path) {
tracing::warn!(
"atomic rename of {} onto {} failed: {e}",
tmp.display(),
path.display()
);
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
Ok(())
}
pub fn expand_and_canonicalize(raw: &str) -> std::io::Result<PathBuf> {
let normalized = if raw.starts_with('~') {
raw.replacen('~', "~", 1)
} else {
raw.to_string()
};
let expanded = if normalized.starts_with("~/") || normalized == "~" {
dirs::home_dir()
.map(|home| {
if normalized == "~" {
home
} else {
home.join(&normalized[2..])
}
})
.unwrap_or_else(|| PathBuf::from(&normalized))
} else {
PathBuf::from(&normalized)
};
dunce::canonicalize(&expanded).or(Ok(expanded))
}
impl WorkspaceRegistry {
pub fn new(salt: String) -> Self {
Self {
inner: RwLock::new(HashMap::new()),
salt,
persist: RwLock::new(None),
}
}
pub fn set_persist_hook(&self, hook: PersistHook) {
*self.persist.write().unwrap() = Some(hook);
}
fn notify_persist(&self) {
let hook = self.persist.read().unwrap().clone();
if let Some(hook) = hook {
hook(self);
}
}
pub fn add(&self, config: WorkspaceConfig) -> String {
let identity = match &config.single_file {
Some(name) => config.path.join(name),
None => config.path.clone(),
};
let id = hash_id(&identity, &self.salt);
if self.inner.read().unwrap().contains_key(&id) {
self.update_flags(&id, config.flags);
self.notify_persist();
return id;
}
let (config_tx, _) = broadcast::channel(4);
let (events_tx, _) = broadcast::channel(100);
let single_file = config.single_file.clone();
let workspace_fs = Arc::new(WorkspaceFs::new(
config.path.clone(),
single_file.as_deref(),
));
let entry = Arc::new(WorkspaceEntry {
id: id.clone(),
fs: workspace_fs,
enable_search: AtomicBool::new(config.flags.enable_search),
enable_viewed: AtomicBool::new(config.flags.enable_viewed),
enable_edit: AtomicBool::new(config.flags.enable_edit),
enable_live: AtomicBool::new(config.flags.enable_live),
enable_chat: AtomicBool::new(config.flags.enable_chat),
shared_annotation: AtomicBool::new(config.flags.shared_annotation),
config_tx,
events_tx,
search_index: ArcSwapOption::empty(),
single_file: single_file.clone(),
pending_edits: Arc::new(PendingEditStore::new()),
collaborator_access_code_hash: RwLock::new(config.collaborator_access_code_hash),
alias: RwLock::new(config.alias),
stopped: Arc::new(AtomicBool::new(false)),
});
self.inner
.write()
.unwrap()
.insert(id.clone(), entry.clone());
match single_file {
Some(name) => {
refresh_allowed_assets(&entry, &name);
if config.flags.enable_search {
spawn_search_indexer(entry.clone());
}
spawn_single_file_watcher(config.path, entry.clone(), name);
}
None => {
if config.flags.enable_search {
spawn_search_indexer(entry.clone());
}
spawn_directory_watcher(config.path, entry.clone());
}
}
self.notify_persist();
id
}
pub fn update_flags(&self, id: &str, flags: WorkspaceFlags) -> bool {
let guard = self.inner.read().unwrap();
let Some(entry) = guard.get(id).cloned() else {
return false;
};
drop(guard);
let was_search = entry
.enable_search
.swap(flags.enable_search, Ordering::Relaxed);
entry
.enable_viewed
.store(flags.enable_viewed, Ordering::Relaxed);
entry
.enable_edit
.store(flags.enable_edit, Ordering::Relaxed);
entry
.enable_live
.store(flags.enable_live, Ordering::Relaxed);
entry
.enable_chat
.store(flags.enable_chat, Ordering::Relaxed);
entry
.shared_annotation
.store(flags.shared_annotation, Ordering::Relaxed);
let _ = entry.config_tx.send(());
if flags.enable_search && !was_search && entry.search_index.load().is_none() {
spawn_search_indexer(entry);
} else if !flags.enable_search && was_search {
entry.search_index.store(None);
}
self.notify_persist();
true
}
pub fn remove(&self, id: &str) -> bool {
let removed = self.inner.write().unwrap().remove(id);
if let Some(entry) = &removed {
entry.stopped.store(true, Ordering::Relaxed);
let _ = entry.config_tx.send(());
self.notify_persist();
}
removed.is_some()
}
pub(crate) fn get(&self, id: &str) -> Option<Arc<WorkspaceEntry>> {
self.inner.read().unwrap().get(id).cloned()
}
pub fn set_collaborator_access_code(&self, id: &str, hash: &str) -> bool {
let guard = self.inner.read().unwrap();
let Some(entry) = guard.get(id) else {
return false;
};
*entry.collaborator_access_code_hash.write().unwrap() = hash.to_string();
let _ = entry.config_tx.send(());
drop(guard);
self.notify_persist();
true
}
pub fn set_alias(&self, id: &str, alias: &str) -> bool {
let guard = self.inner.read().unwrap();
let Some(entry) = guard.get(id) else {
return false;
};
*entry.alias.write().unwrap() = alias.to_string();
drop(guard);
self.notify_persist();
true
}
pub(crate) fn list(&self) -> Vec<Arc<WorkspaceEntry>> {
let mut v: Vec<_> = self.inner.read().unwrap().values().cloned().collect();
v.sort_by(|a, b| {
a.fs.ambient_root()
.cmp(b.fs.ambient_root())
.then_with(|| a.single_file.cmp(&b.single_file))
});
v
}
pub fn info_list(&self) -> Vec<WorkspaceInfo> {
self.list()
.into_iter()
.map(|e| WorkspaceInfo {
id: e.id.clone(),
path: e.fs.ambient_root().to_string_lossy().to_string(),
flags: e.flags(),
search_ready: e.search_ready(),
ephemeral: e.is_ephemeral(),
single_file: e.single_file.clone(),
collaborator_access_code_hash: e.collaborator_access_code_hash(),
alias: e.alias(),
})
.collect()
}
}
fn refresh_allowed_assets(entry: &WorkspaceEntry, file_name: &str) {
let root = entry.fs.ambient_root();
let abs = root.join(file_name);
let new_set = match std::fs::read_to_string(&abs) {
Ok(content) => extract_referenced_assets_for_file(&content, &abs, root),
Err(_) => HashSet::new(),
};
entry.fs.replace_assets(new_set);
}
fn spawn_watch_thread(
root: PathBuf,
mode: RecursiveMode,
stopped: Arc<AtomicBool>,
mut on_events: impl FnMut(Vec<notify::Event>) + Send + 'static,
) {
std::thread::spawn(move || {
let (tx, rx) = std::sync::mpsc::channel();
let Ok(mut watcher) = notify::recommended_watcher(move |res| {
if let Ok(e) = res {
let _ = tx.send(e);
}
}) else {
return;
};
if watcher.watch(&root, mode).is_err() {
return;
}
loop {
if stopped.load(Ordering::Relaxed) {
return;
}
match rx.recv_timeout(WATCH_STOP_POLL) {
Ok(first) => {
let started = std::time::Instant::now();
let mut events = vec![first];
let mut disconnected = false;
loop {
if stopped.load(Ordering::Relaxed) {
return;
}
let remaining = WATCH_MAX_BATCH_DELAY.saturating_sub(started.elapsed());
if remaining.is_zero() {
break;
}
match rx.recv_timeout(WATCH_DEBOUNCE.min(remaining)) {
Ok(event) => events.push(event),
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => break,
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
disconnected = true;
break;
}
}
}
if stopped.load(Ordering::Relaxed) {
return;
}
on_events(events);
if disconnected {
return;
}
}
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => continue,
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return,
}
}
});
}
fn spawn_single_file_watcher(root: PathBuf, entry: Arc<WorkspaceEntry>, file_name: String) {
let target = root.join(&file_name);
let stopped = entry.stopped.clone();
spawn_watch_thread(
root.clone(),
RecursiveMode::NonRecursive,
stopped,
move |events: Vec<notify::Event>| {
let mut pinned_changed = false;
let mut broadcast_paths = BTreeSet::new();
for event in events {
let is_create_or_modify =
matches!(event.kind, EventKind::Create(_) | EventKind::Modify(_));
let is_mutation = is_create_or_modify || matches!(event.kind, EventKind::Remove(_));
if !is_mutation {
continue;
}
for path in event.paths {
let Ok(rel) = path.strip_prefix(&root) else {
continue;
};
let touched_pinned = path == target;
let touched_asset = entry.fs.is_asset(rel);
if !(touched_pinned || touched_asset) {
continue;
}
pinned_changed |= touched_pinned;
if is_create_or_modify {
broadcast_paths.insert(rel.to_string_lossy().to_string());
}
}
}
if pinned_changed {
if target.is_file() {
refresh_allowed_assets(&entry, &file_name);
} else {
entry.fs.clear_assets();
broadcast_paths.remove(&file_name);
}
if let Some(idx) = entry.search_index.load_full() {
if let Err(error) = idx.reconcile_files(std::slice::from_ref(&target)) {
tracing::warn!("single-file search index update failed: {error}");
}
}
}
for rel_str in broadcast_paths {
let file_payload = serde_json::json!({
"type": "file_changed",
"workspace_id": entry.id,
"path": rel_str,
})
.to_string();
let _ = entry.events_tx.send(WorkspaceEvent::Workspace {
payload: file_payload,
});
}
},
);
}
fn spawn_search_indexer(entry: Arc<WorkspaceEntry>) {
std::thread::spawn(move || {
if let Ok(idx) = SearchIndex::for_workspace(entry.fs.clone()) {
entry.search_index.store(Some(Arc::new(idx)));
}
});
}
fn spawn_directory_watcher(root: PathBuf, entry: Arc<WorkspaceEntry>) {
let stopped = entry.stopped.clone();
spawn_watch_thread(
root.clone(),
RecursiveMode::Recursive,
stopped,
move |events: Vec<notify::Event>| {
let search_changes = coalesce_search_changes(&root, &events);
if let Some(idx) = entry.search_index.load_full() {
let result = if search_changes.rebuild {
if search_changes.paths.is_empty() {
idx.rebuild_if_routes_changed()
} else {
idx.rebuild()
}
} else {
idx.reconcile_files(&search_changes.paths)
};
if let Err(error) = result {
tracing::warn!("directory search index update failed: {error}");
}
}
let mut broadcast_paths = BTreeSet::new();
for event in events {
if !matches!(
event.kind,
EventKind::Create(_) | EventKind::Modify(_) | EventKind::Remove(_)
) {
continue;
}
for path in event.paths {
if let Some(rel_str) = directory_live_reload_path(&root, &path) {
broadcast_paths.insert(rel_str);
}
}
}
for rel_str in broadcast_paths {
let payload = serde_json::json!({
"type": "file_changed",
"workspace_id": entry.id,
"path": rel_str,
})
.to_string();
let _ = entry.events_tx.send(WorkspaceEvent::Workspace { payload });
}
},
);
}
#[derive(Debug, Default, PartialEq, Eq)]
struct SearchChangeBatch {
paths: Vec<PathBuf>,
rebuild: bool,
}
fn coalesce_search_changes(root: &Path, events: &[notify::Event]) -> SearchChangeBatch {
let mut paths = BTreeSet::new();
let mut rebuild = false;
for event in events {
if event.need_rescan() {
rebuild = true;
}
if !matches!(
event.kind,
EventKind::Create(_) | EventKind::Modify(_) | EventKind::Remove(_)
) {
continue;
}
let relevant_paths: Vec<_> = event
.paths
.iter()
.filter(|path| !is_search_event_path_ignored(root, path))
.collect();
if relevant_paths.is_empty() {
continue;
}
let explicit_directory_event = matches!(
event.kind,
EventKind::Create(CreateKind::Folder) | EventKind::Remove(RemoveKind::Folder)
);
let created_directory = matches!(
event.kind,
EventKind::Create(CreateKind::Any | CreateKind::Other)
) && relevant_paths.iter().any(|path| path.is_dir());
let ambiguous_removed_directory =
matches!(
event.kind,
EventKind::Remove(RemoveKind::Any | RemoveKind::Other)
) && relevant_paths.iter().any(|path| path.extension().is_none());
let renamed_directory = matches!(event.kind, EventKind::Modify(ModifyKind::Name(_)))
&& relevant_paths.iter().any(|path| path.extension().is_none());
rebuild |= explicit_directory_event
|| created_directory
|| ambiguous_removed_directory
|| renamed_directory;
for path in relevant_paths {
let rel = path.strip_prefix(root).unwrap_or(path);
if is_search_ignore_file(rel) {
rebuild = true;
}
if path.extension().is_some_and(|ext| ext == "md") {
paths.insert(path.clone());
}
}
}
SearchChangeBatch {
paths: paths.into_iter().collect(),
rebuild,
}
}
fn is_search_event_path_ignored(root: &Path, path: &Path) -> bool {
let rel = path.strip_prefix(root).unwrap_or(path);
let components: Vec<_> = rel.components().map(|part| part.as_os_str()).collect();
let retained_rule_suffix = if components.last().is_some_and(|name| {
*name == std::ffi::OsStr::new(".gitignore") || *name == std::ffi::OsStr::new(".ignore")
}) {
1
} else if components.ends_with(&[
std::ffi::OsStr::new(".git"),
std::ffi::OsStr::new("info"),
std::ffi::OsStr::new("exclude"),
]) {
3
} else {
0
};
components[..components.len().saturating_sub(retained_rule_suffix)]
.iter()
.any(|component| {
let name = component.to_string_lossy();
name.starts_with('.')
|| LIVE_RELOAD_IGNORED_DIRS
.iter()
.any(|ignored| name.eq_ignore_ascii_case(ignored))
})
}
fn is_search_ignore_file(rel: &Path) -> bool {
if rel
.file_name()
.is_some_and(|name| name == ".gitignore" || name == ".ignore")
{
return true;
}
let components: Vec<_> = rel.components().map(|part| part.as_os_str()).collect();
components.ends_with(&[
std::ffi::OsStr::new(".git"),
std::ffi::OsStr::new("info"),
std::ffi::OsStr::new("exclude"),
])
}
fn directory_live_reload_path(root: &Path, path: &Path) -> Option<String> {
let rel = path.strip_prefix(root).ok()?;
if rel.as_os_str().is_empty()
|| rel.components().any(|component| {
let name = component.as_os_str().to_string_lossy();
LIVE_RELOAD_IGNORED_DIRS
.iter()
.any(|ignored| name.eq_ignore_ascii_case(ignored))
})
{
return None;
}
let ext = rel.extension()?.to_string_lossy().to_ascii_lowercase();
if !LIVE_RELOAD_EXTENSIONS.contains(&ext.as_str()) {
return None;
}
Some(path_to_forward_slash(rel))
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
pub struct ServerLock {
pub port: u16,
#[serde(default)]
pub control_socket: String,
#[serde(default)]
pub host: String,
#[serde(default)]
pub advertised_host: Option<String>,
#[serde(default)]
pub owner: String,
}
impl ServerLock {
pub(crate) fn path() -> PathBuf {
dirs::home_dir()
.expect("HOME directory required")
.join(".markon")
.join("server.lock")
}
pub(crate) fn write(&self) -> std::io::Result<()> {
Self::with_write_lock(|| {
let path = Self::path();
write_file_user_private(&path, serde_json::to_string(self).unwrap().as_bytes())
})
}
pub fn read() -> Option<Self> {
let path = Self::path();
let content = match std::fs::read_to_string(&path) {
Ok(c) => c,
Err(e) => {
if e.kind() != std::io::ErrorKind::NotFound {
tracing::warn!("cannot read server lock {}: {e}", path.display());
}
return None;
}
};
match serde_json::from_str(&content) {
Ok(v) => Some(v),
Err(e) => {
tracing::warn!(
"corrupted server lock file {}: {e}; ignoring",
path.display()
);
None
}
}
}
fn with_write_lock<T>(operation: impl FnOnce() -> std::io::Result<T>) -> std::io::Result<T> {
let path = Self::path();
let parent = path.parent().unwrap_or_else(|| Path::new("."));
std::fs::create_dir_all(parent)?;
let lock_path = parent.join("server.write.lock");
let mut options = std::fs::OpenOptions::new();
options.read(true).write(true).create(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
options.mode(0o600);
}
let lock = options.open(lock_path)?;
lock.lock()?;
let result = operation();
let unlock_result = lock.unlock();
match result {
Ok(value) => {
unlock_result?;
Ok(value)
}
Err(error) => Err(error),
}
}
pub(crate) fn remove_if_owned(owner: &str) {
let _ = Self::with_write_lock(|| {
let owned = Self::read().is_some_and(|lock| lock.owner == owner);
if owned {
match std::fs::remove_file(Self::path()) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(error),
}
}
Ok(())
});
}
pub fn is_alive(&self) -> bool {
if !self.control_socket.is_empty()
&& crate::control::transport::probe(&crate::control::ControlSocketName::from_raw(
self.control_socket.clone(),
))
{
return true;
}
let connect_host = if crate::net::host_is_wildcard_v6(&self.host) {
"::1"
} else if crate::net::host_is_wildcard_v4(&self.host) {
"127.0.0.1"
} else {
self.host.as_str()
};
let Ok(addr) = crate::net::bind_socket_addr(connect_host, self.port) else {
return false;
};
std::net::TcpStream::connect_timeout(&addr, std::time::Duration::from_millis(500)).is_ok()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn hash_id_is_deterministic() {
let p = std::path::Path::new("/tmp/test");
assert_eq!(hash_id(p, "s"), hash_id(p, "s"));
}
#[test]
fn hash_id_depends_on_salt() {
let p = std::path::Path::new("/tmp/test");
assert_ne!(hash_id(p, "a"), hash_id(p, "b"));
}
#[test]
fn registry_directory_id_matches_hash_contract() {
let temp_dir = tempfile::TempDir::new().unwrap();
let root = temp_dir.path().to_path_buf();
let salt = "contract-salt";
let registry = WorkspaceRegistry::new(salt.into());
let id = registry.add(WorkspaceConfig {
path: root.clone(),
flags: WorkspaceFlags::default(),
single_file: None,
collaborator_access_code_hash: String::new(),
..Default::default()
});
assert_eq!(id, hash_id(&root, salt));
}
#[test]
fn access_code_hash_is_full_width() {
assert_eq!(hash_access_code("s", "abcd12345678").len(), 64);
assert_eq!(hash_access_code("s", "x").len(), 64);
assert_eq!(hash_access_code("s", "").len(), 0);
}
#[test]
fn access_code_matches_current_scheme() {
let stored = hash_access_code("s", "test1234");
assert_eq!(stored.len(), 64);
assert!(access_code_matches("s", "test1234", &stored));
assert!(!access_code_matches("s", "test1235", &stored));
assert!(!access_code_matches("s", "test12345", &stored));
assert!(!access_code_matches("s", "anything", ""));
}
#[test]
fn access_code_rejects_legacy_truncated_hash() {
let full = access_code_digest("s", "test1234");
let legacy = full[..8].to_string();
assert!(!access_code_matches("s", "test1234", &legacy));
assert!(!access_code_matches("s", "wrongone", &legacy));
}
#[test]
fn access_code_matches_rejects_truncation_collision() {
let stored = hash_access_code("s", "test1234");
for guess in ["test1235", "TEST1234", "test123", "test12340"] {
assert!(!access_code_matches("s", guess, &stored));
}
}
#[test]
fn validate_access_code_enforces_floor() {
assert!(validate_access_code("").is_ok());
assert!(validate_access_code(" ").is_ok());
assert!(validate_access_code("short").is_err());
assert!(validate_access_code(&"x".repeat(MIN_ACCESS_CODE_LEN - 1)).is_err());
assert!(validate_access_code(&"x".repeat(MIN_ACCESS_CODE_LEN)).is_ok());
assert!(validate_access_code("a-reasonable-code").is_ok());
}
#[test]
fn directory_live_reload_filter_tracks_docs_and_assets_only() {
let root = Path::new("/repo");
assert_eq!(
directory_live_reload_path(root, &root.join("docs").join("a.md")).as_deref(),
Some("docs/a.md")
);
assert_eq!(
directory_live_reload_path(root, &root.join("assets").join("app.js")).as_deref(),
Some("assets/app.js")
);
assert_eq!(
directory_live_reload_path(root, &root.join("img").join("hero.PNG")).as_deref(),
Some("img/hero.PNG")
);
assert!(directory_live_reload_path(root, &root.join(".git").join("HEAD")).is_none());
assert!(
directory_live_reload_path(root, &root.join("node_modules").join("x.md")).is_none()
);
assert!(directory_live_reload_path(root, &root.join("target").join("x.css")).is_none());
assert!(directory_live_reload_path(root, &root.join("README")).is_none());
assert!(directory_live_reload_path(root, &root.join("notes.txt")).is_none());
}
#[test]
fn search_change_batch_deduplicates_markdown_paths() {
let root = Path::new("/repo");
let first = root.join("docs").join("a.md");
let second = root.join("docs").join("b.md");
let modify_kind = EventKind::Modify(ModifyKind::Data(notify::event::DataChange::Content));
let events = vec![
notify::Event::new(modify_kind).add_path(first.clone()),
notify::Event::new(modify_kind).add_path(first.clone()),
notify::Event::new(EventKind::Create(CreateKind::File)).add_path(second.clone()),
];
let batch = coalesce_search_changes(root, &events);
assert!(!batch.rebuild);
assert_eq!(batch.paths, vec![first, second]);
}
#[test]
fn search_change_batch_rebuilds_for_ignore_and_directory_changes() {
let root = Path::new("/repo");
let ignore_event = notify::Event::new(EventKind::Modify(ModifyKind::Data(
notify::event::DataChange::Content,
)))
.add_path(root.join("docs").join(".gitignore"));
let directory_event = notify::Event::new(EventKind::Remove(RemoveKind::Folder))
.add_path(root.join("old-docs"));
assert!(coalesce_search_changes(root, &[ignore_event]).rebuild);
assert!(coalesce_search_changes(root, &[directory_event]).rebuild);
assert!(is_search_ignore_file(Path::new(".git/info/exclude")));
assert!(!is_search_ignore_file(Path::new(".git/index")));
}
#[test]
fn search_change_batch_drops_intrinsically_ignored_paths() {
let root = Path::new("/repo");
let ignored_file = root.join("node_modules").join("pkg").join("README.md");
let hidden_file = root.join(".cache").join("generated.md");
let ignored_directory = root.join("target").join("generated-docs");
let events = vec![
notify::Event::new(EventKind::Modify(ModifyKind::Data(
notify::event::DataChange::Content,
)))
.add_path(ignored_file),
notify::Event::new(EventKind::Create(CreateKind::File)).add_path(hidden_file),
notify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(ignored_directory),
];
assert_eq!(
coalesce_search_changes(root, &events),
SearchChangeBatch::default()
);
assert!(is_search_event_path_ignored(
root,
&root.join(".hidden").join("note.md")
));
assert!(!is_search_event_path_ignored(
root,
&root.join("docs").join(".gitignore")
));
assert!(!is_search_event_path_ignored(
root,
&root.join(".git").join("info").join("exclude")
));
assert!(is_search_event_path_ignored(
root,
&root.join("node_modules").join("pkg").join(".gitignore")
));
}
#[test]
fn workspace_list_is_deterministically_ordered_by_path() {
let tmp = tempfile::TempDir::new().unwrap();
let base = tmp.path();
std::fs::create_dir_all(base.join("alpha")).unwrap();
std::fs::create_dir_all(base.join("charlie")).unwrap();
std::fs::write(base.join("alpha").join("a.md"), "# a").unwrap();
std::fs::write(base.join("alpha").join("z.md"), "# z").unwrap();
let reg = WorkspaceRegistry::new("salt".into());
let mk = |path: PathBuf, single: Option<&str>| WorkspaceConfig {
path,
flags: WorkspaceFlags::default(),
single_file: single.map(str::to_string),
collaborator_access_code_hash: String::new(),
..Default::default()
};
reg.add(mk(base.join("charlie"), None));
reg.add(mk(base.join("alpha"), Some("z.md")));
reg.add(mk(base.join("alpha"), None));
reg.add(mk(base.join("alpha"), Some("a.md")));
let order: Vec<(PathBuf, Option<String>)> = reg
.list()
.iter()
.map(|e| (e.fs.ambient_root().to_path_buf(), e.single_file.clone()))
.collect();
assert_eq!(
order,
vec![
(base.join("alpha"), None),
(base.join("alpha"), Some("a.md".into())),
(base.join("alpha"), Some("z.md".into())),
(base.join("charlie"), None),
],
"list() must be sorted by (root, single_file)"
);
let a: Vec<String> = reg.list().iter().map(|e| e.id.clone()).collect();
let b: Vec<String> = reg.list().iter().map(|e| e.id.clone()).collect();
assert_eq!(a, b);
}
#[test]
fn server_lock_optional_fields_default_when_absent() {
let old = r#"{"port":6419,"token":"abc"}"#;
let lock: ServerLock = serde_json::from_str(old).unwrap();
assert_eq!(lock.port, 6419);
assert_eq!(lock.control_socket, "");
assert_eq!(lock.host, "");
assert_eq!(lock.advertised_host, None);
}
#[test]
fn server_lock_round_trips() {
let lock = ServerLock {
port: 6419,
control_socket: "/home/u/.markon/control.sock".into(),
host: "0.0.0.0".into(),
advertised_host: Some("192.168.1.20".into()),
owner: "owner-nonce".into(),
};
let json = serde_json::to_string(&lock).unwrap();
let back: ServerLock = serde_json::from_str(&json).unwrap();
assert_eq!(back.host, "0.0.0.0");
assert_eq!(back.port, 6419);
assert_eq!(back.control_socket, "/home/u/.markon/control.sock");
assert_eq!(back.advertised_host.as_deref(), Some("192.168.1.20"));
assert_eq!(back.owner, "owner-nonce");
}
fn wait_for_index(entry: &Arc<WorkspaceEntry>) -> Arc<SearchIndex> {
for _ in 0..200 {
if let Some(idx) = entry.search_index.load_full() {
return idx;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
panic!("search index was not built in time");
}
#[test]
fn single_file_workspace_search_no_sibling_leakage() {
let temp_dir = tempfile::TempDir::new().unwrap();
let dir = temp_dir.path();
std::fs::write(
dir.join("pinned.md"),
"# Pinned\nuniquepinnedtoken is here.",
)
.unwrap();
std::fs::write(
dir.join("sibling.md"),
"# Sibling\nuniquesiblingtoken stays private.",
)
.unwrap();
let registry = WorkspaceRegistry::new("test-salt".into());
let id = registry.add(WorkspaceConfig {
path: dir.to_path_buf(),
flags: WorkspaceFlags {
enable_search: true,
..Default::default()
},
single_file: Some("pinned.md".into()),
collaborator_access_code_hash: String::new(),
..Default::default()
});
let entry = registry.get(&id).unwrap();
assert!(entry.is_ephemeral());
let idx = wait_for_index(&entry);
assert_eq!(
idx.search("uniquesiblingtoken", 10).unwrap().len(),
0,
"single-file workspace leaked a sibling through search"
);
assert_eq!(
idx.search("uniquepinnedtoken", 10).unwrap().len(),
1,
"pinned file should be searchable"
);
}
#[test]
fn single_file_workspace_search_toggle_spawns_and_clears() {
let temp_dir = tempfile::TempDir::new().unwrap();
let dir = temp_dir.path();
std::fs::write(dir.join("note.md"), "# Note\ntoggletoken here.").unwrap();
let registry = WorkspaceRegistry::new("test-salt".into());
let id = registry.add(WorkspaceConfig {
path: dir.to_path_buf(),
flags: WorkspaceFlags::default(),
single_file: Some("note.md".into()),
collaborator_access_code_hash: String::new(),
..Default::default()
});
let entry = registry.get(&id).unwrap();
assert!(entry.search_index.load().is_none());
registry.update_flags(
&id,
WorkspaceFlags {
enable_search: true,
..Default::default()
},
);
let idx = wait_for_index(&entry);
assert_eq!(idx.search("toggletoken", 10).unwrap().len(), 1);
registry.update_flags(&id, WorkspaceFlags::default());
assert!(
entry.search_index.load().is_none(),
"disabling search must clear the single-file index"
);
}
}