diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c3e04ee..e1506bb 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -32,3 +32,16 @@ jobs: - name: Build release binary run: cargo build --locked --release --workspace + + nixos: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v6 + + - uses: cachix/install-nix-action@v31 + + - name: Check Nix formatting + run: nix fmt -- --check flake.nix packaging/nix/*.nix + + - name: Check and build Nix packages + run: nix flake check --print-build-logs diff --git a/.github/workflows/release-please.yml b/.github/workflows/release-please.yml index e046913..246a2ca 100644 --- a/.github/workflows/release-please.yml +++ b/.github/workflows/release-please.yml @@ -35,6 +35,13 @@ jobs: fetch-depth: 0 ref: ${{ steps.release.outputs.tag_name }} + - uses: cachix/install-nix-action@v31 + if: ${{ steps.release.outputs.release_created == 'true' }} + + - name: Verify release Nix package + if: ${{ steps.release.outputs.release_created == 'true' }} + run: nix build .#breakd .#breakd-relay --print-build-logs + - name: Build source archive id: source if: ${{ steps.release.outputs.release_created == 'true' }} diff --git a/Cargo.lock b/Cargo.lock index eaca199..8a708c4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -47,7 +47,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -58,7 +58,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -128,6 +128,15 @@ version = "2.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4388bee8683e3d04af747c73422af53102d2bd24d9eadb6cbc100baef4b43f8" +[[package]] +name = "block-buffer" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d2f6c7dbe95a6ed67ad9f18e57daf93a2f034c524b99fd2b76d18fdfeb6660aa" +dependencies = [ + "hybrid-array", +] + [[package]] name = "breakd" version = "0.1.9" @@ -150,6 +159,7 @@ dependencies = [ name = "breakd-config" version = "0.1.9" dependencies = [ + "breakd-coop", "breakd-core", "nix", "tempfile", @@ -157,6 +167,16 @@ dependencies = [ "toml 0.9.12+spec-1.1.0", ] +[[package]] +name = "breakd-coop" +version = "0.1.9" +dependencies = [ + "breakd-core", + "serde", + "thiserror", + "uuid", +] + [[package]] name = "breakd-core" version = "0.1.9" @@ -173,13 +193,16 @@ version = "0.1.9" dependencies = [ "anyhow", "breakd-config", + "breakd-coop", "breakd-core", "breakd-ipc", "breakd-platform-linux", "breakd-scheduler", "breakd-tray", + "futures-util", "serde_json", "tokio", + "tokio-tungstenite", "tracing", "uuid", ] @@ -221,16 +244,35 @@ dependencies = [ "zbus", ] +[[package]] +name = "breakd-relay" +version = "0.1.9" +dependencies = [ + "anyhow", + "breakd-coop", + "breakd-core", + "clap", + "futures-util", + "serde_json", + "tokio", + "tokio-tungstenite", + "tracing", + "tracing-subscriber", + "uuid", +] + [[package]] name = "breakd-scheduler" version = "0.1.9" dependencies = [ "breakd-config", + "breakd-coop", "breakd-core", "proptest", "serde", "serde_json", "thiserror", + "uuid", ] [[package]] @@ -333,6 +375,17 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" +[[package]] +name = "chacha20" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" +dependencies = [ + "cfg-if", + "cpufeatures", + "rand_core 0.10.1", +] + [[package]] name = "clap" version = "4.6.1" @@ -388,12 +441,53 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "const-oid" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6ef517f0926dd24a1582492c791b6a4818a4d94e789a334894aa15b0d12f55c" + +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + [[package]] name = "crossbeam-utils" version = "0.8.22" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "61803da095bee82a81bb1a452ecc25d3b2f1416d1897eb86430c6159ef717c17" +[[package]] +name = "crypto-common" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce6e4c961d6cd6c9a86db418387425e8bdeaf05b3c8bc1411e6dca4c252f1453" +dependencies = [ + "hybrid-array", +] + +[[package]] +name = "data-encoding" +version = "2.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4ae5f15dda3c708c0ade84bfee31ccab44a3da4f88015ed22f63732abe300c8" + +[[package]] +name = "digest" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" +dependencies = [ + "block-buffer", + "const-oid", + "crypto-common", +] + [[package]] name = "downcast-rs" version = "1.2.1" @@ -440,7 +534,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -558,6 +652,12 @@ dependencies = [ "syn", ] +[[package]] +name = "futures-sink" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c39754e157331b013978ec91992bde1ac089843443c49cbc7f46150b0fad0893" + [[package]] name = "futures-task" version = "0.3.32" @@ -572,6 +672,7 @@ checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ "futures-core", "futures-macro", + "futures-sink", "futures-task", "pin-project-lite", "slab", @@ -635,6 +736,17 @@ dependencies = [ "system-deps", ] +[[package]] +name = "getrandom" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" +dependencies = [ + "cfg-if", + "libc", + "wasi", +] + [[package]] name = "getrandom" version = "0.3.4" @@ -656,6 +768,7 @@ dependencies = [ "cfg-if", "libc", "r-efi 6.0.0", + "rand_core 0.10.1", ] [[package]] @@ -685,7 +798,7 @@ dependencies = [ "gobject-sys", "libc", "system-deps", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -912,12 +1025,37 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" +[[package]] +name = "http" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6970f50e31d6fc17d3fa27329444bfa74e196cf62e95052a3f6fee181dba6425" +dependencies = [ + "bytes", + "itoa", +] + +[[package]] +name = "httparse" +version = "1.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" + [[package]] name = "humantime" version = "2.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "15cdd26707701c53297e2fa6afb323d55fbc1d0810c3aec078ae3ef0424c3c15" +[[package]] +name = "hybrid-array" +version = "0.4.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "818356c5132c1fede50f837ca96afbe78ff42413047f4abb886217845e1b6c8c" +dependencies = [ + "typenum", +] + [[package]] name = "indexmap" version = "2.14.0" @@ -1035,7 +1173,7 @@ checksum = "02bd0af71c67b473010cbbc60715ee815645a4dc942899111f494b4b737d6fda" dependencies = [ "libc", "wasi", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1056,7 +1194,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1197,7 +1335,7 @@ dependencies = [ "bit-vec", "bitflags", "num-traits", - "rand", + "rand 0.9.5", "rand_chacha", "rand_xorshift", "regex-syntax", @@ -1249,7 +1387,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" dependencies = [ "rand_chacha", - "rand_core", + "rand_core 0.9.5", +] + +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "chacha20", + "getrandom 0.4.3", + "rand_core 0.10.1", ] [[package]] @@ -1259,7 +1408,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.9.5", ] [[package]] @@ -1271,13 +1420,19 @@ dependencies = [ "getrandom 0.3.4", ] +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + [[package]] name = "rand_xorshift" version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "513962919efc330f829edb2535844d1b912b0fbe2ca165d613e4e8788bb05a5a" dependencies = [ - "rand_core", + "rand_core 0.9.5", ] [[package]] @@ -1306,6 +1461,20 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" +[[package]] +name = "ring" +version = "0.17.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" +dependencies = [ + "cc", + "cfg-if", + "getrandom 0.2.17", + "libc", + "untrusted", + "windows-sys 0.52.0", +] + [[package]] name = "rustc_version" version = "0.4.1" @@ -1325,7 +1494,40 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys", + "windows-sys 0.61.2", +] + +[[package]] +name = "rustls" +version = "0.23.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" +dependencies = [ + "once_cell", + "rustls-pki-types", + "rustls-webpki", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "764899a24af3980067ee14bc143654f297b22eaebfe3c7b6b211920a5a59b046" +dependencies = [ + "zeroize", +] + +[[package]] +name = "rustls-webpki" +version = "0.103.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +dependencies = [ + "ring", + "rustls-pki-types", + "untrusted", ] [[package]] @@ -1421,6 +1623,17 @@ dependencies = [ "serde_core", ] +[[package]] +name = "sha1" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "aacc4cc499359472b4abe1bf11d0b12e688af9a805fa5e3016f9a386dc2d0214" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sharded-slab" version = "0.1.7" @@ -1465,7 +1678,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" dependencies = [ "libc", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1474,6 +1687,12 @@ version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" +[[package]] +name = "subtle" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" + [[package]] name = "syn" version = "2.0.118" @@ -1514,7 +1733,7 @@ dependencies = [ "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1561,7 +1780,7 @@ dependencies = [ "socket2", "tokio-macros", "tracing", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1575,6 +1794,32 @@ dependencies = [ "syn", ] +[[package]] +name = "tokio-rustls" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" +dependencies = [ + "rustls", + "tokio", +] + +[[package]] +name = "tokio-tungstenite" +version = "0.30.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17a073bfed563fa236697a068031408a93cd9522e08abf9933ead3e73411bd71" +dependencies = [ + "futures-util", + "log", + "rustls", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tungstenite", + "webpki-roots 0.26.11", +] + [[package]] name = "toml" version = "0.9.12+spec-1.1.0" @@ -1724,6 +1969,30 @@ dependencies = [ "tracing-serde", ] +[[package]] +name = "tungstenite" +version = "0.30.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e48ac77174b19c110a50ab2128b24215ac9cb40e0e12e093fb602d175c569d22" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.10.2", + "rustls", + "rustls-pki-types", + "sha1", + "thiserror", +] + +[[package]] +name = "typenum" +version = "1.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" + [[package]] name = "uds_windows" version = "1.2.1" @@ -1732,7 +2001,7 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", - "windows-sys", + "windows-sys 0.61.2", ] [[package]] @@ -1747,6 +2016,12 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "untrusted" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" + [[package]] name = "utf8parse" version = "0.2.2" @@ -1903,6 +2178,24 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "webpki-roots" +version = "0.26.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "521bc38abb08001b01866da9f51eb7c5d647a19260e00054a8c7fd5f9e57f7a9" +dependencies = [ + "webpki-roots 1.0.8", +] + +[[package]] +name = "webpki-roots" +version = "1.0.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf85cb06032201fa7c6f829d7db5a7e5aa45bcc0655327713065f6f0576731bf" +dependencies = [ + "rustls-pki-types", +] + [[package]] name = "winapi" version = "0.3.9" @@ -1931,6 +2224,15 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" +[[package]] +name = "windows-sys" +version = "0.52.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "282be5f36a8ce781fad8c8ae18fa3f9beff57ec1b52cb3de0789201425d9a33d" +dependencies = [ + "windows-targets", +] + [[package]] name = "windows-sys" version = "0.61.2" @@ -1940,6 +2242,70 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows-targets" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b724f72796e036ab90c1021d4780d4d3d648aca59e491e6b98e725b84e99973" +dependencies = [ + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_gnullvm", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" + +[[package]] +name = "windows_i686_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e9b5ad5ab802e97eb8e295ac6720e509ee4c243f69d781394014ebfe8bbfa0b" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" + +[[package]] +name = "windows_i686_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" + [[package]] name = "winnow" version = "0.7.15" @@ -1990,7 +2356,7 @@ dependencies = [ "tracing", "uds_windows", "uuid", - "windows-sys", + "windows-sys 0.61.2", "winnow 1.0.3", "zbus_macros", "zbus_names", @@ -2043,6 +2409,12 @@ dependencies = [ "syn", ] +[[package]] +name = "zeroize" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e" + [[package]] name = "zmij" version = "1.0.22" diff --git a/Cargo.toml b/Cargo.toml index 0c6edb4..9edd9ce 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,6 +8,7 @@ rust-version.workspace = true [workspace] members = [ "crates/core", + "crates/coop", "crates/scheduler", "crates/config", "crates/ipc", @@ -16,6 +17,7 @@ members = [ "crates/tray", "crates/wayland-overlay", "crates/daemon", + "crates/relay", ] resolver = "2" @@ -55,6 +57,7 @@ serde_json = "1.0" tempfile = "3.23" thiserror = "2.0" tokio = { version = "1.49", features = ["full"] } +tokio-tungstenite = { version = "0.30", default-features = false } toml = "0.9" tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } diff --git a/README.md b/README.md index 015ef22..a4cfb6e 100644 --- a/README.md +++ b/README.md @@ -12,6 +12,8 @@ This is why you should use my shit instead of Stretchly. xD ## Install +### Arch Linux + Install the AUR package and start the user service: ```bash @@ -22,6 +24,57 @@ breakd status `breakd` runs as your user. It does not need root access or access to `/dev/input`. +### NixOS + +Add the flake and enable its NixOS module: + +```nix +{ + inputs.breakd.url = "github:simonwinther/breakd"; + + outputs = { nixpkgs, breakd, ... }: { + nixosConfigurations.my-host = nixpkgs.lib.nixosSystem { + system = "x86_64-linux"; + modules = [ + breakd.nixosModules.default + ({ ... }: { services.breakd.enable = true; }) + ]; + }; + }; +} +``` + +Update the `breakd` input when a new release is available, then rebuild: + +```bash +nix flake update breakd +sudo nixos-rebuild switch --flake .#my-host +``` + +For a non-declarative installation, install the package directly and enable its +user service: + +```bash +nix profile install github:simonwinther/breakd +systemctl --user enable --now breakd.service +``` + +The desktop package stays separate from the optional relay package. A NixOS +machine that will host co-op rooms can enable the small system service without +installing the GTK desktop application: + +```nix +services.breakd-relay = { + enable = true; + listen = "127.0.0.1:8787"; + maxRoomSize = 8; + maxRooms = 256; +}; +``` + +Expose that private listener through a TLS reverse proxy before using it over +the internet. The package is also available directly as `.#breakd-relay`. + ## Configure `breakd` works without a configuration file. To change the defaults: @@ -119,6 +172,40 @@ enabled = true The tray shows the current schedule status and provides pause/resume, manual break, skip, postpone, reset, and reload actions. It needs a StatusNotifierItem host, such as Waybar's `tray` module. Disabling the tray does not affect the scheduler or overlays. +### Co-op breaks + +Co-op mode lets one host own the schedule while guests mirror its next break, +active break, pause state, and permitted actions. Each computer still renders +its own native overlay with its own monitor and message settings. + +Start the relay on a server (or use a relay you trust), then create a room: + +```bash +breakd coop host --relay wss://breaks.example.net/ws +``` + +Send the printed invite to your friend. They join by quoting the complete value: + +```bash +breakd coop join 'wss://breaks.example.net/ws#breakd=' +breakd coop status +``` + +The host remains authoritative: a guest's skip, postpone, pause, resume, reset, +or manual-break command is sent to the host, and the resulting host snapshot is +mirrored back. Both systems should have normal network time synchronization +enabled because scheduled starts use absolute timestamps. If snapshots stop for +10 seconds, the guest starts a fresh local schedule; reconnecting adopts the +host again. Leave at any time with `breakd coop leave`. + +The relay has no database, accounts, schedule engine, or desktop dependencies. +It retains only the latest snapshot while a host is connected. Room tokens are +sent in the WebSocket `Authorization` header rather than the request URL; the +invite keeps the token in its URL fragment. Anyone holding an invite can enter +that room, so share it as a secret and create a new room to rotate it. See +[Co-op deployment and protocol](docs/coop.md) for relay setup and operational +details. + ### Pointer and keyboard input `display.pointer_mode` controls where clicks go during a break: @@ -177,6 +264,10 @@ breakd toggle breakd reload breakd outputs [--json] breakd doctor [--json] +breakd coop host --relay +breakd coop join '' +breakd coop status [--json] +breakd coop leave breakd settings breakd example-config ``` @@ -273,6 +364,7 @@ install -Dm644 packaging/systemd/breakd-local.service \ ~/.config/systemd/user/breakd.service install -Dm644 packaging/io.github.simonwinther.breakd.settings.desktop \ ~/.local/share/applications/io.github.simonwinther.breakd.settings.desktop +install -Dm644 crates/platform-linux/assets/*.oga -t ~/.local/share/breakd install -Dm644 THIRD_PARTY_NOTICES.md \ ~/.local/share/licenses/breakd/THIRD_PARTY_NOTICES.md install -Dm600 config.example.toml ~/.config/breakd/config.toml diff --git a/config.example.toml b/config.example.toml index 2356b5b..460f3ec 100644 --- a/config.example.toml +++ b/config.example.toml @@ -189,6 +189,10 @@ submap_fallback = true # enforce blocking through a temporary Hyprland submap [tray] enabled = true +[coop] +mode = "off" # off | host | guest; use `breakd coop host` or `breakd coop join` +disconnect_grace = "10s" # resume a fresh local schedule if the host disappears + [logging] level = "info" format = "journald" diff --git a/crates/config/Cargo.toml b/crates/config/Cargo.toml index 5b57357..ebf4ecd 100644 --- a/crates/config/Cargo.toml +++ b/crates/config/Cargo.toml @@ -6,6 +6,7 @@ license.workspace = true rust-version.workspace = true [dependencies] +breakd-coop = { path = "../coop" } breakd-core = { path = "../core" } nix.workspace = true tempfile.workspace = true diff --git a/crates/config/src/lib.rs b/crates/config/src/lib.rs index 222e210..14d3dac 100644 --- a/crates/config/src/lib.rs +++ b/crates/config/src/lib.rs @@ -7,11 +7,11 @@ use std::{ use breakd_core::{ AppConfig, BreakTiming, CompletionConfig, CompletionSound, ContentConfig, ContentSelector, - DisplayConfig, DisplayMode, DurationMs, FullscreenBehavior, FullscreenConfig, HyprlandConfig, - IdleConfig, KeyboardMode, Layer, LoggingConfig, LongBreakTiming, MissedBreakPolicy, - NotificationsConfig, PointerMode, PostponeConfig, PostponeRule, RecoveryConfig, - RestBreakTiming, ScheduleConfig, SkipConfig, SkipRule, StartupConfig, StrictConfig, StrictMode, - TrayConfig, + CoopConfig, CoopMode, DisplayConfig, DisplayMode, DurationMs, FullscreenBehavior, + FullscreenConfig, HyprlandConfig, IdleConfig, KeyboardMode, Layer, LoggingConfig, + LongBreakTiming, MissedBreakPolicy, NotificationsConfig, PointerMode, PostponeConfig, + PostponeRule, RecoveryConfig, RestBreakTiming, ScheduleConfig, SkipConfig, SkipRule, + StartupConfig, StrictConfig, StrictMode, TrayConfig, }; use nix::unistd::Uid; use thiserror::Error; @@ -194,6 +194,7 @@ pub fn defaults() -> AppConfig { submap_fallback: true, }, tray: TrayConfig { enabled: true }, + coop: CoopConfig::default(), logging: LoggingConfig { level: "info".into(), format: "journald".into(), @@ -410,6 +411,31 @@ pub fn validate(config: &AppConfig) -> Result<(), ConfigError> { "logging.format must be journald, compact, or json".into(), )); } + if config.coop.disconnect_grace.as_millis() == 0 + || config.coop.disconnect_grace.as_millis() > 5 * 60 * 1_000 + { + return Err(ConfigError::Validation( + "coop.disconnect_grace must be between 1ms and 5m".into(), + )); + } + if config.coop.mode != CoopMode::Off { + let relay_url = config.coop.relay_url.as_deref().ok_or_else(|| { + ConfigError::Validation("coop.relay_url is required in host or guest mode".into()) + })?; + if !breakd_coop::valid_relay_url(relay_url) { + return Err(ConfigError::Validation( + "coop.relay_url must be a ws:// or wss:// URL without a fragment".into(), + )); + } + let room_token = config.coop.room_token.as_deref().ok_or_else(|| { + ConfigError::Validation("coop.room_token is required in host or guest mode".into()) + })?; + if !breakd_coop::valid_room_token(room_token) { + return Err(ConfigError::Validation( + "coop.room_token must contain exactly 32 hexadecimal characters".into(), + )); + } + } Ok(()) } @@ -492,6 +518,21 @@ mod tests { assert!(validate(&config).is_err()); } + #[test] + fn enabled_coop_requires_a_valid_relay_and_room_token() { + let mut config = defaults(); + config.coop.mode = CoopMode::Guest; + assert!(validate(&config).is_err()); + + config.coop.relay_url = Some("https://relay.example/ws".into()); + config.coop.room_token = Some("not-secret-enough".into()); + assert!(validate(&config).is_err()); + + config.coop.relay_url = Some("wss://relay.example/ws".into()); + config.coop.room_token = Some("0123456789abcdef0123456789abcdef".into()); + assert!(validate(&config).is_ok()); + } + #[test] fn missing_file_uses_defaults() { let directory = tempfile::tempdir().unwrap(); diff --git a/crates/coop/Cargo.toml b/crates/coop/Cargo.toml new file mode 100644 index 0000000..49fccd3 --- /dev/null +++ b/crates/coop/Cargo.toml @@ -0,0 +1,12 @@ +[package] +name = "breakd-coop" +version = "0.1.9" +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +breakd-core = { path = "../core" } +serde.workspace = true +thiserror.workspace = true +uuid.workspace = true diff --git a/crates/coop/src/lib.rs b/crates/coop/src/lib.rs new file mode 100644 index 0000000..642ef30 --- /dev/null +++ b/crates/coop/src/lib.rs @@ -0,0 +1,301 @@ +use std::fmt; + +use breakd_core::{BreakKind, BreakSessionId, Command, DueBreakId, DurationMs}; +use serde::{Deserialize, Serialize}; +use thiserror::Error; +use uuid::Uuid; + +pub const PROTOCOL_VERSION: u32 = 1; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum CoopRole { + Host, + Guest, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "action", content = "args", rename_all = "kebab-case")] +pub enum CoopAction { + Pause { duration_ms: Option }, + Resume, + ResumeBreak, + Reset, + Skip, + Postpone, + Mini, + Long, + Rest, + Toggle, +} + +impl CoopAction { + pub fn from_command(command: &Command) -> Option { + match command { + Command::Pause { duration } => Some(Self::Pause { + duration_ms: duration.map(DurationMs::as_millis), + }), + Command::Resume => Some(Self::Resume), + Command::ResumeBreak => Some(Self::ResumeBreak), + Command::Reset => Some(Self::Reset), + Command::Skip => Some(Self::Skip), + Command::Postpone => Some(Self::Postpone), + Command::Mini => Some(Self::Mini), + Command::Long => Some(Self::Long), + Command::Rest => Some(Self::Rest), + Command::Toggle => Some(Self::Toggle), + Command::Status + | Command::Reload + | Command::Outputs + | Command::Doctor + | Command::CoopHost { .. } + | Command::CoopJoin { .. } + | Command::CoopLeave + | Command::CoopStatus => None, + } + } + + pub fn into_command(self) -> Command { + match self { + Self::Pause { duration_ms } => Command::Pause { + duration: duration_ms.map(DurationMs::from_millis), + }, + Self::Resume => Command::Resume, + Self::ResumeBreak => Command::ResumeBreak, + Self::Reset => Command::Reset, + Self::Skip => Command::Skip, + Self::Postpone => Command::Postpone, + Self::Mini => Command::Mini, + Self::Long => Command::Long, + Self::Rest => Command::Rest, + Self::Toggle => Command::Toggle, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct CoopSnapshot { + pub host_id: Uuid, + pub revision: u64, + pub generated_unix_ms: u64, + pub paused: bool, + pub resume_at_unix_ms: Option, + pub phase: CoopPhase, + pub minis_since_long: u32, + pub longs_since_rest: u32, + pub postpone_count: u32, + pub can_skip: bool, + pub can_postpone: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "phase", rename_all = "kebab-case")] +pub enum CoopPhase { + Working { next: ScheduledBreak }, + Break { active: SharedBreak }, + Unavailable { reason: String }, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ScheduledBreak { + pub due_id: DueBreakId, + pub kind: BreakKind, + pub starts_unix_ms: u64, + pub duration_ms: u64, + pub strict_duration_ms: u64, + pub strict_entire: bool, + pub manual_resume: bool, + pub can_skip: bool, + pub can_postpone: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SharedBreak { + pub due_id: DueBreakId, + pub session_id: BreakSessionId, + pub kind: BreakKind, + pub started_unix_ms: u64, + pub ends_unix_ms: u64, + pub strict_until_unix_ms: u64, + pub strict_entire: bool, + pub manual_resume: bool, + pub completion_sound_emitted: bool, + pub can_skip: bool, + pub can_postpone: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "kebab-case")] +pub enum ClientMessage { + Hello { + version: u32, + role: CoopRole, + client_id: Uuid, + }, + Snapshot { + snapshot: CoopSnapshot, + }, + ActionRequest { + request_id: Uuid, + action: CoopAction, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "kebab-case")] +pub enum ServerMessage { + Ready { + host_present: bool, + guest_count: usize, + }, + Presence { + host_present: bool, + guest_count: usize, + }, + Snapshot { + snapshot: CoopSnapshot, + }, + ActionRequest { + request_id: Uuid, + action: CoopAction, + }, + Error { + code: String, + message: String, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Invite { + relay_url: String, + room_token: String, +} + +impl Invite { + pub fn new( + relay_url: impl Into, + room_token: impl Into, + ) -> Result { + let relay_url = relay_url.into(); + let room_token = room_token.into().to_ascii_lowercase(); + if !valid_relay_url(&relay_url) { + return Err(InviteError::RelayUrl); + } + if !valid_room_token(&room_token) { + return Err(InviteError::RoomToken); + } + Ok(Self { + relay_url, + room_token, + }) + } + + pub fn parse(value: &str) -> Result { + let (relay_url, fragment) = value.rsplit_once('#').ok_or(InviteError::Fragment)?; + let room_token = fragment + .strip_prefix("breakd=") + .ok_or(InviteError::Fragment)?; + Self::new(relay_url, room_token) + } + + pub fn relay_url(&self) -> &str { + &self.relay_url + } + + pub fn room_token(&self) -> &str { + &self.room_token + } +} + +impl fmt::Display for Invite { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(formatter, "{}#breakd={}", self.relay_url, self.room_token) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum InviteError { + #[error("relay must be a ws:// or wss:// URL without a fragment")] + RelayUrl, + #[error("room token must contain exactly 32 hexadecimal characters")] + RoomToken, + #[error("invite must end in #breakd=")] + Fragment, +} + +pub fn valid_relay_url(value: &str) -> bool { + let Some(rest) = value + .strip_prefix("ws://") + .or_else(|| value.strip_prefix("wss://")) + else { + return false; + }; + let authority = rest.split(['/', '?']).next().unwrap_or_default(); + !authority.is_empty() + && !authority.starts_with(':') + && !authority.contains('@') + && !value.contains('#') + && !value.chars().any(char::is_whitespace) +} + +pub fn valid_room_token(value: &str) -> bool { + value.len() == 32 && value.bytes().all(|byte| byte.is_ascii_hexdigit()) +} + +#[cfg(test)] +mod tests { + use super::*; + + const TOKEN: &str = "0123456789abcdef0123456789abcdef"; + + #[test] + fn invite_round_trips_without_putting_the_token_in_the_request_url() { + let invite = Invite::new("wss://relay.example/ws", TOKEN).unwrap(); + assert_eq!( + invite.to_string(), + format!("wss://relay.example/ws#breakd={TOKEN}") + ); + let parsed = Invite::parse(&invite.to_string()).unwrap(); + assert_eq!(parsed.relay_url(), "wss://relay.example/ws"); + assert!(!parsed.relay_url().contains(TOKEN)); + assert_eq!(parsed.room_token(), TOKEN); + } + + #[test] + fn malformed_invites_are_rejected() { + assert_eq!( + Invite::parse("https://relay.example/ws#breakd=bad"), + Err(InviteError::RelayUrl) + ); + assert_eq!( + Invite::parse("wss://relay.example/ws"), + Err(InviteError::Fragment) + ); + assert_eq!( + Invite::parse("wss://relay.example/ws#token=abc"), + Err(InviteError::Fragment) + ); + assert_eq!( + Invite::parse("ws:///ws#breakd=0123456789abcdef0123456789abcdef"), + Err(InviteError::RelayUrl) + ); + assert_eq!( + Invite::parse( + "wss://user:password@relay.example/ws#breakd=0123456789abcdef0123456789abcdef" + ), + Err(InviteError::RelayUrl) + ); + } + + #[test] + fn supported_commands_round_trip_as_actions() { + let command = Command::Pause { + duration: Some(DurationMs::from_millis(42_000)), + }; + assert_eq!( + CoopAction::from_command(&command).unwrap().into_command(), + command + ); + assert!(CoopAction::from_command(&Command::Status).is_none()); + } +} diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index d1d9621..670376a 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -343,6 +343,38 @@ impl Default for TrayConfig { } } +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum CoopMode { + #[default] + Off, + Host, + Guest, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct CoopConfig { + #[serde(default)] + pub mode: CoopMode, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub relay_url: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub room_token: Option, + #[serde(default = "default_coop_disconnect_grace")] + pub disconnect_grace: DurationMs, +} + +impl Default for CoopConfig { + fn default() -> Self { + Self { + mode: CoopMode::Off, + relay_url: None, + room_token: None, + disconnect_grace: default_coop_disconnect_grace(), + } + } +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct AppConfig { pub schema_version: u32, @@ -363,6 +395,8 @@ pub struct AppConfig { pub hyprland: HyprlandConfig, #[serde(default)] pub tray: TrayConfig, + #[serde(default)] + pub coop: CoopConfig, pub logging: LoggingConfig, } @@ -382,6 +416,10 @@ fn default_rest_postpone() -> PostponeRule { } } +const fn default_coop_disconnect_grace() -> DurationMs { + DurationMs::from_millis(10_000) +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] pub struct ClockSample { pub monotonic_ms: u64, @@ -494,6 +532,10 @@ pub enum Command { Reload, Outputs, Doctor, + CoopHost { relay_url: String }, + CoopJoin { invite: String }, + CoopLeave, + CoopStatus, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] diff --git a/crates/daemon/Cargo.toml b/crates/daemon/Cargo.toml index a5d9359..23b19c6 100644 --- a/crates/daemon/Cargo.toml +++ b/crates/daemon/Cargo.toml @@ -7,13 +7,16 @@ rust-version.workspace = true [dependencies] anyhow.workspace = true +breakd-coop = { path = "../coop" } breakd-config = { path = "../config" } breakd-core = { path = "../core" } breakd-ipc = { path = "../ipc" } breakd-platform-linux = { path = "../platform-linux" } breakd-scheduler = { path = "../scheduler" } breakd-tray = { path = "../tray" } +futures-util.workspace = true serde_json.workspace = true tokio.workspace = true +tokio-tungstenite = { workspace = true, features = ["connect", "rustls-tls-webpki-roots"] } tracing.workspace = true uuid.workspace = true diff --git a/crates/daemon/src/coop.rs b/crates/daemon/src/coop.rs new file mode 100644 index 0000000..af3e5a8 --- /dev/null +++ b/crates/daemon/src/coop.rs @@ -0,0 +1,466 @@ +use std::time::Duration; + +use anyhow::{Context, Result, bail}; +use breakd_coop::{ + ClientMessage, CoopAction, CoopRole, CoopSnapshot, Invite, PROTOCOL_VERSION, ServerMessage, +}; +use breakd_core::{CoopConfig, CoopMode}; +use futures_util::{SinkExt, StreamExt}; +use serde_json::json; +use tokio::{ + sync::{mpsc, watch}, + task::JoinHandle, + time::{Instant, MissedTickBehavior, interval_at, sleep}, +}; +use tokio_tungstenite::{ + connect_async, + tungstenite::{Message, client::IntoClientRequest, http::HeaderValue}, +}; +use uuid::Uuid; + +const MAX_MESSAGE_BYTES: usize = 128 * 1024; + +#[derive(Debug)] +pub struct CoopEvent { + generation: u64, + kind: CoopEventKind, +} + +#[derive(Debug)] +enum CoopEventKind { + Connected, + Disconnected(String), + Presence { + host_present: bool, + guest_count: usize, + }, + Snapshot(CoopSnapshot), + ActionRequest { + request_id: Uuid, + action: CoopAction, + }, + Error(String), +} + +#[derive(Debug)] +pub enum AcceptedEvent { + StateChanged, + Snapshot(CoopSnapshot), + ActionRequest { + request_id: Uuid, + action: CoopAction, + }, + Error(String), + Stale, +} + +pub struct CoopRuntime { + config: CoopConfig, + generation: u64, + event_sender: mpsc::Sender, + task: Option>, + snapshot_sender: watch::Sender>, + action_sender: mpsc::Sender<(Uuid, CoopAction, Instant)>, + connected: bool, + host_present: bool, + guest_count: usize, + started_at: Instant, + last_snapshot_at: Option, + last_snapshot_version: Option<(Uuid, u64)>, + disconnect_reason: Option, + was_holding_local: bool, + host_id: Uuid, + next_revision: u64, +} + +impl CoopRuntime { + pub fn new(config: CoopConfig, event_sender: mpsc::Sender) -> Self { + let (snapshot_sender, _) = watch::channel(None); + let (action_sender, _) = mpsc::channel(1); + let mut runtime = Self { + config: CoopConfig::default(), + generation: 0, + event_sender, + task: None, + snapshot_sender, + action_sender, + connected: false, + host_present: false, + guest_count: 0, + started_at: Instant::now(), + last_snapshot_at: None, + last_snapshot_version: None, + disconnect_reason: None, + was_holding_local: false, + host_id: Uuid::new_v4(), + next_revision: 0, + }; + runtime.reconfigure(config); + runtime + } + + pub fn reconfigure(&mut self, config: CoopConfig) { + if let Some(task) = self.task.take() { + task.abort(); + } + self.generation = self.generation.wrapping_add(1); + self.config = config; + self.connected = false; + self.host_present = false; + self.guest_count = 0; + self.started_at = Instant::now(); + self.last_snapshot_at = None; + self.last_snapshot_version = None; + self.disconnect_reason = None; + self.was_holding_local = self.config.mode == CoopMode::Guest; + self.host_id = Uuid::new_v4(); + self.next_revision = 0; + + let (snapshot_sender, snapshot_receiver) = watch::channel(None); + let (action_sender, action_receiver) = mpsc::channel(32); + self.snapshot_sender = snapshot_sender; + self.action_sender = action_sender; + + let Some(role) = self.role() else { + return; + }; + let (Some(relay_url), Some(room_token)) = ( + self.config.relay_url.clone(), + self.config.room_token.clone(), + ) else { + self.disconnect_reason = Some("co-op configuration is incomplete".into()); + return; + }; + let generation = self.generation; + let events = self.event_sender.clone(); + self.task = Some(tokio::spawn(connection_loop( + ConnectionConfig { + relay_url, + room_token, + role, + client_id: Uuid::new_v4(), + }, + generation, + events, + snapshot_receiver, + action_receiver, + ))); + } + + pub fn accept(&mut self, event: CoopEvent) -> AcceptedEvent { + if event.generation != self.generation { + return AcceptedEvent::Stale; + } + match event.kind { + CoopEventKind::Connected => { + self.connected = true; + self.disconnect_reason = None; + AcceptedEvent::StateChanged + } + CoopEventKind::Disconnected(reason) => { + self.connected = false; + self.host_present = false; + self.disconnect_reason = Some(reason); + AcceptedEvent::StateChanged + } + CoopEventKind::Presence { + host_present, + guest_count, + } => { + self.host_present = host_present; + self.guest_count = guest_count; + AcceptedEvent::StateChanged + } + CoopEventKind::Snapshot(snapshot) => { + let version = (snapshot.host_id, snapshot.revision); + let is_newer = self + .last_snapshot_version + .is_none_or(|(host_id, revision)| { + host_id != snapshot.host_id || snapshot.revision > revision + }); + if !is_newer || self.config.mode != CoopMode::Guest { + return AcceptedEvent::Stale; + } + self.last_snapshot_version = Some(version); + self.last_snapshot_at = Some(Instant::now()); + self.host_present = true; + AcceptedEvent::Snapshot(snapshot) + } + CoopEventKind::ActionRequest { request_id, action } => { + if self.config.mode == CoopMode::Host { + AcceptedEvent::ActionRequest { request_id, action } + } else { + AcceptedEvent::Stale + } + } + CoopEventKind::Error(error) => AcceptedEvent::Error(error), + } + } + + pub fn role(&self) -> Option { + match self.config.mode { + CoopMode::Off => None, + CoopMode::Host => Some(CoopRole::Host), + CoopMode::Guest => Some(CoopRole::Guest), + } + } + + pub fn is_host(&self) -> bool { + self.config.mode == CoopMode::Host + } + + pub fn holds_local_schedule(&self) -> bool { + if self.config.mode != CoopMode::Guest { + return false; + } + self.last_snapshot_at.map_or_else( + || self.started_at.elapsed() <= self.config.disconnect_grace.as_duration(), + |last| last.elapsed() <= self.config.disconnect_grace.as_duration(), + ) + } + + pub fn has_fresh_snapshot(&self) -> bool { + self.config.mode == CoopMode::Guest + && self + .last_snapshot_at + .is_some_and(|last| last.elapsed() <= self.config.disconnect_grace.as_duration()) + } + + pub fn take_fallback_transition(&mut self) -> bool { + let holding = self.holds_local_schedule(); + let transitioned = self.was_holding_local && !holding; + self.was_holding_local = holding; + transitioned + } + + pub fn publish(&self, snapshot: CoopSnapshot) { + if self.is_host() { + self.snapshot_sender.send_replace(Some(snapshot)); + } + } + + pub fn next_snapshot_identity(&mut self) -> (Uuid, u64) { + let revision = self.next_revision; + self.next_revision = self.next_revision.wrapping_add(1); + (self.host_id, revision) + } + + pub fn request_action(&self, action: CoopAction) -> Result { + if self.config.mode != CoopMode::Guest { + bail!("co-op actions can only be forwarded by a guest"); + } + if !self.connected || !self.host_present { + bail!("co-op host is not connected"); + } + let request_id = Uuid::new_v4(); + self.action_sender + .try_send((request_id, action, Instant::now())) + .context("co-op action queue is full")?; + Ok(request_id) + } + + pub fn status_json(&self) -> serde_json::Value { + let invite = match ( + self.config.mode, + self.config.relay_url.as_deref(), + self.config.room_token.as_deref(), + ) { + (CoopMode::Host, Some(relay), Some(token)) => Invite::new(relay, token) + .ok() + .map(|invite| invite.to_string()), + _ => None, + }; + json!({ + "mode": match self.config.mode { + CoopMode::Off => "off", + CoopMode::Host => "host", + CoopMode::Guest => "guest", + }, + "relay_url": self.config.relay_url, + "connected": self.connected, + "host_present": self.host_present, + "guest_count": self.guest_count, + "following_host": self.has_fresh_snapshot(), + "disconnect_grace_ms": self.config.disconnect_grace.as_millis(), + "last_snapshot_age_ms": self.last_snapshot_at.map(|at| { + u64::try_from(at.elapsed().as_millis()).unwrap_or(u64::MAX) + }), + "last_error": self.disconnect_reason, + "invite": invite, + }) + } +} + +impl Drop for CoopRuntime { + fn drop(&mut self) { + if let Some(task) = self.task.take() { + task.abort(); + } + } +} + +struct ConnectionConfig { + relay_url: String, + room_token: String, + role: CoopRole, + client_id: Uuid, +} + +async fn connection_loop( + config: ConnectionConfig, + generation: u64, + events: mpsc::Sender, + mut snapshots: watch::Receiver>, + mut actions: mpsc::Receiver<(Uuid, CoopAction, Instant)>, +) { + let mut retry = Duration::from_secs(1); + loop { + let mut established = false; + let result = connect_once( + &config, + generation, + &events, + &mut snapshots, + &mut actions, + &mut established, + ) + .await; + let reason = result + .err() + .map_or_else(|| "connection closed".into(), |error| format!("{error:#}")); + let _ = events + .send(CoopEvent { + generation, + kind: CoopEventKind::Disconnected(reason), + }) + .await; + if established { + retry = Duration::from_secs(1); + } + sleep(retry).await; + if !established { + retry = (retry * 2).min(Duration::from_secs(30)); + } + } +} + +async fn connect_once( + config: &ConnectionConfig, + generation: u64, + events: &mpsc::Sender, + snapshots: &mut watch::Receiver>, + actions: &mut mpsc::Receiver<(Uuid, CoopAction, Instant)>, + established: &mut bool, +) -> Result<()> { + let mut request = config + .relay_url + .as_str() + .into_client_request() + .context("invalid relay URL")?; + request.headers_mut().insert( + "authorization", + HeaderValue::from_str(&format!("Bearer {}", config.room_token)) + .context("invalid room token")?, + ); + let (mut socket, _) = connect_async(request) + .await + .context("connect to co-op relay")?; + send_json( + &mut socket, + &ClientMessage::Hello { + version: PROTOCOL_VERSION, + role: config.role, + client_id: config.client_id, + }, + ) + .await?; + let initial_snapshot = snapshots.borrow().clone(); + if config.role == CoopRole::Host + && let Some(snapshot) = initial_snapshot + { + send_json(&mut socket, &ClientMessage::Snapshot { snapshot }).await?; + } + let _ = events + .send(CoopEvent { + generation, + kind: CoopEventKind::Connected, + }) + .await; + *established = true; + + let mut heartbeat = interval_at( + Instant::now() + Duration::from_secs(20), + Duration::from_secs(20), + ); + heartbeat.set_missed_tick_behavior(MissedTickBehavior::Skip); + loop { + tokio::select! { + incoming = socket.next() => { + let message = incoming.context("relay closed the WebSocket")??; + match message { + Message::Text(text) => { + if text.len() > MAX_MESSAGE_BYTES { + bail!("relay message exceeds {MAX_MESSAGE_BYTES} bytes"); + } + let message: ServerMessage = serde_json::from_str(text.as_ref()) + .context("invalid relay message")?; + let kind = match message { + ServerMessage::Ready { host_present, guest_count } + | ServerMessage::Presence { host_present, guest_count } => { + CoopEventKind::Presence { host_present, guest_count } + } + ServerMessage::Snapshot { snapshot } => CoopEventKind::Snapshot(snapshot), + ServerMessage::ActionRequest { request_id, action } => { + CoopEventKind::ActionRequest { request_id, action } + } + ServerMessage::Error { code, message } => { + CoopEventKind::Error(format!("{code}: {message}")) + } + }; + let _ = events.send(CoopEvent { generation, kind }).await; + } + Message::Ping(payload) => socket.send(Message::Pong(payload)).await?, + Message::Close(_) => return Ok(()), + Message::Binary(_) | Message::Pong(_) | Message::Frame(_) => {} + } + } + changed = snapshots.changed(), if config.role == CoopRole::Host => { + changed.context("co-op snapshot publisher stopped")?; + let latest_snapshot = snapshots.borrow().clone(); + if let Some(snapshot) = latest_snapshot { + send_json(&mut socket, &ClientMessage::Snapshot { snapshot }).await?; + } + } + action = actions.recv(), if config.role == CoopRole::Guest => { + let (request_id, action, queued_at) = action.context("co-op action publisher stopped")?; + if queued_at.elapsed() > Duration::from_secs(2) { + let _ = events + .send(CoopEvent { + generation, + kind: CoopEventKind::Error(format!( + "action {request_id} expired while reconnecting" + )), + }) + .await; + continue; + } + send_json(&mut socket, &ClientMessage::ActionRequest { request_id, action }).await?; + } + _ = heartbeat.tick() => { + socket.send(Message::Ping(Vec::new().into())).await?; + } + } + } +} + +async fn send_json( + socket: &mut tokio_tungstenite::WebSocketStream, + message: &ClientMessage, +) -> Result<()> +where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin, +{ + let encoded = serde_json::to_string(message)?; + socket.send(Message::Text(encoded.into())).await?; + Ok(()) +} diff --git a/crates/daemon/src/lib.rs b/crates/daemon/src/lib.rs index 957456d..e806704 100644 --- a/crates/daemon/src/lib.rs +++ b/crates/daemon/src/lib.rs @@ -5,7 +5,10 @@ use std::{ }; use anyhow::{Context, Result}; -use breakd_core::{BreakSessionId, Command, CompletionSound, OverlaySpec, Response}; +use breakd_coop::{CoopAction, Invite}; +use breakd_core::{ + BreakSessionId, Command, CompletionSound, CoopConfig, CoopMode, OverlaySpec, Response, +}; use breakd_ipc::{IPC_VERSION, IncomingRequest, Server}; use breakd_platform_linux::{ EventSoundClient, HyprlandClient, IdleCapability, IdleEvent, LinuxClock, NotificationClient, @@ -20,6 +23,10 @@ use tokio::{ time::{Duration, Instant, interval}, }; +mod coop; + +use coop::{AcceptedEvent, CoopRuntime}; + pub async fn run() -> Result<()> { let instance = breakd_config::RuntimeInstance::current(); let mut config = breakd_config::load().context("load configuration")?; @@ -57,6 +64,9 @@ pub async fn run() -> Result<()> { }; state_store.save(&scheduler.snapshot())?; + let (coop_event_sender, mut coop_event_receiver) = mpsc::channel(64); + let mut coop = CoopRuntime::new(config.coop.clone(), coop_event_sender); + let server = Server::bind(&socket_path)?; let (request_sender, mut request_receiver) = mpsc::channel::(32); tokio::spawn(async move { @@ -91,7 +101,11 @@ pub async fn run() -> Result<()> { }; let mut overlay = OverlaySupervisor::default(); apply_effects( - scheduler.startup_effects(now), + if coop.holds_local_schedule() { + Vec::new() + } else { + scheduler.startup_effects(now) + }, &mut overlay, ¬ifications, event_sounds.as_ref(), @@ -119,21 +133,38 @@ pub async fn run() -> Result<()> { ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); let mut recent_requests: VecDeque<(uuid::Uuid, Response)> = VecDeque::new(); let mut terminate = signal(SignalKind::terminate()).context("install SIGTERM handler")?; + let mut next_coop_publish = Instant::now(); loop { tokio::select! { _ = ticker.tick() => { let now = clock.sample()?; - let previous = scheduler.state().clone(); - let effects = scheduler.handle_event(SchedulerEvent::Tick, now); - persist_if_changed(&state_store, &scheduler, &previous)?; - apply_effects( - effects, - &mut overlay, - ¬ifications, - event_sounds.as_ref(), - config.completion.sound, - ).await; + if coop.take_fallback_transition() { + tracing::warn!("co-op host timed out; starting a fresh local schedule"); + let effects = scheduler.reset_after_coop_disconnect(now); + state_store.save(&scheduler.snapshot())?; + apply_effects( + effects, + &mut overlay, + ¬ifications, + event_sounds.as_ref(), + config.completion.sound, + ).await; + } + if !coop.holds_local_schedule() || coop.has_fresh_snapshot() { + let previous = scheduler.state().clone(); + let effects = scheduler.handle_event(SchedulerEvent::Tick, now); + if !coop.has_fresh_snapshot() { + persist_if_changed(&state_store, &scheduler, &previous)?; + } + apply_effects( + effects, + &mut overlay, + ¬ifications, + event_sounds.as_ref(), + config.completion.sound, + ).await; + } if let Some((session_id, reason)) = overlay.poll_exit().await? { let status = scheduler.status(now); if status.active_session == Some(session_id) @@ -170,9 +201,13 @@ pub async fn run() -> Result<()> { &mut overlay, ¬ifications, event_sounds.as_ref(), + &mut coop, idle_capability, tray.available(), ).await; + if coop.is_host() { + next_coop_publish = Instant::now(); + } recent_requests.push_back((incoming.request.request_id, response.clone())); if recent_requests.len() > 128 { recent_requests.pop_front(); @@ -191,12 +226,16 @@ pub async fn run() -> Result<()> { &mut overlay, ¬ifications, event_sounds.as_ref(), + &mut coop, idle_capability, tray.available(), ).await { Ok((message, _)) => tracing::info!(%message, "tray command handled"), Err(error) => tracing::warn!(%error, "tray command failed"), } + if coop.is_host() { + next_coop_publish = Instant::now(); + } } TrayAction::OpenSettings => { if let Err(error) = spawn_settings().await { @@ -206,6 +245,9 @@ pub async fn run() -> Result<()> { } } Some(event) = power_receiver.recv() => { + if coop.holds_local_schedule() { + continue; + } let now = clock.sample()?; let scheduler_event = match event { PowerEvent::PreparingForSleep => SchedulerEvent::SuspendStarted, @@ -223,8 +265,14 @@ pub async fn run() -> Result<()> { event_sounds.as_ref(), config.completion.sound, ).await; + if coop.is_host() { + next_coop_publish = Instant::now(); + } } Some(event) = idle_receiver.recv() => { + if coop.holds_local_schedule() { + continue; + } let now = clock.sample()?; let scheduler_event = match event { IdleEvent::Idled => SchedulerEvent::IdleThresholdReached, @@ -240,6 +288,46 @@ pub async fn run() -> Result<()> { event_sounds.as_ref(), config.completion.sound, ).await; + if coop.is_host() { + next_coop_publish = Instant::now(); + } + } + Some(event) = coop_event_receiver.recv() => { + match coop.accept(event) { + AcceptedEvent::Snapshot(snapshot) => { + let now = clock.sample()?; + let effects = scheduler.adopt_coop_snapshot(&snapshot, now); + apply_effects( + effects, + &mut overlay, + ¬ifications, + event_sounds.as_ref(), + config.completion.sound, + ).await; + } + AcceptedEvent::ActionRequest { request_id, action } => { + let command = action.into_command(); + match execute_command( + &command, + &clock, + &state_store, + &mut scheduler, + &mut config, + &mut overlay, + ¬ifications, + event_sounds.as_ref(), + &mut coop, + idle_capability, + tray.available(), + ).await { + Ok((message, _)) => tracing::info!(%request_id, %message, "co-op action handled"), + Err(error) => tracing::warn!(%request_id, %error, "co-op action rejected"), + } + next_coop_publish = Instant::now(); + } + AcceptedEvent::Error(error) => tracing::warn!(%error, "co-op relay error"), + AcceptedEvent::StateChanged | AcceptedEvent::Stale => {} + } } result = tokio::signal::ctrl_c() => { result?; @@ -252,6 +340,17 @@ pub async fn run() -> Result<()> { } } + if coop.is_host() && Instant::now() >= next_coop_publish { + let now = clock.sample()?; + let (host_id, revision) = coop.next_snapshot_identity(); + let (snapshot, scheduler_changed) = scheduler.coop_snapshot(host_id, revision, now); + if scheduler_changed { + state_store.save(&scheduler.snapshot())?; + } + coop.publish(snapshot); + next_coop_publish = Instant::now() + Duration::from_secs(1); + } + let now = clock.sample()?; let status = scheduler.status(now); let next_tray_state = tray_state(status.clone()); @@ -271,13 +370,24 @@ pub async fn run() -> Result<()> { fn completion_sound_directory(instance: breakd_config::RuntimeInstance) -> PathBuf { match instance { - breakd_config::RuntimeInstance::Production => PathBuf::from("/usr/share/breakd"), + breakd_config::RuntimeInstance::Production => std::env::current_exe() + .ok() + .and_then(|executable| installed_data_directory(&executable)) + .filter(|directory| directory.is_dir()) + .unwrap_or_else(|| PathBuf::from("/usr/share/breakd")), breakd_config::RuntimeInstance::Development => { Path::new(env!("CARGO_MANIFEST_DIR")).join("../platform-linux/assets") } } } +fn installed_data_directory(executable: &Path) -> Option { + executable + .parent()? + .parent() + .map(|prefix| prefix.join("share/breakd")) +} + async fn spawn_settings() -> Result<()> { let executable = std::env::current_exe().context("resolve breakd executable")?; let mut child = TokioCommand::new(executable) @@ -305,6 +415,7 @@ async fn handle_request( overlay: &mut OverlaySupervisor, notifications: &NotificationClient, event_sounds: Option<&EventSoundClient>, + coop: &mut CoopRuntime, idle_capability: IdleCapability, tray_available: bool, ) -> Response { @@ -318,6 +429,7 @@ async fn handle_request( overlay, notifications, event_sounds, + coop, idle_capability, tray_available, ) @@ -351,10 +463,20 @@ async fn execute_command( overlay: &mut OverlaySupervisor, notifications: &NotificationClient, event_sounds: Option<&EventSoundClient>, + coop: &mut CoopRuntime, idle_capability: IdleCapability, tray_available: bool, ) -> Result<(String, Option)> { let now = clock.sample()?; + if coop.holds_local_schedule() + && let Some(action) = CoopAction::from_command(command) + { + let request_id = coop.request_action(action)?; + return Ok(( + "request sent to co-op host".into(), + Some(serde_json::json!({ "request_id": request_id })), + )); + } match command { Command::Status => Ok(( "status".into(), @@ -419,10 +541,81 @@ async fn execute_command( }); Ok(("doctor report".into(), Some(report))) } + Command::CoopHost { relay_url } => { + let token = uuid::Uuid::new_v4().simple().to_string(); + let invite = Invite::new(relay_url, &token)?; + config.coop = CoopConfig { + mode: CoopMode::Host, + relay_url: Some(invite.relay_url().into()), + room_token: Some(invite.room_token().into()), + disconnect_grace: config.coop.disconnect_grace, + }; + breakd_config::save(config)?; + coop.reconfigure(config.coop.clone()); + Ok(( + "co-op room created".into(), + Some(serde_json::json!({ + "invite": invite.to_string(), + "status": coop.status_json(), + })), + )) + } + Command::CoopJoin { invite } => { + let invite = Invite::parse(invite)?; + config.coop = CoopConfig { + mode: CoopMode::Guest, + relay_url: Some(invite.relay_url().into()), + room_token: Some(invite.room_token().into()), + disconnect_grace: config.coop.disconnect_grace, + }; + breakd_config::save(config)?; + coop.reconfigure(config.coop.clone()); + overlay.stop_any().await; + Ok(("joining co-op room".into(), Some(coop.status_json()))) + } + Command::CoopLeave => { + let disconnect_grace = config.coop.disconnect_grace; + config.coop = CoopConfig { + disconnect_grace, + ..CoopConfig::default() + }; + breakd_config::save(config)?; + coop.reconfigure(config.coop.clone()); + let effects = scheduler.reset_after_coop_disconnect(now); + state_store.save(&scheduler.snapshot())?; + apply_effects( + effects, + overlay, + notifications, + event_sounds, + config.completion.sound, + ) + .await; + Ok(("left co-op room; local schedule reset".into(), None)) + } + Command::CoopStatus => Ok(("co-op status".into(), Some(coop.status_json()))), Command::Reload => { let updated = breakd_config::load()?; + let previous_coop = config.coop.clone(); scheduler.replace_config(updated.clone()); *config = updated; + if config.coop != previous_coop { + let was_guest = previous_coop.mode == CoopMode::Guest; + coop.reconfigure(config.coop.clone()); + if config.coop.mode == CoopMode::Guest { + overlay.stop_any().await; + } else if was_guest { + let effects = scheduler.reset_after_coop_disconnect(now); + apply_effects( + effects, + overlay, + notifications, + event_sounds, + config.completion.sound, + ) + .await; + } + } state_store.save(&scheduler.snapshot())?; Ok(( "configuration reloaded; restart for idle-monitor or logging changes".into(), @@ -711,10 +904,30 @@ fn command_message(command: &Command) -> &'static str { Command::Long => "long break started", Command::Rest => "rest break started", Command::Toggle => "schedule toggled", - Command::Status | Command::Reload | Command::Outputs | Command::Doctor => "ok", + Command::Status + | Command::Reload + | Command::Outputs + | Command::Doctor + | Command::CoopHost { .. } + | Command::CoopJoin { .. } + | Command::CoopLeave + | Command::CoopStatus => "ok", } } pub fn socket_exists(path: &Path) -> bool { path.exists() } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn installed_data_directory_follows_the_executable_prefix() { + assert_eq!( + installed_data_directory(Path::new("/nix/store/hash-breakd/bin/breakd")), + Some(PathBuf::from("/nix/store/hash-breakd/share/breakd")) + ); + } +} diff --git a/crates/ipc/src/lib.rs b/crates/ipc/src/lib.rs index 1e044e7..b019041 100644 --- a/crates/ipc/src/lib.rs +++ b/crates/ipc/src/lib.rs @@ -15,7 +15,7 @@ use tokio::{ }; use uuid::Uuid; -pub const IPC_VERSION: u32 = 1; +pub const IPC_VERSION: u32 = 2; const MAX_FRAME_SIZE: usize = 64 * 1024; #[derive(Debug, Error)] diff --git a/crates/relay/Cargo.toml b/crates/relay/Cargo.toml new file mode 100644 index 0000000..4ab765e --- /dev/null +++ b/crates/relay/Cargo.toml @@ -0,0 +1,21 @@ +[package] +name = "breakd-relay" +version = "0.1.9" +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +anyhow.workspace = true +breakd-coop = { path = "../coop" } +clap.workspace = true +futures-util.workspace = true +serde_json.workspace = true +tokio.workspace = true +tokio-tungstenite = { workspace = true, features = ["handshake"] } +tracing.workspace = true +tracing-subscriber.workspace = true +uuid.workspace = true + +[dev-dependencies] +breakd-core = { path = "../core" } diff --git a/crates/relay/src/main.rs b/crates/relay/src/main.rs new file mode 100644 index 0000000..5ee1158 --- /dev/null +++ b/crates/relay/src/main.rs @@ -0,0 +1,575 @@ +use std::{collections::HashMap, net::SocketAddr, sync::Arc, time::Duration}; + +use anyhow::{Context, Result, bail}; +use breakd_coop::{ClientMessage, CoopRole, PROTOCOL_VERSION, ServerMessage, valid_room_token}; +use clap::Parser; +use futures_util::{SinkExt, StreamExt}; +use tokio::{net::TcpListener, sync::mpsc, time::timeout}; +use tokio_tungstenite::{ + accept_hdr_async, + tungstenite::{ + Message, + handshake::server::{ErrorResponse, Request, Response}, + http::{HeaderValue, StatusCode}, + }, +}; +use tracing_subscriber::EnvFilter; +use uuid::Uuid; + +const MAX_MESSAGE_BYTES: usize = 128 * 1024; + +#[derive(Debug, Parser)] +#[command( + name = "breakd-relay", + version, + about = "Small room relay for breakd co-op" +)] +struct Arguments { + /// Address to listen on. Put a TLS reverse proxy in front for public use. + #[arg(long, default_value = "127.0.0.1:8787")] + listen: SocketAddr, + /// Maximum host and guest connections in one room. + #[arg(long, default_value_t = 8)] + max_room_size: usize, + /// Maximum number of simultaneously live rooms. + #[arg(long, default_value_t = 256)] + max_rooms: usize, +} + +#[derive(Default)] +struct Room { + host: Option, + guests: HashMap, + latest_snapshot: Option, +} + +#[derive(Clone)] +struct Peer { + connection_id: Uuid, + sender: mpsc::Sender, +} + +type Rooms = Arc>>; + +#[tokio::main] +async fn main() -> Result<()> { + init_logging(); + let arguments = Arguments::parse(); + if !(2..=64).contains(&arguments.max_room_size) { + bail!("--max-room-size must be between 2 and 64"); + } + if !(1..=65_536).contains(&arguments.max_rooms) { + bail!("--max-rooms must be between 1 and 65536"); + } + let listener = TcpListener::bind(arguments.listen) + .await + .with_context(|| format!("bind {}", arguments.listen))?; + tracing::info!(listen = %arguments.listen, "co-op relay listening"); + let rooms = Rooms::default(); + + loop { + let (stream, remote) = listener.accept().await?; + let rooms = rooms.clone(); + let max_room_size = arguments.max_room_size; + let max_rooms = arguments.max_rooms; + tokio::spawn(async move { + if let Err(error) = handle_connection(stream, rooms, max_room_size, max_rooms).await { + tracing::debug!(%remote, %error, "co-op connection ended"); + } + }); + } +} + +#[allow(clippy::result_large_err)] // tungstenite's handshake callback owns this response type. +async fn handle_connection( + stream: tokio::net::TcpStream, + rooms: Rooms, + max_room_size: usize, + max_rooms: usize, +) -> Result<()> { + let token_slot = Arc::new(std::sync::Mutex::new(None)); + let callback_slot = token_slot.clone(); + let mut socket = + accept_hdr_async( + stream, + move |request: &Request, response: Response| match bearer_token(request) { + Some(token) => { + *callback_slot.lock().expect("token mutex poisoned") = Some(token); + Ok(response) + } + None => Err(handshake_error( + StatusCode::UNAUTHORIZED, + "missing or invalid bearer token", + )), + }, + ) + .await + .context("WebSocket handshake failed")?; + let room_token = token_slot + .lock() + .expect("token mutex poisoned") + .take() + .context("authenticated room token was not retained")?; + + let hello = timeout(Duration::from_secs(5), socket.next()) + .await + .context("hello timed out")? + .context("connection closed before hello")??; + let (role, _client_id) = match parse_client_message(hello)? { + ClientMessage::Hello { + version, + role, + client_id, + } if version == PROTOCOL_VERSION => (role, client_id), + ClientMessage::Hello { version, .. } => { + send_server( + &mut socket, + &ServerMessage::Error { + code: "protocol-version".into(), + message: format!( + "client protocol {version} is unsupported; expected {PROTOCOL_VERSION}" + ), + }, + ) + .await?; + bail!("unsupported protocol version {version}"); + } + _ => bail!("the first message must be hello"), + }; + + let connection_id = Uuid::new_v4(); + let (outgoing_sender, mut outgoing_receiver) = mpsc::channel(32); + let peer = Peer { + connection_id, + sender: outgoing_sender.clone(), + }; + let initial = + match register_peer(&rooms, &room_token, role, peer, max_room_size, max_rooms).await { + Ok(initial) => initial, + Err(error) => { + send_server( + &mut socket, + &ServerMessage::Error { + code: "room-rejected".into(), + message: error.to_string(), + }, + ) + .await?; + return Err(error); + } + }; + for message in initial { + let _ = outgoing_sender.try_send(server_message(&message)?); + } + broadcast_presence(&rooms, &room_token).await; + + let (mut sink, mut source) = socket.split(); + let result = async { + loop { + tokio::select! { + incoming = source.next() => { + let message = incoming.context("WebSocket closed")??; + match message { + Message::Text(text) => { + if text.len() > MAX_MESSAGE_BYTES { + bail!("message exceeds {MAX_MESSAGE_BYTES} bytes"); + } + let message: ClientMessage = serde_json::from_str(text.as_ref()) + .context("invalid client message")?; + route_message( + &rooms, + &room_token, + connection_id, + role, + message, + &outgoing_sender, + ).await?; + } + Message::Ping(payload) => sink.send(Message::Pong(payload)).await?, + Message::Close(_) => break, + Message::Binary(_) | Message::Pong(_) | Message::Frame(_) => {} + } + } + outgoing = outgoing_receiver.recv() => { + let Some(message) = outgoing else { break; }; + sink.send(message).await?; + } + } + } + Ok::<_, anyhow::Error>(()) + } + .await; + + remove_peer(&rooms, &room_token, connection_id, role).await; + broadcast_presence(&rooms, &room_token).await; + result +} + +async fn register_peer( + rooms: &Rooms, + room_token: &str, + role: CoopRole, + peer: Peer, + max_room_size: usize, + max_rooms: usize, +) -> Result> { + let mut rooms = rooms.lock().await; + if role == CoopRole::Guest && !rooms.contains_key(room_token) { + bail!("room host is not connected"); + } + if !rooms.contains_key(room_token) && rooms.len() >= max_rooms { + bail!("relay room limit reached"); + } + let room = rooms.entry(room_token.to_owned()).or_default(); + let size = room.guests.len() + usize::from(room.host.is_some()); + if size >= max_room_size { + bail!("room is full"); + } + match role { + CoopRole::Host if room.host.is_some() => bail!("room already has a host"), + CoopRole::Host => room.host = Some(peer), + CoopRole::Guest => { + room.guests.insert(peer.connection_id, peer); + } + } + let mut messages = vec![ServerMessage::Ready { + host_present: room.host.is_some(), + guest_count: room.guests.len(), + }]; + if role == CoopRole::Guest + && let Some(snapshot) = room.latest_snapshot.clone() + { + messages.push(ServerMessage::Snapshot { snapshot }); + } + Ok(messages) +} + +async fn route_message( + rooms: &Rooms, + room_token: &str, + connection_id: Uuid, + role: CoopRole, + message: ClientMessage, + own_sender: &mpsc::Sender, +) -> Result<()> { + match (role, message) { + (CoopRole::Host, ClientMessage::Snapshot { snapshot }) => { + let recipients = { + let mut rooms = rooms.lock().await; + let room = rooms.get_mut(room_token).context("room disappeared")?; + if room.host.as_ref().map(|peer| peer.connection_id) != Some(connection_id) { + bail!("connection is no longer the room host"); + } + room.latest_snapshot = Some(snapshot.clone()); + room.guests + .values() + .map(|peer| peer.sender.clone()) + .collect::>() + }; + let message = server_message(&ServerMessage::Snapshot { snapshot })?; + for recipient in recipients { + let _ = recipient.try_send(message.clone()); + } + } + (CoopRole::Guest, ClientMessage::ActionRequest { request_id, action }) => { + let host = rooms + .lock() + .await + .get(room_token) + .and_then(|room| room.host.as_ref()) + .map(|peer| peer.sender.clone()); + if let Some(host) = host { + host.try_send(server_message(&ServerMessage::ActionRequest { + request_id, + action, + })?) + .context("host outbound queue is full")?; + } else { + let _ = own_sender.try_send(server_message(&ServerMessage::Error { + code: "host-unavailable".into(), + message: "the room host is not connected".into(), + })?); + } + } + (_, ClientMessage::Hello { .. }) => { + send_protocol_error(own_sender, "hello may only be sent once")?; + } + (CoopRole::Host, ClientMessage::ActionRequest { .. }) => { + send_protocol_error(own_sender, "hosts cannot send action requests")?; + } + (CoopRole::Guest, ClientMessage::Snapshot { .. }) => { + send_protocol_error(own_sender, "guests cannot publish snapshots")?; + } + } + Ok(()) +} + +async fn remove_peer(rooms: &Rooms, room_token: &str, connection_id: Uuid, role: CoopRole) { + let mut rooms = rooms.lock().await; + let Some(room) = rooms.get_mut(room_token) else { + return; + }; + match role { + CoopRole::Host + if room.host.as_ref().map(|peer| peer.connection_id) == Some(connection_id) => + { + room.host = None; + room.latest_snapshot = None; + } + CoopRole::Guest => { + room.guests.remove(&connection_id); + } + CoopRole::Host => {} + } + if room.host.is_none() && room.guests.is_empty() { + rooms.remove(room_token); + } +} + +async fn broadcast_presence(rooms: &Rooms, room_token: &str) { + let (message, recipients) = { + let rooms = rooms.lock().await; + let Some(room) = rooms.get(room_token) else { + return; + }; + let message = ServerMessage::Presence { + host_present: room.host.is_some(), + guest_count: room.guests.len(), + }; + let recipients = room + .host + .iter() + .chain(room.guests.values()) + .map(|peer| peer.sender.clone()) + .collect::>(); + (message, recipients) + }; + let Ok(message) = server_message(&message) else { + return; + }; + for recipient in recipients { + let _ = recipient.try_send(message.clone()); + } +} + +fn bearer_token(request: &Request) -> Option { + let value = request.headers().get("authorization")?.to_str().ok()?; + let token = value.strip_prefix("Bearer ")?; + valid_room_token(token).then(|| token.to_ascii_lowercase()) +} + +fn handshake_error(status: StatusCode, message: &str) -> ErrorResponse { + let mut response = ErrorResponse::new(Some(message.to_owned())); + *response.status_mut() = status; + response.headers_mut().insert( + "content-type", + HeaderValue::from_static("text/plain; charset=utf-8"), + ); + response +} + +fn parse_client_message(message: Message) -> Result { + let Message::Text(text) = message else { + bail!("client message must be JSON text"); + }; + if text.len() > MAX_MESSAGE_BYTES { + bail!("message exceeds {MAX_MESSAGE_BYTES} bytes"); + } + serde_json::from_str(text.as_ref()).context("invalid client message") +} + +fn server_message(message: &ServerMessage) -> Result { + Ok(Message::Text(serde_json::to_string(message)?.into())) +} + +async fn send_server( + socket: &mut tokio_tungstenite::WebSocketStream, + message: &ServerMessage, +) -> Result<()> +where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin, +{ + socket.send(server_message(message)?).await?; + Ok(()) +} + +fn send_protocol_error(sender: &mpsc::Sender, message: &str) -> Result<()> { + let _ = sender.try_send(server_message(&ServerMessage::Error { + code: "protocol".into(), + message: message.into(), + })?); + Ok(()) +} + +fn init_logging() { + let filter = EnvFilter::try_from_default_env() + .unwrap_or_else(|_| EnvFilter::new("breakd_relay=info,warn")); + tracing_subscriber::fmt() + .with_env_filter(filter) + .with_target(false) + .compact() + .init(); +} + +#[cfg(test)] +mod tests { + use breakd_coop::{CoopAction, CoopPhase, CoopSnapshot, ScheduledBreak}; + use breakd_core::{BreakKind, DueBreakId}; + + use super::*; + + const TOKEN: &str = "0123456789abcdef0123456789abcdef"; + + #[test] + fn room_token_only_comes_from_the_authorization_header() { + let request = Request::builder() + .uri("/ws?room=not-a-room-token") + .header("authorization", format!("Bearer {TOKEN}")) + .body(()) + .unwrap(); + assert_eq!(bearer_token(&request).as_deref(), Some(TOKEN)); + } + + #[tokio::test] + async fn a_room_accepts_only_one_host() { + let rooms = Rooms::default(); + let (sender, _) = mpsc::channel(4); + let first = Peer { + connection_id: Uuid::new_v4(), + sender: sender.clone(), + }; + register_peer(&rooms, TOKEN, CoopRole::Host, first, 8, 256) + .await + .unwrap(); + let second = Peer { + connection_id: Uuid::new_v4(), + sender, + }; + assert!( + register_peer(&rooms, TOKEN, CoopRole::Host, second, 8, 256) + .await + .is_err() + ); + } + + #[tokio::test] + async fn guests_cannot_allocate_rooms_without_a_host() { + let rooms = Rooms::default(); + let (sender, _) = mpsc::channel(4); + let guest = Peer { + connection_id: Uuid::new_v4(), + sender, + }; + assert!( + register_peer(&rooms, TOKEN, CoopRole::Guest, guest, 8, 256) + .await + .is_err() + ); + assert!(rooms.lock().await.is_empty()); + } + + #[tokio::test] + async fn snapshots_and_actions_only_travel_in_the_authoritative_direction() { + let rooms = Rooms::default(); + let (host_sender, mut host_receiver) = mpsc::channel(4); + let host_id = Uuid::new_v4(); + register_peer( + &rooms, + TOKEN, + CoopRole::Host, + Peer { + connection_id: host_id, + sender: host_sender.clone(), + }, + 8, + 256, + ) + .await + .unwrap(); + let (guest_sender, mut guest_receiver) = mpsc::channel(4); + let guest_id = Uuid::new_v4(); + register_peer( + &rooms, + TOKEN, + CoopRole::Guest, + Peer { + connection_id: guest_id, + sender: guest_sender.clone(), + }, + 8, + 256, + ) + .await + .unwrap(); + + let snapshot = CoopSnapshot { + host_id, + revision: 1, + generated_unix_ms: 10, + paused: false, + resume_at_unix_ms: None, + phase: CoopPhase::Working { + next: ScheduledBreak { + due_id: DueBreakId(Uuid::new_v4()), + kind: BreakKind::Mini, + starts_unix_ms: 1_000, + duration_ms: 20_000, + strict_duration_ms: 0, + strict_entire: false, + manual_resume: false, + can_skip: true, + can_postpone: true, + }, + }, + minis_since_long: 0, + longs_since_rest: 0, + postpone_count: 0, + can_skip: false, + can_postpone: false, + }; + route_message( + &rooms, + TOKEN, + host_id, + CoopRole::Host, + ClientMessage::Snapshot { + snapshot: snapshot.clone(), + }, + &host_sender, + ) + .await + .unwrap(); + let Message::Text(forwarded) = guest_receiver.recv().await.unwrap() else { + panic!("snapshot was not forwarded as text"); + }; + assert_eq!( + serde_json::from_str::(forwarded.as_ref()).unwrap(), + ServerMessage::Snapshot { snapshot } + ); + + let request_id = Uuid::new_v4(); + route_message( + &rooms, + TOKEN, + guest_id, + CoopRole::Guest, + ClientMessage::ActionRequest { + request_id, + action: CoopAction::Skip, + }, + &guest_sender, + ) + .await + .unwrap(); + let Message::Text(forwarded) = host_receiver.recv().await.unwrap() else { + panic!("action was not forwarded as text"); + }; + assert_eq!( + serde_json::from_str::(forwarded.as_ref()).unwrap(), + ServerMessage::ActionRequest { + request_id, + action: CoopAction::Skip, + } + ); + } +} diff --git a/crates/scheduler/Cargo.toml b/crates/scheduler/Cargo.toml index 5bdd647..21f86a8 100644 --- a/crates/scheduler/Cargo.toml +++ b/crates/scheduler/Cargo.toml @@ -6,9 +6,11 @@ license.workspace = true rust-version.workspace = true [dependencies] +breakd-coop = { path = "../coop" } breakd-core = { path = "../core" } serde.workspace = true thiserror.workspace = true +uuid.workspace = true [dev-dependencies] breakd-config = { path = "../config" } diff --git a/crates/scheduler/src/lib.rs b/crates/scheduler/src/lib.rs index e0aca01..ad8a411 100644 --- a/crates/scheduler/src/lib.rs +++ b/crates/scheduler/src/lib.rs @@ -1,3 +1,4 @@ +use breakd_coop::{CoopPhase, CoopSnapshot, ScheduledBreak, SharedBreak}; use breakd_core::{ AppConfig, BreakKind, BreakSessionId, ClockSample, Command, DueBreakId, DurationMs, MissedBreakPolicy, OverlaySpec, StrictMode, @@ -12,6 +13,18 @@ pub struct PendingBreak { pub id: DueBreakId, pub kind: BreakKind, pub postpone_count: u32, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub mirrored: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct MirroredBreakPolicy { + pub duration_ms: u64, + pub strict_duration_ms: u64, + pub strict_entire: bool, + pub manual_resume: bool, + pub can_skip: bool, + pub can_postpone: bool, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -337,9 +350,14 @@ impl Scheduler { self.pause(None, now) } } - Command::Status | Command::Reload | Command::Outputs | Command::Doctor => { - Err(SchedulerError::DaemonCommand) - } + Command::Status + | Command::Reload + | Command::Outputs + | Command::Doctor + | Command::CoopHost { .. } + | Command::CoopJoin { .. } + | Command::CoopLeave + | Command::CoopStatus => Err(SchedulerError::DaemonCommand), } } @@ -377,12 +395,245 @@ impl Scheduler { can_skip: !paused && !awaiting_resume && active.is_some_and(|active| { - self.skip_available(active.due.kind) && !self.dismissal_locked(active, now) + active.due.mirrored.as_ref().map_or_else( + || self.skip_available(active.due.kind), + |policy| policy.can_skip, + ) && !self.dismissal_locked(active, now) }), can_postpone: !paused && !awaiting_resume - && active.is_some_and(|active| self.postpone_allowed(active, now)), + && active.is_some_and(|active| { + active.due.mirrored.as_ref().map_or_else( + || self.postpone_allowed(active, now), + |policy| policy.can_postpone, + ) + }), + } + } + + pub fn coop_snapshot( + &mut self, + host_id: uuid::Uuid, + revision: u64, + now: ClockSample, + ) -> (CoopSnapshot, bool) { + let prepared_pending = prepare_pending_for_coop(&mut self.state, &self.config); + let status = self.status(now); + let paused = matches!( + self.state, + SchedulerState::PausedIndefinitely { .. } | SchedulerState::PausedUntil { .. } + ); + let resume_at_unix_ms = match &self.state { + SchedulerState::PausedUntil { resume_mono_ms, .. } => { + Some(wall_deadline_from_mono(*resume_mono_ms, now)) + } + _ => None, + }; + let phase = if matches!( + self.state, + SchedulerState::Suspended { .. } + | SchedulerState::IdleReset { .. } + | SchedulerState::Recovering { .. } + ) { + CoopPhase::Unavailable { + reason: state_name(&self.state).into(), + } + } else if let Some(active) = self.active_break() { + CoopPhase::Break { + active: SharedBreak { + due_id: active.due.id, + session_id: active.session_id, + kind: active.due.kind, + started_unix_ms: wall_deadline_from_boot(active.started_boot_ms, now), + ends_unix_ms: wall_deadline_from_boot(active.ends_boot_ms, now), + strict_until_unix_ms: wall_deadline_from_boot(active.strict_until_boot_ms, now), + strict_entire: self.config.strict.mode == StrictMode::Entire, + manual_resume: active.manual_resume, + completion_sound_emitted: active.completion_sound_emitted, + can_skip: self.skip_available(active.due.kind), + can_postpone: self.postpone_available(active), + }, + } + } else if let Some(context) = self.context() { + let due = context + .pending + .as_ref() + .expect("co-op snapshot prepares a pending break"); + let duration_ms = self.break_duration(due.kind); + CoopPhase::Working { + next: ScheduledBreak { + due_id: due.id, + kind: due.kind, + starts_unix_ms: wall_deadline_from_mono(context.next_due_mono_ms, now), + duration_ms, + strict_duration_ms: self.strict_duration(duration_ms), + strict_entire: self.config.strict.mode == StrictMode::Entire, + manual_resume: self.config.completion.manual_resume, + can_skip: self.skip_available(due.kind), + can_postpone: self.postpone_available_for(due.kind, due.postpone_count), + }, + } + } else { + CoopPhase::Unavailable { + reason: state_name(&self.state).into(), + } + }; + + ( + CoopSnapshot { + host_id, + revision, + generated_unix_ms: now.wall_unix_ms, + paused, + resume_at_unix_ms, + phase, + minis_since_long: status.minis_since_long, + longs_since_rest: status.longs_since_rest, + postpone_count: status.postpone_count, + can_skip: status.can_skip, + can_postpone: status.can_postpone, + }, + prepared_pending, + ) + } + + pub fn adopt_coop_snapshot( + &mut self, + snapshot: &CoopSnapshot, + now: ClockSample, + ) -> Vec { + let old_visible = visible_active(&self.state).map(|active| active.session_id); + let context = |pending: Option, next_due_mono_ms: u64| ScheduleContext { + cycle_started_mono_ms: now.monotonic_ms, + next_due_mono_ms, + minis_since_long: snapshot.minis_since_long, + longs_since_rest: snapshot.longs_since_rest, + rest_cycle_started_mono_ms: now.monotonic_ms, + pending, + }; + + let base = match &snapshot.phase { + CoopPhase::Working { next } => { + let policy = MirroredBreakPolicy { + duration_ms: next.duration_ms, + strict_duration_ms: next.strict_duration_ms, + strict_entire: next.strict_entire, + manual_resume: next.manual_resume, + can_skip: next.can_skip, + can_postpone: next.can_postpone, + }; + let due = PendingBreak { + id: next.due_id, + kind: next.kind, + postpone_count: snapshot.postpone_count, + mirrored: Some(policy), + }; + if now.wall_unix_ms >= next.starts_unix_ms { + let active = ActiveBreak { + context: context(None, now.monotonic_ms), + due, + session_id: BreakSessionId(next.due_id.0), + started_boot_ms: boot_deadline_from_wall(next.starts_unix_ms, now), + ends_boot_ms: boot_deadline_from_wall( + next.starts_unix_ms.saturating_add(next.duration_ms), + now, + ), + strict_until_boot_ms: boot_deadline_from_wall( + next.starts_unix_ms.saturating_add(next.strict_duration_ms), + now, + ), + manual_resume: next.manual_resume, + completion_sound_emitted: false, + }; + match next.kind { + BreakKind::Mini => SchedulerState::MiniBreak { active }, + BreakKind::Long => SchedulerState::LongBreak { active }, + BreakKind::Rest => SchedulerState::RestBreak { active }, + } + } else { + SchedulerState::Running { + context: context( + Some(due), + mono_deadline_from_wall(next.starts_unix_ms, now), + ), + } + } + } + CoopPhase::Break { active } => { + let ends_boot_ms = boot_deadline_from_wall(active.ends_unix_ms, now); + let strict_until_boot_ms = + boot_deadline_from_wall(active.strict_until_unix_ms, now); + let duration_ms = active.ends_unix_ms.saturating_sub(active.started_unix_ms); + let strict_duration_ms = active + .strict_until_unix_ms + .saturating_sub(active.started_unix_ms); + let active = ActiveBreak { + context: context(None, now.monotonic_ms), + due: PendingBreak { + id: active.due_id, + kind: active.kind, + postpone_count: snapshot.postpone_count, + mirrored: Some(MirroredBreakPolicy { + duration_ms, + strict_duration_ms, + strict_entire: active.strict_entire, + manual_resume: active.manual_resume, + can_skip: active.can_skip, + can_postpone: active.can_postpone, + }), + }, + session_id: active.session_id, + started_boot_ms: boot_deadline_from_wall(active.started_unix_ms, now), + ends_boot_ms, + strict_until_boot_ms, + manual_resume: active.manual_resume, + completion_sound_emitted: active.completion_sound_emitted, + }; + match active.due.kind { + BreakKind::Mini => SchedulerState::MiniBreak { active }, + BreakKind::Long => SchedulerState::LongBreak { active }, + BreakKind::Rest => SchedulerState::RestBreak { active }, + } + } + CoopPhase::Unavailable { reason } => SchedulerState::Recovering { + reason: format!("co-op host unavailable: {reason}"), + }, + }; + self.state = if snapshot.paused { + match snapshot.resume_at_unix_ms { + Some(resume_at) => SchedulerState::PausedUntil { + inner: Box::new(base), + paused_at: now, + resume_mono_ms: mono_deadline_from_wall(resume_at, now), + }, + None => SchedulerState::PausedIndefinitely { + inner: Box::new(base), + paused_at: now, + }, + } + } else { + base + }; + self.last_clock = now; + + let new_visible = visible_active(&self.state); + let mut effects = Vec::with_capacity(2); + if old_visible != new_visible.map(|active| active.session_id) { + if let Some(session_id) = old_visible { + effects.push(Effect::StopOverlay { session_id }); + } + if let Some(active) = new_visible + && (active.manual_resume || active.ends_boot_ms > now.boottime_ms) + { + effects.push(Effect::StartOverlay(self.overlay_spec(active, now))); + } } + effects + } + + pub fn reset_after_coop_disconnect(&mut self, now: ClockSample) -> Vec { + self.last_clock = now; + self.reset(now) } fn fresh_running(config: &AppConfig, now: ClockSample) -> SchedulerState { @@ -414,6 +665,7 @@ impl Scheduler { id: DueBreakId::new(), kind, postpone_count: 0, + mirrored: None, }); if now.monotonic_ms >= context.next_due_mono_ms { return self.begin_scheduled_break(kind, context, now); @@ -503,6 +755,7 @@ impl Scheduler { id: DueBreakId::new(), kind, postpone_count: 0, + mirrored: None, }); self.begin_break(kind, context, due, now) } @@ -527,6 +780,7 @@ impl Scheduler { id: DueBreakId::new(), kind, postpone_count: 0, + mirrored: None, }; Ok(self.begin_break(kind, context, due, now)) } @@ -538,24 +792,32 @@ impl Scheduler { due: PendingBreak, now: ClockSample, ) -> Vec { - let duration = self.break_duration(kind); + let duration = due + .mirrored + .as_ref() + .map_or_else(|| self.break_duration(kind), |policy| policy.duration_ms); let ends_boot_ms = now.boottime_ms.saturating_add(duration); - let strict_until_boot_ms = match self.config.strict.mode { - StrictMode::Off => now.boottime_ms, - StrictMode::Delay => now - .boottime_ms - .saturating_add(self.config.strict.minimum_visible.as_millis()) - .min(ends_boot_ms), - StrictMode::Entire => ends_boot_ms, - }; + let strict_until_boot_ms = + now.boottime_ms + .saturating_add(due.mirrored.as_ref().map_or_else( + || self.strict_duration(duration), + |policy| policy.strict_duration_ms, + )); + let manual_resume = due + .mirrored + .as_ref() + .map_or(self.config.completion.manual_resume, |policy| { + policy.manual_resume + }); + let session_id = BreakSessionId(due.id.0); let active = ActiveBreak { context, due, - session_id: BreakSessionId::new(), + session_id, started_boot_ms: now.boottime_ms, ends_boot_ms, strict_until_boot_ms, - manual_resume: self.config.completion.manual_resume, + manual_resume, completion_sound_emitted: false, }; let effect = Effect::StartOverlay(self.overlay_spec(&active, now)); @@ -847,6 +1109,19 @@ impl Scheduler { } } + fn strict_duration(&self, break_duration: u64) -> u64 { + match self.config.strict.mode { + StrictMode::Off => 0, + StrictMode::Delay => self + .config + .strict + .minimum_visible + .as_millis() + .min(break_duration), + StrictMode::Entire => break_duration, + } + } + fn overlay_spec(&self, active: &ActiveBreak, now: ClockSample) -> OverlaySpec { let message = if self.config.content.show_message && !self.config.content.messages.is_empty() { @@ -863,14 +1138,27 @@ impl Scheduler { kind: active.due.kind, duration: DurationMs::from_millis(active.ends_boot_ms.saturating_sub(now.boottime_ms)), strict_remaining: DurationMs::from_millis( - if self.config.strict.mode == StrictMode::Entire { + if active + .due + .mirrored + .as_ref() + .map_or(self.config.strict.mode == StrictMode::Entire, |policy| { + policy.strict_entire + }) + { active.strict_until_boot_ms.saturating_sub(now.boottime_ms) } else { 0 }, ), - can_skip: self.skip_available(active.due.kind), - can_postpone: self.postpone_available(active), + can_skip: active.due.mirrored.as_ref().map_or_else( + || self.skip_available(active.due.kind), + |policy| policy.can_skip, + ), + can_postpone: active.due.mirrored.as_ref().map_or_else( + || self.postpone_available(active), + |policy| policy.can_postpone, + ), manual_resume: active.manual_resume, message, socket_path: self.socket_path.clone(), @@ -896,12 +1184,22 @@ impl Scheduler { /// break; the `Delay` mode's minimum-visible window intentionally does not /// gate skip or postpone, so the user can act on the break immediately. fn dismissal_locked(&self, active: &ActiveBreak, now: ClockSample) -> bool { - self.config.strict.mode == StrictMode::Entire + active + .due + .mirrored + .as_ref() + .map_or(self.config.strict.mode == StrictMode::Entire, |policy| { + policy.strict_entire + }) && now.boottime_ms < active.strict_until_boot_ms } fn postpone_available(&self, active: &ActiveBreak) -> bool { - let rule = match active.due.kind { + self.postpone_available_for(active.due.kind, active.due.postpone_count) + } + + fn postpone_available_for(&self, kind: BreakKind, postpone_count: u32) -> bool { + let rule = match kind { BreakKind::Mini => &self.config.postpone.mini, BreakKind::Long => &self.config.postpone.long, BreakKind::Rest => &self.config.postpone.rest, @@ -909,7 +1207,7 @@ impl Scheduler { rule.enabled && rule .max_postponements - .is_none_or(|maximum| active.due.postpone_count < maximum) + .is_none_or(|maximum| postpone_count < maximum) } fn active_break(&self) -> Option<&ActiveBreak> { @@ -933,6 +1231,84 @@ fn active_break_in(state: &SchedulerState) -> Option<&ActiveBreak> { } } +fn visible_active(state: &SchedulerState) -> Option<&ActiveBreak> { + match state { + SchedulerState::MiniBreak { active } + | SchedulerState::LongBreak { active } + | SchedulerState::RestBreak { active } => Some(active), + _ => None, + } +} + +fn prepare_pending_for_coop(state: &mut SchedulerState, config: &AppConfig) -> bool { + match state { + SchedulerState::Running { context } + | SchedulerState::PreMiniBreak { context } + | SchedulerState::PreLongBreak { context } + | SchedulerState::PreRestBreak { context } => { + if context.pending.is_some() { + return false; + } + let rest_elapsed = context + .next_due_mono_ms + .saturating_sub(context.rest_cycle_started_mono_ms) + >= config.schedule.rest.interval.as_millis(); + let long_elapsed = context + .next_due_mono_ms + .saturating_sub(context.cycle_started_mono_ms) + >= config.schedule.long.interval.as_millis(); + let kind = if rest_elapsed + && context.longs_since_rest >= config.schedule.rest.after_longs + { + BreakKind::Rest + } else if long_elapsed && context.minis_since_long >= config.schedule.long.after_minis { + BreakKind::Long + } else { + BreakKind::Mini + }; + context.pending = Some(PendingBreak { + id: DueBreakId::new(), + kind, + postpone_count: 0, + mirrored: None, + }); + true + } + SchedulerState::PausedIndefinitely { inner, .. } + | SchedulerState::PausedUntil { inner, .. } => prepare_pending_for_coop(inner, config), + SchedulerState::MiniBreak { .. } + | SchedulerState::LongBreak { .. } + | SchedulerState::RestBreak { .. } + | SchedulerState::Suspended { .. } + | SchedulerState::IdleReset { .. } + | SchedulerState::Recovering { .. } => false, + } +} + +fn wall_deadline_from_mono(deadline_mono_ms: u64, now: ClockSample) -> u64 { + translate_deadline(deadline_mono_ms, now.monotonic_ms, now.wall_unix_ms) +} + +fn wall_deadline_from_boot(deadline_boot_ms: u64, now: ClockSample) -> u64 { + translate_deadline(deadline_boot_ms, now.boottime_ms, now.wall_unix_ms) +} + +fn mono_deadline_from_wall(deadline_wall_ms: u64, now: ClockSample) -> u64 { + translate_deadline(deadline_wall_ms, now.wall_unix_ms, now.monotonic_ms) +} + +fn boot_deadline_from_wall(deadline_wall_ms: u64, now: ClockSample) -> u64 { + translate_deadline(deadline_wall_ms, now.wall_unix_ms, now.boottime_ms) +} + +fn translate_deadline(deadline: u64, source_now: u64, destination_now: u64) -> u64 { + if deadline >= source_now { + destination_now.saturating_add(deadline - source_now) + } else { + destination_now.saturating_sub(source_now - deadline) + } +} + fn context_in(state: &SchedulerState) -> Option<&ScheduleContext> { match state { SchedulerState::Running { context } @@ -1218,6 +1594,67 @@ mod tests { ); } + #[test] + fn coop_guests_start_the_same_session_with_the_hosts_duration() { + let mut host = test_scheduler(); + let mut guest = test_scheduler(); + let mut guest_config = guest.config.clone(); + guest_config.schedule.mini.duration = DurationMs::from_millis(800); + guest.replace_config(guest_config); + + let (snapshot, _) = host.coop_snapshot(uuid::Uuid::nil(), 1, clock(0)); + assert!(guest.adopt_coop_snapshot(&snapshot, clock(50)).is_empty()); + + let host_effects = host.handle_event(SchedulerEvent::Tick, clock(1_000)); + let guest_effects = guest.handle_event(SchedulerEvent::Tick, clock(1_000)); + let host_session = host.status(clock(1_000)).active_session.unwrap(); + let guest_status = guest.status(clock(1_000)); + + assert_eq!(guest_status.active_session, Some(host_session)); + assert_eq!(guest_status.remaining_ms, Some(100)); + assert!(matches!(host_effects.as_slice(), [Effect::StartOverlay(_)])); + assert!(matches!( + guest_effects.as_slice(), + [Effect::StartOverlay(_)] + )); + } + + #[test] + fn a_working_snapshot_arriving_just_after_deadline_starts_without_flicker() { + let mut host = test_scheduler(); + let mut guest = test_scheduler(); + let (snapshot, _) = host.coop_snapshot(uuid::Uuid::nil(), 1, clock(0)); + + let effects = guest.adopt_coop_snapshot(&snapshot, clock(1_001)); + let [Effect::StartOverlay(spec)] = effects.as_slice() else { + panic!("expected a late working snapshot to start its break"); + }; + let CoopPhase::Working { next } = &snapshot.phase else { + panic!("host should still be working"); + }; + assert_eq!(spec.session_id.0, next.due_id.0); + assert_eq!(spec.duration, DurationMs::from_millis(99)); + } + + #[test] + fn coop_guest_uses_its_own_break_message_when_joining_mid_break() { + let mut host = test_scheduler(); + host.handle_command(&Command::Mini, clock(0)).unwrap(); + let (snapshot, _) = host.coop_snapshot(uuid::Uuid::nil(), 1, clock(20)); + + let mut guest = test_scheduler(); + let mut guest_config = guest.config.clone(); + guest_config.content.messages = vec!["A local message".into()]; + guest.replace_config(guest_config); + let effects = guest.adopt_coop_snapshot(&snapshot, clock(20)); + + let [Effect::StartOverlay(spec)] = effects.as_slice() else { + panic!("expected the mirrored overlay to start"); + }; + assert_eq!(spec.message.as_deref(), Some("A local message")); + assert_eq!(spec.duration, DurationMs::from_millis(80)); + } + #[test] fn idle_does_not_dismiss_a_manual_resume_break() { let mut scheduler = test_scheduler(); diff --git a/docs/coop.md b/docs/coop.md new file mode 100644 index 0000000..a0ae211 --- /dev/null +++ b/docs/coop.md @@ -0,0 +1,107 @@ +# Co-op deployment and protocol + +Co-op deliberately has two small pieces: + +- the existing `breakd` daemon is either a host or a guest; +- `breakd-relay` authenticates room connections and forwards bounded JSON + WebSocket messages. + +The relay never runs a schedule and never decides whether an action is allowed. +It accepts one host per room, forwards snapshots only from that host, and +forwards action requests only from guests. The host's normal scheduler validates +every action. A room disappears when its final connection closes, and its cached +snapshot is cleared as soon as its host disconnects. + +## Run a relay + +For a local test: + +```bash +cargo run -p breakd-relay -- --listen 127.0.0.1:8787 +breakd coop host --relay ws://127.0.0.1:8787/ws +``` + +Plain `ws://` exposes the room token to the network. Use it only on localhost or +another trusted, encrypted network. For internet use, keep the process bound to +localhost and put a TLS reverse proxy in front of it. For example, a Caddy site +can proxy WebSockets without relay-specific headers: + +```caddyfile +breaks.example.net { + reverse_proxy 127.0.0.1:8787 +} +``` + +Then use `wss://breaks.example.net/ws` as the relay URL. The relay accepts any +request path, so a proxy can dedicate `/ws` or an entire hostname to it. + +On NixOS, the flake module can run the isolated relay service: + +```nix +{ + services.breakd-relay = { + enable = true; + listen = "127.0.0.1:8787"; + maxRoomSize = 8; + maxRooms = 256; + }; +} +``` + +The service uses a dynamic user and filesystem hardening. It intentionally does +not open a firewall port or terminate TLS. + +## Room lifecycle + +The host creates a random UUID room token and persists its mode, relay +URL, and token in `~/.config/breakd/config.toml`, which breakd writes with mode +`0600`. The printed invite has this form: + +```text +wss://breaks.example.net/ws#breakd=0123456789abcdef0123456789abcdef +``` + +URL fragments are not included in HTTP or WebSocket requests. The joining daemon +separates the fragment and sends the token as `Authorization: Bearer ...` during +the WebSocket upgrade. Reverse-proxy access logs therefore do not normally +contain the room token. Avoid putting full invites in shell history, screenshots, +or public chat. + +Run `breakd coop host` again to make a new token and invalidate the previous room +from that host. `breakd coop leave` clears the relay URL and token and resets a +fresh local schedule. + +## Synchronization and failure behavior + +The host publishes at most one regular snapshot per second and immediately after +local or guest-requested actions. A working snapshot includes the next break's +absolute Unix start time, duration, type, stable due ID, and strict/manual-resume +policy. That lets both schedulers start the same session locally without waiting +for a round trip at the deadline. Active-break and pause snapshots let a guest +join midway through a room. + +Guests use their own display, content, sound, and monitor configuration. They do +not run idle, lock, or suspend transitions while following the host. Native +WebSocket ping frames keep the one connection alive, and reconnects use bounded +exponential backoff. The client rejects stale revisions and messages larger than +128 KiB. + +When the configured `coop.disconnect_grace` elapses without a snapshot (10 +seconds by default), a guest discards the mirrored state and begins a fresh local +schedule. This avoids leaving the user with a frozen or overdue remote break. +If the host later returns, the next valid snapshot becomes authoritative again. + +Absolute deadlines assume both computers have ordinary NTP-style clock +synchronization. Network and scheduler tick jitter mean the overlays are not a +hard real-time barrier, but under normal clocks they begin within the local +250 ms scheduler tick rather than one relay round trip apart. + +## Resource model + +When co-op is off, no network task or connection is created. Enabling it adds one +Tokio task, one WebSocket, bounded event and action buffers, and a latest-value +snapshot channel. The standalone relay uses one task and small outbound queue per +connection, stores only one snapshot per live room, and caps rooms at 8 +connections by default (configurable from 2 to 64). It has no GTK, Wayland, +database, or account-system dependency. Guests cannot create empty rooms, and +the relay also caps simultaneously live rooms at 256 by default. diff --git a/flake.lock b/flake.lock new file mode 100644 index 0000000..fb043dc --- /dev/null +++ b/flake.lock @@ -0,0 +1,27 @@ +{ + "nodes": { + "nixpkgs": { + "locked": { + "lastModified": 1784007870, + "narHash": "sha256-djcLt/JJphyNt4eDY9XTly+/WbCK5lqWq9lSgCmJkkQ=", + "owner": "NixOS", + "repo": "nixpkgs", + "rev": "18b9261cb3294b6d2a06d03f96872827b8fe2698", + "type": "github" + }, + "original": { + "owner": "NixOS", + "ref": "nixos-unstable", + "repo": "nixpkgs", + "type": "github" + } + }, + "root": { + "inputs": { + "nixpkgs": "nixpkgs" + } + } + }, + "root": "root", + "version": 7 +} diff --git a/flake.nix b/flake.nix new file mode 100644 index 0000000..d54ec28 --- /dev/null +++ b/flake.nix @@ -0,0 +1,46 @@ +{ + description = "Wayland-native break reminder with multi-monitor overlays"; + + inputs.nixpkgs.url = "github:NixOS/nixpkgs/nixos-unstable"; + + outputs = + { self, nixpkgs }: + let + supportedSystems = [ + "x86_64-linux" + "aarch64-linux" + ]; + forAllSystems = nixpkgs.lib.genAttrs supportedSystems; + packagesFor = + system: + let + pkgs = import nixpkgs { inherit system; }; + in + { + breakd = pkgs.callPackage ./packaging/nix/package.nix { }; + breakd-relay = pkgs.callPackage ./packaging/nix/relay.nix { }; + }; + in + { + packages = forAllSystems (system: { + inherit (packagesFor system) breakd breakd-relay; + default = self.packages.${system}.breakd; + }); + + checks = forAllSystems (system: { + inherit (self.packages.${system}) breakd breakd-relay; + }); + + overlays.default = final: _prev: { + breakd = final.callPackage ./packaging/nix/package.nix { }; + breakd-relay = final.callPackage ./packaging/nix/relay.nix { }; + }; + + nixosModules = { + breakd = import ./packaging/nix/module.nix { inherit self; }; + default = self.nixosModules.breakd; + }; + + formatter = forAllSystems (system: nixpkgs.legacyPackages.${system}.nixfmt); + }; +} diff --git a/packaging/nix/module.nix b/packaging/nix/module.nix new file mode 100644 index 0000000..e9f82b6 --- /dev/null +++ b/packaging/nix/module.nix @@ -0,0 +1,96 @@ +{ self }: +{ + config, + lib, + pkgs, + ... +}: + +let + cfg = config.services.breakd; + relayCfg = config.services.breakd-relay; + defaultPackage = self.packages.${pkgs.stdenv.hostPlatform.system}.breakd; + defaultRelayPackage = self.packages.${pkgs.stdenv.hostPlatform.system}.breakd-relay; +in +{ + options.services.breakd = { + enable = lib.mkEnableOption "the breakd Wayland break reminder"; + + package = lib.mkOption { + type = lib.types.package; + default = defaultPackage; + defaultText = lib.literalExpression "inputs.breakd.packages.${pkgs.stdenv.hostPlatform.system}.breakd"; + description = "The breakd package to run."; + }; + }; + + options.services.breakd-relay = { + enable = lib.mkEnableOption "the breakd co-op relay"; + + package = lib.mkOption { + type = lib.types.package; + default = defaultRelayPackage; + defaultText = lib.literalExpression "inputs.breakd.packages.${pkgs.stdenv.hostPlatform.system}.breakd-relay"; + description = "The standalone breakd relay package to run."; + }; + + listen = lib.mkOption { + type = lib.types.str; + default = "127.0.0.1:8787"; + description = "Relay listen address. Keep it private and terminate TLS in a reverse proxy."; + }; + + maxRoomSize = lib.mkOption { + type = lib.types.ints.between 2 64; + default = 8; + description = "Maximum number of connections in one co-op room."; + }; + + maxRooms = lib.mkOption { + type = lib.types.ints.between 1 65536; + default = 256; + description = "Maximum number of simultaneously live co-op rooms."; + }; + }; + + config = lib.mkMerge [ + (lib.mkIf cfg.enable { + environment.systemPackages = [ cfg.package ]; + + systemd.user.services.breakd = { + description = "Wayland-native break reminder"; + partOf = [ "graphical-session.target" ]; + after = [ "graphical-session.target" ]; + wantedBy = [ "graphical-session.target" ]; + environment.GDK_BACKEND = "wayland"; + serviceConfig = { + Type = "simple"; + ExecStart = "${lib.getExe cfg.package} daemon"; + Restart = "on-failure"; + RestartSec = "2s"; + UMask = "0077"; + }; + }; + }) + (lib.mkIf relayCfg.enable { + environment.systemPackages = [ relayCfg.package ]; + + systemd.services.breakd-relay = { + description = "breakd co-op relay"; + wantedBy = [ "multi-user.target" ]; + after = [ "network.target" ]; + serviceConfig = { + Type = "simple"; + ExecStart = "${lib.getExe relayCfg.package} --listen ${relayCfg.listen} --max-room-size ${toString relayCfg.maxRoomSize} --max-rooms ${toString relayCfg.maxRooms}"; + Restart = "on-failure"; + RestartSec = "2s"; + DynamicUser = true; + NoNewPrivileges = true; + PrivateTmp = true; + ProtectHome = true; + ProtectSystem = "strict"; + }; + }; + }) + ]; +} diff --git a/packaging/nix/package.nix b/packaging/nix/package.nix new file mode 100644 index 0000000..78a999d --- /dev/null +++ b/packaging/nix/package.nix @@ -0,0 +1,66 @@ +{ + lib, + rustPlatform, + pkg-config, + wrapGAppsHook4, + gtk4, + gtk4-layer-shell, + libcanberra, + wayland, +}: + +let + manifest = builtins.fromTOML (builtins.readFile ../../Cargo.toml); +in +rustPlatform.buildRustPackage { + pname = "breakd"; + inherit (manifest.package) version; + + src = ../..; + cargoLock.lockFile = ../../Cargo.lock; + + strictDeps = true; + nativeBuildInputs = [ + pkg-config + wrapGAppsHook4 + ]; + buildInputs = [ + gtk4 + gtk4-layer-shell + libcanberra + wayland + ]; + + cargoBuildFlags = [ + "-p" + "breakd" + ]; + cargoTestFlags = [ "--workspace" ]; + + postInstall = '' + install -Dm644 crates/platform-linux/assets/*.oga -t "$out/share/breakd" + install -Dm644 packaging/io.github.simonwinther.breakd.settings.desktop \ + "$out/share/applications/io.github.simonwinther.breakd.settings.desktop" + install -Dm644 config.example.toml "$out/share/doc/breakd/config.example.toml" + install -Dm644 README.md "$out/share/doc/breakd/README.md" + install -Dm644 LICENSE "$out/share/licenses/breakd/LICENSE" + install -Dm644 THIRD_PARTY_NOTICES.md \ + "$out/share/licenses/breakd/THIRD_PARTY_NOTICES.md" + + install -Dm644 packaging/systemd/breakd.service \ + "$out/lib/systemd/user/breakd.service" + substituteInPlace "$out/lib/systemd/user/breakd.service" \ + --replace-fail /usr/bin/breakd "$out/bin/breakd" + ''; + + meta = { + description = "Wayland-native break reminder with multi-monitor overlays"; + homepage = "https://github.com/simonwinther/breakd"; + license = [ + lib.licenses.mit + lib.licenses.bsd2 + ]; + mainProgram = "breakd"; + platforms = lib.platforms.linux; + }; +} diff --git a/packaging/nix/relay.nix b/packaging/nix/relay.nix new file mode 100644 index 0000000..fb0ccd4 --- /dev/null +++ b/packaging/nix/relay.nix @@ -0,0 +1,32 @@ +{ + lib, + rustPlatform, +}: + +let + manifest = builtins.fromTOML (builtins.readFile ../../Cargo.toml); +in +rustPlatform.buildRustPackage { + pname = "breakd-relay"; + inherit (manifest.package) version; + + src = ../..; + cargoLock.lockFile = ../../Cargo.lock; + + cargoBuildFlags = [ + "-p" + "breakd-relay" + ]; + cargoTestFlags = [ + "-p" + "breakd-relay" + ]; + + meta = { + description = "Small authenticated WebSocket relay for breakd co-op rooms"; + homepage = "https://github.com/simonwinther/breakd"; + license = lib.licenses.mit; + mainProgram = "breakd-relay"; + platforms = lib.platforms.linux; + }; +} diff --git a/src/main.rs b/src/main.rs index 3042248..a8713ca 100644 --- a/src/main.rs +++ b/src/main.rs @@ -53,6 +53,11 @@ enum CliCommand { #[arg(long)] json: bool, }, + /// Host, join, or inspect a synchronized co-op room. + Coop { + #[command(subcommand)] + command: CoopCommand, + }, /// Open the graphical configuration window. Settings, #[command(hide = true)] @@ -61,6 +66,25 @@ enum CliCommand { ExampleConfig, } +#[derive(Debug, Subcommand)] +enum CoopCommand { + /// Create a room and print a secret invite. + Host { + /// Public ws:// or wss:// relay endpoint. + #[arg(long)] + relay: String, + }, + /// Join a room using the invite printed by its host. + Join { invite: String }, + /// Disconnect and return to a fresh local schedule. + Leave, + /// Show relay connection and room state. + Status { + #[arg(long)] + json: bool, + }, +} + #[tokio::main(flavor = "multi_thread")] async fn main() -> ExitCode { init_logging(); @@ -101,6 +125,14 @@ async fn execute(arguments: Arguments) -> Result<()> { CliCommand::Reload => send(Command::Reload, false).await, CliCommand::Outputs { json } => send(Command::Outputs, json).await, CliCommand::Doctor { json } => send(Command::Doctor, json).await, + CliCommand::Coop { command } => match command { + CoopCommand::Host { relay } => { + send(Command::CoopHost { relay_url: relay }, false).await + } + CoopCommand::Join { invite } => send(Command::CoopJoin { invite }, false).await, + CoopCommand::Leave => send(Command::CoopLeave, false).await, + CoopCommand::Status { json } => send(Command::CoopStatus, json).await, + }, CliCommand::Settings => breakd_settings::run().map_err(anyhow::Error::msg), } }