Compare commits

...
31 Commits
Author SHA1 Message Date
DuProcess 917d0b74b7 Harden embedded server lifecycle
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 22:10:00 -04:00
DuProcess eb0c1db837 Frame terminal reports across input reads
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 22:01:02 -04:00
DuProcess f8693f08b5 Release v1.0.0-rc49
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 21:39:01 -04:00
DuProcess 26532fc0e1 Stop idle keepalives from repainting terminals
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 21:36:44 -04:00
DuProcess c5f699a6ef Release v1.0.0-rc48
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 21:25:28 -04:00
DuProcess 60403ba4c3 Stop repainting healthy idle TUIs
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 21:19:31 -04:00
DuProcess 833ac1082f Release v1.0.0-rc47
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 21:11:29 -04:00
DuProcess 97cf165527 Recover from sustained terminal backpressure
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 21:06:26 -04:00
DuProcess 8f2d57d95e Exercise terminal reconnects across platforms
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 21:00:20 -04:00
DuProcess 24180c5092 Release v1.0.0-rc46
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 20:54:02 -04:00
DuProcess 0fdfc0ee22 Keep terminal backpressure test MTU safe
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:48:50 -04:00
DuProcess 5dceb2792d Decouple terminal rendering from transport
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:47:42 -04:00
DuProcess 58ac974fe2 Harden Windows terminal byte handling
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:33:44 -04:00
DuProcess 0cdcaaaec1 Release v1.0.0-rc45
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 20:13:08 -04:00
DuProcess 9245c666af Bound reliable stream bursts for terminal fairness
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:08:39 -04:00
DuProcess af11ab889e Keep large file regression test fast
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:03:05 -04:00
DuProcess 43a7a69b9b Keep reliable stream packets below path MTU
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 20:01:12 -04:00
DuProcess 86a1942aa0 Trace terminal loop failures
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:55:16 -04:00
DuProcess 7a9ad38657 Surface service forwarder failures
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:52:19 -04:00
DuProcess 9cec63aeb5 Trace forwarded stream lifecycles
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:50:49 -04:00
DuProcess 3364a7eb7b Release v1.0.0-rc44
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 19:36:21 -04:00
DuProcess 970d54b991 Preserve restart-orphaned sessions for full timeout
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:32:28 -04:00
DuProcess b8b56c0f32 Honor configured shell for remote exec
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:26:08 -04:00
DuProcess 7c915a9fd7 Release v1.0.0-rc43
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
2026-07-17 19:10:26 -04:00
DuProcess d0917c952f Bound terminal output under sustained load
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:08:56 -04:00
DuProcess d4de2915a1 Prevent terminal input from starving reconnects
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 19:05:14 -04:00
DuProcess 7f8d5711ea Unify release selection across clients
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 18:58:58 -04:00
DuProcess 89b81d73b1 Restore reconnect status from snapshots
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 18:42:17 -04:00
DuProcess 25c3358842 Preserve ambiguous terminal input
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 18:31:30 -04:00
DuProcess 367af4de82 Bound SDK integration test waits
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
2026-07-17 18:26:45 -04:00
DuProcess 600f68b63d Prefer configured identities in SDK auth 2026-07-17 18:26:35 -04:00
17 changed files with 3625 additions and 1020 deletions
+3
View File
@@ -88,6 +88,9 @@ jobs:
$errors | ForEach-Object { Write-Error $_ } $errors | ForEach-Object { Write-Error $_ }
exit 1 exit 1
} }
- name: Windows installer runtime test
shell: pwsh
run: scripts/test-install-windows.ps1
- name: cmd installer smoke check - name: cmd installer smoke check
shell: pwsh shell: pwsh
run: | run: |
Generated
+1 -1
View File
@@ -436,7 +436,7 @@ dependencies = [
[[package]] [[package]]
name = "dosh" name = "dosh"
version = "1.0.0-rc42" version = "1.0.0-rc49"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"base64", "base64",
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "dosh" name = "dosh"
version = "1.0.0-rc42" version = "1.0.0-rc49"
edition = "2024" edition = "2024"
license = "MIT" license = "MIT"
+3
View File
@@ -16,6 +16,9 @@ async fn main() -> Result<()> {
client.user, client.session, client.conn_id client.user, client.session, client.conn_id
); );
} }
DoshServerEvent::Disconnected(client) => {
eprintln!("disconnected conn={:?}", client.conn_id);
}
DoshServerEvent::Session { DoshServerEvent::Session {
conn_id, conn_id,
event: SessionEvent::Stream(TransportEvent::Open(open)), event: SessionEvent::Stream(TransportEvent::Open(open)),
+36 -132
View File
@@ -11,14 +11,17 @@ param(
[string]$BinaryBase = $env:DOSH_BINARY_BASE, [string]$BinaryBase = $env:DOSH_BINARY_BASE,
[string]$BinaryName = $env:DOSH_BINARY_NAME, [string]$BinaryName = $env:DOSH_BINARY_NAME,
[string]$BinaryVersion = $(if ($env:DOSH_BINARY_VERSION) { $env:DOSH_BINARY_VERSION } else { "latest" }), [string]$BinaryVersion = $(if ($env:DOSH_BINARY_VERSION) { $env:DOSH_BINARY_VERSION } else { "latest" }),
[switch]$BinaryExact = $($env:DOSH_BINARY_EXACT -and $env:DOSH_BINARY_EXACT -ne "0"),
[string]$UpdateCache = $(if ($env:DOSH_UPDATE_CACHE) { $env:DOSH_UPDATE_CACHE } elseif ($env:LOCALAPPDATA) { Join-Path $env:LOCALAPPDATA "dosh\source" } else { Join-Path $HOME ".cache\dosh\source" }), [string]$UpdateCache = $(if ($env:DOSH_UPDATE_CACHE) { $env:DOSH_UPDATE_CACHE } elseif ($env:LOCALAPPDATA) { Join-Path $env:LOCALAPPDATA "dosh\source" } else { Join-Path $HOME ".cache\dosh\source" }),
[switch]$BinaryRequired = $($env:DOSH_BINARY_REQUIRED -and $env:DOSH_BINARY_REQUIRED -ne "0"), [switch]$BinaryRequired = $($env:DOSH_BINARY_REQUIRED -and $env:DOSH_BINARY_REQUIRED -ne "0"),
[switch]$FromCurrent,
[switch]$ForceConfig [switch]$ForceConfig
) )
$ErrorActionPreference = "Stop" $ErrorActionPreference = "Stop"
$ProgressPreference = "SilentlyContinue" $ProgressPreference = "SilentlyContinue"
$Quiet = $env:DOSH_UPDATE_QUIET -and $env:DOSH_UPDATE_QUIET -ne "0" $Quiet = $env:DOSH_UPDATE_QUIET -and $env:DOSH_UPDATE_QUIET -ne "0"
$script:ResolvedSourceRef = $null
function Write-Info($Message) { function Write-Info($Message) {
if (-not $Quiet) { if (-not $Quiet) {
@@ -121,96 +124,18 @@ function Release-DownloadUrl {
if (-not $web) { if (-not $web) {
return $null return $null
} }
if ($BinaryVersion -eq "latest") { $releaseRef = Source-ReleaseRef
return "$web/releases/latest/download/$name" if (-not $releaseRef) {
}
"$web/releases/download/$BinaryVersion/$name"
}
function Release-LatestTagDownloadUrl {
if ($BinaryUrl -or $BinaryBase -or $BinaryVersion -ne "latest") {
return $null return $null
} }
$web = Repo-WebBase $Repo "$web/releases/download/$releaseRef/$name"
if (-not $web) {
return $null
}
try {
$response = Invoke-WebRequest -UseBasicParsing -Uri "$web/releases/latest" -MaximumRedirection 5
$effective = $null
if ($response.BaseResponse.ResponseUri) {
$effective = $response.BaseResponse.ResponseUri.AbsoluteUri
} elseif ($response.BaseResponse.RequestMessage -and $response.BaseResponse.RequestMessage.RequestUri) {
$effective = $response.BaseResponse.RequestMessage.RequestUri.AbsoluteUri
}
$prefix = "$web/releases/tag/"
if ($effective -and $effective.StartsWith($prefix)) {
$tag = $effective.Substring($prefix.Length)
if ($tag) {
return "$web/releases/download/$tag/$(Release-ArtifactName)"
}
}
}
catch {
return $null
}
return $null
}
function Version-Core($Value) {
$base = $Value.TrimStart("v").Split("+")[0]
$base.Split("-")[0]
}
function Version-Prerelease($Value) {
$base = $Value.TrimStart("v").Split("+")[0]
$dash = $base.IndexOf("-")
if ($dash -lt 0) {
return ""
}
$base.Substring($dash + 1).ToLowerInvariant()
}
function Version-Parts($Value) {
[regex]::Matches((Version-Core $Value), "\d+") | ForEach-Object { [int64]$_.Value }
}
function Compare-Prerelease($Left, $Right) {
if (-not $Left -and -not $Right) {
return 0
}
if (-not $Left) {
return 1
}
if (-not $Right) {
return -1
}
if ($Left -eq $Right) {
return 0
}
if ($Left -match "(?i)^rc(\d+)$" -and $Right -match "(?i)^rc(\d+)$") {
$leftRc = [int64]([regex]::Match($Left, "(?i)^rc(\d+)$").Groups[1].Value)
$rightRc = [int64]([regex]::Match($Right, "(?i)^rc(\d+)$").Groups[1].Value)
return $leftRc.CompareTo($rightRc)
}
return [string]::CompareOrdinal($Left, $Right)
}
function Compare-DoshVersion($Left, $Right) {
$leftParts = @(Version-Parts $Left)
$rightParts = @(Version-Parts $Right)
$width = [Math]::Max($leftParts.Count, $rightParts.Count)
for ($i = 0; $i -lt $width; $i++) {
$l = if ($i -lt $leftParts.Count) { $leftParts[$i] } else { 0 }
$r = if ($i -lt $rightParts.Count) { $rightParts[$i] } else { 0 }
if ($l -lt $r) { return -1 }
if ($l -gt $r) { return 1 }
}
return Compare-Prerelease (Version-Prerelease $Left) (Version-Prerelease $Right)
} }
function Current-SourceVersion { function Current-SourceVersion {
if (Test-Path "Cargo.toml") { if ($FromCurrent) {
if (-not (Test-Path "Cargo.toml")) {
return $null
}
$raw = Get-Content "Cargo.toml" -Raw $raw = Get-Content "Cargo.toml" -Raw
} else { } else {
$web = Repo-WebBase $Repo $web = Repo-WebBase $Repo
@@ -230,47 +155,23 @@ function Current-SourceVersion {
return $null return $null
} }
function Latest-ReleaseTag { function Source-ReleaseRef {
$web = Repo-WebBase $Repo if ($BinaryUrl -or $BinaryBase) {
if (-not $web) {
return $null return $null
} }
try { if ($BinaryVersion -ne "latest" -and $BinaryExact) {
$response = Invoke-WebRequest -UseBasicParsing -Uri "$web/releases/latest" -MaximumRedirection 5 return $BinaryVersion
$effective = $null
if ($response.BaseResponse.ResponseUri) {
$effective = $response.BaseResponse.ResponseUri.AbsoluteUri
} elseif ($response.BaseResponse.RequestMessage -and $response.BaseResponse.RequestMessage.RequestUri) {
$effective = $response.BaseResponse.RequestMessage.RequestUri.AbsoluteUri
}
$prefix = "$web/releases/tag/"
if ($effective -and $effective.StartsWith($prefix)) {
return $effective.Substring($prefix.Length)
}
} }
catch { $current = Current-SourceVersion
return $null if ($current) {
return "v$($current.TrimStart('v'))"
}
if ($BinaryVersion -ne "latest") {
return $BinaryVersion
} }
return $null return $null
} }
function Latest-ReleaseIsStale {
if ($BinaryUrl -or $BinaryBase -or $BinaryVersion -ne "latest") {
return $false
}
$latest = Latest-ReleaseTag
$current = Current-SourceVersion
if (-not $latest -or -not $current) {
return $false
}
$latestVersion = $latest.TrimStart("v")
if ((Compare-DoshVersion $latestVersion $current) -lt 0) {
Write-Warning "latest release $latestVersion is older than source $current; skipping stale prebuilt"
return $true
}
return $false
}
function Verify-ArchiveChecksum($Url, $Archive) { function Verify-ArchiveChecksum($Url, $Archive) {
$checksumPath = "$Archive.sha256" $checksumPath = "$Archive.sha256"
try { try {
@@ -475,16 +376,11 @@ if ($env:DOSH_INSTALL_BINDIR_FILE) {
Apply-PendingBinaryReplacements $bindir Apply-PendingBinaryReplacements $bindir
function Install-Prebuilt { function Install-Prebuilt {
if (Latest-ReleaseIsStale) { $url = Release-DownloadUrl
return $false
}
$url = Release-LatestTagDownloadUrl
if (-not $url) {
$url = Release-DownloadUrl
}
if (-not $url) { if (-not $url) {
return $false return $false
} }
$script:ResolvedSourceRef = Source-ReleaseRef
$tmp = Join-Path ([System.IO.Path]::GetTempPath()) ("dosh-bin-" + [guid]::NewGuid()) $tmp = Join-Path ([System.IO.Path]::GetTempPath()) ("dosh-bin-" + [guid]::NewGuid())
$zip = Join-Path $tmp (Release-ArtifactName) $zip = Join-Path $tmp (Release-ArtifactName)
$extract = Join-Path $tmp "extract" $extract = Join-Path $tmp "extract"
@@ -520,7 +416,10 @@ function Install-Prebuilt {
function Install-FromSource { function Install-FromSource {
Ensure-Cargo Ensure-Cargo
if (Test-Path "Cargo.toml") { if ($FromCurrent) {
if (-not (Test-Path "Cargo.toml")) {
throw "-FromCurrent requires a Dosh checkout with Cargo.toml"
}
$src = (Get-Location).Path $src = (Get-Location).Path
} else { } else {
if (-not $Repo) { if (-not $Repo) {
@@ -530,21 +429,26 @@ function Install-FromSource {
$sourceCache = Assert-SafeUpdateCache $UpdateCache $sourceCache = Assert-SafeUpdateCache $UpdateCache
$parent = Split-Path -Parent $sourceCache $parent = Split-Path -Parent $sourceCache
New-Item -ItemType Directory -Force -Path $parent | Out-Null New-Item -ItemType Directory -Force -Path $parent | Out-Null
$sourceRef = if ($script:ResolvedSourceRef) { $script:ResolvedSourceRef } else { "main" }
if (Test-Path (Join-Path $sourceCache ".git")) { if (Test-Path (Join-Path $sourceCache ".git")) {
git -C $sourceCache remote set-url origin $Repo git -C $sourceCache remote set-url origin $Repo
$fetchArgs = @("-C", $sourceCache, "fetch", "--depth", "1", "origin", "main") $fetchArgs = @("-C", $sourceCache, "fetch", "--depth", "1", "origin", $sourceRef)
if ($Quiet) { if ($Quiet) {
$fetchArgs = @("-C", $sourceCache, "fetch", "-q", "--depth", "1", "origin", "main") $fetchArgs = @("-C", $sourceCache, "fetch", "-q", "--depth", "1", "origin", $sourceRef)
} }
git @fetchArgs git @fetchArgs
git -C $sourceCache checkout -q -B main FETCH_HEAD if ($sourceRef -eq "main") {
git -C $sourceCache checkout -q -B main FETCH_HEAD
} else {
git -C $sourceCache checkout -q --detach FETCH_HEAD
}
} else { } else {
if (Test-Path $sourceCache) { if (Test-Path $sourceCache) {
Remove-Item -Recurse -Force $sourceCache Remove-Item -Recurse -Force $sourceCache
} }
$cloneArgs = @("clone", "--depth", "1", "--branch", "main", $Repo, $sourceCache) $cloneArgs = @("clone", "--depth", "1", "--branch", $sourceRef, $Repo, $sourceCache)
if ($Quiet) { if ($Quiet) {
$cloneArgs = @("clone", "-q", "--depth", "1", "--branch", "main", $Repo, $sourceCache) $cloneArgs = @("clone", "-q", "--depth", "1", "--branch", $sourceRef, $Repo, $sourceCache)
} }
git @cloneArgs | Out-Null git @cloneArgs | Out-Null
} }
+38 -122
View File
@@ -17,6 +17,7 @@ binary_url="${DOSH_BINARY_URL:-}"
binary_base="${DOSH_BINARY_BASE:-}" binary_base="${DOSH_BINARY_BASE:-}"
binary_name="${DOSH_BINARY_NAME:-}" binary_name="${DOSH_BINARY_NAME:-}"
binary_version="${DOSH_BINARY_VERSION:-latest}" binary_version="${DOSH_BINARY_VERSION:-latest}"
binary_exact="${DOSH_BINARY_EXACT:-0}"
binary_required="${DOSH_BINARY_REQUIRED:-0}" binary_required="${DOSH_BINARY_REQUIRED:-0}"
usage() { usage() {
@@ -47,7 +48,9 @@ Environment alternatives:
DOSH_BINARY_NAME NAME DOSH_BINARY_NAME NAME
Release tarball name; defaults to dosh-OS-ARCH.tar.gz Release tarball name; defaults to dosh-OS-ARCH.tar.gz
DOSH_BINARY_VERSION TAG DOSH_BINARY_VERSION TAG
Release tag when deriving DOSH_BINARY_BASE; default latest Requested release tag; repository version wins unless exact
DOSH_BINARY_EXACT=1
Require DOSH_BINARY_VERSION exactly, including downgrades
DOSH_BINARY_REQUIRED=1 DOSH_BINARY_REQUIRED=1
Fail instead of falling back to source when binary install fails Fail instead of falling back to source when binary install fails
EOF EOF
@@ -173,6 +176,7 @@ config_dir="$HOME/.config/dosh"
data_dir="$HOME/.local/share/dosh" data_dir="$HOME/.local/share/dosh"
systemd_user_dir="$HOME/.config/systemd/user" systemd_user_dir="$HOME/.config/systemd/user"
src_dir="" src_dir=""
resolved_source_ref=""
mkdir -p "$bindir" "$config_dir" "$data_dir" mkdir -p "$bindir" "$config_dir" "$data_dir"
@@ -228,97 +232,13 @@ release_download_url() {
return 1 return 1
fi fi
web_base="$(repo_web_base "$repo")" || return 1 web_base="$(repo_web_base "$repo")" || return 1
if [ "$binary_version" = "latest" ]; then release_ref="$(release_source_ref)" || return 1
printf '%s/releases/latest/download/%s\n' "$web_base" "$(release_artifact_name)" printf '%s/releases/download/%s/%s\n' "$web_base" "$release_ref" "$(release_artifact_name)"
else
printf '%s/releases/download/%s/%s\n' "$web_base" "$binary_version" "$(release_artifact_name)"
fi
}
release_latest_tag_download_url() {
if [ -n "$binary_url" ] || [ -n "$binary_base" ] || [ "$binary_version" != "latest" ] || [ -z "$repo" ]; then
return 1
fi
web_base="$(repo_web_base "$repo")" || return 1
latest_url="$(curl -fsSL -o /dev/null -w '%{url_effective}' "$web_base/releases/latest" 2>/dev/null || true)"
case "$latest_url" in
"$web_base"/releases/tag/*)
tag="${latest_url##"$web_base"/releases/tag/}"
[ -n "$tag" ] || return 1
printf '%s/releases/download/%s/%s\n' "$web_base" "$tag" "$(release_artifact_name)"
;;
*)
return 1
;;
esac
}
version_core() {
value="${1#v}"
value="${value%%+*}"
printf '%s\n' "${value%%-*}"
}
version_prerelease() {
value="${1#v}"
value="${value%%+*}"
case "$value" in
*-*) printf '%s\n' "${value#*-}" ;;
*) printf '\n' ;;
esac
}
version_core_part() {
version="$1"
index="$2"
version_core "$version" \
| sed 's/[^0-9][^0-9]*/ /g' \
| awk -v index="$index" '{ value=$index; if (value == "") value=0; print value }'
}
prerelease_less_than() {
left="$1"
right="$2"
if [ -z "$left" ] && [ -z "$right" ]; then
return 1
fi
if [ -z "$left" ]; then
return 1
fi
if [ -z "$right" ]; then
return 0
fi
if [ "$left" = "$right" ]; then
return 1
fi
left_rc="$(printf '%s\n' "$left" | sed -n 's/^[Rr][Cc]\([0-9][0-9]*\)$/\1/p')"
right_rc="$(printf '%s\n' "$right" | sed -n 's/^[Rr][Cc]\([0-9][0-9]*\)$/\1/p')"
if [ -n "$left_rc" ] && [ -n "$right_rc" ]; then
[ "$left_rc" -lt "$right_rc" ]
return
fi
first="$(printf '%s\n%s\n' "$left" "$right" | LC_ALL=C sort | sed -n '1p')"
[ "$first" = "$left" ]
}
version_less_than() {
left="$1"
right="$2"
for index in 1 2 3 4 5 6 7 8; do
lpart="$(version_core_part "$left" "$index")"
rpart="$(version_core_part "$right" "$index")"
if [ "$lpart" -lt "$rpart" ]; then
return 0
fi
if [ "$lpart" -gt "$rpart" ]; then
return 1
fi
done
prerelease_less_than "$(version_prerelease "$left")" "$(version_prerelease "$right")"
} }
current_source_version() { current_source_version() {
if [ -f Cargo.toml ]; then if [ "$from_current" -eq 1 ]; then
[ -f Cargo.toml ] || return 1
sed -n 's/^version = "\(.*\)"/\1/p' Cargo.toml | sed -n '1p' sed -n 's/^version = "\(.*\)"/\1/p' Cargo.toml | sed -n '1p'
return 0 return 0
fi fi
@@ -329,32 +249,21 @@ current_source_version() {
| sed -n '1p' | sed -n '1p'
} }
latest_release_tag() { release_source_ref() {
[ -n "$repo" ] || return 1 if [ -n "$binary_url" ] || [ -n "$binary_base" ]; then
web_base="$(repo_web_base "$repo")" || return 1
latest_url="$(curl -fsSL -o /dev/null -w '%{url_effective}' "$web_base/releases/latest" 2>/dev/null || true)"
case "$latest_url" in
"$web_base"/releases/tag/*)
tag="${latest_url##"$web_base"/releases/tag/}"
[ -n "$tag" ] || return 1
printf '%s\n' "$tag"
;;
*)
return 1
;;
esac
}
latest_release_is_stale() {
if [ -n "$binary_url" ] || [ -n "$binary_base" ] || [ "$binary_version" != "latest" ]; then
return 1 return 1
fi fi
latest="$(latest_release_tag || true)" if [ "$binary_version" != "latest" ] && [ "$binary_exact" = "1" ]; then
printf '%s\n' "$binary_version"
return 0
fi
current="$(current_source_version || true)" current="$(current_source_version || true)"
[ -n "$latest" ] && [ -n "$current" ] || return 1 if [ -n "$current" ]; then
latest="${latest#v}" printf 'v%s\n' "${current#v}"
if version_less_than "$latest" "$current"; then return 0
echo "latest release $latest is older than source $current; skipping stale prebuilt" >&2 fi
if [ "$binary_version" != "latest" ]; then
printf '%s\n' "$binary_version"
return 0 return 0
fi fi
return 1 return 1
@@ -450,10 +359,8 @@ verify_archive_version() {
} }
try_install_prebuilt() { try_install_prebuilt() {
if latest_release_is_stale; then download_url="$(release_download_url)" || return 1
return 1 resolved_source_ref="$(release_source_ref || true)"
fi
download_url="$(release_latest_tag_download_url || release_download_url)" || return 1
tmpdir="$(mktemp -d)" tmpdir="$(mktemp -d)"
archive="$tmpdir/$(release_artifact_name)" archive="$tmpdir/$(release_artifact_name)"
checksum_file="$archive.sha256" checksum_file="$archive.sha256"
@@ -490,7 +397,11 @@ try_install_prebuilt() {
install_from_source() { install_from_source() {
ensure_cargo ensure_cargo
if [ "$from_current" -eq 1 ] || [ -f Cargo.toml ]; then if [ "$from_current" -eq 1 ]; then
[ -f Cargo.toml ] || {
echo "--from-current requires a Dosh checkout with Cargo.toml" >&2
exit 2
}
src_dir="$(pwd)" src_dir="$(pwd)"
else else
if [ -z "$repo" ]; then if [ -z "$repo" ]; then
@@ -500,22 +411,27 @@ install_from_source() {
update_cache="$(safe_update_cache_path "$update_cache")" update_cache="$(safe_update_cache_path "$update_cache")"
need git need git
mkdir -p "$(dirname "$update_cache")" mkdir -p "$(dirname "$update_cache")"
source_ref="${resolved_source_ref:-main}"
if [ -d "$update_cache/.git" ]; then if [ -d "$update_cache/.git" ]; then
[ "$quiet" != "1" ] && echo "Updating Dosh source" [ "$quiet" != "1" ] && echo "Updating Dosh source"
git -C "$update_cache" remote set-url origin "$repo" git -C "$update_cache" remote set-url origin "$repo"
if [ "$quiet" = "1" ]; then if [ "$quiet" = "1" ]; then
git -C "$update_cache" fetch -q --depth 1 origin main git -C "$update_cache" fetch -q --depth 1 origin "$source_ref"
else else
git -C "$update_cache" fetch --depth 1 origin main git -C "$update_cache" fetch --depth 1 origin "$source_ref"
fi
if [ "$source_ref" = "main" ]; then
git -C "$update_cache" checkout -q -B main FETCH_HEAD
else
git -C "$update_cache" checkout -q --detach FETCH_HEAD
fi fi
git -C "$update_cache" checkout -q -B main FETCH_HEAD
else else
[ "$quiet" != "1" ] && echo "Downloading Dosh source" [ "$quiet" != "1" ] && echo "Downloading Dosh source"
rm -rf "$update_cache" rm -rf "$update_cache"
if [ "$quiet" = "1" ]; then if [ "$quiet" = "1" ]; then
git clone -q --depth 1 --branch main "$repo" "$update_cache" >/dev/null git clone -q --depth 1 --branch "$source_ref" "$repo" "$update_cache" >/dev/null
else else
git clone --depth 1 --branch main "$repo" "$update_cache" >/dev/null git clone --depth 1 --branch "$source_ref" "$repo" "$update_cache" >/dev/null
fi fi
fi fi
src_dir="$update_cache" src_dir="$update_cache"
+168
View File
@@ -0,0 +1,168 @@
$ErrorActionPreference = "Stop"
$ProgressPreference = "SilentlyContinue"
$version = "9.8.7-rc3"
$root = Join-Path ([System.IO.Path]::GetTempPath()) ("dosh-installer-test-" + [guid]::NewGuid())
$webRoot = Join-Path $root "web"
$repoRoot = Join-Path $webRoot "repo"
$releaseRoot = Join-Path $repoRoot "releases\download\v$version"
$rawRoot = Join-Path $repoRoot "raw\branch\main"
$stageRoot = Join-Path $root "stage\dosh"
$requestLog = Join-Path $root "requests.log"
$prefix = Join-Path $root "prefix"
$home = Join-Path $root "home"
$unrelated = Join-Path $root "unrelated"
$job = $null
$savedBinaryVersion = $env:DOSH_BINARY_VERSION
$savedBinaryExact = $env:DOSH_BINARY_EXACT
function Invoke-TestInstaller($HomePath, $WorkingDirectory, $InstallPrefix, $Repository) {
$savedHome = $env:HOME
$savedUserProfile = $env:USERPROFILE
try {
$env:HOME = $HomePath
$env:USERPROFILE = $HomePath
Push-Location $WorkingDirectory
try {
$powershell = (Get-Process -Id $PID).Path
& $powershell -NoProfile -ExecutionPolicy Bypass -File (Join-Path $PSScriptRoot "..\install.ps1") -Role client -Repo $Repository -Prefix $InstallPrefix -BinaryRequired
if ($LASTEXITCODE -ne 0) {
throw "Windows installer exited with status $LASTEXITCODE"
}
}
finally {
Pop-Location
}
}
finally {
$env:HOME = $savedHome
$env:USERPROFILE = $savedUserProfile
}
}
try {
New-Item -ItemType Directory -Force -Path $releaseRoot, $rawRoot, (Join-Path $stageRoot "bin"), $home, $unrelated | Out-Null
@"
[package]
name = "dosh"
version = "$version"
"@ | Set-Content -Encoding ascii -NoNewline (Join-Path $rawRoot "Cargo.toml")
@"
[package]
name = "not-dosh"
version = "99.0.0"
"@ | Set-Content -Encoding ascii -NoNewline (Join-Path $unrelated "Cargo.toml")
"$version`n" | Set-Content -Encoding ascii -NoNewline (Join-Path $stageRoot "VERSION")
"fixture" | Set-Content -Encoding ascii -NoNewline (Join-Path $stageRoot "bin\dosh-client.exe")
$archive = Join-Path $releaseRoot "dosh-windows-x86_64.zip"
Compress-Archive -Force -Path $stageRoot -DestinationPath $archive
$digest = (Get-FileHash -Algorithm SHA256 $archive).Hash.ToLowerInvariant()
"$digest dosh-windows-x86_64.zip`n" | Set-Content -Encoding ascii -NoNewline "$archive.sha256"
$probe = [System.Net.Sockets.TcpListener]::new([System.Net.IPAddress]::Loopback, 0)
$probe.Start()
$port = ([System.Net.IPEndPoint]$probe.LocalEndpoint).Port
$probe.Stop()
$listenPrefix = "http://127.0.0.1:$port/"
$job = Start-Job -ArgumentList $webRoot, $listenPrefix, $requestLog -ScriptBlock {
param($Root, $Prefix, $Log)
$listener = [System.Net.HttpListener]::new()
$listener.Prefixes.Add($Prefix)
$listener.Start()
try {
while ($listener.IsListening) {
$context = $listener.GetContext()
$path = [System.Uri]::UnescapeDataString($context.Request.Url.AbsolutePath.TrimStart('/'))
Add-Content -Encoding ascii -Path $Log -Value $context.Request.Url.AbsolutePath
$file = Join-Path $Root ($path.Replace('/', [System.IO.Path]::DirectorySeparatorChar))
if (Test-Path -LiteralPath $file -PathType Leaf) {
$bytes = [System.IO.File]::ReadAllBytes($file)
$context.Response.StatusCode = 200
$context.Response.ContentLength64 = $bytes.Length
if ($context.Request.HttpMethod -ne "HEAD") {
$context.Response.OutputStream.Write($bytes, 0, $bytes.Length)
}
} else {
$context.Response.StatusCode = 404
}
$context.Response.Close()
}
}
finally {
$listener.Stop()
}
}
$ready = $false
for ($attempt = 0; $attempt -lt 50; $attempt++) {
try {
Invoke-WebRequest -UseBasicParsing -Uri "${listenPrefix}repo/raw/branch/main/Cargo.toml" | Out-Null
$ready = $true
break
}
catch {
Start-Sleep -Milliseconds 100
}
}
if (-not $ready) {
throw "installer fixture HTTP server did not start"
}
$env:DOSH_BINARY_VERSION = "v0.1.17"
Remove-Item Env:DOSH_BINARY_EXACT -ErrorAction SilentlyContinue
Invoke-TestInstaller $home $unrelated $prefix "${listenPrefix}repo.git"
if (-not (Test-Path (Join-Path $prefix "bin\dosh-client.exe"))) {
throw "Windows installer did not install dosh-client.exe"
}
$requests = Get-Content $requestLog -Raw
if ($requests -notmatch "/repo/raw/branch/main/Cargo.toml") {
throw "Windows installer did not read the repository release version"
}
if ($requests -notmatch "/repo/releases/download/v$([regex]::Escape($version))/dosh-windows-x86_64.zip") {
throw "Windows installer did not request the repository-declared release artifact"
}
if ($requests -match "/releases/latest" -or $requests -match "99.0.0") {
throw "Windows installer used a stable redirect or unrelated Cargo.toml"
}
"fixture-v2" | Set-Content -Encoding ascii -NoNewline (Join-Path $stageRoot "bin\dosh-client.exe")
Compress-Archive -Force -Path $stageRoot -DestinationPath $archive
$digest = (Get-FileHash -Algorithm SHA256 $archive).Hash.ToLowerInvariant()
"$digest dosh-windows-x86_64.zip`n" | Set-Content -Encoding ascii -NoNewline "$archive.sha256"
$installedAlias = Join-Path $prefix "bin\dosh.exe"
$lock = [System.IO.File]::Open($installedAlias, [System.IO.FileMode]::Open, [System.IO.FileAccess]::Read, [System.IO.FileShare]::None)
try {
Invoke-TestInstaller $home $unrelated $prefix "${listenPrefix}repo.git"
$pending = @(Get-ChildItem (Join-Path $prefix "bin") -Filter "dosh.exe.pending.*")
if ($pending.Count -ne 1) {
throw "locked Windows executable was not staged exactly once"
}
}
finally {
$lock.Dispose()
}
$replaced = $false
for ($attempt = 0; $attempt -lt 100; $attempt++) {
if ((Get-Content $installedAlias -Raw) -eq "fixture-v2" -and -not (Get-ChildItem (Join-Path $prefix "bin") -Filter "dosh.exe.pending.*")) {
$replaced = $true
break
}
Start-Sleep -Milliseconds 100
}
if (-not $replaced) {
throw "staged Windows executable was not applied after its lock was released"
}
Write-Host "Windows installer runtime test passed"
}
finally {
$env:DOSH_BINARY_VERSION = $savedBinaryVersion
$env:DOSH_BINARY_EXACT = $savedBinaryExact
if ($job) {
Stop-Job $job -ErrorAction SilentlyContinue
Remove-Job $job -Force -ErrorAction SilentlyContinue
}
Remove-Item -Recurse -Force $root -ErrorAction SilentlyContinue
}
+1261 -634
View File
File diff suppressed because it is too large Load Diff
+170 -39
View File
@@ -27,7 +27,9 @@ use dosh::protocol::{
StreamClose, StreamData, StreamEof, StreamOpen, StreamOpenOk, StreamOpenReject, StreamClose, StreamData, StreamEof, StreamOpen, StreamOpenOk, StreamOpenReject,
StreamWindowAdjust, TicketAttachBody, TicketAttachEnvelope, TicketAttachOkEnvelope, StreamWindowAdjust, TicketAttachBody, TicketAttachEnvelope, TicketAttachOkEnvelope,
}; };
use dosh::pty::{PtyHandle, PtyOutput, adopt_pty_from_fd, spawn_pty_session}; use dosh::pty::{
PTY_OUTPUT_QUEUE_CAPACITY, PtyHandle, PtyOutput, adopt_pty_from_fd, spawn_pty_session,
};
use dosh::udp::{is_transient_udp_error, is_transient_udp_send_error}; use dosh::udp::{is_transient_udp_error, is_transient_udp_send_error};
use sha2::{Digest, Sha256}; use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, HashMap, HashSet, VecDeque}; use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
@@ -46,7 +48,7 @@ use tokio::net::{TcpListener, TcpStream, UdpSocket, UnixListener};
use tokio::process::Command as TokioCommand; use tokio::process::Command as TokioCommand;
use tokio::sync::mpsc; use tokio::sync::mpsc;
const STREAM_INITIAL_WINDOW: usize = 1024 * 1024; const STREAM_INITIAL_WINDOW: usize = dosh::transport::DEFAULT_INITIAL_WINDOW;
const STREAM_RETIRED_TOMBSTONES: usize = 16 * 1024; const STREAM_RETIRED_TOMBSTONES: usize = 16 * 1024;
const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8; const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
@@ -209,7 +211,7 @@ async fn serve(config_path: Option<std::path::PathBuf>) -> Result<()> {
], ],
); );
let (pty_tx, mut pty_rx) = mpsc::unbounded_channel(); let (pty_tx, mut pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let state = Arc::new(Mutex::new(ServerState::new( let state = Arc::new(Mutex::new(ServerState::new(
config.clone(), config.clone(),
secret, secret,
@@ -242,6 +244,7 @@ async fn serve(config_path: Option<std::path::PathBuf>) -> Result<()> {
let retransmit_socket = Arc::clone(&socket); let retransmit_socket = Arc::clone(&socket);
tokio::spawn(async move { tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_millis(100)); let mut interval = tokio::time::interval(Duration::from_millis(100));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop { loop {
interval.tick().await; interval.tick().await;
if let Err(err) = retransmit_pending(&retransmit_state, &retransmit_socket).await { if let Err(err) = retransmit_pending(&retransmit_state, &retransmit_socket).await {
@@ -256,6 +259,7 @@ async fn serve(config_path: Option<std::path::PathBuf>) -> Result<()> {
let cleanup_state = Arc::clone(&state); let cleanup_state = Arc::clone(&state);
tokio::spawn(async move { tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(5)); let mut interval = tokio::time::interval(Duration::from_secs(5));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop { loop {
interval.tick().await; interval.tick().await;
cleanup_disconnected_clients(&cleanup_state); cleanup_disconnected_clients(&cleanup_state);
@@ -270,6 +274,7 @@ async fn serve(config_path: Option<std::path::PathBuf>) -> Result<()> {
let flush_state = Arc::clone(&state); let flush_state = Arc::clone(&state);
tokio::spawn(async move { tokio::spawn(async move {
let mut interval = tokio::time::interval(SCREEN_PERSIST_MAX_AGE); let mut interval = tokio::time::interval(SCREEN_PERSIST_MAX_AGE);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop { loop {
interval.tick().await; interval.tick().await;
flush_persistent_screens(&flush_state); flush_persistent_screens(&flush_state);
@@ -293,7 +298,7 @@ async fn serve(config_path: Option<std::path::PathBuf>) -> Result<()> {
struct ServerState { struct ServerState {
config: ServerConfig, config: ServerConfig,
secret: [u8; 32], secret: [u8; 32],
pty_tx: mpsc::UnboundedSender<PtyOutput>, pty_tx: mpsc::Sender<PtyOutput>,
sessions: HashMap<String, Session>, sessions: HashMap<String, Session>,
pending_native: HashMap<[u8; 16], PendingNativeAuth>, pending_native: HashMap<[u8; 16], PendingNativeAuth>,
next_server_stream_id: u64, next_server_stream_id: u64,
@@ -393,6 +398,10 @@ struct Session {
holder_control: Option<StdUnixStream>, holder_control: Option<StdUnixStream>,
/// Whether this session's shell lives in a holder process (persistent). /// Whether this session's shell lives in a holder process (persistent).
persistent: bool, persistent: bool,
/// The holder survived a server restart but no client has reattached yet.
/// Such sessions need the full reconnect window because the client may be
/// asleep during an unattended server update.
restart_orphaned: bool,
/// Bytes of session output since the screen was last mirrored to disk, used /// Bytes of session output since the screen was last mirrored to disk, used
/// to throttle the (atomic) screen-persistence writes. /// to throttle the (atomic) screen-persistence writes.
bytes_since_persist: usize, bytes_since_persist: usize,
@@ -504,7 +513,7 @@ impl ServerState {
fn new( fn new(
config: ServerConfig, config: ServerConfig,
secret: [u8; 32], secret: [u8; 32],
pty_tx: mpsc::UnboundedSender<PtyOutput>, pty_tx: mpsc::Sender<PtyOutput>,
) -> Self { ) -> Self {
let per_minute = config.native_auth_rate_limit_per_minute; let per_minute = config.native_auth_rate_limit_per_minute;
Self { Self {
@@ -564,6 +573,7 @@ impl ServerState {
empty_since: None, empty_since: None,
holder_control: control, holder_control: control,
persistent, persistent,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now() - SCREEN_PERSIST_MAX_AGE, last_screen_persist_at: Instant::now() - SCREEN_PERSIST_MAX_AGE,
@@ -709,6 +719,7 @@ impl ServerState {
empty_since: Some(Instant::now()), empty_since: Some(Instant::now()),
holder_control: Some(control), holder_control: Some(control),
persistent: true, persistent: true,
restart_orphaned: true,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: output_seq, last_persisted_seq: output_seq,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -727,6 +738,7 @@ impl ServerState {
if let Some(session) = self.sessions.get_mut(session_name) { if let Some(session) = self.sessions.get_mut(session_name) {
session.clients.insert(client_id, client); session.clients.insert(client_id, client);
session.empty_since = None; session.empty_since = None;
session.restart_orphaned = false;
self.client_index self.client_index
.insert(client_id, session_name.to_string()); .insert(client_id, session_name.to_string());
} }
@@ -2077,6 +2089,14 @@ async fn handle_stream_open(
} }
if open.target_host == FILE_STREAM_SENTINEL { if open.target_host == FILE_STREAM_SENTINEL {
dosh::trace::event(
"server.file_stream_open",
&[
("session", session_name.clone()),
("stream", open.stream_id.to_string()),
("peer", peer.to_string()),
],
);
let (writer_tx, writer_rx) = mpsc::channel::<Vec<u8>>(1024); let (writer_tx, writer_rx) = mpsc::channel::<Vec<u8>>(1024);
register_opened_stream( register_opened_stream(
state, state,
@@ -2101,6 +2121,7 @@ async fn handle_stream_open(
) )
.await .await
{ {
eprintln!("file service stream {stream_id} failed: {err:#}");
let _ = send_file_response_to_client( let _ = send_file_response_to_client(
&state, &state,
&socket, &socket,
@@ -2452,6 +2473,14 @@ async fn handle_stream_eof(
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?; let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?; let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
let eof: StreamEof = protocol::from_body(&body)?; let eof: StreamEof = protocol::from_body(&body)?;
dosh::trace::event(
"server.stream_eof_received",
&[
("session", session_name.clone()),
("stream", eof.stream_id.to_string()),
("peer", peer.to_string()),
],
);
let mut locked = state.lock().expect("server state poisoned"); let mut locked = state.lock().expect("server state poisoned");
let Some(session) = locked.sessions.get_mut(&session_name) else { let Some(session) = locked.sessions.get_mut(&session_name) else {
return Ok(()); return Ok(());
@@ -2483,6 +2512,14 @@ async fn handle_stream_close(
let (key, session_name) = find_client_decrypt_key(state, &packet.header)?; let (key, session_name) = find_client_decrypt_key(state, &packet.header)?;
let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?; let body = protocol::decrypt_body(packet, &key, CLIENT_TO_SERVER)?;
let close: StreamClose = protocol::from_body(&body)?; let close: StreamClose = protocol::from_body(&body)?;
dosh::trace::event(
"server.stream_close_received",
&[
("session", session_name.clone()),
("stream", close.stream_id.to_string()),
("peer", peer.to_string()),
],
);
let mut locked = state.lock().expect("server state poisoned"); let mut locked = state.lock().expect("server state poisoned");
let Some(session) = locked.sessions.get_mut(&session_name) else { let Some(session) = locked.sessions.get_mut(&session_name) else {
return Ok(()); return Ok(());
@@ -2875,6 +2912,7 @@ async fn run_file_stream_service(
) )
.await .await
{ {
eprintln!("file service request on stream {stream_id} failed: {err:#}");
send_file_response_to_client( send_file_response_to_client(
&state, &state,
&socket, &socket,
@@ -2888,6 +2926,10 @@ async fn run_file_stream_service(
} }
} }
} }
dosh::trace::event(
"server.file_stream_input_closed",
&[("stream", stream_id.to_string())],
);
Ok(()) Ok(())
} }
@@ -3266,15 +3308,25 @@ async fn run_exec_stream_service(
continue; continue;
} }
}; };
let output = TokioCommand::new("sh") let shell = {
.arg("-lc") state
.lock()
.expect("server state poisoned")
.config
.shell
.clone()
};
let output = TokioCommand::new(&shell)
.arg("-c")
.arg(&request.command) .arg(&request.command)
.stdin(Stdio::null()) .stdin(Stdio::null())
.stdout(Stdio::piped()) .stdout(Stdio::piped())
.stderr(Stdio::piped()) .stderr(Stdio::piped())
.output() .output()
.await .await
.with_context(|| format!("run command {:?}", request.command))?; .with_context(|| {
format!("run command {:?} with shell {shell}", request.command)
})?;
for chunk in output.stdout.chunks(CHUNK_SIZE) { for chunk in output.stdout.chunks(CHUNK_SIZE) {
send_exec_response_to_client( send_exec_response_to_client(
&state, &state,
@@ -3795,6 +3847,13 @@ async fn send_stream_close_to_client(
client_id: [u8; 16], client_id: [u8; 16],
stream_id: u64, stream_id: u64,
) -> Result<()> { ) -> Result<()> {
dosh::trace::event(
"server.stream_close_sent",
&[
("stream", stream_id.to_string()),
("client", hex_id(client_id)),
],
);
{ {
let mut locked = state.lock().expect("server state poisoned"); let mut locked = state.lock().expect("server state poisoned");
if let Some(client) = locked.client_mut(&client_id) { if let Some(client) = locked.client_mut(&client_id) {
@@ -4442,7 +4501,8 @@ fn cleanup_disconnected_clients(state: &Arc<Mutex<ServerState>>) {
.filter(|(name, session)| { .filter(|(name, session)| {
!prewarm.contains(name.as_str()) !prewarm.contains(name.as_str())
&& session.empty_since.is_some_and(|since| { && session.empty_since.is_some_and(|since| {
now.duration_since(since) >= empty_session_timeout(name, timeout) now.duration_since(since)
>= empty_session_timeout(name, timeout, session.restart_orphaned)
}) })
}) })
.map(|(name, _)| name.clone()) .map(|(name, _)| name.clone())
@@ -4462,8 +4522,12 @@ fn cleanup_disconnected_clients(state: &Arc<Mutex<ServerState>>) {
} }
} }
fn empty_session_timeout(name: &str, configured_timeout: Duration) -> Duration { fn empty_session_timeout(
if protocol::is_implicit_session_name(name) { name: &str,
configured_timeout: Duration,
restart_orphaned: bool,
) -> Duration {
if protocol::is_implicit_session_name(name) && !restart_orphaned {
configured_timeout.min(Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS)) configured_timeout.min(Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS))
} else { } else {
configured_timeout configured_timeout
@@ -4517,7 +4581,7 @@ mod tests {
#[test] #[test]
fn persists_terminal_sessions_when_enabled() { fn persists_terminal_sessions_when_enabled() {
let (pty_tx, _rx) = mpsc::unbounded_channel(); let (pty_tx, _rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let config = ServerConfig { let config = ServerConfig {
persist_sessions: true, persist_sessions: true,
prewarm_sessions: vec!["default".to_string()], prewarm_sessions: vec!["default".to_string()],
@@ -4531,7 +4595,7 @@ mod tests {
#[test] #[test]
fn persist_disabled_never_persists() { fn persist_disabled_never_persists() {
let (pty_tx, _rx) = mpsc::unbounded_channel(); let (pty_tx, _rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let config = ServerConfig { let config = ServerConfig {
persist_sessions: false, persist_sessions: false,
prewarm_sessions: vec!["default".to_string()], prewarm_sessions: vec!["default".to_string()],
@@ -4608,6 +4672,7 @@ mod tests {
empty_since: Some(Instant::now()), empty_since: Some(Instant::now()),
holder_control: None, holder_control: None,
persistent: true, persistent: true,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 7, last_persisted_seq: 7,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -4635,6 +4700,7 @@ mod tests {
empty_since: Some(Instant::now()), empty_since: Some(Instant::now()),
holder_control: None, holder_control: None,
persistent: true, persistent: true,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 7, last_persisted_seq: 7,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -4655,7 +4721,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn no_client_pty_output_before_first_reattach_keeps_restored_screen() { async fn no_client_pty_output_before_first_reattach_keeps_restored_screen() {
let restored = b"\x1b[?1049lRESTORED_AFTER_RESTART".to_vec(); let restored = b"\x1b[?1049lRESTORED_AFTER_RESTART".to_vec();
let (pty_tx, _rx) = mpsc::unbounded_channel(); let (pty_tx, _rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let state = Arc::new(Mutex::new(ServerState::new( let state = Arc::new(Mutex::new(ServerState::new(
ServerConfig { ServerConfig {
persist_sessions: true, persist_sessions: true,
@@ -4682,6 +4748,7 @@ mod tests {
empty_since: Some(Instant::now()), empty_since: Some(Instant::now()),
holder_control: None, holder_control: None,
persistent: true, persistent: true,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 7, last_persisted_seq: 7,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -4733,7 +4800,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn unknown_resume_reject_keeps_client_id() { async fn unknown_resume_reject_keeps_client_id() {
let (pty_tx, _rx) = mpsc::unbounded_channel(); let (pty_tx, _rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let state = Arc::new(Mutex::new(ServerState::new( let state = Arc::new(Mutex::new(ServerState::new(
ServerConfig::default(), ServerConfig::default(),
[0u8; 32], [0u8; 32],
@@ -4957,7 +5024,7 @@ mod tests {
#[test] #[test]
fn client_index_stays_in_sync_with_session_clients() { fn client_index_stays_in_sync_with_session_clients() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
state state
.ensure_session("work", 80, 24, "forward-only", &[]) .ensure_session("work", 80, 24, "forward-only", &[])
@@ -5004,7 +5071,7 @@ mod tests {
#[test] #[test]
fn cleanup_purges_timed_out_clients_from_index() { fn cleanup_purges_timed_out_clients_from_index() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let config = ServerConfig { let config = ServerConfig {
client_timeout_secs: 1, client_timeout_secs: 1,
..ServerConfig::default() ..ServerConfig::default()
@@ -5031,28 +5098,79 @@ mod tests {
} }
#[test] #[test]
fn implicit_empty_session_timeout_is_bounded_for_update_reconnect() { fn implicit_empty_session_timeout_preserves_restart_orphans() {
let configured = Duration::from_secs(2_592_000); let configured = Duration::from_secs(2_592_000);
let implicit = protocol::generate_implicit_session_name(); let implicit = protocol::generate_implicit_session_name();
assert_eq!( assert_eq!(
empty_session_timeout(&implicit, configured), empty_session_timeout(&implicit, configured, false),
Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS) Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS)
); );
assert_eq!( assert_eq!(
empty_session_timeout("work", configured), empty_session_timeout(&implicit, configured, true),
configured,
"a sleeping client must retain its restart-surviving shell"
);
assert_eq!(
empty_session_timeout("work", configured, false),
configured, configured,
"named sessions keep the normal long timeout" "named sessions keep the normal long timeout"
); );
assert_eq!( assert_eq!(
empty_session_timeout(&implicit, Duration::from_secs(1)), empty_session_timeout(&implicit, Duration::from_secs(1), false),
Duration::from_secs(1), Duration::from_secs(1),
"tests/admins can still configure a shorter timeout" "tests/admins can still configure a shorter timeout"
); );
} }
#[test]
fn cleanup_keeps_restart_orphan_then_reaps_after_reattach_disconnect() {
let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(
ServerConfig {
client_timeout_secs: 2_592_000,
prewarm_sessions: Vec::new(),
..ServerConfig::default()
},
[0u8; 32],
pty_tx,
);
let implicit = protocol::generate_implicit_session_name();
state
.ensure_session(&implicit, 80, 24, "forward-only", &[])
.unwrap();
{
let session = state.sessions.get_mut(&implicit).unwrap();
session.empty_since = Some(
Instant::now()
- Duration::from_secs(IMPLICIT_EMPTY_SESSION_GRACE_SECS + 10),
);
session.restart_orphaned = true;
}
let state = Arc::new(Mutex::new(state));
cleanup_disconnected_clients(&state);
assert!(
state.lock().unwrap().sessions.contains_key(&implicit),
"restart orphan was reaped before the reconnect timeout"
);
state
.lock()
.unwrap()
.sessions
.get_mut(&implicit)
.unwrap()
.restart_orphaned = false;
cleanup_disconnected_clients(&state);
assert!(
!state.lock().unwrap().sessions.contains_key(&implicit),
"ordinary abandoned implicit session was not reaped"
);
}
#[tokio::test] #[tokio::test]
async fn forged_plaintext_detach_does_not_remove_client() { async fn forged_plaintext_detach_does_not_remove_client() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
state state
.ensure_session("work", 80, 24, "forward-only", &[]) .ensure_session("work", 80, 24, "forward-only", &[])
@@ -5078,7 +5196,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn authenticated_detach_removes_client() { async fn authenticated_detach_removes_client() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
state state
.ensure_session("work", 80, 24, "forward-only", &[]) .ensure_session("work", 80, 24, "forward-only", &[])
@@ -5109,7 +5227,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn duplicate_stream_open_for_open_stream_resends_ok() { async fn duplicate_stream_open_for_open_stream_resends_ok() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [11u8; 16]; let client_id = [11u8; 16];
let session_key = [12u8; 32]; let session_key = [12u8; 32];
@@ -5131,6 +5249,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5173,7 +5292,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn retired_stream_open_is_rejected_not_reopened() { async fn retired_stream_open_is_rejected_not_reopened() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [15u8; 16]; let client_id = [15u8; 16];
let session_key = [16u8; 32]; let session_key = [16u8; 32];
@@ -5196,6 +5315,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5243,7 +5363,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_open_is_retransmitted() { async fn pending_server_stream_open_is_retransmitted() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [13u8; 16]; let client_id = [13u8; 16];
let session_key = [14u8; 32]; let session_key = [14u8; 32];
@@ -5273,6 +5393,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5298,7 +5419,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn server_stream_retransmit_uses_observed_rtt() { async fn server_stream_retransmit_uses_observed_rtt() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [15u8; 16]; let client_id = [15u8; 16];
let session_key = [16u8; 32]; let session_key = [16u8; 32];
@@ -5332,6 +5453,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5357,7 +5479,7 @@ mod tests {
#[test] #[test]
fn forward_only_session_does_not_allocate_pty() { fn forward_only_session_does_not_allocate_pty() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
state state
@@ -5369,7 +5491,7 @@ mod tests {
#[test] #[test]
fn cleanup_reaps_abandoned_sessions_but_keeps_prewarmed() { fn cleanup_reaps_abandoned_sessions_but_keeps_prewarmed() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let config = ServerConfig { let config = ServerConfig {
client_timeout_secs: 1, client_timeout_secs: 1,
prewarm_sessions: vec!["default".to_string()], prewarm_sessions: vec!["default".to_string()],
@@ -5454,7 +5576,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn stream_data_waits_for_open_and_credit() { async fn stream_data_waits_for_open_and_credit() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [7u8; 16]; let client_id = [7u8; 16];
let session_key = [9u8; 32]; let session_key = [9u8; 32];
@@ -5509,6 +5631,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5558,7 +5681,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn stream_eof_to_client_does_not_retire_stream() { async fn stream_eof_to_client_does_not_retire_stream() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [13u8; 16]; let client_id = [13u8; 16];
let session_key = [14u8; 32]; let session_key = [14u8; 32];
@@ -5613,6 +5736,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5647,7 +5771,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_eof_is_retransmitted() { async fn pending_server_stream_eof_is_retransmitted() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [17u8; 16]; let client_id = [17u8; 16];
let session_key = [18u8; 32]; let session_key = [18u8; 32];
@@ -5678,6 +5802,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5707,7 +5832,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_close_is_retransmitted() { async fn pending_server_stream_close_is_retransmitted() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [19u8; 16]; let client_id = [19u8; 16];
let session_key = [20u8; 32]; let session_key = [20u8; 32];
@@ -5735,6 +5860,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5763,7 +5889,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_close_expires_after_attempt_cap() { async fn pending_server_stream_close_expires_after_attempt_cap() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [21u8; 16]; let client_id = [21u8; 16];
let session_key = [22u8; 32]; let session_key = [22u8; 32];
@@ -5791,6 +5917,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5808,7 +5935,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_window_adjust_is_retransmitted() { async fn pending_server_stream_window_adjust_is_retransmitted() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [23u8; 16]; let client_id = [23u8; 16];
let session_key = [24u8; 32]; let session_key = [24u8; 32];
@@ -5838,6 +5965,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5871,7 +5999,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn pending_server_stream_window_adjust_expires_after_attempt_cap() { async fn pending_server_stream_window_adjust_expires_after_attempt_cap() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [25u8; 16]; let client_id = [25u8; 16];
let session_key = [26u8; 32]; let session_key = [26u8; 32];
@@ -5901,6 +6029,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -5922,7 +6051,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn server_stream_send_splits_large_writes() { async fn server_stream_send_splits_large_writes() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [11u8; 16]; let client_id = [11u8; 16];
let session_key = [12u8; 32]; let session_key = [12u8; 32];
@@ -5977,6 +6106,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -6029,7 +6159,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn blocked_stream_data_does_not_block_terminal_frames() { async fn blocked_stream_data_does_not_block_terminal_frames() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx); let mut state = ServerState::new(ServerConfig::default(), [0u8; 32], pty_tx);
let client_id = [8u8; 16]; let client_id = [8u8; 16];
let session_key = [10u8; 32]; let session_key = [10u8; 32];
@@ -6084,6 +6214,7 @@ mod tests {
empty_since: None, empty_since: None,
holder_control: None, holder_control: None,
persistent: false, persistent: false,
restart_orphaned: false,
bytes_since_persist: 0, bytes_since_persist: 0,
last_persisted_seq: 0, last_persisted_seq: 0,
last_screen_persist_at: Instant::now(), last_screen_persist_at: Instant::now(),
@@ -6134,7 +6265,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn live_terminal_output_is_never_replaced_by_paced_snapshots() { async fn live_terminal_output_is_never_replaced_by_paced_snapshots() {
let (pty_tx, _pty_rx) = mpsc::unbounded_channel(); let (pty_tx, _pty_rx) = mpsc::channel(PTY_OUTPUT_QUEUE_CAPACITY);
let config = ServerConfig { let config = ServerConfig {
output_frame_interval_ms: 1000, output_frame_interval_ms: 1000,
..ServerConfig::default() ..ServerConfig::default()
+35 -26
View File
@@ -201,7 +201,6 @@ impl DoshClientBuilder {
match connect_sdk_peer( match connect_sdk_peer(
peer_addr, peer_addr,
&self.client.config, &self.client.config,
&host_config,
&self.host, &self.host,
&raw_server, &raw_server,
ssh_port, ssh_port,
@@ -244,7 +243,6 @@ impl ConnectedDoshClient {
async fn connect_sdk_peer( async fn connect_sdk_peer(
peer_addr: SocketAddr, peer_addr: SocketAddr,
config: &ClientConfig, config: &ClientConfig,
host_config: &HostConfig,
host: &str, host: &str,
server: &str, server: &str,
ssh_port: Option<u16>, ssh_port: Option<u16>,
@@ -311,18 +309,17 @@ async fn connect_sdk_peer(
&hello, &hello,
&server_hello.hello, &server_hello.hello,
)?; )?;
let auth = sign_auth( let auth = sign_auth(NativeAuthSignRequest {
config, config,
host_config, hello: &hello,
&hello, server_hello: &server_hello.hello,
&server_hello.hello, requested_forwardings: requested_forwardings.to_vec(),
requested_forwardings.to_vec(), explicit_identity_files: identity_files.to_vec(),
identity_files.to_vec(),
use_ssh_agent, use_ssh_agent,
server, server,
ssh_port, ssh_port,
ssh_config, ssh_config,
)?; })?;
let mut pending_id = [0u8; 16]; let mut pending_id = [0u8; 16];
pending_id.copy_from_slice(&server_hello.hello.auth_challenge[..16]); pending_id.copy_from_slice(&server_hello.hello.auth_challenge[..16]);
let auth_packet = protocol::encode_encrypted( let auth_packet = protocol::encode_encrypted(
@@ -387,18 +384,30 @@ fn verify_or_trust_host(
} }
} }
fn sign_auth( struct NativeAuthSignRequest<'a> {
config: &ClientConfig, config: &'a ClientConfig,
_host_config: &HostConfig, hello: &'a NativeClientHello,
hello: &NativeClientHello, server_hello: &'a native::NativeServerHello,
server_hello: &native::NativeServerHello,
requested_forwardings: Vec<ForwardingRequest>, requested_forwardings: Vec<ForwardingRequest>,
explicit_identity_files: Vec<PathBuf>, explicit_identity_files: Vec<PathBuf>,
use_ssh_agent: Option<bool>, use_ssh_agent: Option<bool>,
server: &str, server: &'a str,
ssh_port: Option<u16>, ssh_port: Option<u16>,
ssh_config: &SdkSshConfig, ssh_config: &'a SdkSshConfig,
) -> Result<native::NativeUserAuth> { }
fn sign_auth(request: NativeAuthSignRequest<'_>) -> Result<native::NativeUserAuth> {
let NativeAuthSignRequest {
config,
hello,
server_hello,
requested_forwardings,
explicit_identity_files,
use_ssh_agent,
server,
ssh_port,
ssh_config,
} = request;
let use_agent = use_ssh_agent.unwrap_or(config.use_ssh_agent); let use_agent = use_ssh_agent.unwrap_or(config.use_ssh_agent);
let mut errors = Vec::new(); let mut errors = Vec::new();
if use_agent && !ssh_config.identities_only { if use_agent && !ssh_config.identities_only {
@@ -448,12 +457,6 @@ fn sdk_identity_paths(
for path in explicit_identity_files { for path in explicit_identity_files {
push_identity_path(&mut paths, path); push_identity_path(&mut paths, path);
} }
for path in &ssh_config.identity_files {
push_identity_path(
&mut paths,
expand_tilde(&expand_ssh_path_tokens(path, token_context)),
);
}
if !ssh_config.identities_only { if !ssh_config.identities_only {
for path in &config.identity_files { for path in &config.identity_files {
push_identity_path( push_identity_path(
@@ -462,6 +465,12 @@ fn sdk_identity_paths(
); );
} }
} }
for path in &ssh_config.identity_files {
push_identity_path(
&mut paths,
expand_tilde(&expand_ssh_path_tokens(path, token_context)),
);
}
paths paths
} }
@@ -818,10 +827,10 @@ fn local_username() -> Option<String> {
local_username_from_env(|name| std::env::var(name).ok()) local_username_from_env(|name| std::env::var(name).ok())
} }
fn local_username_from_env(mut get: impl FnMut(&str) -> Option<String>) -> Option<String> { fn local_username_from_env(get: impl FnMut(&str) -> Option<String>) -> Option<String> {
["USER", "USERNAME"] ["USER", "USERNAME"]
.into_iter() .into_iter()
.filter_map(|name| get(name)) .filter_map(get)
.find(|value| !value.is_empty()) .find(|value| !value.is_empty())
} }
@@ -1092,7 +1101,7 @@ mod tests {
let paths = sdk_identity_paths(&config, Vec::new(), &ssh_config, &token_context); let paths = sdk_identity_paths(&config, Vec::new(), &ssh_config, &token_context);
assert_eq!(paths, vec![dir.path().join("ssh"), config_identity]); assert_eq!(paths, vec![config_identity, dir.path().join("ssh")]);
} }
#[test] #[test]
+41 -10
View File
@@ -12,6 +12,10 @@ use tokio::sync::mpsc;
// TUIs often write several KiB on the first draw; sending that as one UDP // TUIs often write several KiB on the first draw; sending that as one UDP
// datagram can fragment and vanish, leaving only a blank alternate screen. // datagram can fragment and vanish, leaving only a blank alternate screen.
const PTY_OUTPUT_CHUNK_BYTES: usize = 1024; const PTY_OUTPUT_CHUNK_BYTES: usize = 1024;
/// Shared server queue bound. At the current chunk size this limits queued PTY
/// output to roughly 1 MiB before the producing shell receives normal PTY
/// backpressure.
pub const PTY_OUTPUT_QUEUE_CAPACITY: usize = 1024;
/// Backing for a PTY master held by the server. /// Backing for a PTY master held by the server.
/// ///
@@ -126,7 +130,7 @@ pub fn spawn_pty_session(
cols: u16, cols: u16,
rows: u16, rows: u16,
env: &[(String, String)], env: &[(String, String)],
tx: mpsc::UnboundedSender<PtyOutput>, tx: mpsc::Sender<PtyOutput>,
) -> Result<PtyHandle> { ) -> Result<PtyHandle> {
let pty_system = NativePtySystem::default(); let pty_system = NativePtySystem::default();
let pair = pty_system let pair = pty_system
@@ -232,7 +236,7 @@ pub fn build_shell_command(shell: &str, env: &[(String, String)]) -> CommandBuil
pub fn adopt_pty_from_fd( pub fn adopt_pty_from_fd(
session: String, session: String,
master_fd: RawFd, master_fd: RawFd,
tx: mpsc::UnboundedSender<PtyOutput>, tx: mpsc::Sender<PtyOutput>,
) -> Result<PtyHandle> { ) -> Result<PtyHandle> {
// Take ownership of the fd. A clone gives us an independent reader so the // Take ownership of the fd. A clone gives us an independent reader so the
// reader thread and the writer/resize side hold separate `File`s and don't // reader thread and the writer/resize side hold separate `File`s and don't
@@ -252,7 +256,7 @@ pub fn adopt_pty_from_fd(
fn spawn_reader_thread( fn spawn_reader_thread(
session: String, session: String,
mut reader: Box<dyn Read + Send>, mut reader: Box<dyn Read + Send>,
tx: mpsc::UnboundedSender<PtyOutput>, tx: mpsc::Sender<PtyOutput>,
) -> Result<()> { ) -> Result<()> {
let reader_session = session.clone(); let reader_session = session.clone();
thread::Builder::new() thread::Builder::new()
@@ -262,7 +266,7 @@ fn spawn_reader_thread(
loop { loop {
match reader.read(&mut buf) { match reader.read(&mut buf) {
Ok(0) => { Ok(0) => {
let _ = tx.send(PtyOutput { let _ = tx.blocking_send(PtyOutput {
session: reader_session.clone(), session: reader_session.clone(),
bytes: Vec::new(), bytes: Vec::new(),
exited: true, exited: true,
@@ -271,15 +275,20 @@ fn spawn_reader_thread(
} }
Ok(n) => { Ok(n) => {
for chunk in buf[..n].chunks(PTY_OUTPUT_CHUNK_BYTES) { for chunk in buf[..n].chunks(PTY_OUTPUT_CHUNK_BYTES) {
let _ = tx.send(PtyOutput { if tx
session: reader_session.clone(), .blocking_send(PtyOutput {
bytes: chunk.to_vec(), session: reader_session.clone(),
exited: false, bytes: chunk.to_vec(),
}); exited: false,
})
.is_err()
{
return;
}
} }
} }
Err(_) => { Err(_) => {
let _ = tx.send(PtyOutput { let _ = tx.blocking_send(PtyOutput {
session: reader_session.clone(), session: reader_session.clone(),
bytes: Vec::new(), bytes: Vec::new(),
exited: true, exited: true,
@@ -298,6 +307,28 @@ mod tests {
use super::*; use super::*;
const _: () = assert!(PTY_OUTPUT_CHUNK_BYTES <= 1200); const _: () = assert!(PTY_OUTPUT_CHUNK_BYTES <= 1200);
const _: () = assert!(PTY_OUTPUT_QUEUE_CAPACITY * PTY_OUTPUT_CHUNK_BYTES <= 1024 * 1024);
#[test]
fn pty_output_queue_capacity_is_memory_bounded() {
let (tx, _rx) = mpsc::channel::<PtyOutput>(PTY_OUTPUT_QUEUE_CAPACITY);
for index in 0..PTY_OUTPUT_QUEUE_CAPACITY {
tx.try_send(PtyOutput {
session: "load".to_string(),
bytes: vec![index as u8; PTY_OUTPUT_CHUNK_BYTES],
exited: false,
})
.unwrap();
}
assert!(matches!(
tx.try_send(PtyOutput {
session: "load".to_string(),
bytes: vec![0; PTY_OUTPUT_CHUNK_BYTES],
exited: false,
}),
Err(mpsc::error::TrySendError::Full(_))
));
}
#[test] #[test]
fn terminfo_available_detects_known_and_unknown() { fn terminfo_available_detects_known_and_unknown() {
+425 -30
View File
@@ -16,8 +16,8 @@ use crate::transport::{
use crate::udp::{is_transient_udp_error, is_transient_udp_send_error}; use crate::udp::{is_transient_udp_error, is_transient_udp_send_error};
use anyhow::{Context, Result, anyhow, bail}; use anyhow::{Context, Result, anyhow, bail};
use ed25519_dalek::SigningKey; use ed25519_dalek::SigningKey;
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet, VecDeque};
use std::net::SocketAddr; use std::net::{IpAddr, SocketAddr};
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
@@ -31,10 +31,13 @@ pub struct DoshServerConfig {
pub transport: TransportConfig, pub transport: TransportConfig,
pub require_current_user: bool, pub require_current_user: bool,
pub auth_timeout: Duration, pub auth_timeout: Duration,
pub max_pending_auth: usize,
pub connection_timeout: Duration,
} }
impl DoshServerConfig { impl DoshServerConfig {
pub fn new(server: ServerConfig) -> Self { pub fn new(server: ServerConfig) -> Self {
let connection_timeout = Duration::from_secs(server.client_timeout_secs.max(1));
Self { Self {
server, server,
bind_addr: None, bind_addr: None,
@@ -42,6 +45,8 @@ impl DoshServerConfig {
transport: TransportConfig::default(), transport: TransportConfig::default(),
require_current_user: true, require_current_user: true,
auth_timeout: Duration::from_secs(30), auth_timeout: Duration::from_secs(30),
max_pending_auth: 1024,
connection_timeout,
} }
} }
@@ -71,6 +76,21 @@ impl DoshServerConfig {
self.transport = transport; self.transport = transport;
self self
} }
pub fn connection_timeout(mut self, timeout: Duration) -> Self {
self.connection_timeout = timeout.max(ADAPTIVE_RETRANSMIT_MIN);
self
}
pub fn auth_timeout(mut self, timeout: Duration) -> Self {
self.auth_timeout = timeout.max(ADAPTIVE_RETRANSMIT_MIN);
self
}
pub fn max_pending_auth(mut self, max_pending: usize) -> Self {
self.max_pending_auth = max_pending.max(1);
self
}
} }
impl Default for DoshServerConfig { impl Default for DoshServerConfig {
@@ -91,6 +111,7 @@ pub struct DoshAccepted {
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub enum DoshServerEvent { pub enum DoshServerEvent {
Accepted(DoshAccepted), Accepted(DoshAccepted),
Disconnected(DoshAccepted),
Session { Session {
conn_id: [u8; 16], conn_id: [u8; 16],
event: SessionEvent, event: SessionEvent,
@@ -107,13 +128,78 @@ struct PendingServerAuth {
created_at: Instant, created_at: Instant,
} }
struct AuthRateLimiter {
per_minute: u32,
max_sources: usize,
buckets: HashMap<IpAddr, AuthTokenBucket>,
}
#[derive(Clone, Copy)]
struct AuthTokenBucket {
tokens: f64,
last_refill: Instant,
}
impl AuthRateLimiter {
fn new(per_minute: u32, max_sources: usize) -> Self {
Self {
per_minute,
max_sources: max_sources.max(1),
buckets: HashMap::new(),
}
}
fn check(&mut self, ip: IpAddr, now: Instant) -> Result<u32, ()> {
if self.per_minute == 0 {
return Err(());
}
self.evict_full(now);
if !self.buckets.contains_key(&ip) && self.buckets.len() >= self.max_sources {
return Err(());
}
let capacity = self.per_minute as f64;
let refill_per_sec = capacity / 60.0;
let bucket = self.buckets.entry(ip).or_insert(AuthTokenBucket {
tokens: capacity,
last_refill: now,
});
let elapsed = now
.saturating_duration_since(bucket.last_refill)
.as_secs_f64();
bucket.tokens = (bucket.tokens + elapsed * refill_per_sec).min(capacity);
bucket.last_refill = now;
if bucket.tokens < 1.0 {
return Err(());
}
bucket.tokens -= 1.0;
Ok(bucket.tokens as u32)
}
fn evict_full(&mut self, now: Instant) {
if self.per_minute == 0 {
self.buckets.clear();
return;
}
let capacity = self.per_minute as f64;
let refill_per_sec = capacity / 60.0;
self.buckets.retain(|_, bucket| {
let elapsed = now
.saturating_duration_since(bucket.last_refill)
.as_secs_f64();
(bucket.tokens + elapsed * refill_per_sec) < capacity
});
}
}
pub struct DoshServer { pub struct DoshServer {
socket: Arc<UdpSocket>, socket: Arc<UdpSocket>,
config: DoshServerConfig, config: DoshServerConfig,
host_signing: SigningKey, host_signing: SigningKey,
auth_limiter: AuthRateLimiter,
pending: HashMap<[u8; 16], PendingServerAuth>, pending: HashMap<[u8; 16], PendingServerAuth>,
transports: HashMap<[u8; 16], DoshTransport>, transports: HashMap<[u8; 16], DoshTransport>,
accepted: HashMap<[u8; 16], DoshAccepted>, accepted: HashMap<[u8; 16], DoshAccepted>,
disconnected: VecDeque<DoshAccepted>,
} }
impl DoshServer { impl DoshServer {
@@ -139,13 +225,19 @@ impl DoshServer {
}; };
let host_signing = load_or_create_host_key(&config.server)?; let host_signing = load_or_create_host_key(&config.server)?;
let socket = Arc::new(UdpSocket::bind(bind_addr).await?); let socket = Arc::new(UdpSocket::bind(bind_addr).await?);
let auth_limiter = AuthRateLimiter::new(
config.server.native_auth_rate_limit_per_minute,
config.max_pending_auth,
);
Ok(Self { Ok(Self {
socket, socket,
config, config,
host_signing, host_signing,
auth_limiter,
pending: HashMap::new(), pending: HashMap::new(),
transports: HashMap::new(), transports: HashMap::new(),
accepted: HashMap::new(), accepted: HashMap::new(),
disconnected: VecDeque::new(),
}) })
} }
@@ -165,10 +257,19 @@ impl DoshServer {
self.transports.get_mut(conn_id) self.transports.get_mut(conn_id)
} }
pub fn remove_connection(&mut self, conn_id: &[u8; 16]) -> Option<DoshAccepted> {
self.transports.remove(conn_id);
self.accepted.remove(conn_id)
}
pub async fn recv(&mut self) -> Result<DoshServerEvent> { pub async fn recv(&mut self) -> Result<DoshServerEvent> {
let mut buf = vec![0u8; 65535]; let mut buf = vec![0u8; 65535];
loop { loop {
self.expire_pending(); self.expire_pending();
self.expire_connections();
if let Some(connection) = self.disconnected.pop_front() {
return Ok(DoshServerEvent::Disconnected(connection));
}
let (n, peer) = match tokio::time::timeout( let (n, peer) = match tokio::time::timeout(
ADAPTIVE_RETRANSMIT_MIN, ADAPTIVE_RETRANSMIT_MIN,
self.socket.recv_from(&mut buf), self.socket.recv_from(&mut buf),
@@ -207,6 +308,9 @@ impl DoshServer {
} }
if let Some(transport) = self.transports.get_mut(&packet.header.conn_id) { if let Some(transport) = self.transports.get_mut(&packet.header.conn_id) {
let event = transport.handle_datagram(datagram, peer).await?; let event = transport.handle_datagram(datagram, peer).await?;
if let Some(accepted) = self.accepted.get_mut(&packet.header.conn_id) {
accepted.peer_addr = transport.peer_addr();
}
return Ok(Some(DoshServerEvent::Session { return Ok(Some(DoshServerEvent::Session {
conn_id: packet.header.conn_id, conn_id: packet.header.conn_id,
event, event,
@@ -270,7 +374,21 @@ impl DoshServer {
self.send_reject(peer, [0u8; 16], &err.to_string()).await?; self.send_reject(peer, [0u8; 16], &err.to_string()).await?;
return Ok(()); return Ok(());
} }
let result = self.build_server_hello(req.hello, peer); self.expire_pending();
if self.pending.len() >= self.config.max_pending_auth {
self.send_reject(peer, [0u8; 16], "native auth server busy")
.await?;
return Ok(());
}
let rate_limit_remaining = match self.auth_limiter.check(peer.ip(), Instant::now()) {
Ok(remaining) => remaining,
Err(()) => {
self.send_reject(peer, [0u8; 16], "native auth rate limit exceeded")
.await?;
return Ok(());
}
};
let result = self.build_server_hello(req.hello, peer, Some(rate_limit_remaining));
let (pending_id, hello) = match result { let (pending_id, hello) = match result {
Ok(value) => value, Ok(value) => value,
Err(err) => { Err(err) => {
@@ -288,6 +406,7 @@ impl DoshServer {
&mut self, &mut self,
client: native::NativeClientHello, client: native::NativeClientHello,
peer: SocketAddr, peer: SocketAddr,
rate_limit_remaining: Option<u32>,
) -> Result<([u8; 16], NativeServerHello)> { ) -> Result<([u8; 16], NativeServerHello)> {
if !self.config.server.native_auth { if !self.config.server.native_auth {
bail!("native auth disabled"); bail!("native auth disabled");
@@ -307,11 +426,14 @@ impl DoshServer {
bail!("native auth requires a supported user key algorithm"); bail!("native auth requires a supported user key algorithm");
} }
if self.config.require_current_user { if self.config.require_current_user {
let current_user = std::env::var("USER").unwrap_or_else(|_| "unknown".to_string()); let current_user = local_username();
if client.requested_user != current_user { if client.requested_user != current_user {
bail!("native auth user mismatch"); bail!("native auth user mismatch");
} }
} }
if self.pending.len() >= self.config.max_pending_auth {
bail!("native auth server busy");
}
let (server_secret, server_public) = generate_native_ephemeral(); let (server_secret, server_public) = generate_native_ephemeral();
let mut server = NativeServerHello { let mut server = NativeServerHello {
@@ -322,7 +444,7 @@ impl DoshServer {
chosen_aead: "chacha20poly1305".to_string(), chosen_aead: "chacha20poly1305".to_string(),
server_key_epoch: 1, server_key_epoch: 1,
auth_challenge: crypto::random_32(), auth_challenge: crypto::random_32(),
rate_limit_remaining: None, rate_limit_remaining,
host_signature: Vec::new(), host_signature: Vec::new(),
}; };
sign_server_hello(&self.host_signing, &client, &mut server)?; sign_server_hello(&self.host_signing, &client, &mut server)?;
@@ -477,7 +599,37 @@ impl DoshServer {
let timeout = self.config.auth_timeout; let timeout = self.config.auth_timeout;
self.pending self.pending
.retain(|_, pending| pending.created_at.elapsed() <= timeout); .retain(|_, pending| pending.created_at.elapsed() <= timeout);
self.auth_limiter.evict_full(Instant::now());
} }
fn expire_connections(&mut self) {
let timeout = self.config.connection_timeout;
let expired = self
.transports
.iter()
.filter_map(|(conn_id, transport)| {
(transport.stale_for() > timeout).then_some(*conn_id)
})
.collect::<Vec<_>>();
for conn_id in expired {
self.transports.remove(&conn_id);
if let Some(connection) = self.accepted.remove(&conn_id) {
self.disconnected.push_back(connection);
}
}
}
}
fn local_username() -> String {
local_username_from_env(|name| std::env::var(name).ok())
}
fn local_username_from_env(get: impl FnMut(&str) -> Option<String>) -> String {
["USER", "USERNAME"]
.into_iter()
.filter_map(get)
.find(|value| !value.is_empty())
.unwrap_or_else(|| "unknown".to_string())
} }
async fn send_udp(socket: &UdpSocket, packet: &[u8], peer: SocketAddr) -> Result<bool> { async fn send_udp(socket: &UdpSocket, packet: &[u8], peer: SocketAddr) -> Result<bool> {
@@ -530,6 +682,24 @@ mod tests {
use crate::transport::TransportEvent; use crate::transport::TransportEvent;
use ed25519_dalek::SigningKey; use ed25519_dalek::SigningKey;
fn native_client_hello(public: [u8; 32]) -> native::NativeClientHello {
native::NativeClientHello {
protocol_version: native::NATIVE_PROTOCOL_VERSION,
client_random: crypto::random_32(),
client_ephemeral_public: public,
requested_host: "127.0.0.1".to_string(),
requested_user: "sdk-user".to_string(),
requested_session: "test".to_string(),
requested_mode: "forward-only".to_string(),
terminal_size: (80, 24),
supported_aead: vec!["chacha20poly1305".to_string()],
supported_user_key_algorithms: vec!["ssh-ed25519".to_string()],
cached_host_key_fingerprint: None,
attach_ticket_envelope: None,
requested_env: Vec::new(),
}
}
#[tokio::test] #[tokio::test]
async fn sdk_client_and_server_exchange_service_stream() { async fn sdk_client_and_server_exchange_service_stream() {
let dir = tempfile::tempdir().unwrap(); let dir = tempfile::tempdir().unwrap();
@@ -578,20 +748,43 @@ mod tests {
.user("sdk-user") .user("sdk-user")
.service("echo") .service("echo")
.connect(); .connect();
let accept = async { let handshake = async {
tokio::pin!(connect);
let mut connected = None;
let mut accepted = None;
loop { loop {
if let DoshServerEvent::Accepted(accepted) = server.recv().await.unwrap() { tokio::select! {
break accepted; result = &mut connect, if connected.is_none() => {
connected = Some(result?);
}
event = server.recv(), if accepted.is_none() => {
if let DoshServerEvent::Accepted(connection) = event? {
accepted = Some(connection);
}
}
}
if connected.is_some() && accepted.is_some() {
return Ok::<_, anyhow::Error>((
connected.take().expect("connected result checked"),
accepted.take().expect("accepted result checked"),
));
} }
} }
}; };
let (connected, accepted) = tokio::join!(connect, accept); let (connected, accepted) = tokio::time::timeout(Duration::from_secs(3), handshake)
let mut client_transport = connected.unwrap().into_transport(); .await
.expect("SDK client/server authentication timed out")
.unwrap();
let mut client_transport = connected.into_transport();
let conn_id = accepted.conn_id; let conn_id = accepted.conn_id;
assert_eq!(accepted.services, vec!["echo".to_string()]); assert_eq!(accepted.services, vec!["echo".to_string()]);
let stream_id = client_transport.open_service("echo").await.unwrap(); let stream_id = client_transport.open_service("echo").await.unwrap();
match server.recv().await.unwrap() { match tokio::time::timeout(Duration::from_secs(3), server.recv())
.await
.expect("server timed out waiting for stream open")
.unwrap()
{
DoshServerEvent::Session { DoshServerEvent::Session {
event: SessionEvent::Stream(TransportEvent::Open(open)), event: SessionEvent::Stream(TransportEvent::Open(open)),
.. ..
@@ -602,7 +795,10 @@ mod tests {
other => panic!("unexpected event {other:?}"), other => panic!("unexpected event {other:?}"),
} }
assert!(matches!( assert!(matches!(
client_transport.recv().await.unwrap(), tokio::time::timeout(Duration::from_secs(3), client_transport.recv())
.await
.expect("client timed out waiting for stream-open acknowledgement")
.unwrap(),
SessionEvent::Stream(TransportEvent::OpenOk { .. }) SessionEvent::Stream(TransportEvent::OpenOk { .. })
)); ));
@@ -611,7 +807,11 @@ mod tests {
.await .await
.unwrap(); .unwrap();
loop { loop {
match server.recv().await.unwrap() { match tokio::time::timeout(Duration::from_secs(3), server.recv())
.await
.expect("server timed out waiting for stream data")
.unwrap()
{
DoshServerEvent::Session { DoshServerEvent::Session {
event: SessionEvent::Stream(TransportEvent::Data(data)), event: SessionEvent::Stream(TransportEvent::Data(data)),
.. ..
@@ -628,7 +828,11 @@ mod tests {
} }
} }
loop { loop {
match client_transport.recv().await.unwrap() { match tokio::time::timeout(Duration::from_secs(3), client_transport.recv())
.await
.expect("client timed out waiting for stream response")
.unwrap()
{
SessionEvent::Stream(TransportEvent::Data(data)) => { SessionEvent::Stream(TransportEvent::Data(data)) => {
assert_eq!(data.chunks, vec![b"pong".to_vec()]); assert_eq!(data.chunks, vec![b"pong".to_vec()]);
break; break;
@@ -718,6 +922,209 @@ mod tests {
assert_eq!(open.target_host, "@dosh-echo"); assert_eq!(open.target_host, "@dosh-echo");
} }
#[tokio::test]
async fn sdk_server_expires_disconnected_transport_and_reports_it() {
let dir = tempfile::tempdir().unwrap();
let server_config = ServerConfig {
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
..ServerConfig::default()
};
let server_config = DoshServerConfig::new(server_config)
.bind_addr("127.0.0.1:0".parse().unwrap())
.require_current_user(false)
.connection_timeout(Duration::from_millis(30));
let mut server = DoshServer::bind(server_config).await.unwrap();
let peer_socket = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let peer_addr = peer_socket.local_addr().unwrap();
let conn_id = [81u8; 16];
let transport = DoshTransport::new(
Arc::clone(&server.socket),
SessionTransportConfig {
role: SessionRole::Server,
conn_id,
session_key: [82u8; 32],
peer_addr,
initial_send_seq: 1,
initial_ack: 0,
stream: server.config.transport.clone(),
},
);
let accepted = DoshAccepted {
conn_id,
user: "sdk-user".to_string(),
session: "mobile".to_string(),
services: vec!["echo".to_string()],
peer_addr,
};
server.transports.insert(conn_id, transport);
server.accepted.insert(conn_id, accepted.clone());
tokio::time::sleep(Duration::from_millis(40)).await;
let event = tokio::time::timeout(Duration::from_secs(1), server.recv())
.await
.expect("server did not report expired connection")
.unwrap();
match event {
DoshServerEvent::Disconnected(disconnected) => {
assert_eq!(disconnected.conn_id, conn_id);
assert_eq!(disconnected.session, "mobile");
}
other => panic!("unexpected expiry event {other:?}"),
}
assert!(server.connection(&conn_id).is_none());
assert!(server.transport(&conn_id).is_none());
}
#[tokio::test]
async fn sdk_server_updates_connection_metadata_after_authenticated_roam() {
let dir = tempfile::tempdir().unwrap();
let server_config = ServerConfig {
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
..ServerConfig::default()
};
let server_config = DoshServerConfig::new(server_config)
.bind_addr("127.0.0.1:0".parse().unwrap())
.require_current_user(false);
let mut server = DoshServer::bind(server_config).await.unwrap();
let original = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let roaming = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let original_addr = original.local_addr().unwrap();
let roaming_addr = roaming.local_addr().unwrap();
let conn_id = [83u8; 16];
let session_key = [84u8; 32];
let transport = DoshTransport::new(
Arc::clone(&server.socket),
SessionTransportConfig {
role: SessionRole::Server,
conn_id,
session_key,
peer_addr: original_addr,
initial_send_seq: 1,
initial_ack: 0,
stream: server.config.transport.clone(),
},
);
server.transports.insert(conn_id, transport);
server.accepted.insert(
conn_id,
DoshAccepted {
conn_id,
user: "sdk-user".to_string(),
session: "mobile".to_string(),
services: Vec::new(),
peer_addr: original_addr,
},
);
let ping = protocol::encode_encrypted(
PacketKind::Ping,
conn_id,
1,
0,
&session_key,
CLIENT_TO_SERVER,
b"",
)
.unwrap();
roaming
.send_to(&ping, server.local_addr().unwrap())
.await
.unwrap();
let event = tokio::time::timeout(Duration::from_secs(1), server.recv())
.await
.expect("server did not receive roaming ping")
.unwrap();
assert!(matches!(
event,
DoshServerEvent::Session {
event: SessionEvent::Ping,
..
}
));
assert_eq!(server.connection(&conn_id).unwrap().peer_addr, roaming_addr);
}
#[tokio::test]
async fn sdk_server_bounds_and_rate_limits_pending_authentication() {
let dir = tempfile::tempdir().unwrap();
let server_config = ServerConfig {
host_key: dir.path().join("host_key").to_string_lossy().to_string(),
native_auth_rate_limit_per_minute: 1,
..ServerConfig::default()
};
let server_config = DoshServerConfig::new(server_config)
.bind_addr("127.0.0.1:0".parse().unwrap())
.require_current_user(false)
.max_pending_auth(1);
let mut server = DoshServer::bind(server_config).await.unwrap();
let client = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let peer = client.local_addr().unwrap();
let (_, public) = native::generate_native_ephemeral();
let body = protocol::to_body(&NativeClientHelloBody {
hello: native_client_hello(public),
})
.unwrap();
server
.handle_client_hello(peer, body.clone())
.await
.unwrap();
let mut packet = [0u8; 65535];
let (n, _) = client.recv_from(&mut packet).await.unwrap();
let response = protocol::decode(&packet[..n]).unwrap();
assert_eq!(response.header.kind, PacketKind::NativeServerHello);
let hello: NativeServerHelloBody = protocol::from_body(&response.body).unwrap();
assert_eq!(hello.hello.rate_limit_remaining, Some(0));
assert_eq!(server.pending.len(), 1);
server
.handle_client_hello(peer, body.clone())
.await
.unwrap();
let (n, _) = client.recv_from(&mut packet).await.unwrap();
let response = protocol::decode(&packet[..n]).unwrap();
assert_eq!(response.header.kind, PacketKind::AttachReject);
let reject: AttachReject = protocol::from_body(&response.body).unwrap();
assert_eq!(reject.reason, "native auth server busy");
assert_eq!(server.pending.len(), 1);
server.pending.clear();
server.handle_client_hello(peer, body).await.unwrap();
let (n, _) = client.recv_from(&mut packet).await.unwrap();
let response = protocol::decode(&packet[..n]).unwrap();
assert_eq!(response.header.kind, PacketKind::AttachReject);
let reject: AttachReject = protocol::from_body(&response.body).unwrap();
assert_eq!(reject.reason, "native auth rate limit exceeded");
assert!(server.pending.is_empty());
}
#[test]
fn sdk_auth_rate_limiter_bounds_source_tracking_and_refills() {
let now = Instant::now();
let mut limiter = AuthRateLimiter::new(2, 1);
let first: IpAddr = "192.0.2.1".parse().unwrap();
let second: IpAddr = "192.0.2.2".parse().unwrap();
assert_eq!(limiter.check(first, now), Ok(1));
assert_eq!(limiter.check(first, now), Ok(0));
assert_eq!(limiter.check(first, now), Err(()));
assert_eq!(limiter.check(second, now), Err(()));
assert_eq!(limiter.check(second, now + Duration::from_secs(60)), Ok(1));
assert_eq!(limiter.buckets.len(), 1);
}
#[test]
fn sdk_server_username_supports_windows_environment() {
assert_eq!(
local_username_from_env(|name| match name {
"USER" => None,
"USERNAME" => Some("palav-win".to_string()),
_ => None,
}),
"palav-win"
);
}
#[tokio::test] #[tokio::test]
async fn bad_native_auth_does_not_consume_pending_challenge() { async fn bad_native_auth_does_not_consume_pending_challenge() {
let dir = tempfile::tempdir().unwrap(); let dir = tempfile::tempdir().unwrap();
@@ -744,22 +1151,10 @@ mod tests {
let mut server = DoshServer::bind(server_config).await.unwrap(); let mut server = DoshServer::bind(server_config).await.unwrap();
let peer: SocketAddr = "127.0.0.1:9".parse().unwrap(); let peer: SocketAddr = "127.0.0.1:9".parse().unwrap();
let (client_secret, client_public) = native::generate_native_ephemeral(); let (client_secret, client_public) = native::generate_native_ephemeral();
let hello = native::NativeClientHello { let hello = native_client_hello(client_public);
protocol_version: native::NATIVE_PROTOCOL_VERSION, let (pending_id, server_hello) = server
client_random: crypto::random_32(), .build_server_hello(hello.clone(), peer, None)
client_ephemeral_public: client_public, .unwrap();
requested_host: "127.0.0.1".to_string(),
requested_user: "sdk-user".to_string(),
requested_session: "test".to_string(),
requested_mode: "forward-only".to_string(),
terminal_size: (80, 24),
supported_aead: vec!["chacha20poly1305".to_string()],
supported_user_key_algorithms: vec!["ssh-ed25519".to_string()],
cached_host_key_fingerprint: None,
attach_ticket_envelope: None,
requested_env: Vec::new(),
};
let (pending_id, server_hello) = server.build_server_hello(hello.clone(), peer).unwrap();
let session_key = native::derive_native_session_key( let session_key = native::derive_native_session_key(
&client_secret, &client_secret,
server_hello.server_ephemeral_public, server_hello.server_ephemeral_public,
+33 -2
View File
@@ -31,7 +31,10 @@ use std::sync::Arc;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use tokio::net::UdpSocket; use tokio::net::UdpSocket;
pub const DEFAULT_INITIAL_WINDOW: usize = 1024 * 1024; /// Per-stream bytes allowed in flight before window credit returns. This bounds
/// each sender burst to 64 MTU-safe packets so bulk streams cannot starve the
/// terminal while retaining useful bandwidth on high-latency links.
pub const DEFAULT_INITIAL_WINDOW: usize = 64 * 1024;
pub const DEFAULT_RETRANSMIT_AFTER: Duration = Duration::from_millis(200); pub const DEFAULT_RETRANSMIT_AFTER: Duration = Duration::from_millis(200);
pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10); pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
pub const ADAPTIVE_RETRANSMIT_MIN: Duration = Duration::from_millis(10); pub const ADAPTIVE_RETRANSMIT_MIN: Duration = Duration::from_millis(10);
@@ -39,7 +42,10 @@ pub const DEFAULT_KEEPALIVE_AFTER: Duration = Duration::from_secs(2);
pub const DEFAULT_RETIRED_STREAM_TOMBSTONES: usize = 16 * 1024; pub const DEFAULT_RETIRED_STREAM_TOMBSTONES: usize = 16 * 1024;
pub const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8; pub const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8;
pub const SERVICE_TARGET_PREFIX: &str = "@dosh-"; pub const SERVICE_TARGET_PREFIX: &str = "@dosh-";
pub const MAX_STREAM_DATA_BYTES: usize = 60 * 1024; /// Maximum application bytes in one encrypted UDP stream packet. Keeping this
/// aligned with terminal output framing avoids IP fragmentation and stays
/// below macOS route MTUs after Dosh, AEAD, UDP, and IP overhead.
pub const MAX_STREAM_DATA_BYTES: usize = 1024;
#[derive(Debug, Clone, PartialEq, Eq)] #[derive(Debug, Clone, PartialEq, Eq)]
pub struct TransportConfig { pub struct TransportConfig {
@@ -1675,6 +1681,31 @@ mod tests {
assert_eq!(third.bytes.len(), 13); assert_eq!(third.bytes.len(), 13);
} }
#[test]
fn maximum_stream_chunk_fits_a_safe_udp_datagram() {
let body = protocol::to_body(&StreamData {
stream_id: u64::MAX,
offset: u64::MAX,
bytes: vec![0xff; MAX_STREAM_DATA_BYTES],
})
.unwrap();
let packet = protocol::encode_encrypted(
PacketKind::StreamData,
[0xff; 16],
u64::MAX,
u64::MAX,
&[0xff; 32],
CLIENT_TO_SERVER,
&body,
)
.unwrap();
assert!(
packet.len() <= 1200,
"encrypted stream datagram is {} bytes",
packet.len()
);
}
#[test] #[test]
fn queued_large_write_flushes_in_window_sized_chunks_after_open() { fn queued_large_write_flushes_in_window_sized_chunks_after_open() {
let mut mux = StreamMux::new(TransportConfig { let mut mux = StreamMux::new(TransportConfig {
File diff suppressed because it is too large Load Diff
+265
View File
@@ -0,0 +1,265 @@
#![cfg(unix)]
use sha2::{Digest, Sha256};
use std::fs;
use std::os::unix::fs::PermissionsExt;
use std::path::{Path, PathBuf};
use std::process::Command;
use tempfile::TempDir;
const VERSION: &str = "9.8.7-rc3";
fn write_executable(path: &Path, body: &str) {
fs::write(path, body).unwrap();
let mut permissions = fs::metadata(path).unwrap().permissions();
permissions.set_mode(0o755);
fs::set_permissions(path, permissions).unwrap();
}
struct InstallerFixture {
root: TempDir,
fake_bin: PathBuf,
archive: PathBuf,
checksum: PathBuf,
curl_log: PathBuf,
git_log: PathBuf,
version: String,
}
impl InstallerFixture {
fn new(version: &str) -> Self {
let root = tempfile::tempdir().unwrap();
let fake_bin = root.path().join("fake-bin");
let stage = root.path().join("stage/dosh");
fs::create_dir_all(fake_bin.as_path()).unwrap();
fs::create_dir_all(stage.join("bin")).unwrap();
fs::write(stage.join("VERSION"), format!("{version}\n")).unwrap();
write_executable(
&stage.join("bin/dosh-client"),
"#!/bin/sh\nprintf 'fixture dosh\\n'\n",
);
let archive = root.path().join("dosh-macos-aarch64.tar.gz");
let status = Command::new("tar")
.args(["-C", root.path().join("stage").to_str().unwrap(), "-czf"])
.arg(&archive)
.arg("dosh")
.status()
.unwrap();
assert!(status.success());
let digest = Sha256::digest(fs::read(&archive).unwrap());
let checksum = root.path().join("dosh-macos-aarch64.tar.gz.sha256");
fs::write(
&checksum,
format!("{digest:x} dosh-macos-aarch64.tar.gz\n"),
)
.unwrap();
write_executable(
&fake_bin.join("uname"),
"#!/bin/sh\ncase \"${1:-}\" in -s) echo Darwin;; -m) echo arm64;; *) echo Darwin;; esac\n",
);
let curl_log = root.path().join("curl.log");
let git_log = root.path().join("git.log");
write_executable(
&fake_bin.join("curl"),
r#"#!/bin/sh
printf '%s\n' "$*" >>"$TEST_CURL_LOG"
out=
previous=
url=
for arg in "$@"; do
if [ "$previous" = "-o" ]; then out="$arg"; fi
previous="$arg"
case "$arg" in http://*|https://*) url="$arg";; esac
done
case "$url" in
*/raw/branch/main/Cargo.toml)
printf '[package]\nname = "dosh"\nversion = "%s"\n' "$TEST_RELEASE_VERSION"
;;
*.sha256)
cp "$TEST_ARCHIVE_CHECKSUM" "$out"
;;
*/releases/download/*)
[ "${TEST_FAIL_ARCHIVE:-0}" = 1 ] && exit 22
cp "$TEST_ARCHIVE" "$out"
;;
*) exit 22;;
esac
"#,
);
Self {
root,
fake_bin,
archive,
checksum,
curl_log,
git_log,
version: version.to_string(),
}
}
fn command(&self, cwd: &Path) -> Command {
let mut command = Command::new("sh");
command
.arg(format!("{}/install.sh", env!("CARGO_MANIFEST_DIR")))
.args([
"client",
"--repo",
"https://example.invalid/Palav/dosh.git",
"--prefix",
])
.arg(self.root.path().join("prefix"))
.current_dir(cwd)
.env("HOME", self.root.path().join("home"))
.env(
"PATH",
format!(
"{}:{}",
self.fake_bin.display(),
std::env::var("PATH").unwrap_or_default()
),
)
.env("DOSH_BINARY_REQUIRED", "1")
.env("TEST_RELEASE_VERSION", &self.version)
.env("TEST_ARCHIVE", &self.archive)
.env("TEST_ARCHIVE_CHECKSUM", &self.checksum)
.env("TEST_CURL_LOG", &self.curl_log)
.env("TEST_GIT_LOG", &self.git_log);
command
}
fn installed_client(&self) -> PathBuf {
self.root.path().join("prefix/bin/dosh-client")
}
}
#[test]
fn unix_installer_uses_repository_version_even_inside_an_unrelated_rust_project() {
let fixture = InstallerFixture::new(VERSION);
let unrelated = fixture.root.path().join("unrelated");
fs::create_dir_all(&unrelated).unwrap();
fs::write(
unrelated.join("Cargo.toml"),
"[package]\nname = \"not-dosh\"\nversion = \"99.0.0\"\n",
)
.unwrap();
let output = fixture.command(&unrelated).output().unwrap();
assert!(
output.status.success(),
"installer failed: {}",
String::from_utf8_lossy(&output.stderr)
);
assert!(fixture.installed_client().is_file());
let requests = fs::read_to_string(&fixture.curl_log).unwrap();
assert!(requests.contains("/raw/branch/main/Cargo.toml"));
assert!(requests.contains(&format!(
"/releases/download/v{VERSION}/dosh-macos-aarch64.tar.gz"
)));
assert!(!requests.contains("/releases/latest"));
assert!(!requests.contains("99.0.0"));
}
#[test]
fn unix_installer_source_fallback_checks_out_the_same_release_tag() {
let fixture = InstallerFixture::new(VERSION);
write_executable(
&fixture.fake_bin.join("git"),
r#"#!/bin/sh
printf '%s\n' "$*" >>"$TEST_GIT_LOG"
if [ "${1:-}" = clone ]; then
for destination in "$@"; do :; done
mkdir -p "$destination/.git"
fi
"#,
);
write_executable(
&fixture.fake_bin.join("cargo"),
"#!/bin/sh\nmkdir -p target/release\nprintf '#!/bin/sh\\n' >target/release/dosh-client\nchmod +x target/release/dosh-client\n",
);
let cwd = fixture.root.path().join("plain");
fs::create_dir_all(&cwd).unwrap();
let mut command = fixture.command(&cwd);
command
.env("DOSH_BINARY_REQUIRED", "0")
.env("TEST_FAIL_ARCHIVE", "1");
let output = command.output().unwrap();
assert!(
output.status.success(),
"installer failed: {}",
String::from_utf8_lossy(&output.stderr)
);
assert!(fixture.installed_client().is_file());
let git = fs::read_to_string(&fixture.git_log).unwrap();
assert!(git.contains(&format!("clone --depth 1 --branch v{VERSION}")));
assert!(!git.contains("--branch main"));
}
#[test]
fn unix_installer_resolves_a_stable_repository_version_without_latest_redirects() {
let version = "1.2.3";
let fixture = InstallerFixture::new(version);
let cwd = fixture.root.path().join("plain");
fs::create_dir_all(&cwd).unwrap();
let output = fixture.command(&cwd).output().unwrap();
assert!(
output.status.success(),
"installer failed: {}",
String::from_utf8_lossy(&output.stderr)
);
let requests = fs::read_to_string(&fixture.curl_log).unwrap();
assert!(requests.contains(&format!(
"/releases/download/v{version}/dosh-macos-aarch64.tar.gz"
)));
assert!(!requests.contains("/releases/latest"));
}
#[test]
fn unix_installer_repairs_a_legacy_updater_that_requests_its_old_tag() {
let fixture = InstallerFixture::new(VERSION);
let cwd = fixture.root.path().join("plain");
fs::create_dir_all(&cwd).unwrap();
let mut command = fixture.command(&cwd);
command.env("DOSH_BINARY_VERSION", "v0.1.17");
let output = command.output().unwrap();
assert!(
output.status.success(),
"installer failed: {}",
String::from_utf8_lossy(&output.stderr)
);
let requests = fs::read_to_string(&fixture.curl_log).unwrap();
assert!(requests.contains(&format!(
"/releases/download/v{VERSION}/dosh-macos-aarch64.tar.gz"
)));
assert!(!requests.contains("/releases/download/v0.1.17/"));
}
#[test]
fn unix_installer_honors_an_explicit_exact_release_pin() {
let pinned = "0.1.17";
let fixture = InstallerFixture::new(pinned);
let cwd = fixture.root.path().join("plain");
fs::create_dir_all(&cwd).unwrap();
let mut command = fixture.command(&cwd);
command
.env("TEST_RELEASE_VERSION", VERSION)
.env("DOSH_BINARY_VERSION", format!("v{pinned}"))
.env("DOSH_BINARY_EXACT", "1");
let output = command.output().unwrap();
assert!(
output.status.success(),
"installer failed: {}",
String::from_utf8_lossy(&output.stderr)
);
let requests = fs::read_to_string(&fixture.curl_log).unwrap();
assert!(requests.contains(&format!(
"/releases/download/v{pinned}/dosh-macos-aarch64.tar.gz"
)));
assert!(!requests.contains(&format!("/releases/download/v{VERSION}/")));
}
+35
View File
@@ -1027,6 +1027,8 @@ fn native_file_copy_recursive_round_trip() {
fs::create_dir_all(src.join("nested")).unwrap(); fs::create_dir_all(src.join("nested")).unwrap();
fs::write(src.join("root.txt"), b"root file\n").unwrap(); fs::write(src.join("root.txt"), b"root file\n").unwrap();
fs::write(src.join("nested/child.txt"), b"child file\n").unwrap(); fs::write(src.join("nested/child.txt"), b"child file\n").unwrap();
let large_payload: Vec<u8> = (0..128 * 1024).map(|index| (index % 251) as u8).collect();
fs::write(src.join("large.bin"), &large_payload).unwrap();
std::os::unix::fs::symlink("root.txt", src.join("root-link")).unwrap(); std::os::unix::fs::symlink("root.txt", src.join("root-link")).unwrap();
std::os::unix::fs::symlink("nested/child.txt", src.join("child-link")).unwrap(); std::os::unix::fs::symlink("nested/child.txt", src.join("child-link")).unwrap();
@@ -1134,6 +1136,10 @@ fn native_file_copy_recursive_round_trip() {
fs::read_to_string(downloaded.join("nested/child.txt")).unwrap(), fs::read_to_string(downloaded.join("nested/child.txt")).unwrap(),
"child file\n" "child file\n"
); );
assert_eq!(
fs::read(downloaded.join("large.bin")).unwrap(),
large_payload
);
assert_eq!( assert_eq!(
fs::read_link(downloaded.join("root-link")).unwrap(), fs::read_link(downloaded.join("root-link")).unwrap(),
PathBuf::from("root.txt") PathBuf::from("root.txt")
@@ -1355,6 +1361,22 @@ fn native_exec_command_smoke() {
let dir = tempfile::tempdir().unwrap(); let dir = tempfile::tempdir().unwrap();
let port = free_udp_port(); let port = free_udp_port();
let config = write_server_config(&dir, port); let config = write_server_config(&dir, port);
let shell_log = dir.path().join("exec-shell.log");
let shell = dir.path().join("exec-shell");
fs::write(
&shell,
format!(
"#!/bin/sh\nprintf '%s\\n' \"$*\" >> '{}'\nexec /bin/sh \"$@\"\n",
shell_log.display()
),
)
.unwrap();
fs::set_permissions(&shell, fs::Permissions::from_mode(0o700)).unwrap();
let raw = fs::read_to_string(&config).unwrap().replace(
"shell = \"/bin/sh\"",
&format!("shell = {:?}", shell.display().to_string()),
);
fs::write(&config, raw).unwrap();
write_native_client_auth(&dir, &config); write_native_client_auth(&dir, &config);
let mut server = start_server(&dir, &config); let mut server = start_server(&dir, &config);
let client_bin = env!("CARGO_BIN_EXE_dosh-client"); let client_bin = env!("CARGO_BIN_EXE_dosh-client");
@@ -1380,6 +1402,19 @@ fn native_exec_command_smoke() {
); );
assert_eq!(String::from_utf8_lossy(&output.stdout), "out"); assert_eq!(String::from_utf8_lossy(&output.stdout), "out");
assert_eq!(String::from_utf8_lossy(&output.stderr), "err"); assert_eq!(String::from_utf8_lossy(&output.stderr), "err");
let shell_invocations = fs::read_to_string(shell_log).unwrap();
assert!(
shell_invocations
.lines()
.any(|line| line.starts_with("-c ")),
"configured shell was not used for exec: {shell_invocations:?}"
);
assert!(
!shell_invocations
.lines()
.any(|line| line.starts_with("-lc ")),
"exec must not start a login shell: {shell_invocations:?}"
);
} }
#[test] #[test]
+26 -23
View File
@@ -94,6 +94,8 @@ fn release_gates_test_all_targets() {
assert!(windows_client.contains("runs-on: windows-latest")); assert!(windows_client.contains("runs-on: windows-latest"));
assert!(windows_client.contains("name: Test Windows client")); assert!(windows_client.contains("name: Test Windows client"));
assert!(windows_client.contains("name: PowerShell syntax check")); assert!(windows_client.contains("name: PowerShell syntax check"));
assert!(windows_client.contains("name: Windows installer runtime test"));
assert!(windows_client.contains("run: scripts/test-install-windows.ps1"));
assert!(windows_client.contains("shell: pwsh")); assert!(windows_client.contains("shell: pwsh"));
assert!( assert!(
windows_client windows_client
@@ -138,28 +140,24 @@ fn one_dot_zero_gates_default_to_full_length_soaks() {
} }
#[test] #[test]
fn installers_skip_stale_latest_release_prebuilts() { fn installers_resolve_the_repository_declared_release_without_stable_redirects() {
let install = include_str!("../install.sh"); let install = include_str!("../install.sh");
assert!(install.contains("latest_release_is_stale"));
assert!(install.contains("current_source_version")); assert!(install.contains("current_source_version"));
assert!(install.contains("latest release $latest is older than source $current")); assert!(install.contains("release_source_ref"));
assert!(install.contains("version_prerelease")); assert!(install.contains("printf 'v%s\\n' \"${current#v}\""));
assert!(install.contains("prerelease_less_than")); assert!(install.contains("DOSH_BINARY_EXACT=1"));
assert!( assert!(!install.contains("/releases/latest"));
install.contains("if latest_release_is_stale; then"), assert!(install.contains("resolved_source_ref=\"$(release_source_ref || true)\""));
"unix installer must skip stale latest before downloading a prebuilt" assert!(install.contains("checkout -q --detach FETCH_HEAD"));
);
let ps1 = include_str!("../install.ps1"); let ps1 = include_str!("../install.ps1");
assert!(ps1.contains("Latest-ReleaseIsStale"));
assert!(ps1.contains("Current-SourceVersion")); assert!(ps1.contains("Current-SourceVersion"));
assert!(ps1.contains("latest release $latestVersion is older than source $current")); assert!(ps1.contains("Source-ReleaseRef"));
assert!(ps1.contains("Version-Prerelease")); assert!(ps1.contains("\"v$($current.TrimStart('v'))\""));
assert!(ps1.contains("Compare-Prerelease")); assert!(ps1.contains("DOSH_BINARY_EXACT"));
assert!( assert!(!ps1.contains("/releases/latest"));
ps1.contains("if (Latest-ReleaseIsStale)"), assert!(ps1.contains("$script:ResolvedSourceRef = Source-ReleaseRef"));
"windows installer must skip stale latest before downloading a prebuilt" assert!(ps1.contains("checkout -q --detach FETCH_HEAD"));
);
} }
#[test] #[test]
@@ -198,15 +196,16 @@ fn windows_installer_reuses_persistent_source_update_cache() {
assert!(ps1.contains("$full -eq $localAppData -or $full -eq $localAppDataDosh")); assert!(ps1.contains("$full -eq $localAppData -or $full -eq $localAppDataDosh"));
assert!(ps1.contains("$sourceCache = Assert-SafeUpdateCache $UpdateCache")); assert!(ps1.contains("$sourceCache = Assert-SafeUpdateCache $UpdateCache"));
assert!(ps1.contains( assert!(ps1.contains(
"$fetchArgs = @(\"-C\", $sourceCache, \"fetch\", \"--depth\", \"1\", \"origin\", \"main\")" "$fetchArgs = @(\"-C\", $sourceCache, \"fetch\", \"--depth\", \"1\", \"origin\", $sourceRef)"
)); ));
assert!(ps1.contains("$fetchArgs = @(\"-C\", $sourceCache, \"fetch\", \"-q\", \"--depth\", \"1\", \"origin\", \"main\")")); assert!(ps1.contains("$fetchArgs = @(\"-C\", $sourceCache, \"fetch\", \"-q\", \"--depth\", \"1\", \"origin\", $sourceRef)"));
assert!(ps1.contains("git @fetchArgs")); assert!(ps1.contains("git @fetchArgs"));
assert!(ps1.contains("git -C $sourceCache checkout -q -B main FETCH_HEAD")); assert!(ps1.contains("git -C $sourceCache checkout -q -B main FETCH_HEAD"));
assert!(ps1.contains("git -C $sourceCache checkout -q --detach FETCH_HEAD"));
assert!(ps1.contains( assert!(ps1.contains(
"$cloneArgs = @(\"clone\", \"--depth\", \"1\", \"--branch\", \"main\", $Repo, $sourceCache)" "$cloneArgs = @(\"clone\", \"--depth\", \"1\", \"--branch\", $sourceRef, $Repo, $sourceCache)"
)); ));
assert!(ps1.contains("$cloneArgs = @(\"clone\", \"-q\", \"--depth\", \"1\", \"--branch\", \"main\", $Repo, $sourceCache)")); assert!(ps1.contains("$cloneArgs = @(\"clone\", \"-q\", \"--depth\", \"1\", \"--branch\", $sourceRef, $Repo, $sourceCache)"));
assert!(ps1.contains("git @cloneArgs | Out-Null")); assert!(ps1.contains("git @cloneArgs | Out-Null"));
assert!( assert!(
!ps1.contains("git clone --depth 1 $Repo $tmp"), !ps1.contains("git clone --depth 1 $Repo $tmp"),
@@ -218,8 +217,12 @@ fn windows_installer_reuses_persistent_source_update_cache() {
fn windows_quiet_source_updates_silence_git_like_unix() { fn windows_quiet_source_updates_silence_git_like_unix() {
let install = include_str!("../install.sh"); let install = include_str!("../install.sh");
let ps1 = include_str!("../install.ps1"); let ps1 = include_str!("../install.ps1");
assert!(install.contains("git -C \"$update_cache\" fetch -q --depth 1 origin main")); assert!(install.contains("git -C \"$update_cache\" fetch -q --depth 1 origin \"$source_ref\""));
assert!(install.contains("git clone -q --depth 1 --branch main \"$repo\" \"$update_cache\"")); assert!(
install.contains(
"git clone -q --depth 1 --branch \"$source_ref\" \"$repo\" \"$update_cache\""
)
);
assert!(ps1.contains("if ($Quiet)")); assert!(ps1.contains("if ($Quiet)"));
assert!(ps1.contains("\"fetch\", \"-q\", \"--depth\"")); assert!(ps1.contains("\"fetch\", \"-q\", \"--depth\""));
assert!(ps1.contains("\"clone\", \"-q\", \"--depth\"")); assert!(ps1.contains("\"clone\", \"-q\", \"--depth\""));