Compare commits

..

No commits in common. "somethings" and "main" have entirely different histories.

37 changed files with 280 additions and 1260 deletions

370
Cargo.lock generated
View File

@ -1,6 +1,6 @@
# This file is automatically @generated by Cargo. # This file is automatically @generated by Cargo.
# It is not intended for manual editing. # It is not intended for manual editing.
version = 4 version = 3
[[package]] [[package]]
name = "addr2line" name = "addr2line"
@ -78,12 +78,8 @@ dependencies = [
name = "atem-connection-rs" name = "atem-connection-rs"
version = "0.1.0" version = "0.1.0"
dependencies = [ dependencies = [
"defmt",
"derive-getters", "derive-getters",
"derive-new", "derive-new",
"embassy-net",
"enumflags2",
"itertools",
"log", "log",
"tokio", "tokio",
"tokio-util", "tokio-util",
@ -140,12 +136,6 @@ version = "1.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a"
[[package]]
name = "byteorder"
version = "1.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b"
[[package]] [[package]]
name = "bytes" name = "bytes"
version = "1.5.0" version = "1.5.0"
@ -195,7 +185,7 @@ dependencies = [
"heck", "heck",
"proc-macro2", "proc-macro2",
"quote", "quote",
"syn 2.0.104", "syn 2.0.48",
] ]
[[package]] [[package]]
@ -237,44 +227,6 @@ version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "acbf1af155f9b9ef647e42cdc158db4b64a1b61f743629225fde6f3e0be2a7c7" checksum = "acbf1af155f9b9ef647e42cdc158db4b64a1b61f743629225fde6f3e0be2a7c7"
[[package]]
name = "critical-section"
version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b"
[[package]]
name = "defmt"
version = "1.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "548d977b6da32fa1d1fda2876453da1e7df63ad0304c8b3dae4dbe7b96f39b78"
dependencies = [
"bitflags",
"defmt-macros",
]
[[package]]
name = "defmt-macros"
version = "1.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3d4fc12a85bcf441cfe44344c4b72d58493178ce635338a3f3b78943aceb258e"
dependencies = [
"defmt-parser",
"proc-macro-error2",
"proc-macro2",
"quote",
"syn 2.0.104",
]
[[package]]
name = "defmt-parser"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "10d60334b3b2e7c9d91ef8150abfb6fa4c1c39ebbcf4a81c2e346aad939fee3e"
dependencies = [
"thiserror",
]
[[package]] [[package]]
name = "derive-getters" name = "derive-getters"
version = "0.2.0" version = "0.2.0"
@ -294,163 +246,7 @@ checksum = "d150dea618e920167e5973d70ae6ece4385b7164e0d799fe7c122dd0a5d912ad"
dependencies = [ dependencies = [
"proc-macro2", "proc-macro2",
"quote", "quote",
"syn 2.0.104", "syn 2.0.48",
]
[[package]]
name = "document-features"
version = "0.2.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "95249b50c6c185bee49034bcb378a49dc2b5dff0be90ff6616d31d64febab05d"
dependencies = [
"litrs",
]
[[package]]
name = "either"
version = "1.15.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "48c757948c5ede0e46177b7add2e67155f70e33c07fea8284df6576da70b3719"
[[package]]
name = "embassy-net"
version = "0.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "940c4b9fe5c1375b09a0c6722c0100d6b2ed46a717a34f632f26e8d7327c4383"
dependencies = [
"document-features",
"embassy-net-driver",
"embassy-sync",
"embassy-time",
"embedded-io-async",
"embedded-nal-async",
"heapless",
"managed",
"smoltcp",
]
[[package]]
name = "embassy-net-driver"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "524eb3c489760508f71360112bca70f6e53173e6fe48fc5f0efd0f5ab217751d"
[[package]]
name = "embassy-sync"
version = "0.6.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8d2c8cdff05a7a51ba0087489ea44b0b1d97a296ca6b1d6d1a33ea7423d34049"
dependencies = [
"cfg-if",
"critical-section",
"embedded-io-async",
"futures-sink",
"futures-util",
"heapless",
]
[[package]]
name = "embassy-time"
version = "0.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f820157f198ada183ad62e0a66f554c610cdcd1a9f27d4b316358103ced7a1f8"
dependencies = [
"cfg-if",
"critical-section",
"document-features",
"embassy-time-driver",
"embedded-hal 0.2.7",
"embedded-hal 1.0.0",
"embedded-hal-async",
"futures-util",
]
[[package]]
name = "embassy-time-driver"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8d45f5d833b6d98bd2aab0c2de70b18bfaa10faf661a1578fd8e5dfb15eb7eba"
dependencies = [
"document-features",
]
[[package]]
name = "embedded-hal"
version = "0.2.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "35949884794ad573cf46071e41c9b60efb0cb311e3ca01f7af807af1debc66ff"
dependencies = [
"nb 0.1.3",
"void",
]
[[package]]
name = "embedded-hal"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "361a90feb7004eca4019fb28352a9465666b24f840f5c3cddf0ff13920590b89"
[[package]]
name = "embedded-hal-async"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0c4c685bbef7fe13c3c6dd4da26841ed3980ef33e841cddfa15ce8a8fb3f1884"
dependencies = [
"embedded-hal 1.0.0",
]
[[package]]
name = "embedded-io"
version = "0.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d"
[[package]]
name = "embedded-io-async"
version = "0.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3ff09972d4073aa8c299395be75161d582e7629cd663171d62af73c8d50dba3f"
dependencies = [
"embedded-io",
]
[[package]]
name = "embedded-nal"
version = "0.9.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c56a28be191a992f28f178ec338a0bf02f63d7803244add736d026a471e6ed77"
dependencies = [
"nb 1.1.0",
]
[[package]]
name = "embedded-nal-async"
version = "0.8.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "76959917cd2b86f40a98c28dd5624eddd1fa69d746241c8257eac428d83cb211"
dependencies = [
"embedded-io-async",
"embedded-nal",
]
[[package]]
name = "enumflags2"
version = "0.7.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1027f7680c853e056ebcec683615fb6fbbc07dbaa13b4d5d9442b146ded4ecef"
dependencies = [
"enumflags2_derive",
]
[[package]]
name = "enumflags2_derive"
version = "0.7.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "67c78a4d8fdf9953a5c9d458f9efe940fd97a0cab0941c075a813ac594733827"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.104",
] ]
[[package]] [[package]]
@ -478,9 +274,9 @@ dependencies = [
[[package]] [[package]]
name = "futures-core" name = "futures-core"
version = "0.3.31" version = "0.3.30"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "05f29059c0c2090612e8d742178b0580d2dc940c837851ad723096f87af6663e" checksum = "dfc6580bb841c5a68e9ef15c77ccc837b40a7504914d52e47b8b0e9bbda25a1d"
[[package]] [[package]]
name = "futures-sink" name = "futures-sink"
@ -488,49 +284,12 @@ version = "0.3.30"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9fb8e00e87438d937621c1c6269e53f536c14d3fbd6a042bb24879e57d474fb5" checksum = "9fb8e00e87438d937621c1c6269e53f536c14d3fbd6a042bb24879e57d474fb5"
[[package]]
name = "futures-task"
version = "0.3.31"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f90f7dce0722e95104fcb095585910c0977252f286e354b5e3bd38902cd99988"
[[package]]
name = "futures-util"
version = "0.3.31"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9fa08315bb612088cc391249efdc3bc77536f16c91f6cf495e6fbe85b20a4a81"
dependencies = [
"futures-core",
"futures-task",
"pin-project-lite",
"pin-utils",
]
[[package]] [[package]]
name = "gimli" name = "gimli"
version = "0.25.0" version = "0.25.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f0a01e0497841a3b2db4f8afa483cce65f7e96a3498bd6c541734792aeac8fe7" checksum = "f0a01e0497841a3b2db4f8afa483cce65f7e96a3498bd6c541734792aeac8fe7"
[[package]]
name = "hash32"
version = "0.3.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "47d60b12902ba28e2730cd37e95b8c9223af2808df9e902d4df49588d1470606"
dependencies = [
"byteorder",
]
[[package]]
name = "heapless"
version = "0.8.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0bfb9eb618601c89945a70e254898da93b13be0388091d42117462b265bb3fad"
dependencies = [
"hash32",
"stable_deref_trait",
]
[[package]] [[package]]
name = "heck" name = "heck"
version = "0.4.1" version = "0.4.1"
@ -564,15 +323,6 @@ version = "0.3.3"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ce23b50ad8242c51a442f3ff322d56b02f08852c77e4c0b4d3fd684abc89c683" checksum = "ce23b50ad8242c51a442f3ff322d56b02f08852c77e4c0b4d3fd684abc89c683"
[[package]]
name = "itertools"
version = "0.14.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2b192c782037fadd9cfa75548310488aabdbf3d2da73885b31bd0abd03351285"
dependencies = [
"either",
]
[[package]] [[package]]
name = "lazy_static" name = "lazy_static"
version = "1.4.0" version = "1.4.0"
@ -585,12 +335,6 @@ version = "0.2.153"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9c198f91728a82281a64e1f4f9eeb25d82cb32a5de251c6bd1b5154d63a8e7bd" checksum = "9c198f91728a82281a64e1f4f9eeb25d82cb32a5de251c6bd1b5154d63a8e7bd"
[[package]]
name = "litrs"
version = "0.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b4ce301924b7887e9d637144fdade93f9dfff9b60981d4ac161db09720d39aa5"
[[package]] [[package]]
name = "lock_api" name = "lock_api"
version = "0.4.6" version = "0.4.6"
@ -609,12 +353,6 @@ dependencies = [
"cfg-if", "cfg-if",
] ]
[[package]]
name = "managed"
version = "0.8.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ca88d725a0a943b096803bd34e73a4437208b6077654cc4ecb2947a5f91618d"
[[package]] [[package]]
name = "memchr" name = "memchr"
version = "2.4.0" version = "2.4.0"
@ -642,21 +380,6 @@ dependencies = [
"windows-sys 0.48.0", "windows-sys 0.48.0",
] ]
[[package]]
name = "nb"
version = "0.1.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "801d31da0513b6ec5214e9bf433a77966320625a37860f910be265be6e18d06f"
dependencies = [
"nb 1.1.0",
]
[[package]]
name = "nb"
version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8d5439c4ad607c3c23abf66de8c8bf57ba8adcd1f129e699851a6e43935d339d"
[[package]] [[package]]
name = "num_cpus" name = "num_cpus"
version = "1.16.0" version = "1.16.0"
@ -717,39 +440,11 @@ version = "0.2.13"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8afb450f006bf6385ca15ef45d71d2288452bc3683ce2e2cacc0d18e4be60b58" checksum = "8afb450f006bf6385ca15ef45d71d2288452bc3683ce2e2cacc0d18e4be60b58"
[[package]]
name = "pin-utils"
version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184"
[[package]]
name = "proc-macro-error-attr2"
version = "2.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96de42df36bb9bba5542fe9f1a054b8cc87e172759a1868aa05c1f3acc89dfc5"
dependencies = [
"proc-macro2",
"quote",
]
[[package]]
name = "proc-macro-error2"
version = "2.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "11ec05c52be0a07b08061f7dd003e7d7092e0472bc731b4af7bb1ef876109802"
dependencies = [
"proc-macro-error-attr2",
"proc-macro2",
"quote",
"syn 2.0.104",
]
[[package]] [[package]]
name = "proc-macro2" name = "proc-macro2"
version = "1.0.95" version = "1.0.78"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "02b3e5e68a3a1a02aad3ec490a98007cbc13c37cbe84a3cd7b8e406d76e7f778" checksum = "e2422ad645d89c99f8f3e6b88a9fdeca7fabeac836b1002371c4367c8f984aae"
dependencies = [ dependencies = [
"unicode-ident", "unicode-ident",
] ]
@ -825,19 +520,6 @@ version = "1.13.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e6ecd384b10a64542d77071bd64bd7b231f4ed5940fba55e98c3de13824cf3d7" checksum = "e6ecd384b10a64542d77071bd64bd7b231f4ed5940fba55e98c3de13824cf3d7"
[[package]]
name = "smoltcp"
version = "0.12.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dad095989c1533c1c266d9b1e8d70a1329dd3723c3edac6d03bbd67e7bf6f4bb"
dependencies = [
"bitflags",
"byteorder",
"cfg-if",
"heapless",
"managed",
]
[[package]] [[package]]
name = "socket2" name = "socket2"
version = "0.5.6" version = "0.5.6"
@ -848,12 +530,6 @@ dependencies = [
"windows-sys 0.52.0", "windows-sys 0.52.0",
] ]
[[package]]
name = "stable_deref_trait"
version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a8f112729512f8e442d81f95a8a7ddf2b7c6b8a1a6f509a95864142b30cab2d3"
[[package]] [[package]]
name = "strsim" name = "strsim"
version = "0.10.0" version = "0.10.0"
@ -873,9 +549,9 @@ dependencies = [
[[package]] [[package]]
name = "syn" name = "syn"
version = "2.0.104" version = "2.0.48"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "17b6f705963418cdb9927482fa304bc562ece2fdd4f616084c50b7023b435a40" checksum = "0f3531638e407dfc0814761abb7c00a5b54992b849452a0646b7f65c9f770f3f"
dependencies = [ dependencies = [
"proc-macro2", "proc-macro2",
"quote", "quote",
@ -891,26 +567,6 @@ dependencies = [
"winapi-util", "winapi-util",
] ]
[[package]]
name = "thiserror"
version = "2.0.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "567b8a2dae586314f7be2a752ec7474332959c6460e02bde30d702a66d488708"
dependencies = [
"thiserror-impl",
]
[[package]]
name = "thiserror-impl"
version = "2.0.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7f7cf42b4507d8ea322120659672cf1b9dbb93f8f2d4ecfd6e51350ff5b17a1d"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.104",
]
[[package]] [[package]]
name = "thread_local" name = "thread_local"
version = "1.1.3" version = "1.1.3"
@ -947,7 +603,7 @@ checksum = "5b8a1e28f2deaa14e508979454cb3a223b10b938b45af148bc0986de36f1923b"
dependencies = [ dependencies = [
"proc-macro2", "proc-macro2",
"quote", "quote",
"syn 2.0.104", "syn 2.0.48",
] ]
[[package]] [[package]]
@ -1034,12 +690,6 @@ version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "711b9620af191e0cdc7468a8d14e709c3dcdb115b36f838e601583af800a370a" checksum = "711b9620af191e0cdc7468a8d14e709c3dcdb115b36f838e601583af800a370a"
[[package]]
name = "void"
version = "1.0.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6a02e4885ed3bc0f2de90ea6dd45ebcbb66dacffe03547fadbb0eeae2770887d"
[[package]] [[package]]
name = "wasi" name = "wasi"
version = "0.11.0+wasi-snapshot-preview1" version = "0.11.0+wasi-snapshot-preview1"

View File

@ -1,22 +1,11 @@
[package] [package]
name = "atem-connection-rs" name = "atem-connection-rs"
version = "0.1.0" version = "0.1.0"
edition = "2024" edition = "2021"
[dependencies] [dependencies]
derive-getters = "0.2.0" derive-getters = "0.2.0"
derive-new = "0.6.0" derive-new = "0.6.0"
itertools = {version = "0.14.0", default-features = false}
log = "0.4.14" log = "0.4.14"
tokio = { version = "1.13.0", features = ["full"], optional = true } tokio = { version = "1.13.0", features = ["full"] }
tokio-util = { version = "0.7.10", optional = true } tokio-util = "0.7.10"
enumflags2 = { version = "0.7.12", default-features = false }
embassy-net = { version = "0.7.0", optional = true, features = ["proto-ipv4", "medium-ip", "udp"]}
defmt = {version = "1.0.0", optional = true}
[features]
default = ["std", "embassy-net", "defmt"]
std = ["dep:tokio", "dep:tokio-util", "itertools/use_std"]
embassy-net = ["dep:embassy-net"]
defmt = ["dep:defmt"]

View File

@ -1,11 +1,9 @@
use std::{ use std::{
boxed::Box,
collections::{HashMap, VecDeque}, collections::{HashMap, VecDeque},
net::SocketAddr, net::SocketAddr,
ops::DerefMut, ops::DerefMut,
sync::Arc, sync::Arc,
time::Duration, time::Duration,
vec::Vec,
}; };
use tokio::{select, sync::Semaphore}; use tokio::{select, sync::Semaphore};

View File

@ -1,87 +0,0 @@
//! Definitions and decoding of ATEM protocol fields
/// An uninterpreted ATEM protocol field - a 4-ascii character type plus variable length data
pub struct RawField<'a> {
pub r#type: &'a str,
pub data: &'a [u8],
}
#[derive(Debug)]
pub struct _Ver {
pub major: u16,
pub minor: u16,
}
impl<'a> Field<'a> for _Ver {
const TYPE: &'static str = "_ver";
fn decode(data: &'a [u8]) -> Result<Self, FieldParsingError> {
let data = checked_len::<4>(data)?;
Ok(Self {
major: u16::from_be_bytes(data[0..=1].try_into().unwrap()),
minor: u16::from_be_bytes(data[2..=3].try_into().unwrap()),
})
}
}
#[derive(Debug)]
pub struct PrvI {
pub m_e_index: u8,
pub source_index: u16,
pub pvw_in_pgm: bool,
}
impl<'a> Field<'a> for PrvI {
const TYPE: &'static str = "PrvI";
fn decode(data: &'a [u8]) -> Result<Self, FieldParsingError> {
let data = checked_len::<8>(data)?;
Ok(Self {
m_e_index: data[0],
source_index: u16::from_be_bytes(data[2..=3].try_into().unwrap()),
pvw_in_pgm: data[4] != 0,
})
}
}
#[derive(Debug)]
pub struct PrgI {
pub m_e_index: u8,
pub source_index: u16,
}
impl<'a> Field<'a> for PrgI {
const TYPE: &'static str = "PrgI";
fn decode(data: &'a [u8]) -> Result<Self, FieldParsingError> {
let data = checked_len::<4>(data)?;
Ok(Self {
m_e_index: data[0],
source_index: u16::from_be_bytes(data[2..=3].try_into().unwrap()),
})
}
}
pub trait Field<'a>: Sized {
const TYPE: &'static str;
fn decode(data: &'a [u8]) -> Result<Self, FieldParsingError>;
fn try_from_raw(raw: RawField<'a>) -> Result<Self, FieldParsingError> {
if Self::TYPE != raw.r#type {
Err(FieldParsingError::MismatchedFieldType)
} else {
Self::decode(&raw.data)
}
}
}
#[derive(Debug)]
pub enum FieldParsingError {
UnexpectedLength { expected: usize, got: usize },
MismatchedFieldType,
}
fn checked_len<'a, const LEN: usize>(data: &'a [u8]) -> Result<&'a [u8; LEN], FieldParsingError> {
data.try_into()
.map_err(|_| FieldParsingError::UnexpectedLength {
expected: LEN,
got: data.len(),
})
}

View File

@ -1,314 +1,115 @@
//! This module contains [`AtemPacket`] which is a zero*ish* copy abstraction over a [`AsRef`]`<[u8]>`.
use core::{fmt::Display, str};
use enumflags2::{BitFlags, bitflags};
use super::atem_field::{Field, FieldParsingError, RawField};
/// The "hello" packet to start communication with the ATEM
pub const COMMAND_CONNECT_HELLO: [u8; 20] = [
0x10, 0x14, 0x53, 0xab, 0x00, 0x00, 0x00, 0x00, 0x00, 0x3a, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
];
/// Maximum atem packet length. This is determined by the maximum value that can be stored
/// in the length field.
pub const MAX_LEN: usize = field::LEN_MASK as usize;
mod field {
use core::ops::{Range, RangeInclusive};
pub(crate) const FLAGS_LEN_H: usize = 0;
pub(crate) const FLAGS_MASK: u8 = 0b1111_1000;
pub(crate) const LEN_L: usize = 1;
pub(crate) const LEN: RangeInclusive<usize> = FLAGS_LEN_H..=LEN_L;
pub(crate) const LEN_MASK: u16 = 0x7ff;
pub(crate) const SESSION_ID: Range<usize> = 2..4;
pub(crate) const ACK_NUMBER: Range<usize> = 4..6;
pub(crate) const REMOTE_SEQ_NUM: Range<usize> = 8..10;
pub(crate) const LOCAL_SEQ_NUM: Range<usize> = 10..12;
}
/// An ATEM protocol packet
#[derive(Debug)] #[derive(Debug)]
pub struct AtemPacket<T: AsRef<[u8]>> { pub struct AtemPacket<'packet_buffer> {
buf: T, flags: u8,
session_id: u16,
remote_packet_id: u16,
retransmit_requested_from_packet_id: Option<u16>,
ack_reply: Option<u16>,
body: Option<&'packet_buffer [u8]>,
} }
impl<'a> TryFrom<&'a [u8]> for AtemPacket<&'a [u8]> { pub enum AtemPacketErr {
type Error = AtemPacketParseError; TooShort(String),
fn try_from(buf: &'a [u8]) -> Result<Self, Self::Error> { LengthDiffers(String),
AtemPacket::new_checked(buf)
}
} }
#[derive(Debug)] #[derive(PartialEq)]
#[cfg_attr(feature = "defmt", derive(defmt::Format))]
pub enum AtemPacketParseError {
/// The packet was shorter than any valid packet should be.
TooShort { got: usize },
/// The packet's stated and actual lengths were different
LengthDiffers { expected: u16, got: usize },
/// The packet had invalid (unknown) flags set
InvalidFlags,
}
/// ATEM packet flags
#[bitflags]
#[repr(u8)]
#[derive(PartialEq, Copy, Clone, Debug)]
pub enum PacketFlag { pub enum PacketFlag {
/// AKA "reliable". This packet should be ACKed by the other end AckRequest,
AckRequest = 0x1, NewSessionId,
/// AKA "syn" - used in the connection handshake IsRetransmit,
NewSessionId = 0x2, RetransmitRequest,
/// This packet is a retransmit of a previous one AckReply,
IsRetransmit = 0x4,
/// Request retransmission of a previous sequence ID
RetransmitRequest = 0x8,
/// This packet is an ACK
AckReply = 0x10,
} }
/// An error while constructing a packet impl From<PacketFlag> for u8 {
#[derive(Debug)] fn from(flag: PacketFlag) -> Self {
pub enum AtemPacketInitError { match flag {
/// The provided buf was too small PacketFlag::AckRequest => 0x01,
BufTooSmall { need: usize, got: usize }, PacketFlag::NewSessionId => 0x02,
/// The provided payload was too long PacketFlag::IsRetransmit => 0x04,
DataTooLong { got: usize, max: usize }, PacketFlag::RetransmitRequest => 0x08,
} PacketFlag::AckReply => 0x10,
impl<'a> AtemPacket<&'a mut [u8]> {
/// Initialise an [`AtemPacket`] into the given `buf`, with optional payload copied from `data`.
/// Returns the constructed packet, and the remaining spare space in `buf`, if there was any.
pub fn init<'b>(
buf: &'a mut [u8],
flags: BitFlags<PacketFlag>,
session_id: u16,
local_seq_num: u16,
data: Option<&'b [u8]>,
) -> Result<(Self, &'a mut [u8]), AtemPacketInitError> {
let len = 12 + data.map_or(0, |d| d.len());
if len > buf.len() {
return Err(AtemPacketInitError::BufTooSmall {
need: len,
got: buf.len(),
});
} }
if len > MAX_LEN {
return Err(AtemPacketInitError::DataTooLong {
got: len - 12,
max: MAX_LEN - 12,
});
}
let (pkt, rem) = buf.split_at_mut(len);
let mut p = AtemPacket { buf: pkt };
p.set_flags(flags)
.set_session_id(session_id)
.set_local_seq_num(local_seq_num)
.set_len(len.try_into().unwrap());
if let Some(d) = data {
p.set_data(d);
}
Ok((p, rem))
} }
} }
impl<T: AsRef<[u8]> + AsMut<[u8]>> AtemPacket<T> { impl<'packet_buffer> AtemPacket<'packet_buffer> {
pub fn set_flags(&mut self, flags: BitFlags<PacketFlag>) -> &mut Self {
let prev = self.buf.as_ref()[field::FLAGS_LEN_H];
self.buf.as_mut()[field::FLAGS_LEN_H] = (prev & !field::FLAGS_MASK) | (flags.bits() << 3);
self
}
pub fn set_len(&mut self, value: u16) -> &mut Self {
let v = value.to_be_bytes();
self.buf.as_mut()[field::FLAGS_LEN_H] &= v[0] | field::FLAGS_MASK;
self.buf.as_mut()[field::LEN_L] = v[1];
self
}
pub fn set_session_id(&mut self, value: u16) -> &mut Self {
self.buf.as_mut()[field::SESSION_ID].copy_from_slice(&value.to_be_bytes());
self
}
pub fn set_local_seq_num(&mut self, value: u16) -> &mut Self {
self.buf.as_mut()[field::LOCAL_SEQ_NUM].copy_from_slice(&value.to_be_bytes());
self
}
pub fn set_remote_seq_num(&mut self, value: u16) -> &mut Self {
self.buf.as_mut()[field::REMOTE_SEQ_NUM].copy_from_slice(&value.to_be_bytes());
self
}
pub fn set_data<'a>(&mut self, value: &'a [u8]) -> &mut Self {
self.buf.as_mut()[12..].copy_from_slice(value);
self
}
pub fn set_ack_num(&mut self, value: u16) -> &mut Self {
self.buf.as_mut()[field::ACK_NUMBER].copy_from_slice(&value.to_be_bytes());
self
}
}
impl<T: AsRef<[u8]>> AtemPacket<T> {
pub fn new_checked(buf: T) -> Result<Self, AtemPacketParseError> {
let len = buf.as_ref().len();
if len < 12 {
return Err(AtemPacketParseError::TooShort {
got: buf.as_ref().len(),
});
}
let p = Self { buf };
if p.length() as usize != len {
return Err(AtemPacketParseError::LengthDiffers {
expected: p.length(),
got: len,
});
}
// Check flags are valid
let _: BitFlags<PacketFlag> = p
.flags_raw()
.try_into()
.map_err(|_| AtemPacketParseError::InvalidFlags)?;
Ok(p)
}
/// Consumes self, returning the inner packet buffer
pub fn inner(self) -> T {
self.buf
}
pub fn length(&self) -> u16 {
u16::from_be_bytes(self.buf.as_ref()[field::LEN].try_into().unwrap()) & field::LEN_MASK
}
fn flags_raw(&self) -> u8 {
self.buf.as_ref()[field::FLAGS_LEN_H] >> 3
}
pub fn flags(&self) -> BitFlags<PacketFlag> {
// We `unwrap` here, but given we check the flags in the constructor this should never panic.
self.flags_raw().try_into().unwrap()
}
pub fn session_id(&self) -> u16 { pub fn session_id(&self) -> u16 {
u16::from_be_bytes(self.buf.as_ref()[field::SESSION_ID].try_into().unwrap()) self.session_id
} }
pub fn ack_number(&self) -> u16 { pub fn remote_packet_id(&self) -> u16 {
u16::from_be_bytes(self.buf.as_ref()[field::ACK_NUMBER].try_into().unwrap()) self.remote_packet_id
} }
pub fn remote_sequence_number(&self) -> u16 { pub fn body(&self) -> Option<&[u8]> {
u16::from_be_bytes(self.buf.as_ref()[field::REMOTE_SEQ_NUM].try_into().unwrap()) self.body
}
pub fn local_sequence_number(&self) -> u16 {
u16::from_be_bytes(self.buf.as_ref()[field::LOCAL_SEQ_NUM].try_into().unwrap())
} }
pub fn retransmit_request(&self) -> Option<u16> { pub fn retransmit_request(&self) -> Option<u16> {
self.flags() self.retransmit_requested_from_packet_id
.contains(PacketFlag::RetransmitRequest)
.then_some(u16::from_be_bytes([
self.buf.as_ref()[6],
self.buf.as_ref()[7],
]))
} }
pub fn ack_reply(&self) -> Option<u16> { pub fn ack_reply(&self) -> Option<u16> {
self.flags() self.ack_reply
.contains(PacketFlag::AckReply)
.then_some(self.ack_number())
} }
/// Get an iterator over the `Field`s in this packet. pub fn has_flag(&self, flag: PacketFlag) -> bool {
/// self.flags & u8::from(flag) > 0
/// Returns None if this is a packet without fields.
pub fn raw_fields(&self) -> RawFields {
// TODO: do we only ever get newsessionid during the handshake (i.e. not in a packet with fields)?
let has_fields = !self.flags().contains(PacketFlag::NewSessionId);
RawFields::new(if has_fields { self.body() } else { &[] })
}
pub fn fields<'a, F: Field<'a>>(
&'a self,
) -> impl Iterator<Item = Result<F, FieldParsingError>> + use<'a, F, T> {
self.raw_fields().filter_map(|f| match F::try_from_raw(f) {
Err(FieldParsingError::MismatchedFieldType) => None,
x => Some(x),
})
}
pub fn body(&self) -> &[u8] {
&self.buf.as_ref()[12..]
} }
} }
impl<T: AsRef<[u8]>> Display for AtemPacket<T> { impl<'packet_buffer> TryFrom<&'packet_buffer [u8]> for AtemPacket<'packet_buffer> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { type Error = AtemPacketErr;
write!(
f,
"len: {}, sid: {}, flags: {}, ack: {}, localseq: {}, remoteseq: {}",
self.length(),
self.session_id(),
self.flags(),
self.ack_number(),
self.local_sequence_number(),
self.remote_sequence_number(),
)
}
}
/// Created by [`AtemPacket::raw_fields`] fn try_from(buffer: &'packet_buffer [u8]) -> Result<Self, Self::Error> {
pub struct RawFields<'a> { if buffer.len() < 12 {
data: &'a [u8], return Err(AtemPacketErr::TooShort(format!(
// The offset of the next field in the packet "Invalid packet from ATEM {:x?}",
offset: usize, buffer
} )));
impl<'a> RawFields<'a> {
pub(crate) fn new(data: &'a [u8]) -> Self {
Self { data, offset: 0 }
}
}
impl<'a> Iterator for RawFields<'a> {
type Item = RawField<'a>;
fn next(&mut self) -> Option<Self::Item> {
let remain = self.data.len() - self.offset;
if remain == 0 {
return None;
} else if remain < 8 {
// TODO: is 8 indeed the minimum size for something here? (i.e. no field data)
panic!("Oh no");
} }
let length = let length = u16::from_be_bytes([buffer[0], buffer[1]]) & 0x07ff;
u16::from_be_bytes(self.data[self.offset..=self.offset + 1].try_into().unwrap()); if length as usize != buffer.len() {
// TODO: sanity check length return Err(AtemPacketErr::LengthDiffers(format!(
let r#type = str::from_utf8(&self.data[self.offset + 4..=self.offset + 7]).unwrap(); "Length of message differs, expected {} got {}",
let data = &self.data[self.offset + 8..self.offset + (length as usize)]; length,
self.offset += length as usize; buffer.len()
Some(RawField { r#type, data }) )));
} }
}
let flags = buffer[0] >> 3;
#[cfg(test)] let session_id = u16::from_be_bytes([buffer[2], buffer[3]]);
mod tests { let remote_packet_id = u16::from_be_bytes([buffer[10], buffer[11]]);
use super::*;
let body = if buffer.len() > 12 {
#[test] Some(&buffer[12..])
fn test_init() { } else {
let mut buf = [0u8; 12]; None
let (mut p, _) = };
AtemPacket::init(&mut buf, PacketFlag::AckReply.into(), 1234, 5678, None).unwrap();
p.set_remote_seq_num(9012); let retransmit_requested_from_packet_id =
assert!(p.flags().contains(PacketFlag::AckReply)); if flags & u8::from(PacketFlag::RetransmitRequest) > 0 {
assert_eq!(p.session_id(), 1234); Some(u16::from_be_bytes([buffer[6], buffer[7]]))
assert_eq!(p.local_sequence_number(), 5678); } else {
assert_eq!(p.remote_sequence_number(), 9012); None
assert_eq!(p.length(), 12); };
let ack_reply = if flags & u8::from(PacketFlag::AckReply) > 0 {
Some(u16::from_be_bytes([buffer[4], buffer[5]]))
} else {
None
};
Ok(AtemPacket {
flags,
session_id,
remote_packet_id,
body,
retransmit_requested_from_packet_id,
ack_reply,
})
} }
} }

View File

@ -1,16 +1,12 @@
use std::{ use std::{
borrow::ToOwned, collections::VecDeque,
fmt::Display, fmt::Display,
io, io,
net::SocketAddr, net::SocketAddr,
string::{String, ToString},
sync::Arc, sync::Arc,
time::{Duration, Instant}, time::{Duration, SystemTime},
vec,
vec::Vec,
}; };
use itertools::Itertools;
use tokio::{ use tokio::{
net::UdpSocket, net::UdpSocket,
select, select,
@ -18,17 +14,20 @@ use tokio::{
}; };
use crate::{ use crate::{
atem_lib::atem_packet::{AtemPacket, COMMAND_CONNECT_HELLO}, atem_lib::{atem_packet::AtemPacket, atem_util},
commands::command_base::{BasicWritableCommand, DeserializedCommand}, commands::{
command_base::{BasicWritableCommand, DeserializedCommand},
parse_commands::deserialize_commands,
},
enums::ProtocolVersion, enums::ProtocolVersion,
}; };
use super::atem_packet::PacketFlag; use super::atem_packet::PacketFlag;
const IN_FLIGHT_TIMEOUT: Duration = Duration::from_millis(60); const IN_FLIGHT_TIMEOUT: u64 = 60;
const CONNECTION_TIMEOUT: Duration = Duration::from_millis(5000); const CONNECTION_TIMEOUT: u64 = 5000;
const CONNECTION_RETRY_INTERVAL: Duration = Duration::from_millis(1000); const CONNECTION_RETRY_INTERVAL: u64 = 1000;
const RETRANSMIT_CHECK_INTERVAL: Duration = Duration::from_millis(1000); const RETRANSMIT_CHECK_INTERVAL: u64 = 1000;
const MAX_PACKET_RETRIES: u16 = 10; const MAX_PACKET_RETRIES: u16 = 10;
const MAX_PACKET_ID: u16 = 1 << 15; const MAX_PACKET_ID: u16 = 1 << 15;
const MAX_PACKET_PER_ACK: u16 = 16; const MAX_PACKET_PER_ACK: u16 = 16;
@ -93,8 +92,8 @@ impl AtemSocketCommand {
pub struct AtemSocket { pub struct AtemSocket {
connection_state: ConnectionState, connection_state: ConnectionState,
reconnect_timer: Option<Instant>, reconnect_timer: Option<SystemTime>,
retransmit_timer: Option<Instant>, retransmit_timer: Option<SystemTime>,
next_tracking_id: u64, next_tracking_id: u64,
@ -106,10 +105,10 @@ pub struct AtemSocket {
protocol_version: ProtocolVersion, protocol_version: ProtocolVersion,
last_received_at: Instant, last_received_at: SystemTime,
last_received_packed_id: u16, last_received_packed_id: u16,
in_flight: Vec<InFlightPacket>, in_flight: Vec<InFlightPacket>,
ack_timer: Option<Instant>, ack_timer: Option<SystemTime>,
received_without_ack: u16, received_without_ack: u16,
atem_message_rx: tokio::sync::mpsc::Receiver<AtemSocketMessage>, atem_message_rx: tokio::sync::mpsc::Receiver<AtemSocketMessage>,
@ -142,7 +141,7 @@ struct InFlightPacket {
packet_id: u16, packet_id: u16,
tracking_id: u64, tracking_id: u64,
payload: Vec<u8>, payload: Vec<u8>,
pub last_sent: Instant, pub last_sent: SystemTime,
pub resent: u16, pub resent: u16,
} }
@ -171,7 +170,7 @@ impl AtemSocket {
protocol_version: ProtocolVersion::V7_2, protocol_version: ProtocolVersion::V7_2,
last_received_at: Instant::now(), last_received_at: SystemTime::now(),
last_received_packed_id: 0, last_received_packed_id: 0,
in_flight: vec![], in_flight: vec![],
ack_timer: None, ack_timer: None,
@ -259,7 +258,7 @@ impl AtemSocket {
self.in_flight = vec![]; self.in_flight = vec![];
log::debug!("Reconnect"); log::debug!("Reconnect");
self.send_packet(&COMMAND_CONNECT_HELLO).await; self.send_packet(&atem_util::COMMAND_CONNECT_HELLO).await;
self.connection_state = ConnectionState::SynSent; self.connection_state = ConnectionState::SynSent;
Ok(()) Ok(())
@ -300,24 +299,28 @@ impl AtemSocket {
self.next_send_packet_id = 0; self.next_send_packet_id = 0;
} }
let opcode = u16::from(u8::from(PacketFlag::AckRequest)) << 11;
let mut buffer = vec![0; 20 + payload.len()]; let mut buffer = vec![0; 20 + payload.len()];
let (p, _) = AtemPacket::init( // Headers
&mut buffer, buffer[0..2].copy_from_slice(&u16::to_be_bytes(opcode | (payload.len() as u16 + 20)));
PacketFlag::AckRequest.into(), buffer[2..4].copy_from_slice(&u16::to_be_bytes(self.session_id));
self.session_id, buffer[10..12].copy_from_slice(&u16::to_be_bytes(packet_id));
packet_id,
Some([raw_name.as_bytes(), payload].concat().as_slice()),
)
.unwrap();
self.send_packet(&p.inner()).await; // Command
buffer[12..14].copy_from_slice(&u16::to_be_bytes(payload.len() as u16 + 8));
buffer[16..20].copy_from_slice(raw_name.as_bytes());
// Body
buffer[20..20 + payload.len()].copy_from_slice(payload);
self.send_packet(&buffer).await;
self.in_flight.push(InFlightPacket { self.in_flight.push(InFlightPacket {
packet_id, packet_id,
tracking_id, tracking_id,
payload: buffer, payload: buffer,
last_sent: Instant::now(), last_sent: SystemTime::now(),
resent: 0, resent: 0,
}) })
} }
@ -335,15 +338,17 @@ impl AtemSocket {
} }
} }
if let Some(ack_time) = self.ack_timer { if let Some(ack_time) = self.ack_timer {
if ack_time <= Instant::now() { if ack_time <= SystemTime::now() {
self.ack_timer = None; self.ack_timer = None;
self.received_without_ack = 0; self.received_without_ack = 0;
self.send_ack(self.last_received_packed_id).await; self.send_ack(self.last_received_packed_id).await;
} }
} }
if let Some(reconnect_time) = self.reconnect_timer { if let Some(reconnect_time) = self.reconnect_timer {
if reconnect_time <= Instant::now() { if reconnect_time <= SystemTime::now() {
if self.last_received_at + CONNECTION_TIMEOUT <= Instant::now() { if self.last_received_at + Duration::from_millis(CONNECTION_TIMEOUT)
<= SystemTime::now()
{
log::debug!("{:?}", self.last_received_at); log::debug!("{:?}", self.last_received_at);
log::debug!("Connection timed out, restarting"); log::debug!("Connection timed out, restarting");
self.restart_connection().await; self.restart_connection().await;
@ -352,7 +357,7 @@ impl AtemSocket {
} }
} }
if let Some(retransmit_time) = self.retransmit_timer { if let Some(retransmit_time) = self.retransmit_timer {
if retransmit_time <= Instant::now() { if retransmit_time <= SystemTime::now() {
self.check_for_retransmit().await; self.check_for_retransmit().await;
self.start_retransmit_timer(); self.start_retransmit_timer();
} }
@ -381,23 +386,18 @@ impl AtemSocket {
} }
async fn recieved_packet(&mut self, packet: &[u8]) { async fn recieved_packet(&mut self, packet: &[u8]) {
let Ok(atem_packet): Result<AtemPacket<_>, _> = packet.try_into() else { let Ok(atem_packet): Result<AtemPacket, _> = packet.try_into() else {
return; return;
}; };
log::debug!("Received {}", atem_packet,); log::debug!("Received {:x?}", atem_packet);
log::debug!(
"fields: {}",
atem_packet.raw_fields().map(|f| f.r#type).join(",")
);
self.last_received_at = Instant::now(); self.last_received_at = SystemTime::now();
self.session_id = atem_packet.session_id(); self.session_id = atem_packet.session_id();
// TODO: bit of a naming clash here let remote_packet_id = atem_packet.remote_packet_id();
let remote_packet_id = atem_packet.local_sequence_number();
if atem_packet.flags().contains(PacketFlag::NewSessionId) { if atem_packet.has_flag(PacketFlag::NewSessionId) {
log::debug!("New session"); log::debug!("New session");
self.connection_state = ConnectionState::Established; self.connection_state = ConnectionState::Established;
self.last_received_packed_id = remote_packet_id; self.last_received_packed_id = remote_packet_id;
@ -413,12 +413,14 @@ impl AtemSocket {
self.retransmit_from(from_packet_id).await; self.retransmit_from(from_packet_id).await;
} }
if atem_packet.flags().contains(PacketFlag::AckRequest) { if atem_packet.has_flag(PacketFlag::AckRequest) {
if remote_packet_id == (self.last_received_packed_id + 1) % MAX_PACKET_ID { if remote_packet_id == (self.last_received_packed_id + 1) % MAX_PACKET_ID {
self.last_received_packed_id = remote_packet_id; self.last_received_packed_id = remote_packet_id;
self.send_or_queue_ack().await; self.send_or_queue_ack().await;
self.on_commands_received(atem_packet.body()); if let Some(body) = atem_packet.body() {
self.on_commands_received(body);
}
} else if self } else if self
.is_packet_covered_by_ack(self.last_received_packed_id, remote_packet_id) .is_packet_covered_by_ack(self.last_received_packed_id, remote_packet_id)
{ {
@ -426,7 +428,7 @@ impl AtemSocket {
} }
} }
if atem_packet.flags().contains(PacketFlag::IsRetransmit) { if atem_packet.has_flag(PacketFlag::IsRetransmit) {
log::debug!("ATEM retransmitted packet {:x?}", remote_packet_id); log::debug!("ATEM retransmitted packet {:x?}", remote_packet_id);
} }
@ -470,23 +472,19 @@ impl AtemSocket {
self.ack_timer = None; self.ack_timer = None;
self.send_ack(self.last_received_packed_id).await; self.send_ack(self.last_received_packed_id).await;
} else if self.ack_timer.is_none() { } else if self.ack_timer.is_none() {
self.ack_timer = Some(Instant::now() + Duration::from_millis(5)); self.ack_timer = Some(SystemTime::now() + Duration::from_millis(5));
} }
} }
async fn send_ack(&mut self, packet_id: u16) { async fn send_ack(&mut self, packet_id: u16) {
log::debug!("Sending ack for packet {:x?}", packet_id); log::debug!("Sending ack for packet {:x?}", packet_id);
let flag: u8 = PacketFlag::AckReply.into();
let opcode = u16::from(flag) << 11;
let mut buffer: [u8; ACK_PACKET_LENGTH as _] = [0; 12]; let mut buffer: [u8; ACK_PACKET_LENGTH as _] = [0; 12];
let (mut p, _) = AtemPacket::init( buffer[0..2].copy_from_slice(&u16::to_be_bytes(opcode | ACK_PACKET_LENGTH));
&mut buffer, buffer[2..4].copy_from_slice(&u16::to_be_bytes(self.session_id));
PacketFlag::AckReply.into(), buffer[4..6].copy_from_slice(&u16::to_be_bytes(packet_id));
self.session_id, self.send_packet(&buffer).await;
0,
None,
)
.unwrap();
p.set_ack_num(packet_id);
self.send_packet(p.inner()).await;
} }
async fn retransmit_from(&mut self, from_id: u16) { async fn retransmit_from(&mut self, from_id: u16) {
@ -507,7 +505,7 @@ impl AtemSocket {
if sent_packet.packet_id == from_id if sent_packet.packet_id == from_id
|| !self.is_packet_covered_by_ack(from_id, sent_packet.packet_id) || !self.is_packet_covered_by_ack(from_id, sent_packet.packet_id)
{ {
sent_packet.last_sent = Instant::now(); sent_packet.last_sent = SystemTime::now();
sent_packet.resent += 1; sent_packet.resent += 1;
self.send_packet(&sent_packet.payload).await; self.send_packet(&sent_packet.payload).await;
@ -521,7 +519,8 @@ impl AtemSocket {
async fn check_for_retransmit(&mut self) { async fn check_for_retransmit(&mut self) {
for sent_packet in self.in_flight.clone() { for sent_packet in self.in_flight.clone() {
if sent_packet.last_sent + IN_FLIGHT_TIMEOUT < Instant::now() { if sent_packet.last_sent + Duration::from_millis(IN_FLIGHT_TIMEOUT) < SystemTime::now()
{
if sent_packet.resent <= MAX_PACKET_RETRIES if sent_packet.resent <= MAX_PACKET_RETRIES
&& self && self
.is_packet_covered_by_ack(self.next_send_packet_id, sent_packet.packet_id) .is_packet_covered_by_ack(self.next_send_packet_id, sent_packet.packet_id)
@ -538,11 +537,9 @@ impl AtemSocket {
} }
fn on_commands_received(&mut self, payload: &[u8]) { fn on_commands_received(&mut self, payload: &[u8]) {
if !payload.is_empty() { let _ = self
let _ = self .atem_event_tx
.atem_event_tx .send(AtemSocketEvent::ReceivedCommands(payload.to_vec()));
.send(AtemSocketEvent::ReceivedCommands(payload.to_vec()));
}
} }
fn on_command_acknowledged(&mut self, packets: Vec<AckedPacket>) { fn on_command_acknowledged(&mut self, packets: Vec<AckedPacket>) {
@ -577,11 +574,13 @@ impl AtemSocket {
} }
fn start_reconnect_timer(&mut self) { fn start_reconnect_timer(&mut self) {
self.reconnect_timer = Some(Instant::now() + CONNECTION_RETRY_INTERVAL); self.reconnect_timer =
Some(SystemTime::now() + Duration::from_millis(CONNECTION_RETRY_INTERVAL));
} }
fn start_retransmit_timer(&mut self) { fn start_retransmit_timer(&mut self) {
self.retransmit_timer = Some(Instant::now() + RETRANSMIT_CHECK_INTERVAL); self.retransmit_timer =
Some(SystemTime::now() + Duration::from_millis(RETRANSMIT_CHECK_INTERVAL));
} }
fn next_packet_tracking_id(&mut self) -> u64 { fn next_packet_tracking_id(&mut self) -> u64 {

View File

@ -1,285 +0,0 @@
use core::net::{Ipv4Addr, SocketAddrV4};
use log::debug;
use super::atem_packet::{AtemPacket, AtemPacketParseError, PacketFlag};
pub const ATEM_PORT: u16 = 9910;
pub const SYN_PAYLOAD: [u8; 8] = [0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00];
#[derive(Debug, PartialEq, Eq)]
pub enum AtemSocketState {
Closed,
StateDump,
Established { session_id: u16 },
}
/// An abstraction over an async UDP socket
pub trait SocketBackend {
type SendError;
type RecvError;
async fn send_to(&mut self, data: &[u8], addr: SocketAddrV4) -> Result<(), Self::SendError>;
async fn recv(&mut self, data: &mut [u8]) -> Result<usize, Self::RecvError>;
}
#[cfg(feature = "std")]
pub mod tokio {
use core::net::Ipv4Addr;
use tokio::net::UdpSocket;
use super::{AtemSocket, SocketBackend};
#[derive(Debug)]
pub struct TokioBackend {
sock: UdpSocket,
}
impl TokioBackend {
pub async fn new(addr: Ipv4Addr) -> Result<AtemSocket<Self>, std::io::Error> {
let b = Self {
sock: UdpSocket::bind("0.0.0.0:0").await?,
};
Ok(AtemSocket::new_with_backend(addr, b))
}
}
impl SocketBackend for TokioBackend {
type SendError = std::io::Error;
type RecvError = std::io::Error;
async fn send_to(
&mut self,
data: &[u8],
addr: core::net::SocketAddrV4,
) -> Result<(), Self::SendError> {
self.sock.send_to(data, addr).await?;
Ok(())
}
async fn recv(&mut self, data: &mut [u8]) -> Result<usize, Self::RecvError> {
let (cnt, _) = self.sock.recv_from(data).await?;
Ok(cnt)
}
}
}
#[cfg(feature = "embassy-net")]
pub mod enet {
use core::net::{Ipv4Addr, SocketAddrV4};
use embassy_net::udp::{RecvError, SendError, UdpSocket};
use super::{AtemSocket, SocketBackend};
pub struct EnetBackend<'a> {
sock: UdpSocket<'a>,
}
impl core::fmt::Debug for EnetBackend<'_> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(f, "EnetBackend")
}
}
impl<'a> EnetBackend<'a> {
pub fn new(addr: Ipv4Addr, sock: UdpSocket<'a>) -> AtemSocket<Self> {
let b = Self { sock };
AtemSocket::new_with_backend(addr, b)
}
}
impl SocketBackend for EnetBackend<'_> {
type SendError = SendError;
type RecvError = RecvError;
async fn send_to(
&mut self,
data: &[u8],
addr: SocketAddrV4,
) -> Result<(), Self::SendError> {
self.sock.send_to(data, addr).await
}
async fn recv(&mut self, data: &mut [u8]) -> Result<usize, Self::RecvError> {
let (count, _) = self.sock.recv_from(data).await?;
Ok(count)
}
}
}
#[derive(Debug)]
#[cfg_attr(feature = "defmt", derive(defmt::Format))]
pub enum AtemError<I: SocketBackend> {
/// Tried to do something but the socket is closed
SocketClosed,
/// There was an error from the underlying socket while receiving
Recv(I::RecvError),
/// There was an error from the underlying socket while sending
Send(I::SendError),
/// We received an invalid ATEM protocol packet
BadPacket(AtemPacketParseError),
/// The switcher denied our connection attempt.
SwitcherDeniedConnect,
/// There was some nonspecific protocol error
ProtocolError,
}
impl<I> From<AtemPacketParseError> for AtemError<I>
where
I: SocketBackend,
{
fn from(value: AtemPacketParseError) -> Self {
Self::BadPacket(value)
}
}
pub struct AtemSocket<I: SocketBackend> {
inner: I,
addr: SocketAddrV4,
state: AtemSocketState,
local_seq_num: u16,
}
impl<I> AtemSocket<I>
where
I: SocketBackend,
{
fn new_with_backend(addr: Ipv4Addr, backend: I) -> Self {
Self {
inner: backend,
addr: SocketAddrV4::new(addr, ATEM_PORT),
state: AtemSocketState::Closed,
local_seq_num: 0,
}
}
/// Connect to the ATEM, performing the connection handshake.
pub async fn connect(&mut self) -> Result<(), AtemError<I>> {
// Reset state
self.state = AtemSocketState::Closed;
self.local_seq_num = 0;
let mut buf = [0u8; 20];
// FIXME: this should be random
let initial_sid = 0x53ab;
// Step 1: Send the "SYN" with a random session id and magic payload
let (mut p, _) = AtemPacket::init(
&mut buf,
PacketFlag::NewSessionId.into(),
initial_sid,
self.local_seq_num,
Some(&SYN_PAYLOAD),
)
.unwrap();
p.set_remote_seq_num(0x3a);
self.inner
.send_to(p.inner(), self.addr)
.await
.map_err(|e| AtemError::Send(e))?;
debug!("Connect step 1: sent SYN");
// Step 2: Wait for the response from the ATEM.
// this should also have the SYN flag set and have the same session id that we sent.
// The first byte of the payload should be 0x02 - if it's 0x04 we should restart
// the connection
self.inner
.recv(&mut buf)
.await
.map_err(|e| AtemError::Recv(e))?;
let ack = AtemPacket::new_checked(buf)?;
debug!("RX: {}", ack);
if ack.flags() != PacketFlag::NewSessionId {
return Err(AtemError::ProtocolError);
}
if ack.session_id() != initial_sid {
return Err(AtemError::ProtocolError);
}
match ack.body()[0] {
0x02 => {
// all good
debug!("Connect step 2: got reply");
}
0x04 => {
// Restart connection
return Err(AtemError::SwitcherDeniedConnect);
}
_ => {
return Err(AtemError::ProtocolError);
}
}
// Step 3: Send an ACK for the packet we just received.
self.ack(initial_sid, ack.local_sequence_number()).await?;
debug!("Connect step 3: sent ACK");
self.state = AtemSocketState::StateDump;
Ok(())
}
async fn ack(&mut self, sid: u16, ack_num: u16) -> Result<(), AtemError<I>> {
let mut buf = [0u8; 12];
let (mut p, _) = AtemPacket::init(
&mut buf,
PacketFlag::AckReply.into(),
sid,
// N.B: local seq num is set to 0 in ACKs
0,
None,
)
.unwrap();
p.set_ack_num(ack_num);
debug!("TX: {}", p);
self.inner
.send_to(p.inner(), self.addr)
.await
.map_err(|e| AtemError::Send(e))?;
Ok(())
}
/// Receive a packet from the ATEM. Automatically handles sending ACKs for received packets.
pub async fn recv<'a>(
&mut self,
buf: &'a mut [u8],
) -> Result<AtemPacket<&'a [u8]>, AtemError<I>> {
if self.state == AtemSocketState::Closed {
return Err(AtemError::SocketClosed);
}
let size = self.inner.recv(buf).await.map_err(|e| AtemError::Recv(e))?;
let p: AtemPacket<_> = buf[0..size].try_into()?;
debug!("RX: {}", p);
match self.state {
AtemSocketState::StateDump => {
// We've just connected and the state's being dumped to us.
if p.body().len() == 0 {
// State dump is done.
self.state = AtemSocketState::Established {
session_id: p.session_id(),
};
debug!("ATEM state dump done, session ID is {}", p.session_id());
self.ack(p.session_id(), p.local_sequence_number()).await?;
}
}
AtemSocketState::Established { session_id } => {
if p.flags().contains(PacketFlag::AckRequest) {
self.ack(session_id, p.local_sequence_number()).await?;
}
}
AtemSocketState::Closed => unreachable!(),
}
Ok(p)
}
/// Send a packet to the ATEM. The ATEM will respond with an ACK - it's up to the caller to handle packet retransmission
pub async fn send<T: AsRef<[u8]>>(&mut self, p: AtemPacket<T>) -> Result<(), AtemError<I>> {
self.inner
.send_to(p.inner().as_ref(), self.addr)
.await
.map_err(|e| AtemError::Send(e))
}
}

View File

@ -0,0 +1,4 @@
pub const COMMAND_CONNECT_HELLO: [u8; 20] = [
0x10, 0x14, 0x53, 0xab, 0x00, 0x00, 0x00, 0x00, 0x00, 0x3a, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
];

View File

@ -1,5 +1,3 @@
pub mod atem_field; mod atem_packet;
pub mod atem_packet;
#[cfg(feature = "std")]
pub mod atem_socket; pub mod atem_socket;
pub mod atem_socket_two; pub mod atem_util;

View File

@ -1,4 +1,4 @@
use std::{boxed::Box, collections::HashMap, fmt::Debug, string::String, sync::Arc, vec::Vec}; use std::{collections::HashMap, fmt::Debug, process::Command, sync::Arc};
use crate::{enums::ProtocolVersion, state::AtemState}; use crate::{enums::ProtocolVersion, state::AtemState};
@ -9,7 +9,7 @@ pub trait DeserializedCommand: Send + Sync + Debug {
pub trait CommandDeserializer: Send + Sync { pub trait CommandDeserializer: Send + Sync {
fn deserialize(&self, buffer: &[u8], version: &ProtocolVersion) fn deserialize(&self, buffer: &[u8], version: &ProtocolVersion)
-> Arc<dyn DeserializedCommand>; -> Arc<dyn DeserializedCommand>;
} }
pub trait SerializableCommand: Send + Sync { pub trait SerializableCommand: Send + Sync {

View File

@ -1,9 +1,4 @@
use std::{ use std::{ffi::CString, sync::Arc};
ffi::CString,
string::{String, ToString},
sync::Arc,
vec,
};
use crate::{ use crate::{
commands::command_base::{CommandDeserializer, DeserializedCommand}, commands::command_base::{CommandDeserializer, DeserializedCommand},

View File

@ -1,4 +1,4 @@
use std::{sync::Arc, vec}; use std::sync::Arc;
use crate::{ use crate::{
commands::command_base::{CommandDeserializer, DeserializedCommand}, commands::command_base::{CommandDeserializer, DeserializedCommand},

View File

@ -1,5 +1,4 @@
use derive_new::new; use std::sync::Arc;
use std::{sync::Arc, vec, vec::Vec};
use crate::{ use crate::{
commands::command_base::{ commands::command_base::{

View File

@ -1,8 +1,7 @@
use std::{boxed::Box, collections::VecDeque, sync::Arc}; use std::{collections::VecDeque, sync::Arc};
use crate::{ use crate::{
atem_lib::atem_packet::RawFields, commands::device_profile::version::{deserialize_version, DESERIALIZE_VERSION_RAW_NAME},
commands::device_profile::version::{DESERIALIZE_VERSION_RAW_NAME, deserialize_version},
enums::ProtocolVersion, enums::ProtocolVersion,
}; };
@ -10,44 +9,59 @@ use super::{
command_base::{CommandDeserializer, DeserializedCommand}, command_base::{CommandDeserializer, DeserializedCommand},
device_profile::{ device_profile::{
audio_mixer_config::{AudioMixerConfigDeserializer, DESERIALIZE_AUDIO_MIXER_CONFIG_NAME}, audio_mixer_config::{AudioMixerConfigDeserializer, DESERIALIZE_AUDIO_MIXER_CONFIG_NAME},
media_pool_config::{DESERIALIZE_MEDIA_POOL_CONFIG_NAME, MediaPoolConfigDeserializer}, media_pool_config::{MediaPoolConfigDeserializer, DESERIALIZE_MEDIA_POOL_CONFIG_NAME},
mix_effect_block_config::{ mix_effect_block_config::{
DESERIALIZE_MIX_EFFECT_BLOCK_CONFIG_NAME, MixEffectBlockConfigDeserializer, MixEffectBlockConfigDeserializer, DESERIALIZE_MIX_EFFECT_BLOCK_CONFIG_NAME,
}, },
multiviewer_config::{DESERIALIZE_MULTIVIEWER_NAME, MultiviewerConfigDeserializer}, multiviewer_config::{MultiviewerConfigDeserializer, DESERIALIZE_MULTIVIEWER_NAME},
product_identifier::{ product_identifier::{
DESERIALIZE_PRODUCT_IDENTIFIER_RAW_NAME, ProductIdentifierDeserializer, ProductIdentifierDeserializer, DESERIALIZE_PRODUCT_IDENTIFIER_RAW_NAME,
}, },
topology::{DESERIALIZE_TOPOLOGY_RAW_NAME, TopologyDeserializer}, topology::{TopologyDeserializer, DESERIALIZE_TOPOLOGY_RAW_NAME},
}, },
init_complete::{DESERIALIZE_INIT_COMPLETE_RAW_NAME, InitCompleteDeserializer}, init_complete::{InitCompleteDeserializer, DESERIALIZE_INIT_COMPLETE_RAW_NAME},
mix_effects::program_input::{DESERIALIZE_PROGRAM_INPUT_RAW_NAME, ProgramInputDeserializer}, mix_effects::program_input::{ProgramInputDeserializer, DESERIALIZE_PROGRAM_INPUT_RAW_NAME},
tally_by_source::{DESERIALIZE_TALLY_BY_SOURCE_RAW_NAME, TallyBySourceDeserializer}, tally_by_source::{TallyBySourceDeserializer, DESERIALIZE_TALLY_BY_SOURCE_RAW_NAME},
time::{DESERIALIZE_TIME_RAW_NAME, TimeDeserializer}, time::{TimeDeserializer, DESERIALIZE_TIME_RAW_NAME},
}; };
pub fn deserialize_commands<T: AsRef<[u8]>>( pub fn deserialize_commands(
payload: T, payload: &[u8],
version: &mut ProtocolVersion, version: &mut ProtocolVersion,
) -> VecDeque<Arc<dyn DeserializedCommand>> { ) -> VecDeque<Arc<dyn DeserializedCommand>> {
let mut parsed_commands: VecDeque<Arc<dyn DeserializedCommand>> = VecDeque::new(); let mut parsed_commands: VecDeque<Arc<dyn DeserializedCommand>> = VecDeque::new();
let mut head = 0;
for field in RawFields::new(payload.as_ref()) { while payload.len() > head + 8 {
let name: &str = field.r#type.try_into().unwrap(); let length = u16::from_be_bytes([payload[head], payload[head + 1]]) as usize;
log::debug!("Received command {} with length {}", name, field.data.len(),); let Ok(name) = String::from_utf8(payload[(head + 4)..(head + 8)].to_vec()) else {
break;
};
if field.r#type == DESERIALIZE_VERSION_RAW_NAME { if length < 8 {
let version_command = deserialize_version(field.data); break;
}
log::debug!("Received command {} with length {}", name, length);
let command_buffer = &payload[head + 8..head + length];
if name == DESERIALIZE_VERSION_RAW_NAME {
let version_command = deserialize_version(command_buffer);
*version = version_command.version.clone(); *version = version_command.version.clone();
log::info!("Switched to protocol version {}", version); log::info!("Switched to protocol version {}", version);
parsed_commands.push_back(Arc::new(version_command)); parsed_commands.push_back(Arc::new(version_command));
} else if let Some(deserializer) = command_deserializer_from_string(name) { } else if let Some(deserializer) = command_deserializer_from_string(name.as_str()) {
let deserialized_command = deserializer.deserialize(field.data, version); let deserialized_command = deserializer.deserialize(command_buffer, version);
log::debug!("Received {:?}", deserialized_command); log::debug!("Received {:?}", deserialized_command);
parsed_commands.push_back(deserialized_command); parsed_commands.push_back(deserialized_command);
} else { } else {
log::warn!("Received command {name} for which there is no deserializer."); log::warn!("Received command {name} for which there is no deserializer.");
// TODO: Remove!
todo!("Write deserializer for {name}.");
} }
head += length;
} }
parsed_commands parsed_commands

View File

@ -1,4 +1,3 @@
use derive_new::new;
use std::{collections::HashMap, sync::Arc}; use std::{collections::HashMap, sync::Arc};
use crate::enums::ProtocolVersion; use crate::enums::ProtocolVersion;

View File

@ -1,16 +1,11 @@
#![no_std] #[macro_use]
extern crate derive_new;
#[macro_use]
extern crate derive_getters;
#[cfg(feature = "std")]
extern crate std;
#[cfg(feature = "std")]
pub mod atem; pub mod atem;
pub mod atem_lib; pub mod atem_lib;
#[cfg(feature = "std")]
pub mod commands; pub mod commands;
#[cfg(feature = "std")]
pub mod enums; pub mod enums;
#[cfg(feature = "std")]
pub mod state; pub mod state;
#[cfg(feature = "std")]
pub mod tally; pub mod tally;

View File

@ -1,7 +1,3 @@
use std::{string::String, vec::Vec};
use derive_getters::Getters;
use derive_new::new;
#[derive(Clone, PartialEq, Getters, new, Default)] #[derive(Clone, PartialEq, Getters, new, Default)]
pub struct MacroPlayerState { pub struct MacroPlayerState {
pub is_running: bool, pub is_running: bool,

View File

@ -1,8 +1,6 @@
use std::collections::HashMap; use std::collections::HashMap;
use crate::enums::{AudioMixOption, AudioSourceType, ExternalPortType}; use crate::enums::{AudioMixOption, AudioSourceType, ExternalPortType};
use derive_getters::Getters;
use derive_new::new;
pub type AudioChannel = ClassicAudioChannel; pub type AudioChannel = ClassicAudioChannel;
pub type AudioMasterChannel = ClassicAudioMasterChannel; pub type AudioMasterChannel = ClassicAudioMasterChannel;

View File

@ -1,5 +1,3 @@
use derive_getters::Getters;
use derive_new::new;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct ColorGeneratorState { pub struct ColorGeneratorState {
pub hue: u64, pub hue: u64,

View File

@ -1,5 +1,3 @@
use derive_getters::Getters;
use derive_new::new;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct Timecode { pub struct Timecode {
pub hours: u64, pub hours: u64,

View File

@ -1,6 +1,4 @@
use derive_getters::Getters; use std::collections::HashMap;
use derive_new::new;
use std::{collections::HashMap, string::String, vec::Vec};
use crate::enums::{ use crate::enums::{
ExternalPortType, FairlightAnalogInputLevel, FairlightAudioMixOption, FairlightAudioSourceType, ExternalPortType, FairlightAnalogInputLevel, FairlightAudioMixOption, FairlightAudioSourceType,

View File

@ -1,8 +1,4 @@
use std::{string::String, vec::Vec};
use crate::enums::{Model, ProtocolVersion}; use crate::enums::{Model, ProtocolVersion};
use derive_getters::Getters;
use derive_new::new;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct AtemCapabilites { pub struct AtemCapabilites {

View File

@ -1,8 +1,4 @@
use crate::enums::{ExternalPortType, InternalPortType, MeAvailability, SourceAvailability}; use crate::enums::{ExternalPortType, InternalPortType, MeAvailability, SourceAvailability};
use derive_getters::Getters;
use derive_new::new;
use std::string::String;
use std::vec::Vec;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct InputChannel { pub struct InputChannel {

View File

@ -1,8 +1,4 @@
use crate::enums; use crate::enums;
use derive_getters::Getters;
use derive_new::new;
use std::string::String;
use std::vec::Vec;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct MediaPlayer { pub struct MediaPlayer {

View File

@ -1,7 +1,4 @@
use derive_getters::Getters;
use derive_new::new;
use std::collections::HashMap; use std::collections::HashMap;
use std::string::String;
use crate::enums::{RecordingDiskStatus, RecordingError, RecordingStatus}; use crate::enums::{RecordingDiskStatus, RecordingError, RecordingStatus};

View File

@ -1,7 +1,4 @@
use crate::enums::{MultiViewerLayout, VideoMode}; use crate::enums::{MultiViewerLayout, VideoMode};
use derive_getters::Getters;
use derive_new::new;
use std::vec::Vec;
pub trait MultiViewerSourceState { pub trait MultiViewerSourceState {
fn get_source(&self) -> u64; fn get_source(&self) -> u64;

View File

@ -1,7 +1,4 @@
use crate::enums::{StreamingError, StreamingStatus}; use crate::enums::{StreamingError, StreamingStatus};
use derive_getters::Getters;
use derive_new::new;
use std::string::String;
use super::common::Timecode; use super::common::Timecode;

View File

@ -1,10 +1,9 @@
use crate::enums::{TransitionSelection, TransitionStyle}; use crate::enums::{TransitionSelection, TransitionStyle};
use std::vec;
use super::{ use super::{
AtemState,
settings::MultiViewer, settings::MultiViewer,
video::{MixEffect, TransitionPosition, TransitionProperties, TransitionSettings}, video::{MixEffect, TransitionPosition, TransitionProperties, TransitionSettings},
AtemState,
}; };
pub fn create() -> AtemState { pub fn create() -> AtemState {

View File

@ -1,5 +1,3 @@
use derive_getters::Getters;
use derive_new::new;
pub trait DownstreamKeyerBase { pub trait DownstreamKeyerBase {
fn get_in_transition(&self) -> bool; fn get_in_transition(&self) -> bool;
fn get_remaining_frames(&self) -> f64; fn get_remaining_frames(&self) -> f64;

View File

@ -1,7 +1,4 @@
use crate::enums; use crate::enums;
use derive_getters::Getters;
use derive_new::new;
use std::vec::Vec;
mod downstream_keyers; mod downstream_keyers;
mod super_source; mod super_source;

View File

@ -1,6 +1,4 @@
use crate::enums; use crate::enums;
use derive_getters::Getters;
use derive_new::new;
#[derive(Clone, PartialEq, Getters, new)] #[derive(Clone, PartialEq, Getters, new)]
pub struct SuperSourceBox { pub struct SuperSourceBox {

View File

@ -1,4 +1,3 @@
use std::string::String;
#[derive(Debug)] #[derive(Debug)]
pub struct TallyEvent { pub struct TallyEvent {
tally_state: TallyState, tally_state: TallyState,

View File

@ -1,8 +1,9 @@
[package] [package]
name = "atem-test" name = "atem-test"
version = "0.1.0" version = "0.1.0"
edition = "2024" edition = "2021"
default-run = "main"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[dependencies] [dependencies]
atem-connection-rs = { path = "../atem-connection-rs" } atem-connection-rs = { path = "../atem-connection-rs" }

View File

@ -1,48 +0,0 @@
use atem_connection_rs::atem_lib::{atem_field::PrgI, atem_socket_two::tokio::TokioBackend};
use clap::Parser;
use color_eyre::Report;
/// ATEM Rust Library Test App
#[derive(Parser, Debug)]
#[command(author, version, about, long_about = None)]
struct Args {
/// IP of the ATEM to connect to
#[arg(short, long)]
ip: String,
}
#[tokio::main]
async fn main() {
let args = Args::parse();
setup_logging().unwrap();
let mut s = TokioBackend::new(args.ip.parse().unwrap()).await.unwrap();
s.connect().await.expect("aa");
let mut buf = [0u8; 2048];
loop {
let p = s.recv(&mut buf).await.unwrap();
let prg: Vec<PrgI> = p.fields::<PrgI>().collect::<Result<_, _>>().unwrap();
if !prg.is_empty() {
println!("{:?}", prg);
}
}
}
fn setup_logging() -> Result<(), Report> {
if std::env::var("RUST_LIB_BACKTRACE").is_err() {
unsafe {
std::env::set_var("RUST_LIB_BACKTRACE", "1");
}
}
color_eyre::install()?;
if std::env::var("RUST_LOG").is_err() {
unsafe {
std::env::set_var("RUST_LOG", "debug");
}
}
env_logger::init();
Ok(())
}

View File

@ -36,7 +36,7 @@ async fn main() {
let (atem_event_tx, atem_event_rx) = tokio::sync::mpsc::unbounded_channel(); let (atem_event_tx, atem_event_rx) = tokio::sync::mpsc::unbounded_channel();
let cancel = CancellationToken::new(); let cancel = CancellationToken::new();
let atem_socket = AtemSocket::new(socket_message_rx, atem_event_tx); let mut atem_socket = AtemSocket::new(socket_message_rx, atem_event_tx);
let atem = Arc::new(Atem::new(atem_socket, socket_message_tx)); let atem = Arc::new(Atem::new(atem_socket, socket_message_tx));
let atem_thread = atem.clone(); let atem_thread = atem.clone();
@ -69,16 +69,12 @@ async fn main() {
fn setup_logging() -> Result<(), Report> { fn setup_logging() -> Result<(), Report> {
if std::env::var("RUST_LIB_BACKTRACE").is_err() { if std::env::var("RUST_LIB_BACKTRACE").is_err() {
unsafe { std::env::set_var("RUST_LIB_BACKTRACE", "1");
std::env::set_var("RUST_LIB_BACKTRACE", "1");
}
} }
color_eyre::install()?; color_eyre::install()?;
if std::env::var("RUST_LOG").is_err() { if std::env::var("RUST_LOG").is_err() {
unsafe { std::env::set_var("RUST_LOG", "debug");
std::env::set_var("RUST_LOG", "debug");
}
} }
env_logger::init(); env_logger::init();

View File

@ -1,9 +1,27 @@
{ {
"nodes": { "nodes": {
"naersk": { "devshell": {
"inputs": { "inputs": {
"nixpkgs": "nixpkgs" "nixpkgs": "nixpkgs"
}, },
"locked": {
"lastModified": 1741473158,
"narHash": "sha256-kWNaq6wQUbUMlPgw8Y+9/9wP0F8SHkjy24/mN3UAppg=",
"owner": "numtide",
"repo": "devshell",
"rev": "7c9e793ebe66bcba8292989a68c0419b737a22a0",
"type": "github"
},
"original": {
"owner": "numtide",
"repo": "devshell",
"type": "github"
}
},
"naersk": {
"inputs": {
"nixpkgs": "nixpkgs_2"
},
"locked": { "locked": {
"lastModified": 1745925850, "lastModified": 1745925850,
"narHash": "sha256-cyAAMal0aPrlb1NgzMxZqeN1mAJ2pJseDhm2m6Um8T0=", "narHash": "sha256-cyAAMal0aPrlb1NgzMxZqeN1mAJ2pJseDhm2m6Um8T0=",
@ -20,11 +38,11 @@
}, },
"nixpkgs": { "nixpkgs": {
"locked": { "locked": {
"lastModified": 1749619289, "lastModified": 1722073938,
"narHash": "sha256-qX6gXVjaCXXbcn6A9eSLUf8Fm07MgPGe5ir3++y2O1Q=", "narHash": "sha256-OpX0StkL8vpXyWOGUD6G+MA26wAXK6SpT94kLJXo6B4=",
"owner": "NixOS", "owner": "NixOS",
"repo": "nixpkgs", "repo": "nixpkgs",
"rev": "f72be405a10668b8b00937b452f2145244103ebc", "rev": "e36e9f57337d0ff0cf77aceb58af4c805472bfae",
"type": "github" "type": "github"
}, },
"original": { "original": {
@ -36,10 +54,25 @@
}, },
"nixpkgs_2": { "nixpkgs_2": {
"locked": { "locked": {
"lastModified": 1733412085, "lastModified": 1749401433,
"narHash": "sha256-FillH0qdWDt/nlO6ED7h4cmN+G9uXwGjwmCnHs0QVYM=", "narHash": "sha256-HXIQzULIG/MEUW2Q/Ss47oE3QrjxvpUX7gUl4Xp6lnc=",
"path": "/nix/store/5wvrlpfrwlmg24ywkybvidyzcvki6bwx-source", "owner": "NixOS",
"rev": "4dc2fc4e62dbf62b84132fe526356fbac7b03541", "repo": "nixpkgs",
"rev": "08fcb0dcb59df0344652b38ea6326a2d8271baff",
"type": "github"
},
"original": {
"owner": "NixOS",
"ref": "nixpkgs-unstable",
"repo": "nixpkgs",
"type": "github"
}
},
"nixpkgs_3": {
"locked": {
"lastModified": 0,
"narHash": "sha256-DDe16FJk18sadknQKKG/9FbwEro7A57tg9vB5kxZ8kY=",
"path": "/nix/store/2d1ahim48jhzg4bbm97mvjlb4p7fpan3-source",
"type": "path" "type": "path"
}, },
"original": { "original": {
@ -47,7 +80,7 @@
"type": "indirect" "type": "indirect"
} }
}, },
"nixpkgs_3": { "nixpkgs_4": {
"locked": { "locked": {
"lastModified": 1744536153, "lastModified": 1744536153,
"narHash": "sha256-awS2zRgF4uTwrOKwwiJcByDzDOdo3Q1rPZbiHQg/N38=", "narHash": "sha256-awS2zRgF4uTwrOKwwiJcByDzDOdo3Q1rPZbiHQg/N38=",
@ -65,22 +98,23 @@
}, },
"root": { "root": {
"inputs": { "inputs": {
"devshell": "devshell",
"naersk": "naersk", "naersk": "naersk",
"nixpkgs": "nixpkgs_2", "nixpkgs": "nixpkgs_3",
"rust-overlay": "rust-overlay", "rust-overlay": "rust-overlay",
"utils": "utils" "utils": "utils"
} }
}, },
"rust-overlay": { "rust-overlay": {
"inputs": { "inputs": {
"nixpkgs": "nixpkgs_3" "nixpkgs": "nixpkgs_4"
}, },
"locked": { "locked": {
"lastModified": 1749695868, "lastModified": 1749436897,
"narHash": "sha256-debjTLOyqqsYOUuUGQsAHskFXH5+Kx2t3dOo/FCoNRA=", "narHash": "sha256-OkDtaCGQQVwVFz5HWfbmrMJR99sFIMXHCHEYXzUJEJY=",
"owner": "oxalica", "owner": "oxalica",
"repo": "rust-overlay", "repo": "rust-overlay",
"rev": "55f914d5228b5c8120e9e0f9698ed5b7214d09cd", "rev": "e7876c387e35dc834838aff254d8e74cf5bd4f19",
"type": "github" "type": "github"
}, },
"original": { "original": {

View File

@ -3,6 +3,7 @@
inputs = { inputs = {
utils.url = "github:numtide/flake-utils"; utils.url = "github:numtide/flake-utils";
devshell.url = "github:numtide/devshell";
naersk.url = "github:nix-community/naersk"; naersk.url = "github:nix-community/naersk";
rust-overlay.url = "github:oxalica/rust-overlay"; rust-overlay.url = "github:oxalica/rust-overlay";
}; };
@ -12,6 +13,7 @@
nixpkgs, nixpkgs,
utils, utils,
naersk, naersk,
devshell,
rust-overlay, rust-overlay,
}: }:
utils.lib.eachDefaultSystem (system: let utils.lib.eachDefaultSystem (system: let
@ -33,10 +35,18 @@
apps.default = utils.lib.mkApp {drv = packages.default;}; apps.default = utils.lib.mkApp {drv = packages.default;};
devShells.default = pkgs.mkShell { # Provide a dev env with rust and rust-analyzer
packages = with pkgs; [(rust.override {extensions = ["rust-src"];}) rust-analyzer gcc]; devShells.default = let
pkgs = import nixpkgs {
inherit system;
overlays = [devshell.overlays.default];
}; };
formatter = pkgs.nixfmt-rfc-style; in
pkgs.devshell.mkShell {
motd = "Hello you wonderful person, I hope you are having a lovely day 💜";
packages = with pkgs; [(rust.override {extensions = ["rust-src"];}) rust-analyzer gcc];
};
formatter = pkgs.alejandra;
}); });
} }