diff --git a/Cargo.lock b/Cargo.lock index 1873abb..46a7d6b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -32,6 +32,12 @@ dependencies = [ "as-slice", ] +[[package]] +name = "allocator-api2" +version = "0.2.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" + [[package]] name = "anstream" version = "0.6.21" @@ -148,7 +154,7 @@ dependencies = [ "smlang", "static_cell", "stm32-fmc", - "strum", + "strum 0.27.2", "uor-high-level", "uor-peripherals", "uor-utils", @@ -256,12 +262,43 @@ version = "1.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1e748733b7cbc798e1434b6ac524f0c1ff2ab456fe201501e6497c8417a4fc33" +[[package]] +name = "cassowary" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df8670b8c7b9dae1793364eafadf7239c40d669904660c5960d74cfd80b46a53" + +[[package]] +name = "castaway" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dec551ab6e7578819132c713a93c022a05d60159dc86e7a7050223577484c55a" +dependencies = [ + "rustversion", +] + +[[package]] +name = "cc" +version = "1.2.67" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e17dd265a7d0f31ef544e1b20e03add05d3b45b491b633b10d67145d2acc1a38" +dependencies = [ + "find-msvc-tools", + "shlex", +] + [[package]] name = "cfg-if" version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9555578bc9e57714c812a1f84e4fc5b4d21fcb063490c624de019f7464c91268" +[[package]] +name = "cfg_aliases" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fd16c4719339c4530435d38e511904438d07cce7950afa3718a84ac36c10e89e" + [[package]] name = "cfg_aliases" version = "0.2.1" @@ -331,6 +368,19 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75" +[[package]] +name = "compact_str" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f86b9c4c00838774a6d902ef931eff7470720c51d90c2e32cfe15dc304737b3f" +dependencies = [ + "castaway", + "cfg-if", + "itoa", + "ryu", + "static_assertions", +] + [[package]] name = "const-default" version = "1.0.0" @@ -407,6 +457,62 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" +[[package]] +name = "crossbeam-channel" +version = "0.5.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e" +dependencies = [ + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-utils" +version = "0.8.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61803da095bee82a81bb1a452ecc25d3b2f1416d1897eb86430c6159ef717c17" + +[[package]] +name = "crossterm" +version = "0.26.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a84cda67535339806297f1b331d6dd6320470d2a0fe65381e79ee9e156dd3d13" +dependencies = [ + "bitflags 1.3.2", + "crossterm_winapi", + "libc", + "mio 0.8.11", + "parking_lot", + "signal-hook", + "signal-hook-mio", + "winapi", +] + +[[package]] +name = "crossterm" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f476fe445d41c9e991fd07515a6f463074b782242ccf4a5b7b1d1012e70824df" +dependencies = [ + "bitflags 2.10.0", + "crossterm_winapi", + "libc", + "mio 0.8.11", + "parking_lot", + "signal-hook", + "signal-hook-mio", + "winapi", +] + +[[package]] +name = "crossterm_winapi" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "acdd7c62a3665c7f6830a51635d9ac9b23ed385797f70a83bb8bafe9c572ab2b" +dependencies = [ + "winapi", +] + [[package]] name = "csv-core" version = "0.1.12" @@ -462,6 +568,19 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "dashmap" +version = "5.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "978747c1d849a7d2ee5e8adc0159961c48fb7e5db2f06af6723b80123bb53856" +dependencies = [ + "cfg-if", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "defmt" version = "0.3.100" @@ -1070,6 +1189,12 @@ version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" +[[package]] +name = "find-msvc-tools" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582" + [[package]] name = "fixedbitset" version = "0.5.7" @@ -1082,12 +1207,29 @@ version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" + [[package]] name = "futures-core" version = "0.3.31" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "05f29059c0c2090612e8d742178b0580d2dc940c837851ad723096f87af6663e" +[[package]] +name = "futures-macro" +version = "0.3.31" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "162ee34ebcb7c64a8abebc059ce0fee27c2262618d7b60ed8faf72fef13c3650" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + [[package]] name = "futures-sink" version = "0.3.31" @@ -1107,9 +1249,22 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9fa08315bb612088cc391249efdc3bc77536f16c91f6cf495e6fbe85b20a4a81" dependencies = [ "futures-core", + "futures-macro", "futures-task", "pin-project-lite", "pin-utils", + "slab", +] + +[[package]] +name = "getrandom" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" +dependencies = [ + "cfg-if", + "libc", + "wasi 0.11.1+wasi-snapshot-preview1", ] [[package]] @@ -1121,7 +1276,7 @@ dependencies = [ "cfg-if", "libc", "r-efi", - "wasi", + "wasi 0.14.2+wasi-0.2.4", ] [[package]] @@ -1166,6 +1321,23 @@ dependencies = [ "ahash", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" + +[[package]] +name = "hashbrown" +version = "0.15.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" +dependencies = [ + "allocator-api2", + "equivalent", + "foldhash", +] + [[package]] name = "hashbrown" version = "0.16.1" @@ -1213,12 +1385,33 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hostname" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c731c3e10504cc8ed35cfe2f1db4c9274c3d35fa486e3b31df46f068ef3e867" +dependencies = [ + "libc", + "match_cfg", + "winapi", +] + [[package]] name = "ident_case" version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b9e0384b61958566e926dc50660321d12159025e767c18e043daf26b70104c39" +[[package]] +name = "if-addrs" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cbc0fa01ffc752e9dbc72818cdb072cd028b86be5e09dd04c5a643704fe101a9" +dependencies = [ + "libc", + "winapi", +] + [[package]] name = "indexmap" version = "2.13.0" @@ -1254,6 +1447,24 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "itertools" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba291022dbbd398a455acf126c1e341954079855bc60dfdda641363bd6922569" +dependencies = [ + "either", +] + +[[package]] +name = "itertools" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "413ee7dfc52ee1a4949ceeb7dbc8a33f2d6c088194d9f922fb8318faf1f01186" +dependencies = [ + "either", +] + [[package]] name = "itertools" version = "0.14.0" @@ -1320,6 +1531,26 @@ version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" +[[package]] +name = "libmdns" +version = "0.7.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b04ae6b56b3b19ade26f0e7e7c1360a1713514f326c5ed0797cf2c109c9e010" +dependencies = [ + "byteorder", + "futures-util", + "hostname", + "if-addrs", + "log", + "multimap 0.8.3", + "nix 0.23.2", + "rand 0.8.7", + "socket2 0.4.10", + "thiserror 1.0.69", + "tokio", + "winapi", +] + [[package]] name = "libudev" version = "0.3.0" @@ -1373,6 +1604,15 @@ version = "0.4.27" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "13dc2df351e3202783a1fe0d44375f7295ffb4049267b0f3018346dc122a1d94" +[[package]] +name = "lru" +version = "0.12.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "234cf4f4a04dc1f57e24b96cc0cd600cf2af460d4161ac5ecdd0af8e1f3b2a38" +dependencies = [ + "hashbrown 0.15.5", +] + [[package]] name = "mach2" version = "0.4.3" @@ -1388,6 +1628,21 @@ version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ca88d725a0a943b096803bd34e73a4437208b6077654cc4ecb2947a5f91618d" +[[package]] +name = "match_cfg" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ffbee8634e0d45d258acb448e7eaab3fce7a0a467395d4d9f228e3c1f01fb2e4" + +[[package]] +name = "matchers" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9" +dependencies = [ + "regex-automata", +] + [[package]] name = "mavlink" version = "0.13.1" @@ -1447,12 +1702,53 @@ version = "2.7.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a282da65faaf38286cf3be983213fcf1d2e2a58700e808f83f4ea9a4804bc0" +[[package]] +name = "memoffset" +version = "0.6.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5aa361d4faea93603064a027415f07bd8e1d5c88c9fbf68bf56a285428fd79ce" +dependencies = [ + "autocfg", +] + [[package]] name = "micromath" version = "2.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c3c8dda44ff03a2f238717214da50f65d5a53b45cd213a7370424ffdb6fae815" +[[package]] +name = "mio" +version = "0.8.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4a650543ca06a924e8b371db273b2756685faae30f8487da1b56505a8f78b0c" +dependencies = [ + "libc", + "log", + "wasi 0.11.1+wasi-snapshot-preview1", + "windows-sys 0.48.0", +] + +[[package]] +name = "mio" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a69bcab0ad47271a0234d9422b131806bf3968021e5dc9328caf2d4cd58557fc" +dependencies = [ + "libc", + "wasi 0.11.1+wasi-snapshot-preview1", + "windows-sys 0.61.2", +] + +[[package]] +name = "multimap" +version = "0.8.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5ce46fe64a9d73be07dcbe690a38ce1b293be448fd8ce1e6c1b8062c9f72c6a" +dependencies = [ + "serde", +] + [[package]] name = "multimap" version = "0.10.1" @@ -1489,6 +1785,19 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8d5439c4ad607c3c23abf66de8c8bf57ba8adcd1f129e699851a6e43935d339d" +[[package]] +name = "nix" +version = "0.23.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f3790c00a0150112de0f4cd161e3d7fc4b2d8a5542ffc35f099a2562aecb35c" +dependencies = [ + "bitflags 1.3.2", + "cc", + "cfg-if", + "libc", + "memoffset", +] + [[package]] name = "nix" version = "0.26.4" @@ -1500,6 +1809,18 @@ dependencies = [ "libc", ] +[[package]] +name = "nix" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ab2156c4fce2f8df6c499cc1c763e4394b7482525bf2a9701c9d79d215f519e4" +dependencies = [ + "bitflags 2.10.0", + "cfg-if", + "cfg_aliases 0.1.1", + "libc", +] + [[package]] name = "nix" version = "0.30.1" @@ -1508,7 +1829,7 @@ checksum = "74523f3a35e05aba87a1d978330aef40f67b0304ac79c1c00b294c9830543db6" dependencies = [ "bitflags 2.10.0", "cfg-if", - "cfg_aliases", + "cfg_aliases 0.2.1", "libc", ] @@ -1523,6 +1844,15 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "nu-ansi-term" +version = "0.50.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "num-complex" version = "0.4.6" @@ -1605,6 +1935,29 @@ dependencies = [ "defmt 0.3.100", ] +[[package]] +name = "parking_lot" +version = "0.12.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93857453250e3077bd71ff98b6a65ea6621a19bb0f559a85248955ac12c45a1a" +dependencies = [ + "lock_api", + "parking_lot_core", +] + +[[package]] +name = "parking_lot_core" +version = "0.9.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2621685985a2ebf1c516881c026032ac7deafcda1a2c9b7850dc81e3dfcb64c1" +dependencies = [ + "cfg-if", + "libc", + "redox_syscall", + "smallvec", + "windows-link 0.2.1", +] + [[package]] name = "paste" version = "1.0.15" @@ -1742,9 +2095,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac6c3320f9abac597dcbc668774ef006702672474aad53c6d596b62e487b40b1" dependencies = [ "heck", - "itertools", + "itertools 0.14.0", "log", - "multimap", + "multimap 0.10.1", "once_cell", "petgraph", "prettyplease", @@ -1762,7 +2115,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9120690fafc389a67ba3803df527d0ec9cbbc9cc45e4cc20b332996dfb672425" dependencies = [ "anyhow", - "itertools", + "itertools 0.14.0", "proc-macro2", "quote", "syn 2.0.106", @@ -1865,16 +2218,37 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" +[[package]] +name = "rand" +version = "0.8.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22f6172bdec972074665ed81ed53b71da00bfc44b65a753cfde883ec4c702a1a" +dependencies = [ + "libc", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + [[package]] name = "rand" version = "0.9.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6db2770f06117d490610c7488547d543617b21bfa07796d7a12f6f1bd53850d1" dependencies = [ - "rand_chacha", + "rand_chacha 0.9.0", "rand_core 0.9.3", ] +[[package]] +name = "rand_chacha" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" +dependencies = [ + "ppv-lite86", + "rand_core 0.6.4", +] + [[package]] name = "rand_chacha" version = "0.9.0" @@ -1890,6 +2264,9 @@ name = "rand_core" version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" +dependencies = [ + "getrandom 0.2.17", +] [[package]] name = "rand_core" @@ -1897,7 +2274,36 @@ version = "0.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" dependencies = [ - "getrandom", + "getrandom 0.3.3", +] + +[[package]] +name = "ratatui" +version = "0.26.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f44c9e68fd46eda15c646fbb85e1040b657a58cdc8c98db1d97a55930d991eef" +dependencies = [ + "bitflags 2.10.0", + "cassowary", + "compact_str", + "crossterm 0.27.0", + "itertools 0.12.1", + "lru", + "paste", + "stability", + "strum 0.26.3", + "unicode-segmentation", + "unicode-truncate", + "unicode-width", +] + +[[package]] +name = "redox_syscall" +version = "0.5.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" +dependencies = [ + "bitflags 2.10.0", ] [[package]] @@ -1972,6 +2378,12 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "rustversion" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f" + [[package]] name = "ryu" version = "1.0.20" @@ -2091,13 +2503,39 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "serde_json" +version = "1.0.150" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" +dependencies = [ + "itoa", + "memchr", + "serde", + "serde_core", + "zmij", +] + [[package]] name = "sergw" version = "1.0.1" dependencies = [ + "anyhow", + "bytes", "clap", + "crossbeam-channel", + "crossterm 0.26.1", "ctrlc", + "dashmap", + "libmdns", + "nix 0.28.0", + "ratatui", + "serde", + "serde_json", "serialport", + "thiserror 1.0.69", + "tracing", + "tracing-subscriber", ] [[package]] @@ -2162,6 +2600,52 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "sharded-slab" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f40ca3c46823713e0d4209592e8d6e826aa57e928f09752619fc696c499637f6" +dependencies = [ + "lazy_static", +] + +[[package]] +name = "shlex" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8fadd59c855ef2080decdef8ff161eb6661b86933c9d82e5ba29dc602a55aba" + +[[package]] +name = "signal-hook" +version = "0.3.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d881a16cf4426aa584979d30bd82cb33429027e42122b169753d6ef1085ed6e2" +dependencies = [ + "libc", + "signal-hook-registry", +] + +[[package]] +name = "signal-hook-mio" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b75a19a7a740b25bc7944bdee6172368f988763b744e3d4dfe753f6b4ece40cc" +dependencies = [ + "libc", + "mio 0.8.11", + "signal-hook", +] + +[[package]] +name = "signal-hook-registry" +version = "1.4.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c4db69cba1110affc0e9f7bcd48bbf87b3f4fc7c61fc9155afd4c469eb3d6c1b" +dependencies = [ + "errno", + "libc", +] + [[package]] name = "simba" version = "0.9.1" @@ -2174,6 +2658,18 @@ dependencies = [ "paste", ] +[[package]] +name = "slab" +version = "0.4.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" + +[[package]] +name = "smallvec" +version = "1.15.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90" + [[package]] name = "smlang" version = "0.8.0" @@ -2209,6 +2705,26 @@ dependencies = [ "managed", ] +[[package]] +name = "socket2" +version = "0.4.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9f7916fc008ca5542385b89a3d3ce689953c143e9304a9bf8beec1de48994c0d" +dependencies = [ + "libc", + "winapi", +] + +[[package]] +name = "socket2" +version = "0.6.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d1e2c7f27f8d4cb10542a02c49005dbd6e93095799d6f3be745fae9f8fedd4" +dependencies = [ + "libc", + "windows-sys 0.61.2", +] + [[package]] name = "spin" version = "0.9.8" @@ -2228,6 +2744,16 @@ dependencies = [ "serde", ] +[[package]] +name = "stability" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d904e7009df136af5297832a3ace3370cd14ff1546a232f4f185036c2736fcac" +dependencies = [ + "quote", + "syn 2.0.106", +] + [[package]] name = "stable_deref_trait" version = "1.2.0" @@ -2280,13 +2806,35 @@ version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" +[[package]] +name = "strum" +version = "0.26.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8fec0f0aef304996cf250b31b5a10dee7980c85da9d759361292b8bca5a18f06" +dependencies = [ + "strum_macros 0.26.4", +] + [[package]] name = "strum" version = "0.27.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" dependencies = [ - "strum_macros", + "strum_macros 0.27.2", +] + +[[package]] +name = "strum_macros" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4c6bee85a5a24955dc440386795aa378cd9cf82acd5f764469152d2270e581be" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "rustversion", + "syn 2.0.106", ] [[package]] @@ -2343,7 +2891,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "655da9c7eb6305c55742045d5a8d2037996d61d8de95806335c7c86ce0f82e9c" dependencies = [ "fastrand", - "getrandom", + "getrandom 0.3.3", "once_cell", "rustix", "windows-sys 0.61.2", @@ -2407,6 +2955,89 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "thread_local" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ad99c4c6d32803332c548b1af0540b357b3f5fc0be8f6c6bfe8b2e6ae784070" +dependencies = [ + "cfg-if", +] + +[[package]] +name = "tokio" +version = "1.50.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "27ad5e34374e03cfffefc301becb44e9dc3c17584f414349ebe29ed26661822d" +dependencies = [ + "libc", + "mio 1.1.1", + "pin-project-lite", + "socket2 0.6.5", + "windows-sys 0.61.2", +] + +[[package]] +name = "tracing" +version = "0.1.44" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" +dependencies = [ + "pin-project-lite", + "tracing-attributes", + "tracing-core", +] + +[[package]] +name = "tracing-attributes" +version = "0.1.31" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7490cfa5ec963746568740651ac6781f701c9c5ea257c58e057f3ba8cf69e8da" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + +[[package]] +name = "tracing-core" +version = "0.1.36" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "db97caf9d906fbde555dd62fa95ddba9eecfd14cb388e4f491a66d74cd5fb79a" +dependencies = [ + "once_cell", + "valuable", +] + +[[package]] +name = "tracing-log" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee855f1f400bd0e5c02d150ae5de3840039a3f54b025156404e34c23c03f47c3" +dependencies = [ + "log", + "once_cell", + "tracing-core", +] + +[[package]] +name = "tracing-subscriber" +version = "0.3.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb7f578e5945fb242538965c2d0b04418d38ec25c79d160cd279bf0731c8d319" +dependencies = [ + "matchers", + "nu-ansi-term", + "once_cell", + "regex-automata", + "sharded-slab", + "smallvec", + "thread_local", + "tracing", + "tracing-core", + "tracing-log", +] + [[package]] name = "ts-rs" version = "11.1.0" @@ -2503,6 +3134,17 @@ version = "1.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f6ccf251212114b54433ec949fd6a7841275f9ada20dddd2f29e9ceea4501493" +[[package]] +name = "unicode-truncate" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b3644627a5af5fa321c95b9b235a72fd24cd29c648c2c379431e6628655627bf" +dependencies = [ + "itertools 0.13.0", + "unicode-segmentation", + "unicode-width", +] + [[package]] name = "unicode-width" version = "0.1.14" @@ -2547,7 +3189,7 @@ dependencies = [ "mavlink-core 0.16.2", "num-derive 0.4.2", "num-traits", - "rand", + "rand 0.9.2", "serde", "serde_arrays 0.2.0", "ts-rs", @@ -2656,6 +3298,12 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "valuable" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" + [[package]] name = "vcell" version = "0.1.3" @@ -2683,6 +3331,12 @@ dependencies = [ "vcell", ] +[[package]] +name = "wasi" +version = "0.11.1+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" + [[package]] name = "wasi" version = "0.14.2+wasi-0.2.4" @@ -2692,6 +3346,22 @@ dependencies = [ "wit-bindgen-rt", ] +[[package]] +name = "winapi" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c839a674fcd7a98952e593242ea400abe93992746761e38641405d28b00f419" +dependencies = [ + "winapi-i686-pc-windows-gnu", + "winapi-x86_64-pc-windows-gnu", +] + +[[package]] +name = "winapi-i686-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" + [[package]] name = "winapi-util" version = "0.1.11" @@ -2701,6 +3371,12 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "winapi-x86_64-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" + [[package]] name = "windows-link" version = "0.1.3" @@ -2713,6 +3389,15 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" +[[package]] +name = "windows-sys" +version = "0.48.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "677d2418bec65e3338edb076e806bc1ec15693c5d0104683f2efe857f61056a9" +dependencies = [ + "windows-targets 0.48.5", +] + [[package]] name = "windows-sys" version = "0.52.0" @@ -2740,6 +3425,21 @@ dependencies = [ "windows-link 0.2.1", ] +[[package]] +name = "windows-targets" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a2fa6e2155d7247be68c096456083145c183cbbbc2764150dda45a87197940c" +dependencies = [ + "windows_aarch64_gnullvm 0.48.5", + "windows_aarch64_msvc 0.48.5", + "windows_i686_gnu 0.48.5", + "windows_i686_msvc 0.48.5", + "windows_x86_64_gnu 0.48.5", + "windows_x86_64_gnullvm 0.48.5", + "windows_x86_64_msvc 0.48.5", +] + [[package]] name = "windows-targets" version = "0.52.6" @@ -2773,6 +3473,12 @@ dependencies = [ "windows_x86_64_msvc 0.53.0", ] +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b38e32f0abccf9987a4e3079dfb67dcd799fb61361e53e2882c3cbaf0d905d8" + [[package]] name = "windows_aarch64_gnullvm" version = "0.52.6" @@ -2785,6 +3491,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "86b8d5f90ddd19cb4a147a5fa63ca848db3df085e25fee3cc10b39b6eebae764" +[[package]] +name = "windows_aarch64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc" + [[package]] name = "windows_aarch64_msvc" version = "0.52.6" @@ -2797,6 +3509,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c7651a1f62a11b8cbd5e0d42526e55f2c99886c77e007179efff86c2b137e66c" +[[package]] +name = "windows_i686_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e" + [[package]] name = "windows_i686_gnu" version = "0.52.6" @@ -2821,6 +3539,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ce6ccbdedbf6d6354471319e781c0dfef054c81fbc7cf83f338a4296c0cae11" +[[package]] +name = "windows_i686_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406" + [[package]] name = "windows_i686_msvc" version = "0.52.6" @@ -2833,6 +3557,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "581fee95406bb13382d2f65cd4a908ca7b1e4c2f1917f143ba16efe98a589b5d" +[[package]] +name = "windows_x86_64_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53d40abd2583d23e4718fddf1ebec84dbff8381c07cae67ff7768bbf19c6718e" + [[package]] name = "windows_x86_64_gnu" version = "0.52.6" @@ -2845,6 +3575,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2e55b5ac9ea33f2fc1716d1742db15574fd6fc8dadc51caab1c16a3d3b4190ba" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b7b52767868a23d5bab768e390dc5f5c55825b6d30b86c844ff2dc7414044cc" + [[package]] name = "windows_x86_64_gnullvm" version = "0.52.6" @@ -2857,6 +3593,12 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0a6e035dd0599267ce1ee132e51c27dd29437f63325753051e71dd9e42406c57" +[[package]] +name = "windows_x86_64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed94fce61571a4006852b7389a063ab983c02eb1bb37b47f8272ce92d06d9538" + [[package]] name = "windows_x86_64_msvc" version = "0.52.6" @@ -2897,3 +3639,9 @@ dependencies = [ "quote", "syn 2.0.106", ] + +[[package]] +name = "zmij" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29666d0abbfad1e3dc4dcf6144730dd3a3ab225bbbdac83319345b1b44ccfc1b" diff --git a/apps/sergw/Cargo.toml b/apps/sergw/Cargo.toml index 250cf93..7abb1c6 100644 --- a/apps/sergw/Cargo.toml +++ b/apps/sergw/Cargo.toml @@ -7,6 +7,28 @@ publish.workspace = true license.workspace = true [dependencies] -clap = { version = "4.5.16", features = ["derive"] } -ctrlc = "3.4" +anyhow = "1" +bytes = "1" +clap = { version = "4", features = ["derive"] } +crossbeam-channel = "0.5" +ctrlc = "3" serialport = "4" +tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] } +dashmap = "5" +serde = { version = "1", features = ["derive"] } +serde_json = "1" +thiserror = "1" +crossterm = "0.26" +ratatui = "0.26" +libmdns = { version = "0.7", optional = true } + +[target.'cfg(target_os = "linux")'.dev-dependencies] +nix = { version = "0.28", features = ["term"] } + +[target.'cfg(target_os = "linux")'.dependencies] +nix = { version = "0.28", features = ["term"] } + +[features] +default = ["mdns"] +mdns = ["libmdns"] diff --git a/apps/sergw/README.md b/apps/sergw/README.md index b9bf192..6a10236 100644 --- a/apps/sergw/README.md +++ b/apps/sergw/README.md @@ -1,118 +1,101 @@ -# SERGW - Serial Gateway +### sergw -## NAME +Simple serial ↔ TCP gateway with a built‑in TUI and optional zero‑config mDNS advertisement. -sergw - Serial-to-TCP gateway for communication between serial devices and TCP clients +sergw opens a local serial device and serves it over TCP, broadcasting serial output to all connected clients and forwarding client input to the serial device. It emphasizes pragmatic reliability (auto‑reconnect) and visibility (TUI overview and hex/ascii/dec inspector). -## SYNOPSIS +### Highlights -**sergw** [*OPTIONS*] _COMMAND_ +- Serial ↔ TCP bridge over raw TCP (no telnet/RFC2217, no TLS) +- Multi‑client fan‑out: serial output is broadcast to all connected TCP clients +- Backpressure handling: slow or disconnected clients are dropped +- TUI: overview (connections, throughput, events) and inspector (hex/ascii/dec) +- Auto‑reconnect for serial (reader/writer) with buffered retry for writes +- Zero‑config mDNS (feature ‘mdns’, enabled by default): `_sergw._tcp` +- Linux mock tools: PTY‑backed mock serial and TCP chat helper (Linux only) -**sergw ports** +### Install -**sergw listen** [*OPTIONS*] --serial _PORT_ +- From crates.io (default features include mDNS): -## DESCRIPTION - -sergw bridges communication between a serial device and multiple TCP clients. Data from the serial port is broadcast to all connected TCP clients. Data from any TCP client is forwarded to the serial port. - -The gateway automatically reconnects to the serial port if the connection is lost. Graceful shutdown is supported via Ctrl-C. - -## COMMANDS - -### ports - -Lists available USB serial ports. Outputs port names separated by spaces. - -```bash -sergw ports -# Output: /dev/ttyUSB0 /dev/ttyUSB1 /dev/ttyACM0 +``` +cargo install sergw ``` -### listen - -Starts the TCP server and establishes serial communication. - -**Required Options:** - -- `--serial *PORT*` - Serial port device path (e.g., `/dev/ttyUSB0`, `COM3`) - -**Optional Options:** - -- `--baud *RATE*` - Serial baud rate (default: 57600) -- `--host *ADDRESS*` - TCP bind address and port (default: 127.0.0.1:5656) -- `-v, --verbose` - Enable verbose output +- Without mDNS: -## OPTIONS +``` +cargo install sergw --no-default-features +``` -| Option | Description | Default | -| --------------- | ------------------------- | -------------- | -| `--serial` | Serial port device path | Required | -| `--baud` | Serial baud rate | 57600 | -| `--host` | TCP bind address and port | 127.0.0.1:5656 | -| `-v, --verbose` | Enable verbose logging | false | +### Quick start -## EXAMPLES +Bridge a serial device on port 5656 (defaults shown): -```bash -sergw listen --serial /dev/ttyUSB0 -sergw listen --serial /dev/ttyUSB0 --baud 115200 --host 0.0.0.0:8080 -sergw listen --serial /dev/ttyUSB0 --verbose -sergw ports +``` +sergw listen --serial /dev/ttyUSB0 --baud 115200 --host 127.0.0.1:5656 ``` -## ARCHITECTURE - -The gateway uses multiple threads: +Connect a TCP client (e.g. `nc 127.0.0.1 5656`) to interact. -1. **Main Reconnect Loop** - Attempts to connect to the serial port and spawns a reader thread. If the connection is lost, it cleans up and retries after a delay. +### CLI -2. **Serial Reader Thread** - Reads data from the serial port and broadcasts to all TCP clients. +``` +sergw + ports [--all] [--verbose] [--format text|json] + listen [--serial ] [--baud ] [--host ] + [--data-bits five|six|seven|eight] + [--parity none|odd|even] + [--stop-bits one|two] + [--buffer ] + mock serial [--alias ] # Linux only + mock listener [--host ] # Linux only +``` -3. **Serial Writer Thread** - Receives data from TCP clients and writes to the serial port. +- `ports`: list serial ports (USB‑only by default). Use `--all` to include non‑USB. `--format json` for machine output. +- `listen`: start the bridge. If `--serial` is omitted and exactly one USB serial is present, it is auto‑selected; otherwise a helpful error is returned. +- `mock serial` (Linux): create a PTY that behaves like a serial device and open a TUI to interact. +- `mock listener` (Linux): connect to a TCP server with a TUI (handy for testing the bridge from the client side). -4. **TCP Listener Thread** - Accepts new TCP connections and spawns handler threads. +### mDNS / Bonjour (optional) -5. **TCP Client Handler Threads** - One per client, handles bidirectional data forwarding. +When built with the `mdns` feature (default), sergw advertises the service using type `_sergw._tcp`. -Thread coordination uses shared state (`Arc>`), message channels (`std::sync::mpsc`), and an atomic shutdown flag. +- Instance name: `sergw:` (e.g. `sergw:ttyUSB0`) +- TXT records: `provider=sergw` -## REQUIREMENTS +Disable mDNS by building without default features: -- Rust toolchain (install via [rustup](https://rustup.rs/)) -- Serial port access permissions - - Linux: Add user to `dialout` group - - Windows: Administrator privileges may be required +``` +cargo build --no-default-features +``` -## TROUBLESHOOTING +### TUI overview -### Permission Denied on Serial Port +- Tabs: Overview (connections, throughput, events), Inspector (live dump) +- Inspector: formats (hex/ascii/dec), per‑device filter, pause/scroll +- Key hints in footer -**Linux:** +### Reliability & behavior -```bash -sudo usermod -a -G dialout $USER -# Log out and back in -``` +- Serial auto‑reconnect on read/write failures; writer retries buffered write after reconnect +- TCP reader/writer per connection; on backpressure the connection is dropped rather than slowing others +- Raw byte forwarding (no framing, no higher protocols) -**Windows:** +### Exit codes -- Run as Administrator +- 2: no serial ports found for auto‑selection +- 3: multiple serial ports detected, explicit `--serial` required +- 4: bind‑like networking error (e.g. address in use) +- 5: serial open/error +- 1: other errors -### Serial Port Not Found +### Development -1. Verify device is connected and powered -2. Check if port is in use: - ```bash - lsof /dev/ttyUSB0 # Linux - ``` -3. Try different baud rates +- Build: `cargo build --all-features` +- Lint: `cargo clippy --all-targets --all-features -- -D warnings` +- Tests: unit tests + Linux PTY integration test -### TCP Connection Issues +### License -1. **Port in use:** - ```bash - netstat -tulpn | grep :5656 # Linux - ``` -2. **Firewall blocking:** Allow incoming connections on specified port -3. **Network access:** Use `0.0.0.0` instead of `127.0.0.1` +GPL‑3.0‑or‑later diff --git a/apps/sergw/src/app/listen.rs b/apps/sergw/src/app/listen.rs new file mode 100644 index 0000000..f283eaa --- /dev/null +++ b/apps/sergw/src/app/listen.rs @@ -0,0 +1 @@ +pub use crate::net::server::run_listen; diff --git a/apps/sergw/src/app/listener.rs b/apps/sergw/src/app/listener.rs new file mode 100644 index 0000000..11fbe3d --- /dev/null +++ b/apps/sergw/src/app/listener.rs @@ -0,0 +1 @@ +pub use crate::net::listener::run_chat; diff --git a/apps/sergw/src/app/mock/mod.rs b/apps/sergw/src/app/mock/mod.rs new file mode 100644 index 0000000..4978a4f --- /dev/null +++ b/apps/sergw/src/app/mock/mod.rs @@ -0,0 +1,4 @@ +pub mod pty; +pub mod serial; +pub mod ui; // orchestrator +pub use serial::run_mock_serial; diff --git a/apps/sergw/src/app/mock/pty.rs b/apps/sergw/src/app/mock/pty.rs new file mode 100644 index 0000000..df0750f --- /dev/null +++ b/apps/sergw/src/app/mock/pty.rs @@ -0,0 +1,18 @@ +#[cfg(target_os = "linux")] +use std::os::fd::AsRawFd; +#[cfg(target_os = "linux")] +use std::os::unix::io::OwnedFd; + +#[cfg(target_os = "linux")] +use anyhow::Result; +#[cfg(target_os = "linux")] +use nix::pty::{openpty, OpenptyResult}; + +#[cfg(target_os = "linux")] +pub fn create_pty_pair() -> Result<(OwnedFd, OwnedFd, String)> { + let OpenptyResult { master, slave, .. } = openpty(None, None)?; + // Resolve stable path to the slave PTY for use as a serial device path + let slave_symlink = format!("/proc/self/fd/{}", slave.as_raw_fd()); + let slave_path = std::fs::read_link(&slave_symlink)?; + Ok((master, slave, slave_path.to_string_lossy().into_owned())) +} diff --git a/apps/sergw/src/app/mock/serial.rs b/apps/sergw/src/app/mock/serial.rs new file mode 100644 index 0000000..e79d955 --- /dev/null +++ b/apps/sergw/src/app/mock/serial.rs @@ -0,0 +1,29 @@ +// Orchestrates PTY creation and UI + +#[cfg(target_os = "linux")] +use anyhow::Result; + +#[cfg(target_os = "linux")] +pub fn run_mock_serial() -> Result<()> { + use super::pty::create_pty_pair; + use super::ui::run_mock_chat_with_title; + + let (master, _slave_fd, slave_path) = create_pty_pair()?; + + // Create a default temporary alias symlink for the slave path for the program duration + let alias_path = "/tmp/sergw-serial"; + // ensure old alias is removed, then create new symlink; cleaned up by guard on exit + let _ = std::fs::remove_file(alias_path); + let _ = std::os::unix::fs::symlink(&slave_path, alias_path); + + struct SymlinkGuard(&'static str); + impl Drop for SymlinkGuard { + fn drop(&mut self) { + let _ = std::fs::remove_file(self.0); + } + } + let _guard = SymlinkGuard(alias_path); + + run_mock_chat_with_title(master, format!("mock serial | {alias_path}"))?; + Ok(()) +} diff --git a/apps/sergw/src/app/mock/ui.rs b/apps/sergw/src/app/mock/ui.rs new file mode 100644 index 0000000..64c9b16 --- /dev/null +++ b/apps/sergw/src/app/mock/ui.rs @@ -0,0 +1,160 @@ +#[cfg(target_os = "linux")] +use std::fs::File; +#[cfg(target_os = "linux")] +use std::io::{Read, Write}; +#[cfg(target_os = "linux")] +use std::os::unix::io::OwnedFd; +#[cfg(target_os = "linux")] +use std::sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + Arc, +}; +#[cfg(target_os = "linux")] +use std::time::{Duration, Instant}; + +#[cfg(target_os = "linux")] +use anyhow::Result; +#[cfg(target_os = "linux")] +use crossbeam_channel as channel; +#[cfg(target_os = "linux")] +use crossterm::{ + event::{self, Event, KeyCode}, + execute, + terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen}, +}; +#[cfg(target_os = "linux")] +use ratatui::{ + backend::CrosstermBackend, + layout::{Constraint, Direction, Layout}, + text::{Line, Span}, + widgets::{Block, Borders, Paragraph, Wrap}, + Terminal, +}; + +#[cfg(target_os = "linux")] +use crate::metrics::ThroughputAverager; + +#[cfg(target_os = "linux")] +pub fn run_mock_chat_with_title( + master: OwnedFd, + title: String, +) -> Result<()> { + let mut master_file: File = master.into(); + + enable_raw_mode()?; + let mut stdout = std::io::stdout(); + execute!(stdout, EnterAlternateScreen)?; + let backend = CrosstermBackend::new(stdout); + let mut terminal = Terminal::new(backend)?; + + let stop = Arc::new(AtomicBool::new(false)); + let rx_bytes = Arc::new(AtomicU64::new(0)); + let tx_bytes = Arc::new(AtomicU64::new(0)); + let (log_tx, log_rx) = channel::unbounded::(); + + // Reader thread from PTY master + let stop_r = stop.clone(); + let rx_b = rx_bytes.clone(); + let mut reader = master_file.try_clone()?; + let log_tx_reader = log_tx.clone(); + std::thread::spawn(move || { + let mut buf = [0u8; 4096]; + loop { + if stop_r.load(Ordering::Relaxed) { + break; + } + match reader.read(&mut buf) { + Ok(0) => std::thread::sleep(Duration::from_millis(20)), + Ok(n) => { + rx_b.fetch_add(n as u64, Ordering::Relaxed); + let s = String::from_utf8_lossy(&buf[..n]).to_string(); + let _ = log_tx_reader.send(format!("< {s}")); + } + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { + std::thread::sleep(Duration::from_millis(20)); + } + Err(_) => break, + } + } + }); + + let mut logs: Vec = Vec::new(); + let mut input = String::new(); + let mut last_rx = 0u64; + let mut last_tx = 0u64; + let mut avg_in = ThroughputAverager::new(5.0); + let mut avg_out = ThroughputAverager::new(5.0); + let mut last_time = Instant::now(); + + loop { + while let Ok(line) = log_rx.try_recv() { + logs.push(line); + if logs.len() > 200 { + logs.remove(0); + } + } + + // Throughput + let now = Instant::now(); + let dt = now.duration_since(last_time).as_secs_f64().max(0.001); + let rx = rx_bytes.load(Ordering::Relaxed); + let tx = tx_bytes.load(Ordering::Relaxed); + let inbound = avg_in.update(rx - last_rx, dt) as u64; // from serial (smoothed) + let outbound = avg_out.update(tx - last_tx, dt) as u64; // to serial (smoothed) + last_rx = rx; + last_tx = tx; + last_time = now; + + terminal.draw(|f| { + let chunks = Layout::default() + .direction(Direction::Vertical) + .constraints([Constraint::Length(1), Constraint::Min(3), Constraint::Length(3)]) + .split(f.size()); + + let header = Paragraph::new(format!("{title} | In: {inbound} B/s Out: {outbound} B/s")); + f.render_widget(header, chunks[0]); + + // Auto-scroll: render only the last lines that fit + let viewport = chunks[1].height.saturating_sub(2) as usize; // minus borders + let start = logs.len().saturating_sub(viewport); + let lines: Vec = logs.iter().skip(start).map(|l| Line::from(Span::raw(l.clone()))).collect(); + let para = Paragraph::new(lines) + .wrap(Wrap { trim: false }) + .block(Block::default().title("Messages").borders(Borders::ALL)); + f.render_widget(para, chunks[1]); + + let input_box = + Paragraph::new(input.clone()).block(Block::default().title("Input (Enter to send, Ctrl+C to quit)").borders(Borders::ALL)); + f.render_widget(input_box, chunks[2]); + })?; + + if event::poll(Duration::from_millis(50))? { + if let Event::Key(k) = event::read()? { + match k.code { + KeyCode::Char('c') if k.modifiers.contains(crossterm::event::KeyModifiers::CONTROL) => break, + KeyCode::Char(c) => input.push(c), + KeyCode::Backspace => { + input.pop(); + } + KeyCode::Enter => { + if !input.is_empty() { + let mut to_send = input.clone(); + to_send.push('\n'); + let _ = master_file.write_all(to_send.as_bytes()); + tx_bytes.fetch_add(to_send.len() as u64, Ordering::Relaxed); + let _ = log_tx.send(format!("> {input}")); + input.clear(); + } + } + KeyCode::Esc => input.clear(), + _ => {} + } + } + } + } + + disable_raw_mode()?; + execute!(terminal.backend_mut(), LeaveAlternateScreen)?; + terminal.show_cursor()?; + Ok(()) +} diff --git a/apps/sergw/src/app/mod.rs b/apps/sergw/src/app/mod.rs new file mode 100644 index 0000000..eda8b0d --- /dev/null +++ b/apps/sergw/src/app/mod.rs @@ -0,0 +1,4 @@ +// High-level app modules +pub mod listen; +pub mod listener; +pub mod mock; diff --git a/apps/sergw/src/cli.rs b/apps/sergw/src/cli.rs new file mode 100644 index 0000000..e93831a --- /dev/null +++ b/apps/sergw/src/cli.rs @@ -0,0 +1,222 @@ +use std::net::SocketAddr; + +use clap::{Parser, Subcommand, ValueEnum}; +use serialport::{DataBits, Parity, StopBits}; + +#[derive(Parser)] +#[command(author, version, about, long_about = None)] +#[command(propagate_version = true)] +pub struct Cli { + #[command(subcommand)] + pub command: Option, +} + +#[derive(Subcommand)] +pub enum Commands { + /// List available serial ports + Ports { + /// Include non-USB ports as well + #[arg(long)] + all: bool, + /// Show detailed metadata + #[arg(long)] + verbose: bool, + /// Output format + #[arg(long, value_enum, default_value_t = PortsFormat::Text)] + format: PortsFormat, + }, + /// Bridge a serial port to TCP + Listen(Listen), + + #[cfg(target_os = "linux")] + /// Mock utilities + Mock { + #[command(subcommand)] + cmd: MockCmd, + }, +} + +#[derive(Parser, Clone, Debug)] +pub struct Listen { + /// Serial port to open (auto-select if exactly one is found and this is omitted) + #[arg(long)] + pub serial: Option, + + /// Baud rate + #[arg(long, default_value_t = 115_200)] + pub baud: u32, + + /// TCP listen address + #[arg(long, default_value = "127.0.0.1:5656")] + pub host: SocketAddr, + + /// Data bits + #[arg(long, value_enum, default_value_t = DataBitsOpt::Eight)] + pub data_bits: DataBitsOpt, + + /// Parity + #[arg(long, value_enum, default_value_t = ParityOpt::None)] + pub parity: ParityOpt, + + /// Stop bits + #[arg(long, value_enum, default_value_t = StopBitsOpt::One)] + pub stop_bits: StopBitsOpt, + + /// Buffer capacity (messages) for internal channels + #[arg(long, default_value_t = 4096)] + pub buffer: usize, + + /// Run without the TUI (daemon mode, logs output to stderr) + #[arg(long)] + pub no_tui: bool, +} + +#[cfg(target_os = "linux")] +#[derive(Subcommand, Clone, Debug)] +pub enum MockCmd { + /// Create a PTY-backed serial device and open a chat UI bound to it + Serial { + /// Optionally create a symlink to the slave PTY at this path (cannot force /dev/pts/N) + #[arg(long)] + alias: Option, + }, + /// Open a chat UI connected to a TCP server (replaces `socat - TCP:host:port`) + Listener { + #[command(flatten)] + chat: Chat, + }, +} + +#[derive(Parser, Clone, Debug)] +pub struct Chat { + /// TCP server to connect to (e.g. 127.0.0.1:5656) + #[arg(long, default_value = "127.0.0.1:5656")] + pub host: std::net::SocketAddr, +} + +#[derive(ValueEnum, Clone, Debug)] +pub enum DataBitsOpt { + Five, + Six, + Seven, + Eight, +} + +impl From for DataBits { + fn from(v: DataBitsOpt) -> Self { + match v { + DataBitsOpt::Five => DataBits::Five, + DataBitsOpt::Six => DataBits::Six, + DataBitsOpt::Seven => DataBits::Seven, + DataBitsOpt::Eight => DataBits::Eight, + } + } +} + +#[derive(ValueEnum, Clone, Debug)] +pub enum ParityOpt { + None, + Odd, + Even, +} + +impl From for Parity { + fn from(v: ParityOpt) -> Self { + match v { + ParityOpt::None => Parity::None, + ParityOpt::Odd => Parity::Odd, + ParityOpt::Even => Parity::Even, + } + } +} + +#[derive(ValueEnum, Clone, Debug)] +pub enum StopBitsOpt { + One, + Two, +} + +impl From for StopBits { + fn from(v: StopBitsOpt) -> Self { + match v { + StopBitsOpt::One => StopBits::One, + StopBitsOpt::Two => StopBits::Two, + } + } +} + +#[derive(ValueEnum, Clone, Debug)] +pub enum PortsFormat { + Text, + Json, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_listen_defaults() { + let cli = Cli::parse_from(["sergw", "listen"]); + match cli.command.unwrap() { + Commands::Listen(l) => { + assert_eq!(l.serial, None); + assert_eq!(l.baud, 115_200); + assert_eq!(l.host, "127.0.0.1:5656".parse().unwrap()); + assert!(matches!(l.data_bits, DataBitsOpt::Eight)); + assert!(matches!(l.parity, ParityOpt::None)); + assert!(matches!(l.stop_bits, StopBitsOpt::One)); + assert_eq!(l.buffer, 4096); + assert!(!l.no_tui); + } + _ => panic!("expected listen"), + } + } + + #[test] + fn parse_listen_values() { + let cli = Cli::parse_from([ + "sergw", + "listen", + "--serial", + "/dev/ttyUSB9", + "--baud", + "57600", + "--host", + "0.0.0.0:9000", + "--data-bits", + "seven", + "--parity", + "even", + "--stop-bits", + "two", + "--buffer", + "123", + "--no-tui", + ]); + match cli.command.unwrap() { + Commands::Listen(l) => { + assert_eq!(l.serial.as_deref(), Some("/dev/ttyUSB9")); + assert_eq!(l.baud, 57_600); + assert_eq!(l.host, "0.0.0.0:9000".parse().unwrap()); + assert!(matches!(l.data_bits, DataBitsOpt::Seven)); + assert!(matches!(l.parity, ParityOpt::Even)); + assert!(matches!(l.stop_bits, StopBitsOpt::Two)); + assert_eq!(l.buffer, 123); + assert!(l.no_tui); + } + _ => panic!("expected listen"), + } + } + + #[test] + fn parse_ports_json() { + let cli = Cli::parse_from(["sergw", "ports", "--format", "json"]); + match cli.command.unwrap() { + Commands::Ports { format, .. } => { + assert!(matches!(format, PortsFormat::Json)); + } + _ => panic!("expected ports"), + } + } +} diff --git a/apps/sergw/src/main.rs b/apps/sergw/src/main.rs index 745a13c..a432e39 100644 --- a/apps/sergw/src/main.rs +++ b/apps/sergw/src/main.rs @@ -1,134 +1,198 @@ -use std::sync::atomic::AtomicBool; -use std::sync::mpsc::channel; -use std::sync::Arc; - -use clap::{Parser, Subcommand}; - +mod app; +mod cli; +mod metrics; +mod net; mod serial; mod state; -mod tcp_server; - -use serial::{list_available_ports, setup_serial_port, spawn_serial_reader, spawn_serial_writer, RECONNECT_DELAY_DURATION}; -use state::SharedState; -use tcp_server::spawn_tcp_listener; - -#[derive(Parser)] -#[command(author, version, about, long_about = None)] -#[command(propagate_version = true, arg_required_else_help = true)] -struct Cli { - #[command(subcommand)] - command: Option, -} +mod ui; -#[derive(Subcommand)] -enum Commands { - Ports, - Listen(Listen), -} +use anyhow::Result; +use clap::{CommandFactory, Parser}; +use serialport::SerialPortType; +use tracing_subscriber::EnvFilter; + +use crate::app::listen::run_listen; +use crate::cli::{Cli, Commands, PortsFormat}; +use crate::serial::list_available_ports; -#[derive(Parser)] -struct Listen { - #[arg(long)] - serial: String, - #[arg(long, default_value_t = 57600)] - baud: u32, - #[arg(long, default_value = "127.0.0.1:5656")] - host: String, - #[arg(short, long, action)] +fn print_ports( + all: bool, verbose: bool, + format: PortsFormat, +) { + let ports = list_available_ports(all); + match format { + PortsFormat::Text => { + if ports.is_empty() { + eprintln!(""); + std::process::exit(2); + } + for p in ports { + if verbose { + match p.port_type { + SerialPortType::UsbPort(info) => { + println!( + "{}\tUSB vid:pid {:04x}:{:04x}\t{:?}\t{:?}", + p.port_name, info.vid, info.pid, info.product, info.manufacturer, + ); + } + other => { + println!("{}\t{:?}", p.port_name, other); + } + } + } else { + println!("{}", p.port_name); + } + } + } + PortsFormat::Json => { + #[derive(serde::Serialize)] + struct PortOut { + name: String, + kind: String, + vid: Option, + pid: Option, + product: Option, + manufacturer: Option, + } + + let out: Vec = ports + .into_iter() + .map(|p| match p.port_type { + SerialPortType::UsbPort(info) => PortOut { + name: p.port_name, + kind: "usb".into(), + vid: Some(info.vid), + pid: Some(info.pid), + product: info.product, + manufacturer: info.manufacturer, + }, + other => PortOut { + name: p.port_name, + kind: format!("{other:?}"), + vid: None, + pid: None, + product: None, + manufacturer: None, + }, + }) + .collect(); + println!("{}", serde_json::to_string_pretty(&out).unwrap()); + } + } } fn main() { let cli = Cli::parse(); - match &cli.command { - Some(Commands::Ports) => match list_available_ports() { - Ok(ports) => println!("{}", ports.join(" ")), - Err(e) => eprintln!("Error listing serial ports: {}", e), - }, - Some(Commands::Listen(listen_opts)) => { - if let Err(e) = run_gateway(listen_opts) { - eprintln!("Application error: {}", e); + let no_tui = match &cli.command { + Some(Commands::Listen(listen)) => listen.no_tui, + _ => false, + }; + + if no_tui { + // Enable normal external logging for background daemon + tracing_subscriber::fmt() + .with_env_filter(EnvFilter::from_default_env()) + .with_target(false) + .try_init() + .ok(); + } else { + // Silence external logging to keep TUI clean; route important status via the UI event log. + tracing_subscriber::fmt() + .with_env_filter(EnvFilter::from_default_env()) + .with_target(false) + .with_writer(std::io::sink) + .try_init() + .ok(); + } + + let result: Result<()> = match cli.command { + Some(Commands::Ports { all, verbose, format }) => { + print_ports(all, verbose, format); + Ok(()) + } + Some(Commands::Listen(listen)) => run_listen(listen), + #[cfg(target_os = "linux")] + Some(Commands::Mock { cmd: sub }) => match sub { + crate::cli::MockCmd::Serial { alias } => { + let _ = alias; + crate::app::mock::run_mock_serial() } + crate::cli::MockCmd::Listener { chat } => crate::app::listener::run_chat(chat), + }, + None => { + Cli::command().print_help().ok(); + println!(); + Ok(()) } - None => unreachable!("Should be covered by arg_required_else_help = true"), - } -} + }; -fn setup_ctrl_c_handler(shutdown_flag: &Arc) -> Result<(), Box> { - let shutdown_flag_clone = Arc::clone(shutdown_flag); - ctrlc::set_handler(move || { - println!("\nReceived Ctrl-C, initiating graceful shutdown..."); - shutdown_flag_clone.store(true, std::sync::atomic::Ordering::Release); - })?; - Ok(()) + if let Err(err) = result { + // Map to stable exit codes + let code = exit_code_for_error(&err); + eprintln!("error: {err:?}"); + std::process::exit(code); + } } -fn run_gateway(opts: &Listen) -> Result<(), Box> { - let shared_state = Arc::new(std::sync::Mutex::new(SharedState::new(opts.verbose))); - let shutdown_flag = Arc::new(AtomicBool::new(false)); - let shared_port_writer = Arc::new(std::sync::Mutex::new(None::>)); - let (serial_writer_tx, serial_writer_rx) = channel::>(); - - setup_ctrl_c_handler(&shutdown_flag)?; - - println!("Starting TCP listener on: {}", opts.host); - let tcp_handle_task = spawn_tcp_listener(&opts.host, &shared_state, &shutdown_flag, serial_writer_tx)?; - - let serial_writer_task = spawn_serial_writer(serial_writer_rx, &shared_port_writer, &shutdown_flag); - - 'reconnect_loop: loop { - if shutdown_flag.load(std::sync::atomic::Ordering::Acquire) { - println!("Shutdown signal received, exiting main loop."); - break 'reconnect_loop; +pub(crate) fn exit_code_for_error(err: &anyhow::Error) -> i32 { + // 2: no ports, 3: multiple ports, 4: bind failure, 5: serial open failure, 1: other + for cause in err.chain() { + if let Some(sel) = cause.downcast_ref::() { + return match sel { + crate::serial::SerialSelectError::NoPorts => 2, + crate::serial::SerialSelectError::MultiplePorts { .. } => 3, + }; } + if let Some(ioe) = cause.downcast_ref::() { + use std::io::ErrorKind::*; + return match ioe.kind() { + AddrInUse | AddrNotAvailable | PermissionDenied | ConnectionAborted | ConnectionReset => 4, + _ => 1, + }; + } + if cause.is::() { + return 5; + } + } + 1 +} - println!("Attempting to connect to serial port {} at {} baud...", opts.serial, opts.baud); - let port = match setup_serial_port(&opts.serial, opts.baud) { - Ok(p) => { - println!("Successfully connected to serial port."); - p - } - Err(e) => { - eprintln!( - "Failed to open serial port: {}. Retrying in {} seconds...", - e, - RECONNECT_DELAY_DURATION.as_secs() - ); - std::thread::sleep(RECONNECT_DELAY_DURATION); - continue 'reconnect_loop; - } - }; - - // The port was opened successfully - let port_writer_clone = port.try_clone()?; - *shared_port_writer.lock().expect("Mutex poisoned") = Some(port_writer_clone); +#[cfg(test)] +mod tests { + use super::*; - let serial_reader_task = spawn_serial_reader(port, &shared_state, &shutdown_flag); + #[test] + fn exit_code_no_ports() { + let err = anyhow::Error::from(crate::serial::SerialSelectError::NoPorts); + assert_eq!(exit_code_for_error(&err), 2); + } - // Block until the reader thread exits - serial_reader_task.join().unwrap_or_else(|e| { - eprintln!("Serial reader thread panicked: {:?}", e); + #[test] + fn exit_code_multiple_ports() { + let err = anyhow::Error::from(crate::serial::SerialSelectError::MultiplePorts { + list: vec!["a".into(), "b".into()], }); - - // Connection is dead. - *shared_port_writer.lock().expect("Mutex poisoned") = None; - - // Retry or exit - if !shutdown_flag.load(std::sync::atomic::Ordering::Acquire) { - eprintln!("Serial connection lost. Attempting to reconnect..."); - std::thread::sleep(RECONNECT_DELAY_DURATION); - } + assert_eq!(exit_code_for_error(&err), 3); } - println!("Shutting down long-running tasks..."); - if let Err(e) = tcp_handle_task.join() { - eprintln!("TCP listener thread panicked: {:?}", e); + #[test] + fn exit_code_bind_like_io_error() { + let err = anyhow::Error::from(std::io::Error::from(std::io::ErrorKind::AddrInUse)); + assert_eq!(exit_code_for_error(&err), 4); } - if let Err(e) = serial_writer_task.join() { - eprintln!("Serial writer thread panicked: {:?}", e); + + #[test] + fn exit_code_serial_error() { + let serr = serialport::Error::new(serialport::ErrorKind::NoDevice, "no device"); + let err = anyhow::Error::from(serr); + assert_eq!(exit_code_for_error(&err), 5); } - println!("Application has shut down."); - Ok(()) + #[test] + fn exit_code_other() { + let err = anyhow::anyhow!("other"); + assert_eq!(exit_code_for_error(&err), 1); + } } diff --git a/apps/sergw/src/metrics.rs b/apps/sergw/src/metrics.rs new file mode 100644 index 0000000..50af718 --- /dev/null +++ b/apps/sergw/src/metrics.rs @@ -0,0 +1,38 @@ +pub struct ThroughputAverager { + tau_secs: f64, + smoothed_bps: f64, +} + +impl ThroughputAverager { + pub fn new(tau_secs: f64) -> Self { + Self { tau_secs, smoothed_bps: 0.0 } + } + + pub fn update( + &mut self, + bytes_delta: u64, + dt_secs: f64, + ) -> f64 { + let dt = dt_secs.max(1e-3); + let alpha = 1.0 - (-dt / self.tau_secs).exp(); + let inst = (bytes_delta as f64) / dt; + self.smoothed_bps = self.smoothed_bps * (1.0 - alpha) + inst * alpha; + self.smoothed_bps + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn ewma_smooths_rate() { + let mut avg = ThroughputAverager::new(5.0); + // 1000 bytes per second over 1 second + let r1 = avg.update(1000, 1.0); + // next second 0 bytes; smoothed should not drop to zero instantly + let r2 = avg.update(0, 1.0); + assert!(r1 > r2); + assert!(r2 > 0.0); + } +} diff --git a/apps/sergw/src/net/listener.rs b/apps/sergw/src/net/listener.rs new file mode 100644 index 0000000..640a645 --- /dev/null +++ b/apps/sergw/src/net/listener.rs @@ -0,0 +1,233 @@ +use std::io::{Read, Write}; +use std::net::TcpStream; +use std::sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + Arc, Mutex, +}; +use std::time::{Duration, Instant}; + +use anyhow::Result; +use crossbeam_channel as channel; +use crossterm::{ + event::{self, Event, KeyCode}, + execute, + terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen}, +}; +use ratatui::{ + backend::CrosstermBackend, + layout::{Constraint, Direction, Layout}, + text::{Line, Span}, + widgets::{Block, Borders, Paragraph, Wrap}, + Terminal, +}; + +use crate::cli::Chat; +use crate::metrics::ThroughputAverager; + +pub fn run_chat(chat: Chat) -> Result<()> { + // Connect TCP (retry until available) + let connect = |host: std::net::SocketAddr| -> TcpStream { + loop { + match TcpStream::connect(host) { + Ok(s) => { + let _ = s.set_nodelay(true); + let _ = s.set_nonblocking(true); + break s; + } + Err(_) => { + std::thread::sleep(Duration::from_millis(800)); + } + } + } + }; + let stream = connect(chat.host); + let stream = Arc::new(Mutex::new(stream)); + + // helper to write with one retry on WouldBlock + let try_send = |s: &mut TcpStream, data: &[u8]| -> bool { + match s.write_all(data) { + Ok(_) => true, + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { + std::thread::sleep(Duration::from_millis(50)); + s.write_all(data).is_ok() + } + Err(_) => false, + } + }; + + // UI setup + enable_raw_mode()?; + let mut stdout = std::io::stdout(); + execute!(stdout, EnterAlternateScreen)?; + let backend = CrosstermBackend::new(stdout); + let mut terminal = Terminal::new(backend)?; + + let stop = Arc::new(AtomicBool::new(false)); + let rx_bytes = Arc::new(AtomicU64::new(0)); + let tx_bytes = Arc::new(AtomicU64::new(0)); + let (log_tx, log_rx) = channel::unbounded::(); + + // Reader thread + let stop_r = stop.clone(); + let rx_b = rx_bytes.clone(); + let rstream = Arc::clone(&stream); + let log_tx_reader = log_tx.clone(); + std::thread::spawn(move || { + let mut buf = [0u8; 4096]; + while !stop_r.load(Ordering::Relaxed) { + // lock the stream for this read iteration + let mut guard = match rstream.lock() { + Ok(g) => g, + Err(_) => { + std::thread::sleep(Duration::from_millis(50)); + continue; + } + }; + match guard.read(&mut buf) { + Ok(0) => { + // EOF: server closed; reconnect proactively + drop(guard); + let new_s = connect(chat.host); + if let Ok(mut g) = rstream.lock() { + *g = new_s; + } + let _ = log_tx_reader.send("! reconnected".to_string()); + std::thread::sleep(Duration::from_millis(100)); + } + Ok(n) => { + drop(guard); + rx_b.fetch_add(n as u64, Ordering::Relaxed); + let s = String::from_utf8_lossy(&buf[..n]).to_string(); + let _ = log_tx_reader.send(format!("< {s}")); + } + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { + drop(guard); + std::thread::sleep(Duration::from_millis(20)); + } + Err(_) => { + drop(guard); + // attempt immediate reconnect and notify + let new_s = connect(chat.host); + if let Ok(mut g) = rstream.lock() { + *g = new_s; + } + let _ = log_tx_reader.send("! reconnected".to_string()); + std::thread::sleep(Duration::from_millis(100)); + } + } + } + }); + + let mut logs: Vec = Vec::new(); + let mut input = String::new(); + let mut last_sent: Option> = None; + let mut last_rx = 0u64; + let mut last_tx = 0u64; + let mut avg_in = ThroughputAverager::new(5.0); + let mut avg_out = ThroughputAverager::new(5.0); + let mut last_time = Instant::now(); + + loop { + while let Ok(line) = log_rx.try_recv() { + logs.push(line); + if logs.len() > 200 { + logs.remove(0); + } + } + + // Throughput calc + let now = Instant::now(); + let dt = now.duration_since(last_time).as_secs_f64().max(0.001); + let rx = rx_bytes.load(Ordering::Relaxed); + let tx = tx_bytes.load(Ordering::Relaxed); + let inbound = avg_in.update(rx - last_rx, dt) as u64; // from TCP (smoothed) + let outbound = avg_out.update(tx - last_tx, dt) as u64; // to TCP (smoothed) + last_rx = rx; + last_tx = tx; + last_time = now; + + terminal.draw(|f| { + let chunks = Layout::default() + .direction(Direction::Vertical) + .constraints([Constraint::Length(1), Constraint::Min(3), Constraint::Length(3)]) + .split(f.size()); + + let header = Paragraph::new(format!("listener | {} | In: {} B/s Out: {} B/s", chat.host, inbound, outbound)); + f.render_widget(header, chunks[0]); + + // Auto-scroll: render only the last lines that fit + let viewport = chunks[1].height.saturating_sub(2) as usize; // minus borders + let start = logs.len().saturating_sub(viewport); + let lines: Vec = logs.iter().skip(start).map(|l| Line::from(Span::raw(l.clone()))).collect(); + let para = Paragraph::new(lines) + .wrap(Wrap { trim: false }) + .block(Block::default().title("Messages").borders(Borders::ALL)); + f.render_widget(para, chunks[1]); + + let input_box = + Paragraph::new(input.clone()).block(Block::default().title("Input (Enter to send, Ctrl+C to quit)").borders(Borders::ALL)); + f.render_widget(input_box, chunks[2]); + })?; + + if event::poll(Duration::from_millis(50))? { + if let Event::Key(k) = event::read()? { + match k.code { + KeyCode::Char('c') if k.modifiers.contains(crossterm::event::KeyModifiers::CONTROL) => break, + KeyCode::Char(c) => input.push(c), + KeyCode::Backspace => { + input.pop(); + } + KeyCode::Enter => { + if !input.is_empty() { + let mut to_send = input.clone(); + to_send.push('\n'); + let mut wrote = false; + // try write with reconnect on failure + if let Ok(mut g) = stream.lock() { + if let Ok(Some(_)) = g.take_error() { + // immediate reconnect if socket error present + let new_s = connect(chat.host); + if let Ok(mut gg) = stream.lock() { + *gg = new_s; + } + } + wrote = try_send(&mut g, to_send.as_bytes()); + if !wrote { + let _ = log_tx.send("! write error: Broken pipe".to_string()); + } + } + if !wrote { + // reconnect and retry once + let new_s = connect(chat.host); + if let Ok(mut g) = stream.lock() { + *g = new_s; + } + if let Ok(mut g) = stream.lock() { + if let Some(prev) = &last_sent { + let _ = try_send(&mut g, prev.as_slice()); + } + std::thread::sleep(Duration::from_millis(150)); + wrote = try_send(&mut g, to_send.as_bytes()); + } + } + if wrote { + tx_bytes.fetch_add(to_send.len() as u64, Ordering::Relaxed); + let _ = log_tx.send(format!("> {input}")); + last_sent = Some(to_send.as_bytes().to_vec()); + } + input.clear(); + } + } + KeyCode::Esc => input.clear(), + _ => {} + } + } + } + } + + stop.store(true, Ordering::Relaxed); + disable_raw_mode()?; + execute!(terminal.backend_mut(), LeaveAlternateScreen)?; + terminal.show_cursor()?; + Ok(()) +} diff --git a/apps/sergw/src/net/mod.rs b/apps/sergw/src/net/mod.rs new file mode 100644 index 0000000..6a4f32e --- /dev/null +++ b/apps/sergw/src/net/mod.rs @@ -0,0 +1,2 @@ +pub mod listener; +pub mod server; diff --git a/apps/sergw/src/net/server.rs b/apps/sergw/src/net/server.rs new file mode 100644 index 0000000..f3e56ce --- /dev/null +++ b/apps/sergw/src/net/server.rs @@ -0,0 +1,459 @@ +use std::io::{Read, Write}; +use std::net::TcpListener; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::thread; +use std::time::Duration; + +use anyhow::{Context, Result}; +use bytes::Bytes; +use crossbeam_channel as channel; +#[cfg(feature = "mdns")] +use libmdns as _mdns; +use tracing::{info, warn}; + +use crate::cli::Listen; +use crate::serial::{configure_serial, select_serial_port}; +use crate::state::SharedState; +use crate::ui::inspector::{DirectionTag, Sample}; +use crate::ui::overview::{run_tui, Counters}; + +pub fn run_listen(listen: Listen) -> Result<()> { + let stop_flag = Arc::new(AtomicBool::new(false)); + { + let stop = stop_flag.clone(); + let _ = ctrlc::set_handler(move || { + stop.store(true, Ordering::Relaxed); + }); + } + + run_listen_with_shutdown(listen, stop_flag) +} + +pub(crate) fn run_listen_with_shutdown( + listen: Listen, + stop_flag: Arc, +) -> Result<()> { + let serial_path = select_serial_port(&listen.serial)?; + info!(serial = %serial_path, baud = listen.baud, host = %listen.host, "Starting sergw"); + let (status_tx, status_rx) = channel::unbounded::(); + let status_tx_reader = status_tx.clone(); + let status_tx_writer = status_tx.clone(); + + // Open serial with auto-reconnect loop for writer and reader handles + let (mut serial_port, mut serial_writer_port) = open_serial_pair(&serial_path, &listen)?; + + // Channels + // - to_serial_rx: buffers from TCP -> serial writer + let (to_serial_tx, to_serial_rx) = channel::bounded::(listen.buffer); + + // - shared state for broadcasting serial -> TCP + let shared_state = Arc::new(SharedState::new()); + let counters = Arc::new(Counters::default()); + let (event_tx_base, event_rx) = channel::unbounded::(); + let event_tx = Some(event_tx_base); + + // TUI thread(s) + let shared_for_tui = Arc::clone(&shared_state); + let counters_for_tui = Arc::clone(&counters); + let stop_for_tui = stop_flag.clone(); + // Inspector UI: channel + let (insp_tx, insp_rx) = channel::bounded::(1024); + let status_rx_tui = status_rx.clone(); + let tui_handle = if listen.no_tui { + None + } else { + Some(thread::spawn(move || { + // Merge status messages into events + let (tx, merged_rx) = channel::unbounded::(); + let status_rx_tui_clone = status_rx_tui.clone(); + std::thread::spawn(move || loop { + crossbeam_channel::select! { + recv(event_rx) -> msg => if let Ok(m)=msg { let _=tx.send(m); } else { break; }, + recv(status_rx_tui_clone) -> msg => if let Ok(m)=msg { let _=tx.send(m); } else { break; }, + } + }); + let _ = run_tui(shared_for_tui, counters_for_tui, merged_rx, insp_rx, stop_for_tui); + })) + }; + + // Inspector receiver is moved into the TUI above; keep tx for sampling below + + // Metrics reporter (always on; logs to info every 5 seconds) + { + let counters_for_metrics = Arc::clone(&counters); + let stop_for_metrics = stop_flag.clone(); + std::thread::spawn(move || { + let mut last_in: u64 = 0; // bytes_in (TCP -> serial) + let mut last_out: u64 = 0; // bytes_out (serial -> TCP) + let mut last = std::time::Instant::now(); + while !stop_for_metrics.load(std::sync::atomic::Ordering::Relaxed) { + std::thread::sleep(std::time::Duration::from_secs(5)); + let now = std::time::Instant::now(); + let dt = now.duration_since(last).as_secs_f64().max(0.001); + last = now; + let bi = counters_for_metrics.bytes_in.load(std::sync::atomic::Ordering::Relaxed); + let bo = counters_for_metrics.bytes_out.load(std::sync::atomic::Ordering::Relaxed); + let outbound = ((bi - last_in) as f64 / dt) as u64; // to serial + let inbound = ((bo - last_out) as f64 / dt) as u64; // from serial + last_in = bi; + last_out = bo; + info!(inbound_bps = inbound, outbound_bps = outbound, "Throughput"); + } + }); + } + + // Serial reader thread: serial -> broadcast + let shared_state_for_reader = Arc::clone(&shared_state); + let stop_reader = stop_flag.clone(); + let serial_path_for_reader = serial_path.clone(); + let listen_for_reader = listen.clone(); + let counters_reader = Arc::clone(&counters); + let insp_tx_reader = insp_tx.clone(); + let serial_reader = thread::spawn(move || -> Result<()> { + let mut buffer = vec![0u8; 4096]; + loop { + while !stop_reader.load(Ordering::Relaxed) { + match serial_port.read(&mut buffer) { + Ok(n) if n > 0 => { + counters_reader.bytes_out.fetch_add(n as u64, Ordering::Relaxed); + let _ = insp_tx_reader.try_send(Sample { + dir: DirectionTag::Inbound, + data: Bytes::copy_from_slice(&buffer[..n]), + }); + let bytes = Bytes::copy_from_slice(&buffer[..n]); + shared_state_for_reader.broadcast(bytes); + } + Ok(_) => {} + Err(e) if e.kind() == std::io::ErrorKind::TimedOut => {} + Err(e) if e.kind() == std::io::ErrorKind::BrokenPipe => { + warn!("Serial: disconnected (BrokenPipe), attempting reconnect..."); + let _ = status_tx_reader.send("Serial: disconnected, attempting reconnect...".into()); + break; + } + Err(e) => { + warn!(?e, "Error reading from serial"); + break; + } + } + } + if stop_reader.load(Ordering::Relaxed) { + break; + } + // Attempt reconnect every second + match open_serial_pair(&serial_path_for_reader, &listen_for_reader) { + Ok((sp, spw)) => { + serial_port = sp; + // serial writer port is owned by writer thread; we keep only reader here + drop(spw); + // Quiet console; status sent to UI + let _ = status_tx_reader.send("Serial: reconnected (reader)".into()); + } + Err(e) => { + warn!(?e, "Reconnect failed (reader), retrying in 1s"); + std::thread::sleep(Duration::from_secs(1)); + } + } + } + Ok(()) + }); + + // Serial writer thread: TCP -> serial + let stop_writer = stop_flag.clone(); + let serial_path_for_writer = serial_path.clone(); + let listen_for_writer = listen.clone(); + let serial_writer = thread::spawn(move || -> Result<()> { + loop { + if stop_writer.load(Ordering::Relaxed) { + break; + } + match to_serial_rx.recv_timeout(Duration::from_millis(200)) { + Ok(buf) => { + if let Err(e) = serial_writer_port.write_all(&buf) { + warn!(?e, "Serial: write failed, reconnecting writer..."); + let _ = status_tx_writer.send("Serial: write failed, reconnecting writer...".into()); + // try to reconnect serial writer and send a priming + // zero-length write to ensure OS queues are ready + loop { + if stop_writer.load(Ordering::Relaxed) { + return Ok(()); + } + match open_serial_pair(&serial_path_for_writer, &listen_for_writer) { + Ok((sp, spw)) => { + // keep writer + serial_writer_port = spw; + drop(sp); // reader will reconnect separately + // Quiet console; status sent to UI + let _ = status_tx_writer.send("Serial: reconnected (writer)".into()); + // After successful reconnect, retry the buffered write once + let _ = serial_writer_port.write_all(&buf); + let _ = serial_writer_port.flush(); + break; + } + Err(err) => { + warn!(?err, "Reconnect failed (writer), retrying in 1s"); + std::thread::sleep(Duration::from_secs(1)); + } + } + } + } + } + Err(channel::RecvTimeoutError::Timeout) => {} + Err(channel::RecvTimeoutError::Disconnected) => break, + } + } + Ok(()) + }); + + // TCP acceptor + let listener = TcpListener::bind(listen.host).with_context(|| format!("Binding TCP listener at {}", listen.host))?; + listener.set_nonblocking(true).context("Setting TCP listener non-blocking mode")?; + + // mDNS/Bonjour advertisement (zero-config), optional via feature flag + #[cfg(feature = "mdns")] + let _mdns_guard: Option<(_mdns::Responder, _mdns::Service)> = { + // Derive a friendly instance name from the serial device + let instance = serial_path + .rsplit('/') + .next() + .map(|s| format!("sergw:{s}")) + .unwrap_or_else(|| "sergw".to_string()); + match _mdns::Responder::new() { + Ok(responder) => { + let port = listen.host.port(); + let txt: [&str; 1] = ["provider=sergw"]; + let service = responder.register("_sergw._tcp".to_string(), instance, port, &txt); + Some((responder, service)) + } + Err(e) => { + warn!(error = ?e, "mDNS responder init failed; continuing without mDNS"); + None + } + } + }; + + loop { + if stop_flag.load(Ordering::Relaxed) { + break; + } + let (stream, addr) = match listener.accept() { + Ok(conn) => conn, + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { + // avoid busy loop + std::thread::sleep(Duration::from_millis(50)); + continue; + } + Err(e) => { + warn!(?e, "Accept failed"); + continue; + } + }; + let mut stream_reader = stream.try_clone().context("Cloning TCP stream (reader)")?; + let mut stream_writer = stream; + if let Err(e) = stream_reader.set_nodelay(true) { + warn!(?e, %addr, "Failed to set TCP_NODELAY on reader"); + } + if let Err(e) = stream_writer.set_nodelay(true) { + warn!(?e, %addr, "Failed to set TCP_NODELAY on writer"); + } + info!(%addr, "Accepted connection"); + if let Some(tx) = &event_tx { + let _ = tx.send(format!("Connected: {addr}")); + } + + let to_serial_tx_conn = to_serial_tx.clone(); + let (to_tcp_tx, to_tcp_rx) = channel::bounded::(listen.buffer); + + // Register connection for broadcasts + shared_state.insert(addr, to_tcp_tx); + + // Introduce connection-specific stop flag to cleanly tear down both threads if either fails + let conn_stop = Arc::new(AtomicBool::new(false)); + + // TCP reader: TCP -> to_serial + let stop_conn = stop_flag.clone(); + let conn_stop_reader = conn_stop.clone(); + let reader_addr = addr; + let counters_in = Arc::clone(&counters); + let insp_tx_reader = insp_tx.clone(); + let tcp_reader = thread::spawn(move || -> Result<()> { + let mut buffer = [0u8; 4096]; + while !stop_conn.load(Ordering::Relaxed) && !conn_stop_reader.load(Ordering::Relaxed) { + match stream_reader.read(&mut buffer) { + Ok(0) => break, + Ok(n) => { + counters_in.bytes_in.fetch_add(n as u64, Ordering::Relaxed); + let buf = Bytes::copy_from_slice(&buffer[..n]); + let _ = insp_tx_reader.try_send(Sample { + dir: DirectionTag::Outbound(reader_addr), + data: buf.clone(), + }); + if let Err(e) = to_serial_tx_conn.send(buf) { + warn!(?e, "Dropping data to serial, backpressure or shutdown"); + break; + } + } + Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { + // Prevent 100% CPU busy loop + std::thread::sleep(Duration::from_millis(20)); + } + Err(e) => { + warn!(?e, addr = %reader_addr, "TCP read error"); + break; + } + } + } + conn_stop_reader.store(true, Ordering::Relaxed); + Ok(()) + }); + + // TCP writer: from broadcast -> TCP + let stop_conn = stop_flag.clone(); + let conn_stop_writer = conn_stop.clone(); + let writer_addr = addr; + let tcp_writer = thread::spawn(move || -> Result<()> { + while !stop_conn.load(Ordering::Relaxed) && !conn_stop_writer.load(Ordering::Relaxed) { + match to_tcp_rx.recv_timeout(Duration::from_millis(200)) { + Ok(buf) => { + if let Err(e) = stream_writer.write_all(&buf) { + warn!(?e, addr = %writer_addr, "TCP write error"); + break; + } + } + Err(channel::RecvTimeoutError::Timeout) => {} + Err(_e) => break, + } + } + conn_stop_writer.store(true, Ordering::Relaxed); + Ok(()) + }); + + // Detach a supervisor for the connection + let shared_state_remove = Arc::clone(&shared_state); + let event_tx_conn = event_tx.clone(); + let conn_stop_supervisor = conn_stop.clone(); + thread::spawn(move || { + // Wait for reader to complete (client closed or error) + let _ = tcp_reader.join(); + // Stop writer thread immediately + conn_stop_supervisor.store(true, Ordering::Relaxed); + // Remove connection immediately so writers drop their sender and exit + shared_state_remove.remove(&addr); + if let Some(tx) = &event_tx_conn { + let _ = tx.send(format!("Disconnected: {addr}")); + } + // Now wait for writer to finish draining/exit + let _ = tcp_writer.join(); + info!(%addr, "Closed connection"); + }); + } + + // Shutdown + info!("Shutting down"); + if let Err(e) = serial_reader.join().unwrap_or(Ok(())) { + warn!(?e, "Serial reader error on shutdown"); + } + if let Err(e) = serial_writer.join().unwrap_or(Ok(())) { + warn!(?e, "Serial writer error on shutdown"); + } + shared_state.dispose(); + + if let Some(handle) = tui_handle { + let _ = handle.join(); + } + + Ok(()) +} + +fn open_serial_pair( + serial_path: &str, + listen: &Listen, +) -> Result<(Box, Box)> { + let builder = serialport::new(serial_path, listen.baud); + let port = configure_serial(builder, listen).with_context(|| format!("Opening serial port {serial_path}"))?; + let writer = port + .try_clone() + .with_context(|| format!("Cloning serial port {serial_path} for writer"))?; + Ok((port, writer)) +} + +#[cfg(all(test, target_os = "linux"))] +mod itests { + use std::fs::File; + use std::io::{Read, Write}; + use std::net::TcpStream; + use std::os::unix::io::AsRawFd; + use std::os::unix::io::OwnedFd; + use std::sync::atomic::AtomicBool; + use std::thread::JoinHandle; + + use super::*; + + // Use a PTY pair to simulate a serial device. The slave path behaves like a tty device. + fn create_pty() -> anyhow::Result<(OwnedFd, String)> { + use nix::pty::{openpty, OpenptyResult, Winsize}; + let OpenptyResult { master, slave, .. } = openpty(None::<&Winsize>, None)?; + // Resolve the symlink to an actual tty path + let slave_path = format!("/proc/self/fd/{}", slave.as_raw_fd()); + let path = std::fs::read_link(&slave_path)?; + // Drop slave (closes fd). Keep master for test I/O. + drop(slave); + Ok((master, path.to_string_lossy().into_owned())) + } + + fn spawn_server( + serial_path: String, + host: &str, + buffer: usize, + ) -> (JoinHandle>, Arc) { + let listen = Listen { + serial: Some(serial_path), + baud: 115_200, + host: host.parse().unwrap(), + data_bits: crate::cli::DataBitsOpt::Eight, + parity: crate::cli::ParityOpt::None, + stop_bits: crate::cli::StopBitsOpt::One, + buffer, + no_tui: true, + }; + let stop = Arc::new(AtomicBool::new(false)); + let stop_clone = stop.clone(); + let handle = std::thread::spawn(move || run_listen_with_shutdown(listen, stop_clone)); + (handle, stop) + } + + #[test] + fn tcp_to_serial_and_back() { + // Arrange: create PTY and start server bound to localhost ephemeral port + let (master_fd, slave_path) = create_pty().expect("pty"); + let mut master: File = master_fd.into(); + let host = "127.0.0.1:6767"; // fixed test port + let (handle, stop) = spawn_server(slave_path, host, 64); + + // connect TCP client + std::thread::sleep(Duration::from_millis(100)); + let mut tcp = loop { + match TcpStream::connect(host) { + Ok(s) => break s, + Err(_) => std::thread::sleep(Duration::from_millis(50)), + } + }; + tcp.set_nodelay(true).ok(); + + // TCP -> serial: write to TCP, read from PTY master + tcp.write_all(b"hello").unwrap(); + let mut serial_buf = [0u8; 5]; + master.read_exact(&mut serial_buf).unwrap(); + assert_eq!(&serial_buf, b"hello"); + + // Serial -> TCP: write to PTY master, read from TCP + master.write_all(b"world").unwrap(); + let mut tcp_buf = [0u8; 5]; + tcp.read_exact(&mut tcp_buf).unwrap(); + assert_eq!(&tcp_buf, b"world"); + + // Shutdown + stop.store(true, Ordering::Relaxed); + let _ = handle.join().unwrap(); + } +} diff --git a/apps/sergw/src/serial.rs b/apps/sergw/src/serial.rs deleted file mode 100644 index 92d2652..0000000 --- a/apps/sergw/src/serial.rs +++ /dev/null @@ -1,99 +0,0 @@ -use std::io::Write; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::Receiver; -use std::sync::{Arc, Mutex}; -use std::time::Duration; - -use serialport::{available_ports, SerialPortType}; - -const BUFFER_SIZE: usize = 1024; -const RECONNECT_DELAY: Duration = Duration::from_secs(3); - -pub fn list_available_ports() -> Result, Box> { - let ports = available_ports()? - .iter() - .filter(|port| matches!(port.port_type, SerialPortType::UsbPort(_))) - .map(|port| port.port_name.clone()) - .collect(); - Ok(ports) -} - -pub fn setup_serial_port( - serial_path: &str, - baud_rate: u32, -) -> Result, Box> { - let port = serialport::new(serial_path, baud_rate) - .timeout(std::time::Duration::from_millis(100)) - .open()?; - Ok(port) -} - -pub fn spawn_serial_reader( - mut port: Box, - shared_state: &Arc>, - shutdown_flag: &Arc, -) -> std::thread::JoinHandle<()> { - let shared_state_clone = Arc::clone(shared_state); - let shutdown_flag_clone = Arc::clone(shutdown_flag); - - std::thread::spawn(move || { - let mut buffer = vec![0; BUFFER_SIZE]; - loop { - if shutdown_flag_clone.load(Ordering::Acquire) { - println!("Serial reader shutting down..."); - return; - } - match port.read(&mut buffer) { - Ok(bytes_read) => { - if bytes_read > 0 { - let data = Arc::from(&buffer[..bytes_read]); - crate::tcp_server::broadcast_data(&shared_state_clone, data); - } - } - Err(e) if e.kind() == std::io::ErrorKind::TimedOut => (), - Err(e) => { - eprintln!("Error reading from serial port: {:?}. Closing reader thread.", e); - return; // Exit on any other error to trigger reconnect. - } - } - } - }) -} - -pub fn spawn_serial_writer( - rx: Receiver>, - port_writer_handle: &Arc>>>, - shutdown_flag: &Arc, -) -> std::thread::JoinHandle<()> { - let shutdown_flag_clone = Arc::clone(shutdown_flag); - let port_writer_handle_clone = Arc::clone(port_writer_handle); - - std::thread::spawn(move || { - println!("Serial writer started."); - loop { - if shutdown_flag_clone.load(Ordering::Acquire) { - println!("Serial writer shutting down..."); - break; - } - - match rx.recv_timeout(Duration::from_millis(100)) { - Ok(data) => { - let mut port_guard = port_writer_handle_clone.lock().expect("Mutex poisoned"); - if let Some(port) = port_guard.as_mut() { - if let Err(e) = port.write_all(&data) { - eprintln!("Serial write failed: {:?}. Data dropped.", e); - } - } - } - Err(std::sync::mpsc::RecvTimeoutError::Timeout) => continue, - Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => { - println!("Serial writer channel disconnected. This should not happen in normal shutdown."); - break; - } - } - } - println!("Serial writer finished."); - }) -} - -pub const RECONNECT_DELAY_DURATION: Duration = RECONNECT_DELAY; diff --git a/apps/sergw/src/serial/io.rs b/apps/sergw/src/serial/io.rs new file mode 100644 index 0000000..2634524 --- /dev/null +++ b/apps/sergw/src/serial/io.rs @@ -0,0 +1,87 @@ +use std::time::Duration; + +use anyhow::Result; +use serialport::{available_ports, SerialPort, SerialPortBuilder, SerialPortInfo, SerialPortType}; +use thiserror::Error; + +use crate::cli::Listen; + +pub fn list_available_ports(include_all: bool) -> Vec { + available_ports() + .unwrap_or_default() + .into_iter() + .filter(|port| include_all || matches!(port.port_type, SerialPortType::UsbPort(_))) + .collect::>() +} + +pub fn select_serial_port(explicit: &Option) -> Result { + if let Some(p) = explicit { + return Ok(p.clone()); + } + let ports = list_available_ports(false).into_iter().map(|p| p.port_name).collect::>(); + decide_port(None, ports) +} + +// Pure decision function for easier testing +pub(crate) fn decide_port( + explicit: Option, + available: Vec, +) -> Result { + if let Some(p) = explicit { + return Ok(p); + } + match available.len() { + 0 => Err(SerialSelectError::NoPorts.into()), + 1 => Ok(available[0].clone()), + _ => Err(SerialSelectError::MultiplePorts { list: available }.into()), + } +} + +#[derive(Debug, Error)] +pub enum SerialSelectError { + #[error("No serial ports found. Re-run with --serial or use --all in 'ports' to inspect.")] + NoPorts, + #[error("Multiple serial ports detected: {list:?}. Please specify --serial .")] + MultiplePorts { list: Vec }, +} + +pub fn configure_serial( + builder: SerialPortBuilder, + listen: &Listen, +) -> serialport::Result> { + builder + .data_bits(listen.data_bits.clone().into()) + .parity(listen.parity.clone().into()) + .stop_bits(listen.stop_bits.clone().into()) + .timeout(Duration::from_millis(200)) + .open() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_decide_port_explicit() { + let r = decide_port(Some("/dev/ttyUSB9".into()), vec!["/dev/ttyUSB0".into()]).unwrap(); + assert_eq!(r, "/dev/ttyUSB9"); + } + + #[test] + fn test_decide_port_none_single() { + let r = decide_port(None, vec!["/dev/ttyUSB0".into()]).unwrap(); + assert_eq!(r, "/dev/ttyUSB0"); + } + + #[test] + fn test_decide_port_none_zero() { + let err = decide_port(None, vec![]).unwrap_err(); + assert!(err.to_string().contains("No serial ports")); + } + + #[test] + fn test_decide_port_none_multiple() { + let err = decide_port(None, vec!["/dev/ttyUSB0".into(), "/dev/ttyUSB1".into()]).unwrap_err(); + assert!(err.to_string().contains("Multiple serial ports")); + } +} diff --git a/apps/sergw/src/serial/mod.rs b/apps/sergw/src/serial/mod.rs new file mode 100644 index 0000000..99dd6fa --- /dev/null +++ b/apps/sergw/src/serial/mod.rs @@ -0,0 +1,2 @@ +pub mod io; +pub use io::*; diff --git a/apps/sergw/src/state.rs b/apps/sergw/src/state.rs index e964223..a3e527b 100644 --- a/apps/sergw/src/state.rs +++ b/apps/sergw/src/state.rs @@ -1,18 +1,144 @@ -use std::collections::HashMap; use std::net::SocketAddr; -use std::sync::mpsc::Sender; -use std::sync::Arc; + +use bytes::Bytes; +use crossbeam_channel as channel; +use dashmap::DashMap; pub struct SharedState { - pub connections: HashMap>>, - pub verbose: bool, + // outbound to TCP, concurrent map to avoid global mutex during broadcast + pub tcp_connections: DashMap>, } impl SharedState { - pub fn new(verbose: bool) -> Self { - SharedState { - connections: HashMap::new(), - verbose, + pub fn new() -> Self { + Self { + tcp_connections: DashMap::new(), + } + } + + pub fn insert( + &self, + addr: SocketAddr, + tx: channel::Sender, + ) { + self.tcp_connections.insert(addr, tx); + } + + pub fn remove( + &self, + addr: &SocketAddr, + ) { + self.tcp_connections.remove(addr); + } + + pub fn dispose(&self) { + self.tcp_connections.clear(); + } + + pub fn broadcast( + &self, + data: Bytes, + ) { + // Clone senders without holding any global lock; DashMap provides + // per-bucket locking which is brief during iteration. + let snapshot: Vec<(SocketAddr, channel::Sender)> = self.tcp_connections.iter().map(|e| (*e.key(), e.value().clone())).collect(); + + let mut to_remove: Vec = Vec::new(); + for (addr, tx) in snapshot.into_iter() { + match tx.try_send(data.clone()) { + Ok(()) => {} + Err(channel::TrySendError::Full(_)) => { + // Slow client: drop this client to enforce backpressure + to_remove.push(addr); + } + Err(channel::TrySendError::Disconnected(_)) => { + to_remove.push(addr); + } + } } + + for addr in to_remove { + self.remove(&addr); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn broadcast_removes_dead_receivers() { + let (tx_alive, rx_alive) = channel::bounded::(1); + let (tx_dead, _rx_dead) = channel::bounded::(1); + drop(_rx_dead); // drop to simulate dead receiver + + let state = SharedState::new(); + let a1: SocketAddr = "127.0.0.1:10000".parse().unwrap(); + let a2: SocketAddr = "127.0.0.1:10001".parse().unwrap(); + state.insert(a1, tx_alive); + state.insert(a2, tx_dead); + + state.broadcast(Bytes::from_static(b"hello")); + + // Alive should receive + assert_eq!(rx_alive.recv().unwrap(), Bytes::from_static(b"hello")); + // Dead should be removed + assert!(!state.tcp_connections.contains_key(&a2)); + } + + #[test] + fn broadcast_removes_slow_receivers_on_full() { + let (tx_alive, rx_alive) = channel::bounded::(1); + let (tx_slow, _rx_slow) = channel::bounded::(1); + + let state = SharedState::new(); + let a_alive: SocketAddr = "127.0.0.1:11000".parse().unwrap(); + let a_slow: SocketAddr = "127.0.0.1:11001".parse().unwrap(); + state.insert(a_alive, tx_alive); + state.insert(a_slow, tx_slow); + + // First broadcast fills both queues + state.broadcast(Bytes::from_static(b"one")); + + // Drain the alive receiver so it won't be full for the next broadcast + assert_eq!(rx_alive.recv().unwrap(), Bytes::from_static(b"one")); + + // Second broadcast: slow stays full and should be removed; alive receives + state.broadcast(Bytes::from_static(b"two")); + + assert_eq!(rx_alive.recv().unwrap(), Bytes::from_static(b"two")); + assert!(!state.tcp_connections.contains_key(&a_slow)); + } + + #[test] + fn broadcast_delivers_to_multiple_alive_receivers() { + let (tx1, rx1) = channel::unbounded::(); + let (tx2, rx2) = channel::unbounded::(); + + let state = SharedState::new(); + let a1: SocketAddr = "127.0.0.1:12000".parse().unwrap(); + let a2: SocketAddr = "127.0.0.1:12001".parse().unwrap(); + state.insert(a1, tx1); + state.insert(a2, tx2); + + state.broadcast(Bytes::from_static(b"abc")); + + assert_eq!(rx1.recv().unwrap(), Bytes::from_static(b"abc")); + assert_eq!(rx2.recv().unwrap(), Bytes::from_static(b"abc")); + } + + #[test] + fn dispose_clears_all_connections() { + let (tx1, _rx1) = channel::unbounded::(); + let (tx2, _rx2) = channel::unbounded::(); + let state = SharedState::new(); + let a1: SocketAddr = "127.0.0.1:13000".parse().unwrap(); + let a2: SocketAddr = "127.0.0.1:13001".parse().unwrap(); + state.insert(a1, tx1); + state.insert(a2, tx2); + + state.dispose(); + assert!(state.tcp_connections.is_empty()); } } diff --git a/apps/sergw/src/tcp_server.rs b/apps/sergw/src/tcp_server.rs deleted file mode 100644 index d0d37cb..0000000 --- a/apps/sergw/src/tcp_server.rs +++ /dev/null @@ -1,118 +0,0 @@ -use std::io::{Read, Write}; -use std::net::TcpListener; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::mpsc::{channel, Sender}; -use std::sync::{Arc, Mutex}; -use std::time::Duration; - -use crate::state::SharedState; - -const BUFFER_SIZE: usize = 1024; -const TCP_LOOP_SLEEP_DURATION: Duration = Duration::from_millis(10); - -pub fn broadcast_data( - shared_state: &Arc>, - data: Arc<[u8]>, -) { - let state = shared_state.lock().expect("Mutex was poisoned"); - if state.verbose { - println!( - "Broadcasting {} bytes to {} TCP connections: {:?}", - data.len(), - state.connections.len(), - data - ); - } - for (addr, sender) in state.connections.iter() { - if sender.send(data.clone()).is_err() && state.verbose { - eprintln!("Failed to send to {}, client handler will clean it up.", addr); - } - } -} - -pub fn spawn_tcp_listener( - host: &str, - shared_state: &Arc>, - shutdown_flag: &Arc, - serial_writer_tx: Sender>, -) -> Result, Box> { - let listener = TcpListener::bind(host)?; - listener.set_nonblocking(true)?; - - let shared_state_clone = Arc::clone(shared_state); - let shutdown_flag_clone = Arc::clone(shutdown_flag); - - let handle = std::thread::spawn(move || { - for stream_result in listener.incoming() { - if shutdown_flag_clone.load(Ordering::Acquire) { - println!("TCP listener shutting down..."); - break; - } - - match stream_result { - Ok(mut stream) => { - let addr = stream.peer_addr().expect("Could not get peer address"); - println!("Accepted connection from: {}", addr); - - let (serial_sender, serial_receiver) = channel::>(); - shared_state_clone - .lock() - .expect("Mutex was poisoned") - .connections - .insert(addr, serial_sender); - - let shared_state_for_cleanup = Arc::clone(&shared_state_clone); - let serial_writer_tx_clone = serial_writer_tx.clone(); - - std::thread::spawn(move || { - stream.set_nonblocking(true).expect("Failed to set stream to non-blocking"); - let mut tcp_buffer = [0; BUFFER_SIZE]; - - loop { - // Read from TCP to serial - match stream.read(&mut tcp_buffer) { - Ok(0) => break, // Connection closed by client - Ok(n) => { - let data = Arc::from(&tcp_buffer[..n]); - if serial_writer_tx_clone.send(data).is_err() { - eprintln!("Serial writer channel closed, closing connection."); - break; - } - } - Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {} - Err(_) => break, // Connection error - } - - // Read from serial to TCP - match serial_receiver.try_recv() { - Ok(data) => { - if stream.write_all(&data).is_err() { - break; // Connection error - } - } - Err(std::sync::mpsc::TryRecvError::Empty) => { - std::thread::sleep(TCP_LOOP_SLEEP_DURATION); - } - Err(std::sync::mpsc::TryRecvError::Disconnected) => break, - } - } - - println!("Closing connection from: {}", addr); - shared_state_for_cleanup.lock().expect("Mutex was poisoned").connections.remove(&addr); - }); - } - Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { - std::thread::sleep(TCP_LOOP_SLEEP_DURATION); - continue; - } - Err(e) => { - eprintln!("TCP accept error: {}. Shutting down listener.", e); - break; - } - } - } - println!("TCP listener thread finished."); - }); - - Ok(handle) -} diff --git a/apps/sergw/src/ui/inspector.rs b/apps/sergw/src/ui/inspector.rs new file mode 100644 index 0000000..462ff42 --- /dev/null +++ b/apps/sergw/src/ui/inspector.rs @@ -0,0 +1,116 @@ +use std::collections::VecDeque; +use std::net::SocketAddr; + +// no time imports needed here +use bytes::Bytes; +use ratatui::layout::Rect; +use ratatui::text::{Line, Span}; +use ratatui::widgets::{Paragraph, Wrap}; + +#[derive(Copy, Clone, Debug, PartialEq, Eq)] +pub enum DirectionTag { + Inbound, + Outbound(SocketAddr), +} + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub enum DeviceId { + Serial, + Client(SocketAddr), +} + +#[derive(Copy, Clone, Debug, PartialEq, Eq)] +pub enum DumpFormat { + Hex, + Ascii, + Dec, +} + +#[derive(Clone, Debug)] +pub struct Sample { + pub dir: DirectionTag, + pub data: Bytes, +} + +pub struct InspectorState { + pub format: DumpFormat, + pub paused: bool, + pub devices: Vec, + pub selected: usize, + pub scroll: usize, + pub capture: VecDeque, +} + +impl InspectorState { + pub fn new() -> Self { + Self { + format: DumpFormat::Hex, + paused: false, + devices: vec![DeviceId::Serial], + selected: 0, + scroll: 0, + capture: VecDeque::with_capacity(2048), + } + } +} + +pub fn dump_bytes( + buf: &[u8], + fmt: DumpFormat, + max: usize, +) -> String { + let slice = &buf[..buf.len().min(max)]; + match fmt { + DumpFormat::Hex => slice.iter().map(|b| format!("{b:02x} ")).collect(), + DumpFormat::Ascii => { + let mut s = String::new(); + for &b in slice { + if b == b'\n' || b == b'\r' { + continue; + } + if b.is_ascii_graphic() || b == b' ' { + s.push(b as char); + } else { + s.push('.'); + } + } + s + } + DumpFormat::Dec => slice.iter().map(|b| format!("{b:03} ")).collect(), + } +} + +// Render wrapped text for inspector messages. Returns a Paragraph with Wrap enabled. +pub fn inspector_paragraph( + state: &InspectorState, + area: Rect, +) -> Paragraph<'static> { + let filter = state.devices.get(state.selected); + // Build lines as strings first + let lines: Vec = state + .capture + .iter() + .filter_map(|s| { + let dev = match s.dir { + DirectionTag::Inbound => DeviceId::Serial, + DirectionTag::Outbound(a) => DeviceId::Client(a), + }; + if let Some(sel) = filter { + if &dev != sel { + return None; + } + } + Some(dump_bytes(&s.data, state.format, 4096)) + }) + .collect(); + + if state.scroll > 0 { + let _ = &lines; + } + let start = lines.len().saturating_sub(area.height.saturating_sub(2) as usize + state.scroll); + let visible = lines.into_iter().skip(start); + + let text_lines: Vec = visible.map(|s| Line::from(Span::raw(s))).collect(); + + Paragraph::new(text_lines).wrap(Wrap { trim: false }) +} diff --git a/apps/sergw/src/ui/mod.rs b/apps/sergw/src/ui/mod.rs new file mode 100644 index 0000000..2557331 --- /dev/null +++ b/apps/sergw/src/ui/mod.rs @@ -0,0 +1,2 @@ +pub mod inspector; +pub mod overview; diff --git a/apps/sergw/src/ui/overview.rs b/apps/sergw/src/ui/overview.rs new file mode 100644 index 0000000..dd3dbec --- /dev/null +++ b/apps/sergw/src/ui/overview.rs @@ -0,0 +1,291 @@ +use std::sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + Arc, +}; +use std::time::{Duration, Instant}; + +use anyhow::Result; +use crossbeam_channel::Receiver; +use crossterm::{ + event::{self, Event, KeyCode, KeyModifiers}, + execute, + terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen}, +}; +use ratatui::{ + backend::CrosstermBackend, + layout::{Constraint, Direction, Layout}, + widgets::{Block, Borders, List, ListItem, Paragraph, Tabs}, + Terminal, +}; + +use crate::metrics::ThroughputAverager; +use crate::state::SharedState; +use crate::ui::inspector::{DeviceId, DumpFormat, InspectorState}; + +#[derive(Default)] +pub struct Counters { + pub bytes_in: AtomicU64, + pub bytes_out: AtomicU64, +} + +pub fn run_tui( + shared: Arc, + counters: Arc, + events: Receiver, + insp_rx: Receiver, + stop: Arc, +) -> Result<()> { + enable_raw_mode()?; + let mut stdout = std::io::stdout(); + execute!(stdout, EnterAlternateScreen)?; + let backend = CrosstermBackend::new(stdout); + let mut terminal = Terminal::new(backend)?; + + let mut logs: Vec = Vec::new(); + let mut log_scroll: usize = 0; + let mut active_tab: usize = 0; // 0: Overview, 1: Inspector + let mut _prev_tab: usize = active_tab; + let mut insp = InspectorState::new(); + let mut last_in = 0u64; + let mut last_out = 0u64; + let mut avg_in = ThroughputAverager::new(5.0); + let mut avg_out = ThroughputAverager::new(5.0); + let mut last_time = Instant::now(); + + while !stop.load(Ordering::Relaxed) { + while let Ok(ev) = events.try_recv() { + logs.push(ev); + if logs.len() > 100 { + logs.remove(0); + } + } + + let now = Instant::now(); + let dt = now.duration_since(last_time).as_secs_f64().max(0.001); + let bi = counters.bytes_in.load(Ordering::Relaxed); + let bo = counters.bytes_out.load(Ordering::Relaxed); + let tin = avg_out.update(bi - last_in, dt) as u64; // TCP -> serial (outbound, smoothed) + let tout = avg_in.update(bo - last_out, dt) as u64; // serial -> TCP (inbound, smoothed) + last_in = bi; + last_out = bo; + last_time = now; + + // Pull inspector samples; skip if paused + while let Ok(s) = insp_rx.try_recv() { + if !insp.paused { + // Track devices + match s.dir { + crate::ui::inspector::DirectionTag::Inbound => { + if !insp.devices.iter().any(|d| matches!(d, DeviceId::Serial)) { + insp.devices.insert(0, DeviceId::Serial); + } + } + crate::ui::inspector::DirectionTag::Outbound(addr) => { + if !insp.devices.iter().any(|d| matches!(d, DeviceId::Client(a) if *a == addr)) { + insp.devices.push(DeviceId::Client(addr)); + } + } + } + insp.capture.push_back(s); + if insp.capture.len() > 4096 { + insp.capture.pop_front(); + } + } + } + + terminal.draw(|f| { + // Top-level: header tabs, main, footer + let outer = Layout::default() + .direction(Direction::Vertical) + .constraints( + [ + Constraint::Length(1), // Tabs header + Constraint::Min(0), // Main + Constraint::Length(1), // Footer + ] + .as_ref(), + ) + .split(f.size()); + + // Tabs header + let titles = ["Overview", "Inspector"].iter().map(|t| (*t).to_string()); + let tabs = Tabs::new(titles).select(active_tab); + f.render_widget(tabs, outer[0]); + + if active_tab == 0 { + // Overview: connections, throughput, events + let main = outer[1]; + let sub = Layout::default() + .direction(Direction::Vertical) + .constraints( + [ + Constraint::Length(5), // Connections + Constraint::Length(4), // Throughput + Constraint::Min(0), // Events + ] + .as_ref(), + ) + .split(main); + + let items: Vec = shared.tcp_connections.iter().map(|e| ListItem::new(e.key().to_string())).collect(); + let list = List::new(items).block(Block::default().title("Connections").borders(Borders::ALL)); + f.render_widget(list, sub[0]); + + let throughput = Paragraph::new(format!("Inbound: {tout} B/s\nOutbound: {tin} B/s")) + .block(Block::default().title("Throughput").borders(Borders::ALL)); + f.render_widget(throughput, sub[1]); + + let viewport = sub[2].height.saturating_sub(2) as usize; + let start = logs.len().saturating_sub(viewport + log_scroll); + let log_items: Vec = logs.iter().skip(start).map(|l| ListItem::new(l.clone())).collect(); + let log_list = List::new(log_items).block(Block::default().title("Events").borders(Borders::ALL)); + f.render_widget(log_list, sub[2]); + } else { + // Inspector tab: header summary + dump list + let main = outer[1]; + // Sidebar + main list + let columns = Layout::default() + .direction(Direction::Horizontal) + .constraints( + [ + Constraint::Length(24), // sidebar + Constraint::Min(0), // inspector content + ] + .as_ref(), + ) + .split(main); + + // Sidebar devices + let dev_labels: Vec = insp + .devices + .iter() + .map(|d| match d { + DeviceId::Serial => "serial".to_string(), + DeviceId::Client(a) => format!("{a}"), + }) + .collect(); + let dev_items: Vec = dev_labels + .iter() + .enumerate() + .map(|(i, s)| { + let prefix = if i == insp.selected { "> " } else { " " }; + ListItem::new(format!("{prefix}{s}")) + }) + .collect(); + let dev_list = List::new(dev_items).block(Block::default().title("Devices").borders(Borders::ALL)); + f.render_widget(dev_list, columns[0]); + + let sub = Layout::default() + .direction(Direction::Vertical) + .constraints( + [ + Constraint::Length(1), // header + Constraint::Min(0), // list + ] + .as_ref(), + ) + .split(columns[1]); + + let header = Paragraph::new(format!( + "fmt: {:?} | status: {}", + insp.format, + if insp.paused { "paused" } else { "resumed" } + )); + f.render_widget(header, sub[0]); + + let para = crate::ui::inspector::inspector_paragraph(&insp, sub[1]); + let block = Block::default().title("Messages").borders(Borders::ALL); + f.render_widget(para.block(block), sub[1]); + } + + // Sticky footer with keybinds + let footer = if active_tab == 0 { + Paragraph::new("Tab: inspector | q: quit | ↑/↓/Home: scroll events | c: clear events") + } else { + Paragraph::new( + "Tab: overview | q: quit | t: toggle type | p: pause/resume | ↑/↓: select device | PgUp/PgDn: scroll | Home: top | c: clear", + ) + }; + f.render_widget(footer, outer[2]); + })?; + + if event::poll(Duration::from_millis(200))? { + if let Event::Key(key) = event::read()? { + if key.code == KeyCode::Char('q') || (key.code == KeyCode::Char('c') && key.modifiers.contains(KeyModifiers::CONTROL)) { + stop.store(true, Ordering::Relaxed); + } else if key.code == KeyCode::Tab { + _prev_tab = active_tab; + active_tab = (active_tab + 1) % 2; + if active_tab == 0 { + // leaving inspector: clear state + insp.capture.clear(); + insp.devices = vec![DeviceId::Serial]; + insp.selected = 0; + insp.scroll = 0; + insp.paused = false; + } + } else if active_tab == 0 { + match key.code { + KeyCode::Up => { + log_scroll = log_scroll.saturating_add(1); + } + KeyCode::Down => { + log_scroll = log_scroll.saturating_sub(1); + } + KeyCode::Home => { + log_scroll = 0; + } + KeyCode::Char('c') => { + logs.clear(); + log_scroll = 0; + } + _ => {} + } + } else { + match key.code { + KeyCode::Char('t') => { + insp.format = match insp.format { + DumpFormat::Hex => DumpFormat::Ascii, + DumpFormat::Ascii => DumpFormat::Dec, + DumpFormat::Dec => DumpFormat::Hex, + }; + } + KeyCode::Char('p') => { + insp.paused = !insp.paused; + } + KeyCode::Char('c') => { + insp.capture.clear(); + insp.scroll = 0; + } + KeyCode::Up => { + if insp.selected > 0 { + insp.selected -= 1; + } + } + KeyCode::Down => { + if insp.selected + 1 < insp.devices.len() { + insp.selected += 1; + } + } + KeyCode::PageUp => { + let limit = insp.capture.len(); + if insp.scroll < limit { + insp.scroll = insp.scroll.saturating_add(5); + } + } + KeyCode::PageDown => { + insp.scroll = insp.scroll.saturating_sub(5); + } + KeyCode::Home => insp.scroll = 0, + _ => {} + } + } + } + } + } + + disable_raw_mode()?; + execute!(terminal.backend_mut(), LeaveAlternateScreen)?; + terminal.show_cursor()?; + Ok(()) +}