Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
9 changes: 9 additions & 0 deletions docs/user/getting-started-and-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,7 @@ provider = "openai-subscription" # or "openrouter" or "speakeasy"
model = "gpt-5.4"
reasoning_effort = "medium" # low, medium, or high
a2a = "127.0.0.1:7331"
capture_error_spans = false # optional local context in fatal error logs
otel_endpoint = "http://localhost:4317"
otel_protocol = "grpc" # grpc, http/protobuf, or http/json
otel_capture_message_content = false
Expand Down Expand Up @@ -210,6 +211,14 @@ For settings exposed by a command, precedence is:
2. values in `~/.kit/config.toml`;
3. built-in defaults.

Set `capture_error_spans = true` when troubleshooting unexpected failures or
preparing a bug report. Kit adds diagnostic context about recent operations to
error logs in `~/.kit/errors/<session-id>/`, which can help explain a failure.

It is disabled by default to avoid additional collection overhead and does not
require OpenTelemetry. The extra context excludes prompts and tool inputs and
outputs; it is not a complete execution history.

The OpenTelemetry endpoint follows the same CLI-over-TOML precedence, then falls
back to the standard `OTEL_EXPORTER_OTLP_ENDPOINT` environment variable. If none
is set, trace export is disabled. The trace protocol precedence is
Expand Down
102 changes: 101 additions & 1 deletion src/fatal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,22 @@ struct FatalRecord {
message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
diagnostics: Option<TransportDiagnostics>,
#[serde(
default,
skip_serializing_if = "Option::is_none",
deserialize_with = "deserialize_span_context"
)]
span_context: Option<crate::telemetry::error_spans::Snapshot>,
}

fn deserialize_span_context<'de, D: serde::Deserializer<'de>>(
deserializer: D,
) -> Result<Option<crate::telemetry::error_spans::Snapshot>, D::Error> {
let snapshot = Option::<crate::telemetry::error_spans::Snapshot>::deserialize(deserializer)?;
if snapshot.as_ref().is_some_and(|snapshot| !snapshot.valid()) {
return Err(serde::de::Error::custom("invalid span context"));
}
Ok(snapshot)
}

#[derive(Clone, Copy, Debug, Deserialize, Serialize)]
Expand Down Expand Up @@ -447,7 +463,7 @@ fn write_in_with_diagnostics(
std::process::id(),
NEXT_EVENT.fetch_add(1, Ordering::Relaxed)
);
let record = FatalRecord {
let mut record = FatalRecord {
schema_version: SCHEMA_VERSION,
event_id: event_id.clone(),
occurred_at_ms,
Expand All @@ -458,9 +474,15 @@ fn write_in_with_diagnostics(
code: canonical_code(code).into(),
message: bounded(message),
diagnostics: diagnostics.filter(|value| value.valid()).cloned(),
span_context: crate::telemetry::error_spans::snapshot(&tracing::Span::current()),
};
let mut bytes = serde_json::to_vec_pretty(&record)
.map_err(|error| format!("could not encode fatal error log: {error}"))?;
if bytes.len() >= MAX_RECORD_BYTES && record.span_context.take().is_some() {
// Optional diagnostics must not displace an otherwise valid ordinary error.
bytes = serde_json::to_vec_pretty(&record)
.map_err(|error| format!("could not encode fatal error log: {error}"))?;
}
bytes.push(b'\n');
if bytes.len() > MAX_RECORD_BYTES {
return Err("fatal error log exceeds size limit".into());
Expand Down Expand Up @@ -620,6 +642,84 @@ mod tests {
assert_eq!(record.code, "stream_transport");
}

#[test]
fn schema_two_readers_preserve_optional_span_context() {
use tracing_subscriber::prelude::*;
// Frozen shipped schema-v2 shape: unknown top-level fields are ignored.
#[derive(serde::Deserialize, serde::Serialize)]
struct ShippedV2 {
schema_version: u64,
event_id: String,
occurred_at_ms: u64,
kit_version: String,
session_id: String,
surface: String,
kind: String,
code: String,
message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
diagnostics: Option<super::TransportDiagnostics>,
}
let root = tempfile::tempdir().unwrap();
tracing::subscriber::with_default(
tracing_subscriber::registry().with(crate::telemetry::error_spans::ErrorSpanLayer),
|| {
let operation = crate::telemetry::error_spans::operation("prompt");
operation.in_scope(|| {
{ let _child = tracing::info_span!(target: "agentkit_loop", "agent.execute_tool", launch_kind = "plain"); }
let path = write_in_with_diagnostics(root.path(), "session-context", Surface::Prompt, "provider", "stream_transport", "openai-subscription stream transport failed", Some(&sample_diagnostics())).unwrap();
let bytes = fs::read(&path).unwrap();
let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert_eq!(value["schema_version"], 2);
assert_eq!(value["span_context"]["fragments"][1]["fields"]["launch_kind"], "plain");
let current: FatalRecord = serde_json::from_slice(&bytes).unwrap();
let legacy: ShippedV2 = serde_json::from_slice(&bytes).unwrap();
let mut known = serde_json::to_value(&current).unwrap();
known.as_object_mut().unwrap().remove("span_context");
assert_eq!(serde_json::to_value(&legacy).unwrap(), known);
assert_eq!(current.message, "openai-subscription stream transport failed");
let supplied = value["span_context"].clone();
for marker in [1, 2, 3, 4] {
value["schema_version"] = json!(marker);
let parsed: FatalRecord = serde_json::from_value(value.clone()).unwrap();
assert_eq!(serde_json::to_value(parsed.span_context).unwrap(), supplied);
}
value.as_object_mut().unwrap().remove("span_context");
assert!(serde_json::from_value::<FatalRecord>(value.clone()).unwrap().span_context.is_none());
value["span_context"] = supplied;
value["span_context"]["fragments"][1]["fields"]["launch_kind"] = json!("SECRET");
assert!(serde_json::from_value::<FatalRecord>(value).is_err());
assert_eq!(fs::read(path).unwrap(), bytes);
assert!(bytes.len() < super::MAX_RECORD_BYTES);
});
},
);
}

#[test]
fn ordinary_writer_omits_disabled_context() {
tracing::subscriber::with_default(tracing_subscriber::registry(), || {
let root = tempfile::tempdir().unwrap();
let operation = crate::telemetry::error_spans::operation("prompt");
let path = operation
.in_scope(|| {
write_in(
root.path(),
"session-disabled",
Surface::Prompt,
"runtime",
"runtime_error",
"ordinary error",
)
})
.unwrap();
let value: serde_json::Value =
serde_json::from_slice(&fs::read(path).unwrap()).unwrap();
assert!(value.get("span_context").is_none());
assert_eq!(value["message"], "ordinary error");
});
}

#[test]
fn schema_v1_records_remain_readable() {
let record: FatalRecord = serde_json::from_value(json!({
Expand Down
66 changes: 66 additions & 0 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,9 @@ fn resolve_openrouter_api_key(

#[derive(Args)]
struct TelemetryArgs {
/// Resolved local diagnostic setting inherited by built-in Kit children.
#[arg(long, hide = true, global = true, value_name = "BOOL", action = clap::ArgAction::Set)]
internal_capture_error_spans: Option<bool>,
/// OTLP collector endpoint for OpenTelemetry trace export.
#[arg(long, global = true)]
otel_endpoint: Option<String>,
Expand Down Expand Up @@ -232,6 +235,7 @@ struct Config {
provider: Option<kit::ProviderKind>,
reasoning_effort: Option<kit::ReasoningEffort>,
a2a: Option<String>,
capture_error_spans: Option<bool>,
otel_endpoint: Option<String>,
otel_protocol: Option<kit::telemetry::Protocol>,
otel_capture_message_content: Option<bool>,
Expand Down Expand Up @@ -409,6 +413,13 @@ impl Config {
max_messages,
max_bytes,
)
.map(|mut settings| {
settings.capture_error_spans = args
.internal_capture_error_spans
.or(self.capture_error_spans)
.unwrap_or(false);
settings
})
.map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error))
}

Expand Down Expand Up @@ -1679,6 +1690,61 @@ credential_store = "keychain"
assert!(toml::from_str::<Config>("otel_protocol = 'http'").is_err());
}

#[test]
fn inherited_error_capture_overrides_config_without_rewriting_it() {
let root = tempfile::tempdir().unwrap();
let path = root.path().join("config.toml");
for inherited in [false, true] {
let text = format!(
"# User-owned settings\ncapture_error_spans = {}\n",
!inherited
);
fs::write(&path, &text).unwrap();
let config = Config::load(&path).unwrap();
let cli = Cli::try_parse_from([
"kit",
"prompt",
"--internal-capture-error-spans",
&inherited.to_string(),
"hello",
])
.unwrap();
let settings = config
.telemetry_settings(&cli.telemetry, None, None, None, None)
.unwrap();
assert_eq!(settings.capture_error_spans, inherited);
assert_eq!(fs::read_to_string(&path).unwrap(), text);
}
}

#[test]
fn error_span_capture_defaults_off_and_is_independent_of_export() {
let cli = Cli::try_parse_from(["kit", "prompt", "hello"]).unwrap();
for (text, expected) in [
("", false),
("capture_error_spans = false", false),
("capture_error_spans = true", true),
] {
let config: Config = toml::from_str(text).unwrap();
for endpoint in [None, Some("http://localhost:4317".to_owned())] {
let settings = config
.telemetry_settings(
&cli.telemetry,
endpoint.clone(),
Some("true".into()),
None,
None,
)
.unwrap();
assert_eq!(settings.capture_error_spans, expected);
assert_eq!(settings.endpoint, endpoint);
assert!(settings.capture_message_content);
}
}
assert!(toml::from_str::<Config>("capture_error_spans = 'true'").is_err());
assert!(toml::from_str::<Config>("capture_error_spans = 1").is_err());
}

#[test]
fn telemetry_environment_is_strict_and_settings_are_bounded() {
let config = Config::default();
Expand Down
120 changes: 62 additions & 58 deletions src/protocols/a2a.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use a2a_protocol_types::{
};

use sha2::{Digest as _, Sha256};
use tracing::Instrument as _;

use crate::runtime::Runtime;

Expand All @@ -24,67 +25,70 @@ impl AgentExecutor for KitAgent {
context: &'a RequestContext,
queue: &'a dyn EventQueueWriter,
) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>> {
Box::pin(async move {
let emit = EventEmitter::new(context, queue);
emit.status(TaskState::Working).await?;
let prompt = context
.message
.parts
.iter()
.filter_map(Part::text_content)
.collect::<Vec<_>>()
.join("\n");
if prompt.trim().is_empty() {
emit.artifact(
"error",
vec![Part::text("A2A request must contain a text part")],
None,
Some(true),
)
.await?;
emit.status(TaskState::Failed).await?;
return Ok(());
}
match self
.0
.run_cancelled(prompt, 0, Some(context.cancellation_token.clone()))
.await
{
Ok(output) => {
emit.artifact("result", vec![Part::text(output)], None, Some(true))
.await?;
emit.status(TaskState::Completed).await?;
}
Err(error) => {
let session_id = a2a_session_id(context);
let rendered = crate::fatal::render_loop_error(&error);
let rendered = match crate::fatal::record_loop_error(
&session_id,
crate::fatal::Surface::A2a,
&error,
) {
Ok(Some(path)) => {
eprintln!(
"stored fatal error log for {session_id}: {}",
path.display()
);
rendered
}
Ok(None) => rendered,
Err(log_error) => {
eprintln!(
"could not store fatal error log for {session_id}: {log_error}"
);
rendered
}
};
emit.artifact("error", vec![Part::text(rendered)], None, Some(true))
.await?;
Box::pin(
async move {
let emit = EventEmitter::new(context, queue);
emit.status(TaskState::Working).await?;
let prompt = context
.message
.parts
.iter()
.filter_map(Part::text_content)
.collect::<Vec<_>>()
.join("\n");
if prompt.trim().is_empty() {
emit.artifact(
"error",
vec![Part::text("A2A request must contain a text part")],
None,
Some(true),
)
.await?;
emit.status(TaskState::Failed).await?;
return Ok(());
}
match self
.0
.run_cancelled(prompt, 0, Some(context.cancellation_token.clone()))
.await
{
Ok(output) => {
emit.artifact("result", vec![Part::text(output)], None, Some(true))
.await?;
emit.status(TaskState::Completed).await?;
}
Err(error) => {
let session_id = a2a_session_id(context);
let rendered = crate::fatal::render_loop_error(&error);
let rendered = match crate::fatal::record_loop_error(
&session_id,
crate::fatal::Surface::A2a,
&error,
) {
Ok(Some(path)) => {
eprintln!(
"stored fatal error log for {session_id}: {}",
path.display()
);
rendered
}
Ok(None) => rendered,
Err(log_error) => {
eprintln!(
"could not store fatal error log for {session_id}: {log_error}"
);
rendered
}
};
emit.artifact("error", vec![Part::text(rendered)], None, Some(true))
.await?;
emit.status(TaskState::Failed).await?;
}
}
Ok(())
}
Ok(())
})
.instrument(crate::telemetry::error_spans::operation("a2a")),
)
}
}

Expand Down
Loading
Loading