From 29cdc8c8d8a0a80e3de8f078568f3fc6f60616a0 Mon Sep 17 00:00:00 2001 From: Will Killian Date: Fri, 2 Oct 2026 15:15:24 -0400 Subject: [PATCH] refactor: address all open Sonar findings on main Signed-off-by: Will Killian --- crates/cli/src/daemon/broker/server/socket.rs | 80 +++--- crates/cli/src/daemon/mcp/mod.rs | 125 +++++----- .../tests/fixtures/native_plugin/src/lib.rs | 232 +++++++++--------- go/nemo_relay/llm_test.go | 34 ++- 4 files changed, 254 insertions(+), 217 deletions(-) diff --git a/crates/cli/src/daemon/broker/server/socket.rs b/crates/cli/src/daemon/broker/server/socket.rs index 5f1ec8cab..e73fa8453 100644 --- a/crates/cli/src/daemon/broker/server/socket.rs +++ b/crates/cli/src/daemon/broker/server/socket.rs @@ -786,45 +786,55 @@ async fn disconnected( log_control_disconnected(role, &id, reason); state.sockets.changed.notify_waiters(); drop(_transaction); - tokio::spawn(async move { - tokio::time::sleep_until(deadline).await; - let _transaction = state.sockets.transaction(role, &id).await; - let current = lock(&state.sockets.peers) - .get(&key(role, &id)) - .is_some_and(|p| p.generation == generation && p.sender.is_none()); - if !current { - return; - } - if role == ComponentRole::Mcp { - let session = lock(&state.mcp_sessions).remove(&id); - if let Some(session) = session - && let Ok(session_id) = McpSessionId::new(id.clone()) - && let Ok(action) = state.registry.release_mcp( + tokio::spawn(expire_disconnected_peer( + state, role, id, generation, deadline, + )); +} + +async fn expire_disconnected_peer( + state: Arc, + role: ComponentRole, + id: String, + generation: String, + deadline: tokio::time::Instant, +) { + tokio::time::sleep_until(deadline).await; + let _transaction = state.sockets.transaction(role, &id).await; + let current = lock(&state.sockets.peers) + .get(&key(role, &id)) + .is_some_and(|p| p.generation == generation && p.sender.is_none()); + if !current { + return; + } + if role == ComponentRole::Mcp { + let session = lock(&state.mcp_sessions).remove(&id); + if let Some(session) = session + && let Ok(session_id) = McpSessionId::new(id.clone()) + && let Ok(action) = state.registry.release_mcp( + session.fingerprint, + &session_id, + now_unix_ms().saturating_add(DRAIN_LIFETIME_MS), + ) + { + if !session.released { + log_mcp_removed( session.fingerprint, &session_id, - now_unix_ms().saturating_add(DRAIN_LIFETIME_MS), - ) - { - if !session.released { - log_mcp_removed( - session.fingerprint, - &session_id, - "control_disconnect_timeout", - &action, - ); - } - handle_release_action(state.clone(), session.fingerprint, action); - } - lock(&state.pending_directives).remove(&id); - } else { - let session = lock(&state.worker_sessions).remove(&id); - if let Some(session) = session { - cleanup_worker_session(Arc::clone(&state), id.clone(), session).await; + "control_disconnect_timeout", + &action, + ); } + handle_release_action(state.clone(), session.fingerprint, action); } - lock(&state.sockets.peers).remove(&key(role, &id)); - state.sockets.changed.notify_waiters(); - }); + lock(&state.pending_directives).remove(&id); + } else { + let session = lock(&state.worker_sessions).remove(&id); + if let Some(session) = session { + cleanup_worker_session(Arc::clone(&state), id.clone(), session).await; + } + } + lock(&state.sockets.peers).remove(&key(role, &id)); + state.sockets.changed.notify_waiters(); } /// Complete worker cleanup after its session has been removed from admission state. diff --git a/crates/cli/src/daemon/mcp/mod.rs b/crates/cli/src/daemon/mcp/mod.rs index c9f52e8b4..8da5f21de 100644 --- a/crates/cli/src/daemon/mcp/mod.rs +++ b/crates/cli/src/daemon/mcp/mod.rs @@ -434,64 +434,10 @@ async fn supervise_route( BrokerDirective::LaunchWorker { .. } => { let bootstrap = WorkerBootstrap::from_directive(directive.clone()) .expect("launch directive was matched"); - let already_launched = launched - .as_ref() - .is_some_and(|(id, _, _)| id == &bootstrap.activation_id); - if !already_launched { - cleanup_pending_launch(lease, &mut launched).await?; - match launch_worker(&lease.daemon_origin, &bootstrap).await { - Ok(child) => { - launched = Some(( - bootstrap.activation_id.clone(), - child, - tokio::time::Instant::now(), - )); - } - Err(error) => { - report_activation_failed( - lease, - &bootstrap.activation_id, - error.failure_reason, - &error.source, - ) - .await?; - directive = refresh_registration(lease).await?.directive; - continue; - } - } - } - if let Some((activation_id, child, _)) = launched.as_mut() - && activation_id == &bootstrap.activation_id - && let Some(status) = child.child.try_wait().map_err(CliError::Io)? - { - let error = CliError::Launch(format!( - "activated worker exited before readiness with {status}" - )); - report_activation_failed( - lease, - &bootstrap.activation_id, - WorkerActivationFailureReason::WorkerExitedBeforeReady, - &error, - ) - .await?; - directive = refresh_registration(lease).await?.directive; - continue; - } - if launched.as_ref().is_some_and(|(id, _, _)| id == &bootstrap.activation_id) - && super::common::control::now_unix_ms() >= bootstrap.deadline_unix_ms + if let Some(next_directive) = + supervise_worker_launch(lease, &bootstrap, &mut launched).await? { - let error = CliError::Launch( - "activated worker exceeded its startup deadline".into(), - ); - cleanup_pending_launch(lease, &mut launched).await?; - report_activation_failed( - lease, - &bootstrap.activation_id, - WorkerActivationFailureReason::WorkerReadinessTimeout, - &error, - ) - .await?; - directive = refresh_registration(lease).await?.directive; + directive = next_directive; continue; } } @@ -514,6 +460,71 @@ async fn supervise_route( result.and(cleanup) } +async fn supervise_worker_launch( + lease: &mut McpSession, + bootstrap: &WorkerBootstrap, + launched: &mut PendingLaunch, +) -> Result, CliError> { + let already_launched = launched + .as_ref() + .is_some_and(|(id, _, _)| id == &bootstrap.activation_id); + if !already_launched { + cleanup_pending_launch(lease, launched).await?; + match launch_worker(&lease.daemon_origin, bootstrap).await { + Ok(child) => { + *launched = Some(( + bootstrap.activation_id.clone(), + child, + tokio::time::Instant::now(), + )); + } + Err(error) => { + report_activation_failed( + lease, + &bootstrap.activation_id, + error.failure_reason, + &error.source, + ) + .await?; + return Ok(Some(refresh_registration(lease).await?.directive)); + } + } + } + if let Some((activation_id, child, _)) = launched.as_mut() + && activation_id == &bootstrap.activation_id + && let Some(status) = child.child.try_wait().map_err(CliError::Io)? + { + let error = CliError::Launch(format!( + "activated worker exited before readiness with {status}" + )); + report_activation_failed( + lease, + &bootstrap.activation_id, + WorkerActivationFailureReason::WorkerExitedBeforeReady, + &error, + ) + .await?; + return Ok(Some(refresh_registration(lease).await?.directive)); + } + if launched + .as_ref() + .is_some_and(|(id, _, _)| id == &bootstrap.activation_id) + && super::common::control::now_unix_ms() >= bootstrap.deadline_unix_ms + { + let error = CliError::Launch("activated worker exceeded its startup deadline".into()); + cleanup_pending_launch(lease, launched).await?; + report_activation_failed( + lease, + &bootstrap.activation_id, + WorkerActivationFailureReason::WorkerReadinessTimeout, + &error, + ) + .await?; + return Ok(Some(refresh_registration(lease).await?.directive)); + } + Ok(None) +} + async fn next_directive(lease: &mut McpSession) -> Result { let event = lease.client.next().await; receive_directive(lease, event).await diff --git a/crates/core/tests/fixtures/native_plugin/src/lib.rs b/crates/core/tests/fixtures/native_plugin/src/lib.rs index 92be56a2e..62e44a32a 100644 --- a/crates/core/tests/fixtures/native_plugin/src/lib.rs +++ b/crates/core/tests/fixtures/native_plugin/src/lib.rs @@ -220,118 +220,7 @@ impl NativePlugin for FixtureNativePlugin { } } })?; - ctx.register_tool_execution_intercept("fixture_tool_execution", 0, { - let runtime = runtime.clone(); - move |context, next| { - let runtime = runtime.clone(); - async move { - let args = context.args; - let args = mark_json(args, "native_plugin_tool_execution_request"); - let result = if args - .get("use_scoped_next") - .and_then(Json::as_bool) - .unwrap_or(false) - { - let mut scope = runtime.scope( - "fixture.native.scoped.next", - ScopeType::Custom, - None, - None, - Some(&Json::String("scoped-next-input".into())), - )?; - let call_result = next.call(args).await; - let close_result = - scope.close(Some(&Json::String("scoped-next-output".into())), None); - close_result?; - call_result? - } else if args - .get("use_isolated_next") - .and_then(Json::as_bool) - .unwrap_or(false) - { - let isolated = runtime.create_scope_stack()?; - let previous = runtime.capture_scope_stack_thread()?; - if isolated.set_thread() != NemoRelayStatus::Ok { - return Err("failed to install isolated scope stack".into()); - } - let mut scope = runtime.scope( - "fixture.native.isolated.next", - ScopeType::Custom, - None, - None, - Some(&Json::String("isolated-next-input".into())), - )?; - let call_result = next.call(args).await; - let close_result = - scope.close(Some(&Json::String("isolated-next-output".into())), None); - if previous.restore() != NemoRelayStatus::Ok { - return Err("failed to restore callback scope stack".into()); - } - close_result?; - call_result? - } else if args - .get("use_concurrent_next") - .and_then(Json::as_bool) - .unwrap_or(false) - { - let first_next = next.clone(); - let (first, second) = - tokio::join!(first_next.call(args.clone()), next.call(args),); - let result = first?; - second?; - result - } else { - next.call(args).await? - }; - let mut result = result; - result.result = mark_json(result.result, "native_plugin_tool_execution"); - Ok( - ToolExecutionInterceptOutcome::from(result).with_pending_mark( - PendingMarkSpec::builder() - .name("fixture.native.tool_execution.mark") - .category(EventCategory::custom()) - .category_profile(CategoryProfile { - subtype: Some("fixture.native.tool_execution".into()), - ..CategoryProfile::default() - }) - .data(json!({ "source": "native_tool_execution" })) - .metadata(json!({ "fixture": true })) - .build(), - ), - ) - } - } - })?; - ctx.register_tool_execution_intercept("fixture_tool_execution_nested", 1, { - let runtime = runtime.clone(); - move |context, next| { - let runtime = runtime.clone(); - async move { - let args = context.args; - if !args - .get("use_scoped_next") - .and_then(Json::as_bool) - .unwrap_or(false) - { - return next.call(args).await.map(Into::into); - } - let mut scope = runtime.scope( - "fixture.native.scoped.next.downstream", - ScopeType::Custom, - None, - None, - Some(&Json::String("scoped-next-downstream-input".into())), - )?; - let call_result = next.call(args).await; - let close_result = scope.close( - Some(&Json::String("scoped-next-downstream-output".into())), - None, - ); - close_result?; - call_result.map(Into::into) - } - } - })?; + register_fixture_tool_execution(ctx, &runtime)?; ctx.register_llm_sanitize_request_guardrail( "fixture_llm_sanitize_request", @@ -469,6 +358,125 @@ impl NativePlugin for FixtureNativePlugin { } } +fn register_fixture_tool_execution( + ctx: &mut PluginContext<'_>, + runtime: &PluginRuntime, +) -> nemo_relay_plugin::Result<()> { + ctx.register_tool_execution_intercept("fixture_tool_execution", 0, { + let runtime = runtime.clone(); + move |context, next| { + let runtime = runtime.clone(); + async move { + let args = context.args; + let args = mark_json(args, "native_plugin_tool_execution_request"); + let result = if args + .get("use_scoped_next") + .and_then(Json::as_bool) + .unwrap_or(false) + { + let mut scope = runtime.scope( + "fixture.native.scoped.next", + ScopeType::Custom, + None, + None, + Some(&Json::String("scoped-next-input".into())), + )?; + let call_result = next.call(args).await; + let close_result = + scope.close(Some(&Json::String("scoped-next-output".into())), None); + close_result?; + call_result? + } else if args + .get("use_isolated_next") + .and_then(Json::as_bool) + .unwrap_or(false) + { + let isolated = runtime.create_scope_stack()?; + let previous = runtime.capture_scope_stack_thread()?; + if isolated.set_thread() != NemoRelayStatus::Ok { + return Err("failed to install isolated scope stack".into()); + } + let mut scope = runtime.scope( + "fixture.native.isolated.next", + ScopeType::Custom, + None, + None, + Some(&Json::String("isolated-next-input".into())), + )?; + let call_result = next.call(args).await; + let close_result = + scope.close(Some(&Json::String("isolated-next-output".into())), None); + if previous.restore() != NemoRelayStatus::Ok { + return Err("failed to restore callback scope stack".into()); + } + close_result?; + call_result? + } else if args + .get("use_concurrent_next") + .and_then(Json::as_bool) + .unwrap_or(false) + { + let first_next = next.clone(); + let (first, second) = + tokio::join!(first_next.call(args.clone()), next.call(args),); + let result = first?; + second?; + result + } else { + next.call(args).await? + }; + let mut result = result; + result.result = mark_json(result.result, "native_plugin_tool_execution"); + Ok( + ToolExecutionInterceptOutcome::from(result).with_pending_mark( + PendingMarkSpec::builder() + .name("fixture.native.tool_execution.mark") + .category(EventCategory::custom()) + .category_profile(CategoryProfile { + subtype: Some("fixture.native.tool_execution".into()), + ..CategoryProfile::default() + }) + .data(json!({ "source": "native_tool_execution" })) + .metadata(json!({ "fixture": true })) + .build(), + ), + ) + } + } + })?; + ctx.register_tool_execution_intercept("fixture_tool_execution_nested", 1, { + let runtime = runtime.clone(); + move |context, next| { + let runtime = runtime.clone(); + async move { + let args = context.args; + if !args + .get("use_scoped_next") + .and_then(Json::as_bool) + .unwrap_or(false) + { + return next.call(args).await.map(Into::into); + } + let mut scope = runtime.scope( + "fixture.native.scoped.next.downstream", + ScopeType::Custom, + None, + None, + Some(&Json::String("scoped-next-downstream-input".into())), + )?; + let call_result = next.call(args).await; + let close_result = scope.close( + Some(&Json::String("scoped-next-downstream-output".into())), + None, + ); + close_result?; + call_result.map(Into::into) + } + } + })?; + Ok(()) +} + fn register_fixture_event_metadata_injector( ctx: &mut PluginContext<'_>, return_error: bool, diff --git a/go/nemo_relay/llm_test.go b/go/nemo_relay/llm_test.go index 16398ef67..acc37cc63 100644 --- a/go/nemo_relay/llm_test.go +++ b/go/nemo_relay/llm_test.go @@ -774,6 +774,24 @@ func TestLlmExecutionInterceptRegisterDeregister(t *testing.T) { } } +func resolveDirectionalExecutionTestCodecs(context LLMExecutionContext) (*LLMRequestSanitizeCodec, *LLMResponseSanitizeCodec, error) { + if context.RequestCodec.Codec.CodecKind != LLMCodecOpaque { + return nil, nil, fmt.Errorf("unexpected request codec identity: %#v", context.RequestCodec.Codec) + } + if context.ResponseCodec == nil || + context.ResponseCodec.Codec.CodecKind != LLMCodecBuiltin || + context.ResponseCodec.Codec.CodecID == nil || + *context.ResponseCodec.Codec.CodecID != "openai_chat" { + return nil, nil, fmt.Errorf("unexpected response codec context: %#v", context.ResponseCodec) + } + requestCodec := context.RequestCodec.ResolveCodec() + responseCodec := context.ResponseCodec.ResolveCodec() + if requestCodec == nil || responseCodec == nil { + return nil, nil, errors.New("execution codec capability did not resolve") + } + return requestCodec, responseCodec, nil +} + func TestLlmExecutionInterceptResolvesDirectionalCodecs(t *testing.T) { const interceptName = "go_llm_execution_codec_context" _ = DeregisterLlmExecutionIntercept(interceptName) @@ -788,19 +806,9 @@ func TestLlmExecutionInterceptResolvesDirectionalCodecs(t *testing.T) { interceptName, 1, func(nativeJSON json.RawMessage, context LLMExecutionContext, next func(json.RawMessage) (json.RawMessage, error)) (json.RawMessage, error) { - if context.RequestCodec.Codec.CodecKind != LLMCodecOpaque { - return nil, fmt.Errorf("unexpected request codec identity: %#v", context.RequestCodec.Codec) - } - if context.ResponseCodec == nil || - context.ResponseCodec.Codec.CodecKind != LLMCodecBuiltin || - context.ResponseCodec.Codec.CodecID == nil || - *context.ResponseCodec.Codec.CodecID != "openai_chat" { - return nil, fmt.Errorf("unexpected response codec context: %#v", context.ResponseCodec) - } - requestCodec := context.RequestCodec.ResolveCodec() - responseCodec := context.ResponseCodec.ResolveCodec() - if requestCodec == nil || responseCodec == nil { - return nil, errors.New("execution codec capability did not resolve") + requestCodec, responseCodec, err := resolveDirectionalExecutionTestCodecs(context) + if err != nil { + return nil, err } var request LLMRequestDTO if err := json.Unmarshal(nativeJSON, &request); err != nil {