diff --git a/.cargo/config.toml b/.cargo/config.toml index b43d287b..e63b40a9 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -26,3 +26,37 @@ rustflags = [ [alias] pctx = "run -p pctx --" +# Hot reload: rebuild and restart `pctx start` on source changes +dev = [ + "watch", + "-w", + "crates", + "-w", + "Cargo.toml", + "-w", + "Cargo.lock", + "-x", + "run -p pctx -- start", +] +dev-v = [ + "watch", + "-w", + "crates", + "-w", + "Cargo.toml", + "-w", + "Cargo.lock", + "-x", + "run -p pctx -- start -v", +] +dev-vv = [ + "watch", + "-w", + "crates", + "-w", + "Cargo.toml", + "-w", + "Cargo.lock", + "-x", + "run -p pctx -- start -vv", +] diff --git a/CHANGELOG.md b/CHANGELOG.md index 8c795ddd..34d96a76 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,11 +19,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- `-v`/`-vv` flags apply globally (previously was only applied to `pctx mcp ` commands) +- `ExecuteBashOutput`'s `Display` printed `stdout` in the `# STDERR` section, so bash stderr was never shown. ([#141](https://github.com/portofcontext/pctx/issues/141)) - `/register/tools` no longer fails the whole batch when a single tool cannot be registered. Each tool is registered independently: a genuinely bad tool (name clash, unparseable schema) is skipped and returned in the response's `failed` list, and a tool our codegen cannot type degrades to an `any` signature rather than being dropped. Previously one bad tool from an upstream server — such as a recursive `$ref` in a federated schema — took down registration for the entire batch, forcing clients into ~135 sequential per-tool calls per session. -- Declared `rmcp` minimum raised from 1.2.0 to 1.8.0, the version the code - actually requires. With the understated minimum, downstream consumers of the - git-dep crates could resolve an older rmcp and fail to compile - `pctx_session_server`. +- Declared `rmcp` minimum raised from 1.2.0 to 1.8.0, the version the code actually requires. With the understated minimum, downstream consumers of the git-dep crates could resolve an older rmcp and fail to compile `pctx_session_server`. ## [v0.7.3] - 2026-07-22 diff --git a/crates/pctx/src/lib.rs b/crates/pctx/src/lib.rs index f8ad2e75..b3f32ac0 100644 --- a/crates/pctx/src/lib.rs +++ b/crates/pctx/src/lib.rs @@ -6,7 +6,10 @@ use clap::{Parser, Subcommand}; use serde_json::json; use std::io::{self, Write}; -use crate::utils::{logger::init_cli_logger, telemetry::init_telemetry}; +use crate::utils::{ + logger::{self, init_cli_logger}, + telemetry::init_telemetry, +}; use pctx_config::Config; #[derive(Parser)] @@ -42,52 +45,58 @@ pub struct Cli { } impl Cli { - fn cli_logger(&self) -> bool { - !matches!( - &self.command, - Commands::Mcp(McpCommands::Start(_) | McpCommands::Dev(_)) - ) - } - - fn json_l(&self) -> Option { - if let Commands::Mcp(McpCommands::Dev(dev)) = &self.command { - Some(dev.log_file.clone()) - } else { - None - } - } + /// `-v` and `-q` are global, so every command resolves its logging here. + async fn init_logging(&self, cfg: &Config) -> anyhow::Result<()> { + let level = logger::flag_level(self.verbose, self.quiet); - #[allow(clippy::missing_errors_doc)] - pub async fn handle(&self) -> anyhow::Result<()> { match &self.command { - Commands::Mcp(mcp_cmd) => self.handle_mcp(mcp_cmd).await, - Commands::Start(start_cmd) => { - let cfg = Config::load(&self.config).unwrap_or_default(); - // Session server uses stdout for logs (not stdio protocol) - init_telemetry(&cfg, None, false).await?; - - start_cmd.handle().await + // Short-lived commands print for a human, the rest emit structured logs + Commands::Mcp( + McpCommands::Init(_) + | McpCommands::List(_) + | McpCommands::Add(_) + | McpCommands::Remove(_), + ) => { + init_cli_logger(self.verbose, self.quiet); + Ok(()) } + // Dev writes JSONL for its TUI to tail + Commands::Mcp(McpCommands::Dev(dev)) => { + init_telemetry(cfg, Some(dev.log_file.clone()), false, level).await + } + // Stdio mode keeps stdout clean for JSON-RPC + Commands::Mcp(McpCommands::Start(start_cmd)) => { + init_telemetry(cfg, None, start_cmd.stdio, level).await + } + Commands::Start(_) => init_telemetry(cfg, None, false, level).await, } } - async fn handle_mcp(&self, cmd: &McpCommands) -> anyhow::Result<()> { + #[allow(clippy::missing_errors_doc)] + pub async fn handle(&self) -> anyhow::Result<()> { let cfg = Config::load(&self.config); - if let (McpCommands::Start(start_cmd), Err(err)) = (cmd, &cfg) + if let (Commands::Mcp(McpCommands::Start(start_cmd)), Err(err)) = (&self.command, &cfg) && start_cmd.stdio { return Self::handle_stdio_config_error(err); } - if self.cli_logger() { - init_cli_logger(self.verbose, self.quiet); - } else if let Ok(c) = &cfg { - // Use stderr for stdio mode to keep stdout clean for JSON-RPC - let use_stderr = matches!(cmd, McpCommands::Start(start_cmd) if start_cmd.stdio); - init_telemetry(c, self.json_l(), use_stderr).await?; + // A broken config still gets logging, so the error is reported the usual way + let fallback = Config::default(); + self.init_logging(cfg.as_ref().unwrap_or(&fallback)).await?; + + match &self.command { + Commands::Start(start_cmd) => start_cmd.handle().await, + Commands::Mcp(mcp_cmd) => self.handle_mcp(mcp_cmd, cfg).await, } + } + async fn handle_mcp( + &self, + cmd: &McpCommands, + cfg: anyhow::Result, + ) -> anyhow::Result<()> { let _updated_cfg = match cmd { McpCommands::Init(cmd) => cmd.handle(&self.config).await?, McpCommands::List(cmd) => cmd.handle(cfg?).await?, diff --git a/crates/pctx/src/utils/logger.rs b/crates/pctx/src/utils/logger.rs index bda31176..ef39b9d4 100644 --- a/crates/pctx/src/utils/logger.rs +++ b/crates/pctx/src/utils/logger.rs @@ -2,11 +2,16 @@ use std::io::Write; const WHITELISTED_CRATES: &[&str] = &[ "pctx", - "pctx_mcp_server", - "pctx_session_server", + "pctx_code_execution_runtime", + "pctx_code_mode", + "pctx_codegen", "pctx_config", + "pctx_deno_transpiler", "pctx_executor", - "pctx_codegen", + "pctx_mcp_server", + "pctx_registry", + "pctx_session_server", + "pctx_type_check_runtime", ]; pub(crate) fn default_env_filter(level: &str) -> String { @@ -21,16 +26,18 @@ pub(crate) fn default_env_filter(level: &str) -> String { filters.join(",") } +/// Level named by the global `-v`/`-q` flags, or `None` when neither was passed. +pub(crate) fn flag_level(verbose: u8, quiet: bool) -> Option<&'static str> { + match (quiet, verbose) { + (true, _) => Some("warn"), + (false, 0) => None, + (false, 1) => Some("debug"), + (false, _) => Some("trace"), + } +} + pub(crate) fn init_cli_logger(verbose: u8, quiet: bool) { - let level_str = if quiet { - "warn" - } else if verbose == 0 { - "info" - } else if verbose == 1 { - "debug" - } else { - "trace" - }; + let level_str = flag_level(verbose, quiet).unwrap_or("info"); let mut builder = env_logger::Builder::from_env( env_logger::Env::default().default_filter_or(default_env_filter(level_str)), diff --git a/crates/pctx/src/utils/telemetry.rs b/crates/pctx/src/utils/telemetry.rs index 019dbde5..1b6757fc 100644 --- a/crates/pctx/src/utils/telemetry.rs +++ b/crates/pctx/src/utils/telemetry.rs @@ -16,7 +16,15 @@ pub(crate) async fn init_telemetry( cfg: &Config, json_l: Option, use_stderr: bool, + flag_level: Option<&str>, ) -> Result<()> { + // An explicit -v/-q outranks RUST_LOG, which outranks the config file. + let env_filter = |default: &str| match flag_level { + Some(level) => EnvFilter::new(logger::default_env_filter(level)), + None => EnvFilter::try_from_default_env() + .unwrap_or_else(|_| EnvFilter::new(logger::default_env_filter(default))), + }; + // Set global text map propagator for trace context propagation (W3C Trace Context) // This enables parsing of traceparent/tracestate headers in distributed tracing opentelemetry::global::set_text_map_propagator(TraceContextPropagator::new()); @@ -70,17 +78,13 @@ pub(crate) async fn init_telemetry( let write_to = fs::File::create(&log_file).context(format!("failed creating log file: {log_file}"))?; - let env_filter = EnvFilter::try_from_default_env() - .unwrap_or(EnvFilter::new(logger::default_env_filter("debug"))); layers.push( init_tracing_layer(write_to, &LoggerFormat::Json, false) - .with_filter(env_filter) + .with_filter(env_filter("debug")) .boxed(), ); } else if cfg.logger.enabled { - let env_filter = EnvFilter::try_from_default_env().unwrap_or(EnvFilter::new( - logger::default_env_filter(cfg.logger.level.as_str()), - )); + let env_filter = env_filter(cfg.logger.level.as_str()); // Determine log destination based on config and mode: // 1. If file is specified in config, use it (all modes) @@ -209,7 +213,7 @@ mod tests { ..Default::default() }); - let result = init_telemetry(&cfg, None, false).await; + let result = init_telemetry(&cfg, None, false, None).await; assert!(result.is_ok(), "Telemetry initialization should succeed"); assert!( log_path.exists(), diff --git a/crates/pctx_code_execution_runtime/src/mcp_ops.rs b/crates/pctx_code_execution_runtime/src/mcp_ops.rs deleted file mode 100644 index 342465f9..00000000 --- a/crates/pctx_code_execution_runtime/src/mcp_ops.rs +++ /dev/null @@ -1,27 +0,0 @@ -//! Deno ops for MCP client functionality -//! -//! These ops expose the Rust MCP client to JavaScript - -use deno_core::OpState; -use deno_core::op2; -use rmcp::model::JsonObject; -use std::cell::RefCell; -use std::rc::Rc; - -use pctx_registry::{MCPRegistry, RegistryError}; - -/// Call an MCP tool (async op) -#[op2(async)] -#[serde] -pub(crate) async fn op_call_mcp_tool( - state: Rc>, - #[string] server_name: String, - #[string] tool_name: String, - #[serde] args: Option, -) -> Result { - let registry = { - let borrowed = state.borrow(); - borrowed.borrow::().clone() - }; - pctx_registry::call_mcp_tool(®istry, &server_name, &tool_name, args).await -} diff --git a/crates/pctx_code_mode/src/code_mode.rs b/crates/pctx_code_mode/src/code_mode.rs index fe7e8c27..c042d8ed 100644 --- a/crates/pctx_code_mode/src/code_mode.rs +++ b/crates/pctx_code_mode/src/code_mode.rs @@ -7,7 +7,7 @@ use std::{ collections::{HashMap, HashSet}, time::Duration, }; -use tracing::{debug, info, instrument, warn}; +use tracing::{debug, info, instrument, trace, warn}; use crate::{ Error, Result, @@ -602,12 +602,16 @@ export default result;"#, self.default_registry()? }; - // Format for logging only - let formatted_code = pctx_codegen::format::format_ts(code); + let timer = std::time::Instant::now(); + info!( + code_length = code.len(), + disclosure = ?disclosure, + actions = registry.ids().len(), + "Executing TypeScript" + ); debug!( code_from_llm = %code, - formatted_code = %formatted_code, code_length = code.len(), callbacks =? registry.ids(), disclosure =? disclosure, @@ -701,7 +705,7 @@ export default result;"#, } }; - debug!("Executing TypeScript in sandbox:\n{to_execute}"); + trace!("Executing TypeScript in sandbox:\n{to_execute}"); let execution_res = pctx_executor::execute( &to_execute, @@ -710,9 +714,18 @@ export default result;"#, .await?; if execution_res.success { - debug!("TypeScript execution completed successfully"); + info!( + duration_ms = timer.elapsed().as_millis(), + trace_events = execution_res.trace.events.len(), + "TypeScript execution succeeded" + ); } else { - warn!("TypeScript execution failed: {:?}", execution_res.stderr); + warn!( + duration_ms = timer.elapsed().as_millis(), + trace_events = execution_res.trace.events.len(), + stderr = %execution_res.stderr, + "TypeScript execution failed" + ); } let output = ExecuteTypescriptOutput { diff --git a/crates/pctx_code_mode/src/model.rs b/crates/pctx_code_mode/src/model.rs index e2141dda..6dc3652a 100644 --- a/crates/pctx_code_mode/src/model.rs +++ b/crates/pctx_code_mode/src/model.rs @@ -155,7 +155,7 @@ impl Display for ExecuteBashOutput { write!( f, "Exit Code: {}\n\n# STDOUT\n{}\n\n# STDERR\n{}", - &self.exit_code, &self.stdout, &self.stdout + &self.exit_code, &self.stdout, &self.stderr ) } } diff --git a/crates/pctx_executor/src/lib.rs b/crates/pctx_executor/src/lib.rs index 58c3b330..018ded0b 100644 --- a/crates/pctx_executor/src/lib.rs +++ b/crates/pctx_executor/src/lib.rs @@ -224,7 +224,7 @@ pub async fn execute(code: &str, options: ExecuteOptions) -> Result Result { let mut check_result = type_check(code)?; @@ -312,7 +312,10 @@ struct InternalExecuteResult { /// /// # Errors /// * Returns error only if internal Deno runtime initialization fails -#[tracing::instrument(fields(runtime = "execution"))] +#[tracing::instrument( + skip(code, options), + fields(runtime = "execution", code_length = code.len()) +)] async fn execute_code( code: &str, options: ExecuteOptions, diff --git a/crates/pctx_registry/src/registry.rs b/crates/pctx_registry/src/registry.rs index ec529075..ecda5c47 100644 --- a/crates/pctx_registry/src/registry.rs +++ b/crates/pctx_registry/src/registry.rs @@ -13,9 +13,9 @@ use std::{ future::Future, pin::Pin, sync::{Arc, RwLock}, - time::SystemTime, + time::{Instant, SystemTime}, }; -use tracing::{debug, instrument, warn}; +use tracing::{debug, info, instrument, warn}; pub type CallbackFn = Arc< dyn Fn( @@ -251,6 +251,10 @@ impl PctxRegistry { RegistryAction::Callback(callback_fn) => { let args_json = args.as_ref().map(|a| json!(a)); let started_at = SystemTime::now(); + let timer = Instant::now(); + + info!("Invoking callback"); + debug!(args = ?args_json, "Callback arguments"); let result = callback_fn(args.map(|a| json!(a))).await.map_err(|e| { RegistryError::ExecutionError(format!( @@ -258,6 +262,18 @@ impl PctxRegistry { )) }); + match &result { + Ok(_) => info!( + duration_ms = timer.elapsed().as_millis(), + "Callback succeeded" + ), + Err(e) => warn!( + duration_ms = timer.elapsed().as_millis(), + error = %e, + "Callback failed" + ), + } + self.trace .push(RegistryEvent::CallbackInvocation(CallbackInvocationEvent { id: id.to_string(), @@ -293,6 +309,10 @@ impl PctxRegistry { let args_json = args.as_ref().map(|a| json!(a)); let started_at = SystemTime::now(); + let timer = Instant::now(); + + info!("Calling MCP tool"); + debug!(args = ?args_json, "MCP tool arguments"); let (client, cached_client) = self .connection_pool @@ -364,6 +384,20 @@ impl PctxRegistry { Ok(val) })(); + match &result { + Ok(_) => info!( + cached_client, + duration_ms = timer.elapsed().as_millis(), + "MCP tool call succeeded" + ), + Err(e) => warn!( + cached_client, + duration_ms = timer.elapsed().as_millis(), + error = %e, + "MCP tool call failed" + ), + } + self.trace .push(RegistryEvent::McpToolCall(McpToolCallEvent { server: mcp_id.sever_name.clone(), diff --git a/crates/pctx_session_server/src/state/ws_manager.rs b/crates/pctx_session_server/src/state/ws_manager.rs index f8b8cda9..87aa4359 100644 --- a/crates/pctx_session_server/src/state/ws_manager.rs +++ b/crates/pctx_session_server/src/state/ws_manager.rs @@ -2,7 +2,7 @@ use std::{collections::HashMap, sync::Arc, time::Duration}; use rmcp::model::RequestId; use tokio::sync::{RwLock, mpsc as tokio_mpsc}; -use tracing::{debug, info, warn}; +use tracing::{debug, warn}; use uuid::Uuid; use crate::model::{ExecuteToolParams, ExecuteToolResult, PctxJsonRpcRequest, WsJsonRpcMessage}; @@ -159,6 +159,13 @@ impl WsSession { .await .insert(req_id.clone(), response_tx); + debug!( + request_id = ?req_id, + tool = %tool_id, + ws_session_id = %self.id, + "Dispatching callback request to client", + ); + // Send message to client self.sender .send(WsJsonRpcMessage::request( @@ -196,9 +203,11 @@ impl WsSession { result: Result, ) -> Result<(), ()> { let pending_read = self.pending_executions.read().await; - info!( + debug!( + request_id = ?request_id, + ws_session_id = %self.id, pending_count = pending_read.len(), - "Handling execution response for request_id: {request_id:?}", + "Handling callback response from client", ); if let Some(response_tx) = pending_read.get(&request_id) { debug!("Found pending execution, sending result"); diff --git a/crates/pctx_session_server/src/websocket/handler.rs b/crates/pctx_session_server/src/websocket/handler.rs index 6532cacc..25d901e1 100644 --- a/crates/pctx_session_server/src/websocket/handler.rs +++ b/crates/pctx_session_server/src/websocket/handler.rs @@ -29,7 +29,7 @@ use rmcp::{ }; use serde_json::json; use tokio::sync::mpsc; -use tracing::{debug, error, info, warn}; +use tracing::{debug, error, info, trace, warn}; use uuid::Uuid; use crate::AppState; @@ -84,8 +84,6 @@ async fn handle_socket( state: AppState, code_mode_session: Uuid, ) { - info!(session_id =? code_mode_session, "New WebSocket connection"); - // Split socket into sender and receiver let (sender, receiver) = socket.split(); @@ -96,11 +94,10 @@ async fn handle_socket( let session = WsSession::new(tx.clone(), code_mode_session); let ws_session = session.id; - debug!( - session_id =? code_mode_session, - ws_session =? ws_session, - "Created session {ws_session} connected to code mode session {}", - session.code_mode_session_id + info!( + session_id = %code_mode_session, + ws_session_id = %ws_session, + "New WebSocket connection" ); state.ws_manager.add(session).await; @@ -125,7 +122,11 @@ async fn handle_socket( state.ws_manager.remove_session(ws_session).await; - info!("WebSocket connection closed for session {ws_session}"); + info!( + session_id = %code_mode_session, + ws_session_id = %ws_session, + "WebSocket connection closed" + ); } /// Handle outgoing WebSocket messages (`execute_tool` requests from server) @@ -370,7 +371,7 @@ async fn handle_message( ) -> Result<(), String> { match msg { Message::Text(text) => { - debug!("Received text message from {ws_session}: {text}"); + trace!("Received text message from {ws_session}: {text}"); let jrpc_msg = serde_json::from_str::(&text) .map_err(|e| format!("Received invalid JsonRpc message from websocket: {e}"))?; @@ -420,7 +421,9 @@ async fn handle_message( Ok(()) } Message::Close(_) => { - info!("Received close message for session {ws_session}"); + // The "WebSocket connection closed" line follows immediately, so + // this one is only interesting when debugging the handshake. + debug!(ws_session_id = %ws_session, "Received close message"); Ok(()) } Message::Ping(_) | Message::Pong(_) => Ok(()), diff --git a/examples/telemetry/docker-compose.yml b/examples/telemetry/docker-compose.yml index 3560db56..f610998a 100644 --- a/examples/telemetry/docker-compose.yml +++ b/examples/telemetry/docker-compose.yml @@ -9,7 +9,7 @@ volumes: services: # OTLP Collector - receives OTLP and exposes metrics for Prometheus otel-collector: - image: otel/opentelemetry-collector-contrib:latest + image: otel/opentelemetry-collector-contrib:0.157.0 command: ["--config=/etc/otel-collector-config.yml"] ports: - "4317:4317" # OTLP gRPC receiver @@ -22,7 +22,7 @@ services: # Traces tempo: - image: grafana/tempo:latest + image: grafana/tempo:2.10.7 command: -config.file=/etc/tempo.yaml ports: - "3200:3200" @@ -47,7 +47,7 @@ services: # Metrics prometheus: - image: prom/prometheus:latest + image: prom/prometheus:v3.13.2 command: - "--config.file=/etc/prometheus/prometheus.yml" - "--storage.tsdb.path=/prometheus" diff --git a/pctx-py/tests/scripts/manual_code_mode.py b/pctx-py/tests/scripts/manual_code_mode.py index 2eed0028..9682ae5d 100755 --- a/pctx-py/tests/scripts/manual_code_mode.py +++ b/pctx-py/tests/scripts/manual_code_mode.py @@ -56,7 +56,7 @@ async def main(): "url": "https://mcp.stripe.com", "auth": { "type": "bearer", - "token": getenv("STRIPE_MCP_KEY"), + "token": getenv("STRIPE_MCP_KEY", ""), }, } ], @@ -73,7 +73,7 @@ async def main(): let subval = await MyMath.subtract({a: addval, b: 2}); let multval = await MyMath.multiply({a: subval, b: 2}); let now = await Tools.nowTimestamp(); - let customers = await Stripe.listCustomers({}); + let info = await Stripe.getStripeAccountInfo(); let logs = await Tools.searchLogs(); let logs2 = await Tools.searchLogs({}); let logs3 = await Tools.searchLogs({ query: "custom query" });