diff --git a/boot/init/src/boot.rs b/boot/init/src/boot.rs index 22eb62b0..2583208d 100644 --- a/boot/init/src/boot.rs +++ b/boot/init/src/boot.rs @@ -1,5 +1,4 @@ -//! Boot sequence: early mounts → resolve disks → overlay → network persist → -//! switch_root → exec. +//! Boot sequence: early mounts, resolve disks, overlay, persist network, switch_root, exec. use std::fs; use std::path::Path; @@ -12,8 +11,7 @@ const LAYER_DIR: &str = "/l"; const COW_DIR: &str = "/cow"; const NEWROOT: &str = "/newroot"; const POLL_INTERVAL: Duration = Duration::from_millis(2); -// NICs probe in single-digit ms; a missing one must degrade to the DHCP -// fallback, not stall rootfs handoff for the full disk budget (10s). +// a missing NIC must degrade to the DHCP fallback, not stall handoff for the disk budget. const NIC_TIMEOUT: Duration = Duration::from_millis(200); /// Cumulative µs checkpoints since sandbox-init start. @@ -44,8 +42,7 @@ impl Marks { pub fn run() -> ! { let mut marks = Marks::new(); - // Best-effort: if devtmpfs fails there is no console either; later - // failures then power off silently, which is still the right end state. + // best-effort: without devtmpfs there is no console, so a later failure powers off silently. let _ = sys::mount( "devtmpfs", "/dev", @@ -57,8 +54,7 @@ pub fn run() -> ! { let _ = sys::mount("proc", "/proc", Some("proc"), sys::MNT_SECURE, None); let _ = sys::mount("sysfs", "/sys", Some("sysfs"), sys::MNT_SECURE, None); - // Start marker: kernel-relative and visible at production loglevel, - // where the kernel's own boot lines are suppressed. One console write. + // start marker, visible at production loglevel where the kernel's own boot lines are suppressed. println!("sandbox-init: start at {}s", uptime()); let cmdline = fs::read_to_string("/proc/cmdline").unwrap_or_default(); @@ -71,13 +67,11 @@ pub fn run() -> ! { sys::fatal(&err, cfg.debug); } - // One deferred trace line (µs, cumulative since sandbox-init start): - // per-phase console writes would perturb exactly what they measure. + // one deferred line: per-phase console writes would perturb what they measure. if cfg.trace { println!("sandbox-init: trace{}", marks.render()); } - // Single marker line; boot-bench.sh keys on it. Uptime is kernel-relative, - // directly comparable with printk timestamps on the serial log. + // boot-bench.sh keys on this line; the uptime is comparable with printk timestamps. println!( "sandbox-init: rootfs ready at {}s, handing off to {}", uptime(), @@ -145,11 +139,7 @@ fn assemble(cfg: &BootCfg, marks: &mut Marks) -> Result<(), String> { Ok(()) } -/// Materializes kernel ip= params (cocoon CNI static flow) as MAC-matched -/// networkd units in the new root — persistence only, nothing is configured -/// in the initramfs; networkd applies them once the real init is up. A NIC -/// that never shows up degrades that interface to the DHCP fallback instead -/// of failing the boot, matching the old init-bottom hook. +/// Persists kernel ip= params as MAC-matched networkd units in the new root; a missing NIC degrades to the DHCP fallback. fn persist_network(cfg: &BootCfg) { if cfg.ips.is_empty() { return; @@ -175,8 +165,7 @@ fn persist_network(cfg: &BootCfg) { } } -/// Resolves every NIC's MAC in one sysfs sweep per poll iteration, against one -/// shared deadline — a missing NIC must cost the timeout once, not once each. +/// Resolves every NIC MAC against one shared deadline, so a missing NIC costs the timeout once. fn wait_nic_macs(devices: &[&str], timeout: Duration) -> Vec> { let deadline = Instant::now() + timeout; let mut found: Vec> = vec![None; devices.len()]; diff --git a/boot/init/src/cfg.rs b/boot/init/src/cfg.rs index aa24d1ef..23cfee09 100644 --- a/boot/init/src/cfg.rs +++ b/boot/init/src/cfg.rs @@ -1,13 +1,11 @@ -//! Kernel cmdline contract, shared with cocoon (hypervisor/utils.go builds -//! the cmdline; this module is the consuming end). +//! Kernel cmdline contract shared with cocoon, whose hypervisor/utils.go builds the cmdline. use std::fmt::Write as _; use std::time::Duration; pub const DEFAULT_INIT: &str = "/sbin/init"; const DEFAULT_TIMEOUT: Duration = Duration::from_secs(10); -// Cap: `Instant + Duration` panics on overflow, and a panic in PID 1 with -// panic=abort is a kernel panic — a hung VM instead of fatal()'s poweroff. +// cap the timeout: `Instant + Duration` panics on overflow, and a PID 1 panic is a kernel panic. const MAX_TIMEOUT_SECS: u64 = 86_400; #[derive(Debug, PartialEq, Eq)] @@ -19,8 +17,7 @@ pub struct BootCfg { /// Per-device wait budget. pub timeout: Duration, pub hostname: Option, - /// Static per-NIC config from kernel ip= params (cocoon CNI flow). - /// Persisted as systemd-networkd units, never applied in the initramfs. + /// Static per-NIC config, persisted as networkd units and never applied in the initramfs. pub ips: Vec, /// Handoff target inside the assembled rootfs. pub init: String, @@ -30,16 +27,14 @@ pub struct BootCfg { pub trace: bool, } -/// One `ip=::::::off[:dns0[:dns1]]` param -/// (the shape cocoon's BuildIPParams emits, one per NIC). +/// One `ip=::::::off[:dns0[:dns1]]` param, as cocoon emits it per NIC. #[derive(Debug, PartialEq, Eq)] pub struct IpParam { pub addr: String, pub prefix: u8, pub gateway: Option, pub dns: Vec, - /// Initramfs-time device name (ethN) — only used to look up the MAC; - /// the persisted unit matches by MAC so rootfs udev renames don't matter. + /// Initramfs-time device name, used only to look up the MAC the persisted unit matches on. pub device: String, } @@ -54,8 +49,7 @@ pub fn parse(cmdline: &str) -> Result { debug: false, trace: false, }; - for tok in cmdline.split_ascii_whitespace() { - let (key, val) = tok.split_once('=').unwrap_or((tok, "")); + for (key, val) in params(cmdline) { match key { "cocoon.layers" => { cfg.layers = val @@ -92,22 +86,14 @@ pub fn parse(cmdline: &str) -> Result { Ok(cfg) } -/// Debug check for the path where parse() itself failed and cfg.debug is -/// unavailable. Token handling mirrors parse() exactly: same value -/// predicate, last occurrence wins. +/// Debug check for the path where parse() itself failed and cfg.debug is unavailable. pub fn debug_requested(cmdline: &str) -> bool { - let mut debug = false; - for tok in cmdline.split_ascii_whitespace() { - let (key, val) = tok.split_once('=').unwrap_or((tok, "")); - if key == "sandbox.debug" { - debug = debug_token(val); - } - } - debug + params(cmdline) + .rfind(|&(key, _)| key == "sandbox.debug") + .is_some_and(|(_, val)| debug_token(val)) } -/// Overlay mount data. Layer mountpoints are index-based (/l/0, /l/1, …) so -/// arbitrary serial strings can never break lowerdir parsing (':' or ','). +/// Overlay mount data; index-based layer mountpoints keep a serial string out of lowerdir parsing. pub fn overlay_data(lower: &[String], cow_dir: &str) -> String { format!( "lowerdir={},upperdir={cow_dir}/upper,workdir={cow_dir}/work,index=on,redirect_dir=on,metacopy=on,xino=on", @@ -115,9 +101,7 @@ pub fn overlay_data(lower: &[String], cow_dir: &str) -> String { ) } -/// systemd-networkd unit for one static NIC. MAC-matched (device names may -/// change once rootfs udev renames), named 10-.network so it sorts -/// before the image's 20-wired.network DHCP fallback. +/// systemd-networkd unit for one static NIC, named to sort before the image's DHCP fallback. pub fn network_unit(ip: &IpParam, mac: &str) -> String { let mut unit = format!( "[Match]\nMACAddress={mac}\n\n[Network]\nAddress={}/{}\n", @@ -132,6 +116,13 @@ pub fn network_unit(ip: &IpParam, mac: &str) -> String { unit } +/// Kernel cmdline tokens as (key, value); a bare token carries an empty value. +fn params(cmdline: &str) -> impl DoubleEndedIterator { + cmdline + .split_ascii_whitespace() + .map(|tok| tok.split_once('=').unwrap_or((tok, ""))) +} + /// sandbox.debug value semantics, shared by parse() and debug_requested(). fn debug_token(val: &str) -> bool { val.is_empty() || val == "1" @@ -247,31 +238,17 @@ mod tests { assert!(debug_requested("sandbox.debug=0 x sandbox.debug")); } - #[test] - fn debug_requested_matches_parse() { - for tail in [ - "", - "sandbox.debug", - "sandbox.debug=", - "sandbox.debug=1", - "sandbox.debug=0", - "sandbox.debug=1 sandbox.debug=0", - ] { - let cmdline = format!("cocoon.layers=l0 cocoon.cow=cow {tail}"); - assert_eq!( - parse(&cmdline).unwrap().debug, - debug_requested(&cmdline), - "divergence for {tail:?}" - ); - } - } - #[test] fn parse_debug_forms() { let base = "cocoon.layers=l0 cocoon.cow=cow"; assert!(parse(&format!("{base} sandbox.debug")).unwrap().debug); assert!(parse(&format!("{base} sandbox.debug=1")).unwrap().debug); assert!(!parse(&format!("{base} sandbox.debug=0")).unwrap().debug); + assert!( + !parse(&format!("{base} sandbox.debug=1 sandbox.debug=0")) + .unwrap() + .debug + ); } #[test] diff --git a/boot/init/src/main.rs b/boot/init/src/main.rs index 012bd63b..2672d36f 100644 --- a/boot/init/src/main.rs +++ b/boot/init/src/main.rs @@ -1,8 +1,6 @@ -//! sandbox-init: the entire initramfs userland. Assembles the EROFS + overlay -//! rootfs described on the kernel cmdline and hands off to the real init. +//! sandbox-init: assembles the EROFS+overlay rootfs named on the kernel cmdline, then execs the real init. -// Off Linux only cfg's own tests use it; the bin compiles it dead so -// `cargo test` still covers the cmdline parsing on dev hosts. +// compiled dead off Linux so `cargo test` still covers cmdline parsing on a dev host. #[cfg_attr(not(target_os = "linux"), allow(dead_code))] mod cfg; diff --git a/boot/init/src/sys.rs b/boot/init/src/sys.rs index 76d15862..bf3848af 100644 --- a/boot/init/src/sys.rs +++ b/boot/init/src/sys.rs @@ -1,5 +1,4 @@ -//! Thin wrappers over the handful of syscalls the boot path needs; all the -//! unsafe lives here. +//! Thin wrappers over the syscalls the boot path needs; all the unsafe lives here. use std::ffi::CString; use std::io; @@ -56,12 +55,10 @@ pub fn sethostname(name: &str) -> Result<(), String> { Ok(()) } -/// Re-root into the assembled overlay. The old-root layer mounts stay pinned -/// by the overlay's references; the ~1MiB initramfs is deliberately not freed -/// (recursive delete would cost more than the memory is worth). +/// Re-roots into the assembled overlay; the ~1MiB initramfs is deliberately never freed. pub fn switch_root(newroot: &str) -> Result<(), String> { std::env::set_current_dir(newroot).map_err(|err| format!("chdir {newroot}: {err}"))?; - mount(".", "/", None, libc::MS_MOVE, None)?; + move_mount(".", "/")?; // SAFETY: the literal is a live 'static CStr for the duration of the call. if unsafe { libc::chroot(c".".as_ptr()) } != 0 { return Err(format!("chroot: {}", io::Error::last_os_error())); @@ -86,8 +83,7 @@ pub fn exec_init(path: &str) -> String { io::Error::last_os_error().to_string() } -/// Route PID 1 stdio to /dev/console. The kernel already did this when the -/// cpio carries the console node; this also covers cpios built without it. +/// Routes PID 1 stdio to /dev/console, covering a cpio built without the console node. pub fn claim_console() { // SAFETY: the literal is a live 'static CStr; the fd is validated (>=0) // before any dup2/close and only closed when it is not already a std stream. @@ -104,9 +100,7 @@ pub fn claim_console() { } } -/// Terminal state: report, then either hand the operator a shell (debug -/// initramfs) or power off so the orchestrator sees a dead VM immediately -/// instead of a hung boot. +/// Terminal state: hands the operator a debug shell, else powers off so the VM dies visibly. pub fn fatal(msg: &str, debug: bool) -> ! { eprintln!("sandbox-init: FATAL: {msg}"); if debug { diff --git a/e2e/e2e_test.go b/e2e/e2e_test.go index b4ae082a..fc49d0c2 100644 --- a/e2e/e2e_test.go +++ b/e2e/e2e_test.go @@ -379,8 +379,6 @@ func TestWritableVolumeEndToEnd(t *testing.T) { } } -// The read-only leg runs second: it is admitted only once the writer's -// release has cleared the dirty marker. func TestVolumeModeWireShape(t *testing.T) { scratch := writeVolumeImage(t, "scratch.img", "scratch-bytes") stack := startTenantStack(t, "node-token", nil, @@ -523,8 +521,6 @@ func TestAttachOnlyVolumeWireShape(t *testing.T) { } } -// TestDirtyVolumeRefusesReader: the marker a crashed writer leaves behind -// (pre-created here) turns read-only claims into 409s over the wire. func TestDirtyVolumeRefusesReader(t *testing.T) { scratch := writeVolumeImage(t, "scratch.img", "scratch-bytes") if err := os.WriteFile(scratch+".dirty", nil, 0o600); err != nil { @@ -601,8 +597,6 @@ func writeVolumeImage(t *testing.T, name, content string) string { return image } -// assertNoDirtyMarker fails if the image carries the write-ahead marker: an -// attach-only claim makes no consistency promise, so it must never write one. func assertNoDirtyMarker(t *testing.T, image, when string) { t.Helper() if _, err := os.Stat(image + ".dirty"); !errors.Is(err, os.ErrNotExist) { @@ -610,8 +604,6 @@ func assertNoDirtyMarker(t *testing.T, image, when string) { } } -// rawClaimResponse decodes the volume entries generically, so the assertion -// is the server's own JSON rather than the SDK's mirror of it. type rawClaimResponse struct { ID string `json:"id"` Token string `json:"token"` diff --git a/e2e/fakeengine_test.go b/e2e/fakeengine_test.go index 4ebb5cc8..f5f05914 100644 --- a/e2e/fakeengine_test.go +++ b/e2e/fakeengine_test.go @@ -16,8 +16,6 @@ import ( "github.com/cocoonstack/sandbox/sdk/go/silkd/silkdtest" ) -// fakeEngine replaces only the cocoon CLI: every "VM" is a silkdtest daemon -// behind a real hybrid-vsock UDS, so the data plane runs production code. type fakeEngine struct { real *engine.Engine dir string @@ -81,8 +79,6 @@ func (f *fakeEngine) SnapshotRemove(_ context.Context, _ string) error { return func (f *fakeEngine) SnapshotList(_ context.Context) ([]string, error) { return nil, nil } -// Hibernate closes the VM's silkd listener and Restore brings a fresh one up, -// mirroring the stop/resume the control plane observes. func (f *fakeEngine) Hibernate(ctx context.Context, name, _ string) error { return f.Remove(ctx, name) } @@ -143,8 +139,6 @@ func (f *fakeEngine) SyncGuest(_ context.Context, _ string) error { return nil } -// Removals join the trace only for VMs that carried a volume, so warm-pool -// churn cannot perturb the order. func (f *fakeEngine) volumeOpsLog() []string { f.mu.Lock() defer f.mu.Unlock() diff --git a/protocol/wire/frame.go b/protocol/wire/frame.go index 499809f2..a7e0cc9b 100644 --- a/protocol/wire/frame.go +++ b/protocol/wire/frame.go @@ -518,6 +518,16 @@ type InfoResp struct { func (InfoResp) RespType() string { return "info" } +// ProcInfo is one entry of Procs; ExitCode is absent while running. +type ProcInfo struct { + PID uint32 `json:"pid"` + Argv []string `json:"argv"` + Detached bool `json:"detached"` + State string `json:"state"` + ExitCode *int32 `json:"exit_code,omitempty"` + StartedAtEpochSecs uint64 `json:"started_at_epoch_secs"` +} + // Procs answers Ps. type Procs struct { Procs []ProcInfo `json:"procs"` @@ -532,6 +542,13 @@ type DataResp struct { func (DataResp) RespType() string { return "data" } +// DirEntry is one entry of Entries; Kind is one of the FileKind* consts. +type DirEntry struct { + Name string `json:"name"` + Kind string `json:"kind"` + Size uint64 `json:"size"` +} + // Entries answers FsList. type Entries struct { Entries []DirEntry `json:"entries"` @@ -539,6 +556,14 @@ type Entries struct { func (Entries) RespType() string { return "entries" } +// FileInfo is the Stat payload; Mode carries permission bits only. +type FileInfo struct { + Kind string `json:"kind"` + Size uint64 `json:"size"` + Mode uint32 `json:"mode"` + MtimeEpochSecs uint64 `json:"mtime_epoch_secs"` +} + // Stat answers FsStat. type Stat struct { Info FileInfo `json:"info"` @@ -625,31 +650,6 @@ type GitBranches struct { func (GitBranches) RespType() string { return "git_branches" } -// ProcInfo is one entry of Procs; ExitCode is absent while running. -type ProcInfo struct { - PID uint32 `json:"pid"` - Argv []string `json:"argv"` - Detached bool `json:"detached"` - State string `json:"state"` - ExitCode *int32 `json:"exit_code,omitempty"` - StartedAtEpochSecs uint64 `json:"started_at_epoch_secs"` -} - -// DirEntry is one entry of Entries; Kind is one of the FileKind* consts. -type DirEntry struct { - Name string `json:"name"` - Kind string `json:"kind"` - Size uint64 `json:"size"` -} - -// FileInfo is the Stat payload; Mode carries permission bits only. -type FileInfo struct { - Kind string `json:"kind"` - Size uint64 `json:"size"` - Mode uint32 `json:"mode"` - MtimeEpochSecs uint64 `json:"mtime_epoch_secs"` -} - // EncodeRequest renders {"v":1,"op":...,fields} without a trailing newline. func EncodeRequest(r Request) ([]byte, error) { return encodeTagged(requestHead+r.Op()+`"`, r) diff --git a/protocol/wire/frame_test.go b/protocol/wire/frame_test.go index 4efef372..19a7e47c 100644 --- a/protocol/wire/frame_test.go +++ b/protocol/wire/frame_test.go @@ -158,7 +158,6 @@ func TestEveryVerbHasAFixture(t *testing.T) { } } -// Twin of silkd's enum_value_sets_match_fixture; the lists are order-sensitive. func TestEnumValueSetsMatchFixture(t *testing.T) { raw, err := os.ReadFile(filepath.Join(fixtureDir, "enums.json")) if err != nil { @@ -244,8 +243,6 @@ func TestBulkDecodeStrictShape(t *testing.T) { } } -// TestTagAfterOtherKeys pins the tokenizer fallback: producers emit the tag -// first, but the protocol never promised order. func TestTagAfterOtherKeys(t *testing.T) { resp, err := DecodeResponse([]byte(`{"data":"aGk=","type":"stdout"}`)) if err != nil { diff --git a/sandboxd/ca.go b/sandboxd/ca.go index 840ee6cf..7547b32c 100644 --- a/sandboxd/ca.go +++ b/sandboxd/ca.go @@ -102,7 +102,6 @@ func writeCAFiles(dir, name string, certPEM, keyPEM []byte, force bool) error { return nil } -// writeKeyMaterial writes path with O_EXCL unless force. func writeKeyMaterial(path string, data []byte, perm os.FileMode, force bool) error { flags := os.O_WRONLY | os.O_CREATE | os.O_TRUNC if !force { diff --git a/sandboxd/main.go b/sandboxd/main.go index a866bc9b..44ddff2d 100644 --- a/sandboxd/main.go +++ b/sandboxd/main.go @@ -207,7 +207,6 @@ func startMesh(ctx context.Context, cfg *config.Config, mgr *pool.Manager) (*mes return msh, nil } -// gossipNodeState republishes this node's counts, templates, and volumes every tick. func gossipNodeState(ctx context.Context, msh *mesh.Mesh, mgr *pool.Manager) { t := time.NewTicker(gossipInterval) defer t.Stop() diff --git a/sandboxd/mesh/mesh.go b/sandboxd/mesh/mesh.go index d915d02c..31138dc3 100644 --- a/sandboxd/mesh/mesh.go +++ b/sandboxd/mesh/mesh.go @@ -39,6 +39,8 @@ type NodeState struct { Digest string `json:"digest,omitempty"` // cluster-invariant config digest } +type nodeMatch func(NodeState) bool + // Mesh is the node's view of the cluster and its own gossiped state. type Mesh struct { ml *memberlist.Memberlist @@ -223,7 +225,7 @@ func (m *Mesh) Shutdown() error { return m.ml.Shutdown() } -func (m *Mesh) warmCandidates(keyHash string, match func(NodeState) bool) []string { +func (m *Mesh) warmCandidates(keyHash string, match nodeMatch) []string { m.mu.Lock() type cand struct { addr string @@ -262,7 +264,7 @@ func (m *Mesh) persistEpoch(epoch uint64) error { return storeEpoch(m.epochPath, epoch) } -func (m *Mesh) owners(match func(NodeState) bool) []string { +func (m *Mesh) owners(match nodeMatch) []string { m.mu.Lock() var owners []string for id, st := range m.view { diff --git a/sandboxd/pool/archive.go b/sandboxd/pool/archive.go index 39a24bc5..4768d74e 100644 --- a/sandboxd/pool/archive.go +++ b/sandboxd/pool/archive.go @@ -348,12 +348,12 @@ func (m *Manager) retryArchiveDelete(ctx context.Context, ckID string) { l.Unlock() if err != nil { m.recDone(ckID) - log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "delete %s: %v", ckID, err) + log.WithFunc("pool.retryArchiveDelete").Warnf(ctx, "delete %s: %v", ckID, err) return } m.recDoneEvict(ckID) if err := m.clearArchiveCk(ckID); err != nil { - log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "clear %s: %v", ckID, err) + log.WithFunc("pool.retryArchiveDelete").Warnf(ctx, "clear %s: %v", ckID, err) } } diff --git a/sandboxd/store/peer/transport.go b/sandboxd/store/peer/transport.go index 9ec9c035..50dcb97d 100644 --- a/sandboxd/store/peer/transport.go +++ b/sandboxd/store/peer/transport.go @@ -10,6 +10,7 @@ import ( "errors" "fmt" "io" + "io/fs" "net/http" "net/url" "os" @@ -196,7 +197,7 @@ func safeJoin(root, name string) (string, error) { // tarInto emits only regular files and directories: a symlink would steer the reader outside dst. func tarInto(src, prefix string, tw *tar.Writer) error { - err := filepath.Walk(src, func(path string, fi os.FileInfo, err error) error { + err := filepath.WalkDir(src, func(path string, entry fs.DirEntry, err error) error { if err != nil { return err } @@ -207,33 +208,36 @@ func tarInto(src, prefix string, tw *tar.Writer) error { if rel == "." { return nil } + if !entry.IsDir() && !entry.Type().IsRegular() { + return nil + } + fi, err := entry.Info() + if err != nil { + return err + } name := prefix + rel - switch { - case fi.IsDir(): + if entry.IsDir() { return tw.WriteHeader(&tar.Header{ Name: name + "/", Mode: int64(fi.Mode().Perm()), Typeflag: tar.TypeDir, }) - case fi.Mode().IsRegular(): - if err := tw.WriteHeader(&tar.Header{ - Name: name, - Mode: int64(fi.Mode().Perm()), - Size: fi.Size(), - Typeflag: tar.TypeReg, - }); err != nil { - return err - } - f, err := os.Open(path) //nolint:gosec // path comes from Walk over src - if err != nil { - return err - } - defer func() { _ = f.Close() }() - _, err = io.Copy(tw, f) + } + if err = tw.WriteHeader(&tar.Header{ + Name: name, + Mode: int64(fi.Mode().Perm()), + Size: fi.Size(), + Typeflag: tar.TypeReg, + }); err != nil { return err - default: - return nil } + f, err := os.Open(path) //nolint:gosec // path comes from WalkDir over src + if err != nil { + return err + } + defer func() { _ = f.Close() }() + _, err = io.Copy(tw, f) + return err }) if err != nil { return fmt.Errorf("tar %s: %w", src, err) diff --git a/sandboxd/types/types.go b/sandboxd/types/types.go index 5d3fb180..74fc4747 100644 --- a/sandboxd/types/types.go +++ b/sandboxd/types/types.go @@ -196,6 +196,16 @@ type Checkpoint struct { Archive bool `json:"archive,omitempty"` } +// VMNetConfig is the per-NIC host tap the egress-lane nft lock binds. +type VMNetConfig struct { + TAP string `json:"tap"` +} + +// VMConfig is the config subset of VMRecord. +type VMConfig struct { + Name string `json:"name"` +} + // VMRecord is the subset of cocoon's VM records the control plane reads. type VMRecord struct { State string `json:"state"` @@ -213,16 +223,6 @@ func (r VMRecord) TapDevice() string { return r.NetworkConfigs[0].TAP } -// VMNetConfig is the per-NIC host tap the egress-lane nft lock binds. -type VMNetConfig struct { - TAP string `json:"tap"` -} - -// VMConfig is the config subset of VMRecord. -type VMConfig struct { - Name string `json:"name"` -} - // Volume is one dataset mount; Mode is normalized to "" (read-only) or VolumeModeRW. type Volume struct { Name string `json:"name"` diff --git a/sdk/go/client.go b/sdk/go/client.go index aef31e58..7ca8920c 100644 --- a/sdk/go/client.go +++ b/sdk/go/client.go @@ -306,10 +306,14 @@ func retryTransient(err error) bool { } } +type claimEncoder func(noRedirect, requirePromoted bool) ([]byte, error) + +type claimPoster func(addr string, body []byte) (claimResponse, error) + // claimFollow runs the claim protocol from origin: claim there, and on a // redirect re-encode with no_redirect and follow via redirectFallback. Only // the fallback error carries the verb — first-contact errors return raw. -func claimFollow(origin, verb string, encode func(noRedirect, requirePromoted bool) ([]byte, error), claimAt func(addr string, body []byte) (claimResponse, error)) (string, claimResponse, error) { +func claimFollow(origin, verb string, encode claimEncoder, claimAt claimPoster) (string, claimResponse, error) { body, err := encode(false, false) if err != nil { return "", claimResponse{}, err diff --git a/sdk/python/cocoonsandbox/checkpoint.py b/sdk/python/cocoonsandbox/checkpoint.py index 0a52b866..5dce1d2d 100644 --- a/sdk/python/cocoonsandbox/checkpoint.py +++ b/sdk/python/cocoonsandbox/checkpoint.py @@ -27,8 +27,7 @@ def new(self, ttl_seconds: int = 0) -> Sandbox: redirect to the node that actually holds it; if every candidate fails transiently, the claim falls back to the origin once so it heals (pulls the checkpoint) locally.""" - # Local import: a top-level one would close the client -> sandbox -> - # checkpoint cycle. + # local import: a top-level one closes the client -> sandbox -> checkpoint cycle. from .client import _redirect_fallback claim = {"ttl_seconds": ttl_seconds} if ttl_seconds else {} diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index 4f3d2b65..7ee377bd 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -49,8 +49,7 @@ def delete_template(self, template: str, net: str = "", size: str = "") -> None: candidates = (reply or {}).get("redirect") or [] if not candidates: return - # The retry carries no_redirect, mirroring the claim protocol: the - # owner answers for itself, never a second hop. + # the owner answers for itself under no_redirect, never a second hop. query["no_redirect"] = "1" path = "/v1/templates?" + urllib.parse.urlencode(query) _try_each(candidates, lambda peer: self._request(peer, "DELETE", path, None, "delete template")) @@ -60,8 +59,7 @@ def lookup(self, id: str, token: str) -> Sandbox: mesh peer concurrently, binding to whichever confirms ownership first — one dead peer must not cost its full timeout.""" def probe(addr: str) -> Sandbox: - # Bounded like _peers: a scatter loser must not hold a socket for - # the full client timeout after the winner has answered. + # bounded like _peers: a scatter loser must not hold a socket for the full timeout. reply = self._request(addr, "GET", f"/v1/sandboxes/{id}/owner", None, "owner", bearer=token, timeout=min(_PEERS_TIMEOUT, self.timeout)) return Sandbox(client=self, id=id, token=token, @@ -128,11 +126,7 @@ def post(peer): return self._handle_from(owner, reply) def _peers(self) -> list: - # /v1/peers is tenant-accessible (cluster topology); /v1/info is - # operator-only, so a tenant lookup cannot read peers from it. A - # single node has none — degrade to just the entry node. Bounded - # tighter than the client default (mirrors the Go SDK's peersTimeout) - # so one slow entry node cannot stall the scatter it feeds. + # /v1/peers is tenant-accessible; /v1/info is operator-only, so a tenant cannot read peers from it. try: return self._request(self.addr, "GET", "/v1/peers", None, "peers", timeout=min(_PEERS_TIMEOUT, self.timeout)).get("peers") or [] @@ -297,8 +291,7 @@ def attempt(addr): try: return attempt(origin) except APIError as origin_exc: - # Both halves matter to whoever reads this: the peers' failure says - # why the claim left the origin, the origin's why returning did not help. + # both halves matter: why the claim left the origin, and why returning did not help. combined = f"{origin_exc.message} (after redirect targets failed: {exc.message})" raise APIError(verb, origin_exc.status, combined) from origin_exc diff --git a/sdk/python/cocoonsandbox/conn.py b/sdk/python/cocoonsandbox/conn.py index 7daa9baf..c9fe2e67 100644 --- a/sdk/python/cocoonsandbox/conn.py +++ b/sdk/python/cocoonsandbox/conn.py @@ -84,7 +84,10 @@ def dial_agent(addr: str, sandbox_id: str, token: str, timeout: float) -> Conn: if any(c in value for c in "\r\n\0"): raise APIError("agent upgrade", 0, f"{name} contains a control character") host, port = addr.rsplit(":", 1) - sock = socket.create_connection((host, int(port)), timeout=timeout) + try: + sock = socket.create_connection((host, int(port)), timeout=timeout) + except OSError as exc: + raise ProtocolError(f"dial {addr}: {exc}") from exc # Nagle off: exec/write send small back-to-back frames before the first read. sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) reader = None diff --git a/sdk/python/cocoonsandbox/frames.py b/sdk/python/cocoonsandbox/frames.py index 2d52d48d..68760364 100644 --- a/sdk/python/cocoonsandbox/frames.py +++ b/sdk/python/cocoonsandbox/frames.py @@ -11,8 +11,7 @@ PROTO_VERSION = 1 MAX_FRAME = 8 * 1024 * 1024 FS_CHUNK = 256 * 1024 # silkd's per-frame chunk size, distinct from BULK_CHUNK below -# Bulk streams (push tars, port bytes) chunk larger — fewer frames for the -# same bytes, still far under MAX_FRAME after base64; mirrors the Go SDK. +# bulk streams chunk larger: fewer frames per byte, still under MAX_FRAME after base64. BULK_CHUNK = 1 << 20 @@ -32,10 +31,7 @@ def encode_request(op: str, **fields) -> bytes: def decode_response(line: bytes) -> dict: """Parses one response frame; the returned dict carries its tag under "type" and any binary payload decoded under "data".""" - # base64 is JSON-escape-free, so a frame shaped exactly {"type":..., - # "data":...} can be sliced directly, skipping json.loads; any other - # shape (extra fields, trailing bytes, non-alphabet bytes) falls through - # to the full parse instead of risking a silently wrong slice. + # base64 is JSON-escape-free, so an exactly-shaped data frame slices without json.loads. if line.startswith(b'{"type":"'): te = line.find(b'"', 9) if te > 0 and line[9:te] in (b"stdout", b"stderr", b"data") and line.startswith(b'","data":"', te): diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index 99f7a472..fe04a921 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -13,7 +13,7 @@ from .checkpoint import Checkpoint from .conn import Conn, _Closeable, dial_agent -from .errors import APIError, ExitError, ProtocolError, StreamTimeout +from .errors import APIError, ExitError, ProtocolError, SandboxError, StreamTimeout from .frames import BULK_CHUNK, FS_CHUNK from .template import Template @@ -66,8 +66,7 @@ def run(self, argv: list[str], cwd: str = "", env: dict | None = None, with self._dial() as conn: conn.send("exec", argv=argv, cwd=cwd or None, env=env, user=user or None, detach=False, session=session or None) - # The guest stops draining stdin while blocked writing stdout, so - # feeding it to completion before reading deadlocks. + # the guest stops draining stdin while blocked on stdout, so feeding it fully first deadlocks. pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) pump.start() code = _pump_stdio(conn, on_stdout, on_stderr) @@ -335,12 +334,12 @@ def _proxy_accept_loop(self, listener: socket.socket, port: int) -> None: def _proxy_conn(self, local: socket.socket, port: int) -> None: try: guest = self.dial_port(port) - except Exception: + except (SandboxError, OSError): local.close() return def pump_out(): - with contextlib.suppress(Exception): + with contextlib.suppress(SandboxError, OSError): while True: chunk = guest.recv() if not chunk: @@ -357,8 +356,7 @@ def pump_out(): if not chunk: break guest.send(chunk) - # Half-close, then let the guest finish answering: closing here - # would cut the reply the local client is still waiting for. + # half-close, not close: the guest's reply is still in flight. guest.close_write() pump.join() finally: @@ -401,9 +399,7 @@ def __init__(self, conn: Conn): self.error: Exception | None = None def __iter__(self) -> Iterator[dict]: - # Connection-bound: a close, drop, or undecodable frame ends iteration; - # a real server error frame (SilkdError) propagates. error tells a - # clean close (None) from a relay that dropped mid-stream. + # a transport failure ends iteration and sets error; a SilkdError frame propagates. while True: try: frame = self._conn.recv() @@ -503,7 +499,7 @@ def _send_chunks(conn: Conn, data: bytes, op: str = "data", chunk: int = FS_CHUN def _feed_stdin(conn: Conn, stdin: bytes) -> None: - with contextlib.suppress(Exception): # the reader reports the real failure + with contextlib.suppress(SandboxError, OSError): # the reader reports the real failure if stdin: _send_chunks(conn, stdin, op="stdin") conn.send("stdin_close") diff --git a/sdk/python/cocoonsandbox/template.py b/sdk/python/cocoonsandbox/template.py index ad9c9cd3..c9863555 100644 --- a/sdk/python/cocoonsandbox/template.py +++ b/sdk/python/cocoonsandbox/template.py @@ -26,8 +26,7 @@ def new(self, ttl_seconds: int = 0, volumes: list[str | Mapping[str, str]] | Non mount: bool = True) -> Sandbox: """Claims the template, following placement when volumes require it. mount=False attaches the volumes without mounting them.""" - # Local import: a top-level one would close the client → sandbox → - # template cycle. + # local import: a top-level one closes the client -> sandbox -> template cycle. from .client import _claim_body claim = _claim_body(self.name, self.net, self.size, ttl_seconds, volumes, mount) diff --git a/sdk/python/tests/test_hardening.py b/sdk/python/tests/test_hardening.py index 78311274..a9161831 100644 --- a/sdk/python/tests/test_hardening.py +++ b/sdk/python/tests/test_hardening.py @@ -17,6 +17,11 @@ def test_dial_agent_rejects_control_chars_in_identity(): dial_agent("127.0.0.1:1", "sb_1", "tok\r\nX-Evil: 1", 0.5) +def test_dial_agent_wraps_refused_connection(dead_addr): + with pytest.raises(ProtocolError): + dial_agent(dead_addr, "sb_1", "tok", 0.5) + + def test_watcher_propagates_silkd_error(): client_sock, guest_sock = socket.socketpair() guest_sock.sendall(json.dumps({"type": "error", "kind": "not_found", "message": "gone"}).encode() + b"\n") diff --git a/silkd/src/exec.rs b/silkd/src/exec.rs index d764249b..96163f0f 100644 --- a/silkd/src/exec.rs +++ b/silkd/src/exec.rs @@ -1,6 +1,4 @@ -//! exec handler: spawn a child, register it in the process table, stream -//! stdout/stderr as frames, and pump client stdin frames into it. Detached -//! execs return after `started` and keep running for later attach/logs. +//! exec handler: a detached exec returns after `started` and stays in the table for attach/logs. use std::process::Stdio; use std::sync::{Arc, OnceLock}; @@ -15,22 +13,14 @@ use crate::proc::{Chunk, Proc, Table, synth_pid}; use crate::proto::{ErrorKind, ExecReq, Request, Response}; use crate::sysutil; -/// After the child is reaped, wait at most this long for the pump to finish -/// before publishing Exit. It bounds a daemonizing child that leaves stdout -/// open in a surviving grandchild; a foreground client stalled past this -/// window can also lose the not-yet-drained pipe tail (bounded, rare). +/// Post-exit drain window; it bounds a daemonizing grandchild, and a stalled client loses the undrained tail. const POST_EXIT_DRAIN: Duration = Duration::from_secs(2); -/// How long an exited detached process stays in the table, so a late -/// `logs`/`attach` still finds its ring instead of a bare not_found. +/// How long an exited detached process stays in the table for a late `logs`/`attach`. const REAP_DELAY: Duration = Duration::from_secs(300); -/// Foreground output buffer before the child is backpressured. Bounds how far -/// ahead of a slow client the child may run, not a loss threshold. +/// Foreground output depth; while the child runs it backpressures here instead of dropping output. const FG_CAP: usize = 256; -/// Runs an exec request to completion (or to `started` when detached), -/// writing response frames to `out`. `client` yields further client frames -/// (stdin / stdin_close) for the foreground case. argv is non-empty — the -/// dispatcher validates it before the session/process split. +/// Runs an exec request, writing response frames to `out`; the dispatcher guarantees argv is non-empty. pub async fn run( table: &Table, now_secs: u64, @@ -70,9 +60,7 @@ where let pid = child.id().unwrap_or_else(synth_pid); let proc = table.register(pid, req.argv, req.detach, now_secs); if let Err(e) = crate::proto::write_frame(out, &Response::Started { pid }).await { - // The relay never learned this pid, so nothing will ever supervise or - // reap it: drop the table entry and kill the just-spawned child rather - // than leave a permanent Running ghost (and, when detached, an orphan). + // the client never learned this pid, so nothing else will ever reap the child. table.remove_if(pid, &proc); let _ = child.start_kill(); return Err(e); @@ -82,12 +70,7 @@ where let stdout = child.stdout.take(); let stderr = child.stderr.take(); - // Foreground consumes a backpressured mpsc: when the client falls behind, - // the sends block, the pipe fills, and the child slows down — nothing - // drops while the child lives (POST_EXIT_DRAIN bounds only the post-exit - // tail). Attachers still ride the best-effort broadcast (a secondary - // observer may drop under extreme lag). Detached has no client to pace - // it, so it gets no foreground sender. + // a backpressured mpsc paces the child to the foreground client; attachers ride the lossy broadcast. let (fg_tx, fg_rx) = mpsc::channel::(FG_CAP); let pump_fg = (!req.detach).then(|| fg_tx.clone()); let sup_fg = (!req.detach).then(|| fg_tx.clone()); @@ -97,12 +80,7 @@ where let pump_abort = pump.abort_handle(); tokio::spawn(pump_stdin(stdin, client, req.detach)); - // One supervisor per exec: reap the child, drain output within a grace - // window (so a daemonizer holding the pipe can't wedge us), then publish - // the terminal Exit — to the broadcast (attachers) and the foreground mpsc. - // The reaped code is recorded before the drain: a disconnect can abort the - // supervisor mid-drain, and its fallback must publish the real code, not - // fabricate -1 for a child that already exited cleanly. + // record the reaped code before the drain: an abort mid-drain must publish the real code, not -1. let reaped: Arc> = Arc::new(OnceLock::new()); let sup_reaped = Arc::clone(&reaped); let sup_proc = Arc::clone(&proc); @@ -135,10 +113,7 @@ where let mut rx = fg_rx; if let Err(e) = stream_to_client(&mut rx, out).await { - // Client gone: abort supervise (kills the child via kill_on_drop) and - // the pump (else its handle only detaches, leaking task/fds behind a - // silent pipe-holder), then publish the terminal state the abort may - // have pre-empted — else an attacher on this pid waits forever. + // client gone: publish the terminal state the abort pre-empts, else an attacher waits forever. supervise.abort(); pump_abort.abort(); let _ = supervise.await; @@ -174,11 +149,7 @@ where Ok(()) } -/// Reads stdout and stderr concurrently within this one task, so aborting the -/// pump (POST_EXIT_DRAIN timeout) cancels both. A nested spawn would survive -/// the abort and keep emitting a daemonizer's output after Exit. Each chunk -/// goes to the ring+broadcast (attachers/replay) and, for a foreground exec, -/// to the backpressured `fg` sender. +/// Reads both streams in this one task, so a nested spawn cannot outlive the pump abort. async fn pump_out( proc: Arc, stdout: Option, @@ -212,15 +183,11 @@ where Ok(n) => n, }; let Some(fg) = fg else { - // Detached: no foreground consumer, so the ring takes the slice - // and no owned chunk is built unless an attacher listens. proc.emit_bytes(stderr, &buf[..n]); continue; }; let chunk = make(buf[..n].to_vec()); proc.emit(&chunk); - // The foreground client backpressures here; emit above is - // best-effort fan-out to attachers. if fg.send(chunk).await.is_err() { break; } @@ -234,7 +201,7 @@ async fn pump_stdin( ) { let Some(mut sink) = stdin else { return }; if detach { - return; // detached execs take no client stdin + return; } while let Some(frame) = client.recv().await { match frame { diff --git a/silkd/src/find.rs b/silkd/src/find.rs index d8cc679b..f1c84808 100644 --- a/silkd/src/find.rs +++ b/silkd/src/find.rs @@ -1,6 +1,4 @@ -//! Search verbs: `fs.find` walks a tree and streams regex matches; `fs.replace` -//! rewrites matches in named files. Regex handling lives here so agents pass a -//! pattern as data rather than shell-quoting it through exec. +//! Search verbs: a pattern travels as data, never shell-quoted through exec. use std::io::Read; use std::path::{Path, PathBuf}; @@ -14,8 +12,7 @@ use tokio::sync::{Semaphore, SemaphorePermit, mpsc}; use crate::proto::{self, ErrorKind, Response, err_frame}; -/// The number of bytes a file may have before find skips it as binary/huge — -/// grep-scale line scanning is for source trees, not blobs. +/// Size above which find skips a file as binary or huge. const FIND_MAX_FILE: u64 = 8 * 1024 * 1024; /// Match frames in flight between the walking thread and the writer. @@ -67,9 +64,7 @@ struct Walk<'a> { } impl Walk<'_> { - /// Walks `root` depth-first, sending a `match` frame per matching line; false - /// means the receiver went away. A failure on the root propagates; deeper - /// ones skip only that directory. + /// Walks `root` depth-first per matching line; a failure below the root skips only that directory. fn run(&self, root: PathBuf) -> std::io::Result { let mut stack = vec![root]; let mut root = true; @@ -103,8 +98,7 @@ impl Walk<'_> { Ok(true) } - /// Scans one file, reporting whether the receiver is still listening. The size - /// bound comes off the open handle, so the check and the read see one file. + /// Scans one file; the size bound comes off the open handle, so check and read see one file. fn scan_file(&self, path: &Path) -> bool { let Ok(mut file) = std::fs::File::open(path) else { return true; @@ -134,9 +128,7 @@ impl Walk<'_> { } } -/// Streams `match` frames for every line under `path` matching `pattern`, -/// terminated by `done`. `glob` narrows the walk to file names matching it -/// (`*` and `?` wildcards); an invalid pattern is a bad-request error. +/// Streams `match` frames for every line under `path` matching `pattern`, terminated by `done`. pub async fn find( reader: &mut R, w: &mut W, @@ -151,13 +143,7 @@ where find_bounded(reader, w, path, pattern, glob, MATCH_QUEUE_BYTES).await } -/// Rewrites every `pattern` match to `replacement` in each of `files`, -/// streaming one `replaced` frame per file (with its match count) and a -/// terminal `done`. A file over the find size bound is skipped with a zero -/// count rather than read into memory. A read/write failure on one file ends -/// the stream with an error — files whose `replaced` frame already went out -/// are committed (each file is atomic; the list is not). An invalid pattern is -/// rejected before any file is touched. +/// Rewrites `pattern` to `replacement` in each of `files`, one `replaced` frame each; per-file atomic, not per-list. pub async fn replace( w: &mut W, files: Vec, @@ -297,9 +283,7 @@ fn name_matches(path: &Path, name_re: Option<&Regex>) -> bool { .is_some_and(|n| re.is_match(&n.to_string_lossy())) } -/// Compiles a `*`/`?` glob into an anchored regex over the whole file name; -/// every other character matches literally. escape() turns the wildcards -/// into exactly `\*`/`\?`, which the replaces then rewrite. +/// Compiles a `*`/`?` glob into an anchored regex over the whole file name. fn glob_regex(glob: &str) -> Result { let pat = regex::escape(glob).replace(r"\*", ".*").replace(r"\?", "."); Regex::new(&format!("^{pat}$")) diff --git a/silkd/src/forward.rs b/silkd/src/forward.rs index e9436ab4..52c3fc30 100644 --- a/silkd/src/forward.rs +++ b/silkd/src/forward.rs @@ -1,9 +1,4 @@ -//! `port_forward`: relay a guest TCP port over the RPC connection, so the -//! host reaches in-guest servers (dev servers, sshd) on either lane — the -//! vsock relay is the only transport the no-network lane has. After `ready`, -//! client `data` frames feed the socket (`data_end` half-closes it) and -//! socket bytes stream back as `data` frames; the guest server closing ends -//! the stream with `done`. +//! `port_forward` relays a guest TCP port: the no-network lane's only host reach. use std::io; @@ -14,11 +9,7 @@ use tokio::sync::mpsc; use crate::proto::{self, BULK_CHUNK, ErrorKind, Request, Response}; -/// Connects to 127.0.0.1:port and relays until either side finishes. The two -/// directions run on separate tasks: a blocking `write_all` into the guest -/// socket must never stall draining the guest's output (a single select loop -/// would head-of-line deadlock any bidirectional bulk transfer), mirroring -/// exec's split pump_stdin/pump_out. +/// Relays 127.0.0.1:port with one task per direction, against a head-of-line deadlock. pub async fn run( port: u16, client: mpsc::Receiver, @@ -46,10 +37,8 @@ pub async fn run( let res = loop { tokio::select! { r = tr.read(&mut buf) => match r { - // Guest closed its write side: the stream is done. Ok(0) => break proto::write_frame(w, &Response::Done).await, - // A read failure (RST, server crash) must not pass for a - // clean close — the client would mistake truncation for EOF. + // a read failure must not reach the client as a clean close. Err(e) => { break proto::error_frame( w, @@ -65,9 +54,7 @@ pub async fn run( } } }, - // A clean feeder end (data_end, client gone) keeps draining the - // socket; a protocol violation terminates the relay with its - // error instead of masquerading as a clean close. + // only a protocol violation from the feeder ends the relay. f = &mut feed, if !feed_done => { feed_done = true; if let Ok(Err(resp)) = f { @@ -82,11 +69,6 @@ pub async fn run( res } -/// Writes client `data` frames into the guest socket until `data_end` -/// (half-close) or the client disconnects; a stray frame is a protocol -/// violation, reported like the fs feeders do. Owns the write half so its -/// `write_all` can block on a slow guest without stalling the output -/// direction. async fn feed_socket( mut client: mpsc::Receiver, mut tw: OwnedWriteHalf, diff --git a/silkd/src/fs.rs b/silkd/src/fs.rs index a804cb46..511ef93e 100644 --- a/silkd/src/fs.rs +++ b/silkd/src/fs.rs @@ -1,10 +1,8 @@ -//! Filesystem verbs: the guest is the sandbox, so paths are taken as-is -//! against the guest root (the vsock boundary is the trust boundary). write -//! consumes a `data` frame stream from the client; read streams `data` frames -//! back; the rest are one-shot. +//! Filesystem verbs: the vsock boundary is the trust boundary, so paths are taken as-is. use std::io; use std::os::unix::fs::PermissionsExt; +use std::path::{Path, PathBuf}; use std::time::UNIX_EPOCH; use tokio::fs; @@ -13,15 +11,10 @@ use tokio::sync::mpsc; use crate::proto::{self, DirEntry, FileInfo, FileKind, Response, err_frame}; -/// Entries per `entries` frame. Worst-case entry (255-byte name, fully -/// JSON-escaped) is ~1.6KiB, so a full batch stays under MAX_FRAME. +/// Entries per `entries` frame; a worst-case 1.6KiB entry keeps a full batch under MAX_FRAME. pub const LIST_BATCH: usize = 4096; -/// Streams `data` frames from the client into `path`, applying `mode` if -/// given, until a `data_end` frame. Writes go to a sibling temp file that is -/// renamed over `path` only on clean completion, so a mid-stream failure or -/// early disconnect never leaves a truncated file at the destination and -/// `done` means "path is exactly what you sent". +/// Streams client `data` frames into `path` via a temp file renamed only on clean completion. pub async fn write( mut reader: R, w: &mut W, @@ -32,7 +25,8 @@ where R: AsyncBufRead + Unpin, W: AsyncWrite + Unpin, { - let tmp = tmp_name(&path); + let path = Path::new(&path); + let tmp = tmp_name(path); let mut file = match fs::File::create(&tmp).await { Ok(f) => f, Err(e) => return err_frame(w, &e, "create").await, @@ -48,7 +42,7 @@ where let _ = fs::remove_file(&tmp).await; return proto::write_feed_error(w, fail).await; } - if let Err(e) = commit_tmp(&tmp, &path, mode).await { + if let Err(e) = commit_tmp(&tmp, path, mode).await { return err_frame(w, &e, "commit").await; } proto::write_frame(w, &Response::Done).await @@ -66,8 +60,7 @@ pub async fn read(w: &mut W, path: String) -> io::Result< proto::write_frame(w, &Response::Done).await } -/// Lists a directory as a stream of `entries` frames terminated by `done`, -/// batched so an arbitrarily large directory can never exceed the frame cap. +/// Lists a directory as `entries` frames terminated by `done`, batched under the frame cap. pub async fn list(w: &mut W, path: String) -> io::Result<()> { // One blocking-pool dispatch for the whole directory, not one per entry. let (tx, mut rx) = mpsc::channel::>(2); @@ -90,16 +83,14 @@ pub async fn list(w: &mut W, path: String) -> io::Result< } } -/// Writes `bytes` to `path` via a sibling temp file committed into place, so -/// a crash never leaves a truncated file. -pub async fn write_atomic(path: &std::path::Path, bytes: &[u8]) -> io::Result<()> { - let path = path.display().to_string(); - let tmp = tmp_name(&path); +/// Writes `bytes` to `path` via a sibling temp file, so a crash never leaves a truncated file. +pub async fn write_atomic(path: &Path, bytes: &[u8]) -> io::Result<()> { + let tmp = tmp_name(path); if let Err(e) = fs::write(&tmp, bytes).await { let _ = fs::remove_file(&tmp).await; return Err(e); } - commit_tmp(&tmp, &path, None).await + commit_tmp(&tmp, path, None).await } /// Reports metadata for `path` (following symlinks). @@ -177,7 +168,10 @@ fn scan_dir(path: &str, tx: &mpsc::Sender>) -> io::Result<()> { Err(_) => (FileKind::Other, 0), }; entries.push(DirEntry { - name: ent.file_name().to_string_lossy().into_owned(), + name: ent + .file_name() + .into_string() + .unwrap_or_else(|os| os.to_string_lossy().into_owned()), kind, size, }); @@ -191,16 +185,16 @@ fn scan_dir(path: &str, tx: &mpsc::Sender>) -> io::Result<()> { Ok(()) } -fn tmp_name(path: &str) -> String { - format!("{path}.silkd-{}.tmp", crate::sysutil::tmp_suffix()) +fn tmp_name(path: &Path) -> PathBuf { + let mut tmp = path.as_os_str().to_os_string(); + tmp.push(".silkd-"); + tmp.push(crate::sysutil::tmp_suffix()); + tmp.push(".tmp"); + PathBuf::from(tmp) } -/// Commits a fully-written temp file over `path`. An explicit mode wins; -/// otherwise an overwrite inherits the destination's permission bits so -/// replacing an executable script doesn't silently strip its exec bit -/// (rename alone would leave the temp's create default). The temp is -/// removed on any failure. -async fn commit_tmp(tmp: &str, path: &str, mode: Option) -> io::Result<()> { +/// Commits a temp file over `path`; without an explicit mode an overwrite inherits the destination's bits. +async fn commit_tmp(tmp: &Path, path: &Path, mode: Option) -> io::Result<()> { let outcome = async { let effective = match mode { Some(m) => Some(m), diff --git a/silkd/src/git.rs b/silkd/src/git.rs index a0e32250..2c0364a3 100644 --- a/silkd/src/git.rs +++ b/silkd/src/git.rs @@ -1,9 +1,6 @@ -//! git verbs: wrap the guest git binary and return structured results -//! (branch, ahead/behind, per-file status, commit hash) rather than stdout to -//! scrape. clone/push/pull need the network, so on the none lane they fail -//! with a typed error pointing at fs.push; the rest are local and always work. -//! An auth token rides in an in-memory `http.extraHeader`, never the guest disk. +//! git verbs wrap the guest git binary into structured results; the none lane fails clone/push/pull with a typed error. +use std::borrow::Cow; use std::process::Stdio; use tokio::io::AsyncWrite; @@ -65,16 +62,14 @@ pub async fn add( terminal(w, "add", &out).await } -/// Commits staged changes with `message` and `author` ("Name "), -/// returning the new commit hash. +/// Commits staged changes and returns the new commit hash. pub async fn commit( w: &mut W, path: String, message: String, author: String, ) -> std::io::Result<()> { - // A fresh guest has no committer identity, so git commit would fail to - // auto-detect one; derive the committer from the author. + // a fresh guest has no committer identity, so derive one from the author. let (name, email) = split_author(&author); let mut cmd = git_cmd(&path, None); cmd.env("GIT_COMMITTER_NAME", name) @@ -165,11 +160,7 @@ async fn net_verb( terminal(w, verb, &out).await } -/// Builds a `git -C dir` command with config and stdio policy applied. Config -/// (an auth token, and quotePath=false so paths come back raw) rides in -/// `GIT_CONFIG_*` env vars, not `-c` args: the process environ is root-only, -/// whereas argv is world-readable via /proc//cmdline — a de-escalated -/// exec could otherwise scrape the token. +/// Builds a `git -C dir` command; config rides in `GIT_CONFIG_*` env vars because argv is world-readable via /proc. fn git_cmd(dir: &str, auth: Option<&str>) -> Command { let mut cmd = Command::new("git"); sysutil::align_proxy_env(&mut cmd); @@ -193,14 +184,18 @@ async fn git( /// Injects git config as GIT_CONFIG_COUNT/KEY_n/VALUE_n env pairs. fn apply_config(cmd: &mut Command, auth: Option<&str>) { - let mut pairs: Vec<(&str, String)> = vec![("core.quotePath", "false".to_string())]; + let mut pairs: Vec<(&str, Cow<'static, str>)> = + vec![("core.quotePath", Cow::Borrowed("false"))]; if let Some(token) = auth { - pairs.push(("http.extraHeader", format!("Authorization: Bearer {token}"))); + pairs.push(( + "http.extraHeader", + Cow::Owned(format!("Authorization: Bearer {token}")), + )); } cmd.env("GIT_CONFIG_COUNT", pairs.len().to_string()); for (i, (key, value)) in pairs.iter().enumerate() { cmd.env(format!("GIT_CONFIG_KEY_{i}"), key); - cmd.env(format!("GIT_CONFIG_VALUE_{i}"), value); + cmd.env(format!("GIT_CONFIG_VALUE_{i}"), value.as_ref()); } } @@ -262,13 +257,7 @@ fn parse_ahead_behind(rest: &str) -> (u32, u32) { (ahead, behind) } -/// Parses one porcelain-v2 entry. XY is field 2; the path is the last -/// space-field (kept intact even with spaces — core.quotePath=false keeps it -/// raw). Ordinary changes ("1") have the path at field 9; renames/copies ("2") -/// add an Xscore field, so the path is field 10 and carries "\t", -/// of which we keep the new path; unmerged ("u") entries carry four modes and -/// three hashes, putting the bare path at field 11. Untracked ("?") is a -/// bare path. +/// Parses one porcelain-v2 entry: XY is field 2, the path the last space-field at 9 ("1"), 10 ("2") or 11 ("u"). fn parse_file_line(line: &str) -> Option { let kind = line.split(' ').next()?; match kind { @@ -276,8 +265,8 @@ fn parse_file_line(line: &str) -> Option { let mut fields = line.split(' '); let xy = fields.nth(1)?; // field 2 let mut chars = xy.chars(); - let staged = chars.next()?.to_string(); - let unstaged = chars.next()?.to_string(); + let staged = chars.next()?; + let unstaged = chars.next()?; let skip = match kind { "2" => 7, "u" => 8, @@ -285,29 +274,30 @@ fn parse_file_line(line: &str) -> Option { }; // to reach the path field let path = fields.nth(skip)?; // Rejoin any spaces the split consumed, then drop a rename's \t. - let rest: Vec<&str> = fields.collect(); - let full = if rest.is_empty() { - path.to_string() - } else { - format!("{path} {}", rest.join(" ")) - }; + let mut full = String::from(path); + for tok in fields { + full.push(' '); + full.push_str(tok); + } + if let Some(tab) = full.find('\t') { + full.truncate(tab); + } Some(GitFileStatus { - path: full.split('\t').next()?.to_string(), + path: full, staged, unstaged, }) } "?" => Some(GitFileStatus { path: line.get(2..)?.to_string(), - staged: "?".to_string(), - unstaged: "?".to_string(), + staged: '?', + unstaged: '?', }), _ => None, } } -/// Splits an author "Name " into (name, email); a missing angle form -/// leaves the whole string as the name. +/// Splits an author "Name " into (name, email); a missing angle form is all name. fn split_author(author: &str) -> (&str, &str) { if let Some(open) = author.find('<') && let Some(close) = author[open..].find('>') @@ -344,10 +334,7 @@ mod tests { fn parses_ordinary_rename_untracked_and_conflict() { let ordinary = parse_file_line("1 .M N... 100644 100644 100644 h1 h2 src/main.rs").unwrap(); assert_eq!(ordinary.path, "src/main.rs"); - assert_eq!( - (ordinary.staged.as_str(), ordinary.unstaged.as_str()), - (".", "M") - ); + assert_eq!((ordinary.staged, ordinary.unstaged), ('.', 'M')); let rename = parse_file_line("2 R. N... 100644 100644 100644 h1 h2 R100 new.rs\told.rs").unwrap(); @@ -359,9 +346,6 @@ mod tests { let conflict = parse_file_line("u UU N... 100644 100644 100644 100644 h1 h2 h3 conflict.rs").unwrap(); assert_eq!(conflict.path, "conflict.rs"); - assert_eq!( - (conflict.staged.as_str(), conflict.unstaged.as_str()), - ("U", "U") - ); + assert_eq!((conflict.staged, conflict.unstaged), ('U', 'U')); } } diff --git a/silkd/src/lib.rs b/silkd/src/lib.rs index ffd4d64e..f3407f1b 100644 --- a/silkd/src/lib.rs +++ b/silkd/src/lib.rs @@ -1,6 +1,4 @@ -//! silkd library surface, exposed so integration tests (and, later, the -//! sandboxd relay's conformance checks) can drive the server in-process. -//! The binary in `main.rs` is a thin wrapper over `vsock::serve`. +//! silkd library surface, exposed so integration tests can drive the server in-process. pub mod exec; pub mod find; diff --git a/silkd/src/lsp.rs b/silkd/src/lsp.rs index 3b94b1f5..1c171ff7 100644 --- a/silkd/src/lsp.rs +++ b/silkd/src/lsp.rs @@ -1,21 +1,8 @@ -//! LSP broker: silkd spawns the language server a flavor image ships for a -//! language (its argv named in `/etc/silkd/lsp.d/`), keeps it -//! running addressed by a server id, and relays JSON-RPC bytes between the -//! client and the server's stdio. silkd is a broker, not a gateway: it never -//! parses LSP method semantics — the server multiplexes, silkd just pipes. -//! -//! `lsp_start` spawns and returns an id; `lsp_request` attaches the byte -//! stream and pumps until either side closes (`data_end` half-closes the -//! server's stdin, as in port_forward); `lsp_stop` kills it. v1 is -//! single-shot per server: the stream ending — clean or dropped — reaps the -//! server, since an LSP stream loses frame sync on any mid-request cut and a -//! resynced reattach is not worth the failure surface. A server no -//! `lsp_request` ever attaches is reaped on a TTL. The id and `lsp_stop` -//! still earn their place: `lsp_start` for several languages yields several -//! ids attached concurrently, and `lsp_stop` tears one down early. The base -//! image ships no manifests, so `lsp_start` for any language there answers a -//! typed `not_found` naming the flavor that provides one. +//! LSP broker: silkd spawns the language server named by +//! `/etc/silkd/lsp.d/` and pipes JSON-RPC bytes, never parsing LSP +//! semantics. A session is single-shot: the stream ending reaps the server. +use std::borrow::Cow; use std::collections::HashMap; use std::process::Stdio; use std::sync::atomic::{AtomicU64, Ordering}; @@ -31,8 +18,7 @@ use crate::sysutil; const MANIFEST_DIR: &str = "/etc/silkd/lsp.d"; -/// A started server nothing attaches is reaped after this: a client that dies -/// between the calls must not pin a 1-2 GB language server for the VM's life. +/// Reap TTL for an unattached server: a dead client must not pin a 1-2 GB server. const ATTACH_TTL: Duration = Duration::from_secs(300); static NEXT_ID: AtomicU64 = AtomicU64::new(1); @@ -48,8 +34,7 @@ impl Broker { Self::default() } - /// start spawns the manifested language server for `language`, rooted at - /// `root`, and returns its id. Absent manifest → typed not_found. + /// start spawns the manifested language server for `language` and returns its id. pub async fn start( &self, w: &mut W, @@ -95,7 +80,7 @@ impl Broker { let id = format!("lsp-{}", NEXT_ID.fetch_add(1, Ordering::Relaxed)); sysutil::lock(&self.inner).insert(id.clone(), server); - // If the client is already gone, don't leak the server we just spawned. + // if the client is already gone, do not leak the server just spawned. if let Err(e) = proto::write_frame( w, &Response::LspStarted { @@ -115,10 +100,7 @@ impl Broker { Ok(()) } - /// request attaches this connection to the server's stdio: client `data` - /// frames feed its stdin, its stdout streams back as `data` frames. When - /// either side closes the session is over — the server is reaped, so a - /// dropped connection never leaves a half-written stdin for a next caller. + /// request attaches this connection to the server's stdio and reaps the server on close. pub async fn request( &self, client: mpsc::Receiver, @@ -136,8 +118,7 @@ impl Broker { } Some(Some(stdio)) => stdio, }; - // The stdio halves are out of the table, so the attach TTL no longer - // covers this server: a client already gone must reap it here. + // the attach TTL no longer covers a server whose stdio halves are taken. if let Err(e) = proto::write_frame(w, &Response::Ready).await { self.reap(server_id).await; return Err(e); @@ -145,9 +126,7 @@ impl Broker { let feed = tokio::spawn(feed_stdin(client, stdin)); let res = pump_stdout(stdout, w).await; - // The session ends with this connection: stop feeding and reap the - // server. Aborting is safe here — we are killing the child, so a - // partially-written stdin frame goes nowhere. + // aborting mid-write is safe: the child is killed next. feed.abort(); let _ = self.reap(server_id).await; res @@ -166,8 +145,6 @@ impl Broker { } } - /// reap removes a server from the table and kills + waits its child (no - /// zombie survives), reporting whether it was there. async fn reap(&self, server_id: &str) -> bool { let removed = sysutil::lock(&self.inner).remove(server_id); match removed { @@ -179,8 +156,7 @@ impl Broker { } } - /// reap for a server no `lsp_request` ever attached; one that is attached - /// (its stdout taken) belongs to that connection and is left alone. + /// reap for a server no `lsp_request` ever attached. async fn reap_unattached(&self, server_id: &str) { let removed = { let mut table = sysutil::lock(&self.inner); @@ -195,8 +171,7 @@ impl Broker { } } -/// One running language server. `lsp_request` takes both stdio halves, so a -/// present stdout is also the flag that no connection has attached yet. +/// One running language server; a present `stdout` means no connection has attached. struct Server { child: Child, stdin: Option, @@ -210,10 +185,7 @@ async fn kill(mut server: Server) { async fn feed_stdin(mut client: mpsc::Receiver, mut stdin: ChildStdin) { while let Some(req) = client.recv().await { - // Anything but data half-closes by dropping stdin, aligned with - // port_forward: the server sees EOF, flushes, and exits (stdout then - // EOFs -> Done -> reap). A stray frame gets the same close, since the - // writer belongs to the stdout pump. + // a stray frame closes rather than errors: this task has no writer. let Request::Data { data } = req else { return }; if stdin.write_all(&data).await.is_err() || stdin.flush().await.is_err() { return; @@ -231,18 +203,16 @@ async fn pump_stdout( proto::write_frame(w, &Response::Done).await } -/// manifest_dir allows tests (and an operator) to relocate the manifest dir -/// via SILKD_LSP_DIR, mirroring silkd's other env overrides. -fn manifest_dir() -> String { - std::env::var("SILKD_LSP_DIR").unwrap_or_else(|_| MANIFEST_DIR.to_string()) +/// SILKD_LSP_DIR relocates the manifest dir for tests and operators. +fn manifest_dir() -> Cow<'static, str> { + std::env::var("SILKD_LSP_DIR").map_or(Cow::Borrowed(MANIFEST_DIR), Cow::Owned) } async fn read_manifest(language: &str) -> Option> { if language.is_empty() || language.contains(['/', '\\', '\0']) || language.contains("..") { return None; // never let a language name escape the manifest dir } - // O_NOFOLLOW: a symlink planted inside lsp.d must not read a target - // outside it. + // O_NOFOLLOW: a symlink in lsp.d must not read a target outside it. let mut f = tokio::fs::OpenOptions::new() .read(true) .custom_flags(libc::O_NOFOLLOW) diff --git a/silkd/src/main.rs b/silkd/src/main.rs index 125320d2..b970b192 100644 --- a/silkd/src/main.rs +++ b/silkd/src/main.rs @@ -1,7 +1,4 @@ -//! silkd: the in-guest sandbox daemon. Listens on a hybrid-vsock port for -//! newline-JSON RPC frames from the host (relayed by sandboxd) and runs -//! commands with context, tracks processes, moves files, and holds sessions. -//! `proto::Request` is the authoritative verb list. +//! silkd: the in-guest daemon serving newline-JSON RPC frames over hybrid-vsock; `proto::Request` is the verb list. use std::sync::Arc; @@ -23,8 +20,7 @@ async fn main() { session::IDLE_TTL, session::REAP_INTERVAL, )); - // Loopback→host egress relay: harmless when the host wired no proxy (the - // per-conn vsock dial is refused), so it needs no lane gate here. + // no lane gate needed: an unwired host refuses the per-conn vsock dial. tokio::spawn(net_egress::serve( net_egress::LOOPBACK_PORT, net_egress::HOST_VSOCK_PORT, diff --git a/silkd/src/net.rs b/silkd/src/net.rs index b035fac3..6d94b4f8 100644 --- a/silkd/src/net.rs +++ b/silkd/src/net.rs @@ -1,21 +1,11 @@ -//! Guest network-lane detection. The egress lane has a device-backed NIC; -//! the none lane carries only virtual interfaces (lo, plus the tunnels an -//! all-builtin kernel auto-creates). Network verbs (git clone/push/pull) and -//! proxy-env forwarding consult this. +//! Guest network-lane detection: the egress lane has a device-backed NIC, the none lane only virtual interfaces. use std::sync::LazyLock; use std::sync::atomic::{AtomicI8, Ordering}; static LANE_OVERRIDE: AtomicI8 = AtomicI8::new(-1); -/// Reports whether the guest can reach a network. `SILKD_NET` overrides the -/// probe (`none` / `egress`) for operators; otherwise an interface under -/// /sys/class/net backed by a real device (a `device` symlink) counts. -/// Name filtering is not enough: the all-builtin sandbox kernel auto-creates -/// virtual tunnels (sit0 and friends) even on the no-NIC lane, and `lo` is -/// virtual too — only a virtio/physical NIC has a device backing. On -/// non-Linux dev hosts (no /sys/class/net) it defaults to true so git verbs -/// are testable. +/// Reports whether the guest can reach a network; only a `device`-backed interface counts, since the kernel auto-creates virtual tunnels. pub fn has_egress() -> bool { match LANE_OVERRIDE.load(Ordering::Relaxed) { 0 => return false, diff --git a/silkd/src/net_egress.rs b/silkd/src/net_egress.rs index 0e2d102d..4c6ace43 100644 --- a/silkd/src/net_egress.rs +++ b/silkd/src/net_egress.rs @@ -1,10 +1,4 @@ -//! Guest→host egress relay. Binds a loopback proxy port; for each accepted TCP -//! connection it opens one guest→host vsock connection to sandboxd's egress -//! proxy and splices raw bytes. One vsock connection per proxied connection, -//! mirroring the host→guest `port_forward` model — the proxy speaks HTTP on top -//! of this byte pipe, so nothing here frames or inspects the payload. The none -//! lane's only route out; when the host has not wired a proxy, the per-conn -//! vsock dial is refused and the guest client sees a closed proxy (default-deny). +//! Guest→host egress relay, one vsock connection per proxied connection; an unwired host refuses the dial (default-deny). use std::io; diff --git a/silkd/src/proc.rs b/silkd/src/proc.rs index b6ca1388..9e31fe65 100644 --- a/silkd/src/proc.rs +++ b/silkd/src/proc.rs @@ -1,8 +1,6 @@ -//! Process table: every exec is registered so `ps`/`kill`/`attach`/`logs` -//! work against a guest pid regardless of which connection started it. -//! Detached processes keep a bounded output ring so a later `logs`/`attach` -//! can replay what already streamed. +//! Process table: every exec is registered so `ps`/`kill`/`attach`/`logs` work across connections. +use std::borrow::Cow; use std::collections::{HashMap, VecDeque}; use std::os::fd::{AsRawFd, OwnedFd}; use std::sync::atomic::{AtomicU64, Ordering}; @@ -17,8 +15,7 @@ use crate::sysutil; const LOG_RING_BYTES: usize = 256 * 1024; const OUTPUT_FANOUT: usize = 256; -/// Monotonic fallback pids for a spawned child with no reported OS pid; kept -/// below i32::MAX so a later signal cast to pid_t stays a valid pid. +/// Fallback pids for a child with no OS pid, masked below i32::MAX so a pid_t cast stays valid. static SYNTH: AtomicU64 = AtomicU64::new(1 << 30); /// A chunk of process output tagged by stream, shared with live attachers. @@ -30,7 +27,6 @@ pub enum Chunk { } impl Chunk { - /// The response frame that carries this chunk to a client. pub fn into_response(self) -> Response { match self { Chunk::Stdout(data) => Response::Stdout { data }, @@ -74,8 +70,7 @@ impl Table { sysutil::lock(&self.inner).get(&pid).cloned() } - /// Looks up a pid, writing a NotFound frame and returning None on a miss — - /// the shared prelude of kill/logs/attach/pty.resize. + /// Looks up a pid, writing a NotFound frame and returning None on a miss. pub async fn get_or_not_found( &self, w: &mut W, @@ -97,9 +92,7 @@ impl Table { .collect() } - /// Removes the entry only if it is still `proc`. OS pids are recycled, so - /// a deferred cleanup that removed by pid alone could evict a newer - /// process that reused the number. + /// Removes the entry only if it is still `proc`, because a recycled pid could name a newer process. pub fn remove_if(&self, pid: u32, proc: &Arc) { let mut map = sysutil::lock(&self.inner); if map.get(&pid).is_some_and(|cur| Arc::ptr_eq(cur, proc)) { @@ -116,10 +109,7 @@ impl Table { } } -/// One tracked process. `tx` fans out live output; `ring` retains a bounded -/// tail for replay; `state` flips to exited when the child is reaped. -/// `pty_master` holds a dup of a pty's master fd (None for a plain exec) so -/// `pty.resize` can ioctl it without racing the I/O task's own fd. +/// One tracked process; `pty_master` is a dup, so `resize` cannot race the I/O task's fd. pub struct Proc { pub pid: u32, pub argv: Vec, @@ -132,12 +122,7 @@ pub struct Proc { } impl Proc { - /// Records output for replay and fans it out to live attachers, holding - /// the ring lock across the broadcast send so an `attach` racing this - /// emit sees the chunk in exactly one of replay or the live stream. - /// The broadcast clone is skipped when nobody listens (the common case - /// on the output hot path); `attach_stream` subscribes under this same - /// ring lock, so the receiver count cannot change mid-emit. + /// Records output and fans it out; the ring lock spans the send, so a racing `attach` cannot double a chunk. pub fn emit(&self, chunk: &Chunk) { let mut ring = sysutil::lock(&self.ring); if let Chunk::Stdout(d) | Chunk::Stderr(d) = chunk { @@ -148,9 +133,7 @@ impl Proc { } } - /// `emit` for produced bytes: the ring takes the borrowed slice, and the - /// owned Chunk is built only when an attacher is actually listening — a - /// pty or detached exec pays no allocation on its output path otherwise. + /// `emit` for borrowed bytes; the owned Chunk is built only when an attacher listens. pub fn emit_bytes(&self, stderr: bool, data: &[u8]) { let mut ring = sysutil::lock(&self.ring); ring.push(stderr, data); @@ -185,8 +168,6 @@ impl Proc { } /// The exit code if the process has already exited; None while running. - /// Lets `logs`/`attach`/`kill` act on terminal state instead of waiting - /// on (or signalling against) a pid whose child is already reaped. pub fn exit_code(&self) -> Option { match *sysutil::lock(&self.state) { State::Running => None, @@ -199,9 +180,7 @@ impl Proc { sysutil::lock(&self.ring).drain_view() } - /// Atomically snapshots retained output and subscribes to live output - /// under one lock, so a chunk emitted concurrently lands in the replay or - /// the receiver but never both (no duplicate in the client's stream). + /// Snapshots retained output and subscribes to live output under one lock. pub fn attach_stream(&self) -> (Vec, broadcast::Receiver) { let mut ring = sysutil::lock(&self.ring); (ring.drain_view(), self.tx.subscribe()) @@ -209,8 +188,8 @@ impl Proc { fn info(&self) -> ProcInfo { let (state, exit_code) = match *sysutil::lock(&self.state) { - State::Running => ("running".to_string(), None), - State::Exited(c) => ("exited".to_string(), Some(c)), + State::Running => (Cow::Borrowed("running"), None), + State::Exited(c) => (Cow::Borrowed("exited"), Some(c)), }; ProcInfo { pid: self.pid, @@ -235,9 +214,7 @@ impl Ring { self.trim(); } - /// Drops whole oldest segments until the buffer fits the cap. VecDeque - /// front-drains are O(dropped), not O(remaining), so a chatty process - /// does not pay a full-buffer memmove per chunk. + /// Drops whole oldest segments until the buffer fits the cap; a front-drain is O(dropped). fn trim(&mut self) { let mut over = self.buf.len().saturating_sub(LOG_RING_BYTES); while over > 0 { @@ -249,8 +226,7 @@ impl Ring { } } - /// Retained output, coalescing adjacent same-stream segments: reads leave - /// tens of thousands of tiny ones and the client concatenates them anyway. + /// Retained output, with adjacent same-stream segments coalesced into one chunk. fn drain_view(&mut self) -> Vec { let bytes = self.buf.make_contiguous(); let mut out: Vec = Vec::new(); diff --git a/silkd/src/proto.rs b/silkd/src/proto.rs index d462235d..635bd137 100644 --- a/silkd/src/proto.rs +++ b/silkd/src/proto.rs @@ -4,6 +4,7 @@ //! (exit / done / error). Binary payloads ride as base64 (`data` fields), //! matching Go's default []byte JSON encoding for the SDK side. +use std::borrow::Cow; use std::collections::HashMap; use std::io; use std::sync::Arc; @@ -21,13 +22,10 @@ pub const PROTO_VERSION: u32 = 1; /// Chunk size for streaming a file back over `fs.read`. pub const READ_CHUNK: usize = 32 * 1024; -/// Bulk streams (fs read/pull, port bytes) chunk larger: fewer frames and -/// fewer flushes for the same bytes, still far under MAX_FRAME after base64. -/// Reads return what is available, so interactivity is unaffected. +/// Bulk streams chunk larger: fewer frames and flushes per byte, still under MAX_FRAME after base64. pub const BULK_CHUNK: usize = 256 * 1024; -/// Client → server frames. Unknown JSON fields (e.g. a future `v`) are -/// ignored by construction, which is the forward-compatibility story. +/// Client → server frames; unknown JSON fields are ignored, which is forward compatibility. #[derive(Debug, Deserialize, Serialize)] #[serde(tag = "op", rename_all = "snake_case")] pub enum Request { @@ -188,8 +186,7 @@ pub struct ExecReq { pub user: Option, #[serde(default)] pub detach: bool, - /// When set, run inside the named persistent shell session (cwd/env/state - /// persist across calls) instead of spawning a fresh process. + /// When set, run inside the named persistent shell session instead of a fresh process. #[serde(default)] pub session: Option, } @@ -225,15 +222,14 @@ pub enum Response { code: i32, }, Done, - /// Acknowledges an armed watch (events after it are guaranteed captured) - /// or a connected port_forward. + /// Acknowledges an armed watch (later events are guaranteed captured) or a connected port_forward. Ready, Error { kind: ErrorKind, message: String, }, Info { - version: String, + version: Cow<'static, str>, proto: u32, uptime_secs: u64, procs: usize, @@ -318,13 +314,12 @@ pub enum GitBranchOp { Checkout, } -/// One porcelain-v2 file entry: `staged`/`unstaged` are the XY status codes -/// (e.g. "M", "A", "?"), `path` the working-tree path. +/// One porcelain-v2 file entry; `staged`/`unstaged` are the XY status codes. #[derive(Debug, Deserialize, Serialize)] pub struct GitFileStatus { pub path: String, - pub staged: String, - pub unstaged: String, + pub staged: char, + pub unstaged: char, } #[derive(Clone, Copy, Debug, Deserialize, PartialEq, Serialize)] @@ -365,7 +360,7 @@ pub struct ProcInfo { pub pid: u32, pub argv: Vec, pub detached: bool, - pub state: String, + pub state: Cow<'static, str>, #[serde(default, skip_serializing_if = "Option::is_none")] pub exit_code: Option, pub started_at_epoch_secs: u64, @@ -379,17 +374,12 @@ pub enum FeedError { } /// Reads one newline-terminated frame (newline stripped); None on clean EOF. -/// One-shot callers (the leading request frame) use this; streaming loops use -/// `read_frame_into` so bulk uploads pay no per-frame growth reallocations. pub async fn read_frame(r: &mut R) -> io::Result>> { let mut line = Vec::new(); Ok(read_frame_into(r, &mut line).await?.then_some(line)) } -/// Reads one frame into `line` (cleared first, reused across calls — the -/// inbound twin of `write_chunk_frame`'s buffer reuse); false on clean EOF. -/// Scans the buffered reader chunk by chunk so a peer that never sends a -/// newline is cut off at MAX_FRAME instead of growing the line unbounded. +/// Reads one frame into `line` (cleared first, reused across calls), capped at MAX_FRAME; false on clean EOF. pub async fn read_frame_into( r: &mut R, line: &mut Vec, @@ -420,9 +410,7 @@ pub async fn write_frame(w: &mut W, resp: &Response) -> i w.flush().await } -/// Writes `frames` as one buffered batch, so a burst costs one write syscall -/// instead of one per frame; a batch past one frame cap flushes early to bound -/// the staging buffer. `buf` is reused across calls. +/// Writes `frames` as one buffered batch, so a burst costs one write syscall; `buf` is reused. pub async fn write_frames( w: &mut W, buf: &mut Vec, @@ -440,9 +428,7 @@ pub async fn write_frames( w.flush().await } -/// Renders `{"type":KIND,"data":""}` into `buf` (reused across calls) -/// and writes it — the bulk path skips serde's owned Vec + String per chunk. -/// base64's alphabet needs no JSON escaping. +/// Hand-renders `{"type":KIND,"data":""}` into `buf`; base64 needs no JSON escaping. pub async fn write_chunk_frame( w: &mut W, buf: &mut Vec, @@ -455,8 +441,7 @@ pub async fn write_chunk_frame( let b64_len = base64::encoded_len(data.len(), true) .ok_or_else(|| io::Error::other("chunk too large to encode"))?; let total = PRE.len() + kind.len() + MID.len() + b64_len + END.len(); - // Grow-only, never cleared: resize's zero-fill runs once at high-water - // instead of per chunk; buf[..total] is fully overwritten below. + // grow-only: buf[..total] is fully overwritten, so the zero-fill runs once at high-water. if buf.len() < total { buf.resize(total, 0); } @@ -474,8 +459,7 @@ pub async fn write_chunk_frame( w.flush().await } -/// Maps an io error to an Error frame, classifying NotFound so a client can -/// tell a missing path from a real failure. Shared by the fs and tree verbs. +/// Maps an io error to an Error frame, classifying NotFound apart from a real failure. pub async fn err_frame( w: &mut W, e: &io::Error, @@ -489,8 +473,7 @@ pub async fn err_frame( write_frame(w, &Response::error(kind, format!("{op}: {e}"))).await } -/// Writes the terminal frame of a finished subprocess verb: Done on success, -/// else `label: trimmed stderr` as an Internal error. Shared by git and tree. +/// Writes a finished subprocess verb's terminal frame: Done, else `label: stderr` as Internal. pub async fn subprocess_result( w: &mut W, success: bool, @@ -513,10 +496,7 @@ pub async fn error_frame( write_frame(w, &Response::error(kind, message)).await } -/// Streams `reader` back as `data` frames until EOF — the outbound twin of -/// `feed_data_frames`, shared by fs.read and fs.pull. A frame-write error -/// propagates as the outer error; a source read error comes back as the inner -/// one so the caller can clean up (reap a tar child) before mapping it. +/// Streams `reader` back as `data` frames until EOF; a source read error returns as the inner error so the caller can reap. pub async fn stream_data_frames( reader: &mut R, w: &mut W, @@ -536,9 +516,7 @@ where } } -/// Feeds client `data` frames into `sink` until `data_end`. Shared by fs.write -/// (sink = temp file) and fs.push (sink = tar stdin); the caller flushes/closes -/// the sink and, on error, reports via `write_feed_error`. +/// Feeds client `data` frames into `sink` until `data_end`; the caller flushes and closes `sink`. pub async fn feed_data_frames(reader: &mut R, sink: &mut S) -> Result<(), FeedError> where R: AsyncBufRead + Unpin, @@ -594,8 +572,7 @@ fn cap_check(len: usize) -> io::Result<()> { Ok(()) } -/// Renders one newline-terminated frame onto `buf`, holding the write side to -/// the same cap the read side enforces — every client rejects a larger frame. +/// Renders one newline-terminated frame onto `buf`, capped because every client rejects a larger one. fn render_frame(buf: &mut Vec, resp: &Response) -> io::Result<()> { let start = buf.len(); serde_json::to_writer(&mut *buf, resp).map_err(io::Error::other)?; @@ -613,9 +590,7 @@ mod b64 { s.serialize_str(&STANDARD.encode(data)) } - /// Decodes inside the visitor so bulk `data` frames (43KB-1.3MB on the - /// wire) never allocate an intermediate String — serde hands the borrowed - /// slice straight to the base64 decoder. + /// Decodes inside the visitor so a bulk `data` frame never allocates an intermediate String. pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result, D::Error> { struct B64Visitor; impl serde::de::Visitor<'_> for B64Visitor { diff --git a/silkd/src/pty.rs b/silkd/src/pty.rs index 37d0f10d..f2c406b6 100644 --- a/silkd/src/pty.rs +++ b/silkd/src/pty.rs @@ -1,7 +1,4 @@ -//! `pty.open` / `pty.resize`: run a shell under a pseudo-terminal. A PTY is a -//! process, so it registers in the proc table and `ps`/`kill`/`attach`/`logs` -//! apply unchanged — output rides the same `stdout` chunk stream, input arrives -//! as `stdin` frames, and the child's exit is the terminal frame. +//! `pty.open` / `pty.resize`: a pty is a process, so it registers in the proc table and `ps`/`kill`/`attach`/`logs` apply unchanged. use std::os::fd::{AsRawFd, OwnedFd}; use std::process::Stdio; @@ -18,15 +15,12 @@ use crate::proc::{Chunk, Proc, Table, synth_pid}; use crate::proto::{ErrorKind, PtyReq, READ_CHUNK, Request, Response}; use crate::sysutil; -/// After the shell exits, drain the master's buffered tail for at most this -/// long so the last screenful isn't lost, without wedging on a stuck fd. +/// Bounds the post-exit drain of the master's buffered tail so a stuck fd cannot wedge teardown. const POST_EXIT_DRAIN: Duration = Duration::from_secs(1); type Master = Arc>; -/// Opens a pty running the guest shell ($SHELL, else bash), streams its output -/// as `stdout` frames and the exit as `exit`, and consumes client `stdin` -/// frames as pty input until the shell exits or the client disconnects. +/// Opens a pty on the guest shell, streams its output as `stdout` frames, and feeds client `stdin` frames to it. pub async fn open( table: &Table, now: u64, @@ -76,9 +70,7 @@ pub async fn open( Ok(c) => c, Err(e) => return crate::proto::err_frame(out, &e, "spawn shell").await, }; - // Drop cmd so its Stdio dups of the slave close: otherwise silkd keeps the - // slave open and the master never sees the shell's exit (no EOF/EIO), so - // the pump would hang forever. + // drop cmd so its slave dups close, else the master never sees the shell's exit and the pump hangs. drop(cmd); let pid = child.id().unwrap_or_else(synth_pid); @@ -97,7 +89,6 @@ pub async fn open( .await .is_err() { - // Client vanished before we streamed anything — tear down cleanly. let _ = child.start_kill(); finish(&proc, -1); table.remove_if(pid, &proc); @@ -127,9 +118,7 @@ pub async fn resize( } } -/// Streams master output to the client and client stdin to the master until -/// the shell exits or the client disconnects, returning the exit code (-1 on -/// a kill or a client-gone teardown). The exit frame is written by the caller. +/// Pumps master output to the client and client stdin to the master; the caller writes the exit frame. async fn pump( master: &Master, proc: &Arc, @@ -137,8 +126,7 @@ async fn pump( out: &mut W, child: &mut tokio::process::Child, ) -> i32 { - // Stdin rides its own task so a slow write can't stall output or reaping; - // a client disconnect (recv None) signals teardown via disc. + // stdin rides its own task so a slow write cannot stall output or reaping. let (disc_tx, mut disc_rx) = oneshot::channel::<()>(); tokio::spawn(pump_stdin(Arc::clone(master), client, disc_tx)); @@ -153,8 +141,7 @@ async fn pump( drain(master, proc, out, &mut buf, &mut frame).await; return code; } - // Client closed the terminal: kill the shell and let child.wait - // publish its real (signalled) code on the next iteration. + // kill the shell so child.wait publishes its real signalled code. _ = &mut disc_rx, if !disc => { disc = true; let _ = child.start_kill(); @@ -165,13 +152,11 @@ async fn pump( Err(_) => { eof = true; continue } }; match guard.try_io(|fd| sysutil::read_fd(fd.get_ref().as_raw_fd(), &mut buf)) { - // A read of 0 (BSD) or any error (EIO on Linux) means the - // slave is fully closed. + // a read of 0 (BSD) or any error (EIO on Linux) means the slave is fully closed. Ok(Ok(0)) | Ok(Err(_)) => eof = true, Ok(Ok(n)) => { proc.emit_bytes(false, &buf[..n]); if crate::proto::write_chunk_frame(out, &mut frame, "stdout", &buf[..n]).await.is_err() { - // Client gone mid-output: kill and reap for the real code. let _ = child.start_kill(); return sysutil::wait_code(child).await; } @@ -183,9 +168,7 @@ async fn pump( } } -/// Writes client stdin frames into the master until the client disconnects, -/// then signals teardown. StdinClose does not close the master — a pty stays -/// alive after EOF on its input. +/// Writes client stdin frames into the master; StdinClose is ignored because a pty outlives input EOF. async fn pump_stdin( master: Master, mut client: mpsc::Receiver, @@ -207,8 +190,7 @@ async fn pump_stdin( } } -/// Reads whatever the master buffered after the shell exits, bounded so a -/// wedged fd can't hang teardown. +/// Reads the master's buffered tail after the shell exits, bounded by POST_EXIT_DRAIN. async fn drain( master: &Master, proc: &Arc, @@ -239,8 +221,7 @@ async fn drain( .await; } -/// finish publishes the terminal state once, so attachers on this pid always -/// see an Exit (and never hang waiting for one). +/// Publishes the terminal state once so attachers always see an Exit. fn finish(proc: &Arc, code: i32) { if proc.exit_code().is_none() { proc.mark_exited(code); diff --git a/silkd/src/server.rs b/silkd/src/server.rs index 563160b0..54b04b22 100644 --- a/silkd/src/server.rs +++ b/silkd/src/server.rs @@ -1,6 +1,6 @@ -//! Connection handling: read the leading request frame, dispatch to a -//! handler, and for exec forward subsequent client frames over a channel. +//! Connection handling: read the leading request frame, dispatch to a handler, and forward later client frames over a channel. +use std::borrow::Cow; use std::time::{Instant, SystemTime, UNIX_EPOCH}; use tokio::io::{AsyncBufRead, AsyncWrite}; @@ -55,9 +55,7 @@ impl State { Request::Logs { pid } => self.logs(&mut writer, pid).await, Request::Attach { pid } => self.attach(&mut writer, pid).await, Request::Exec(e) => { - // Validated here, ahead of the session/process split, so an - // empty argv answers the same error on both paths (the session - // shell would otherwise run it as a successful no-op). + // validated ahead of the session/process split so an empty argv errors on both paths. if e.argv.is_empty() { return proto::error_frame( &mut writer, @@ -199,7 +197,7 @@ impl State { proto::write_frame( w, &Response::Info { - version: env!("CARGO_PKG_VERSION").to_string(), + version: Cow::Borrowed(env!("CARGO_PKG_VERSION")), proto: proto::PROTO_VERSION, uptime_secs: self.started.elapsed().as_secs(), procs: self.table.len(), @@ -224,9 +222,7 @@ impl State { let Some(proc) = self.table.get_or_not_found(w, pid).await? else { return Ok(()); }; - // An exited process's pid may already be recycled by the guest OS — - // signalling it would hit an unrelated process, so treat kill as a - // no-op success once the child is reaped. + // an exited pid may already be recycled, so signalling it would hit an unrelated process. if proc.exit_code().is_none() { crate::sysutil::signal_pid(proc.pid, signal.unwrap_or(libc::SIGKILL)); } @@ -251,9 +247,7 @@ impl State { let Some(proc) = self.table.get_or_not_found(w, pid).await? else { return Ok(()); }; - // Snapshot replay and subscribe atomically so a chunk emitted right - // now lands in one or the other, never both; if already exited, report - // the code instead of waiting on an Exit that fired before we listened. + // snapshot replay and subscribe atomically so a chunk lands in exactly one of them. let (replay, mut rx) = proc.attach_stream(); let mut frame = Vec::new(); for chunk in replay { @@ -283,8 +277,7 @@ impl Default for State { } } -/// Aborts the spawned feeder on drop, so no streaming arm can leak the task -/// (and the reader half it owns) on any return path. +/// Aborts the spawned feeder on drop so no streaming arm leaks the task and its reader half. struct Feeder(tokio::task::JoinHandle<()>); impl Drop for Feeder { @@ -301,8 +294,7 @@ where (rx, Feeder(tokio::spawn(feed_client(reader, tx)))) } -/// Stdout/Stderr ride the reused-buffer bulk path (serde's per-chunk Vec + -/// base64 String otherwise dominate replay/attach); Exit stays a one-off frame. +/// Stdout/Stderr ride the reused-buffer bulk path; serde's per-chunk allocations dominate replay otherwise. async fn write_chunk( w: &mut W, buf: &mut Vec, @@ -315,8 +307,7 @@ async fn write_chunk( } } -/// Forwards post-request client frames (stdin/stdin_close) to the exec -/// handler until the connection half-closes. +/// Forwards post-request client frames to the handler until the connection half-closes. async fn feed_client(mut reader: R, tx: mpsc::Sender) where R: AsyncBufRead + Unpin, diff --git a/silkd/src/session.rs b/silkd/src/session.rs index 6466eced..86ecff4a 100644 --- a/silkd/src/session.rs +++ b/silkd/src/session.rs @@ -1,11 +1,7 @@ -//! Persistent shell sessions: each session owns a long-lived bash whose cwd, -//! environment, and shell state survive across `exec {session}` calls. -//! A command is injected into the shell and -//! delimited by a unique sentinel that also carries its exit code; stderr is -//! merged into stdout (`exec 2>&1`) so one stream frames cleanly without a -//! pipe deadlock, which is the conventional interactive-shell behaviour. +//! Persistent shell sessions: each owns a long-lived bash whose cwd, env, and shell state survive across `exec {session}` calls. use std::collections::HashMap; +use std::fmt::Write as _; use std::io; use std::process::Stdio; use std::sync::{Arc, Mutex}; @@ -17,14 +13,11 @@ use tokio::process::{Child, ChildStdin, ChildStdout, Command}; use crate::proto::{self, ErrorKind, READ_CHUNK, Response}; use crate::sysutil; -/// Idle sessions (no command run within this window) are reaped so an -/// abandoned session's shell + fds don't accumulate. +/// Idle sessions are reaped so an abandoned session's shell and fds do not accumulate. pub const IDLE_TTL: Duration = Duration::from_secs(30 * 60); pub const REAP_INTERVAL: Duration = Duration::from_secs(60); -/// Bounds the post-marker scan for the exit-code line's newline: the line is -/// tiny and atomic with the marker, so this only caps a backgrounded child -/// flooding stdout without a newline. +/// Caps the post-marker scan for the exit-code newline against a backgrounded child flooding stdout. const EXIT_TAIL_MAX: usize = 64 * 1024; /// Registry of live shell sessions, addressed by id across connections. @@ -38,8 +31,7 @@ impl Table { Self::default() } - /// Spawns a bash session, applies cwd/env, and registers it. A blank id - /// yields a generated one; a duplicate id is refused. + /// Spawns a bash session, applies cwd/env, and registers it; a blank id is generated, a duplicate refused. pub async fn create( &self, id: Option, @@ -51,17 +43,14 @@ impl Table { _ => format!("sess-{}", sysutil::tmp_suffix()), }; let mut child = Command::new("bash") - // The same sanitized baseline as exec/pty ("nothing inherited - // from silkd"); the request's cwd/env layer on via the init line. + // same sanitized baseline as exec/pty; the request's cwd/env layer on via the init line. .env_clear() .envs(sysutil::base_env()) .stdin(Stdio::piped()) .stdout(Stdio::piped()) - // Null, not piped: nothing ever reads a stderr pipe, and the init - // line `exec 2>&1` repoints fd 2 at the stdout pipe anyway. + // null, not piped: the init line `exec 2>&1` repoints fd 2 at the stdout pipe. .stderr(Stdio::null()) - // Own process group (leader pgid == pid) so teardown can group-kill - // the shell together with whatever external command it is running. + // own process group so teardown can group-kill the shell with its running command. .process_group(0) .kill_on_drop(true) .spawn()?; @@ -74,6 +63,7 @@ impl Table { stdin, stdout, _child: child, + cmd_buf: String::new(), buf: vec![0u8; READ_CHUNK], acc: Vec::new(), frame: Vec::new(), @@ -81,8 +71,7 @@ impl Table { last_active: Mutex::new(Instant::now()), }); - // Merge stderr, then apply cwd/env; sync on a sentinel so the first - // real command reads a clean stream. + // sync on a sentinel so the first real command reads a clean stream. let mut init = String::from("exec 2>&1\n"); if let Some(dir) = cwd { init.push_str("cd "); @@ -101,9 +90,7 @@ impl Table { io.converse::(&init, None).await?; } - // Reserve the id atomically: a concurrent create with the same id - // must not silently replace (and orphan) the loser. Dropping the loser - // Arc here fires kill_on_drop on its shell. + // reserve the id atomically: a concurrent create must not orphan the loser's shell. let mut map = sysutil::lock(&self.inner); if map.contains_key(&id) { return Err(io::Error::new( @@ -125,9 +112,7 @@ impl Table { ids } - /// Removes and kills a session's shell. The explicit kill (not just - /// kill_on_drop) unwedges a session whose command is blocked while a - /// run() still holds the io lock and its Arc. + /// Removes and kills a session's shell; the explicit kill unwedges a blocked command holding the io lock. pub fn remove(&self, id: &str) -> bool { let removed = sysutil::lock(&self.inner).remove(id); if let Some(s) = &removed { @@ -144,9 +129,7 @@ impl Table { self.len() == 0 } - /// Removes and kills sessions idle longer than `ttl`. A session currently - /// running a command (io lock held) is never idle, so it is skipped even - /// if its last stamp is old (a long-running command). Returns the count. + /// Removes and kills sessions idle longer than `ttl`; a session holding the io lock is never idle. pub fn reap_idle(&self, ttl: Duration) -> usize { let mut map = sysutil::lock(&self.inner); let before = map.len(); @@ -161,9 +144,7 @@ impl Table { } } -/// One persistent shell. `io` serializes commands so a session runs one at a -/// time; the child is killed on drop, and `pid` allows an explicit kill to -/// unwedge a session whose command is blocked while the io lock is held. +/// One persistent shell; `io` serializes commands so a session runs one at a time. pub struct Session { pid: u32, io: tokio::sync::Mutex, @@ -171,11 +152,7 @@ pub struct Session { } impl Session { - /// Runs `argv` in the session shell, streaming stdout frames to `w` and - /// ending with an exit frame. argv is shell-quoted and joined so a fresh - /// `cd`/`export` persists to later calls while ordinary argv still runs. - /// Returns whether the session shell is still alive; a dead shell is - /// removed by the caller so its id stops resolving. + /// Runs `argv` in the session shell, shell-quoted so `cd`/`export` persist; false once the shell is dead. pub async fn run(&self, argv: &[String], w: &mut W) -> bool { let mut cmdline = String::new(); for (i, a) in argv.iter().enumerate() { @@ -186,10 +163,7 @@ impl Session { } let mut io = self.io.lock().await; let outcome = io.converse(&cmdline, Some(w)).await; - // Stamp while io is still held: the reaper skips a session whose io lock - // is held, so it can't mis-reap in the gap between the command finishing - // and the fresh stamp landing (a long-running command isn't idle either - // — io stays held for its whole duration). + // stamp while io is held so the reaper cannot mis-reap in the gap before the stamp lands. *sysutil::lock(&self.last_active) = Instant::now(); drop(io); match outcome { @@ -205,42 +179,37 @@ impl Session { } } -/// The shell's pipes plus converse's scratch buffers, which are per-session -/// and reused across commands rather than reallocated per call. +/// The shell's pipes plus converse's scratch buffers, reused across commands. struct Io { stdin: ChildStdin, stdout: ChildStdout, _child: Child, + cmd_buf: String, buf: Vec, acc: Vec, frame: Vec, } impl Io { - /// Writes `cmd` to the shell followed by a unique sentinel printf, then - /// reads output up to the sentinel — emitting the preceding bytes as - /// stdout frames on `out` (None discards, used for init) and returning the - /// command's exit code. A tiny tail is held back so a marker split across - /// reads is still found. + /// Runs `cmd` up to a unique sentinel, streaming the bytes before it to `out` and returning the exit code. async fn converse( &mut self, cmd: &str, mut out: Option<&mut W>, ) -> io::Result { - // The marker is unforgeable: the shell runs untrusted code that must - // not be able to fake the sentinel and desync the stream. `{ …; } - // ":", c => c, }; + // the marker must be unforgeable: the shell runs untrusted code that could fake a sentinel. let marker = format!("__SILK_{}__", sysutil::rand_token()); - self.stdin - .write_all( - format!("{{ {body} ; }} (out: &mut Option<&mut W>, frame: &mut Vec, data: &[u8]) { if data.is_empty() { return; @@ -319,8 +283,7 @@ fn take_pipe(pipe: Option, name: &str) -> io::Result { pipe.ok_or_else(|| io::Error::other(format!("child {name} not piped"))) } -/// POSIX single-quote quoting: wraps `s` in '…' onto the end of `out`, -/// closing/reopening around any embedded single quote. +/// POSIX single-quote quoting: wraps `s` in '…' onto `out`, reopening around embedded quotes. fn shell_quote_into(out: &mut String, s: &str) { out.push('\''); for c in s.chars() { diff --git a/silkd/src/sysutil.rs b/silkd/src/sysutil.rs index 4e5e8c35..71749f33 100644 --- a/silkd/src/sysutil.rs +++ b/silkd/src/sysutil.rs @@ -1,7 +1,4 @@ -//! Small OS helpers: the base environment for spawned commands, best-effort -//! user de-escalation, and the one signal syscall — the crate's unsafe work -//! lives here (pty.rs holds the one other unsafe block, its pre_exec -//! registration). +//! Small OS helpers; the crate's unsafe work lives here, except pty.rs's pre_exec registration. use std::io::Read; use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd}; @@ -11,14 +8,12 @@ use std::sync::{LazyLock, Mutex}; use tokio::process::Command; -/// Serializes getpwnam, whose result points into a process-global static -/// buffer the next call clobbers — not safe across tokio's worker threads. +/// Serializes getpwnam, whose result points into a static buffer the next call clobbers. static NSS_LOCK: Mutex<()> = Mutex::new(()); static TMP_SEQ: AtomicU64 = AtomicU64::new(0); -/// Forwarded from silkd's own environment into every exec: the loopback proxy -/// relay is the no-NIC lane's only way out, and nothing else may leak. +/// Forwarded into every exec: the loopback proxy relay is the no-NIC lane's only way out. const FORWARDED: [&str; 6] = [ "http_proxy", "https_proxy", @@ -30,9 +25,7 @@ const FORWARDED: [&str; 6] = [ const HEX: [u8; 16] = *b"0123456789abcdef"; -/// A suffix unique within this process (pid + monotonic counter), enough to -/// name a temp file a concurrent write in the same directory won't collide -/// with. Not cryptographic — just unique. +/// A suffix unique within this process, enough to name a temp file; not cryptographic. pub fn tmp_suffix() -> String { format!( "{}-{}", @@ -41,10 +34,7 @@ pub fn tmp_suffix() -> String { ) } -/// A 128-bit unpredictable hex token from the OS CSPRNG. Used where an -/// in-sandbox command must not be able to guess or forge the value (the -/// session-command sentinel). Falls back to the predictable tmp_suffix only if -/// the CSPRNG is unreadable, which never happens on a normal guest. +/// A 128-bit CSPRNG hex token, used where an in-sandbox command must not forge the value. pub fn rand_token() -> String { let mut b = [0u8; 16]; if !fill_random(&mut b) { @@ -56,32 +46,9 @@ pub fn rand_token() -> String { .collect() } -/// Fills `b` from the OS CSPRNG, sparing an async caller the open/read/close a -/// tokio worker would block on. GRND_NONBLOCK never waits on an early-boot -/// pool; /dev/urandom is the same source and covers the fallback. -fn fill_random(b: &mut [u8]) -> bool { - #[cfg(target_os = "linux")] - { - // SAFETY: getrandom writes at most b.len() bytes into b's live buffer. - let n = unsafe { libc::getrandom(b.as_mut_ptr().cast(), b.len(), libc::GRND_NONBLOCK) }; - if n == b.len() as isize { - return true; - } - } - std::fs::File::open("/dev/urandom") - .and_then(|mut f| f.read_exact(b)) - .is_ok() -} - -/// SIGKILLs the process group led by `pgid` — a session's shell plus its -/// foreground children, so tearing down a session that is running an external -/// command actually stops the command (killing only the shell would leave the -/// child holding the stdout pipe open). The shell is spawned as its own group -/// leader (pgid == its pid). +/// SIGKILLs the group led by `pgid`, so a session's external command dies with its shell. pub fn kill_group(pgid: u32) { - // kill(-0) targets the CALLER's group (silkd itself). Synthetic ids pass - // the guard: synth_pid() keeps them ≤ i32::MAX and above pid_max, so they - // reach kill() and miss with ESRCH. + // kill(-0) would target silkd's own group; a synthetic id passes the guard and misses with ESRCH. if !valid_pid(pgid) { return; } @@ -90,11 +57,9 @@ pub fn kill_group(pgid: u32) { unsafe { libc::kill(-(pgid as libc::pid_t), libc::SIGKILL) }; } -/// Sends `sig` to `pid`, ignoring the result: an ESRCH against a -/// just-exited pid already satisfies the caller's goal (the process is gone). +/// Sends `sig` to `pid`, ignoring ESRCH against a just-exited pid. pub fn signal_pid(pid: u32, sig: i32) { - // pid 0 means "my whole process group" to kill(2) — signalling it would - // take down silkd and every child (see kill_group for synthetic ids). + // pid 0 means silkd's whole process group to kill(2). if !valid_pid(pid) { return; } @@ -103,68 +68,17 @@ pub fn signal_pid(pid: u32, sig: i32) { unsafe { libc::kill(pid as libc::pid_t, sig) }; } -/// Rejects pid 0 and anything that would go negative through the pid_t cast -/// (and so hit a process group instead of a pid). -fn valid_pid(id: u32) -> bool { - id != 0 && id <= i32::MAX as u32 -} - -/// The environment every exec starts from before the request's env is layered -/// on — a sane PATH and TERM, plus the proxy snapshot on the no-network lane -/// (a lane with its own NIC would be steered into a closed relay). +/// The environment every exec starts from; the proxy snapshot rides only on the no-network lane. pub fn base_env() -> Vec<(&'static str, &'static str)> { compose_env(crate::net::has_egress(), proxy_vars()) } -/// Applies base_env's lane rule to an env-inheriting child (git, language -/// servers): with a network of its own, the proxy variables come off. +/// Applies base_env's lane rule to an env-inheriting child: an own NIC drops the proxy variables. pub fn align_proxy_env(cmd: &mut Command) { align_proxy_env_for(cmd, crate::net::has_egress()); } -fn align_proxy_env_for(cmd: &mut Command, nic: bool) { - if !nic { - return; - } - for key in FORWARDED { - cmd.env_remove(key); - } -} - -fn compose_env<'a>(nic: bool, proxy: &'a [(&'static str, String)]) -> Vec<(&'static str, &'a str)> { - let mut env = vec![ - ( - "PATH", - "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", - ), - ("TERM", "xterm-256color"), - ]; - if !nic { - env.extend(proxy.iter().map(|(k, v)| (*k, v.as_str()))); - } - env -} - -/// The FORWARDED values present in silkd's own environment, read once: the -/// unit environment is fixed at service start and silkd never mutates it. -fn proxy_vars() -> &'static [(&'static str, String)] { - static PROXY_VARS: LazyLock> = - LazyLock::new(|| snapshot_proxy(|key| std::env::var(key).ok())); - &PROXY_VARS -} - -fn snapshot_proxy(get: impl Fn(&str) -> Option) -> Vec<(&'static str, String)> { - FORWARDED - .iter() - .filter_map(|&key| match get(key) { - Some(v) if !v.is_empty() => Some((key, v)), - _ => None, - }) - .collect() -} - -/// Resolves a username to uid/gid via getpwnam and sets them on the command. -/// silkd runs as root, so this de-escalates; unknown users are rejected. +/// Resolves a username via getpwnam and de-escalates the command onto it. pub fn apply_user(cmd: &mut Command, user: &str) -> Result<(), String> { let (uid, gid, home) = lookup_user(user)?; cmd.uid(uid).gid(gid); @@ -173,31 +87,7 @@ pub fn apply_user(cmd: &mut Command, user: &str) -> Result<(), String> { Ok(()) } -fn lookup_user(user: &str) -> Result<(u32, u32, String), String> { - use std::ffi::CString; - let cname = CString::new(user).map_err(|_| format!("invalid user {user:?}"))?; - // Hold NSS_LOCK across the call and every read of the returned static - // buffer, so a concurrent lookup on another worker thread cannot clobber - // it mid-read (which would de-escalate to the wrong uid). - let _guard = lock(&NSS_LOCK); - // SAFETY: cname is a live NUL-terminated CString for the call. - let pw = unsafe { libc::getpwnam(cname.as_ptr()) }; - if pw.is_null() { - return Err(format!("unknown user {user:?}")); - } - // SAFETY: pw is non-null (checked) and, with NSS_LOCK still held, points to - // a valid passwd whose pw_dir is a NUL-terminated string it owns. - let pw = unsafe { &*pw }; - // SAFETY: pw_dir is NUL-terminated and stays valid while NSS_LOCK is held. - let home = unsafe { std::ffi::CStr::from_ptr(pw.pw_dir) } - .to_string_lossy() - .into_owned(); - Ok((pw.pw_uid, pw.pw_gid, home)) -} - -/// Opens a pseudo-terminal, returning the (master, slave) fds. The master is -/// set non-blocking for async I/O; the slave becomes the child's controlling -/// terminal after setsid (see `pty::open`). +/// Opens a pty, returning the (master, slave) fds; the master is non-blocking for async I/O. pub fn openpty(cols: u16, rows: u16) -> std::io::Result<(OwnedFd, OwnedFd)> { let mut master: libc::c_int = 0; let mut slave: libc::c_int = 0; @@ -249,9 +139,7 @@ pub fn set_winsize(fd: RawFd, cols: u16, rows: u16) -> std::io::Result<()> { Ok(()) } -/// Makes the calling process a session leader and adopts fd 0 as its -/// controlling terminal. Called in the child between fork and exec, so it must -/// touch nothing but the two syscalls. +/// Makes the calling process a session leader and adopts fd 0 as its controlling terminal. /// /// # Safety /// Only async-signal-safe syscalls; valid in a post-fork child. @@ -269,8 +157,7 @@ pub unsafe fn make_controlling_tty() -> std::io::Result<()> { } } -/// Reads from a raw fd (the pty master), returning bytes read; a WouldBlock -/// error is surfaced to the caller's readiness loop. +/// Reads from a raw fd, surfacing WouldBlock to the caller's readiness loop. pub fn read_fd(fd: RawFd, buf: &mut [u8]) -> std::io::Result { // SAFETY: fd is a live open fd; buf is a valid mutable slice of buf.len(). let n = unsafe { libc::read(fd, buf.as_mut_ptr().cast(), buf.len()) }; @@ -290,21 +177,7 @@ pub fn write_fd(fd: RawFd, buf: &[u8]) -> std::io::Result { Ok(n as usize) } -fn set_nonblocking(fd: RawFd) -> std::io::Result<()> { - // SAFETY: fd is open; F_GETFL/F_SETFL only read and set its flags. - let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) }; - if flags < 0 { - return Err(std::io::Error::last_os_error()); - } - // SAFETY: same fd; F_SETFL only sets the just-read flags plus O_NONBLOCK. - if unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0 { - return Err(std::io::Error::last_os_error()); - } - Ok(()) -} - -/// Locks a std mutex, panicking on poisoning: silkd's critical sections never -/// panic, so a poisoned lock is unreachable. +/// Locks a std mutex, panicking on poisoning: silkd's critical sections never panic. #[allow(clippy::unwrap_used)] pub fn lock(m: &Mutex) -> std::sync::MutexGuard<'_, T> { m.lock().unwrap() @@ -319,12 +192,104 @@ pub fn exit_code(status: ExitStatus) -> i32 { status.signal().map_or(-1, |s| 128 + s) } -/// Reaps the child and maps its status via `exit_code`; -1 when the wait -/// itself fails. +/// Reaps the child and maps its status via `exit_code`; -1 when the wait fails. pub async fn wait_code(child: &mut tokio::process::Child) -> i32 { child.wait().await.map_or(-1, exit_code) } +/// Fills `b` from the OS CSPRNG without blocking a tokio worker; /dev/urandom covers the fallback. +fn fill_random(b: &mut [u8]) -> bool { + #[cfg(target_os = "linux")] + { + // SAFETY: getrandom writes at most b.len() bytes into b's live buffer. + let n = unsafe { libc::getrandom(b.as_mut_ptr().cast(), b.len(), libc::GRND_NONBLOCK) }; + if n == b.len() as isize { + return true; + } + } + std::fs::File::open("/dev/urandom") + .and_then(|mut f| f.read_exact(b)) + .is_ok() +} + +/// Rejects pid 0 and anything that would go negative through the pid_t cast. +fn valid_pid(id: u32) -> bool { + id != 0 && id <= i32::MAX as u32 +} + +fn align_proxy_env_for(cmd: &mut Command, nic: bool) { + if !nic { + return; + } + for key in FORWARDED { + cmd.env_remove(key); + } +} + +fn compose_env<'a>(nic: bool, proxy: &'a [(&'static str, String)]) -> Vec<(&'static str, &'a str)> { + let mut env = vec![ + ( + "PATH", + "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", + ), + ("TERM", "xterm-256color"), + ]; + if !nic { + env.extend(proxy.iter().map(|(k, v)| (*k, v.as_str()))); + } + env +} + +/// The FORWARDED values in silkd's own environment, read once: the unit environment is fixed. +fn proxy_vars() -> &'static [(&'static str, String)] { + static PROXY_VARS: LazyLock> = + LazyLock::new(|| snapshot_proxy(|key| std::env::var(key).ok())); + &PROXY_VARS +} + +fn snapshot_proxy(get: impl Fn(&str) -> Option) -> Vec<(&'static str, String)> { + FORWARDED + .iter() + .filter_map(|&key| match get(key) { + Some(v) if !v.is_empty() => Some((key, v)), + _ => None, + }) + .collect() +} + +fn lookup_user(user: &str) -> Result<(u32, u32, String), String> { + use std::ffi::CString; + let cname = CString::new(user).map_err(|_| format!("invalid user {user:?}"))?; + // hold NSS_LOCK across every read of the static buffer, else a concurrent lookup clobbers it. + let _guard = lock(&NSS_LOCK); + // SAFETY: cname is a live NUL-terminated CString for the call. + let pw = unsafe { libc::getpwnam(cname.as_ptr()) }; + if pw.is_null() { + return Err(format!("unknown user {user:?}")); + } + // SAFETY: pw is non-null (checked) and, with NSS_LOCK still held, points to + // a valid passwd whose pw_dir is a NUL-terminated string it owns. + let pw = unsafe { &*pw }; + // SAFETY: pw_dir is NUL-terminated and stays valid while NSS_LOCK is held. + let home = unsafe { std::ffi::CStr::from_ptr(pw.pw_dir) } + .to_string_lossy() + .into_owned(); + Ok((pw.pw_uid, pw.pw_gid, home)) +} + +fn set_nonblocking(fd: RawFd) -> std::io::Result<()> { + // SAFETY: fd is open; F_GETFL/F_SETFL only read and set its flags. + let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) }; + if flags < 0 { + return Err(std::io::Error::last_os_error()); + } + // SAFETY: same fd; F_SETFL only sets the just-read flags plus O_NONBLOCK. + if unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0 { + return Err(std::io::Error::last_os_error()); + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; diff --git a/silkd/src/tree.rs b/silkd/src/tree.rs index 5a5c5e33..0e771654 100644 --- a/silkd/src/tree.rs +++ b/silkd/src/tree.rs @@ -1,14 +1,6 @@ -//! Whole-tree transfer: `fs.push` extracts a client tar stream into a -//! directory, `fs.pull` streams a path back as a tar. The archive format is -//! the contract and tar(1) is its authoritative implementation, so these -//! stream through the guest tar binary rather than buffering a whole project -//! in memory — the only project-ingestion path on the no-network lane. -//! -//! Push is network-failure atomic: the stream extracts into a staging dir -//! inside `dest` (same filesystem) and merges into place only after the tar -//! terminates cleanly, so a dropped connection or truncated stream leaves -//! `dest` untouched. The merge itself is local renames — a residual crash -//! window of microseconds instead of the whole transfer. +//! Whole-tree transfer through the guest tar binary: `fs.push` extracts a client +//! tar stream into a directory, `fs.pull` streams a path back as a tar. Push is +//! network-failure atomic: it stages inside `dest`, then merges by rename. use std::io; use std::path::Path; @@ -20,9 +12,7 @@ use tokio::process::Command; use crate::proto::{self, ErrorKind, Response, err_frame}; use crate::sysutil; -/// Extracts a client tar stream (`data` frames until `data_end`) into `dest`, -/// creating it if needed. A non-zero tar exit, a stray frame, or a truncated -/// stream is reported as an error with `dest` unchanged. +/// Extracts a client tar stream into `dest`, creating it; a stream failure leaves `dest` unchanged. pub async fn push(mut reader: R, w: &mut W, dest: String) -> io::Result<()> where R: AsyncBufRead + Unpin, @@ -36,29 +26,24 @@ where Err(e) => return err_frame(w, &e, "stage push").await, }; let res = push_staged(&mut reader, w, &staging, dest).await; - // One cleanup point for every exit path. A failed push can leave a large - // partially-extracted tree; tokio walks it off the async runtime. After a - // successful merge the renames have all but drained it. + // a failed push can leave a large partial tree, so cleanup runs on every exit path. let _ = tokio::fs::remove_dir_all(&staging).await; res } -/// Streams `path` back as a tar (`data` frames then `done`). The archive -/// carries the entry under its own name, so the client extracts it back to -/// the same basename. +/// Streams `path` back as a tar carrying the entry under its own basename. pub async fn pull(w: &mut W, path: String) -> io::Result<()> { let p = Path::new(&path); let (parent, name) = match (p.parent(), p.file_name()) { - (Some(par), Some(n)) if !n.is_empty() => (par.to_path_buf(), n.to_os_string()), + (Some(par), Some(n)) if !n.is_empty() => (par, n), _ => return proto::error_frame(w, ErrorKind::BadRequest, "invalid path").await, }; - // symlink_metadata, not metadata: a dangling symlink is a valid tar source - // (tar archives the link itself), and following it would wrongly 404 it. + // symlink_metadata: a dangling symlink is still a valid tar source. if let Err(e) = tokio::fs::symlink_metadata(p).await { return err_frame(w, &e, "stat source").await; } let cwd = if parent.as_os_str().is_empty() { - Path::new(".").to_path_buf() + Path::new(".") } else { parent }; @@ -66,9 +51,9 @@ pub async fn pull(w: &mut W, path: String) -> io::Result< // `--` so a basename starting with `-` is a path, not a tar option. cmd.arg("-c") .arg("-C") - .arg(&cwd) + .arg(cwd) .arg("--") - .arg(&name) + .arg(name) .stdout(Stdio::piped()); let (mut child, err_task) = match spawn_tar(cmd) { Ok(pair) => pair, @@ -87,8 +72,7 @@ pub async fn pull(w: &mut W, path: String) -> io::Result< tar_result(w, status, &msg, "tar create").await } -/// Runs the tar extraction into `staging` and merges into `dest`; the caller -/// owns removing `staging` on every path. +/// The caller owns removing `staging` on every path. async fn push_staged(reader: &mut R, w: &mut W, staging: &str, dest: String) -> io::Result<()> where R: AsyncBufRead + Unpin, @@ -120,8 +104,7 @@ where return tar_result(w, status, &msg, "tar extract").await; } let src = staging.to_owned(); - // The merge walks and renames on the blocking pool: pure local fs - // metadata work, but a big tree is thousands of syscalls. + // spawn_blocking: a big tree is thousands of local fs syscalls. let merged = tokio::task::spawn_blocking(move || merge_tree(Path::new(&src), Path::new(&dest))) .await .map_err(io::Error::other) @@ -132,8 +115,6 @@ where } } -/// Terminal frame for a finished tar: `done` on success, else the captured -/// stderr under `label`. async fn tar_result( w: &mut W, status: ExitStatus, @@ -143,9 +124,7 @@ async fn tar_result( proto::subprocess_result(w, status.success(), label, msg.as_bytes()).await } -/// Spawns a configured tar with stderr piped and drained on its own task — -/// a noisy tar filling the stderr pipe must not deadlock against the stdio -/// end we are servicing. +/// Spawns tar draining stderr on its own task, so a full stderr pipe cannot deadlock. fn spawn_tar( mut cmd: Command, ) -> io::Result<(tokio::process::Child, tokio::task::JoinHandle)> { @@ -157,18 +136,14 @@ fn spawn_tar( Ok((child, tokio::spawn(drain(stderr)))) } -/// Creates a unique staging dir inside `dest` — the same filesystem, so the -/// post-extract merge is pure renames. +/// Creates the staging dir inside `dest` so the merge is same-filesystem renames. async fn stage_dir(dest: &str) -> io::Result { let staging = format!("{dest}/.silkd-push-{}", sysutil::tmp_suffix()); tokio::fs::create_dir(&staging).await?; Ok(staging) } -/// Overlays `src` into `dst` with tar's semantics: directories merge -/// recursively, files and symlinks rename atomically over an existing file. -/// A subtree new to `dst` moves in one rename; a file landing on a directory -/// fails like tar does. +/// Overlays `src` into `dst` with tar's semantics: dirs merge, a file over a dir fails. fn merge_tree(src: &Path, dst: &Path) -> io::Result<()> { for entry in std::fs::read_dir(src)? { let entry = entry?; @@ -177,8 +152,7 @@ fn merge_tree(src: &Path, dst: &Path) -> io::Result<()> { let to_meta = std::fs::symlink_metadata(&to).ok(); match (entry.file_type()?.is_dir(), to_meta) { (true, Some(m)) if m.is_dir() => merge_tree(&from, &to)?, - // A dir replacing a file/symlink drops it first, mirroring tar's - // unlink-before-extract. + // mirrors tar's unlink-before-extract. (true, Some(_)) => { std::fs::remove_file(&to)?; std::fs::rename(&from, &to)?; @@ -189,8 +163,7 @@ fn merge_tree(src: &Path, dst: &Path) -> io::Result<()> { Ok(()) } -/// Reads a child's stderr to a string, keeping the first CAP bytes so a -/// pathological tar can't balloon memory — tar's first error is the useful one. +/// Reads a child's stderr, capped at CAP bytes; tar's first error is the useful one. async fn drain(mut stderr: tokio::process::ChildStderr) -> String { const CAP: usize = 16 * 1024; let mut out = Vec::new(); @@ -203,5 +176,5 @@ async fn drain(mut stderr: tokio::process::ChildStderr) -> String { out.extend_from_slice(&buf[..(n.min(CAP - out.len()))]); } } - String::from_utf8_lossy(&out).into_owned() + String::from_utf8(out).unwrap_or_else(|e| String::from_utf8_lossy(e.as_bytes()).into_owned()) } diff --git a/silkd/src/vsock.rs b/silkd/src/vsock.rs index 0fca4b11..901baf33 100644 --- a/silkd/src/vsock.rs +++ b/silkd/src/vsock.rs @@ -1,8 +1,4 @@ -//! vsock listener. On Linux silkd binds the guest's hybrid-vsock port and -//! serves each connection; the host (sandboxd) reaches it through the VMM -//! muxer's `CONNECT ` handshake, which the VMM answers — silkd sees a -//! plain byte stream. Off Linux the crate still builds (host tooling, tests) -//! with a stub so `cargo test` runs everywhere. +//! vsock listener: silkd sees a plain byte stream, the VMM muxer having answered the `CONNECT ` handshake. use std::sync::Arc; @@ -14,9 +10,7 @@ pub async fn serve(port: u32, state: Arc) -> std::io::Result<()> { let listener = VsockListener::bind(VsockAddr::new(VMADDR_CID_ANY, port))?; loop { - // A transient accept error (e.g. EMFILE under load) must not tear down - // the daemon and lose the process table and every session — only a bind - // failure above is fatal. + // a transient accept error must not tear down the daemon and lose every session. let conn = match listener.accept().await { Ok((conn, _)) => conn, Err(e) => { @@ -28,8 +22,7 @@ pub async fn serve(port: u32, state: Arc) -> std::io::Result<()> { let state = Arc::clone(&state); tokio::spawn(async move { let (read, write) = tokio::io::split(conn); - // Data/stdin frames run 43KB–1.3MB on the wire; the 8KB BufReader - // default would cost several read syscalls per frame on bulk paths. + // bulk frames run 43KB-1.3MB, so the 8KB BufReader default costs several reads per frame. let reader = tokio::io::BufReader::with_capacity(64 * 1024, read); if let Err(e) = state.serve(reader, write).await { eprintln!("silkd: connection: {e}"); diff --git a/silkd/src/watch.rs b/silkd/src/watch.rs index c618325d..52429804 100644 --- a/silkd/src/watch.rs +++ b/silkd/src/watch.rs @@ -1,6 +1,4 @@ -//! `fs.watch`: stream filesystem events under a path until the client -//! disconnects. Like every connection-bound verb, an event feed has no -//! meaningful detached state, so it lives only as long as its connection. +//! `fs.watch` streams filesystem events until the client disconnects; the feed has no detached state. use notify::{RecursiveMode, Watcher}; use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite}; @@ -10,8 +8,7 @@ use crate::proto::{self, ErrorKind, EventKind, Response}; const OVERFLOW_MESSAGE: &str = "watch event queue overflow"; -/// Events written per syscall. A build tool emits them in tight bursts, so -/// draining the queue in batches also keeps it far from its 256-slot cap. +/// Events written per syscall; batching keeps the queue far from its 256-slot cap. const EVENT_BATCH: usize = 64; /// Watches `path`, writing `ready`, ordered events, or a terminal error until disconnect. @@ -62,8 +59,7 @@ where return Ok(()); } }, - // The client sends nothing during a watch, so any readable state — - // EOF (disconnect), a stray frame, or an error — ends the watch. + // the client sends nothing during a watch, so any readable state ends it. _ = reader.fill_buf() => return Ok(()), } } @@ -82,8 +78,7 @@ fn forward_frames(tx: &mut Option>, frames: &mut Vec, out: &mut Vec) { use notify::EventKind as N; let event = match res { diff --git a/silkd/tests/common/mod.rs b/silkd/tests/common/mod.rs index c4eb9884..28901d11 100644 --- a/silkd/tests/common/mod.rs +++ b/silkd/tests/common/mod.rs @@ -1,8 +1,5 @@ -//! Shared harness for the integration tests: drive silkd's server over an -//! in-memory duplex exactly as a relayed host connection would. -//! -//! Compiled into each test binary separately, so not every helper is used by -//! every binary — allow the resulting dead_code rather than fragment this. +//! Shared harness: drives silkd's server over an in-memory duplex as a relayed host connection would. +//! Each test binary compiles it separately, so #![allow(dead_code)] covers the helpers one binary skips. #![allow(clippy::unwrap_used, clippy::expect_used)] #![allow(dead_code)] @@ -19,8 +16,6 @@ use tokio::task::JoinHandle; pub type FrameWriter = WriteHalf; pub type FrameLines = Lines>>; -/// Opens a framed connection to a task serving `state`: write half, response -/// lines, and the serve task handle, split so streaming verbs can interleave. pub fn connect(state: &Arc) -> (FrameWriter, FrameLines, JoinHandle>) { let (client, server) = tokio::io::duplex(1 << 20); let state = Arc::clone(state); @@ -30,12 +25,10 @@ pub fn connect(state: &Arc) -> (FrameWriter, FrameLines, JoinHandle Vec { request_on(&Arc::new(State::new()), &[request_line.to_string()]).await } -/// Runs request lines on one connection against `state`, returning the frames. pub async fn request_on(state: &Arc, lines: &[String]) -> Vec { let (mut cw, mut out, handle) = connect(state); for line in lines { @@ -52,12 +45,10 @@ pub async fn request_on(state: &Arc, lines: &[String]) -> Vec { frames } -/// Convenience for the many single-line-on-shared-state calls. pub async fn one(state: &Arc, request_line: &str) -> Vec { request_on(state, &[request_line.to_string()]).await } -/// Runs multiple request lines on one connection against a fresh server. pub async fn exchange(lines: &[String]) -> Vec { request_on(&Arc::new(State::new()), lines).await } @@ -76,7 +67,6 @@ pub fn decode(frame: &Value) -> Vec { .unwrap() } -/// Concatenated `data` of every frame of the given type (e.g. "stdout", "data"). pub fn payload(frames: &[Value], frame_type: &str) -> Vec { frames .iter() @@ -85,12 +75,10 @@ pub fn payload(frames: &[Value], frame_type: &str) -> Vec { .collect() } -/// UTF-8 text of the concatenated stdout frames. pub fn stdout_body(frames: &[Value]) -> String { String::from_utf8_lossy(&payload(frames, "stdout")).into_owned() } -/// A `data`-frame stream (16K chunks) terminated by `data_end`. pub fn data_frames(bytes: &[u8]) -> Vec { let mut lines: Vec = bytes .chunks(16 * 1024) diff --git a/silkd/tests/exec_e2e.rs b/silkd/tests/exec_e2e.rs index 7606fb9d..fb4e5b9e 100644 --- a/silkd/tests/exec_e2e.rs +++ b/silkd/tests/exec_e2e.rs @@ -1,6 +1,4 @@ -//! exec/procs verb E2E over the in-memory duplex. Spawns real processes -//! (echo, sh), so it runs on any Unix — keep it that way so exec-path bugs -//! surface on the dev host, not only on Linux CI. +//! exec/procs verb E2E over the in-memory duplex, spawning real processes so it runs on any Unix. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common; diff --git a/silkd/tests/forward_e2e.rs b/silkd/tests/forward_e2e.rs index cdc4979e..c959d721 100644 --- a/silkd/tests/forward_e2e.rs +++ b/silkd/tests/forward_e2e.rs @@ -1,6 +1,4 @@ -//! port_forward integration: a real TCP echo server behind the verb. Every -//! test runs under a deadline — the regressions this suite exists to catch -//! (relay deadlocks) would otherwise hang CI instead of failing it. +//! port_forward integration against a real TCP echo server; every test runs under a deadline so a relay deadlock fails CI. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common; @@ -118,7 +116,6 @@ async fn forward_bidirectional_bulk_no_deadlock() { .expect("test deadline"); } -/// Extracts the base64 body of a `data` response frame; None for ready/done. fn data_payload(line: &str) -> Option<&str> { let start = line.find("\"data\":\"")? + "\"data\":\"".len(); let rest = &line[start..]; diff --git a/silkd/tests/git_e2e.rs b/silkd/tests/git_e2e.rs index 6c758db6..de36626f 100644 --- a/silkd/tests/git_e2e.rs +++ b/silkd/tests/git_e2e.rs @@ -1,5 +1,4 @@ -//! git verb E2E against a real git binary in temp repos. Local verbs run on -//! any host; the none-lane guard is exercised via the test lane override. +//! git verb E2E against a real git binary in temp repos; the none-lane guard uses the test lane override. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common; diff --git a/silkd/tests/lsp_e2e.rs b/silkd/tests/lsp_e2e.rs index e45c91d6..fbadd53f 100644 --- a/silkd/tests/lsp_e2e.rs +++ b/silkd/tests/lsp_e2e.rs @@ -1,8 +1,4 @@ -//! LSP broker integration: a fake "language server" (a shell script that -//! echoes a canned reply) stands in for pylsp, so the broker mechanism — -//! manifest lookup, spawn, bidirectional relay, stop — is exercised without a -//! real language server. Every test runs under a deadline so a relay deadlock -//! fails CI instead of hanging it. +//! LSP broker integration against a fake language server, each test under a deadline so a relay deadlock fails CI. #![allow(clippy::unwrap_used, clippy::expect_used)] use std::io::Write; @@ -28,7 +24,6 @@ IFS= read -r line printf 'reply:%s\n' "$line" "#; -/// Creates a fake language server and points SILKD_LSP_DIR at its dir. fn manifest_env(server_body: &str) -> TempDir { let dir = tempfile::tempdir().unwrap(); let bin = dir.path().join("fake-lsp"); diff --git a/silkd/tests/pty_e2e.rs b/silkd/tests/pty_e2e.rs index 571ca87d..08ff0750 100644 --- a/silkd/tests/pty_e2e.rs +++ b/silkd/tests/pty_e2e.rs @@ -1,5 +1,4 @@ //! pty.open E2E: a shell under a pseudo-terminal, driven over a duplex. -//! openpty exists on macOS and Linux, so this runs on the dev host. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common; @@ -17,7 +16,6 @@ async fn send(cw: &mut FrameWriter, frame: Value) { cw.write_all(b"\n").await.unwrap(); } -/// Reads frames until one satisfies `pred` or 5s elapses. async fn read_until(lines: &mut FrameLines, pred: impl Fn(&Value) -> bool) -> Value { tokio::time::timeout(Duration::from_secs(5), async { loop { diff --git a/silkd/tests/session_e2e.rs b/silkd/tests/session_e2e.rs index 677ec2a7..ac0dce5f 100644 --- a/silkd/tests/session_e2e.rs +++ b/silkd/tests/session_e2e.rs @@ -1,6 +1,4 @@ -//! Session verb E2E: a persistent shell whose cwd/env/state survive across -//! exec calls. Spawns real bash, so Unix-only in practice; kept ungated so -//! the persistence contract is checked on the dev host too. +//! Session verb E2E: a persistent shell whose cwd/env/state survive across exec calls. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common; diff --git a/silkd/tests/tree_e2e.rs b/silkd/tests/tree_e2e.rs index 6dd4ddc7..4fc45f83 100644 --- a/silkd/tests/tree_e2e.rs +++ b/silkd/tests/tree_e2e.rs @@ -1,5 +1,4 @@ -//! Whole-tree transfer E2E: push a tar stream in, pull a tar stream out, -//! using the system tar to build/extract the reference archives. +//! Whole-tree transfer E2E: push a tar stream in, pull a tar stream out, via the system tar. #![allow(clippy::unwrap_used, clippy::expect_used)] mod common;