Skip to content
Open
Changes from 1 commit
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
201 changes: 195 additions & 6 deletions src/anthropic/stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,10 @@
//! 实现 Kiro → Anthropic 流式响应转换和 SSE 状态管理

use std::collections::{HashMap, VecDeque};
use std::sync::OnceLock;

use serde_json::json;
use sha2::{Digest, Sha256};
use uuid::Uuid;

use crate::kiro::model::events::Event;
Expand Down Expand Up @@ -1318,15 +1320,94 @@ pub(crate) fn canonicalize_structured_json(
schema: &serde_json::Value,
) -> Result<String, (&'static str, String)> {
let trimmed = text.trim();
let candidate = if let Some(inner) = trimmed
.strip_prefix("```json\n")
.and_then(|value| value.strip_suffix("\n```"))
{
inner.trim()
let normalized;
let (candidate, fence_class) = if trimmed.starts_with("```") {
normalized = trimmed.replace("\r\n", "\n");
let Some((opening, remainder)) = normalized.split_once('\n') else {
log_structured_json_failure(
"kiro_structured_json_parse_failed",
text,
trimmed,
trimmed,
"malformed",
None,
);
return Err((
"kiro_structured_json_parse_failed",
"Kiro structured output was not valid JSON: malformed fenced block".to_string(),
));
};
let opening = opening.trim();
let class = if opening == "```" {
"unlabelled"
} else if opening.eq_ignore_ascii_case("```json") {
"json"
} else {
"unsupported"
};
let (inner, closing) = match remainder.rsplit_once('\n') {
Some(parts) => parts,
None if remainder.trim() == "```" => ("", remainder),
None => {
log_structured_json_failure(
"kiro_structured_json_parse_failed",
text,
trimmed,
trimmed,
class,
None,
);
return Err((
"kiro_structured_json_parse_failed",
"Kiro structured output was not valid JSON: malformed fenced block"
.to_string(),
));
}
};
if class == "unsupported" || closing.trim() != "```" {
log_structured_json_failure(
"kiro_structured_json_parse_failed",
text,
trimmed,
trimmed,
class,
None,
);
return Err((
"kiro_structured_json_parse_failed",
"Kiro structured output was not valid JSON: unsupported or malformed fenced block"
.to_string(),
));
}
(inner.trim(), class)
} else {
trimmed
(trimmed, "none")
};

if candidate.is_empty() {
log_structured_json_failure(
"kiro_structured_empty_output",
text,
trimmed,
candidate,
fence_class,
None,
);
return Err((
"kiro_structured_empty_output",
"Kiro structured output was empty".to_string(),
));
}

let value: serde_json::Value = serde_json::from_str(candidate).map_err(|error| {
log_structured_json_failure(
"kiro_structured_json_parse_failed",
text,
trimmed,
candidate,
fence_class,
Some((error.line(), error.column())),
);
(
"kiro_structured_json_parse_failed",
format!("Kiro structured output was not valid JSON: {error}"),
Expand All @@ -1335,6 +1416,14 @@ pub(crate) fn canonicalize_structured_json(
let mut errors = Vec::new();
validate_json_schema(&value, schema, "$", &mut errors);
if !errors.is_empty() {
log_structured_json_failure(
"kiro_structured_json_schema_failed",
text,
trimmed,
candidate,
fence_class,
None,
);
return Err((
"kiro_structured_json_schema_failed",
format!(
Expand All @@ -1344,13 +1433,75 @@ pub(crate) fn canonicalize_structured_json(
));
}
serde_json::to_string(&value).map_err(|error| {
log_structured_json_failure(
"kiro_structured_json_parse_failed",
text,
trimmed,
candidate,
fence_class,
None,
);
(
"kiro_structured_json_parse_failed",
format!("Kiro structured output could not be serialized: {error}"),
)
})
}

fn log_structured_json_failure(
error_kind: &'static str,
text: &str,
trimmed: &str,
candidate: &str,
fence_class: &'static str,
parse_position: Option<(usize, usize)>,
) {
static PROCESS_SALT: OnceLock<[u8; 16]> = OnceLock::new();
let salt = PROCESS_SALT.get_or_init(|| *Uuid::new_v4().as_bytes());
let mut hasher = Sha256::new();
hasher.update(salt);
hasher.update(candidate.as_bytes());
let candidate_process_digest = hex::encode(hasher.finalize());
let first_non_whitespace_class = candidate
.chars()
.find(|value| !value.is_whitespace())
.map(classify_structured_json_lead)
.unwrap_or("none");
let (parse_line, parse_column) = parse_position.unwrap_or((0, 0));

tracing::warn!(
error_kind,
input_bytes = text.len(),
input_chars = text.chars().count(),
trimmed_bytes = trimmed.len(),
candidate_bytes = candidate.len(),
candidate_chars = candidate.chars().count(),
empty = candidate.is_empty(),
fence_class,
first_non_whitespace_class,
candidate_process_digest,
parse_line,
parse_column,
"Kiro structured output validation failed; response content was not logged"
);
}

fn classify_structured_json_lead(value: char) -> &'static str {
match value {
'{' => "object",
'[' => "array",
'"' => "string",
'-' => "minus",
'0'..='9' => "digit",
't' | 'f' | 'n' => "literal",
'`' => "fence",
value if value.is_ascii_alphabetic() => "ascii_alpha",
value if value.is_control() => "control",
value if !value.is_ascii() => "non_ascii",
_ => "other",
}
}

/// SSE 事件
#[derive(Debug, Clone)]
pub struct SseEvent {
Expand Down Expand Up @@ -3674,6 +3825,17 @@ mod tests {
canonicalize_structured_json("```json\n{ \"ok\": true }\n```", &schema).unwrap(),
"{\"ok\":true}"
);
for fenced in [
"```\n{ \"ok\": true }\n```",
"```JSON\n{ \"ok\": true }\n```",
"```Json\r\n{ \"ok\": true }\r\n```",
" ```json\n{ \"ok\": true }\n``` \n",
] {
assert_eq!(
canonicalize_structured_json(fenced, &schema).unwrap(),
"{\"ok\":true}"
);
}
assert_eq!(
canonicalize_structured_json("{\"ok\":\"yes\"}", &schema)
.unwrap_err()
Expand All @@ -3686,6 +3848,33 @@ mod tests {
.0,
"kiro_structured_json_parse_failed"
);
for empty in [
"",
" \r\n\t",
"```json\n```",
"```\r\n```",
"```json\n\n```",
"```\r\n \r\n```",
] {
assert_eq!(
canonicalize_structured_json(empty, &schema).unwrap_err().0,
"kiro_structured_empty_output"
);
}
for invalid in [
"Here is the JSON:\n```json\n{\"ok\":true}\n```",
"```json\n{\"ok\":true}\n```\ndone",
"```yaml\n{\"ok\":true}\n```",
"```json\n```json\n{\"ok\":true}\n```\n```",
"```json {\"ok\":true}```",
] {
assert_eq!(
canonicalize_structured_json(invalid, &schema)
.unwrap_err()
.0,
"kiro_structured_json_parse_failed"
);
}
}

/// 防回归:统一管道的两个去向(流式 emit_completed_tool_use 与非流式 to_anthropic_block)
Expand Down
Loading