diff --git a/packages/client-v2/.gitignore b/packages/client-v2/.gitignore new file mode 100644 index 0000000..87c6086 --- /dev/null +++ b/packages/client-v2/.gitignore @@ -0,0 +1,2 @@ +*.mp3 +*.wav \ No newline at end of file diff --git a/packages/client-v2/Cargo.lock b/packages/client-v2/Cargo.lock new file mode 100644 index 0000000..57a53c0 --- /dev/null +++ b/packages/client-v2/Cargo.lock @@ -0,0 +1,547 @@ +# This file is automatically @generated by Cargo. +# It is not intended for manual editing. +version = 4 + +[[package]] +name = "alsa" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "812947049edcd670a82cd5c73c3661d2e58468577ba8489de58e1a73c04cbd5d" +dependencies = [ + "alsa-sys", + "bitflags", + "cfg-if", + "libc", +] + +[[package]] +name = "alsa-sys" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ad7569085a265dd3f607ebecce7458eaab2132a84393534c95b18dcbc3f31e04" +dependencies = [ + "libc", + "pkg-config", +] + +[[package]] +name = "anyhow" +version = "1.0.100" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61" + +[[package]] +name = "atomic-polyfill" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8cf2bce30dfe09ef0bfaef228b9d414faaf7e563035494d7fe092dba54b300f4" +dependencies = [ + "critical-section", +] + +[[package]] +name = "audiopus_sys" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "62314a1546a2064e033665d658e88c620a62904be945f8147e6b16c3db9f8651" +dependencies = [ + "cmake", + "log", + "pkg-config", +] + +[[package]] +name = "bitflags" +version = "2.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "812e12b5285cc515a9c72a5c1d3b6d46a19dac5acfef5265968c166106e31dd3" + +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + +[[package]] +name = "bytes" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b35204fbdc0b3f4446b89fc1ac2cf84a8a68971995d0bf2e925ec7cd960f9cb3" + +[[package]] +name = "cc" +version = "1.2.51" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a0aeaff4ff1a90589618835a598e545176939b97874f7abc7851caa0618f203" +dependencies = [ + "find-msvc-tools", + "shlex", +] + +[[package]] +name = "cfg-if" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" + +[[package]] +name = "cmake" +version = "0.1.57" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75443c44cd6b379beb8c5b45d85d0773baf31cce901fe7bb252f4eff3008ef7d" +dependencies = [ + "cc", +] + +[[package]] +name = "cobs" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fa961b519f0b462e3a3b4a34b64d119eeaca1d59af726fe450bbba07a9fc0a1" +dependencies = [ + "thiserror", +] + +[[package]] +name = "critical-section" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" + +[[package]] +name = "embedded-io" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef1a6892d9eef45c8fa6b9e0086428a2cca8491aca8f787c534a3d6d0bcb3ced" + +[[package]] +name = "embedded-io" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d" + +[[package]] +name = "errno" +version = "0.3.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" +dependencies = [ + "libc", + "windows-sys 0.60.2", +] + +[[package]] +name = "find-msvc-tools" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "645cbb3a84e60b7531617d5ae4e57f7e27308f6445f5abf653209ea76dec8dff" + +[[package]] +name = "hash32" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0c35f58762feb77d74ebe43bdbc3210f09be9fe6742234d573bacc26ed92b67" +dependencies = [ + "byteorder", +] + +[[package]] +name = "heapless" +version = "0.7.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cdc6457c0eb62c71aac4bc17216026d8410337c4126773b9c5daba343f17964f" +dependencies = [ + "atomic-polyfill", + "hash32", + "rustc_version", + "serde", + "spin", + "stable_deref_trait", +] + +[[package]] +name = "libc" +version = "0.2.178" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37c93d8daa9d8a012fd8ab92f088405fb202ea0b6ab73ee2482ae66af4f42091" + +[[package]] +name = "lock_api" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "224399e74b87b5f3557511d98dff8b14089b3dadafcab6bb93eab67d3aace965" +dependencies = [ + "scopeguard", +] + +[[package]] +name = "log" +version = "0.4.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" + +[[package]] +name = "mio" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a69bcab0ad47271a0234d9422b131806bf3968021e5dc9328caf2d4cd58557fc" +dependencies = [ + "libc", + "wasi", + "windows-sys 0.61.2", +] + +[[package]] +name = "opus" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6526409b274a7e98e55ff59d96aafd38e6cd34d46b7dbbc32ce126dffcd75e8e" +dependencies = [ + "audiopus_sys", + "libc", +] + +[[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", +] + +[[package]] +name = "pin-project-lite" +version = "0.2.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3b3cff922bd51709b605d9ead9aa71031d81447142d828eb4a6eba76fe619f9b" + +[[package]] +name = "pkg-config" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" + +[[package]] +name = "postcard" +version = "1.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6764c3b5dd454e283a30e6dfe78e9b31096d9e32036b5d1eaac7a6119ccb9a24" +dependencies = [ + "cobs", + "embedded-io 0.4.0", + "embedded-io 0.6.1", + "heapless", + "serde", +] + +[[package]] +name = "proc-macro2" +version = "1.0.104" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9695f8df41bb4f3d222c95a67532365f569318332d03d5f3f67f37b20e6ebdf0" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "quote" +version = "1.0.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a338cc41d27e6cc6dce6cefc13a0729dfbb81c262b1f519331575dd80ef3067f" +dependencies = [ + "proc-macro2", +] + +[[package]] +name = "redox_syscall" +version = "0.5.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" +dependencies = [ + "bitflags", +] + +[[package]] +name = "rustc_version" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfcb3a22ef46e85b45de6ee7e79d063319ebb6594faafcf1c225ea92ab6e9b92" +dependencies = [ + "semver", +] + +[[package]] +name = "scopeguard" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" + +[[package]] +name = "semver" +version = "1.0.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d767eb0aabc880b29956c35734170f26ed551a859dbd361d140cdbeca61ab1e2" + +[[package]] +name = "serde" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +dependencies = [ + "serde_core", + "serde_derive", +] + +[[package]] +name = "serde_core" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" +dependencies = [ + "serde_derive", +] + +[[package]] +name = "serde_derive" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "shlex" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" + +[[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 = "smallvec" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" + +[[package]] +name = "socket2" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17129e116933cf371d018bb80ae557e889637989d8638274fb25622827b03881" +dependencies = [ + "libc", + "windows-sys 0.60.2", +] + +[[package]] +name = "spin" +version = "0.9.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] + +[[package]] +name = "stable_deref_trait" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" + +[[package]] +name = "syn" +version = "2.0.111" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "390cc9a294ab71bdb1aa2e99d13be9c753cd2d7bd6560c77118597410c4d2e87" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + +[[package]] +name = "thiserror" +version = "2.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f63587ca0f12b72a0600bcba1d40081f830876000bb46dd2337a3051618f4fc8" +dependencies = [ + "thiserror-impl", +] + +[[package]] +name = "thiserror-impl" +version = "2.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3ff15c8ecd7de3849db632e14d18d2571fa09dfc5ed93479bc4485c7a517c913" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "tokio" +version = "1.48.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff360e02eab121e0bc37a2d3b4d4dc622e6eda3a8e5253d5435ecf5bd4c68408" +dependencies = [ + "bytes", + "libc", + "mio", + "parking_lot", + "pin-project-lite", + "signal-hook-registry", + "socket2", + "tokio-macros", + "windows-sys 0.61.2", +] + +[[package]] +name = "tokio-macros" +version = "2.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af407857209536a95c8e56f8231ef2c2e2aff839b22e07a1ffcbc617e9db9fa5" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "unicode-ident" +version = "1.0.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9312f7c4f6ff9069b165498234ce8be658059c6728633667c526e27dc2cf1df5" + +[[package]] +name = "wasi" +version = "0.11.1+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" + +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + +[[package]] +name = "windows-sys" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" +dependencies = [ + "windows-targets", +] + +[[package]] +name = "windows-sys" +version = "0.61.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" +dependencies = [ + "windows-link", +] + +[[package]] +name = "windows-targets" +version = "0.53.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4945f9f551b88e0d65f3db0bc25c33b8acea4d9e41163edf90dcd0b19f9069f3" +dependencies = [ + "windows-link", + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_gnullvm", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006" + +[[package]] +name = "windows_i686_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "960e6da069d81e09becb0ca57a65220ddff016ff2d6af6a223cf372a506593a3" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c" + +[[package]] +name = "windows_i686_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" + +[[package]] +name = "xiao" +version = "0.1.0" +dependencies = [ + "alsa", + "anyhow", + "opus", + "postcard", + "serde", + "tokio", +] diff --git a/packages/client-v2/Cargo.toml b/packages/client-v2/Cargo.toml new file mode 100644 index 0000000..15fee32 --- /dev/null +++ b/packages/client-v2/Cargo.toml @@ -0,0 +1,22 @@ +[package] +name = "xiao" +version = "0.1.0" +edition = "2024" + +[profile.release] +lto = true +opt-level = "s" +codegen-units = 1 +panic = "abort" +strip = true +debug = false + +[dependencies] +opus = "0.3" +anyhow = "1.0" +tokio = { version = "1.48", features = ["full"] } +serde = { version = "1.0", features = ["derive"] } +postcard = { version = "1.0", features = ["alloc", "use-std"] } + +[target.'cfg(target_os = "linux")'.dependencies] +alsa = "0.11" diff --git a/packages/client-v2/Makefile b/packages/client-v2/Makefile new file mode 100644 index 0000000..0ecbf3b --- /dev/null +++ b/packages/client-v2/Makefile @@ -0,0 +1,10 @@ +build: + docker run --rm -v $(shell pwd):/app idootop/open-xiaoai-runtime:oh2p \ + cargo build --target armv7-unknown-linux-gnueabihf --release + +# 部署到小爱音箱(调试自用) +deploy: + dd if=target/armv7-unknown-linux-gnueabihf/release/xiao \ + | sshpass -p open-xiaoai ssh -o HostKeyAlgorithms=+ssh-rsa root@192.168.31.153 "dd of=/data/xiao" + dd if=target/armv7-unknown-linux-gnueabihf/release/xiao \ + | sshpass -p open-xiaoai ssh -o HostKeyAlgorithms=+ssh-rsa root@192.168.31.235 "dd of=/data/xiao" \ No newline at end of file diff --git a/packages/client-v2/README.md b/packages/client-v2/README.md new file mode 100644 index 0000000..3162c40 --- /dev/null +++ b/packages/client-v2/README.md @@ -0,0 +1,7 @@ +# Open-XiaoAI Client V2 + +> 开发中,敬请期待 + +## License + +MIT License © 2026-PRESENT [Del Wang](https://del.wang) diff --git a/packages/client-v2/src/app/entry.rs b/packages/client-v2/src/app/entry.rs new file mode 100644 index 0000000..8d281c3 --- /dev/null +++ b/packages/client-v2/src/app/entry.rs @@ -0,0 +1,35 @@ +#![cfg(target_os = "linux")] + +use crate::app::master::run_master; +use crate::app::slave::run_slave; +use crate::net::protocol::ChannelRole; +use anyhow::Result; +use std::env; + +pub async fn run_xiao() -> Result<()> { + let args: Vec = env::args().collect(); + if args.len() < 3 { + eprintln!("用法: {} [master|slave] [left|right]", args[0]); + return Ok(()); + } + + let mode = if args[1].to_lowercase() == "master" { + "主节点" + } else { + "从节点" + }; + + let role = if args[2].to_lowercase() == "left" { + ChannelRole::Left + } else { + ChannelRole::Right + }; + + println!("🚗 当前为: {} {}", mode, role.to_string()); + + if mode == "主节点" { + run_master(role).await + } else { + run_slave(role).await + } +} diff --git a/packages/client-v2/src/app/master.rs b/packages/client-v2/src/app/master.rs new file mode 100644 index 0000000..4d29be4 --- /dev/null +++ b/packages/client-v2/src/app/master.rs @@ -0,0 +1,365 @@ +#![cfg(target_os = "linux")] + +use crate::audio::codec::OpusCodec; +use crate::audio::config::AudioConfig; +use crate::audio::player::AudioPlayer; +use crate::net::discovery::Discovery; +use crate::net::network::{ControlConnection, MasterNetwork}; +use crate::net::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::utils::alsa::AlsaRedirector; +use crate::utils::sync::now_us; +use anyhow::{Result, anyhow}; +use std::net::SocketAddr; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::signal::unix::{SignalKind, signal}; +use tokio::sync::Mutex; + +pub const SERVER_TCP_PORT: u16 = 53531; + +#[derive(Clone)] +struct SlaveSession { + udp_addr: SocketAddr, + role: ChannelRole, +} + +pub async fn run_master(master_role: ChannelRole) -> Result<()> { + // 0. 设置 ALSA 重定向 + println!("🔥 启动中,请稍等..."); + let _alsa_guard = AlsaRedirector::new()?; + + // 1. 设置网络 (UDP + TCP) + let network = MasterNetwork::setup(SERVER_TCP_PORT).await?; + let audio_socket = network.audio_socket().clone_inner(); + + // 2. 启动服务发现广播 + Discovery::start_broadcast(SERVER_TCP_PORT).await?; + + println!("✅ 服务已启动,等待连接..."); + + let shutdown_flag = Arc::new(AtomicBool::new(false)); + let slaves = Arc::new(Mutex::new(Vec::::new())); + + // 3. 启动连接监听任务 + let slaves_clone = slaves.clone(); + let audio_socket_clone = audio_socket.clone(); + tokio::spawn(async move { + loop { + match network.accept().await { + Ok((control_conn, client_addr)) => { + let slaves_for_session = slaves_clone.clone(); + let audio_socket_for_session = audio_socket_clone.clone(); + tokio::spawn(async move { + if let Err(e) = handle_master_session( + control_conn, + audio_socket_for_session, + slaves_for_session, + client_addr.to_string(), + ) + .await + { + eprintln!("❌ 会话错误: {:?}", e); + } + }); + } + Err(e) => { + eprintln!("❌ Accept 错误: {:?}", e); + } + } + } + }); + + // 4. 音频处理主循环 + let config = AudioConfig::music(); + let encode_config = AudioConfig { + channels: 1, + vbr: true, + ..AudioConfig::music() + }; + let player = AudioPlayer::new(&AudioConfig { + channels: 2, + playback_device: "plug:original_default".into(), + ..config.clone() + })?; + + let mut raw_buf = vec![0u8; config.frame_size * 2 * 2]; + let mut pcm_out = vec![0i16; config.frame_size * 2]; + let mut left_pcm = vec![0i16; config.frame_size]; + let mut right_pcm = vec![0i16; config.frame_size]; + let mut opus_out = vec![0u8; 1500]; + let mut seq = 0u32; + + // 播放延迟: 只需覆盖网络延迟 + 时钟偏移 + let delay_us = 100_000; // 100ms 基础延迟 + let frame_duration_us = + (config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128; + + let mut stream_start_ts = 0; + let mut stream_start_seq = 0; + + let shutdown_flag_clone = shutdown_flag.clone(); + let audio_loop = async move { + loop { + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; + } + + // 打开 FIFO + let mut fifo = match tokio::fs::File::open(AlsaRedirector::fifo_path()).await { + Ok(f) => f, + Err(_) => { + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; + } + tokio::time::sleep(Duration::from_secs(1)).await; + continue; + } + }; + + // 每个新流开始时,重置编码器状态以避免残留音频导致爆音 + let mut left_encoder = OpusCodec::new(&encode_config)?; + let mut right_encoder = OpusCodec::new(&encode_config)?; + + loop { + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; + } + + // 从 FIFO 读取 + if let Err(_) = fifo.read_exact(&mut raw_buf).await { + break; // FIFO 关闭,重新打开 + } + + let active_slaves = { + let s = slaves.lock().await; + if s.is_empty() { None } else { Some(s.clone()) } + }; + + let now = now_us(); + if stream_start_ts == 0 { + stream_start_ts = now; + stream_start_seq = seq; + } + + // 计算该帧应当播放的目标时间 + // target_ts = 数据包发送时间 + 播放延迟 + let target_ts = stream_start_ts + + ((seq - stream_start_seq) as u128 * frame_duration_us) + + delay_us; + + // 提取 PCM 数据 + for i in 0..config.frame_size { + left_pcm[i] = i16::from_le_bytes([raw_buf[i * 4], raw_buf[i * 4 + 1]]); + right_pcm[i] = i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]); + } + + if let Some(slaves_list) = active_slaves { + // 情况 1: 有从节点,进行网络传输,并本地构造静音声道回放 + + // 1. 检查各声道是否有从节点需要 + let needs_left = slaves_list.iter().any(|s| s.role == ChannelRole::Left); + let needs_right = slaves_list.iter().any(|s| s.role == ChannelRole::Right); + + // 2. 编码需要的声道 + let mut left_bytes = None; + let mut right_bytes = None; + + if needs_left { + let len = left_encoder.encode(&left_pcm, &mut opus_out)?; + let packet = AudioPacket { + seq, + timestamp: target_ts, + data: opus_out[..len].to_vec(), + }; + left_bytes = Some(postcard::to_allocvec(&packet)?); + } + + if needs_right { + let len = right_encoder.encode(&right_pcm, &mut opus_out)?; + let packet = AudioPacket { + seq, + timestamp: target_ts, + data: opus_out[..len].to_vec(), + }; + right_bytes = Some(postcard::to_allocvec(&packet)?); + } + + // 3. 发送给对应的从节点 + for slave in &slaves_list { + let bytes = match slave.role { + ChannelRole::Left => left_bytes.as_ref(), + ChannelRole::Right => right_bytes.as_ref(), + }; + if let Some(b) = bytes { + let _ = audio_socket.send_to(b, slave.udp_addr).await; + } + } + + // 4. 将非本节点的声道置为静音 + for i in 0..config.frame_size { + match master_role { + ChannelRole::Left => { + pcm_out[i * 2] = left_pcm[i]; + pcm_out[i * 2 + 1] = 0; + } + ChannelRole::Right => { + pcm_out[i * 2] = 0; + pcm_out[i * 2 + 1] = right_pcm[i]; + } + } + } + + // 5. 等待播放 + let now = now_us(); + if now < target_ts { + let wait = target_ts - now; + if wait > 1000 { + tokio::time::sleep(Duration::from_micros(wait as u64)).await; + } else { + // 小于 1ms,直接播放,让播放时机稍微早一点点 + } + } + } else { + // 情况 2: 没有从节点,本地立体声播放 + for i in 0..config.frame_size { + pcm_out[i * 2] = left_pcm[i]; + pcm_out[i * 2 + 1] = right_pcm[i]; + } + } + + // 统一写入播放器 (始终是立体声) + if let Err(_) = player.write(&pcm_out) { + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; + } + } + + seq += 1; + } + + // 重置流计时 + stream_start_ts = 0; + } + + Ok::<(), anyhow::Error>(()) + }; + + tokio::select! { + res = audio_loop => { + if let Err(e) = res { + eprintln!("❌ 音频循环错误: {:?}", e); + } + }, + _ = shutdown_signal() => { + // 设置退出标志,通知音频循环停止 + shutdown_flag.store(true, Ordering::Relaxed); + }, + } + + // 显式清理 + println!("👋 正在退出..."); + AlsaRedirector::cleanup(); + + // 强制退出 + std::process::exit(0); +} + +/// 监听系统退出信号 (SIGINT, SIGTERM, SIGQUIT) +async fn shutdown_signal() { + let mut sigint = signal(SignalKind::interrupt()).expect("无法注册 SIGINT 处理器"); + let mut sigterm = signal(SignalKind::terminate()).expect("无法注册 SIGTERM 处理器"); + let mut sigquit = signal(SignalKind::quit()).expect("无法注册 SIGQUIT 处理器"); + + tokio::select! { + _ = sigint.recv() => {}, + _ = sigterm.recv() => {}, + _ = sigquit.recv() => {}, + } +} + +/// 处理主节点与从节点的会话 +async fn handle_master_session( + mut control: ControlConnection, + audio_socket: Arc, + slaves: Arc>>, + client_tcp_addr: String, +) -> Result<()> { + let mut buf = [0u8; 1024]; + + // 握手 + let pkt = control.recv_packet(&mut buf).await?; + let slave_role = match pkt { + ControlPacket::ClientIdentify { role } => role, + _ => return Err(anyhow!("无效的握手协议")), + }; + + let xiao = ControlPacket::ServerHello { + udp_port: audio_socket.local_addr()?.port(), + }; + control.send_packet(&xiao).await?; + + // 等待 UDP 打洞/确认 + let mut buf = [0u8; 128]; + let (_, client_udp_addr) = audio_socket.recv_from(&mut buf).await?; + + println!( + "✅ 从节点已连接: {} {}", + client_tcp_addr, + slave_role.to_string(), + ); + + // 添加到从节点列表 + let session = SlaveSession { + udp_addr: client_udp_addr, + role: slave_role, + }; + { + let mut s = slaves.lock().await; + s.push(session.clone()); + } + + // 分离 TCP 读写,处理控制消息和心跳 + let (mut tcp_rx, mut tcp_tx) = control.split(); + + let mut buf = [0u8; 1024]; + loop { + match tcp_rx.read(&mut buf).await { + Ok(0) | Err(_) => { + break; + } + Ok(n) => { + if let Ok(ControlPacket::Ping { client_ts, seq }) = postcard::from_bytes(&buf[..n]) + { + let pong = ControlPacket::Pong { + client_ts, + server_ts: now_us(), + seq, + }; + if tcp_tx + .write_all(&postcard::to_allocvec(&pong).unwrap()) + .await + .is_err() + { + break; + } + } + } + } + } + + println!( + "❌ 从节点已断开: {} {}", + client_tcp_addr, + slave_role.to_string(), + ); + + // 从列表中移除 + { + let mut s = slaves.lock().await; + s.retain(|x| x.udp_addr != client_udp_addr); + } + + Ok(()) +} diff --git a/packages/client-v2/src/app/mod.rs b/packages/client-v2/src/app/mod.rs new file mode 100644 index 0000000..8af1333 --- /dev/null +++ b/packages/client-v2/src/app/mod.rs @@ -0,0 +1,3 @@ +pub mod entry; +pub mod master; +pub mod slave; diff --git a/packages/client-v2/src/app/slave.rs b/packages/client-v2/src/app/slave.rs new file mode 100644 index 0000000..7ab3d6e --- /dev/null +++ b/packages/client-v2/src/app/slave.rs @@ -0,0 +1,199 @@ +#![cfg(target_os = "linux")] + +use crate::audio::codec::OpusCodec; +use crate::audio::config::AudioConfig; +use crate::audio::player::AudioPlayer; +use crate::net::discovery::Discovery; +use crate::net::network::SlaveNetwork; +use crate::net::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::utils::sync::{ClockSync, now_us}; +use anyhow::{Result, anyhow}; +use std::sync::Arc; +use std::time::Duration; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::sync::{Mutex, mpsc}; + +/// 运行从节点模式 +pub async fn run_slave(role: ChannelRole) -> Result<()> { + loop { + match handle_connection(role.clone()).await { + Err(e) => { + eprintln!("❌ {:?}", e); + tokio::time::sleep(Duration::from_secs(3)).await; + } + Ok(_) => {} + } + } +} + +async fn handle_connection(role: ChannelRole) -> Result<()> { + // 1. 发现主节点 + println!("🔍 正在扫描主节点..."); + let (master_ip, master_tcp_port) = Discovery::discover_master().await?; + let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port); + + // 2. 建立 TCP 连接 + println!("🔥 发现主节点: {}", master_tcp_addr); + let network = SlaveNetwork::connect(master_tcp_addr.parse()?).await?; + let (mut control, audio) = network.split(); + + // 3. 身份认证 + control + .send_packet(&ControlPacket::ClientIdentify { role: role.clone() }) + .await?; + + let mut buf = [0u8; 1024]; + let pkt = control.recv_packet(&mut buf).await?; + let server_udp_port = match pkt { + ControlPacket::ServerHello { udp_port } => udp_port, + _ => return Err(anyhow!("身份认证应答异常")), + }; + + // 4. UDP 打洞 + audio + .punch(format!("{}:{}", master_ip, server_udp_port).parse()?) + .await?; + + // 5. 初始化音频与同步组件 + let config = AudioConfig { + channels: 1, + ..AudioConfig::music() + }; + let player = AudioPlayer::new(&config)?; + let mut codec = OpusCodec::new(&config)?; + let clock = Arc::new(Mutex::new(ClockSync::new(100))); + + // 用于通知主循环 TCP 已断开的消息通道 + let (disconnect_tx, mut disconnect_rx) = mpsc::channel::<()>(1); + + // 6. 分离 TCP 读写 + let (mut tcp_rx, mut tcp_tx) = control.split(); + let clock_updater = clock.clone(); + let d_tx_ping = disconnect_tx.clone(); + let d_tx_pong = disconnect_tx.clone(); + + // 定时发送 Ping (心跳 & 时间同步) + let _sync_handle = tokio::spawn(async move { + let mut seq = 0; + loop { + let t1 = now_us(); + let msg = ControlPacket::Ping { client_ts: t1, seq }; + let data = postcard::to_allocvec(&msg).unwrap(); + if tcp_tx.write_all(&data).await.is_err() { + let _ = d_tx_ping.send(()).await; // 通知主线程 TCP 失败 + break; + } + tokio::time::sleep(Duration::from_millis(200)).await; + seq += 1; + } + }); + + // 接收 Pong + tokio::spawn(async move { + let mut buf = [0u8; 1024]; + loop { + match tcp_rx.read(&mut buf).await { + Ok(n) if n > 0 => { + if let Ok(ControlPacket::Pong { + client_ts, + server_ts, + .. + }) = postcard::from_bytes(&buf[..n]) + { + let t4 = now_us(); + clock_updater.lock().await.update(client_ts, server_ts, t4); + } + } + _ => { + let _ = d_tx_pong.send(()).await; // TCP 断开 + break; + } + } + } + }); + + // 7. 接收音频数据包 (UDP) + let (audio_tx, mut audio_rx) = mpsc::channel(100); + let audio_socket = audio.clone_inner(); + tokio::spawn(async move { + let mut buf = [0u8; 2048]; + loop { + if let Ok((len, _)) = audio_socket.recv_from(&mut buf).await { + if let Ok(packet) = postcard::from_bytes::(&buf[..len]) { + if audio_tx.send(packet).await.is_err() { + break; + } + } + } + } + }); + + // 8. 播放提示 + println!("✅ 主节点已连接,音频串流中..."); + let role_str = role.to_string(); + tokio::spawn(async move { + let _ = tokio::process::Command::new("sh") + .arg("-c") + .arg(format!( + "/usr/sbin/tts_play.sh \"主节点已连接,{}\" >/dev/null 2>&1", + role_str + )) + .status() + .await; + }); + + // 9. 播放主循环 + let mut pcm_buf = vec![0i16; config.frame_size]; + let mut last_seq: Option = None; + + loop { + // 检查 TCP 是否已断开 + if let Ok(_) = disconnect_rx.try_recv() { + return Err(anyhow!("主节点已断开: {}", master_tcp_addr)); + } + + // 接收数据包 + if let Ok(pkt) = audio_rx.try_recv() { + let now = now_us(); + let current_server_time = clock.lock().await.to_server_time(now); + + // 检查包是否迟到(目标时间已过) + if current_server_time > pkt.timestamp { + let late_ms = (current_server_time - pkt.timestamp) / 1000; + if late_ms > 50 { + // 迟到超过 50ms,直接丢弃 + continue; + } + // 轻微迟到(<50ms),尝试播放 + } + + last_seq = Some(pkt.seq); + + // 精确等待到播放时间 + loop { + let now = now_us(); + let current_server_time_precise = clock.lock().await.to_server_time(now); + + if current_server_time_precise >= pkt.timestamp { + break; + } + + let wait_us = (pkt.timestamp - current_server_time_precise) as u64; + + if wait_us > 1000 { + tokio::time::sleep(Duration::from_micros(wait_us as u64)).await; + } else { + // 小于 1ms,直接播放,让播放时机稍微早一点点 + break; + } + } + + // 解码并播放 + if let Ok(len) = codec.decode(&pkt.data, &mut pcm_buf) { + let _ = player.write(&pcm_buf[..len]); + } + } else { + tokio::time::sleep(Duration::from_micros(100)).await; + } + } +} diff --git a/packages/client-v2/src/audio/codec.rs b/packages/client-v2/src/audio/codec.rs new file mode 100644 index 0000000..641d127 --- /dev/null +++ b/packages/client-v2/src/audio/codec.rs @@ -0,0 +1,70 @@ +use crate::audio::config::{AudioConfig, AudioScene}; +use anyhow::{Context, Result}; +use opus::{Application, Bitrate, Channels, Decoder, Encoder}; + +pub struct OpusCodec { + encoder: Encoder, + decoder: Decoder, +} + +impl OpusCodec { + pub fn new(config: &AudioConfig) -> Result { + let channels = match config.channels { + 1 => Channels::Mono, + 2 => Channels::Stereo, + _ => return Err(anyhow::anyhow!("Invalid channels: {}", config.channels)), + }; + let mode = match config.audio_scene { + AudioScene::Music => Application::Audio, + AudioScene::Voice => Application::Voip, + }; + let bitrate = match config.bitrate { + -1 => Bitrate::Max, + 0 => Bitrate::Auto, + _ => Bitrate::Bits(config.bitrate), + }; + + let mut encoder = Encoder::new(config.sample_rate, channels, mode) + .context("Failed to create Opus encoder")?; + + encoder.set_bitrate(bitrate)?; + if config.vbr { + encoder.set_vbr(true)?; + } + if config.fec { + encoder.set_inband_fec(true)?; // 内联前向纠错 + encoder.set_packet_loss_perc(20)?; // 预期丢包率20% + } + + let decoder = + Decoder::new(config.sample_rate, channels).context("Failed to create Opus decoder")?; + + Ok(Self { encoder, decoder }) + } + + pub fn encode(&mut self, pcm: &[i16], out: &mut [u8]) -> Result { + self.encoder + .encode(pcm, out) + .context("Opus encoding failed") + } + + pub fn decode(&mut self, opus: &[u8], out: &mut [i16]) -> Result { + self.decoder + .decode(opus, out, false) + .context("Opus decoding failed") + } + + /// 前向纠错(FEC) + pub fn decode_fec(&mut self, opus: &[u8], out: &mut [i16]) -> Result { + self.decoder + .decode(opus, out, true) + .context("Opus FEC decoding failed") + } + + /// 丢包补偿(PLC) + pub fn decode_loss(&mut self, out: &mut [i16]) -> Result { + self.decoder + .decode(&[], out, false) + .context("Opus PLC (decode_loss) failed") + } +} diff --git a/packages/client-v2/src/audio/config.rs b/packages/client-v2/src/audio/config.rs new file mode 100644 index 0000000..1c982ef --- /dev/null +++ b/packages/client-v2/src/audio/config.rs @@ -0,0 +1,61 @@ +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub enum AudioScene { + Music, + Voice, +} + +#[derive(Debug, Clone)] +pub struct AudioConfig { + // ALSA 设备参数,用于录音和播放 + pub capture_device: String, + pub playback_device: String, + pub sample_rate: u32, + pub channels: u16, + pub frame_size: usize, // 帧大小,单位为采样点 + + // Opus 编解码参数,用于音频传输 + pub audio_scene: AudioScene, + pub bitrate: i32, + pub vbr: bool, // 是否启用 VBR(动态比特率) + pub fec: bool, // 是否启用 FEC(内联前向纠错) +} + +impl AudioConfig { + pub fn music() -> Self { + Self { + audio_scene: AudioScene::Music, + sample_rate: 48_000, // 48kHz + channels: 2, + frame_size: 960, // 20ms at 48kHz + bitrate: 320_000, // 320 kbps + ..Default::default() + } + } + + pub fn voice() -> Self { + Self { + audio_scene: AudioScene::Voice, + sample_rate: 16_000, // 16kHz + channels: 1, + frame_size: 320, // 20ms at 16kHz + bitrate: 32_000, // 32 kbps + ..Default::default() + } + } +} + +impl Default for AudioConfig { + fn default() -> Self { + Self { + audio_scene: AudioScene::Voice, + capture_device: "plug:Capture".to_string(), + playback_device: "default".to_string(), + sample_rate: 16_000, + channels: 1, + frame_size: 320, // 20ms at 16kHz + bitrate: 32_000, + vbr: false, + fec: false, + } + } +} diff --git a/packages/client-v2/src/audio/mod.rs b/packages/client-v2/src/audio/mod.rs new file mode 100644 index 0000000..25d9907 --- /dev/null +++ b/packages/client-v2/src/audio/mod.rs @@ -0,0 +1,4 @@ +pub mod codec; +pub mod config; +pub mod player; +pub mod recorder; diff --git a/packages/client-v2/src/audio/player.rs b/packages/client-v2/src/audio/player.rs new file mode 100644 index 0000000..fe38f5f --- /dev/null +++ b/packages/client-v2/src/audio/player.rs @@ -0,0 +1,71 @@ +#![cfg(target_os = "linux")] + +use crate::audio::config::AudioConfig; +use alsa::Direction; +use alsa::pcm::{Access, Format, HwParams, PCM}; +use anyhow::{Context, Result}; + +pub struct AudioPlayer { + pcm: PCM, +} + +impl AudioPlayer { + pub fn new(config: &AudioConfig) -> Result { + let pcm = PCM::new(&config.playback_device, Direction::Playback, false) + .context("Failed to open playback PCM device")?; + + setup_pcm(&pcm, config.sample_rate, config.channels)?; + Ok(Self { pcm }) + } + + pub fn write(&self, buffer: &[i16]) -> Result { + let res = self.pcm.io_i16()?.writei(buffer); + + match res { + Ok(written) => Ok(written), + Err(e) => { + // Buffer Underrun,即播放缓冲区的数据被耗尽,导致音频流中断 + if e.errno() == 32 { + // 恢复音频流状态 + self.pcm.prepare()?; + // 重新获取 IO 对象并尝试写入数据 + self.pcm + .io_i16()? + .writei(buffer) + .context("Failed to write to playback device after recovery") + } else { + Err(e).context("Failed to write to playback device") + } + } + } + } + + pub fn prepare(&self) -> Result<()> { + self.pcm.prepare().context("Failed to prepare PCM") + } +} + +fn setup_pcm(pcm: &PCM, sample_rate: u32, channels: u16) -> Result<()> { + let hwp = HwParams::any(pcm).context("Failed to get HwParams")?; + hwp.set_access(Access::RWInterleaved)?; + hwp.set_format(Format::s16())?; + hwp.set_rate(sample_rate, alsa::ValueOr::Nearest)?; + hwp.set_channels(channels as u32)?; + + // 设置较大的缓冲区以减少由于调度抖动和设备重初始化导致的断音/卡顿 + // 使用 100ms 缓冲区,既能防止 underrun,又不会引入过大延迟 + let buffer_size = (sample_rate as f64 * 0.1) as u32; // 100ms 缓冲 + let period_size = buffer_size / 4; // 25ms 周期 + hwp.set_buffer_size_near(buffer_size as alsa::pcm::Frames)?; + hwp.set_period_size_near(period_size as alsa::pcm::Frames, alsa::ValueOr::Nearest)?; + + pcm.hw_params(&hwp).context("Failed to set HwParams")?; + + let swp = pcm.sw_params_current()?; + // 设置 start_threshold,当缓冲区有 1 个 period 数据时就开始播放 + // 这样可以快速启动,同时保持足够的缓冲余量 + swp.set_start_threshold(period_size as alsa::pcm::Frames)?; + pcm.sw_params(&swp)?; + pcm.prepare()?; + Ok(()) +} diff --git a/packages/client-v2/src/audio/recorder.rs b/packages/client-v2/src/audio/recorder.rs new file mode 100644 index 0000000..952e73b --- /dev/null +++ b/packages/client-v2/src/audio/recorder.rs @@ -0,0 +1,41 @@ +#![cfg(target_os = "linux")] + +use crate::audio::config::AudioConfig; +use alsa::Direction; +use alsa::pcm::{Access, Format, HwParams, PCM}; +use anyhow::{Context, Result}; + +pub struct AudioRecorder { + pcm: PCM, +} + +impl AudioRecorder { + pub fn new(config: &AudioConfig) -> Result { + let pcm = PCM::new(&config.capture_device, Direction::Capture, false) + .context("Failed to open capture PCM device")?; + + setup_pcm(&pcm, config.sample_rate, config.channels)?; + Ok(Self { pcm }) + } + + pub fn read(&self, buffer: &mut [i16]) -> Result { + self.pcm + .io_i16()? + .readi(buffer) + .context("Failed to read from capture device") + } +} + +fn setup_pcm(pcm: &PCM, sample_rate: u32, channels: u16) -> Result<()> { + let hwp = HwParams::any(pcm).context("Failed to get HwParams")?; + hwp.set_access(Access::RWInterleaved)?; + hwp.set_format(Format::s16())?; + hwp.set_rate(sample_rate, alsa::ValueOr::Nearest)?; + hwp.set_channels(channels as u32)?; + pcm.hw_params(&hwp).context("Failed to set HwParams")?; + + let swp = pcm.sw_params_current()?; + pcm.sw_params(&swp)?; + pcm.prepare()?; + Ok(()) +} diff --git a/packages/client-v2/src/bin/client.rs b/packages/client-v2/src/bin/client.rs new file mode 100644 index 0000000..37d36f8 --- /dev/null +++ b/packages/client-v2/src/bin/client.rs @@ -0,0 +1,11 @@ +use anyhow::Result; + +#[tokio::main] +async fn main() -> Result<()> { + #[cfg(target_os = "linux")] + { + xiao::app::entry::run_xiao().await.unwrap(); + } + println!("Only support Linux"); + Ok(()) +} diff --git a/packages/client-v2/src/lib.rs b/packages/client-v2/src/lib.rs new file mode 100644 index 0000000..bc714b2 --- /dev/null +++ b/packages/client-v2/src/lib.rs @@ -0,0 +1,4 @@ +pub mod app; +pub mod audio; +pub mod net; +pub mod utils; diff --git a/packages/client-v2/src/net/discovery.rs b/packages/client-v2/src/net/discovery.rs new file mode 100644 index 0000000..e4f069f --- /dev/null +++ b/packages/client-v2/src/net/discovery.rs @@ -0,0 +1,45 @@ +use crate::net::protocol::ControlPacket; +use anyhow::Result; +use std::net::{IpAddr, SocketAddr}; +use std::time::Duration; +use tokio::net::UdpSocket; + +pub const DISCOVERY_PORT: u16 = 53530; + +/// 服务发现模块,用于主从节点的自动发现 +pub struct Discovery; + +impl Discovery { + /// 主节点:启动广播,告知从节点自己的 TCP 端口 + pub async fn start_broadcast(tcp_port: u16) -> Result<()> { + let socket = UdpSocket::bind("0.0.0.0:0").await?; + socket.set_broadcast(true)?; + + let target: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT).parse()?; + let msg = postcard::to_allocvec(&ControlPacket::ServerHello { udp_port: tcp_port })?; + + tokio::spawn(async move { + loop { + let _ = socket.send_to(&msg, target).await; + tokio::time::sleep(Duration::from_secs(1)).await; + } + }); + + Ok(()) + } + + /// 从节点:监听广播,发现主节点的 IP 和 TCP 端口 + pub async fn discover_master() -> Result<(IpAddr, u16)> { + let socket = UdpSocket::bind(format!("0.0.0.0:{}", DISCOVERY_PORT)).await?; + let mut buf = [0u8; 1024]; + + loop { + let (len, addr) = socket.recv_from(&mut buf).await?; + if let Ok(ControlPacket::ServerHello { udp_port }) = + postcard::from_bytes::(&buf[..len]) + { + return Ok((addr.ip(), udp_port)); + } + } + } +} diff --git a/packages/client-v2/src/net/mod.rs b/packages/client-v2/src/net/mod.rs new file mode 100644 index 0000000..3bba08b --- /dev/null +++ b/packages/client-v2/src/net/mod.rs @@ -0,0 +1,3 @@ +pub mod discovery; +pub mod network; +pub mod protocol; diff --git a/packages/client-v2/src/net/network.rs b/packages/client-v2/src/net/network.rs new file mode 100644 index 0000000..69dd573 --- /dev/null +++ b/packages/client-v2/src/net/network.rs @@ -0,0 +1,126 @@ +use crate::net::protocol::{AudioPacket, ControlPacket}; +use anyhow::{Context, Result}; +use std::net::SocketAddr; +use std::sync::Arc; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream, UdpSocket}; + +/// UDP 音频传输 +pub struct AudioSocket { + socket: Arc, +} + +impl AudioSocket { + pub async fn bind() -> Result { + let socket = UdpSocket::bind("0.0.0.0:0").await?; + Ok(Self { + socket: Arc::new(socket), + }) + } + + pub fn local_port(&self) -> Result { + Ok(self.socket.local_addr()?.port()) + } + + pub async fn send_packet(&self, packet: &AudioPacket, target: SocketAddr) -> Result<()> { + let bytes = postcard::to_allocvec(packet)?; + self.socket.send_to(&bytes, target).await?; + Ok(()) + } + + pub async fn recv_packet(&self, buf: &mut [u8]) -> Result<(AudioPacket, SocketAddr)> { + let (len, addr) = self.socket.recv_from(buf).await?; + let packet = postcard::from_bytes(&buf[..len])?; + Ok((packet, addr)) + } + + pub async fn punch(&self, target: SocketAddr) -> Result<()> { + self.socket.send_to(&[0u8; 1], target).await?; + Ok(()) + } + + pub fn clone_inner(&self) -> Arc { + self.socket.clone() + } +} + +/// TCP 控制连接 +pub struct ControlConnection { + stream: TcpStream, +} + +impl ControlConnection { + pub fn new(stream: TcpStream) -> Self { + Self { stream } + } + + pub async fn send_packet(&mut self, packet: &ControlPacket) -> Result<()> { + let bytes = postcard::to_allocvec(packet)?; + self.stream.write_all(&bytes).await?; + Ok(()) + } + + pub async fn recv_packet(&mut self, buf: &mut [u8]) -> Result { + let len = self.stream.read(buf).await?; + if len == 0 { + return Err(anyhow::anyhow!("连接已关闭")); + } + let packet = postcard::from_bytes(&buf[..len])?; + Ok(packet) + } + + pub fn split( + self, + ) -> ( + tokio::net::tcp::OwnedReadHalf, + tokio::net::tcp::OwnedWriteHalf, + ) { + self.stream.into_split() + } +} + +/// 主节点网络管理器 +pub struct MasterNetwork { + listener: TcpListener, + audio: AudioSocket, +} + +impl MasterNetwork { + pub async fn setup(port: u16) -> Result { + let listener = TcpListener::bind(format!("0.0.0.0:{}", port)).await?; + let audio = AudioSocket::bind().await?; + Ok(Self { listener, audio }) + } + + pub async fn accept(&self) -> Result<(ControlConnection, SocketAddr)> { + let (stream, addr) = self.listener.accept().await?; + Ok((ControlConnection::new(stream), addr)) + } + + pub fn audio_socket(&self) -> &AudioSocket { + &self.audio + } +} + +/// 从节点网络管理器 +pub struct SlaveNetwork { + control: ControlConnection, + audio: AudioSocket, +} + +impl SlaveNetwork { + pub async fn connect(master_addr: SocketAddr) -> Result { + let stream = TcpStream::connect(master_addr) + .await + .context(format!("无法连接到主节点 TCP 地址: {}", master_addr))?; + let audio = AudioSocket::bind().await?; + Ok(Self { + control: ControlConnection::new(stream), + audio, + }) + } + + pub fn split(self) -> (ControlConnection, AudioSocket) { + (self.control, self.audio) + } +} diff --git a/packages/client-v2/src/net/protocol.rs b/packages/client-v2/src/net/protocol.rs new file mode 100644 index 0000000..c699019 --- /dev/null +++ b/packages/client-v2/src/net/protocol.rs @@ -0,0 +1,45 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, Copy)] +pub enum ChannelRole { + Left, + Right, +} + +impl ChannelRole { + pub fn to_string(&self) -> String { + match self { + ChannelRole::Left => "左声道".to_string(), + ChannelRole::Right => "右声道".to_string(), + } + } +} + +#[derive(Serialize, Deserialize, Debug, Clone)] +pub enum ControlPacket { + // 发现协议 + ServerHello { + udp_port: u16, // UDP 音频流端口 + }, + // 握手协议 + ClientIdentify { + role: ChannelRole, + }, + // 时间同步 (持续进行) + Ping { + client_ts: u128, + seq: u32, + }, + Pong { + client_ts: u128, + server_ts: u128, + seq: u32, + }, +} + +#[derive(Serialize, Deserialize, Debug, Clone)] +pub struct AudioPacket { + pub seq: u32, // 序列号,用于丢包检测 + pub timestamp: u128, // 目标播放时间 (主节点时间) + pub data: Vec, // Opus 编码数据 +} diff --git a/packages/client-v2/src/utils/alsa.rs b/packages/client-v2/src/utils/alsa.rs new file mode 100644 index 0000000..42a876c --- /dev/null +++ b/packages/client-v2/src/utils/alsa.rs @@ -0,0 +1,85 @@ +#![cfg(target_os = "linux")] + +use anyhow::{Context, Result}; +use std::fs; +use std::process::Command; + +const FIFO_PATH: &str = "/tmp/xiao_out.fifo"; +const REAL_ASOUND_CONF: &str = "/etc/asound.conf"; +const TEMP_ASOUND_CONF: &str = "/tmp/asound.xiao.conf"; + +/// ALSA 音频重定向器,用于拦截系统音频输出到 FIFO 管道 +pub struct AlsaRedirector; + +impl AlsaRedirector { + pub fn new() -> Result { + Self::cleanup(); // 确保环境干净 + + let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default(); + + if !original_conf.contains("pcm.original_default") { + // 重命名原有的 default 逻辑,插入拦截器 + let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default"); + new_conf.push_str(&format!( + "\npcm.!default {{ type plug slave {{ pcm \"xiao_interceptor\" format S16_LE rate 48000 channels 2 }} }}\n\ + pcm.xiao_interceptor {{ type file slave.pcm \"null\" file \"{}\" format \"raw\" }}\n", + FIFO_PATH + )); + + fs::write(TEMP_ASOUND_CONF, new_conf)?; + + // 挂载覆盖 /etc/asound.conf + let status = Command::new("mount") + .arg("--bind") + .arg(TEMP_ASOUND_CONF) + .arg(REAL_ASOUND_CONF) + .status() + .context("执行 mount 命令失败")?; + + if !status.success() { + return Err(anyhow::anyhow!("挂载 asound.conf 失败")); + } + + Self::restart_applications(); + } + + // 创建 FIFO 管道 + let _ = Command::new("mkfifo").arg(FIFO_PATH).status(); + let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status(); + + Ok(Self) + } + + pub fn cleanup() { + let _ = Command::new("sh") + .arg("-c") + .arg(format!("umount -l {} >/dev/null 2>&1", REAL_ASOUND_CONF)) + .status(); + let _ = fs::remove_file(TEMP_ASOUND_CONF); + let _ = fs::remove_file(FIFO_PATH); + Self::restart_applications(); + } + + pub fn fifo_path() -> &'static str { + FIFO_PATH + } + + pub fn restart_applications() { + // 重启媒体播放器 + let _ = Command::new("sh") + .arg("-c") + .arg("/etc/init.d/mediaplayer restart >/dev/null 2>&1") + .status(); + // 重启蓝牙 + let _ = Command::new("sh") + .arg("-c") + .arg("/etc/init.d/bluetooth restart >/dev/null 2>&1") + .status(); + } +} + +impl Drop for AlsaRedirector { + fn drop(&mut self) { + Self::cleanup(); + } +} diff --git a/packages/client-v2/src/utils/mod.rs b/packages/client-v2/src/utils/mod.rs new file mode 100644 index 0000000..8eaa5c4 --- /dev/null +++ b/packages/client-v2/src/utils/mod.rs @@ -0,0 +1,2 @@ +pub mod alsa; +pub mod sync; diff --git a/packages/client-v2/src/utils/sync.rs b/packages/client-v2/src/utils/sync.rs new file mode 100644 index 0000000..354ab8d --- /dev/null +++ b/packages/client-v2/src/utils/sync.rs @@ -0,0 +1,224 @@ +use std::collections::VecDeque; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// 获取当前微秒级时间戳 +pub fn now_us() -> u128 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("时间倒流") + .as_micros() +} + +/// 时钟同步管理器,用于计算主从节点间的时钟偏移 +/// 采用改进的 NTP 算法 + Kalman 滤波思想 +pub struct ClockSync { + /// 偏移量样本窗口 + offsets: VecDeque, + /// 当前估计的时钟偏移 (server_time - client_time) + pub current_offset: i128, + /// RTT 样本窗口 + rtts: VecDeque, + /// 当前估计的最小 RTT + min_rtt: i128, + /// 窗口大小 + window_size: usize, + /// 时钟漂移率 (ppm: parts per million) + /// 正值表示从节点时钟比主节点快 + drift_rate: f64, + /// 上次更新时间 + last_update_time: u128, + /// 漂移率估计窗口 + drift_samples: VecDeque, +} + +#[derive(Clone, Copy)] +struct OffsetSample { + offset: i128, + rtt: i128, + timestamp: u128, +} + +#[derive(Clone, Copy)] +struct DriftSample { + offset: i128, + timestamp: u128, +} + +impl ClockSync { + pub fn new(window_size: usize) -> Self { + Self { + offsets: VecDeque::with_capacity(window_size), + current_offset: 0, + rtts: VecDeque::with_capacity(window_size), + min_rtt: i128::MAX, + window_size, + drift_rate: 0.0, + last_update_time: now_us(), + drift_samples: VecDeque::with_capacity(60), // 保留 60 秒的样本 + } + } + + /// 更新时钟偏移估计 (NTP 算法) + /// + /// NTP 时间戳标记: + /// t1 = client_send_ts : 客户端发送 Ping 的时间 + /// t2 = server_ts : 服务器接收 Ping 的时间 + /// t3 = server_ts : 服务器发送 Pong 的时间 (假设处理时间忽略不计) + /// t4 = client_recv_ts : 客户端接收 Pong 的时间 + /// + /// RTT = (t4 - t1) - (t3 - t2) = (t4 - t1) (因为 t3 = t2) + /// Offset = ((t2 - t1) + (t3 - t4)) / 2 = ((t2 - t1) + (t2 - t4)) / 2 + /// = t2 - (t1 + t4) / 2 + pub fn update(&mut self, client_send_ts: u128, server_ts: u128, client_recv_ts: u128) { + let t1 = client_send_ts as i128; + let t2 = server_ts as i128; + let t4 = client_recv_ts as i128; + + let rtt = t4 - t1; + + // 过滤异常 RTT (局域网内 > 100ms 视为异常) + if rtt < 0 || rtt > 100_000 { + return; + } + + // 计算时钟偏移: offset = server_time - client_time + // offset = t2 - (t1 + t4) / 2 + let offset = t2 - (t1 + t4) / 2; + + // 更新 RTT 窗口 + self.rtts.push_back(rtt); + if self.rtts.len() > self.window_size { + self.rtts.pop_front(); + } + self.min_rtt = *self.rtts.iter().min().unwrap_or(&rtt); + + // 更新偏移量窗口 + let sample = OffsetSample { + offset, + rtt, + timestamp: client_recv_ts, + }; + self.offsets.push_back(sample); + if self.offsets.len() > self.window_size { + self.offsets.pop_front(); + } + + // 偏移量估计: 使用低 RTT 样本的中位数 + // 原理: RTT 较小的样本受网络抖动影响小,时间测量更准确 + let mut low_rtt_offsets: Vec = self + .offsets + .iter() + .filter(|s| s.rtt <= self.min_rtt + 5000) // 5ms 容差 + .map(|s| s.offset) + .collect(); + + if !low_rtt_offsets.is_empty() { + low_rtt_offsets.sort_unstable(); + let new_offset = low_rtt_offsets[low_rtt_offsets.len() / 2]; + + // 漂移率估计 + self.estimate_drift(new_offset, client_recv_ts); + + // 平滑更新偏移量 (避免突变) + let alpha = 0.3; // 低通滤波系数 + self.current_offset = + (alpha * new_offset as f64 + (1.0 - alpha) * self.current_offset as f64) as i128; + } + + self.last_update_time = client_recv_ts; + } + + /// 估计时钟漂移率 + /// 时钟漂移率 = d(offset) / dt + fn estimate_drift(&mut self, offset: i128, timestamp: u128) { + self.drift_samples.push_back(DriftSample { offset, timestamp }); + if self.drift_samples.len() > 60 { + self.drift_samples.pop_front(); + } + + // 至少需要 10 秒的数据才能估计漂移 + if self.drift_samples.len() < 10 { + return; + } + + // 使用线性回归估计漂移率 + let first = self.drift_samples.front().unwrap(); + let last = self.drift_samples.back().unwrap(); + + let dt = (last.timestamp - first.timestamp) as f64; + let d_offset = (last.offset - first.offset) as f64; + + if dt > 10_000_000.0 { + // 超过 10 秒 + // drift_rate 单位: 微秒/秒 = ppm + let new_drift = d_offset / (dt / 1_000_000.0); + + // 平滑更新漂移率 + let beta = 0.1; + self.drift_rate = beta * new_drift + (1.0 - beta) * self.drift_rate; + } + } + + /// 将本地时间转换为服务器(主节点)时间 + /// 考虑时钟漂移补偿 + pub fn to_server_time(&self, client_time: u128) -> u128 { + let base_server_time = (client_time as i128 + self.current_offset) as u128; + + // 漂移补偿: 根据距离上次同步的时间,补偿时钟漂移 + let elapsed_since_update = client_time.saturating_sub(self.last_update_time) as f64; + let drift_correction = (self.drift_rate * elapsed_since_update / 1_000_000.0) as i128; + + (base_server_time as i128 + drift_correction) as u128 + } + + /// 将服务器(主节点)时间转换为本地时间 + pub fn to_client_time(&self, server_time: u128) -> u128 { + // 简化版本,不考虑漂移补偿 (播放时主要用 to_server_time) + (server_time as i128 - self.current_offset) as u128 + } + + /// 获取当前估计的 RTT (微秒) + pub fn get_rtt(&self) -> i128 { + self.min_rtt + } + + /// 获取当前时钟漂移率 (ppm) + pub fn get_drift_rate(&self) -> f64 { + self.drift_rate + } + + /// 获取同步质量评估 (0-100, 越高越好) + pub fn get_sync_quality(&self) -> u8 { + if self.offsets.is_empty() { + return 0; + } + + // 基于 RTT 稳定性和偏移量方差评估 + let rtt_variance = self.calculate_variance(&self.rtts.iter().copied().collect::>()); + let offset_variance = self.calculate_variance( + &self.offsets.iter().map(|s| s.offset).collect::>(), + ); + + // RTT 越稳定,方差越小,质量越高 + let rtt_score = ((100_000.0 - rtt_variance.min(100_000.0)) / 100_000.0 * 50.0) as u8; + let offset_score = ((50_000.0 - offset_variance.min(50_000.0)) / 50_000.0 * 50.0) as u8; + + rtt_score + offset_score + } + + fn calculate_variance(&self, samples: &[i128]) -> f64 { + if samples.is_empty() { + return 0.0; + } + let mean = samples.iter().sum::() as f64 / samples.len() as f64; + let variance = samples + .iter() + .map(|&x| { + let diff = x as f64 - mean; + diff * diff + }) + .sum::() + / samples.len() as f64; + variance.sqrt() + } +}