1use 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
40const DEFAULT_IDLE_TIMEOUT_SECS: u64 = 60;
42
43const 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: Option<String>,
58 #[arg(long, value_name = "DIR")]
63 root: Option<PathBuf>,
64 #[arg(long, value_name = "HEX", value_parser = parse_principal)]
70 principal: Option<[u8; 32]>,
71 #[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 #[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
94const 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
109fn 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
120const REPOSITORY: &str = "default";
123
124fn 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
148fn 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
172fn 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
209const 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 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 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 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 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 serve_stdio(
315 &target,
316 std::io::stdin(),
317 std::io::stdout(),
318 idle,
319 stop_after_hello,
320 )
321}
322
323pub(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
341fn 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
367fn sweep_crashed_uploads(repo_root: &Path) {
371 let _ = FsBlobStore::new(repo_root).sweep_stale_uploads(STALE_UPLOAD_AGE);
372}
373
374fn 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
396fn 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
404pub(crate) struct ServeTarget {
406 pub root: PathBuf,
408 pub repo: RepoId,
411 pub write_policy: WritePolicy,
413 pub principal: Principal,
415}
416
417pub(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 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;