magi-code 0.77.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::*;
use crate::providers::error::provider_stream_trace_from_error;

fn push(parser: &mut AnthropicStreamParser, event: Value) -> StreamParseOutcome {
    parser
        .push_chunk_outcome(&format!("data: {event}\n\n"))
        .unwrap()
}

#[test]
fn fragmented_frames_emit_text_usage_and_terminal_done_in_order() {
    let mut parser = AnthropicStreamParser::default();
    let chunks = [
        "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":2}}}\n",
        "\ndata: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n\nda",
        "ta: {\"type\":\"message_delta\",\"usage\":{\"output_tokens\":3}}\n\ndata: {\"type\":\"message_stop\"}\n\n",
    ];
    let mut events = Vec::new();
    for chunk in chunks {
        events.extend(parser.push_chunk_outcome(chunk).unwrap().events);
    }
    assert_eq!(
        events,
        vec![
            ProviderEvent::UsageObserved(crate::providers::stream::UsageObservation {
                usage: Usage {
                    input: 2,
                    total: 2,
                    ..Usage::default()
                },
                presence: crate::providers::stream::UsagePresence {
                    input: true,
                    ..Default::default()
                },
            }),
            ProviderEvent::TextDelta("hi".to_string()),
            ProviderEvent::UsageObserved(crate::providers::stream::UsageObservation {
                usage: Usage {
                    input: 2,
                    output: 3,
                    total: 5,
                    ..Usage::default()
                },
                presence: crate::providers::stream::UsagePresence {
                    output: true,
                    ..Default::default()
                },
            }),
            ProviderEvent::Done,
        ]
    );
    assert!(parser.finish().unwrap().is_empty());
}

#[test]
fn assembles_tool_arguments_and_marks_unsafe_progress() {
    let mut parser = AnthropicStreamParser::default();
    let start = push(
        &mut parser,
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
    );
    let first = push(
        &mut parser,
        json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"path\":"}}),
    );
    let second = push(
        &mut parser,
        json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"\"a.txt\"}"}}),
    );
    let stop = push(&mut parser, json!({"type":"content_block_stop","index":1}));
    assert!(start.semantic_progress && start.unsafe_recovery_progress);
    assert!(first.semantic_progress && first.unsafe_recovery_progress);
    assert!(second.semantic_progress && second.unsafe_recovery_progress);
    assert_eq!(
        stop.events,
        vec![
            ProviderEvent::ResponseItem(
                json!({"type":"function_call","call_id":"call_1","name":"read","arguments":"{\"path\":\"a.txt\"}","status":"completed"})
            ),
            ProviderEvent::ToolCall(ToolCall {
                id: "call_1".into(),
                name: "read".into(),
                arguments: json!({"path":"a.txt"}),
            }),
        ]
    );
}

#[test]
fn emits_thinking_and_redacted_thinking_blocks() {
    let mut parser = AnthropicStreamParser::default();
    assert!(
        push(
            &mut parser,
            json!({"type":"content_block_start","index":0,"content_block":{"type":"thinking"}}),
        )
        .semantic_progress
    );
    push(
        &mut parser,
        json!({"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"private"}}),
    );
    push(
        &mut parser,
        json!({"type":"content_block_delta","index":0,"delta":{"type":"signature_delta","signature":"sig_1"}}),
    );
    assert_eq!(
        push(&mut parser, json!({"type":"content_block_stop","index":0}),).events,
        vec![ProviderEvent::ResponseItem(
            json!({"type":"thinking","thinking":"private","signature":"sig_1"})
        )]
    );

    push(
        &mut parser,
        json!({"type":"content_block_start","index":2,"content_block":{"type":"redacted_thinking","data":"encrypted_1"}}),
    );
    assert_eq!(
        push(&mut parser, json!({"type":"content_block_stop","index":2}),).events,
        vec![ProviderEvent::ResponseItem(
            json!({"type":"redacted_thinking","data":"encrypted_1"})
        )]
    );
}

#[test]
fn ping_has_no_semantic_progress_but_text_does() {
    let mut parser = AnthropicStreamParser::default();
    let ping = push(&mut parser, json!({"type":"ping"}));
    assert!(ping.events.is_empty());
    assert!(!ping.semantic_progress);
    let text = push(
        &mut parser,
        json!({"type":"content_block_delta","delta":{"type":"text_delta","text":"hi"}}),
    );
    assert!(text.semantic_progress);
    assert!(!text.unsafe_recovery_progress);
}

#[test]
fn usage_total_saturates() {
    let mut parser = AnthropicStreamParser::default();
    let events = push(
        &mut parser,
        json!({"type":"message_start","message":{"usage":{"input_tokens":u64::MAX,"output_tokens":1,"cache_read_input_tokens":1,"cache_creation_input_tokens":1}}}),
    )
    .events;
    assert!(matches!(
        events.as_slice(),
        [ProviderEvent::UsageObserved(
            crate::providers::stream::UsageObservation {
                usage: Usage {
                    input: u64::MAX,
                    output: 1,
                    total: u64::MAX,
                    cache_read: 1,
                    cache_write: 1,
                    ..
                },
                ..
            }
        )]
    ));
}

#[test]
fn completion_and_incomplete_eof_are_validated() {
    let mut complete = AnthropicStreamParser::default();
    assert_eq!(
        push(&mut complete, json!({"type":"message_stop"})).events,
        vec![ProviderEvent::Done]
    );
    assert!(complete.finish().unwrap().is_empty());

    let mut incomplete = AnthropicStreamParser::default();
    push(
        &mut incomplete,
        json!({"type":"content_block_delta","delta":{"type":"text_delta","text":"hi"}}),
    );
    assert!(
        incomplete
            .finish()
            .unwrap_err()
            .to_string()
            .contains("missing provider stream completion")
    );
}

#[test]
fn bounded_diagnostics_redact_and_escape_sensitive_fragments() {
    let secret = "sk-ant-api03-secret";
    let long = format!("{}tail-marker", "x".repeat(240));
    let mut malformed = AnthropicStreamParser::default();
    let error = malformed
        .push_chunk_outcome(&format!(
            "data: {{\"type\":\"message_start\",\"token\":\"{secret}\",\"text\":\"line\\n{long}\"\n\n"
        ))
        .unwrap_err()
        .to_string();
    assert!(
        error.contains("malformed provider SSE data JSON"),
        "{error}"
    );
    assert!(!error.contains(secret), "{error}");
    assert!(error.contains("<redacted>"), "{error}");
    assert!(error.contains("\\n"), "{error}");
    assert!(error.contains("..."), "{error}");
    assert!(!error.contains("tail-marker"), "{error}");

    let mut incomplete = AnthropicStreamParser::default();
    incomplete
        .push_chunk_outcome(&format!("data: token={secret}\n{long}"))
        .unwrap();
    let error = incomplete.finish().unwrap_err().to_string();
    assert!(!error.contains(secret), "{error}");
    assert!(error.contains("<redacted>"), "{error}");
    assert!(error.contains("\\n"), "{error}");
    assert!(error.contains("..."), "{error}");
    assert!(!error.contains("tail-marker"), "{error}");
}

#[test]
fn provider_error_and_bad_tool_arguments_use_bounded_redacted_snippets() {
    let secret = "sk-ant-api03-secret";
    let long = format!("{}tail-marker", "x".repeat(240));
    let mut parser = AnthropicStreamParser::default();
    let error = push_error(
        &mut parser,
        json!({"type":"error","error":{"message":format!("{secret}\n{long}")}}),
    );
    assert!(error.contains("anthropic provider stream error"), "{error}");
    assert!(!error.contains(secret), "{error}");
    assert!(error.contains("<redacted>"), "{error}");
    assert!(error.contains("\\n"), "{error}");
    assert!(error.contains("..."), "{error}");
    assert!(!error.contains("tail-marker"), "{error}");

    let mut parser = AnthropicStreamParser::default();
    push(
        &mut parser,
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
    );
    push(
        &mut parser,
        json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":format!("{{\"token\":\"{secret}\",\"text\":\"{long}")}}),
    );
    let error = push_error(&mut parser, json!({"type":"content_block_stop","index":1}));
    assert!(error.contains("malformed non-empty provider tool call arguments"));
    assert!(!error.contains(secret), "{error}");
    assert!(error.contains("<redacted>"), "{error}");
    assert!(error.contains("..."), "{error}");
    assert!(!error.contains("tail-marker"), "{error}");
}

fn push_error(parser: &mut AnthropicStreamParser, event: Value) -> String {
    parser
        .push_chunk_outcome(&format!("data: {event}\n\n"))
        .unwrap_err()
        .to_string()
}

#[test]
fn rejects_oversized_event_and_tool_arguments() {
    let mut event_parser = AnthropicStreamParser::default();
    let error = event_parser
        .push_chunk_outcome(&"x".repeat(MAX_SSE_EVENT_BUFFER_BYTES + 1))
        .unwrap_err()
        .to_string();
    assert!(error.contains("maximum buffered size"), "{error}");

    let mut tool_parser = AnthropicStreamParser::default();
    push(
        &mut tool_parser,
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
    );
    let error = push_error(
        &mut tool_parser,
        json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"x".repeat(MAX_TOOL_ARGUMENT_BYTES + 1)}}),
    );
    assert!(error.contains("maximum size"), "{error}");
}

#[test]
fn validates_tool_indexes_metadata_and_active_state() {
    let invalid = [
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","name":"read"}}),
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1"}}),
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"","name":"read"}}),
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"   "}}),
    ];
    for event in invalid {
        let mut parser = AnthropicStreamParser::default();
        assert!(push_error(&mut parser, event).contains("missing non-empty"));
        assert!(parser.pending_tools.is_empty());
    }

    let mut missing_index = AnthropicStreamParser::default();
    assert!(
        push_error(
            &mut missing_index,
            json!({"type":"content_block_start","content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
        )
        .contains("missing required index")
    );

    let mut duplicate = AnthropicStreamParser::default();
    push(
        &mut duplicate,
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
    );
    assert!(
        push_error(
            &mut duplicate,
            json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_2","name":"write"}}),
        )
        .contains("duplicate active tool content block index 1")
    );
}

#[test]
fn pending_tool_error_has_bounded_hashed_trace_without_raw_arguments() {
    let mut parser = AnthropicStreamParser::default();
    push(
        &mut parser,
        json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"toolu_trace_full_id","name":"read"}}),
    );
    push(
        &mut parser,
        json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"path\":\"/tmp"}}),
    );
    push(
        &mut parser,
        json!({"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":7}}),
    );
    let error = parser
        .push_chunk_outcome("data: {\"type\":\"message_stop\"}\n\n")
        .unwrap_err();
    let trace = provider_stream_trace_from_error(&error).unwrap();

    assert_eq!(trace.provider, "anthropic");
    assert_eq!(trace.failure_context, "message_stop");
    assert_eq!(trace.message_delta_stop_reason.as_deref(), Some("tool_use"));
    let pending = trace.pending_tools.first().unwrap();
    assert_eq!(pending.argument_bytes, 13);
    assert_eq!(
        pending.argument_sha256,
        "ac2a055274ea75011d1fc6958da49b826e7a101c25258bb537ba6c36fed6a56c"
    );
    let delta = trace
        .recent_events
        .iter()
        .find(|event| event.delta_type.as_deref() == Some("input_json_delta"))
        .unwrap();
    assert_eq!(delta.partial_json_bytes, Some(13));
    assert_eq!(
        delta.partial_json_sha256.as_deref(),
        Some("ac2a055274ea75011d1fc6958da49b826e7a101c25258bb537ba6c36fed6a56c")
    );
    let serialized = serde_json::to_string(&trace).unwrap();
    assert!(!serialized.contains("/tmp"));
    assert!(!serialized.contains(r#"{\"path\":\"/tmp"#));
}

#[test]
fn stream_trace_bounds_events_and_pending_tools() {
    let mut parser = AnthropicStreamParser::default();
    for index in 0..20 {
        push(
            &mut parser,
            json!({"type":"content_block_start","index":index,"content_block":{"type":"tool_use","id":format!("toolu_{index}"),"name":"read"}}),
        );
        push(
            &mut parser,
            json!({"type":"content_block_delta","index":index,"delta":{"type":"input_json_delta","partial_json":format!("raw-/tmp-marker-{index}")}}),
        );
    }
    let error = parser.finish().unwrap_err();
    let trace = provider_stream_trace_from_error(&error).unwrap();
    let serialized = serde_json::to_string(&trace).unwrap();
    assert!(trace.recent_events.len() <= ANTHROPIC_STREAM_TRACE_EVENT_LIMIT);
    assert_eq!(trace.pending_tool_count, 20);
    assert_eq!(
        trace.pending_tools.len(),
        ANTHROPIC_STREAM_TRACE_PENDING_TOOL_LIMIT
    );
    assert!(trace.pending_tools_truncated);
    assert!(!serialized.contains("raw-/tmp-marker"));
    assert!(!serialized.contains("/tmp"));
    assert!(serialized.contains("partial_json_sha256"));
}

#[test]
fn sse_data_strips_only_one_optional_space() {
    assert_eq!(
        sse_data("data:  hello  \ndata:\tthere\t").as_deref(),
        Some(" hello  \n\tthere\t")
    );
}