Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
b919738
refactor(plugins): drop unused sherpa-onnx wave FFI bindings
streamkit-devin Sep 19, 2026
ee4a619
refactor(client): drop unread pipeline_path from DynamicSession
streamkit-devin Sep 19, 2026
644f617
refactor(client): drop unread total_success snapshot field
streamkit-devin Sep 19, 2026
314c68d
refactor(server): drop unused permission helpers
streamkit-devin Sep 19, 2026
aa3d8fe
refactor(nodes): drop stale allow(unused_variables) in blit
streamkit-devin Sep 19, 2026
ee9b363
refactor(server): use MoQ route session_id in gateway logs
streamkit-devin Sep 19, 2026
0777cb8
refactor(server): gate CreateMoqTokenRequest behind moq feature
streamkit-devin Sep 19, 2026
026638a
refactor(server): cfg-gate MoQ claim types instead of dead_code
streamkit-devin Sep 19, 2026
b8b8c8f
refactor(server): tighten auth-core dead_code suppressions
streamkit-devin Sep 19, 2026
dc6f926
refactor(server): drop dead MaybeAuth scaffolding in auth extractor
streamkit-devin Sep 19, 2026
1bd3b2e
refactor(server): drop never-constructed KeyNotFound store error
streamkit-devin Sep 19, 2026
c8c35d9
refactor(server): drop unread root field from MoqAuthContext
streamkit-devin Sep 19, 2026
6c38b88
docs(nodes): add rationale for js_to_packet lint suppression
streamkit-devin Sep 19, 2026
d1feb1f
refactor(plugins): share SentenceSplitter via streamkit_core::text
streamkit-devin Sep 19, 2026
2826ead
refactor(plugins): dedup sherpa-onnx FFI bindings into shared crate
streamkit-devin Sep 19, 2026
fb5f9ef
refactor(plugins): dedup SileroVAD into shared crate
streamkit-devin Sep 19, 2026
5298909
chore(deps): bump rustls to 0.23.45 for RUSTSEC-2026-0285
streamkit-devin Sep 19, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 0 additions & 3 deletions apps/skit-cli/src/load_test/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,6 @@ impl MetricsCollector {
MetricsSnapshot {
elapsed,
total_ops,
total_success,
total_failures,
throughput,
success_rate,
Expand Down Expand Up @@ -255,8 +254,6 @@ impl MetricsCollector {
pub struct MetricsSnapshot {
pub elapsed: Duration,
pub total_ops: usize,
#[allow(dead_code)]
pub total_success: usize,
pub total_failures: usize,
pub throughput: f64,
pub success_rate: f64,
Expand Down
3 changes: 0 additions & 3 deletions apps/skit-cli/src/load_test/workers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,8 +85,6 @@ pub async fn oneshot_worker(

pub struct DynamicSession {
pub session_id: String,
#[allow(dead_code)]
pub pipeline_path: String,
pub tunable_node_ids: Vec<String>,
}

Expand Down Expand Up @@ -448,7 +446,6 @@ pub async fn session_creator_worker(
let _ = session_tx
.send(DynamicSession {
session_id,
pipeline_path: pipeline_path.clone(),
tunable_node_ids,
})
.await;
Expand Down
11 changes: 6 additions & 5 deletions apps/skit/src/auth/claims.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use serde::{Deserialize, Serialize};
pub const AUD_API: &str = "skit-api";

/// Audience value for MoQ tokens.
#[allow(dead_code)]
#[cfg(feature = "moq")]
pub const AUD_MOQ: &str = "skit-moq";

/// JWT claims for API tokens (HTTP API and WebSocket control plane).
Expand Down Expand Up @@ -79,8 +79,8 @@ impl ApiClaims {
/// - `[""]` (empty string in array) = all broadcasts allowed
/// - `[]` (empty array) = no broadcasts allowed
/// - `["foo", "bar"]` = broadcasts starting with "foo" or "bar" allowed
#[cfg(feature = "moq")]
#[derive(Debug, Clone, Serialize, Deserialize)]
#[allow(dead_code)]
pub struct MoqClaims {
/// Must be [`AUD_MOQ`].
pub aud: String,
Expand All @@ -105,7 +105,7 @@ pub struct MoqClaims {
pub jti: String,
}

#[allow(dead_code)]
#[cfg(feature = "moq")]
impl MoqClaims {
/// Validate claims structure (not cryptographic verification).
///
Expand Down Expand Up @@ -141,8 +141,8 @@ pub enum ClaimsValidationError {
#[error("Missing role claim")]
MissingRole,

#[cfg(feature = "moq")]
#[error("Missing root claim")]
#[allow(dead_code)]
MissingRoot,
}

Expand All @@ -163,7 +163,7 @@ mod tests {
assert!(valid.validate().is_ok());

// Wrong audience
let wrong_aud = ApiClaims { aud: AUD_MOQ.to_string(), ..valid.clone() };
let wrong_aud = ApiClaims { aud: "skit-moq".to_string(), ..valid.clone() };
assert!(matches!(wrong_aud.validate(), Err(ClaimsValidationError::InvalidAudience { .. })));

// Missing jti
Expand All @@ -175,6 +175,7 @@ mod tests {
assert!(matches!(no_role.validate(), Err(ClaimsValidationError::MissingRole)));
}

#[cfg(feature = "moq")]
#[test]
fn test_moq_claims_validation() {
let valid = MoqClaims {
Expand Down
58 changes: 0 additions & 58 deletions apps/skit/src/auth/extractor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,53 +15,9 @@ use axum::http::{HeaderMap, StatusCode};
pub struct AuthContext {
pub claims: ApiClaims,
pub role: String,
#[allow(dead_code)]
pub permissions: Permissions,
}

#[allow(dead_code)]
impl AuthContext {
/// JWT ID (used for revocation tracking).
pub fn jti(&self) -> &str {
&self.claims.jti
}

/// Subject (token holder identifier).
pub fn sub(&self) -> &str {
&self.claims.sub
}
}

/// Optional auth context — never fails; contains `None` when unauthenticated.
#[derive(Debug, Clone)]
#[allow(dead_code)]
pub struct MaybeAuth(pub Option<AuthContext>);

#[allow(dead_code)]
impl MaybeAuth {
pub const fn context(&self) -> Option<&AuthContext> {
self.0.as_ref()
}

pub const fn is_authenticated(&self) -> bool {
self.0.is_some()
}

#[allow(clippy::ref_option)]
pub const fn as_option(&self) -> &Option<AuthContext> {
&self.0
}

/// Unwrap or return an unauthorized error.
///
/// # Errors
///
/// Returns `(StatusCode::UNAUTHORIZED, ...)` if not authenticated.
pub fn require(self) -> Result<AuthContext, (StatusCode, String)> {
self.0.ok_or_else(|| (StatusCode::UNAUTHORIZED, "Authentication required".to_string()))
}
}

/// Extract token from Authorization header or cookie.
///
/// Checks the Authorization header first (Bearer token format),
Expand Down Expand Up @@ -214,18 +170,4 @@ mod tests {
let token = extract_token(&headers, &config);
assert!(token.is_none());
}

#[test]
fn test_maybe_auth_require() {
let auth = MaybeAuth(None);
assert!(auth.require().is_err());

let ctx = AuthContext {
claims: ApiClaims::anonymous("admin"),
role: "admin".to_string(),
permissions: crate::permissions::Permissions::admin(),
};
let auth = MaybeAuth(Some(ctx));
assert!(auth.require().is_ok());
}
}
2 changes: 1 addition & 1 deletion apps/skit/src/auth/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,8 @@ pub struct CreateApiTokenRequest {
#[serde(default)]
pub ttl_secs: Option<u64>,
}
#[cfg(feature = "moq")]
#[derive(Debug, Deserialize, Serialize)]
#[allow(dead_code)]
pub struct CreateMoqTokenRequest {
pub root: String,
#[serde(default)]
Expand Down
14 changes: 1 addition & 13 deletions apps/skit/src/auth/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,6 @@ use tracing::{debug, info};

/// Errors that can occur during authentication.
#[derive(Debug, thiserror::Error)]
#[allow(dead_code)]
pub enum AuthError {
#[error("Authentication is disabled")]
Disabled,
Expand All @@ -78,15 +77,6 @@ pub enum AuthError {
#[error("Token not found in metadata store (not minted by this server)")]
UnknownToken,

#[error("Token has been revoked")]
Revoked,

#[error("Token expired")]
Expired,

#[error("Invalid audience: expected {expected}, got {actual}")]
InvalidAudience { expected: String, actual: String },

#[error("TTL exceeds maximum allowed ({max} seconds)")]
TtlExceedsMax { max: u64 },

Expand Down Expand Up @@ -215,7 +205,7 @@ impl AuthState {
}

/// Returns false if auth is disabled or the revocation store is unavailable.
#[allow(dead_code)]
#[cfg(feature = "moq")]
pub fn is_revoked(&self, token_hash: &str) -> bool {
self.revocation_store.as_ref().is_some_and(|store| store.is_revoked(token_hash))
}
Expand All @@ -224,7 +214,6 @@ impl AuthState {
self.token_metadata_store.as_ref()
}

#[allow(dead_code)]
pub fn key_provider(&self) -> Option<&Arc<dyn KeyProvider>> {
self.key_provider.as_ref()
}
Expand Down Expand Up @@ -472,7 +461,6 @@ impl AuthState {
}

/// Whether auth should be enabled based on config and bind address.
#[allow(dead_code)]
pub const fn should_enable(config: &AuthConfig, bind_addr: &std::net::SocketAddr) -> bool {
match config.mode {
AuthMode::Auto => !bind_addr.ip().is_loopback(),
Expand Down
5 changes: 1 addition & 4 deletions apps/skit/src/auth/moq.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,6 @@ use streamkit_core::moq_gateway::MoqAuthChecker;
/// Verified MoQ auth context with permissions reduced by connection path depth.
#[derive(Debug, Clone)]
pub struct MoqAuthContext {
/// The actual connection path (after root validation)
#[allow(dead_code)]
pub root: PathOwned,
/// Reduced subscribe permissions (broadcast paths relative to connection)
pub subscribe: Vec<PathOwned>,
/// Reduced publish permissions (broadcast paths relative to connection)
Expand Down Expand Up @@ -86,7 +83,7 @@ pub fn verify_moq_token(
let subscribe = claims.subscribe.iter().filter_map(|p| reduce_permission(p, &suffix)).collect();
let publish = claims.publish.iter().filter_map(|p| reduce_permission(p, &suffix)).collect();

Ok(MoqAuthContext { root: url_path.to_owned(), subscribe, publish })
Ok(MoqAuthContext { subscribe, publish })
}

/// Reduce a permission path based on connection suffix.
Expand Down
4 changes: 0 additions & 4 deletions apps/skit/src/auth/stores/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,10 +44,6 @@ pub enum AuthStoreError {
#[error("Base64 decode error: {0}")]
Base64(#[from] base64::DecodeError),

#[error("Key not found: {0}")]
#[allow(dead_code)]
KeyNotFound(String),

#[error("Invalid file permissions on {path}: expected 0600, got {actual:o}")]
InsecurePermissions { path: String, actual: u32 },

Expand Down
27 changes: 15 additions & 12 deletions apps/skit/src/moq_gateway.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ use tracing::{debug, error, info, warn};
/// A route registration from a path pattern to a connection receiver
struct Route {
/// The session ID that owns this route
#[allow(dead_code)]
session_id: String,

/// Channel to send accepted connections to the node
Expand Down Expand Up @@ -94,17 +93,17 @@ impl MoqGateway {
// The UI currently connects WebTransport before creating the dynamic session; when the
// session starts, it registers its routes. Without waiting here, the server would drop
// the connection before the route exists.
let connection_tx: Option<mpsc::UnboundedSender<MoqConnection>> = {
let route: Option<(String, mpsc::UnboundedSender<MoqConnection>)> = {
const MAX_WAIT: Duration = Duration::from_secs(30);
const POLL_INTERVAL: Duration = Duration::from_millis(200);

let mut waited = Duration::from_secs(0);
loop {
if let Some(tx) = {
if let Some(entry) = {
let routes = self.routes.read().await;
routes.get(&path).map(|r| r.connection_tx.clone())
routes.get(&path).map(|r| (r.session_id.clone(), r.connection_tx.clone()))
} {
break Some(tx);
break Some(entry);
}

if waited >= MAX_WAIT {
Expand All @@ -119,7 +118,7 @@ impl MoqGateway {
}
};

if let Some(connection_tx) = connection_tx {
if let Some((session_id, connection_tx)) = route {
let (response_tx, response_rx) = oneshot::channel();

// Type-erase the moq-native Request
Expand All @@ -129,22 +128,22 @@ impl MoqGateway {
MoqConnection { path: path.clone(), session: session_boxed, response_tx, auth };

if connection_tx.send(conn).is_err() {
error!(path = %path, "Failed to send connection to node (channel closed)");
error!(path = %path, session_id = %session_id, "Failed to send connection to node (channel closed)");
return Err("Node disconnected".to_string());
}

// Wait for node to accept or reject
match response_rx.await {
Ok(MoqConnectionResult::Accepted) => {
info!(path = %path, "Connection accepted by node");
info!(path = %path, session_id = %session_id, "Connection accepted by node");
Ok(())
},
Ok(MoqConnectionResult::Rejected(reason)) => {
warn!(path = %path, reason = %reason, "Connection rejected by node");
warn!(path = %path, session_id = %session_id, reason = %reason, "Connection rejected by node");
Err(reason)
},
Err(_) => {
error!(path = %path, "Node dropped connection without responding");
error!(path = %path, session_id = %session_id, "Node dropped connection without responding");
Err("Node did not respond".to_string())
},
}
Expand Down Expand Up @@ -209,9 +208,13 @@ impl MoqGatewayTrait for MoqGateway {

async fn unregister_route(&self, path_pattern: &str) {
let mut routes = self.routes.write().await;
if routes.remove(path_pattern).is_some() {
if let Some(route) = routes.remove(path_pattern) {
self.route_notify.notify_waiters();
info!(path_pattern = %path_pattern, "Unregistered MoQ route");
info!(
path_pattern = %path_pattern,
session_id = %route.session_id,
"Unregistered MoQ route"
);
}
}
}
Expand Down
24 changes: 0 additions & 24 deletions apps/skit/src/permissions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -436,28 +436,13 @@ impl PermissionsConfig {
)
}

/// Get the default role permissions
#[allow(dead_code)]
pub fn get_default(&self) -> Permissions {
self.get_role(&self.default_role)
}

/// Check if we can accept a new session (global limit check)
pub const fn can_accept_session(&self, current_count: usize) -> bool {
match self.max_concurrent_sessions {
None => true,
Some(max) => current_count < max,
}
}

/// Check if we can accept a new oneshot pipeline (global limit check)
#[allow(dead_code)]
pub const fn can_accept_oneshot(&self, current_count: usize) -> bool {
match self.max_concurrent_oneshots {
None => true,
Some(max) => current_count < max,
}
}
}

#[cfg(test)]
Expand Down Expand Up @@ -544,15 +529,6 @@ mod tests {
assert!(!config.can_accept_session(11));
}

#[test]
fn test_global_oneshot_limits() {
let config = PermissionsConfig { max_concurrent_oneshots: Some(5), ..Default::default() };

assert!(config.can_accept_oneshot(0));
assert!(config.can_accept_oneshot(4));
assert!(!config.can_accept_oneshot(5));
}

#[test]
fn test_user_role_defaults() {
let user = Permissions::user();
Expand Down
Loading
Loading