Skip to main content

mkit_cli/commands/serve/
mod.rs

1//! `mkit serve <path>` — the `mkit+ssh://` forced-command server: the
2//! `mkit.rpc.v1.ssh` frame protocol (SPEC-TRANSPORT §4.2) on stdin/stdout
3//! against a local repository.
4//!
5//! The engine is `mkit-server`'s: [`mkit_server::ssh::serve_session`] over
6//! the pipeline, with the `.mkit` layout stores (`FsBlobStore` for
7//! `<root>/packs/`, `FsLayoutStore` for `<root>/refs/`), so the files are
8//! the ones `FileTransport`, `mkit+file://` remotes and any
9//! `mkit-server` fs-layout deployment read and write. It runs under a blocking executor
10//! (`futures::executor::block_on`): the CLI builds no async runtime, and
11//! the stores do blocking I/O inside their futures. A reader thread feeds
12//! stdin frames to the session and enforces the idle timeout
13//! (`--idle-timeout-secs`, `stdio.rs`).
14
15use std::io::{Read, Write};
16use std::path::{Path, PathBuf};
17use std::sync::Arc;
18use std::time::Duration;
19
20use clap::Parser;
21use mkit_core::repo_identity::RepositoryIdentity;
22use mkit_core::repo_lock::{self, LockError, RepoLock};
23use mkit_rpc::mkit::rpc::v1::ErrorCode;
24use mkit_server::fs::{FsBlobStore, FsLayoutStore};
25use mkit_server::pipeline::{AuthMode, Hooks, Pipeline, PipelineConfig};
26use mkit_server::policy::WritePolicy;
27use mkit_server::ssh::{SessionConfig, SessionEnd, WriteFrames, serve_session, upload_limits};
28use mkit_server::{
29    Addressing, NamespaceKey, NoopMetrics, Principal, RepoId, RepoName, SystemClock,
30};
31
32use crate::clap_shim;
33use crate::cli::CLI_VERSION;
34use crate::exit;
35
36mod stdio;
37
38use stdio::StdioFrameSource;
39
40/// The default of `--idle-timeout-secs` (planner decision Q12).
41const DEFAULT_IDLE_TIMEOUT_SECS: u64 = 60;
42
43/// The largest `--idle-timeout-secs` and `--max-session-secs`: 7 days.
44const MAX_TIMEOUT_SECS: u64 = 7 * 24 * 60 * 60;
45
46#[derive(Debug, Parser)]
47#[command(
48    name = "mkit serve",
49    about = "Speak the mkit-rpc SSH-frame protocol on stdin/stdout (the \
50             mkit+ssh:// forced-command server). It has no HTTP or \
51             mkit+enc:// listener."
52)]
53struct ServeOpts {
54    /// Path to the repository to serve. Under `--root` it is the
55    /// `<NAMESPACE>/<NAME>` the session binds; when it is omitted there,
56    /// `SSH_ORIGINAL_COMMAND` supplies it.
57    path: Option<String>,
58    /// Serve the repositories under DIR: the path names one
59    /// `<NAMESPACE>/<NAME>`, resolved to the directory
60    /// `<DIR>/<NAMESPACE>/<NAME>`, and only the namespace's owner may
61    /// write (SPEC-TRANSPORT §4 root mode).
62    #[arg(long, value_name = "DIR")]
63    root: Option<PathBuf>,
64    /// The Ed25519 public key the sshd forced-command configuration
65    /// asserts for this session: 64 lowercase hex characters (the raw
66    /// 32-byte key, no prefix). It is a trust assertion — sshd attaches it
67    /// to the key it already verified — and is never read from the
68    /// environment or `SSH_ORIGINAL_COMMAND`.
69    #[arg(long, value_name = "HEX", value_parser = parse_principal)]
70    principal: Option<[u8; 32]>,
71    /// End the session after this many seconds without a byte from the
72    /// client; 0 disables the timeout (at most 604800, 7 days). A slow
73    /// upload that keeps sending never trips it.
74    #[arg(
75        long,
76        value_name = "SECS",
77        default_value_t = DEFAULT_IDLE_TIMEOUT_SECS,
78        value_parser = clap::value_parser!(u64).range(..=MAX_TIMEOUT_SECS)
79    )]
80    idle_timeout_secs: u64,
81    /// End the process this many seconds after it starts, whatever the
82    /// client is doing; 0 (the default) disables the cap (at most 604800,
83    /// 7 days). It also bounds a client that trickles bytes or stops
84    /// reading, which the idle timeout does not.
85    #[arg(
86        long,
87        value_name = "SECS",
88        default_value_t = 0,
89        value_parser = clap::value_parser!(u64).range(..=MAX_TIMEOUT_SECS)
90    )]
91    max_session_secs: u64,
92}
93
94/// The listener flags `mkit serve` used to take, removed along with the HTTP
95/// (`--http`) and encrypted (`--listen-enc`) listeners. Clap rejects them as
96/// unknown arguments; [`run`] then adds a hint.
97const REMOVED_LISTENER_FLAGS: &[&str] = &[
98    "--http",
99    "--http-token",
100    "--unsafe-allow-any-http-peer",
101    "--listen-enc",
102    "--enc-authorized-peers",
103    "--enc-server-key",
104    "--unsafe-allow-any-enc-peer",
105    "--enc-idle-timeout-secs",
106    "--enc-handshake-timeout-secs",
107];
108
109/// The first removed listener flag in `args` (as `--flag` or
110/// `--flag=value`), stopping at a `--` separator.
111fn removed_listener_flag(args: &[String]) -> Option<&'static str> {
112    args.iter()
113        .take_while(|a| a.as_str() != "--")
114        .find_map(|a| {
115            let name = a.split_once('=').map_or(a.as_str(), |(n, _)| n);
116            REMOVED_LISTENER_FLAGS.iter().copied().find(|f| *f == name)
117        })
118}
119
120/// The repository identity `mkit serve` serves its root as: `mkit-server`'s
121/// default `--repository`, so both address one root's refs alike.
122const REPOSITORY: &str = "default";
123
124/// `--principal`'s value: exactly the raw 32-byte Ed25519 public key as
125/// 64 lowercase hex characters — no `ed25519-` prefix, no `0x`, nothing
126/// else. Clap maps a refusal to USAGE before any frame is read.
127fn parse_principal(text: &str) -> Result<[u8; 32], String> {
128    let bad =
129        || "expected 64 lowercase hex characters (a raw 32-byte Ed25519 public key)".to_owned();
130    let hex = |b: u8| b.is_ascii_digit() || (b'a'..=b'f').contains(&b);
131    if text.len() != 64 || !text.bytes().all(&hex) {
132        return Err(bad());
133    }
134    let mut out = [0u8; 32];
135    let nibble = |b: u8| {
136        if b.is_ascii_digit() {
137            b - b'0'
138        } else {
139            b - b'a' + 10
140        }
141    };
142    for (byte, pair) in out.iter_mut().zip(text.as_bytes().chunks_exact(2)) {
143        *byte = (nibble(pair[0]) << 4) | nibble(pair[1]);
144    }
145    Ok(out)
146}
147
148/// The repository path an `SSH_ORIGINAL_COMMAND` of exactly
149/// `mkit serve <path>` carries (SPEC-TRANSPORT §4's forced-command form):
150/// every byte is `[A-Za-z0-9._/-]` or a single space, the command splits
151/// on single spaces into exactly `["mkit", "serve", path]`, and `path`
152/// does not start with `-` (so no flag can be smuggled). Any other form —
153/// extra arguments, `sh -c`, quotes, tabs or doubled spaces, NUL, CR/LF,
154/// `;`, `$` — is refused.
155fn parse_original_command(command: &str) -> Option<&str> {
156    let ok_byte =
157        |b: u8| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b'/' | b'-' | b' ');
158    if !command.bytes().all(ok_byte) {
159        return None;
160    }
161    let mut parts = command.split(' ');
162    match (parts.next(), parts.next(), parts.next(), parts.next()) {
163        (Some("mkit"), Some("serve"), Some(path), None)
164            if !path.is_empty() && !path.starts_with('-') =>
165        {
166            Some(path)
167        }
168        _ => None,
169    }
170}
171
172/// The identity `path` names in root mode, validated and resolved under
173/// `root`: `<NAMESPACE>/<NAME>` in the §7.4 grammar (a bare name, an
174/// uppercase byte, `..` or another component is refused), then the
175/// canonical `<ROOT>/<NAMESPACE>/<NAME>` must be that path exactly — a
176/// surviving symlink component would serve one repository under
177/// another's identity, which the owner check would then attribute
178/// wrongly. `MKIT_SERVE_ROOT` pins the result like plain mode.
179fn resolve_root_repo(root: &Path, path: &str) -> Result<(PathBuf, RepoId), u8> {
180    let trimmed = path.trim_matches('/');
181    let identity = RepositoryIdentity::parse(trimmed).map_err(|_| exit::USAGE)?;
182    let Some(namespace) = identity.namespace() else {
183        return Err(exit::USAGE);
184    };
185    let root = std::fs::canonicalize(root).map_err(|_| exit::NOINPUT)?;
186    let expected = root.join(namespace.to_string()).join(identity.name());
187    let resolved = std::fs::canonicalize(&expected).map_err(|_| exit::NOINPUT)?;
188    if resolved != expected {
189        return Err(exit::NOPERM);
190    }
191    if !resolved.is_dir() || !resolved.join(".mkit").is_dir() {
192        return Err(exit::DATAERR);
193    }
194    if let Ok(pinned) = std::env::var("MKIT_SERVE_ROOT") {
195        let pinned = std::fs::canonicalize(&pinned).map_err(|_| exit::NOPERM)?;
196        if !resolved.starts_with(&pinned) {
197            return Err(exit::NOPERM);
198        }
199    }
200    Ok((
201        resolved,
202        RepoId {
203            namespace: NamespaceKey::from_namespace(namespace),
204            name: RepoName::new(identity.name()).map_err(|_| exit::USAGE)?,
205        },
206    ))
207}
208
209/// How old an upload's temp file must be before the startup sweep removes
210/// it. A live upload rewrites its temp file continuously; see
211/// [`sweep_crashed_uploads`].
212const STALE_UPLOAD_AGE: Duration = Duration::from_hours(1);
213
214#[must_use]
215pub fn run(args: &[String]) -> u8 {
216    let opts = match clap_shim::parse::<ServeOpts>("mkit serve", args) {
217        Ok(o) => o,
218        Err(code) => {
219            if let Some(flag) = removed_listener_flag(args) {
220                eprintln!(
221                    "hint: `mkit serve` no longer takes `{flag}`; it only speaks the ssh-frame \
222                     protocol on stdin/stdout.\n\
223                     \x20     Use SSH (`mkit serve`) or a Connect server such as vcs-worker.\n\
224                     \x20     The mkit+enc:// transport is deprecated: no maintained server."
225                );
226            }
227            return code;
228        }
229    };
230
231    let principal = Principal::SshForcedCommand {
232        key: opts.principal,
233    };
234    let target = match &opts.root {
235        // Plain mode, unchanged: one repository directory, open writes.
236        None => {
237            let Some(path) = &opts.path else {
238                eprintln!(
239                    "error: the following required arguments were not provided:\n  <PATH>\n\n\
240                     Usage: mkit serve <PATH>\n\n\
241                     For more information, try '--help'."
242                );
243                return exit::USAGE;
244            };
245            match resolve_repo_path(path) {
246                Ok(root) => ServeTarget {
247                    root,
248                    repo: repo_id(),
249                    write_policy: WritePolicy::Open,
250                    principal,
251                },
252                Err(code) => return code,
253            }
254        }
255        // Root mode: the path (or `SSH_ORIGINAL_COMMAND`) names the one
256        // `<NAMESPACE>/<NAME>` this process serves under `--root`.
257        Some(root) => {
258            let path = match &opts.path {
259                Some(path) => Some(path.clone()),
260                None => std::env::var("SSH_ORIGINAL_COMMAND")
261                    .ok()
262                    .and_then(|command| parse_original_command(&command).map(str::to_owned)),
263            };
264            let Some(path) = path else {
265                eprintln!(
266                    "mkit serve: --root serves <NAMESPACE>/<NAME>, from the path or \
267                     `SSH_ORIGINAL_COMMAND` `mkit serve <NAMESPACE>/<NAME>`"
268                );
269                return exit::USAGE;
270            };
271            match resolve_root_repo(root, &path) {
272                Ok((root, repo)) => ServeTarget {
273                    root,
274                    repo,
275                    write_policy: WritePolicy::Owner,
276                    principal,
277                },
278                Err(code) => return code,
279            }
280        }
281    };
282
283    // Held for the whole lifetime of this `serve` process
284    // (SPEC-CONCURRENCY §3.1, MKIT-11/#655): local
285    // worktree-mutating commands and `gc` probe this same lock (see
286    // `commands::warn_if_served`) to detect a live `serve` and warn.
287    // SHARED, not exclusive — SPEC-TRANSPORT documents multiple
288    // concurrent `serve` processes against one root (e.g. one per SSH
289    // forced-command connection) as a supported deployment, so `serve`
290    // instances must not exclude each other.
291    let _serve_guard = match lock_and_sweep(&target.root) {
292        Ok(g) => g,
293        Err(e) => {
294            eprintln!("mkit serve: serve lock: {e}");
295            return exit::TEMPFAIL;
296        }
297    };
298
299    if opts.max_session_secs > 0 {
300        spawn_session_cap(Duration::from_secs(opts.max_session_secs));
301    }
302
303    // Test-only fault injection for the mkit#703 SSH retry regression
304    // test (`tests/ssh_retry_e2e.rs`): end right after a successful
305    // `Hello`/`HelloResponse`, before answering any verb, so the process
306    // exits and the child pipe closes — simulating a mid-session
307    // connection drop that the client's `SshTransport` retry/reconnect
308    // path (SPEC-TRANSPORT §7) must recover from. Never set in
309    // production.
310    let stop_after_hello = std::env::var_os("MKIT_SERVE_TEST_DIE_AFTER_HELLO").is_some();
311    let idle = (opts.idle_timeout_secs > 0).then(|| Duration::from_secs(opts.idle_timeout_secs));
312    // `Stdin`/`Stdout`, not their locks: the reader thread needs `Send`,
313    // and each call locks. `WriteFrames` flushes after every frame.
314    serve_stdio(
315        &target,
316        std::io::stdin(),
317        std::io::stdout(),
318        idle,
319        stop_after_hello,
320    )
321}
322
323/// Resolve and validate the on-disk path supplied to `mkit serve`.
324pub(crate) fn resolve_repo_path(path: &str) -> Result<PathBuf, u8> {
325    let resolved = std::fs::canonicalize(path).map_err(|_| exit::NOINPUT)?;
326    if !resolved.is_dir() {
327        return Err(exit::DATAERR);
328    }
329    if !resolved.join(".mkit").is_dir() {
330        return Err(exit::DATAERR);
331    }
332    if let Ok(root) = std::env::var("MKIT_SERVE_ROOT") {
333        let pinned = std::fs::canonicalize(&root).map_err(|_| exit::NOPERM)?;
334        if !resolved.starts_with(&pinned) {
335            return Err(exit::NOPERM);
336        }
337    }
338    Ok(resolved)
339}
340
341/// Take the shared `serve.lock` for this process's lifetime, first
342/// sweeping crashed uploads when no other server holds it.
343///
344/// The sweep runs only while this process holds `serve.lock`
345/// EXCLUSIVELY, taken without waiting: every `mkit serve` and
346/// `mkit-server` holds it shared for its whole life, so while it is held
347/// exclusively no streaming upload of theirs can be in flight, and a busy
348/// lock (another server is up) skips the sweep. `FileTransport` writers
349/// (`mkit push` to a `mkit+file://` remote) take no such lock, but write
350/// each temp file in one go, so the age bound keeps theirs safe. The
351/// exclusive hold is then released and the shared one taken, as before.
352fn lock_and_sweep(repo_root: &Path) -> Result<RepoLock, LockError> {
353    let dot_mkit = repo_root.join(".mkit");
354    if let Ok(exclusive) =
355        repo_lock::acquire(&dot_mkit, crate::commands::SERVE_LOCK, Duration::ZERO)
356    {
357        sweep_crashed_uploads(repo_root);
358        drop(exclusive);
359    }
360    repo_lock::acquire_shared(
361        &dot_mkit,
362        crate::commands::SERVE_LOCK,
363        repo_lock::DEFAULT_TIMEOUT,
364    )
365}
366
367/// Remove the temp files (`packs/.<hex>.tmp.<pid>.<seq>`) that uploads
368/// left when their process crashed, once they are [`STALE_UPLOAD_AGE`]
369/// old. Best effort: a failure is ignored, and the next start retries.
370fn sweep_crashed_uploads(repo_root: &Path) {
371    let _ = FsBlobStore::new(repo_root).sweep_stale_uploads(STALE_UPLOAD_AGE);
372}
373
374/// `--max-session-secs`: a thread that ends the process once `max` has
375/// passed. It must work while the main thread is blocked in a write to a
376/// client that stopped reading, so it ends the process rather than
377/// signalling the session. That is safe at any point: an upload's temp
378/// file is never visible (and is swept later), every ref write is one
379/// atomic rename, and the kernel releases the ref and serve locks.
380fn spawn_session_cap(max: Duration) {
381    let spawned = std::thread::Builder::new()
382        .name("mkit-serve-session-cap".to_owned())
383        .spawn(move || {
384            std::thread::sleep(max);
385            eprintln!(
386                "mkit serve: session exceeded --max-session-secs {}; closing",
387                max.as_secs()
388            );
389            std::process::exit(i32::from(exit::PROTOCOL_ERROR));
390        });
391    if let Err(e) = spawned {
392        eprintln!("mkit serve: --max-session-secs is not enforced: {e}");
393    }
394}
395
396/// The repository `mkit serve` serves in plain mode (see [`REPOSITORY`]).
397fn repo_id() -> RepoId {
398    RepoId {
399        namespace: NamespaceKey::deployment_default(),
400        name: RepoName::new(REPOSITORY).unwrap_or_else(|_| unreachable!("a valid repo name")),
401    }
402}
403
404/// What `serve_stdio` serves and as whom.
405pub(crate) struct ServeTarget {
406    /// The repository's canonical on-disk root.
407    pub root: PathBuf,
408    /// Its storage identity: `default` in the reserved `root` namespace
409    /// in plain mode; the path's `<NAMESPACE>/<NAME>` in root mode.
410    pub repo: RepoId,
411    /// `Open` in plain mode, `Owner` in root mode.
412    pub write_policy: WritePolicy,
413    /// The session's principal; `key` only from `--principal`.
414    pub principal: Principal,
415}
416
417/// Serve one ssh-frame session over `input` and `output` against
418/// `target`, returning the exit code: [`exit::OK`] for a clean end or
419/// a failed write (the client went away), [`exit::PROTOCOL_ERROR`] for a
420/// protocol error or an idle timeout. `idle` bounds the time without a byte
421/// from the client (`None` disables it); a timeout is answered, best
422/// effort, with `Error{INVALID_REQUEST, "idle timeout"}`.
423///
424/// The ref store is opened with [`FsLayoutStore::open`], which refuses a
425/// root whose refs live in `mkit-server`'s `SQLite` database.
426pub(crate) fn serve_stdio<R, W>(
427    target: &ServeTarget,
428    input: R,
429    output: W,
430    idle: Option<Duration>,
431    stop_after_hello: bool,
432) -> u8
433where
434    R: Read + Send + 'static,
435    W: Write + Send,
436{
437    let meta = match FsLayoutStore::open(&target.root, &target.repo) {
438        Ok(meta) => meta,
439        Err(e) => {
440            eprintln!("mkit serve: {e}");
441            return exit::CONFIG_ERROR;
442        }
443    };
444    let mut cfg = PipelineConfig::new(
445        Addressing::Single {
446            repo: target.repo.clone(),
447        },
448        AuthMode::TransportIdentity,
449        upload_limits(),
450    );
451    cfg.write_policy = target.write_policy;
452    let pipeline = match Pipeline::new(
453        FsBlobStore::new(&target.root),
454        meta,
455        Hooks::new(),
456        cfg,
457        Arc::new(SystemClock),
458        Arc::new(NoopMetrics),
459    ) {
460        Ok(p) => p,
461        Err(e) => {
462            eprintln!("mkit serve: {}", e.public_message());
463            return exit::SOFTWARE;
464        }
465    };
466    let mut src = match StdioFrameSource::spawn(input, idle) {
467        Ok(src) => src,
468        Err(e) => {
469            eprintln!("mkit serve: stdin reader: {e}");
470            return exit::SOFTWARE;
471        }
472    };
473    let mut sink = WriteFrames(output);
474    let mut session = SessionConfig::new(format!("mkit serve/{CLI_VERSION}"));
475    session.stop_after_hello = stop_after_hello;
476    let end = futures::executor::block_on(serve_session(
477        &pipeline,
478        target.principal.clone(),
479        &mut src,
480        &mut sink,
481        &session,
482    ));
483    // The reader thread may still be blocked on stdin; the process exits
484    // under it.
485    match end {
486        SessionEnd::Clean | SessionEnd::IoError => exit::OK,
487        SessionEnd::ProtocolError => exit::PROTOCOL_ERROR,
488        SessionEnd::Timeout => {
489            let frame = mkit_rpc::ssh_error_frame(ErrorCode::InvalidRequest, "idle timeout");
490            let _ = mkit_rpc::write_frame(&mut sink.0, &frame);
491            let _ = sink.0.flush();
492            exit::PROTOCOL_ERROR
493        }
494    }
495}
496
497#[cfg(test)]
498mod tests;