From b76c07ece2af64ef54591fd3267bae8068f12e1c Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 24 Jul 2026 08:22:22 +0200 Subject: [PATCH] feat(pi-natives): added LiveWebRtcPeer and deviceCheckGenerateToken bindings - Replaced puppeteer-based WebRTC with native LiveWebRtcPeer for cross-platform live audio delivery. - Added cross-platform microphone capture via miniaudio and Opus codec integration for live encoding/decoding. - Added Apple DeviceCheck attestation token generation via raw Objective-C FFI for macOS. - Updated live session model to "gpt-live-1-codex" and default voice to "sol" across protocol and controller. - Added LiveWebRtcPeer and deviceCheckGenerateToken to the public native bindings API. --- Cargo.lock | 1323 ++++++++++++++++- Cargo.toml | 8 + crates/pi-natives/Cargo.toml | 5 + crates/pi-natives/src/audio.rs | 427 ++++++ crates/pi-natives/src/devicecheck.rs | 348 +++++ crates/pi-natives/src/lib.rs | 3 + crates/pi-natives/src/live.rs | 769 ++++++++++ packages/coding-agent/CHANGELOG.md | 2 + packages/coding-agent/src/live/attestation.ts | 91 ++ .../coding-agent/src/live/browser-runtime.txt | 21 +- packages/coding-agent/src/live/controller.ts | 4 +- .../coding-agent/src/live/protocol.test.ts | 6 +- packages/coding-agent/src/live/protocol.ts | 2 +- packages/coding-agent/src/live/transport.ts | 215 +-- packages/natives/CHANGELOG.md | 2 + packages/natives/native/index.d.ts | 60 + packages/natives/native/index.js | 4 + packages/natives/test/devicecheck.test.ts | 32 + 18 files changed, 3152 insertions(+), 170 deletions(-) create mode 100644 crates/pi-natives/src/audio.rs create mode 100644 crates/pi-natives/src/devicecheck.rs create mode 100644 crates/pi-natives/src/live.rs create mode 100644 packages/coding-agent/src/live/attestation.ts create mode 100644 packages/natives/test/devicecheck.test.ts diff --git a/Cargo.lock b/Cargo.lock index 15b1a3f98..d01a8b3d6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8,6 +8,41 @@ version = "2.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa" +[[package]] +name = "aead" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d122413f284cf2d62fb1b7db97e02edb8cda96d769b16e443a4f6195e35662b0" +dependencies = [ + "crypto-common 0.1.7", + "generic-array", +] + +[[package]] +name = "aes" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b169f7a6d4742236a0a00c541b845991d0ac43e546831af1249753ab4c3aa3a0" +dependencies = [ + "cfg-if", + "cipher", + "cpufeatures 0.2.17", +] + +[[package]] +name = "aes-gcm" +version = "0.10.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "831010a0f742e1209b3bcea8fab6a8e149051ba6099432c8cb2cc117dec3ead1" +dependencies = [ + "aead", + "aes", + "cipher", + "ctr", + "ghash", + "subtle", +] + [[package]] name = "ahash" version = "0.8.12" @@ -142,6 +177,15 @@ dependencies = [ "x11rb", ] +[[package]] +name = "arc-swap" +version = "1.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c049c0be4daef0b145cb3555416b3b8ef5b7888a38aea1a3a155801fe7b0810b" +dependencies = [ + "rustversion", +] + [[package]] name = "archery" version = "1.2.2" @@ -174,6 +218,45 @@ version = "0.7.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3fb67a6e08acf24fdeccbac2cb6ac4305825bd1f117462e0e6f2f193345ad56" +[[package]] +name = "asn1-rs" +version = "0.6.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5493c3bedbacf7fd7382c6346bbd66687d12bbaad3a89a2d2c303ee6cf20b048" +dependencies = [ + "asn1-rs-derive", + "asn1-rs-impl", + "displaydoc", + "nom 7.1.3", + "num-traits", + "rusticata-macros", + "thiserror 1.0.69", + "time", +] + +[[package]] +name = "asn1-rs-derive" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "965c2d33e53cb6b267e148a4cb0760bc01f4904c1cd4bb4002a085bb016d1490" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", + "synstructure", +] + +[[package]] +name = "asn1-rs-impl" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b18050c2cd6fe86c3a76584ef5e0baf286d038cda203eb6223df2cc413565f7" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "ast-grep-core" version = "0.39.9" @@ -332,12 +415,29 @@ version = "1.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" +[[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 = "autocfg" version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" +[[package]] +name = "base16ct" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4c7f02d4ea65f2c1853089ffd8d2787bdbc63de2f0d29dedbcf8ccdfa0ccd4cf" + [[package]] name = "base64" version = "0.22.1" @@ -354,6 +454,12 @@ dependencies = [ "vsimd", ] +[[package]] +name = "base64ct" +version = "1.8.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06" + [[package]] name = "bigdecimal" version = "0.4.10" @@ -385,6 +491,24 @@ dependencies = [ "serde", ] +[[package]] +name = "bindgen" +version = "0.71.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5f58bf3d7db68cfbac37cfc485a8d711e87e064c3d0fe0435b92f7a407f9d6b3" +dependencies = [ + "bitflags 2.13.1", + "cexpr", + "clang-sys", + "itertools 0.13.0", + "proc-macro2", + "quote", + "regex", + "rustc-hash 2.1.3", + "shlex 1.3.0", + "syn 2.0.119", +] + [[package]] name = "bindgen" version = "0.72.1" @@ -486,6 +610,15 @@ dependencies = [ "hybrid-array", ] +[[package]] +name = "block-padding" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8894febbff9f758034a5b8e12d87918f56dfc64a8e1fe757d65e29041538d93" +dependencies = [ + "generic-array", +] + [[package]] name = "block2" version = "0.6.2" @@ -651,6 +784,29 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "64fa3c856b712db6612c019f14756e64e4bcea13337a6b33b696333a9eaa2d06" +[[package]] +name = "bytecheck" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0caa33a2c0edca0419d15ac723dff03f1956f7978329b1e3b5fdaaaed9d3ca8b" +dependencies = [ + "bytecheck_derive", + "ptr_meta", + "rancor", + "simdutf8", +] + +[[package]] +name = "bytecheck_derive" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "89385e82b5d1821d2219e0b095efa2cc1f246cbf99080f3be46a1a85c0d392d9" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "bytecount" version = "0.6.9" @@ -677,6 +833,12 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + [[package]] name = "byteorder-lite" version = "0.1.0" @@ -760,6 +922,15 @@ dependencies = [ "displaydoc", ] +[[package]] +name = "cbc" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26b52a9543ae338f279b96b0b9fed9c8093744685043739079ce85cd58f289a6" +dependencies = [ + "cipher", +] + [[package]] name = "cc" version = "1.2.67" @@ -772,6 +943,18 @@ dependencies = [ "shlex 2.0.1", ] +[[package]] +name = "ccm" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ae3c82e4355234767756212c570e29833699ab63e6ffd161887314cc5b43847" +dependencies = [ + "aead", + "cipher", + "ctr", + "subtle", +] + [[package]] name = "cexpr" version = "0.6.0" @@ -809,6 +992,17 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f079e83a288787bcd14a6aea84cee5c87a67c5a3e660c30f557a3d24761b3527" +[[package]] +name = "chacha20" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3613f74bd2eac03dad61bd53dbe620703d4371614fe0bc3b9f04dd36fe4e818" +dependencies = [ + "cfg-if", + "cipher", + "cpufeatures 0.2.17", +] + [[package]] name = "chacha20" version = "0.10.1" @@ -820,6 +1014,19 @@ dependencies = [ "rand_core 0.10.1", ] +[[package]] +name = "chacha20poly1305" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "10cd79432192d1c0f4e1a0fef9527696cc039165d729fb41b3f4f4f354c2dc35" +dependencies = [ + "aead", + "chacha20 0.9.1", + "cipher", + "poly1305", + "zeroize", +] + [[package]] name = "check_elevation" version = "0.2.7" @@ -852,6 +1059,17 @@ dependencies = [ "phf 0.12.1", ] +[[package]] +name = "cipher" +version = "0.4.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad" +dependencies = [ + "crypto-common 0.1.7", + "inout", + "zeroize", +] + [[package]] name = "clang-sys" version = "1.8.1" @@ -913,6 +1131,15 @@ dependencies = [ "error-code", ] +[[package]] +name = "cmake" +version = "0.1.58" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c0f78a02292a74a88ac736019ab962ece0bc380e3f977bf72e376c5d78ff0678" +dependencies = [ + "cc", +] + [[package]] name = "codesnake" version = "0.2.1" @@ -989,6 +1216,12 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "const-oid" +version = "0.9.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" + [[package]] name = "const-oid" version = "0.10.2" @@ -1112,6 +1345,21 @@ dependencies = [ "libc", ] +[[package]] +name = "crc" +version = "3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5eb8a2a1cd12ab0d987a5d5e825195d372001a4094a0376319d5a0ad71c1ba0d" +dependencies = [ + "crc-catalog", +] + +[[package]] +name = "crc-catalog" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853" + [[package]] name = "crc-fast" version = "1.10.0" @@ -1162,6 +1410,18 @@ version = "0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" +[[package]] +name = "crypto-bigint" +version = "0.5.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76" +dependencies = [ + "generic-array", + "rand_core 0.6.4", + "subtle", + "zeroize", +] + [[package]] name = "crypto-common" version = "0.1.7" @@ -1169,6 +1429,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a" dependencies = [ "generic-array", + "rand_core 0.6.4", "typenum", ] @@ -1187,6 +1448,41 @@ version = "1.0.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e2e30e509674ef0ec91e21a7735766db37d163d46151b6a361d8b83dd79116bd" +[[package]] +name = "ctr" +version = "0.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0369ee1ad671834580515889b80f2ea915f23b8be8d0daa4bbaf2ac5c7590835" +dependencies = [ + "cipher", +] + +[[package]] +name = "curve25519-dalek" +version = "4.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97fb8b7c4503de7d6ae7b42ab72a5a59857b4c937ec27a3d4539dba95b5ab2be" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "curve25519-dalek-derive", + "fiat-crypto", + "rustc_version", + "subtle", + "zeroize", +] + +[[package]] +name = "curve25519-dalek-derive" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f46882e17999c6cc590af592290432be3bce0428cb0d5f8b6715e4dc7b383eb3" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "darling" version = "0.20.11" @@ -1327,6 +1623,37 @@ dependencies = [ "thiserror 2.0.19", ] +[[package]] +name = "der" +version = "0.7.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb" +dependencies = [ + "const-oid 0.9.6", + "pem-rfc7468", + "zeroize", +] + +[[package]] +name = "der-parser" +version = "9.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5cd0a5c643689626bec213c4d8bd4d96acc8ffdb4ad4bb6bc16abf27d5f4b553" +dependencies = [ + "asn1-rs", + "displaydoc", + "nom 7.1.3", + "num-bigint", + "num-traits", + "rusticata-macros", +] + +[[package]] +name = "deranged" +version = "0.5.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7cd812cc2bc1d69d4764bd80df88b4317eaef9e773c75226407d9bc0876b211c" + [[package]] name = "digest" version = "0.10.7" @@ -1334,7 +1661,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" dependencies = [ "block-buffer 0.10.4", + "const-oid 0.9.6", "crypto-common 0.1.7", + "subtle", ] [[package]] @@ -1344,7 +1673,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" dependencies = [ "block-buffer 0.12.1", - "const-oid", + "const-oid 0.10.2", "crypto-common 0.2.2", ] @@ -1388,7 +1717,7 @@ checksum = "6e39034cee21a2f5bbb66ba0e3689819c4bb5d00382a282006e802a7ffa6c41d" dependencies = [ "cfg-if", "libc", - "socket2", + "socket2 0.6.5", "windows-sys 0.60.2", ] @@ -1453,6 +1782,42 @@ dependencies = [ "linux-raw-sys 0.9.4", ] +[[package]] +name = "dtls" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "01a431e87fc386bd5e02deb554a013f97bc406e47a7d8e97efb7c6366b980e9c" +dependencies = [ + "aes", + "aes-gcm", + "async-trait", + "bytecheck", + "byteorder", + "cbc", + "ccm", + "chacha20poly1305", + "der-parser", + "hmac", + "log", + "p256", + "p384", + "portable-atomic", + "rand 0.9.5", + "rand_core 0.6.4", + "rcgen", + "ring", + "rkyv", + "rustls", + "sec1", + "sha1", + "sha2", + "thiserror 1.0.69", + "tokio", + "webrtc-util", + "x25519-dalek", + "x509-parser", +] + [[package]] name = "dunce" version = "1.0.5" @@ -1465,12 +1830,47 @@ version = "1.0.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555" +[[package]] +name = "ecdsa" +version = "0.16.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee27f32b5c5292967d2d4a9d7f1e0b0aed2c15daded5a60300e4abb9d8020bca" +dependencies = [ + "der", + "digest 0.10.7", + "elliptic-curve", + "rfc6979", + "signature", + "spki", +] + [[package]] name = "either" version = "1.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "91622ff5e7162018101f2fea40d6ebf4a78bbe5a49736a2020649edf9693679e" +[[package]] +name = "elliptic-curve" +version = "0.13.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b5e6043086bf7973472e0c7dff2142ea0b680d30e18d9cc40f267efbf222bd47" +dependencies = [ + "base16ct", + "crypto-bigint", + "digest 0.10.7", + "ff", + "generic-array", + "group", + "hkdf", + "pem-rfc7468", + "pkcs8", + "rand_core 0.6.4", + "sec1", + "subtle", + "zeroize", +] + [[package]] name = "encode_unicode" version = "1.0.0" @@ -1655,6 +2055,22 @@ dependencies = [ "simd-adler32", ] +[[package]] +name = "ff" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c0b50bfb653653f9ca9095b427bed08ab8d75a137839d9ad64eb11810d5b6393" +dependencies = [ + "rand_core 0.6.4", + "subtle", +] + +[[package]] +name = "fiat-crypto" +version = "0.2.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d" + [[package]] name = "filedescriptor" version = "0.8.3" @@ -1989,6 +2405,7 @@ checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a" dependencies = [ "typenum", "version_check", + "zeroize", ] [[package]] @@ -2040,6 +2457,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "ghash" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0d8a4362ccb29cb0b265253fb0a2728f592895ee6854fd9bc13f2ffda266ff1" +dependencies = [ + "opaque-debug", + "polyval", +] + [[package]] name = "gif" version = "0.14.2" @@ -2166,6 +2593,17 @@ dependencies = [ "memmap2", ] +[[package]] +name = "group" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0f9ef7462f7c099f518d754361858f86d8a07af53ba9af0fe635bbccb151a63" +dependencies = [ + "ff", + "rand_core 0.6.4", + "subtle", +] + [[package]] name = "half" version = "2.7.1" @@ -2235,6 +2673,24 @@ version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0a7763b98ba8a24f59e698bf9ab197e7676c640d6455d1580b4ce7dc560f0f0d" +[[package]] +name = "hkdf" +version = "0.12.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b5f8eb2ad728638ea2c7d47a21db23b7b58a72ed6a38256b8a1849f15fbbdf7" +dependencies = [ + "hmac", +] + +[[package]] +name = "hmac" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" +dependencies = [ + "digest 0.10.7", +] + [[package]] name = "hostname" version = "0.4.2" @@ -2734,6 +3190,16 @@ dependencies = [ "libc", ] +[[package]] +name = "inout" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01" +dependencies = [ + "block-padding", + "generic-array", +] + [[package]] name = "insta" version = "1.48.0" @@ -2749,6 +3215,27 @@ dependencies = [ "tempfile", ] +[[package]] +name = "interceptor" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "88c11a956a48159f7fe539b8198f12b4db9b709ae5f94385b840db38f97fed74" +dependencies = [ + "async-trait", + "bytes", + "futures", + "log", + "portable-atomic", + "rand 0.9.5", + "rtcp", + "rtp", + "thiserror 1.0.69", + "tokio", + "waitgroup", + "webrtc-srtp", + "webrtc-util", +] + [[package]] name = "intl-memoizer" version = "0.5.3" @@ -2768,6 +3255,12 @@ dependencies = [ "unic-langid", ] +[[package]] +name = "ipnet" +version = "2.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2" + [[package]] name = "is_terminal_polyfill" version = "1.70.2" @@ -3055,7 +3548,7 @@ version = "0.9.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "901049455d2eb6decf9058235d745237952f4804bc584c5fcb41412e6adcc6e0" dependencies = [ - "bindgen", + "bindgen 0.72.1", "cc", "system-deps", ] @@ -3181,6 +3674,15 @@ dependencies = [ "libc", ] +[[package]] +name = "memoffset" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5de893c32cde5f383baa4c04c5d6dbdd735cfd4a794b0debdb2bb1b421da5ff4" +dependencies = [ + "autocfg", +] + [[package]] name = "memoffset" version = "0.9.1" @@ -3228,6 +3730,26 @@ dependencies = [ "pxfm", ] +[[package]] +name = "munge" +version = "0.4.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e17401f259eba956ca16491461b6e8f72913a0a114e39736ce404410f915a0c" +dependencies = [ + "munge_macro", +] + +[[package]] +name = "munge_macro" +version = "0.4.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4568f25ccbd45ab5d5603dc34318c1ec56b117531781260002151b8530a9f931" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "nanorand" version = "0.7.0" @@ -3312,6 +3834,19 @@ dependencies = [ "libc", ] +[[package]] +name = "nix" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "598beaf3cc6fdd9a5dfb1630c2800c7acd31df7aaf0f565796fba2b53ca1af1b" +dependencies = [ + "bitflags 1.3.2", + "cfg-if", + "libc", + "memoffset 0.7.1", + "pin-utils", +] + [[package]] name = "nix" version = "0.28.0" @@ -3437,6 +3972,12 @@ dependencies = [ "num-traits", ] +[[package]] +name = "num-conv" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "521739c6d2bac4aa25192232afe6841231376b2b26d4d9fae5ecf8ca5772e441" + [[package]] name = "num-format" version = "0.4.4" @@ -3744,6 +4285,37 @@ dependencies = [ "objc2-core-foundation", ] +[[package]] +name = "oid-registry" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8d8034d9489cdaf79228eb9f6a3b8d7bb32ba00d6645ebd48eef4077ceb5bd9" +dependencies = [ + "asn1-rs", +] + +[[package]] +name = "om-fork-ep-miniaudio-sys" +version = "2.6.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b15275d28e354c1e3665940e8bd668817d6a6bafb6bdcb3860b49acf9358281f" +dependencies = [ + "bindgen 0.71.1", + "bitflags 2.13.1", + "cc", + "libc", +] + +[[package]] +name = "om-fork-miniaudio" +version = "0.12.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dec79816e1fe9db8c435cb3efac5cf470c4e7bde239ebcac22742412878aaeff" +dependencies = [ + "bitflags 2.13.1", + "om-fork-ep-miniaudio-sys", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -3778,6 +4350,21 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "opaque-debug" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c08d65885ee38876c4f86fa503fb49d7b507c2b62552df7c70b2fce627e06381" + +[[package]] +name = "opus" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4d3809943dff6fbad5f0484449ea26bdb9cb7d8efdf26ed50d3c7f227f69eb5c" +dependencies = [ + "audiopus_sys", +] + [[package]] name = "ordered-float" version = "5.3.0" @@ -3822,6 +4409,30 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1a80800c0488c3a21695ea981a54918fbb37abf04f4d0720c453632255e2ff0e" +[[package]] +name = "p256" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c9863ad85fa8f4460f9c48cb909d38a0d689dba1f6f6988a5e3e0d31071bcd4b" +dependencies = [ + "ecdsa", + "elliptic-curve", + "primeorder", + "sha2", +] + +[[package]] +name = "p384" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fe42f1670a52a47d448f14b6a5c61dd78fce51856e68edaa38f7ae3a46b8d6b6" +dependencies = [ + "ecdsa", + "elliptic-curve", + "primeorder", + "sha2", +] + [[package]] name = "palette" version = "0.7.6" @@ -3935,6 +4546,25 @@ version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7011d97b484a5ebdc4b1fdb3b12d5e4bbbea56e9d22b688f2e79e04b65a7d8a6" +[[package]] +name = "pem" +version = "3.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d30c53c26bc5b31a98cd02d20f25a7c8567146caf63ed593a9d87b2775291be" +dependencies = [ + "base64", + "serde_core", +] + +[[package]] +name = "pem-rfc7468" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "88b39c9bfcfc231068454382784bb460aae594343fb030d46e9f50a645418412" +dependencies = [ + "base64ct", +] + [[package]] name = "percent-encoding" version = "2.3.2" @@ -4191,7 +4821,9 @@ dependencies = [ "anyhow", "arboard", "ast-grep-core", + "audiopus_sys", "base64", + "bytes", "clap", "clipboard-win", "core-graphics", @@ -4212,6 +4844,8 @@ dependencies = [ "napi", "napi-build", "napi-derive", + "om-fork-miniaudio", + "opus", "parking_lot", "phf 0.13.1", "pi-ast", @@ -4234,6 +4868,7 @@ dependencies = [ "toml", "unicode-segmentation", "unicode-width 0.2.2", + "webrtc", "windows-sys 0.61.2", "winreg 0.56.0", "x11rb", @@ -4376,6 +5011,12 @@ version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" +[[package]] +name = "pin-utils" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" + [[package]] name = "piper" version = "0.2.5" @@ -4410,11 +5051,21 @@ version = "0.9.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cb028afee0d6ca17020b090e3b8fa2d7de23305aef975c7e5192a5050246ea36" dependencies = [ - "bindgen", + "bindgen 0.72.1", "libspa-sys", "system-deps", ] +[[package]] +name = "pkcs8" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f950b2377845cebe5cf8b5165cb3cc1a5e0fa5cfa3e1f7f55707d8fd82e0a7b7" +dependencies = [ + "der", + "spki", +] + [[package]] name = "pkg-config" version = "0.3.33" @@ -4458,6 +5109,29 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "poly1305" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8159bd90725d2df49889a078b54f4f79e87f1f8a8444194cdca81d38f5393abf" +dependencies = [ + "cpufeatures 0.2.17", + "opaque-debug", + "universal-hash", +] + +[[package]] +name = "polyval" +version = "0.6.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9d1fe60d06143b2430aa532c94cfe9e29783047f06c0d7fd359a9a51b729fa25" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "opaque-debug", + "universal-hash", +] + [[package]] name = "portable-atomic" version = "1.14.0" @@ -4505,6 +5179,12 @@ dependencies = [ "zerovec", ] +[[package]] +name = "powerfmt" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" + [[package]] name = "ppv-lite86" version = "0.2.21" @@ -4530,6 +5210,15 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "primeorder" +version = "0.13.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "353e1ca18966c16d9deb1c69278edbc5f194139612772bd9537af60ac231e1e6" +dependencies = [ + "elliptic-curve", +] + [[package]] name = "proc-macro-crate" version = "3.5.0" @@ -4572,6 +5261,26 @@ dependencies = [ "hex", ] +[[package]] +name = "ptr_meta" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b9a0cf95a1196af61d4f1cbdab967179516d9a4a4312af1f31948f8f6224a79" +dependencies = [ + "ptr_meta_derive", +] + +[[package]] +name = "ptr_meta_derive" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7347867d0a7e1208d93b46767be83e2b8f978c3dad35f775ac8d8847551d6fe1" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "pxfm" version = "0.1.30" @@ -4649,6 +5358,15 @@ version = "0.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dc33ff2d4973d518d823d61aa239014831e521c75da58e3df4840d3f47749d09" +[[package]] +name = "rancor" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "daff8b7b3ccf5f7ba270b3e7a0a4d4c701c5797e38dec27c7e2c3dbb830fed1c" +dependencies = [ + "ptr_meta", +] + [[package]] name = "rand" version = "0.8.7" @@ -4674,7 +5392,7 @@ version = "0.10.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" dependencies = [ - "chacha20", + "chacha20 0.10.1", "getrandom 0.4.3", "rand_core 0.10.1", ] @@ -4694,6 +5412,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" @@ -4739,6 +5460,20 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "rcgen" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75e669e5202259b5314d1ea5397316ad400819437857b90861765f24c4cf80a2" +dependencies = [ + "pem", + "ring", + "rustls-pki-types", + "time", + "x509-parser", + "yasna", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -4809,6 +5544,25 @@ version = "1.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba39f3699c378cd8970968dcbff9c43159ea4cfbd88d43c00b22f2ef10a435d2" +[[package]] +name = "rend" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "663ba70707f96e871406fe10d68128412e619b06d1d47cb91c3a4c6501176240" +dependencies = [ + "bytecheck", +] + +[[package]] +name = "rfc6979" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8dd2a808d456c4a54e300a23e9f5a67e122c3024119acbfd73e3bf664491cb2" +dependencies = [ + "hmac", + "subtle", +] + [[package]] name = "rgb" version = "0.8.53" @@ -4818,6 +5572,50 @@ dependencies = [ "bytemuck", ] +[[package]] +name = "ring" +version = "0.17.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" +dependencies = [ + "cc", + "cfg-if", + "getrandom 0.2.17", + "libc", + "untrusted", + "windows-sys 0.52.0", +] + +[[package]] +name = "rkyv" +version = "0.8.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "815cc8a37159a463064825246cadb07961e25cd9885908606f6d08a98d8f8874" +dependencies = [ + "bytecheck", + "bytes", + "hashbrown 0.17.1", + "indexmap", + "munge", + "ptr_meta", + "rancor", + "rend", + "rkyv_derive", + "tinyvec", + "uuid", +] + +[[package]] +name = "rkyv_derive" +version = "0.8.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c0ed1a78a1b19d184b0daa629dd9a024573173ec7d485b287cb369fb3607cc1c" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "rlimit" version = "0.11.0" @@ -4866,6 +5664,32 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "rtcp" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "adad7f6a501162881032fc84b4bc78ae11f1b30748180b9f5cbe8810bf613aab" +dependencies = [ + "bytes", + "thiserror 1.0.69", + "webrtc-util", +] + +[[package]] +name = "rtp" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "149329e78ada26b5e174a4c281a7a7e5c8bda6008adc90daddd6f98fce56db29" +dependencies = [ + "bytes", + "memchr", + "portable-atomic", + "rand 0.9.5", + "serde", + "thiserror 1.0.69", + "webrtc-util", +] + [[package]] name = "rustc-hash" version = "1.1.0" @@ -4887,6 +5711,15 @@ dependencies = [ "semver", ] +[[package]] +name = "rusticata-macros" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "faf0c4a6ece9950b9abdb62b1cfcf2a68b3b67a10ba445b3bb85be2a293d0632" +dependencies = [ + "nom 7.1.3", +] + [[package]] name = "rustix" version = "0.38.44" @@ -4913,6 +5746,40 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "rustls" +version = "0.23.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" +dependencies = [ + "once_cell", + "ring", + "rustls-pki-types", + "rustls-webpki", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96" +dependencies = [ + "zeroize", +] + +[[package]] +name = "rustls-webpki" +version = "0.103.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +dependencies = [ + "ring", + "rustls-pki-types", + "untrusted", +] + [[package]] name = "rustversion" version = "1.0.23" @@ -4949,6 +5816,32 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" +[[package]] +name = "sdp" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "66b6eecfa5151edef84d544ff3885dff98e5b5fe97586757a4b5118bde51c958" +dependencies = [ + "rand 0.9.5", + "substring", + "thiserror 1.0.69", + "url", +] + +[[package]] +name = "sec1" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3e97a565f76233a6003f9f5c54be1d9c5bdfa3eccfb189469f11ec4901c47dc" +dependencies = [ + "base16ct", + "der", + "generic-array", + "pkcs8", + "subtle", + "zeroize", +] + [[package]] name = "self_cell" version = "1.3.0" @@ -5106,12 +5999,28 @@ dependencies = [ "libc", ] +[[package]] +name = "signature" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" +dependencies = [ + "digest 0.10.7", + "rand_core 0.6.4", +] + [[package]] name = "simd-adler32" version = "0.3.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3a219298ac11a56ea9a6d2120044824d6f01aeb034955e7af7bc16858527deea" +[[package]] +name = "simdutf8" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e3a9fe34e3e7a50316060351f37187a3f546bce95496156754b601a5fa71b76e" + [[package]] name = "similar" version = "2.7.0" @@ -5157,6 +6066,25 @@ dependencies = [ "serde", ] +[[package]] +name = "smol_str" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dd538fb6910ac1099850255cf94a94df6551fbdd602454387d0adb2d1ca6dead" +dependencies = [ + "serde", +] + +[[package]] +name = "socket2" +version = "0.5.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e22376abed350d73dd1cd119b57ffccad95b4e585a7cda43e286245ce23c0678" +dependencies = [ + "libc", + "windows-sys 0.52.0", +] + [[package]] name = "socket2" version = "0.6.5" @@ -5182,6 +6110,16 @@ version = "0.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "023a211cb3138dbc438680b32560ad89f699977624c9f8dbb95a47d5b4c07dd3" +[[package]] +name = "spki" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d91ed6c858b01f942cd56b37a94b3e0a1798290327d1236e4d9cf4eaca44d29d" +dependencies = [ + "base64ct", + "der", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -5248,6 +6186,40 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "stun" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e3b8d5ad1dd4c8cd0595440f26521a71dd593cb5873ee8e06b91510b7d97269" +dependencies = [ + "base64", + "crc", + "lazy_static", + "md-5", + "rand 0.9.5", + "ring", + "subtle", + "thiserror 1.0.69", + "tokio", + "url", + "webrtc-util", +] + +[[package]] +name = "substring" +version = "1.4.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "42ee6433ecef213b2e72f587ef64a2f5943e7cd16fbd82dbe8bc07486c534c86" +dependencies = [ + "autocfg", +] + +[[package]] +name = "subtle" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" + [[package]] name = "syn" version = "2.0.119" @@ -5447,6 +6419,36 @@ dependencies = [ "rustc-hash 1.1.0", ] +[[package]] +name = "time" +version = "0.3.54" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3e1d5e639ff6bab73cb6885cc7e7b1de96c3f32c68ec55f3952614bec1092244" +dependencies = [ + "deranged", + "num-conv", + "powerfmt", + "serde_core", + "time-core", + "time-macros", +] + +[[package]] +name = "time-core" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9e1c906769ad99c88eaa54e728060edef082f8e358ff32030cb7c7d315e81109" + +[[package]] +name = "time-macros" +version = "0.2.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e689342a48d2ea927c87ea50cabf8594854bf940e9310208848d680d668ed85" +dependencies = [ + "num-conv", + "time-core", +] + [[package]] name = "tiny-keccak" version = "2.0.2" @@ -5467,6 +6469,21 @@ dependencies = [ "zerovec", ] +[[package]] +name = "tinyvec" +version = "1.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bb4ebadaa0af04fab11ae01eb5f9fdb5f9c5b875506e210e71c07873528baa7f" +dependencies = [ + "tinyvec_macros", +] + +[[package]] +name = "tinyvec_macros" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" + [[package]] name = "tokio" version = "1.53.1" @@ -5479,7 +6496,7 @@ dependencies = [ "parking_lot", "pin-project-lite", "signal-hook-registry", - "socket2", + "socket2 0.6.5", "tokio-macros", "windows-sys 0.61.2", ] @@ -6200,6 +7217,27 @@ version = "0.21.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2c591d83f69777866b9126b24c6dd9a18351f177e49d625920d19f989fd31cf8" +[[package]] +name = "turn" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d99249a493335eb44c4d7943a8751b22561a82c9903d8eb2d29780b3927880a7" +dependencies = [ + "async-trait", + "base64", + "futures", + "log", + "md-5", + "portable-atomic", + "rand 0.9.5", + "ring", + "stun", + "thiserror 1.0.69", + "tokio", + "tokio-util", + "webrtc-util", +] + [[package]] name = "type-map" version = "0.5.1" @@ -6233,7 +7271,7 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ - "memoffset", + "memoffset 0.9.1", "tempfile", "windows-sys 0.61.2", ] @@ -6256,6 +7294,12 @@ dependencies = [ "tinystr", ] +[[package]] +name = "unicase" +version = "2.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dbc4bc3a9f746d862c45cb89d705aa10f187bb96c76001afab07a0d35ce60142" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -6286,6 +7330,22 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "81e544489bf3d8ef66c953931f56617f423cd4b5494be343d9b9d3dda037b9a3" +[[package]] +name = "universal-hash" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc1de2c688dc15305988b563c3854064043356019f97a4b46276fe734c4f07ea" +dependencies = [ + "crypto-common 0.1.7", + "subtle", +] + +[[package]] +name = "untrusted" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" + [[package]] name = "url" version = "2.5.8" @@ -7024,6 +8084,7 @@ version = "1.24.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "bf3923a6f5c4c6382e0b653c4117f48d631ea17f38ed86e2a828e6f7412f5239" dependencies = [ + "getrandom 0.4.3", "js-sys", "serde_core", "wasm-bindgen", @@ -7066,6 +8127,15 @@ version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5c3082ca00d5a5ef149bb8b555a72ae84c9c59f7250f013ac822ac2e49b19c64" +[[package]] +name = "waitgroup" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d1f50000a783467e6c0200f9d10642f4bc424e39efc1b770203e88b488f79292" +dependencies = [ + "atomic-waker", +] + [[package]] name = "walkdir" version = "2.5.0" @@ -7238,7 +8308,7 @@ dependencies = [ "dlib", "libc", "log", - "memoffset", + "memoffset 0.9.1", "pkg-config", ] @@ -7274,6 +8344,175 @@ dependencies = [ "string_cache_codegen", ] +[[package]] +name = "webrtc" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "baaacdf9d96224d7b6e2872ba065578f38775f7634367e3ac2cc87c7271da433" +dependencies = [ + "arc-swap", + "async-trait", + "bytes", + "dtls", + "hex", + "interceptor", + "lazy_static", + "log", + "portable-atomic", + "rand 0.9.5", + "rcgen", + "regex", + "ring", + "rtcp", + "rtp", + "sdp", + "serde", + "serde_json", + "sha2", + "smol_str", + "stun", + "thiserror 1.0.69", + "tokio", + "turn", + "unicase", + "url", + "waitgroup", + "webrtc-data", + "webrtc-ice", + "webrtc-mdns", + "webrtc-media", + "webrtc-sctp", + "webrtc-srtp", + "webrtc-util", +] + +[[package]] +name = "webrtc-data" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd470286275809f2fcfcdb1e73ef5f1500be82eff7fe98150ce81b20aad5a2a4" +dependencies = [ + "bytes", + "log", + "portable-atomic", + "thiserror 1.0.69", + "tokio", + "webrtc-sctp", + "webrtc-util", +] + +[[package]] +name = "webrtc-ice" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b7fd30f52e6fda8664779b84b7904b2553b76fee24d9ca665e774ae32b13f53" +dependencies = [ + "arc-swap", + "async-trait", + "crc", + "log", + "portable-atomic", + "rand 0.9.5", + "serde", + "serde_json", + "stun", + "thiserror 1.0.69", + "tokio", + "turn", + "url", + "uuid", + "waitgroup", + "webrtc-mdns", + "webrtc-util", +] + +[[package]] +name = "webrtc-mdns" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91ffa0ea00c0fae979aafa5db285fec5eeb360e1ea246ebec529ba4e09f8e03e" +dependencies = [ + "log", + "socket2 0.5.10", + "thiserror 1.0.69", + "tokio", + "webrtc-util", +] + +[[package]] +name = "webrtc-media" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26a6c7335bdd03dc023cb9a3bd7966866d969c92c86987b73aa16dfa39c7ca29" +dependencies = [ + "byteorder", + "bytes", + "rand 0.9.5", + "rtp", + "thiserror 1.0.69", +] + +[[package]] +name = "webrtc-sctp" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1c4b637f0d8eb96d900ac0f79b3060ddd21ca88dfefa467e96e67ab0016f2574" +dependencies = [ + "arc-swap", + "async-trait", + "bytes", + "crc", + "log", + "portable-atomic", + "rand 0.9.5", + "thiserror 1.0.69", + "tokio", + "webrtc-util", +] + +[[package]] +name = "webrtc-srtp" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45d2b667a0b5d04eebcb7cb22fd51b2af3721f418097b496e72de16386c8d0c3" +dependencies = [ + "aead", + "aes", + "aes-gcm", + "byteorder", + "bytes", + "ctr", + "hmac", + "log", + "rtcp", + "rtp", + "sha1", + "subtle", + "thiserror 1.0.69", + "tokio", + "webrtc-util", +] + +[[package]] +name = "webrtc-util" +version = "0.17.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1beae0b4f24969741c26ca282a257dc5823bd9cbe608f7e3ca3775102d323ad9" +dependencies = [ + "async-trait", + "bitflags 1.3.2", + "bytes", + "ipnet", + "lazy_static", + "log", + "nix 0.26.4", + "portable-atomic", + "rand 0.9.5", + "thiserror 1.0.69", + "tokio", + "winapi", +] + [[package]] name = "weezl" version = "0.1.12" @@ -7549,6 +8788,15 @@ dependencies = [ "windows-link 0.2.1", ] +[[package]] +name = "windows-sys" +version = "0.52.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "282be5f36a8ce781fad8c8ae18fa3f9beff57ec1b52cb3de0789201425d9a33d" +dependencies = [ + "windows-targets 0.52.6", +] + [[package]] name = "windows-sys" version = "0.59.0" @@ -7882,6 +9130,36 @@ version = "0.13.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ea6fc2961e4ef194dcbfe56bb845534d0dc8098940c7e5c012a258bfec6701bd" +[[package]] +name = "x25519-dalek" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7e468321c81fb07fa7f4c636c3972b9100f0346e5b6a9f2bd0603a52f7ed277" +dependencies = [ + "curve25519-dalek", + "rand_core 0.6.4", + "serde", + "zeroize", +] + +[[package]] +name = "x509-parser" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fcbc162f30700d6f3f82a24bf7cc62ffe7caea42c0b2cba8bf7f3ae50cf51f69" +dependencies = [ + "asn1-rs", + "data-encoding", + "der-parser", + "lazy_static", + "nom 7.1.3", + "oid-registry", + "ring", + "rusticata-macros", + "thiserror 1.0.69", + "time", +] + [[package]] name = "xattr" version = "1.6.1" @@ -7978,6 +9256,15 @@ version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cfe53a6657fd280eaa890a3bc59152892ffa3e30101319d168b781ed6529b049" +[[package]] +name = "yasna" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e17bb3549cc1321ae1296b9cdc2698e2b6cb1992adfa19a8c72e5b7a738f44cd" +dependencies = [ + "time", +] + [[package]] name = "yoke" version = "0.8.3" @@ -8109,6 +9396,26 @@ dependencies = [ "synstructure", ] +[[package]] +name = "zeroize" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e" +dependencies = [ + "zeroize_derive", +] + +[[package]] +name = "zeroize_derive" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c50655cbb0fe3fc43170059e702f1ce5e19b84cec58dc87b037a09935c2f328" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "zerotrie" version = "0.2.4" diff --git a/Cargo.toml b/Cargo.toml index a3629b610..ca0b335c6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -238,6 +238,14 @@ smallvec = { version = "1.15.1", features = [ xxhash-rust = { version = "0.8", features = ["xxh64"] } +# ────────────────────────────────────────────────────────────────────────────── +# Audio & Realtime Media +# ────────────────────────────────────────────────────────────────────────────── +audiopus_sys = { version = "0.2.2", features = ["static"] } +miniaudio = { package = "om-fork-miniaudio", version = "0.12.2" } +opus = "0.3.1" +webrtc = "0.17.2" + # ────────────────────────────────────────────────────────────────────────────── # System & Platform # ────────────────────────────────────────────────────────────────────────────── diff --git a/crates/pi-natives/Cargo.toml b/crates/pi-natives/Cargo.toml index 2d8907311..615950534 100644 --- a/crates/pi-natives/Cargo.toml +++ b/crates/pi-natives/Cargo.toml @@ -15,9 +15,11 @@ workspace = true [dependencies] anyhow.workspace = true +audiopus_sys.workspace = true arboard.workspace = true ast-grep-core.workspace = true base64.workspace = true +bytes.workspace = true clap.workspace = true globset.workspace = true fontdue.workspace = true @@ -30,8 +32,10 @@ icy_sixel.workspace = true ignore.workspace = true image = { workspace = true, features = ["bmp"] } inferno.workspace = true +miniaudio.workspace = true napi.workspace = true napi-derive.workspace = true +opus.workspace = true parking_lot.workspace = true phf.workspace = true flume.workspace = true @@ -54,6 +58,7 @@ tokio-util.workspace = true toml.workspace = true unicode-segmentation.workspace = true unicode-width.workspace = true +webrtc.workspace = true xxhash-rust.workspace = true [target.'cfg(target_os = "linux")'.dependencies] diff --git a/crates/pi-natives/src/audio.rs b/crates/pi-natives/src/audio.rs new file mode 100644 index 000000000..fe0e62c8e --- /dev/null +++ b/crates/pi-natives/src/audio.rs @@ -0,0 +1,427 @@ +//! Cross-platform microphone capture and streaming speaker playback. +//! +//! miniaudio owns platform device discovery, format conversion, channel mixing, +//! and resampling. The N-API classes expose one stable mono `f32` contract to +//! TypeScript while the internal playback stream is shared with native WebRTC. + +use std::sync::{ + Arc, + atomic::{AtomicBool, AtomicU32, Ordering}, +}; + +use flume::TryRecvError; +use miniaudio::{Device, DeviceConfig, DeviceType, Format, PerformanceProfile}; +use napi::{ + bindgen_prelude::{Float32Array, Result}, + threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode, UnknownReturnValue}, +}; +use napi_derive::napi; +use parking_lot::Mutex; +use tokio::sync::Notify; + +const AUDIO_CHANNELS: u32 = 1; +const AUDIO_PERIOD_MS: u32 = 20; +const PLAYBACK_DRAIN_CALLBACKS: usize = 2; + +type CaptureCallback = ThreadsafeFunction; +type NativeResult = std::result::Result; + +struct PlaybackState { + gain_bits: AtomicU32, + drained: AtomicBool, + stopped: AtomicBool, + notify: Notify, +} + +impl PlaybackState { + fn new() -> Self { + Self { + gain_bits: AtomicU32::new(1.0f32.to_bits()), + drained: AtomicBool::new(false), + stopped: AtomicBool::new(false), + notify: Notify::new(), + } + } + + fn gain(&self) -> f32 { + f32::from_bits(self.gain_bits.load(Ordering::Acquire)) + } + + fn set_gain(&self, gain: f32) { + self.gain_bits.store(gain.to_bits(), Ordering::Release); + } + + fn mark_drained(&self) { + if !self.drained.swap(true, Ordering::AcqRel) { + self.notify.notify_waiters(); + } + } + + fn mark_stopped(&self) { + self.stopped.store(true, Ordering::Release); + self.notify.notify_waiters(); + } + + async fn wait_for_drain(&self) { + loop { + let notified = self.notify.notified(); + if self.drained.load(Ordering::Acquire) || self.stopped.load(Ordering::Acquire) { + return; + } + notified.await; + } + } +} + +/// Producer endpoint for one native playback device. +#[derive(Clone)] +pub(crate) struct PlaybackWriter { + tx: flume::Sender>, + state: Arc, +} + +impl PlaybackWriter { + /// Queue mono floating-point samples without blocking the caller. + pub(crate) fn write(&self, samples: &[f32]) -> NativeResult<()> { + if samples.is_empty() { + return Ok(()); + } + if self.state.stopped.load(Ordering::Acquire) || self.state.drained.load(Ordering::Acquire) { + return Err("Native audio playback is closed".to_owned()); + } + self.tx + .send(samples.to_vec()) + .map_err(|_| "Native audio playback is closed".to_owned()) + } +} + +/// Running mono playback stream shared by N-API playback and native WebRTC. +pub(crate) struct PlaybackStream { + device: Option, + writer: Option, + state: Arc, +} + +impl PlaybackStream { + /// Open and start the default speaker at the requested logical sample rate. + pub(crate) fn start(sample_rate: u32) -> NativeResult { + validate_sample_rate(sample_rate)?; + let state = Arc::new(PlaybackState::new()); + let (tx, rx) = flume::unbounded::>(); + let mut config = audio_config(DeviceType::Playback, sample_rate); + config.playback_mut().set_format(Format::F32); + config.playback_mut().set_channels(AUDIO_CHANNELS); + let mut device = Device::new(None, &config) + .map_err(|error| format!("Failed to open the default speaker: {error}"))?; + + let callback_state = Arc::clone(&state); + let mut current = Vec::new(); + let mut cursor = 0; + let mut empty_callbacks = 0; + device.set_data_callback(move |_device, output, _input| { + fill_playback( + &rx, + &mut current, + &mut cursor, + output.as_samples_mut::(), + &callback_state, + &mut empty_callbacks, + ); + }); + let stop_state = Arc::clone(&state); + device.set_stop_callback(move |_device| stop_state.mark_stopped()); + device + .start() + .map_err(|error| format!("Failed to start speaker playback: {error}"))?; + + Ok(Self { + device: Some(device), + writer: Some(PlaybackWriter { tx, state: Arc::clone(&state) }), + state, + }) + } + + /// Clone the producer endpoint used by the remote-audio decoder. + pub(crate) fn writer(&self) -> NativeResult { + self.writer + .clone() + .ok_or_else(|| "Native audio playback is closed".to_owned()) + } + + fn state(&self) -> Arc { + Arc::clone(&self.state) + } + + fn finish_input(&mut self) { + self.writer.take(); + } + + fn set_gain(&self, gain: f32) -> NativeResult<()> { + if !gain.is_finite() { + return Err("Audio playback gain must be finite".to_owned()); + } + self.state.set_gain(gain.max(0.0)); + Ok(()) + } + + /// Stop playback immediately and release the default speaker. + pub(crate) fn stop(&mut self) -> NativeResult<()> { + self.writer.take(); + self.state.mark_stopped(); + let Some(device) = self.device.take() else { + return Ok(()); + }; + device + .stop() + .map_err(|error| format!("Failed to stop speaker playback: {error}")) + } +} + +impl Drop for PlaybackStream { + fn drop(&mut self) { + let _ = self.stop(); + } +} + +fn audio_config(device_type: DeviceType, sample_rate: u32) -> DeviceConfig { + let mut config = DeviceConfig::new(device_type); + config.set_sample_rate(sample_rate); + config.set_period_size_in_milliseconds(AUDIO_PERIOD_MS); + config.set_performance_profile(PerformanceProfile::LowLatency); + config +} + +fn validate_sample_rate(sample_rate: u32) -> NativeResult<()> { + if sample_rate == 0 { + return Err("Audio sample rate must be greater than zero".to_owned()); + } + Ok(()) +} + +fn fill_playback( + rx: &flume::Receiver>, + current: &mut Vec, + cursor: &mut usize, + output: &mut [f32], + state: &PlaybackState, + empty_callbacks: &mut usize, +) { + output.fill(0.0); + if state.stopped.load(Ordering::Acquire) { + return; + } + + let gain = state.gain(); + let mut output_offset = 0; + while output_offset < output.len() { + if *cursor == current.len() { + match rx.try_recv() { + Ok(next) => { + *current = next; + *cursor = 0; + *empty_callbacks = 0; + }, + Err(TryRecvError::Empty) => { + *empty_callbacks = 0; + break; + }, + Err(TryRecvError::Disconnected) => { + *empty_callbacks += 1; + if *empty_callbacks >= PLAYBACK_DRAIN_CALLBACKS { + state.mark_drained(); + } + break; + }, + } + } + + let count = (current.len() - *cursor).min(output.len() - output_offset); + let source = ¤t[*cursor..*cursor + count]; + let destination = &mut output[output_offset..output_offset + count]; + if gain == 1.0 { + destination.copy_from_slice(source); + } else { + for (destination, source) in destination.iter_mut().zip(source) { + *destination = *source * gain; + } + } + *cursor += count; + output_offset += count; + } +} + +/// Default-microphone capture converted to mono `f32` at the requested sample rate. +#[napi] +pub struct AudioCapture { + device: Mutex>, +} + +#[napi] +impl AudioCapture { + /// Open the default microphone and deliver low-latency mono PCM chunks. + #[napi(constructor)] + pub fn new( + sample_rate: u32, + #[napi(ts_arg_type = "(error: Error | null, samples: Float32Array) => void")] + on_audio: CaptureCallback, + ) -> Result { + validate_sample_rate(sample_rate).map_err(napi::Error::from_reason)?; + let mut config = audio_config(DeviceType::Capture, sample_rate); + config.capture_mut().set_format(Format::F32); + config.capture_mut().set_channels(AUDIO_CHANNELS); + let mut device = Device::new(None, &config) + .map_err(|error| napi::Error::from_reason(format!("Failed to open the default microphone: {error}")))?; + device.set_data_callback(move |_device, _output, input| { + if input.sample_count() == 0 { + return; + } + on_audio.call( + Ok(Float32Array::new(input.as_samples::().to_vec())), + ThreadsafeFunctionCallMode::NonBlocking, + ); + }); + device.start().map_err(|error| { + napi::Error::from_reason(format!("Failed to start microphone capture: {error}")) + })?; + Ok(Self { device: Mutex::new(Some(device)) }) + } + + /// Stop capture immediately and release the microphone. + #[napi] + pub fn stop(&self) -> Result<()> { + let device = self.device.lock().take(); + let Some(device) = device else { + return Ok(()); + }; + device + .stop() + .map_err(|error| napi::Error::from_reason(format!("Failed to stop microphone capture: {error}"))) + } +} + +impl Drop for AudioCapture { + fn drop(&mut self) { + if let Some(device) = self.device.get_mut().take() { + let _ = device.stop(); + } + } +} + +/// Gapless mono `f32` playback through the default speaker. +#[napi] +pub struct AudioPlayback { + stream: Mutex>, + state: Arc, +} + +#[napi] +impl AudioPlayback { + /// Open the default speaker at the requested logical sample rate. + #[napi(constructor)] + pub fn new(sample_rate: u32) -> Result { + let stream = PlaybackStream::start(sample_rate).map_err(napi::Error::from_reason)?; + let state = stream.state(); + Ok(Self { stream: Mutex::new(Some(stream)), state }) + } + + /// Queue mono floating-point PCM in playback order. + #[napi] + pub fn write(&self, samples: Float32Array) -> Result<()> { + let stream = self.stream.lock(); + let stream = stream + .as_ref() + .ok_or_else(|| napi::Error::from_reason("Native audio playback is closed"))?; + stream + .writer() + .and_then(|writer| writer.write(&samples)) + .map_err(napi::Error::from_reason) + } + + /// Scale audio at render time so gain changes affect already queued samples. + #[napi] + pub fn set_gain(&self, gain: f64) -> Result<()> { + let stream = self.stream.lock(); + let stream = stream + .as_ref() + .ok_or_else(|| napi::Error::from_reason("Native audio playback is closed"))?; + stream.set_gain(gain as f32).map_err(napi::Error::from_reason) + } + + /// Close input, wait until queued samples reach the speaker, then release it. + #[napi] + pub async fn end(&self) -> Result<()> { + { + let mut stream = self.stream.lock(); + let Some(stream) = stream.as_mut() else { + return Ok(()); + }; + stream.finish_input(); + } + self.state.wait_for_drain().await; + let stream = self.stream.lock().take(); + if let Some(mut stream) = stream { + stream.stop().map_err(napi::Error::from_reason)?; + } + Ok(()) + } + + /// Stop immediately and discard all queued samples. + #[napi] + pub fn stop(&self) -> Result<()> { + let stream = self.stream.lock().take(); + if let Some(mut stream) = stream { + stream.stop().map_err(napi::Error::from_reason)?; + } + Ok(()) + } +} + +impl Drop for AudioPlayback { + fn drop(&mut self) { + if let Some(mut stream) = self.stream.get_mut().take() { + let _ = stream.stop(); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn playback_preserves_chunk_order_and_applies_render_gain() { + let state = PlaybackState::new(); + state.set_gain(0.5); + let (tx, rx) = flume::unbounded(); + tx.send(vec![1.0, -1.0]).expect("receiver is live"); + tx.send(vec![0.5, -0.5]).expect("receiver is live"); + drop(tx); + let mut current = Vec::new(); + let mut cursor = 0; + let mut empty_callbacks = 0; + let mut output = [9.0; 5]; + + fill_playback( + &rx, + &mut current, + &mut cursor, + &mut output, + &state, + &mut empty_callbacks, + ); + + assert_eq!(output, [0.5, -0.5, 0.25, -0.25, 0.0]); + assert!(!state.drained.load(Ordering::Acquire)); + let mut silence = [1.0; 2]; + fill_playback( + &rx, + &mut current, + &mut cursor, + &mut silence, + &state, + &mut empty_callbacks, + ); + assert_eq!(silence, [0.0, 0.0]); + assert!(state.drained.load(Ordering::Acquire)); + } +} diff --git a/crates/pi-natives/src/devicecheck.rs b/crates/pi-natives/src/devicecheck.rs new file mode 100644 index 000000000..8cd4261f6 --- /dev/null +++ b/crates/pi-natives/src/devicecheck.rs @@ -0,0 +1,348 @@ +//! Apple DeviceCheck token generation (`DCDevice.generateToken`). +//! +//! Reimplements the flow the ChatGPT desktop app's `devicecheck.node` addon +//! uses to mint attestation tokens: resolve `DCDevice.currentDevice`, check +//! `isSupported`, then call `generateTokenWithCompletionHandler:` and wait up +//! to one second for the completion block, reporting the base64-encoded token +//! or the failure reason. +//! +//! Uses raw Objective-C runtime FFI and a hand-built block literal — no +//! `objc2`/`block2` dependency. +//! +//! # Platform +//! - **macOS**: Full implementation via `DeviceCheck.framework`. +//! - **Other**: Returns `supported: false` without touching the network. + +use napi_derive::napi; + +use crate::task; + +/// Outcome of a single `DCDevice.generateToken` request. +#[napi(object)] +pub struct DeviceCheckTokenResult { + /// Whether `DCDevice.isSupported` reported attestation support. + pub supported: bool, + /// Base64-encoded DeviceCheck token; present only when generation succeeded. + pub token_base64: Option, + /// Human-readable failure reason when no token was produced. + pub error: Option, + /// Wall-clock time spent in the native call, in milliseconds. + pub latency_ms: f64, +} + +/// Generate an Apple DeviceCheck attestation token. +/// +/// Resolves with the token (or the error reason) after at most a 1-second +/// wait, matching the upstream `devicecheck.node` addon contract. +#[napi] +pub fn device_check_generate_token() -> task::Promise { + task::blocking("devicecheck.generate_token", (), move |_| Ok(platform::generate_token())) +} + +// --------------------------------------------------------------------------- +// macOS implementation +// --------------------------------------------------------------------------- + +#[cfg(target_os = "macos")] +mod platform { + use std::{ + ffi::{CStr, c_char, c_void}, + panic::{AssertUnwindSafe, catch_unwind}, + ptr, + sync::mpsc::{self, SyncSender}, + time::{Duration, Instant}, + }; + + use super::DeviceCheckTokenResult; + + /// How long to wait for the DeviceCheck completion handler before giving + /// up, matching the timeout in the upstream `devicecheck.node` addon. + const TOKEN_TIMEOUT: Duration = Duration::from_secs(1); + + type Id = *mut c_void; + type Sel = *mut c_void; + + // `objc_msgSend` is typed per call signature via `#[link_name]` aliases, + // the standard idiom for raw ObjC messaging without an objc crate. + #[allow( + clashing_extern_declarations, + reason = "objc_msgSend is an assembly trampoline that forwards to the method IMP; each alias types the same symbol for a distinct call signature" + )] + #[link(name = "objc")] + unsafe extern "C" { + fn objc_getClass(name: *const c_char) -> Id; + fn sel_registerName(name: *const c_char) -> Sel; + fn objc_retain(obj: Id) -> Id; + fn objc_release(obj: Id); + fn objc_autoreleasePoolPush() -> *mut c_void; + fn objc_autoreleasePoolPop(pool: *mut c_void); + + #[link_name = "objc_msgSend"] + fn msg_send_noarg(receiver: Id, selector: Sel) -> Id; + #[link_name = "objc_msgSend"] + fn msg_send_bool(receiver: Id, selector: Sel) -> u8; + #[link_name = "objc_msgSend"] + fn msg_send_u64(receiver: Id, selector: Sel, options: u64) -> Id; + #[link_name = "objc_msgSend"] + fn msg_send_block(receiver: Id, selector: Sel, block: *const c_void); + } + + // Linking DeviceCheck.framework registers `DCDevice` with the ObjC + // runtime when the addon image loads. + #[link(name = "DeviceCheck", kind = "framework")] + unsafe extern "C" {} + + unsafe extern "C" { + /// Stack-block class from libsystem_blocks; used as the literal's isa. + static _NSConcreteStackBlock: *const c_void; + } + + /// Outcome delivered once from the completion block to the waiting worker. + enum Completion { + Token(String), + Error(String), + } + + /// Objective-C block ABI: the 32-byte literal header followed by the + /// captured context (a raw pointer to the channel sender). + #[repr(C)] + struct CompletionBlock { + isa: *const c_void, + flags: i32, + reserved: i32, + invoke: unsafe extern "C" fn(*mut CompletionBlock, Id, Id), + descriptor: *const CompletionBlockDescriptor, + sender: *const SyncSender, + } + + /// `Block_descriptor_1` followed immediately by `Block_descriptor_3`. + /// No `Block_descriptor_2` (copy/dispose helpers) is emitted because the + /// captured sender pointer is plain-old-data and needs no retain/release. + #[repr(C)] + struct CompletionBlockDescriptor { + reserved: usize, + size: usize, + signature: *const c_char, + } + + /// `BLOCK_HAS_SIGNATURE` — the only flag needed for a POD stack block. + const BLOCK_HAS_SIGNATURE: i32 = 1 << 30; + + /// Type encoding for `void (^)(NSData *token, NSError *error)`: + /// void return, 24 bytes of arguments (block at 0, token at 8, error at 16). + const BLOCK_SIGNATURE: &CStr = c"v24@?0@8@16"; + + /// Immutable, process-lifetime data; the raw signature pointer is never + /// mutated, so shared access from the ObjC runtime is race-free. + unsafe impl Sync for CompletionBlockDescriptor {} + + static COMPLETION_DESCRIPTOR: CompletionBlockDescriptor = CompletionBlockDescriptor { + reserved: 0, + size: size_of::(), + signature: BLOCK_SIGNATURE.as_ptr(), + }; + + /// Resolve a selector by name; `sel_registerName` is idempotent and cheap. + /// + /// # Safety + /// The returned selector is valid for the lifetime of the process. + unsafe fn selector(name: &CStr) -> Sel { + // SAFETY: `name` is a valid null-terminated C string. + unsafe { sel_registerName(name.as_ptr()) } + } + + /// Copy a C string owned by an autoreleased `NSString` into a Rust `String`. + /// + /// # Safety + /// `ptr` must be null or point to a valid null-terminated UTF-8 string that + /// outlives the call. + unsafe fn copy_c_string(ptr: *const c_char) -> String { + if ptr.is_null() { + return String::new(); + } + // SAFETY: upheld by the caller; `CStr::from_ptr` only reads. + unsafe { CStr::from_ptr(ptr) } + .to_string_lossy() + .into_owned() + } + + /// Read the UTF-8 payload of an `NSString` into a Rust `String`. + /// + /// # Safety + /// `string` must be a live `NSString` for the duration of the call. + unsafe fn ns_string(string: Id) -> String { + // SAFETY: `string` is a live NSString; the returned pointer stays valid + // until the enclosing autorelease pool drains. + unsafe { copy_c_string(msg_send_noarg(string, selector(c"UTF8String")).cast()) } + } + + /// Completion block body. Runs on DeviceCheck's XPC reply queue, which is + /// why the result travels over a channel instead of a return value. + /// + /// # Safety + /// Called by the Objective-C runtime with a valid block literal; `token` + /// and `error` are live `NSData`/`NSError` objects (or null) for the + /// duration of the call. + unsafe extern "C" fn completion_invoke(block: *mut CompletionBlock, token: Id, error: Id) { + let completion = catch_unwind(AssertUnwindSafe(|| { + if !token.is_null() { + // SAFETY: `token` is a live NSData for the duration of the callback. + let encoded = + unsafe { msg_send_u64(token, selector(c"base64EncodedStringWithOptions:"), 0) }; + if encoded.is_null() { + return Completion::Error("DeviceCheck returned no token".to_owned()); + } + // SAFETY: `encoded` is a live NSString. + return Completion::Token(unsafe { ns_string(encoded) }); + } + if !error.is_null() { + // SAFETY: `error` is a live NSError for the duration of the callback. + let description = unsafe { msg_send_noarg(error, selector(c"localizedDescription")) }; + if description.is_null() { + return Completion::Error("DeviceCheck token request failed".to_owned()); + } + // SAFETY: `description` is a live NSString. + return Completion::Error(unsafe { ns_string(description) }); + } + Completion::Error("DeviceCheck returned no token".to_owned()) + })); + let completion = match completion { + Ok(completion) => completion, + Err(payload) => { + // Never let a panic escape into the ObjC runtime; mirror the + // bounded-leak disposal used by `task::Blocking` instead of + // dropping a potentially panicking payload type here. + std::mem::forget(payload); + Completion::Error("DeviceCheck completion panicked".to_owned()) + }, + }; + // SAFETY: the owner keeps the sender alive until the block has fired + // (and leaks it on timeout), so the captured pointer is always valid. + // `try_send` never blocks the XPC queue, even if the runtime were to + // invoke the block more than once. + unsafe { + _ = (*(*block).sender).try_send(completion); + } + } + + /// Build the result for a supported device by driving + /// `generateTokenWithCompletionHandler:` and waiting on the channel. + /// + /// # Safety + /// `device` must be a live, retained `DCDevice` instance. + unsafe fn run_token_request(device: Id) -> DeviceCheckTokenResult { + let (sender, receiver) = mpsc::sync_channel::(1); + let sender = Box::into_raw(Box::new(sender)); + let block = CompletionBlock { + isa: ptr::addr_of!(_NSConcreteStackBlock).cast::(), + flags: BLOCK_HAS_SIGNATURE, + reserved: 0, + invoke: completion_invoke, + descriptor: &raw const COMPLETION_DESCRIPTOR, + sender, + }; + // SAFETY: `device` is a live DCDevice and `block` follows the block ABI; + // the runtime copies the literal, so the stack frame may die after the call. + unsafe { + msg_send_block( + device, + selector(c"generateTokenWithCompletionHandler:"), + (&raw const block).cast(), + ) + }; + + let mut result = DeviceCheckTokenResult { + supported: true, + token_base64: None, + error: None, + latency_ms: 0.0, + }; + match receiver.recv_timeout(TOKEN_TIMEOUT) { + Ok(Completion::Token(token)) => { + result.token_base64 = Some(token); + // SAFETY: the block has fired and will not fire again, so the + // sender is unreachable from the runtime and can be reclaimed. + drop(unsafe { Box::from_raw(sender) }); + }, + Ok(Completion::Error(message)) => { + result.error = Some(message); + // SAFETY: same as above — the single-shot block already fired. + drop(unsafe { Box::from_raw(sender) }); + }, + Err(_) => { + // Timeout (or a vanished sender): the block may still fire on + // the XPC queue, so deliberately leak the sender to keep the + // captured pointer valid. Bounded to one leak per timeout. + result.error = Some("timed out waiting for DeviceCheck token".to_owned()); + }, + } + result + } + + fn generate_token_inner() -> DeviceCheckTokenResult { + let mut result = DeviceCheckTokenResult { + supported: false, + token_base64: None, + error: None, + latency_ms: 0.0, + }; + // SAFETY: `c"DCDevice"` is a valid null-terminated class name. + let class = unsafe { objc_getClass(c"DCDevice".as_ptr()) }; + if class.is_null() { + result.error = Some("DeviceCheck framework unavailable".to_owned()); + return result; + } + // SAFETY: `class` is a registered ObjC class; `currentDevice` is a + // documented DCDevice class method returning an autoreleased instance. + let device = unsafe { msg_send_noarg(class, selector(c"currentDevice")) }; + if device.is_null() { + result.error = Some("DeviceCheck currentDevice unavailable".to_owned()); + return result; + } + // SAFETY: `device` is a live object; retain balances the release below. + let device = unsafe { objc_retain(device) }; + // SAFETY: `device` is a live DCDevice; `isSupported` returns BOOL. + let supported = unsafe { msg_send_bool(device, selector(c"isSupported")) } != 0; + if supported { + // SAFETY: `device` is live and retained for the duration of the call. + return unsafe { + let mut token_result = run_token_request(device); + objc_release(device); + token_result.supported = true; + token_result + }; + } + // SAFETY: balances the retain above. + unsafe { objc_release(device) }; + result + } + + pub fn generate_token() -> DeviceCheckTokenResult { + let start = Instant::now(); + // SAFETY: pool push/pop are balanced within this scope. + let pool = unsafe { objc_autoreleasePoolPush() }; + let mut result = generate_token_inner(); + result.latency_ms = start.elapsed().as_secs_f64() * 1000.0; + // SAFETY: balances the push above. + unsafe { objc_autoreleasePoolPop(pool) }; + result + } +} + +// --------------------------------------------------------------------------- +// Non-macOS stub +// --------------------------------------------------------------------------- + +#[cfg(not(target_os = "macos"))] +mod platform { + use super::DeviceCheckTokenResult; + + pub fn generate_token() -> DeviceCheckTokenResult { + DeviceCheckTokenResult { + supported: false, + token_base64: None, + error: None, + latency_ms: 0.0, + } + } +} diff --git a/crates/pi-natives/src/lib.rs b/crates/pi-natives/src/lib.rs index faa86fc6b..50d6d719a 100644 --- a/crates/pi-natives/src/lib.rs +++ b/crates/pi-natives/src/lib.rs @@ -23,6 +23,7 @@ #![feature(alloc_error_hook)] pub mod appearance; +pub mod audio; pub mod ast; pub mod block; pub mod clipboard; @@ -34,6 +35,7 @@ pub mod desktop; /// pure conversion helpers stay unit-testable without a live X server. #[cfg(any(target_os = "linux", test))] pub mod desktop_x11; +pub mod devicecheck; pub mod diff; pub mod fd; pub mod glob; @@ -43,6 +45,7 @@ pub mod highlight; pub mod html; pub mod iofs; pub mod keys; +pub mod live; pub mod sixel; pub mod snapcompact; pub use pi_ast::language; diff --git a/crates/pi-natives/src/live.rs b/crates/pi-natives/src/live.rs new file mode 100644 index 000000000..2beeaa39f --- /dev/null +++ b/crates/pi-natives/src/live.rs @@ -0,0 +1,769 @@ +//! Native WebRTC media transport for Codex live conversations. +//! +//! The TypeScript host owns authenticated signaling and the sideband protocol; +//! this module owns the realtime WebRTC peer, Opus media, and speaker playback. + +use std::{ + sync::{ + Arc, Weak, + atomic::{AtomicBool, AtomicUsize, Ordering}, + }, + time::Duration, +}; + +use bytes::Bytes; +use napi::{ + bindgen_prelude::{Float32Array, Result}, + threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode, UnknownReturnValue}, +}; +use napi_derive::napi; +use opus::{Application, Channels, Decoder, Encoder}; +use parking_lot::Mutex; +use tokio::{sync::watch, task::JoinHandle}; +use crate::audio::{PlaybackStream, PlaybackWriter}; +use webrtc::{ + api::{ + APIBuilder, + interceptor_registry::register_default_interceptors, + media_engine::{MIME_TYPE_OPUS, MediaEngine}, + }, + data_channel::{RTCDataChannel, data_channel_message::DataChannelMessage}, + interceptor::registry::Registry, + media::Sample, + peer_connection::{ + RTCPeerConnection, + configuration::RTCConfiguration, + peer_connection_state::RTCPeerConnectionState, + sdp::session_description::RTCSessionDescription, + }, + rtp_transceiver::{ + rtp_codec::{RTCRtpCodecCapability, RTCRtpCodecParameters, RTPCodecType}, + rtp_sender::RTCRtpSender, + }, + track::{ + track_local::{TrackLocal, track_local_static_sample::TrackLocalStaticSample}, + track_remote::TrackRemote, + }, +}; + +const DATA_CHANNEL_LABEL: &str = "oai-events"; +const INPUT_SAMPLE_RATE: u32 = 16_000; +const INPUT_FRAME_SAMPLES: usize = 320; +const INPUT_FRAME_DURATION: Duration = Duration::from_millis(20); +const MAX_ENCODED_OPUS_BYTES: usize = 1_275; +const MAX_QUEUED_INPUT_SAMPLES: usize = 32_000; +const OUTPUT_SAMPLE_RATE: u32 = 48_000; +const MAX_DECODED_OPUS_SAMPLES: usize = 5_760; +const OUTPUT_LEVEL_SAMPLES: usize = 2_400; +const OUTPUT_FRAME_SAMPLES: usize = 960; +const DEFAULT_OPEN_TIMEOUT_MS: u32 = 20_000; +const DISCONNECT_GRACE: Duration = Duration::from_secs(2); +const CLOSE_TASK_TIMEOUT: Duration = Duration::from_secs(1); + +const OPUS_CAPABILITY: RTCRtpCodecCapability = RTCRtpCodecCapability { + mime_type: String::new(), + clock_rate: OUTPUT_SAMPLE_RATE, + channels: 2, + sdp_fmtp_line: String::new(), + rtcp_feedback: Vec::new(), +}; + +type StringCallback = ThreadsafeFunction; +type LevelCallback = ThreadsafeFunction; +type NativeResult = std::result::Result; + +#[derive(Clone, Debug)] +enum PeerSignal { + Connecting, + Open, + Failed(String), + Closed, +} + +enum InputCommand { + Audio(Vec), + Muted(bool), + Close, +} + +struct LiveCallbacks { + event: StringCallback, + level: LevelCallback, + failure: StringCallback, +} + +struct LiveResources { + peer: Arc, + data_channel: Arc, + input_tx: flume::Sender, + input_task: JoinHandle<()>, + rtcp_task: JoinHandle<()>, + playback: PlaybackStream, +} + +struct LivePeerCore { + callbacks: LiveCallbacks, + resources: Mutex>, + signal_tx: watch::Sender, + started: AtomicBool, + closing: AtomicBool, + muted: AtomicBool, + failure_reported: AtomicBool, + queued_samples: AtomicUsize, +} + +impl LivePeerCore { + fn new(callbacks: LiveCallbacks) -> Self { + let (signal_tx, _) = watch::channel(PeerSignal::Connecting); + Self { + callbacks, + resources: Mutex::new(None), + signal_tx, + started: AtomicBool::new(false), + closing: AtomicBool::new(false), + muted: AtomicBool::new(false), + failure_reported: AtomicBool::new(false), + queued_samples: AtomicUsize::new(0), + } + } + + async fn create_offer(self: &Arc) -> NativeResult { + if self.started.swap(true, Ordering::AcqRel) { + return Err("Native live WebRTC peer has already started".to_owned()); + } + if self.closing.load(Ordering::Acquire) { + return Err("Native live WebRTC peer is closed".to_owned()); + } + + let playback = PlaybackStream::start(OUTPUT_SAMPLE_RATE)?; + let playback_tx = playback.writer()?; + let mut media_engine = MediaEngine::default(); + let capability = opus_capability(); + media_engine + .register_codec( + RTCRtpCodecParameters { + capability: capability.clone(), + payload_type: 111, + ..Default::default() + }, + RTPCodecType::Audio, + ) + .map_err(|error| format!("Failed to register the live Opus codec: {error}"))?; + let registry = register_default_interceptors(Registry::new(), &mut media_engine) + .map_err(|error| format!("Failed to configure live WebRTC interceptors: {error}"))?; + let api = APIBuilder::new() + .with_media_engine(media_engine) + .with_interceptor_registry(registry) + .build(); + let peer = Arc::new( + api.new_peer_connection(RTCConfiguration::default()) + .await + .map_err(|error| format!("Failed to create the live WebRTC peer: {error}"))?, + ); + + let track = Arc::new(TrackLocalStaticSample::new( + capability, + "audio".to_owned(), + "omp-live".to_owned(), + )); + let sender = match peer + .add_track(Arc::clone(&track) as Arc) + .await + { + Ok(sender) => sender, + Err(error) => { + let _ = peer.close().await; + return Err(format!("Failed to add the live audio track: {error}")); + }, + }; + + install_peer_callbacks(&peer, Arc::downgrade(self), playback_tx); + let data_channel = match peer.create_data_channel(DATA_CHANNEL_LABEL, None).await { + Ok(channel) => channel, + Err(error) => { + let _ = peer.close().await; + return Err(format!("Failed to create the live data channel: {error}")); + }, + }; + install_data_channel_callbacks(&data_channel, Arc::downgrade(self)); + + let offer = match peer.create_offer(None).await { + Ok(offer) => offer, + Err(error) => { + let _ = peer.close().await; + return Err(format!("Failed to create the live SDP offer: {error}")); + }, + }; + if let Err(error) = peer.set_local_description(offer.clone()).await { + let _ = peer.close().await; + return Err(format!("Failed to install the live SDP offer: {error}")); + } + if self.closing.load(Ordering::Acquire) { + let _ = peer.close().await; + return Err("Native live WebRTC peer was closed while starting".to_owned()); + } + + let (input_tx, input_rx) = flume::unbounded(); + let input_task = tokio::spawn(run_input_audio( + track, + input_rx, + Arc::downgrade(self), + )); + let rtcp_task = tokio::spawn(drain_rtcp(sender)); + let resources = LiveResources { + peer, + data_channel, + input_tx, + input_task, + rtcp_task, + playback, + }; + *self.resources.lock() = Some(resources); + Ok(offer.sdp) + } + + async fn accept_answer(&self, sdp: String) -> NativeResult<()> { + let peer = self + .resources + .lock() + .as_ref() + .map(|resources| Arc::clone(&resources.peer)) + .ok_or_else(|| "Native live WebRTC peer has not started".to_owned())?; + let answer = RTCSessionDescription::answer(sdp) + .map_err(|error| format!("Codex returned an invalid live SDP answer: {error}"))?; + peer.set_remote_description(answer) + .await + .map_err(|error| format!("Failed to install the live SDP answer: {error}")) + } + + async fn wait_for_open(&self, timeout_ms: u32) -> NativeResult<()> { + let mut signal_rx = self.signal_tx.subscribe(); + let wait = async { + loop { + match signal_rx.borrow().clone() { + PeerSignal::Open => return Ok(()), + PeerSignal::Failed(message) => return Err(message), + PeerSignal::Closed => return Err("Native live WebRTC peer closed before opening".to_owned()), + PeerSignal::Connecting => {}, + } + signal_rx + .changed() + .await + .map_err(|_| "Native live WebRTC peer stopped before opening".to_owned())?; + } + }; + tokio::time::timeout(Duration::from_millis(u64::from(timeout_ms)), wait) + .await + .map_err(|_| "Timed out waiting for the live data channel to open".to_owned())? + } + + fn push_audio(&self, samples: &[f32]) -> NativeResult<()> { + if samples.is_empty() || self.muted.load(Ordering::Acquire) { + return Ok(()); + } + let input_tx = self + .resources + .lock() + .as_ref() + .map(|resources| resources.input_tx.clone()) + .ok_or_else(|| "Native live WebRTC peer has not started".to_owned())?; + let sample_count = samples.len().min(MAX_QUEUED_INPUT_SAMPLES); + let retained = &samples[samples.len() - sample_count..]; + let queued = self.queued_samples.fetch_add(sample_count, Ordering::AcqRel); + if queued.saturating_add(sample_count) > MAX_QUEUED_INPUT_SAMPLES { + self.queued_samples.fetch_sub(sample_count, Ordering::AcqRel); + return Ok(()); + } + if input_tx.send(InputCommand::Audio(retained.to_vec())).is_err() { + self.queued_samples.fetch_sub(sample_count, Ordering::AcqRel); + return Err("Native live audio input is closed".to_owned()); + } + Ok(()) + } + + fn set_muted(&self, muted: bool) -> NativeResult<()> { + self.muted.store(muted, Ordering::Release); + let input_tx = self + .resources + .lock() + .as_ref() + .map(|resources| resources.input_tx.clone()); + if let Some(input_tx) = input_tx { + input_tx + .send(InputCommand::Muted(muted)) + .map_err(|_| "Native live audio input is closed".to_owned())?; + } + Ok(()) + } + + fn report_event(&self, payload: String) { + self.callbacks + .event + .call(Ok(payload), ThreadsafeFunctionCallMode::NonBlocking); + } + + fn report_level(&self, level: f64) { + self.callbacks + .level + .call(Ok(level.clamp(0.0, 1.0)), ThreadsafeFunctionCallMode::NonBlocking); + } + + fn mark_open(&self) { + if !self.closing.load(Ordering::Acquire) { + self.signal_tx.send_replace(PeerSignal::Open); + } + } + + fn report_failure(&self, message: String) { + if self.closing.load(Ordering::Acquire) || self.failure_reported.swap(true, Ordering::AcqRel) { + return; + } + self.signal_tx.send_replace(PeerSignal::Failed(message.clone())); + self.callbacks + .failure + .call(Ok(message), ThreadsafeFunctionCallMode::NonBlocking); + } + + async fn close(&self) { + if self.closing.swap(true, Ordering::AcqRel) { + let mut signal_rx = self.signal_tx.subscribe(); + while !matches!(*signal_rx.borrow(), PeerSignal::Closed) { + if signal_rx.changed().await.is_err() { + break; + } + } + return; + } + + let resources = self.resources.lock().take(); + if let Some(mut resources) = resources { + let _ = resources.input_tx.send(InputCommand::Close); + let _ = resources.peer.close().await; + let _ = resources.playback.stop(); + let _ = tokio::time::timeout(CLOSE_TASK_TIMEOUT, resources.input_task).await; + resources.rtcp_task.abort(); + let _ = resources.rtcp_task.await; + drop(resources.data_channel); + } + self.queued_samples.store(0, Ordering::Release); + self.signal_tx.send_replace(PeerSignal::Closed); + } +} + +/// WebRTC peer that accepts 16 kHz mono PCM and renders remote Opus audio. +#[napi] +pub struct LiveWebRtcPeer { + inner: Arc, +} + +#[napi] +impl LiveWebRtcPeer { + /// Create an idle peer and register its event, output-level, and failure callbacks. + #[napi(constructor)] + pub fn new( + #[napi(ts_arg_type = "(error: Error | null, payload: string) => void")] + on_event: StringCallback, + #[napi(ts_arg_type = "(error: Error | null, level: number) => void")] + on_level: LevelCallback, + #[napi(ts_arg_type = "(error: Error | null, message: string) => void")] + on_failure: StringCallback, + ) -> Self { + Self { + inner: Arc::new(LivePeerCore::new(LiveCallbacks { + event: on_event, + level: on_level, + failure: on_failure, + })), + } + } + + /// Start the native media peer and return its SDP offer. + #[napi] + pub async fn create_offer(&self) -> Result { + self.inner.create_offer().await.map_err(napi::Error::from_reason) + } + + /// Apply the remote SDP answer returned by Codex signaling. + #[napi] + pub async fn accept_answer(&self, sdp: String) -> Result<()> { + self.inner.accept_answer(sdp).await.map_err(napi::Error::from_reason) + } + + /// Wait until the `oai-events` data channel is open. + #[napi] + pub async fn wait_for_open(&self, timeout_ms: Option) -> Result<()> { + self.inner + .wait_for_open(timeout_ms.unwrap_or(DEFAULT_OPEN_TIMEOUT_MS)) + .await + .map_err(napi::Error::from_reason) + } + + /// Queue 16 kHz mono floating-point PCM for Opus transmission. + #[napi] + pub fn push_audio(&self, samples: Float32Array) -> Result<()> { + self.inner.push_audio(&samples).map_err(napi::Error::from_reason) + } + + /// Enable or disable microphone transmission, discarding partial muted frames. + #[napi] + pub fn set_muted(&self, muted: bool) -> Result<()> { + self.inner.set_muted(muted).map_err(napi::Error::from_reason) + } + + /// Close media, the data channel, the peer connection, and speaker playback. + #[napi] + pub async fn close(&self) { + self.inner.close().await; + } +} + +impl Drop for LiveWebRtcPeer { + fn drop(&mut self) { + if self.inner.closing.load(Ordering::Acquire) { + return; + } + let inner = Arc::clone(&self.inner); + if let Ok(runtime) = tokio::runtime::Handle::try_current() { + runtime.spawn(async move { + inner.close().await; + }); + } + } +} + +fn opus_capability() -> RTCRtpCodecCapability { + RTCRtpCodecCapability { + mime_type: MIME_TYPE_OPUS.to_owned(), + clock_rate: OPUS_CAPABILITY.clock_rate, + channels: OPUS_CAPABILITY.channels, + sdp_fmtp_line: "minptime=10;useinbandfec=1".to_owned(), + rtcp_feedback: Vec::new(), + } +} + + +fn install_peer_callbacks( + peer: &Arc, + core: Weak, + playback_tx: PlaybackWriter, +) { + let output_sender = Arc::new(Mutex::new(Some(playback_tx))); + let output_sender_for_track = Arc::clone(&output_sender); + let core_for_track = core.clone(); + peer.on_track(Box::new(move |track, _receiver, _transceiver| { + let output_sender = output_sender_for_track.lock().take(); + let core = core_for_track.clone(); + Box::pin(async move { + if track.kind() != RTPCodecType::Audio { + return; + } + let Some(output_sender) = output_sender else { + if let Some(core) = core.upgrade() { + core.report_failure("Codex live returned more than one remote audio track".to_owned()); + } + return; + }; + tokio::spawn(receive_output_audio(track, output_sender, core)); + }) + })); + + let peer_for_state = Arc::downgrade(peer); + peer.on_peer_connection_state_change(Box::new(move |state| { + let core = core.clone(); + let peer = peer_for_state.clone(); + Box::pin(async move { + let Some(core) = core.upgrade() else { + return; + }; + match state { + RTCPeerConnectionState::Failed => { + core.report_failure("Live WebRTC peer connection failed".to_owned()); + }, + RTCPeerConnectionState::Closed => { + if !core.closing.load(Ordering::Acquire) { + core.report_failure("Live WebRTC peer connection closed unexpectedly".to_owned()); + } + }, + RTCPeerConnectionState::Disconnected => { + tokio::time::sleep(DISCONNECT_GRACE).await; + if peer + .upgrade() + .is_some_and(|peer| peer.connection_state() == RTCPeerConnectionState::Disconnected) + { + core.report_failure("Live WebRTC peer connection disconnected".to_owned()); + } + }, + _ => {}, + } + }) + })); +} + +fn install_data_channel_callbacks(data_channel: &Arc, core: Weak) { + let core_for_open = core.clone(); + data_channel.on_open(Box::new(move || { + let core = core_for_open.clone(); + Box::pin(async move { + if let Some(core) = core.upgrade() { + core.mark_open(); + } + }) + })); + + let core_for_message = core.clone(); + data_channel.on_message(Box::new(move |message: DataChannelMessage| { + let core = core_for_message.clone(); + Box::pin(async move { + if !message.is_string { + return; + } + if let (Some(core), Ok(payload)) = (core.upgrade(), String::from_utf8(message.data.to_vec())) { + core.report_event(payload); + } + }) + })); + + let core_for_close = core.clone(); + data_channel.on_close(Box::new(move || { + let core = core_for_close.clone(); + Box::pin(async move { + if let Some(core) = core.upgrade() { + core.report_failure("Live data channel closed unexpectedly".to_owned()); + } + }) + })); + + data_channel.on_error(Box::new(move |error| { + let core = core.clone(); + Box::pin(async move { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Live data channel failed: {error}")); + } + }) + })); +} + +async fn run_input_audio( + track: Arc, + input_rx: flume::Receiver, + core: Weak, +) { + let mut encoder = match Encoder::new(INPUT_SAMPLE_RATE, Channels::Mono, Application::Voip) { + Ok(encoder) => encoder, + Err(error) => { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to initialize the live Opus encoder: {error}")); + } + return; + }, + }; + if let Err(error) = encoder.set_inband_fec(true) { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to configure the live Opus encoder: {error}")); + } + return; + } + + let mut muted = false; + let mut pending = Vec::with_capacity(INPUT_FRAME_SAMPLES * 2); + let mut encoded = [0u8; MAX_ENCODED_OPUS_BYTES]; + let mut ticker = tokio::time::interval(INPUT_FRAME_DURATION); + ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Burst); + ticker.tick().await; + loop { + tokio::select! { + biased; + command = input_rx.recv_async() => { + let Ok(command) = command else { + break; + }; + match command { + InputCommand::Audio(samples) => { + if let Some(core) = core.upgrade() { + core.queued_samples.fetch_sub(samples.len(), Ordering::AcqRel); + } + if muted { + continue; + } + if samples.len() >= MAX_QUEUED_INPUT_SAMPLES { + pending.clear(); + pending.extend_from_slice(&samples[samples.len() - MAX_QUEUED_INPUT_SAMPLES..]); + continue; + } + let overflow = pending + .len() + .saturating_add(samples.len()) + .saturating_sub(MAX_QUEUED_INPUT_SAMPLES); + if overflow > 0 { + pending.drain(..overflow); + } + pending.extend_from_slice(&samples); + }, + InputCommand::Muted(next_muted) => { + muted = next_muted; + pending.clear(); + }, + InputCommand::Close => break, + } + }, + _ = ticker.tick() => { + let mut frame = [0.0f32; INPUT_FRAME_SAMPLES]; + if !muted { + let consumed = pending.len().min(INPUT_FRAME_SAMPLES); + frame[..consumed].copy_from_slice(&pending[..consumed]); + if consumed > 0 { + pending.copy_within(consumed.., 0); + pending.truncate(pending.len() - consumed); + } + } + let encoded_len = match encoder.encode_float(&frame, &mut encoded) { + Ok(encoded_len) => encoded_len, + Err(error) => { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to encode live microphone audio: {error}")); + } + return; + }, + }; + let sample = Sample { + data: Bytes::copy_from_slice(&encoded[..encoded_len]), + duration: INPUT_FRAME_DURATION, + ..Default::default() + }; + if let Err(error) = track.write_sample(&sample).await { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to send live microphone audio: {error}")); + } + return; + } + }, + } + } +} + +async fn drain_rtcp(sender: Arc) { + while sender.read_rtcp().await.is_ok() {} +} + +async fn receive_output_audio( + track: Arc, + playback_tx: PlaybackWriter, + core: Weak, +) { + if !track.codec().capability.mime_type.eq_ignore_ascii_case(MIME_TYPE_OPUS) { + if let Some(core) = core.upgrade() { + core.report_failure(format!( + "Codex live negotiated unsupported audio codec {}", + track.codec().capability.mime_type + )); + } + return; + } + let mut decoder = match Decoder::new(OUTPUT_SAMPLE_RATE, Channels::Mono) { + Ok(decoder) => decoder, + Err(error) => { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to initialize the live Opus decoder: {error}")); + } + return; + }, + }; + let mut decoded = [0.0f32; MAX_DECODED_OPUS_SAMPLES]; + let mut expected_sequence: Option = None; + let mut level = OutputLevel::default(); + + loop { + let packet = match track.read_rtp().await { + Ok((packet, _attributes)) => packet, + Err(error) => { + if let Some(core) = core.upgrade() + && !core.closing.load(Ordering::Acquire) + { + core.report_failure(format!("Live remote audio track failed: {error}")); + } + return; + }, + }; + let sequence = packet.header.sequence_number; + if let Some(expected) = expected_sequence { + let gap = sequence.wrapping_sub(expected); + if gap >= u16::MAX / 2 { + continue; + } + if gap > 0 { + for _ in 1..gap.min(5) { + if let Ok(samples) = + decoder.decode_float(&[], &mut decoded[..OUTPUT_FRAME_SAMPLES], false) + { + if !write_output(&playback_tx, &decoded[..samples], &core) { + return; + } + level.observe(&decoded[..samples], &core); + } + } + if let Ok(samples) = decoder.decode_float(&packet.payload, &mut decoded, true) { + if !write_output(&playback_tx, &decoded[..samples], &core) { + return; + } + level.observe(&decoded[..samples], &core); + } + } + } + expected_sequence = Some(sequence.wrapping_add(1)); + match decoder.decode_float(&packet.payload, &mut decoded, false) { + Ok(samples) => { + if !write_output(&playback_tx, &decoded[..samples], &core) { + return; + } + level.observe(&decoded[..samples], &core); + }, + Err(error) => { + if let Some(core) = core.upgrade() { + core.report_failure(format!("Failed to decode live speaker audio: {error}")); + } + return; + }, + } + } +} + +fn write_output(playback_tx: &PlaybackWriter, samples: &[f32], core: &Weak) -> bool { + match playback_tx.write(samples) { + Ok(()) => true, + Err(error) => { + if let Some(core) = core.upgrade() + && !core.closing.load(Ordering::Acquire) + { + core.report_failure(format!("Live speaker playback failed: {error}")); + } + false + }, + } +} + +#[derive(Default)] +struct OutputLevel { + sum_squares: f64, + samples: usize, +} + +impl OutputLevel { + fn observe(&mut self, decoded: &[f32], core: &Weak) { + let mut offset = 0; + while offset < decoded.len() { + let take = (OUTPUT_LEVEL_SAMPLES - self.samples).min(decoded.len() - offset); + for &sample in &decoded[offset..offset + take] { + self.sum_squares += f64::from(sample) * f64::from(sample); + } + self.samples += take; + offset += take; + if self.samples == OUTPUT_LEVEL_SAMPLES { + if let Some(core) = core.upgrade() { + core.report_level((self.sum_squares / self.samples as f64).sqrt()); + } + self.sum_squares = 0.0; + self.samples = 0; + } + } + } +} diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 4b9e2a659..183f28f04 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -12,7 +12,9 @@ ### Fixed +- Fixed live-call attestation depending on the ChatGPT desktop app being installed: `generateLiveAttestation` now mints DeviceCheck tokens in-process through the `@oh-my-pi/pi-natives` `deviceCheckGenerateToken` binding instead of probing `/Applications` for the app's `devicecheck.node` addon, so the `x-oai-attestation` header works on hosts without the desktop app and drops the `createRequire` addon probing. - Fixed `xd://` device execution failures rendering as `write` errors instead of using the mounted tool's own error renderer. +- Fixed custom tools without bespoke renderers losing the default state-tinted card when mounted under `xd://`; dispatched calls now keep their label, arguments, status, output preview, and expansion affordance instead of dumping a bare result line into the transcript. - Fixed the clipboard image-paste keybind mangling copied URL text into a bogus path error on macOS (e.g. `Image not found at /https/::i.can.ac:CE4Ek3.png` for a copied `https://i.can.ac/CE4Ek3.png`). AppleScript's `the clipboard as «class furl»` coerces plain *text* into a file URL by treating the string as an HFS path (`:`↔`/` swap), so `readMacFileUrlsFromClipboard` returned a garbage path that dead-ended in `handleImagePathPaste` instead of falling through to the text paste. The script now bails early via `clipboard info for «class furl»` unless the pasteboard actually carries a `public.file-url` representation, so URL/text clipboards paste as text. - Fixed spilled tool-output artifact descriptors leaking on error/abort paths. `OutputSink.dump()` was the only path that closed the spill `Bun.FileSink`, but the bash and Python executors re-throw on failure and their `finally` blocks never closed the sink, so a large-output command that errored leaked the artifact descriptor until an unrelated read (e.g. a `SKILL.md` load) hit `EMFILE`. `OutputSink` now exposes an idempotent `dispose()` that closes the sink exactly once, wired into every executor's `finally` ([#6463](https://github.com/can1357/oh-my-pi/issues/6463)). - Fixed the first submitted prompt stalling while the local tiny-title worker started: the interactive submit handler now paints the pending user row before starting title generation, and startup prewarms an idle, unref'd worker so the first submit reuses a live subprocess instead of paying spawn latency ahead of the first frame ([#6462](https://github.com/can1357/oh-my-pi/issues/6462)). diff --git a/packages/coding-agent/src/live/attestation.ts b/packages/coding-agent/src/live/attestation.ts new file mode 100644 index 000000000..000490bb6 --- /dev/null +++ b/packages/coding-agent/src/live/attestation.ts @@ -0,0 +1,91 @@ +import { deviceCheckGenerateToken } from "@oh-my-pi/pi-natives"; + +const CHATGPT_BUNDLE_ID = "com.openai.codex"; +const APP_SESSION_ID = crypto.randomUUID(); + +type DeviceCheckResult = { + supported: boolean; + tokenBase64?: string | null; + latencyMs?: number | null; +}; + +function cborHeader(major: number, value: number): Buffer { + if (!Number.isSafeInteger(value) || value < 0) throw new Error(`Invalid CBOR length: ${value}`); + if (value < 24) return Buffer.from([major + value]); + if (value <= 0xff) return Buffer.from([major + 24, value]); + if (value <= 0xffff) { + const output = Buffer.allocUnsafe(3); + output[0] = major + 25; + output.writeUInt16BE(value, 1); + return output; + } + if (value <= 0xffff_ffff) { + const output = Buffer.allocUnsafe(5); + output[0] = major + 26; + output.writeUInt32BE(value, 1); + return output; + } + throw new Error(`CBOR length is too large: ${value}`); +} + +function cborUnsigned(value: number): Buffer { + return cborHeader(0, value); +} + +function cborText(value: string): Buffer { + const text = Buffer.from(value, "utf8"); + return Buffer.concat([cborHeader(96, text.length), text]); +} + +function cborMap(entries: ReadonlyArray): Buffer { + const values: Buffer[] = [cborHeader(160, entries.length)]; + for (const [key, value] of entries) values.push(key, value); + return Buffer.concat(values); +} + +function attestationSignals(): Buffer { + const resolved = Intl.DateTimeFormat().resolvedOptions(); + const locale = (resolved.locale || "unknown").slice(0, 64); + const timezone = (resolved.timeZone || "unknown").slice(0, 64); + const preferredLanguages = Buffer.concat([cborHeader(128, 1), cborText(locale)]); + return cborMap([ + [cborUnsigned(0), cborUnsigned(1)], + [cborUnsigned(1), preferredLanguages], + [cborUnsigned(2), cborText(locale)], + [cborUnsigned(3), cborText(timezone)], + [cborUnsigned(4), cborUnsigned(0)], + [cborUnsigned(5), cborUnsigned(1)], + [cborUnsigned(6), cborText(APP_SESSION_ID.slice(0, 128))], + ]); +} + +function buildClientAttestation(result: DeviceCheckResult): string { + const entries: Array = []; + if (result.supported && result.tokenBase64) { + entries.push([cborText("token"), cborText(result.tokenBase64)]); + } else { + entries.push([cborText("error_code"), cborUnsigned(result.supported ? 4 : 3)]); + } + entries.push([cborText("bundle_id"), cborText(CHATGPT_BUNDLE_ID)]); + const signals = attestationSignals(); + entries.push([cborText("f"), Buffer.concat([cborHeader(64, signals.length), signals])]); + if (result.latencyMs !== undefined && result.latencyMs !== null) { + const latency = Buffer.allocUnsafe(9); + latency[0] = 0xfb; + latency.writeDoubleBE(result.latencyMs, 1); + entries.push([cborText("t"), latency]); + } + return `v1.${cborMap(entries).toString("base64url")}`; +} + +/** Generates the Codex Desktop DeviceCheck header accepted by the private live-call endpoints. */ +export async function generateLiveAttestation(): Promise { + if (process.platform !== "darwin" || process.arch !== "arm64") return undefined; + let result: DeviceCheckResult; + try { + result = await deviceCheckGenerateToken(); + } catch { + return undefined; + } + return JSON.stringify({ v: 1, s: 0, t: buildClientAttestation(result) }); +} diff --git a/packages/coding-agent/src/live/browser-runtime.txt b/packages/coding-agent/src/live/browser-runtime.txt index 6347fd4b4..2a0fd4db0 100644 --- a/packages/coding-agent/src/live/browser-runtime.txt +++ b/packages/coding-agent/src/live/browser-runtime.txt @@ -34,22 +34,6 @@ Promise.resolve(window.__ompLiveFailure(message)).catch(() => {}); } - function waitForIceGathering() { - if (peer.iceGatheringState === "complete") return Promise.resolve(); - return new Promise((resolve, reject) => { - const timeout = setTimeout(() => { - peer.removeEventListener("icegatheringstatechange", onChange); - reject(new Error("Timed out gathering WebRTC ICE candidates")); - }, 10_000); - function onChange() { - if (peer.iceGatheringState !== "complete") return; - clearTimeout(timeout); - peer.removeEventListener("icegatheringstatechange", onChange); - resolve(); - } - peer.addEventListener("icegatheringstatechange", onChange); - }); - } function decodePcm(base64) { const encoded = atob(base64); @@ -154,10 +138,9 @@ }; const offer = await peer.createOffer(); + if (!offer.sdp) throw new Error("Browser produced an empty WebRTC offer"); await peer.setLocalDescription(offer); - await waitForIceGathering(); - if (!peer.localDescription || !peer.localDescription.sdp) throw new Error("Browser produced an empty WebRTC offer"); - return peer.localDescription.sdp; + return offer.sdp; } async function acceptAnswer(sdp) { diff --git a/packages/coding-agent/src/live/controller.ts b/packages/coding-agent/src/live/controller.ts index a5639a597..d1c3a2a23 100644 --- a/packages/coding-agent/src/live/controller.ts +++ b/packages/coding-agent/src/live/controller.ts @@ -17,7 +17,7 @@ import { import { CodexLiveTransport } from "./transport"; import type { LivePhase, LiveTranscript } from "./visualizer"; -const DEFAULT_VOICE = "marin"; +const DEFAULT_VOICE = "sol"; const OUTPUT_ACTIVE_LEVEL = 0.015; const MIN_BARGE_IN_LEVEL = 0.04; const OUTPUT_ECHO_RATIO = 0.65; @@ -42,7 +42,7 @@ export interface LiveSessionControllerOptions { callbacks: LiveSessionCallbacks; /** Extracts visible assistant text using the caller's normal UI rules. */ extractAssistantText(message: AssistantMessage): string; - /** Realtime output voice, defaulting to marin. */ + /** Realtime output voice, defaulting to sol. */ voice?: string; } diff --git a/packages/coding-agent/src/live/protocol.test.ts b/packages/coding-agent/src/live/protocol.test.ts index 9b28bf957..2511de9a0 100644 --- a/packages/coding-agent/src/live/protocol.test.ts +++ b/packages/coding-agent/src/live/protocol.test.ts @@ -95,12 +95,12 @@ describe("Frameless Bidi server events", () => { }); describe("Frameless Bidi client payloads", () => { - test("builds the exact multipart call session JSON", () => { + test("builds the exact live call session JSON", () => { const payload = buildLiveSessionPayload("Be concise.", "marin"); - expect(LIVE_MODEL).toBe("gpt-live-1-boulder-alpha"); + expect(LIVE_MODEL).toBe("gpt-live-1-codex"); expect(JSON.stringify(payload)).toBe( - '{"model":"gpt-live-1-boulder-alpha","instructions":"Be concise.","audio":{"output":{"voice":"marin"}},"delegation":{"type":"client"}}', + '{"model":"gpt-live-1-codex","instructions":"Be concise.","audio":{"output":{"voice":"marin"}},"delegation":{"type":"client"}}', ); }); diff --git a/packages/coding-agent/src/live/protocol.ts b/packages/coding-agent/src/live/protocol.ts index 84b53f21f..a4164b628 100644 --- a/packages/coding-agent/src/live/protocol.ts +++ b/packages/coding-agent/src/live/protocol.ts @@ -1,5 +1,5 @@ /** Frameless Bidi model used by Codex Desktop live calls. */ -export const LIVE_MODEL: "gpt-live-1-boulder-alpha" = "gpt-live-1-boulder-alpha"; +export const LIVE_MODEL: "gpt-live-1-codex" = "gpt-live-1-codex"; /** Maximum UTF-8 payload size accepted by each context append. */ export const CONTEXT_CHUNK_BYTES = 500; diff --git a/packages/coding-agent/src/live/transport.ts b/packages/coding-agent/src/live/transport.ts index 5cb20d41b..9f0ffae3d 100644 --- a/packages/coding-agent/src/live/transport.ts +++ b/packages/coding-agent/src/live/transport.ts @@ -4,13 +4,10 @@ import { CODEX_BASE_URL, CODEX_CLIENT_VERSION, getCodexAccountId, - OPENAI_HEADER_VALUES, OPENAI_HEADERS, } from "@oh-my-pi/pi-catalog/wire/codex"; -import type { Browser, HTTPRequest, Page } from "puppeteer-core"; -import { launchHeadlessBrowser } from "../tools/browser/launch"; -import audioWorkletSource from "./audio-worklet.txt" with { type: "text" }; -import browserRuntimeSource from "./browser-runtime.txt" with { type: "text" }; +import { LiveWebRtcPeer } from "@oh-my-pi/pi-natives"; +import { generateLiveAttestation } from "./attestation"; import { buildLiveSessionPayload, type LiveClientMessage, @@ -20,37 +17,20 @@ import { const SIGNALING_URL = `${CODEX_BASE_URL}/codex/realtime/calls?intent=quicksilver&architecture=avas`; const MAX_ERROR_BODY_LENGTH = 2_048; -const MAX_HOST_AUDIO_SAMPLES = 32_000; const SIDEBAND_CONNECT_ATTEMPTS = 5; const SIDEBAND_CONNECT_TIMEOUT_MS = 15_000; const LIVE_PROVIDER = "openai-codex"; -const LIVE_CALL_ID_PATTERN = /^rtc_[\da-f]{8}-[\da-f]{4}-[\da-f]{4}-[\da-f]{4}-[\da-f]{12}$/i; +const LIVE_ORIGINATOR = "Codex Desktop"; +const LIVE_CALL_ID_PATTERN = /^rtc_[\w-]+$/; type Lifecycle = "idle" | "connecting" | "connected" | "closing" | "closed"; -type BrowserLiveApi = { - start(workletSource: string): Promise; - acceptAnswer(sdp: string): Promise; - waitForOpen(): Promise; - send(payload: string): void; - pushAudio(payload: string): void; - setMuted(muted: boolean): void; - close(): Promise; -}; - -declare global { - var ompCodexLive: BrowserLiveApi; -} - -interface QueuedAudio { - payload: string; - sampleCount: number; -} interface LiveSignalingResult { answer: string; callId: string; access: OAuthAccess; + attestation: string | undefined; } class LiveSignalingError extends Error { @@ -81,7 +61,7 @@ export interface LiveTransportOptions { signal?: AbortSignal; } -/** Extracts the server-assigned `rtc_` call ID from a signaling Location header. */ +/** Extracts the server-assigned `rtc_*` call ID from a signaling Location header. */ export function parseLiveCallId(location: string | null): string | undefined { if (!location) return undefined; return location @@ -92,23 +72,30 @@ export function parseLiveCallId(location: string | null): string | undefined { /** Builds the Frameless Bidi sideband WebSocket URL for an accepted Codex call. */ export function buildLiveSidebandUrl(callId: string): string { - const url = new URL(`${CODEX_BASE_URL}/codex/${encodeURIComponent(callId)}`); - url.protocol = url.protocol === "http:" ? "ws:" : "wss:"; + const url = new URL(`https://api.openai.com/v1/live/${encodeURIComponent(callId)}`); + url.protocol = "wss:"; return url.toString(); } -function liveSessionHeaders(access: OAuthAccess, sessionId: string): Record { +function liveSessionHeaders( + access: OAuthAccess, + sessionId: string, + realtimeSessionId: string, + attestation: string | undefined, +): Record { const headers: Record = { Authorization: `Bearer ${access.accessToken}`, "OpenAI-Alpha": "quicksilver=v2", - "x-session-id": sessionId, - [OPENAI_HEADERS.ORIGINATOR]: OPENAI_HEADER_VALUES.ORIGINATOR_CODEX, + "User-Agent": `Codex Desktop/${CODEX_CLIENT_VERSION}`, + "x-session-id": realtimeSessionId, + [OPENAI_HEADERS.ORIGINATOR]: LIVE_ORIGINATOR, [OPENAI_HEADERS.VERSION]: CODEX_CLIENT_VERSION, [OPENAI_HEADERS.SCOPED_SESSION_ID]: sessionId, [OPENAI_HEADERS.THREAD_ID]: sessionId, }; const accountId = access.accountId ?? getCodexAccountId(access.accessToken); if (accountId) headers[OPENAI_HEADERS.ACCOUNT_ID] = accountId; + if (attestation) headers["x-oai-attestation"] = attestation; return headers; } @@ -128,23 +115,17 @@ function abortReason(signal: AbortSignal | undefined): Error { return new DOMException("Live connection aborted", "AbortError"); } -function encodePcm(samples: Float32Array): string { - return Buffer.from(samples.buffer, samples.byteOffset, samples.byteLength).toString("base64"); -} -/** Headless-Chromium WebRTC transport for a Codex Frameless Bidi live session. */ +/** Native WebRTC transport for a Codex Frameless Bidi live session. */ export class CodexLiveTransport { readonly #options: LiveTransportOptions; - #browser: Browser | undefined; - #page: Page | undefined; + #peer: LiveWebRtcPeer | undefined; + readonly #realtimeSessionId = crypto.randomUUID(); #sideband: Bun.WebSocket | undefined; #state: Lifecycle = "idle"; #connectPromise: Promise | undefined; #closePromise: Promise | undefined; #sendTail: Promise = Promise.resolve(); - #audioQueue: QueuedAudio[] = []; - #queuedAudioSamples = 0; - #audioPumpRunning = false; #muted = false; #unexpectedFailureReported = false; readonly #abortListener: () => void; @@ -174,49 +155,42 @@ export class CodexLiveTransport { } async #connect(): Promise { - const browser = await launchHeadlessBrowser({ - headless: true, - args: ["--autoplay-policy=no-user-gesture-required"], - ignoreDefaultArgs: ["--mute-audio"], - }); - this.#browser = browser; - if (this.#state !== "connecting") throw abortReason(this.#options.signal); - const page = await browser.newPage(); - this.#page = page; - await page.setRequestInterception(true); - const serveBlankPage = (request: HTTPRequest): void => { - void request - .respond({ status: 200, contentType: "text/html", body: "Codex Live" }) - .catch(() => {}); - }; - page.on("request", serveBlankPage); - try { - await page.goto("http://localhost/", { waitUntil: "domcontentloaded" }); - } finally { - page.off("request", serveBlankPage); - await page.setRequestInterception(false); - } - await page.exposeFunction("__ompLiveServerEvent", (payload: string) => this.#handleServerEvent(payload)); - await page.exposeFunction("__ompLiveOutputLevel", (level: number) => this.#handleOutputLevel(level)); - await page.exposeFunction("__ompLiveFailure", (message: string) => this.#handleBrowserFailure(message)); - await page.evaluate(source => Function(source)(), browserRuntimeSource); - const offer = await page.evaluate(source => globalThis.ompCodexLive.start(source), audioWorkletSource); + const peer = new LiveWebRtcPeer( + (error, payload) => { + if (error) { + this.#handlePeerFailure(error.message); + } else { + this.#handleServerEvent(payload); + } + }, + (error, level) => { + if (error) { + this.#handlePeerFailure(error.message); + } else { + this.#handleOutputLevel(level); + } + }, + (error, message) => this.#handlePeerFailure(error?.message ?? message), + ); + this.#peer = peer; + const offer = await peer.createOffer(); if (this.#state !== "connecting") throw abortReason(this.#options.signal); const signaling = await this.#signal(offer); - await page.evaluate(sdp => globalThis.ompCodexLive.acceptAnswer(sdp), signaling.answer); - await page.evaluate(muted => globalThis.ompCodexLive.setMuted(muted), this.#muted); - await page.evaluate(() => globalThis.ompCodexLive.waitForOpen()); + await peer.acceptAnswer(signaling.answer); + peer.setMuted(this.#muted); + await peer.waitForOpen(); if (this.#state !== "connecting") throw abortReason(this.#options.signal); - await this.#connectSideband(signaling.callId, signaling.access); + await this.#connectSideband(signaling.callId, signaling.access, signaling.attestation); if (this.#state !== "connecting") throw abortReason(this.#options.signal); this.#state = "connected"; } async #signal(offer: string): Promise { + const attestation = await generateLiveAttestation(); return await withOAuthAccess( this.#options.authStorage, LIVE_PROVIDER, - access => this.#signalWithAccess(offer, access), + access => this.#signalWithAccess(offer, access, attestation), { sessionId: this.#options.sessionId, signal: this.#options.signal, @@ -226,10 +200,14 @@ export class CodexLiveTransport { ); } - async #signalWithAccess(offer: string, access: OAuthAccess): Promise { + async #signalWithAccess( + offer: string, + access: OAuthAccess, + attestation: string | undefined, + ): Promise { const headers = new Headers({ - ...liveSessionHeaders(access, this.#options.sessionId), - Accept: "application/sdp", + ...liveSessionHeaders(access, this.#options.sessionId, this.#realtimeSessionId, attestation), + Accept: "*/*", "Content-Type": "application/json", }); const fetchImpl = wrapFetchForProxy(fetch, LIVE_PROVIDER); @@ -247,20 +225,24 @@ export class CodexLiveTransport { const detail = boundedErrorBody(responseBody, response.statusText); throw new LiveSignalingError(response.status, `Codex live signaling failed (${response.status}): ${detail}`); } - const answer = responseBody.trim(); - if (!answer) throw new LiveSignalingError(response.status, "Codex live signaling returned an empty SDP answer"); + const answer = responseBody; + if (!answer.trim()) throw new LiveSignalingError(response.status, "Codex live signaling returned an empty SDP answer"); const callId = parseLiveCallId(response.headers.get("location")); if (!callId) { throw new LiveSignalingError(response.status, "Codex live signaling returned no valid call ID"); } - return { answer, callId, access }; + return { answer, callId, access, attestation }; } - async #connectSideband(callId: string, access: OAuthAccess): Promise { + async #connectSideband( + callId: string, + access: OAuthAccess, + attestation: string | undefined, + ): Promise { let failure = new Error("Codex live sideband connection failed"); for (let attempt = 0; attempt < SIDEBAND_CONNECT_ATTEMPTS; attempt++) { try { - await this.#openSideband(callId, access); + await this.#openSideband(callId, access, attestation); return; } catch (cause) { failure = cause instanceof Error ? cause : new Error(String(cause)); @@ -271,10 +253,14 @@ export class CodexLiveTransport { throw failure; } - async #openSideband(callId: string, access: OAuthAccess): Promise { + async #openSideband( + callId: string, + access: OAuthAccess, + attestation: string | undefined, + ): Promise { const url = buildLiveSidebandUrl(callId); const options = { - headers: liveSessionHeaders(access, this.#options.sessionId), + headers: liveSessionHeaders(access, this.#options.sessionId, this.#realtimeSessionId, attestation), proxy: getProxyForProvider(LIVE_PROVIDER), } satisfies Bun.WebSocketOptions; const socket: Bun.WebSocket = Reflect.construct(WebSocket, [url, options]); @@ -377,7 +363,7 @@ export class CodexLiveTransport { } catch {} } - #handleBrowserFailure(message: string): void { + #handlePeerFailure(message: string): void { this.#reportFailure(message); } @@ -405,52 +391,19 @@ export class CodexLiveTransport { return operation; } - /** Queue 16 kHz mono Float32 PCM for continuous browser-side resampling and playback. */ + /** Queue 16 kHz mono Float32 PCM for native Opus transmission. */ pushAudio(samples: Float32Array): void { if (this.#state !== "connected" || this.#muted || samples.length === 0) return; - const retained = - samples.length > MAX_HOST_AUDIO_SAMPLES ? samples.subarray(samples.length - MAX_HOST_AUDIO_SAMPLES) : samples; - const queued = { payload: encodePcm(retained), sampleCount: retained.length }; - this.#audioQueue.push(queued); - this.#queuedAudioSamples += queued.sampleCount; - while (this.#queuedAudioSamples > MAX_HOST_AUDIO_SAMPLES && this.#audioQueue.length > 1) { - const stale = this.#audioQueue.shift(); - if (stale) this.#queuedAudioSamples -= stale.sampleCount; - } - if (!this.#audioPumpRunning) void this.#pumpAudio(); + this.#peer?.pushAudio(samples); } - async #pumpAudio(): Promise { - this.#audioPumpRunning = true; - try { - while (this.#state === "connected" && !this.#muted) { - const queued = this.#audioQueue.shift(); - if (!queued) break; - this.#queuedAudioSamples -= queued.sampleCount; - const page = this.#page; - if (!page) break; - await page.evaluate(payload => globalThis.ompCodexLive.pushAudio(payload), queued.payload); - } - } catch { - this.#audioQueue.length = 0; - this.#queuedAudioSamples = 0; - } finally { - this.#audioPumpRunning = false; - if (this.#audioQueue.length > 0 && this.#state === "connected" && !this.#muted) void this.#pumpAudio(); - } - } - - /** Enable or disable the browser audio source and discard queued input when muted. */ + /** Enable or disable the native audio source and discard partial input when muted. */ async setMuted(muted: boolean): Promise { this.#muted = muted; - this.#audioQueue.length = 0; - this.#queuedAudioSamples = 0; - const page = this.#page; - if (!page || this.#state !== "connected") return; - await page.evaluate(value => globalThis.ompCodexLive.setMuted(value), muted); + if (this.#state === "connected") this.#peer?.setMuted(muted); } - /** Stop audio, WebRTC, the page, and Chromium. Safe to call repeatedly. */ + /** Stop sideband signaling and the native WebRTC media peer. Safe to call repeatedly. */ close(): Promise { if (this.#closePromise) return this.#closePromise; this.#state = "closing"; @@ -461,28 +414,16 @@ export class CodexLiveTransport { async #close(): Promise { this.#options.signal?.removeEventListener("abort", this.#abortListener); - this.#audioQueue.length = 0; - this.#queuedAudioSamples = 0; const sideband = this.#sideband; - const page = this.#page; - const browser = this.#browser; + const peer = this.#peer; this.#sideband = undefined; - this.#page = undefined; - this.#browser = undefined; + this.#peer = undefined; if (sideband && (sideband.readyState === WebSocket.OPEN || sideband.readyState === WebSocket.CONNECTING)) { sideband.close(1000, "done"); } - if (page) { + if (peer) { try { - await page.evaluate(() => globalThis.ompCodexLive?.close()); - } catch {} - try { - await page.close(); - } catch {} - } - if (browser) { - try { - await browser.close(); + await peer.close(); } catch {} } this.#state = "closed"; diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index 804004ea2..56166baf0 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -4,6 +4,8 @@ ### Added +- Added LiveWebRtcPeer class for WebRTC live streaming with offer/answer negotiation and audio push capabilities +- Added a macOS `deviceCheckGenerateToken` export that generates Apple DeviceCheck attestation tokens natively: it drives `DCDevice.generateToken` through raw Objective-C runtime FFI with a hand-built completion block literal and a bounded one-second wait, resolving `{ supported, tokenBase64, error, latencyMs }` to mirror the ChatGPT desktop app's `devicecheck.node` addon contract. Non-macOS builds resolve `supported: false` without touching the network. - Added a genuine native desktop backend for computer use, bundled in the core addon on every published platform: macOS Quartz/CGEvent, Windows Win32/`SendInput`, and a pure-Rust Linux X11 backend (`x11rb` capture over the display socket, XTest input with keysym mapping) that links no GUI system libraries — so Linux x64/arm64, glibc and musl are all supported and headless hosts are unaffected. Wayland sessions work through XWayland. Execute batches enforce a 60-second native deadline (`DESKTOP_DEADLINE_EXCEEDED`) and never emit input after it expires; unsupported pure-Wayland capture and out-of-XTest-range or negative-origin coordinate layouts fail closed. ### Fixed diff --git a/packages/natives/native/index.d.ts b/packages/natives/native/index.d.ts index 6e7060b3e..60a3641c5 100644 --- a/packages/natives/native/index.d.ts +++ b/packages/natives/native/index.d.ts @@ -1,5 +1,27 @@ /* auto-generated by NAPI-RS */ /* eslint-disable */ +/** Default-microphone capture converted to mono `f32` at the requested sample rate. */ +export declare class AudioCapture { + /** Open the default microphone and deliver low-latency mono PCM chunks. */ + constructor(sampleRate: number, onAudio: (error: Error | null, samples: Float32Array) => void) + /** Stop capture immediately and release the microphone. */ + stop(): void +} + +/** Gapless mono `f32` playback through the default speaker. */ +export declare class AudioPlayback { + /** Open the default speaker at the requested logical sample rate. */ + constructor(sampleRate: number) + /** Queue mono floating-point PCM in playback order. */ + write(samples: Float32Array): void + /** Scale audio at render time so gain changes affect already queued samples. */ + setGain(gain: number): void + /** Close input, wait until queued samples reach the speaker, then release it. */ + end(): Promise + /** Stop immediately and discard all queued samples. */ + stop(): void +} + /** Persistent, serialized native desktop capture/input session. */ export declare class DesktopSession { constructor(options?: DesktopSessionOptions | undefined | null) @@ -16,6 +38,24 @@ export declare class DesktopSession { close(): Promise } +/** WebRTC peer that accepts 16 kHz mono PCM and renders remote Opus audio. */ +export declare class LiveWebRtcPeer { + /** Create an idle peer and register its event, output-level, and failure callbacks. */ + constructor(onEvent: (error: Error | null, payload: string) => void, onLevel: (error: Error | null, level: number) => void, onFailure: (error: Error | null, message: string) => void) + /** Start the native media peer and return its SDP offer. */ + createOffer(): Promise + /** Apply the remote SDP answer returned by Codex signaling. */ + acceptAnswer(sdp: string): Promise + /** Wait until the `oai-events` data channel is open. */ + waitForOpen(timeoutMs?: number | undefined | null): Promise + /** Queue 16 kHz mono floating-point PCM for Opus transmission. */ + pushAudio(samples: Float32Array): void + /** Enable or disable microphone transmission, discarding partial muted frames. */ + setMuted(muted: boolean): void + /** Close media, the data channel, the peer connection, and speaker playback. */ + close(): Promise +} + /** * Long-lived macOS appearance observer. * @@ -612,6 +652,26 @@ export interface DesktopSessionOptions { */ export declare function detectMacOSAppearance(): MacOSAppearance | null +/** + * Generate an Apple DeviceCheck attestation token. + * + * Resolves with the token (or the error reason) after at most a 1-second + * wait, matching the upstream `devicecheck.node` addon contract. + */ +export declare function deviceCheckGenerateToken(): Promise + +/** Outcome of a single `DCDevice.generateToken` request. */ +export interface DeviceCheckTokenResult { + /** Whether `DCDevice.isSupported` reported attestation support. */ + supported: boolean + /** Base64-encoded DeviceCheck token; present only when generation succeeded. */ + tokenBase64?: string + /** Human-readable failure reason when no token was produced. */ + error?: string + /** Wall-clock time spent in the native call, in milliseconds. */ + latencyMs: number +} + /** One jsdiff change object: a run of added, removed, or common tokens. */ export interface DiffChange { /** Joined token text for this run (lines keep their ` diff --git a/packages/natives/native/index.js b/packages/natives/native/index.js index 170460b7d..7f482448d 100644 --- a/packages/natives/native/index.js +++ b/packages/natives/native/index.js @@ -16,7 +16,10 @@ import { loadNative } from "./loader-state.js"; const nativeBindings = loadNative(); // --- generated native exports (do not edit) --- // classes +export const AudioCapture = nativeBindings.AudioCapture; +export const AudioPlayback = nativeBindings.AudioPlayback; export const DesktopSession = nativeBindings.DesktopSession; +export const LiveWebRtcPeer = nativeBindings.LiveWebRtcPeer; export const MacAppearanceObserver = nativeBindings.MacAppearanceObserver; export const MacOSPowerAssertion = nativeBindings.MacOSPowerAssertion; export const Process = nativeBindings.Process; @@ -34,6 +37,7 @@ export const copyToClipboard = nativeBindings.copyToClipboard; export const cosineSimilarityPairs = nativeBindings.cosineSimilarityPairs; export const countTokens = nativeBindings.countTokens; export const detectMacOSAppearance = nativeBindings.detectMacOSAppearance; +export const deviceCheckGenerateToken = nativeBindings.deviceCheckGenerateToken; export const diffLineRuns = nativeBindings.diffLineRuns; export const diffLines = nativeBindings.diffLines; export const diffWords = nativeBindings.diffWords; diff --git a/packages/natives/test/devicecheck.test.ts b/packages/natives/test/devicecheck.test.ts new file mode 100644 index 000000000..e8a24b55d --- /dev/null +++ b/packages/natives/test/devicecheck.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it } from "bun:test"; +import { deviceCheckGenerateToken } from "../native/index.js"; + +// A locally built addon can predate the DeviceCheck binding; skip instead of +// failing on stale artifacts, mirroring the desktop test gate. +const deviceCheckTest = typeof deviceCheckGenerateToken === "function" ? it : it.skip; + +describe("deviceCheckGenerateToken", () => { + deviceCheckTest("resolves the DCDevice.generateToken contract", async () => { + const result = await deviceCheckGenerateToken(); + expect(typeof result.supported).toBe("boolean"); + expect(typeof result.latencyMs).toBe("number"); + expect(result.latencyMs).toBeGreaterThanOrEqual(0); + if (process.platform !== "darwin") { + expect(result.supported).toBe(false); + return; + } + if (!result.supported) { + expect(result.tokenBase64 ?? null).toBeNull(); + return; + } + // A supported device yields exactly one of token or error reason; + // network/Apple-service failures surface as `error`, never a throw. + expect(typeof result.tokenBase64 === "string").toBe(result.tokenBase64 !== undefined); + if (result.tokenBase64 !== undefined) { + expect(result.tokenBase64.length).toBeGreaterThan(0); + expect(result.error ?? null).toBeNull(); + } else { + expect(typeof result.error).toBe("string"); + } + }); +});