markon-core 0.15.23

markon core - Mark it on.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
//! Control plane — the privileged, same-user management channel for an
//! **already-running** markon server.
//!
//! The "only one server per machine" invariant is shared behavior, not a
//! front-end detail: whoever starts second should hand its workspaces to the
//! server that is already up rather than start a competing one (which would open
//! a second connection to the same annotation database — see the single-instance
//! discussion around [`crate::workspace::ServerLock`]). Keeping the client here
//! lets the CLI (forwarding `markon <dir>` to a running daemon) and the GUI
//! (attaching as a controller when a daemon is already up) go through one
//! implementation.
//!
//! **Transport** (see [`transport`]): a cross-platform local socket — a Unix
//! domain socket (`~/.markon/control.sock`, `0600`) on unix, a named pipe on
//! Windows — carrying length-prefixed JSON frames. Connecting over that socket
//! *is* the authorization: it is reachable only by the same local user, so there
//! is no token. Privilege is "which listener you arrived on".

pub mod proto;
pub mod transport;

pub use proto::{ControlRequest, ControlResponse};
pub use transport::{
    bind, dispatch, serve, AdminBootstrapCodeFn, AdminBootstrapFn, ControlContext, ControlServer,
    ControlSocketName,
};

use crate::data_maintenance::{DataCleanupResult, DataCleanupStats};
use crate::workspace::{
    expand_and_canonicalize, WorkspaceFlags, WorkspaceInfo, WorkspaceOpenTarget,
};

/// Error talking to a running server's control socket.
#[derive(Debug, thiserror::Error)]
pub enum ControlError {
    /// The socket could not be reached / read / written, or a frame was
    /// malformed.
    #[error("control socket transport error: {0}")]
    Transport(#[from] std::io::Error),
    /// The server processed the request but reported a failure.
    #[error("running markon server rejected the request: {0}")]
    Server(String),
    /// The server answered, but with a response variant the caller didn't expect
    /// for this request (a protocol mismatch).
    #[error("unexpected control response for this request")]
    Unexpected,
}

/// A handle to a running server, addressed by its control-socket name (a
/// filesystem path on unix, a pipe name on Windows). Each management method
/// opens a fresh connection, sends one framed request, and reads one framed
/// response. It also carries the server's web TCP port so a front-end can build
/// browser/QR URLs without a second discovery step.
#[derive(Clone)]
pub struct RunningServer {
    socket: ControlSocketName,
    /// The server's web TCP port (for building browser URLs). `0` when unknown —
    /// e.g. a socket-only handle built directly in a test.
    web_port: u16,
    /// The server's bind host (from the discovery lock), for building the
    /// printed / opened URLs. Empty when unknown (a socket-only handle).
    web_host: String,
    /// The advertised host active in the *owning* server process (from the lock).
    /// A controller (GUI/CLI) attached to a daemon another process started with a
    /// different `--entry` must build featured/QR URLs from this, not its own
    /// preference. `None` for a socket-only handle or a pre-field lock.
    web_advertised_host: Option<String>,
    /// `markond` package version from the discovery lock. Empty for a legacy
    /// daemon that predates versioned discovery.
    service_version: String,
    /// The public entry URL prefix (`--entry`/`--qr`) active in the *owning*
    /// server process (from the lock). A controller that got no `--entry` of its
    /// own uses this to build featured/QR URLs against the daemon's public
    /// address instead of loopback. `None` for a socket-only handle or a
    /// pre-field lock.
    web_entry: Option<String>,
}

impl RunningServer {
    /// Build a handle for an explicit control-socket name (e.g. a test's temp
    /// socket). The web port is unknown (`0`); use [`RunningServer::from_lock`]
    /// or [`RunningServer::discover`] when a browser port is needed.
    pub fn new(socket: ControlSocketName) -> Self {
        Self {
            socket,
            web_port: 0,
            web_host: String::new(),
            web_advertised_host: None,
            service_version: String::new(),
            web_entry: None,
        }
    }

    /// Build a handle from a discovery [`ServerLock`](crate::workspace::ServerLock):
    /// resolves the recorded control socket (falling back to the default name for
    /// pre-split locks that predate the field) and captures the web port.
    pub fn from_lock(lock: &crate::workspace::ServerLock) -> Self {
        let socket = if lock.control_socket.is_empty() {
            ControlSocketName::default_name()
                .unwrap_or_else(|_| ControlSocketName::from_raw(String::new()))
        } else {
            ControlSocketName::from_raw(lock.control_socket.clone())
        };
        Self {
            socket,
            web_port: lock.port,
            web_host: lock.host.clone(),
            web_advertised_host: lock.advertised_host.clone(),
            service_version: lock.service_version.clone(),
            web_entry: lock.entry.clone(),
        }
    }

    /// Discover the machine's running server via the on-disk lock, returning a
    /// handle only when that server is actually live (the lock's liveness probe
    /// now prefers the control socket). `None` means "no server to attach to".
    pub fn discover() -> Option<Self> {
        let lock = crate::workspace::ServerLock::read()?;
        lock.is_alive().then(|| Self::from_lock(&lock))
    }

    /// The control-socket name this handle targets.
    pub fn socket(&self) -> &ControlSocketName {
        &self.socket
    }

    /// The server's web TCP port, for building browser/QR URLs. `0` if this
    /// handle was built without a known port (see [`RunningServer::new`]).
    pub fn port(&self) -> u16 {
        self.web_port
    }

    /// The server's bind host, for building browser/QR URLs. Empty if this handle
    /// was built without a known host (see [`RunningServer::new`]).
    pub fn host(&self) -> &str {
        &self.web_host
    }

    /// The advertised host active in the *owning* server process (from the lock),
    /// or `None` when this handle predates the field or was built socket-only. A
    /// controller must prefer this over its own `--entry` preference when building
    /// featured / QR URLs, so the address it prints matches the address the daemon
    /// actually serves under.
    pub fn advertised_host(&self) -> Option<&str> {
        self.web_advertised_host.as_deref()
    }

    /// Version of the discovered service, or an empty string for legacy locks.
    pub fn service_version(&self) -> &str {
        &self.service_version
    }

    /// The public entry URL prefix (`--entry`/`--qr`) the *owning* server was
    /// started with, or `None` when this handle predates the field or was built
    /// socket-only. A controller with no `--entry` of its own prefers this so
    /// featured / QR URLs point at the daemon's public address, not loopback.
    pub fn entry(&self) -> Option<&str> {
        self.web_entry.as_deref()
    }

    /// Best-effort liveness probe, mirroring
    /// [`ServerLock::is_alive`](crate::workspace::ServerLock::is_alive): `true`
    /// while this server still answers on its control socket or its web TCP port.
    ///
    /// A respawn uses this to wait for a shut-down daemon to fully release *both*
    /// before starting a replacement — `shutdown()` returns as soon as the daemon
    /// ACKs the request, long before it frees the port and removes its discovery
    /// lock, so the replacement would otherwise race the old process for the fixed
    /// port (`EADDRINUSE`) or latch the stale lock during readiness polling.
    pub fn is_reachable(&self) -> bool {
        // The control socket is the authoritative "server is up" signal; a
        // same-user connect proves liveness. A missing/stale socket refuses.
        if !self.socket.as_str().is_empty() && transport::probe(&self.socket) {
            return true;
        }
        // Fall back to the web TCP port so we only report "down" once the port is
        // actually free — the daemon removes its lock before the listener drops,
        // so the socket probe alone can't guarantee the port was released.
        if self.web_port == 0 {
            return false;
        }
        let connect_host = if crate::net::host_is_wildcard_v6(&self.web_host) {
            "::1"
        } else if crate::net::host_is_wildcard_v4(&self.web_host) || self.web_host.is_empty() {
            "127.0.0.1"
        } else {
            self.web_host.as_str()
        };
        let Ok(addr) = crate::net::bind_socket_addr(connect_host, self.web_port) else {
            return false;
        };
        std::net::TcpStream::connect_timeout(&addr, std::time::Duration::from_millis(500)).is_ok()
    }

    async fn call(&self, req: ControlRequest) -> Result<ControlResponse, ControlError> {
        match transport::request(&self.socket, &req).await? {
            ControlResponse::Err(msg) => Err(ControlError::Server(msg)),
            other => Ok(other),
        }
    }

    /// The running server's live workspace list.
    pub async fn list_workspaces(&self) -> Result<Vec<WorkspaceInfo>, ControlError> {
        match self.call(ControlRequest::ListWorkspaces).await? {
            ControlResponse::Workspaces(w) => Ok(w),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Register a workspace, returning its id. The collaborator hash is already
    /// salted by the caller (empty = inherit / no per-workspace code).
    pub async fn add_workspace(
        &self,
        path: &str,
        flags: WorkspaceFlags,
        collaborator_access_code_hash: &str,
    ) -> Result<String, ControlError> {
        match self
            .call(ControlRequest::AddWorkspace {
                path: path.to_string(),
                flags,
                collaborator_access_code_hash: collaborator_access_code_hash.to_string(),
                single_file: None,
                alias: String::new(),
            })
            .await?
        {
            ControlResponse::WorkspaceId(id) => Ok(id),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Register a temporary single-file (Open-With) workspace, returning its id.
    /// `path` is the file's parent directory and `single_file` the file name; the
    /// resulting workspace exposes only that file (plus locally referenced
    /// assets). Mirrors the in-process single-file add so both backends behave
    /// identically. The collaborator hash is already salted (empty = inherit).
    pub async fn add_single_file(
        &self,
        path: &str,
        single_file: &str,
        flags: WorkspaceFlags,
        collaborator_access_code_hash: &str,
    ) -> Result<String, ControlError> {
        match self
            .call(ControlRequest::AddWorkspace {
                path: path.to_string(),
                flags,
                collaborator_access_code_hash: collaborator_access_code_hash.to_string(),
                single_file: Some(single_file.to_string()),
                alias: String::new(),
            })
            .await?
        {
            ControlResponse::WorkspaceId(id) => Ok(id),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Register a user-opened path using the shared file-vs-directory scope.
    ///
    /// Native frontends resolve paths through [`WorkspaceOpenTarget`] and call
    /// this method instead of choosing between directory and single-file control
    /// requests themselves. That keeps CLI and GUI registration semantics tied
    /// to the same core implementation.
    pub async fn add_or_update_open_target(
        &self,
        target: &WorkspaceOpenTarget,
        flags: WorkspaceFlags,
        collaborator_access_code_hash: Option<&str>,
    ) -> Result<String, ControlError> {
        self.add_or_update_workspace_scoped(
            &target.root.to_string_lossy(),
            flags,
            target.single_file.as_deref(),
            collaborator_access_code_hash,
            None,
        )
        .await
    }

    /// Register a workspace, or — if `path` is already registered — update its
    /// flags (and access code, if supplied), returning the existing id. Mirrors
    /// the CLI's forward semantics so both front-ends behave identically.
    pub async fn add_or_update_workspace(
        &self,
        path: &str,
        flags: WorkspaceFlags,
        collaborator_access_code_hash: Option<&str>,
    ) -> Result<String, ControlError> {
        self.add_or_update_workspace_scoped(path, flags, None, collaborator_access_code_hash, None)
            .await
    }

    /// Register a directory or single-file workspace while preserving its full
    /// persisted identity and metadata. GUI attach/reconnect uses this to replay
    /// the settings snapshot without collapsing a single-file workspace into its
    /// parent directory or dropping its alias.
    pub async fn add_or_update_workspace_scoped(
        &self,
        path: &str,
        flags: WorkspaceFlags,
        single_file: Option<&str>,
        collaborator_access_code_hash: Option<&str>,
        alias: Option<&str>,
    ) -> Result<String, ControlError> {
        // The server stores canonical roots. Normalize before matching so macOS
        // `/var` -> `/private/var`, symlinks, and `..` cannot make an existing
        // identity look new and bypass the explicit metadata updates below.
        let canonical_path = expand_and_canonicalize(path)
            .map(|value| value.to_string_lossy().into_owned())
            .unwrap_or_else(|_| path.to_string());
        let existing = self
            .list_workspaces()
            .await?
            .into_iter()
            .find(|w| w.path == canonical_path && w.single_file.as_deref() == single_file);
        if let Some(existing) = existing {
            // Mirror the embedded registry's `add`, which refreshes the flags of
            // an already-registered identity: re-adding the same path applies the
            // supplied flags in both backends so they stay observably identical.
            self.update_flags(&existing.id, flags).await?;
            if let Some(hash) = collaborator_access_code_hash {
                self.set_access_code(&existing.id, Some(hash)).await?;
            }
            if let Some(alias) = alias {
                self.set_alias(&existing.id, alias).await?;
            }
            return Ok(existing.id);
        }
        match self
            .call(ControlRequest::AddWorkspace {
                path: canonical_path,
                flags,
                collaborator_access_code_hash: collaborator_access_code_hash
                    .unwrap_or("")
                    .to_string(),
                single_file: single_file.map(str::to_string),
                alias: alias.unwrap_or("").to_string(),
            })
            .await?
        {
            ControlResponse::WorkspaceId(id) => Ok(id),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Replace a workspace's feature flags wholesale.
    pub async fn update_flags(&self, id: &str, flags: WorkspaceFlags) -> Result<(), ControlError> {
        match self
            .call(ControlRequest::UpdateFlags {
                id: id.to_string(),
                flags,
            })
            .await?
        {
            ControlResponse::Ok => Ok(()),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Set (or clear, with an empty string) a workspace's display alias.
    pub async fn set_alias(&self, id: &str, alias: &str) -> Result<(), ControlError> {
        match self
            .call(ControlRequest::SetAlias {
                id: id.to_string(),
                alias: alias.to_string(),
            })
            .await?
        {
            ControlResponse::Ok => Ok(()),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Detach a workspace.
    pub async fn remove_workspace(&self, id: &str) -> Result<(), ControlError> {
        match self
            .call(ControlRequest::RemoveWorkspace { id: id.to_string() })
            .await?
        {
            ControlResponse::Ok => Ok(()),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Inspect persistent rows outside all currently registered workspaces.
    pub async fn data_cleanup_stats(&self) -> Result<DataCleanupStats, ControlError> {
        match self.call(ControlRequest::DataCleanupStats).await? {
            ControlResponse::DataCleanupStats(stats) => Ok(stats),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Permanently remove persistent rows outside all registered workspaces.
    pub async fn cleanup_orphaned_data(&self) -> Result<DataCleanupResult, ControlError> {
        match self.call(ControlRequest::CleanupOrphanedData).await? {
            ControlResponse::DataCleanupResult(result) => Ok(result),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Set (`Some(hash)`) or leave (`None`) a workspace's collaborator access
    /// code. The hash must already be salted with the shared per-install salt.
    pub async fn set_access_code(
        &self,
        id: &str,
        collaborator_access_code_hash: Option<&str>,
    ) -> Result<(), ControlError> {
        match self
            .call(ControlRequest::SetAccessCode {
                id: id.to_string(),
                collaborator_access_code_hash: collaborator_access_code_hash.map(str::to_string),
            })
            .await?
        {
            ControlResponse::Ok => Ok(()),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Mint a one-time administrator bootstrap URL that redirects to `redirect`.
    pub async fn admin_bootstrap(&self, redirect: &str) -> Result<String, ControlError> {
        match self
            .call(ControlRequest::AdminBootstrap {
                redirect: redirect.to_string(),
            })
            .await?
        {
            ControlResponse::Url(url) => Ok(url),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Mint a one-time administrator pairing code and its manual-entry URL.
    pub async fn admin_bootstrap_code(
        &self,
        redirect: &str,
    ) -> Result<(String, String), ControlError> {
        match self
            .call(ControlRequest::AdminBootstrapCode {
                redirect: redirect.to_string(),
            })
            .await?
        {
            ControlResponse::AdminCode { url, code } => Ok((url, code)),
            _ => Err(ControlError::Unexpected),
        }
    }

    /// Ask the running server to exit.
    pub async fn shutdown(&self) -> Result<(), ControlError> {
        match self.call(ControlRequest::Shutdown).await? {
            ControlResponse::Ok => Ok(()),
            _ => Err(ControlError::Unexpected),
        }
    }
}

#[cfg(test)]
mod tests;