mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
* feat: Add secure prompt-based skills system (Phase 1 MVP) Implement a skills system that extends the agent with prompt-level instructions from local directories. Skills declare activation criteria, tool permissions, and trust tiers that determine authority attenuation. Core security model: the minimum trust level of any active skill determines a tool ceiling -- tools above the ceiling are removed from the LLM's tool list entirely at the API level, preventing prompt-based manipulation. New modules: - skills/mod.rs: Core types (SkillTrust, SkillManifest, LoadedSkill) - skills/scanner.rs: Content scanner for manipulation detection - skills/registry.rs: Filesystem discovery and manifest parsing - skills/selector.rs: Deterministic two-phase prefilter (no LLM) - skills/attenuation.rs: Trust-based tool filtering Integration: - Agent loop selects skills per-turn and applies tool attenuation - Reasoning engine injects skill context with structural isolation - Config supports SKILLS_ENABLED, SKILLS_DIR, SKILLS_MAX_ACTIVE, SKILLS_MAX_CONTEXT_TOKENS environment variables - Disabled by default (SKILLS_ENABLED=false) 41 new tests covering all modules. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: Address all adversarial review findings for skills system Security fixes: - Escape skill name/version in XML attributes to prevent trust spoofing - Escape prompt content to prevent </skill> tag breakout - Require integrity hash for Verified/Community tier skills - Validate skill names against [a-zA-Z0-9][a-zA-Z0-9._-]{0,63} - Add 64 KiB file size limit on prompt.md Bug fixes: - Use actual SkillsConfig from AgentDeps instead of SkillsConfig::default() - Add skills_config field to AgentDeps, wired through from main.rs Performance: - Pre-compile regex patterns at load time (cached on LoadedSkill) - Selector uses pre-compiled patterns instead of recompiling per message - Switch all std::fs to tokio::fs for non-blocking async I/O Hardening: - Cap keyword score at 30 points to prevent keyword stuffing attacks - Enforce max 20 keywords and 5 patterns per skill - Normalize line endings (CRLF/CR to LF) before hashing - Also includes cargo fmt formatting fixes for adjacent code Tests: 54 skills tests pass (up from 41), zero new clippy warnings. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: Address medium/low severity findings from adversarial review Fixes all 18 medium/low severity findings identified by the security review: - mod.rs: Add MAX_TAGS_PER_SKILL cap (10) in enforce_limits(); use RegexBuilder with 64 KiB size_limit to prevent ReDoS; replace case-enumerated escape_skill_content with regex matching all case variants plus whitespace/null byte injection between </ and skill; document allowed_patterns as unenforced until Phase 2; document Marketplace URL validation as Phase 3 concern - registry.rs: Add MAX_MANIFEST_FILE_SIZE (16 KiB) check before reading; add symlink detection via symlink_metadata to reject symlinks in discover_local; add MAX_DISCOVERED_SKILLS (100) cap; validate prompt_hash format (sha256: + 64 hex chars); warn on name collision before overwriting; accept SkillSource parameter in load_skill instead of always using Local; add InvalidHashFormat, ManifestTooLarge, SymlinkDetected error variants - selector.rs: Add MAX_TAG_SCORE (15) cap parallel to keyword cap; warn when declared max_context_tokens diverges >2x from actual prompt size - scanner.rs: Add mixed-script homoglyph detection (Cyrillic, Greek, Armenian unicode ranges); document token-boundary bypass and semantic paraphrasing as known limitations - attenuation.rs: Document READ_ONLY_TOOLS maintenance requirements - agent_loop.rs: Surface scan warnings via structured tracing; add structured audit events for skill activation and tool attenuation 61 tests pass, 0 new clippy warnings. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: Address 12 findings from second adversarial security review HIGH: - Escape opening <skill tags in prompt content (prevents fake skill block injection) - Scan manifest metadata fields (description, author, tags, reasons) not just prompt - Block trust downgrade on name collision (existing Local can't be replaced by Community) MEDIUM: - Eliminate TOCTOU gap: read files then check size instead of metadata-then-read - Reject file-level symlinks in load_skill (prompt.md, skill.toml) - Truncate and filter manifest.skill.tags (prevent unlimited tag scoring) - Cap regex pattern score at 40 (prevent 5x20=100 dominating keyword+tag) - Add doc comment about skill_list tool exposing metadata (sanitization required) - Move Community disclaimer inside <skill> tags (not outside structural boundary) - Filter keywords/tags shorter than 3 chars (prevent broad matching) LOW: - Enforce minimum token_cost of 1 (max_context_tokens=0 can't bypass budget) - Remove redundant try_exists checks in discover_local (let load_skill handle errors) 70 skills tests passing. Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: Add HTTP endpoint scoping for skills (Phase 1) Skills that declare an [http] section in skill.toml now have their HTTP requests constrained to declared endpoints at runtime. This addresses the gap where allowed_patterns was parsed but never enforced -- once the http tool was visible via attenuation, the LLM could reach any URL. Enforcement reuses EndpointPattern/AllowlistValidator from the WASM capability system. Semantics: if no active skill declares [http], all requests pass through (backward compat). If any skill declares [http], URLs must match at least one skill's allowlist (union). Community skills' [http] declarations are silently ignored (defense in depth). Shell commands using curl/wget are also validated against scopes. Scanner gains detection for known exfiltration domains (webhook.site, ngrok.io, etc.), overly broad wildcards, and credential/host mismatches. Closes #38 Co-Authored-By: Claude Opus 4.6 <[email protected]> * style: Apply cargo fmt to http_scoping.rs Co-Authored-By: Claude Opus 4.6 <[email protected]> * style: Apply cargo fmt across codebase Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: Add parameter-level permission enforcement for skills (Phase 2) Activates enforcement of `allowed_patterns` in skill.toml permissions. Previously these patterns were parsed but not enforced -- a Verified skill declaring `permissions.shell` with `allowed_patterns = [{command = "cargo *"}]` could still run any shell command. Now the enforcer validates tool parameters against declared glob patterns before execution. Key changes: - New `enforcer.rs` module with `SkillPermissionEnforcer`, `glob_to_regex()`, and `validate_tool_call()` with union semantics across active skills - Typed pattern enums (`ShellPattern`, `FilePathPattern`, `MemoryTargetPattern`) replace the previous `Vec<serde_json::Value>` in `ToolPermissionDeclaration` - Scanner gains `scan_permission_patterns()` detecting dangerous patterns (rm, sudo, curl, bare wildcards, command chaining, sensitive paths, identity files) - Registry blocks non-Local skills with critical permission pattern warnings - Agent loop threads enforcer into `execute_chat_tool` alongside HTTP scoping Trust interaction: Community patterns ignored, Verified enforced, Local without patterns unrestricted, Local with patterns enforced as guidance. Union semantics across skills -- tool call allowed if ANY skill's patterns permit it. 34 new tests. All 818 library tests pass. Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: Add worker permission enforcement and LLM behavioral analysis (Phase 3+4) Phase 3 - Worker-side permission enforcement: - Add SerializedToolPermission/SerializedPattern DTOs for HTTP boundary crossing - Extend JobDescription, ContainerHandle, and orchestrator API to carry permissions - CreateJobTool snapshots and forwards skill permissions to spawned workers - Worker runtime builds SkillPermissionEnforcer and checks before tool execution - Load-time token budget enforcement rejects prompts exceeding 2x declared budget - Deduplicate enforcer construction: from_active_skills() delegates to from_serialized() Phase 4 - LLM behavioral analysis: - BehavioralAnalyzer with cached, LLM-based semantic content analysis - Structured output parsing (FINDING|CATEGORY|SEVERITY|DESCRIPTION or CLEAN) - Content-hash caching with bounded size (MAX_CACHE_ENTRIES=256) - Graceful degradation when LLM unavailable - Integrated into load_skill() for non-Local skills; critical findings block loading Review fixes: - Real cache tests with CountingLlm mock (test_cache_hit, test_cache_miss, test_cache_bounded) - UTF-8-safe truncate() in worker runtime - Few-shot examples in behavioral analysis prompt - Documented max_context_tokens=0 opt-out and create_job() permission gap 848 tests passing, no new clippy warnings. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: Address review feedback from serrrfirat on skills-phase2 - Fix truncate_cmd UTF-8 panic: use char-boundary-aware slicing - Remove redundant effective_tools branching in reasoning.rs - Document cache eviction as known limitation (arbitrary, not LRU) - Add safety comment on SkillTrust enum ordering (security-critical) - Simplify active_skills selection (prefilter_skills handles empty input) Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address remaining skills review feedback * refactor: replace skills system with OpenClaw SKILL.md format + 2-state trust Replace the 5-gate, 3-tier trust hierarchy (scanner, behavioral analyzer, parameter-level enforcer, HTTP endpoint scoping) with a simplified 3-layer security model: gating -> attenuation -> Docker confinement. Key changes: - SKILL.md format (YAML frontmatter + markdown prompt) replaces skill.toml + prompt.md - 2-state trust (Installed/Trusted) replaces 3-tier (Community/Verified/Local) - New parser.rs for SKILL.md parsing with serde_yaml - New gating.rs for requirements checking (bins/env/config) - Simplified registry with 2-location discovery (workspace + user dirs) - Removed scanner, behavioral_analyzer, enforcer, http_scoping (~4,100 lines) - Removed skill_permissions propagation through job/orchestrator/worker pipeline - Added serde_yaml dependency for YAML frontmatter parsing Net: -5,298 lines, 59 skills tests pass, 907 total tests pass. Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: add in-app skill management tools and ClawHub catalog integration Add 4 chat-callable tools (skill_list, skill_search, skill_install, skill_remove) plus matching web gateway endpoints for managing skills at runtime. The catalog fetches from ClawHub's public registry API at runtime rather than bundling entries at compile time. Key changes: - SkillRegistry gains mutation methods (install_skill, remove_skill, reload, find_by_name) with Arc<RwLock> for concurrent access - New catalog module queries ClawHub /api/v1/search with in-memory caching (5-min TTL, configurable via CLAWHUB_REGISTRY env var) - skill_list and skill_search added to READ_ONLY_TOOLS for safe use under Installed trust ceiling - Web gateway gets /api/skills, /api/skills/search, /api/skills/install, and /api/skills/{name} DELETE endpoints Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address PR #51 review feedback from ilblackdragon Security: - Add SSRF protection to fetch_skill_content: require HTTPS, reject private/loopback/link-local IPs and internal hostnames, disable redirects. Gateway install handler now reuses the same validation. - URL-encode slug in skill_download_url to prevent query injection. - Require X-Confirm-Action header on gateway skill install/remove endpoints (equivalent to chat tool requires_approval gate). Correctness: - Eliminate all block_in_place/block_on usage in skill tools and gateway handlers. Split install into prepare_install_to_disk (static async, no lock) + commit_install (sync, brief write lock). Same pattern for remove: validate_remove + delete_skill_files + commit_remove. - Write normalized content to disk in install_skill (was writing original un-normalized content, causing hash mismatch on re-read). - Fix token estimation from 0.75 to 0.25 tokens/byte (~4 chars per token) in registry.rs, selector.rs, and standalone loader. Dependencies: - Replace deprecated serde_yaml 0.9 with serde_yml 0.0.12. - Remove unused toml dependency. Co-Authored-By: Claude Opus 4.6 <[email protected]> --------- Co-authored-by: Claude Opus 4.6 <[email protected]>
335 lines
11 KiB
Rust
335 lines
11 KiB
Rust
//! End-to-end integration tests for the WebSocket gateway.
|
|
//!
|
|
//! These tests start a real Axum server on a random port, connect a WebSocket
|
|
//! client, and verify the full message flow:
|
|
//! - WebSocket upgrade with auth
|
|
//! - Ping/pong
|
|
//! - Client message → agent msg_tx
|
|
//! - Broadcast SSE event → WebSocket client
|
|
//! - Connection tracking (counter increment/decrement)
|
|
//! - Gateway status endpoint
|
|
|
|
use std::net::SocketAddr;
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use futures::{SinkExt, StreamExt};
|
|
use tokio::sync::mpsc;
|
|
use tokio::time::timeout;
|
|
use tokio_tungstenite::tungstenite::Message;
|
|
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
|
|
|
|
use ironclaw::channels::IncomingMessage;
|
|
use ironclaw::channels::web::server::{GatewayState, start_server};
|
|
use ironclaw::channels::web::sse::SseManager;
|
|
use ironclaw::channels::web::types::SseEvent;
|
|
use ironclaw::channels::web::ws::WsConnectionTracker;
|
|
|
|
const AUTH_TOKEN: &str = "test-token-12345";
|
|
const TIMEOUT: Duration = Duration::from_secs(5);
|
|
|
|
/// Start a gateway server on a random port and return the bound address + agent
|
|
/// message receiver.
|
|
async fn start_test_server() -> (
|
|
SocketAddr,
|
|
Arc<GatewayState>,
|
|
mpsc::Receiver<IncomingMessage>,
|
|
) {
|
|
let (agent_tx, agent_rx) = mpsc::channel(64);
|
|
|
|
let state = Arc::new(GatewayState {
|
|
msg_tx: tokio::sync::RwLock::new(Some(agent_tx)),
|
|
sse: SseManager::new(),
|
|
workspace: None,
|
|
session_manager: None,
|
|
log_broadcaster: None,
|
|
extension_manager: None,
|
|
tool_registry: None,
|
|
store: None,
|
|
job_manager: None,
|
|
prompt_queue: None,
|
|
user_id: "test-user".to_string(),
|
|
shutdown_tx: tokio::sync::RwLock::new(None),
|
|
ws_tracker: Some(Arc::new(WsConnectionTracker::new())),
|
|
llm_provider: None,
|
|
skill_registry: None,
|
|
skill_catalog: None,
|
|
chat_rate_limiter: ironclaw::channels::web::server::RateLimiter::new(30, 60),
|
|
});
|
|
|
|
let addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
|
|
let bound_addr = start_server(addr, state.clone(), AUTH_TOKEN.to_string())
|
|
.await
|
|
.expect("Failed to start test server");
|
|
|
|
(bound_addr, state, agent_rx)
|
|
}
|
|
|
|
/// Connect a WebSocket client with auth token in query parameter.
|
|
async fn connect_ws(
|
|
addr: SocketAddr,
|
|
) -> tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>> {
|
|
let url = format!("ws://{}/api/chat/ws?token={}", addr, AUTH_TOKEN);
|
|
let mut request = url.into_client_request().unwrap();
|
|
// Server requires an Origin header from localhost to prevent cross-site WS hijacking.
|
|
request.headers_mut().insert(
|
|
"Origin",
|
|
format!("http://127.0.0.1:{}", addr.port()).parse().unwrap(),
|
|
);
|
|
let (stream, _response) = tokio_tungstenite::connect_async(request)
|
|
.await
|
|
.expect("Failed to connect WebSocket");
|
|
stream
|
|
}
|
|
|
|
/// Read the next text frame from the WebSocket, with a timeout.
|
|
async fn recv_text(
|
|
stream: &mut (impl StreamExt<Item = Result<Message, tokio_tungstenite::tungstenite::Error>> + Unpin),
|
|
) -> String {
|
|
let msg = timeout(TIMEOUT, stream.next())
|
|
.await
|
|
.expect("Timed out waiting for WS message")
|
|
.expect("Stream ended")
|
|
.expect("WS error");
|
|
match msg {
|
|
Message::Text(text) => text.to_string(),
|
|
other => panic!("Expected Text frame, got {:?}", other),
|
|
}
|
|
}
|
|
|
|
// ============================================================================
|
|
// Tests
|
|
// ============================================================================
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_ping_pong() {
|
|
let (addr, _state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
|
|
// Send ping
|
|
let ping = r#"{"type":"ping"}"#;
|
|
ws.send(Message::Text(ping.into())).await.unwrap();
|
|
|
|
// Expect pong
|
|
let text = recv_text(&mut ws).await;
|
|
let parsed: serde_json::Value = serde_json::from_str(&text).unwrap();
|
|
assert_eq!(parsed["type"], "pong");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_message_reaches_agent() {
|
|
let (addr, _state, mut agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
|
|
// Send a chat message
|
|
let msg = r#"{"type":"message","content":"hello from ws","thread_id":"t42"}"#;
|
|
ws.send(Message::Text(msg.into())).await.unwrap();
|
|
|
|
// Verify it arrives on the agent's msg_tx
|
|
let incoming = timeout(TIMEOUT, agent_rx.recv())
|
|
.await
|
|
.expect("Timed out waiting for agent message")
|
|
.expect("Agent channel closed");
|
|
|
|
assert_eq!(incoming.content, "hello from ws");
|
|
assert_eq!(incoming.thread_id.as_deref(), Some("t42"));
|
|
assert_eq!(incoming.channel, "gateway");
|
|
assert_eq!(incoming.user_id, "test-user");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_broadcast_event_received() {
|
|
let (addr, state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
|
|
// Give the connection a moment to fully establish
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
|
|
// Broadcast an SSE event (simulates agent sending a response)
|
|
state.sse.broadcast(SseEvent::Response {
|
|
content: "agent says hi".to_string(),
|
|
thread_id: "t1".to_string(),
|
|
});
|
|
|
|
// The WS client should receive it
|
|
let text = recv_text(&mut ws).await;
|
|
let parsed: serde_json::Value = serde_json::from_str(&text).unwrap();
|
|
assert_eq!(parsed["type"], "event");
|
|
assert_eq!(parsed["event_type"], "response");
|
|
assert_eq!(parsed["data"]["content"], "agent says hi");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_thinking_event() {
|
|
let (addr, state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
|
|
state.sse.broadcast(SseEvent::Thinking {
|
|
message: "analyzing...".to_string(),
|
|
thread_id: None,
|
|
});
|
|
|
|
let text = recv_text(&mut ws).await;
|
|
let parsed: serde_json::Value = serde_json::from_str(&text).unwrap();
|
|
assert_eq!(parsed["type"], "event");
|
|
assert_eq!(parsed["event_type"], "thinking");
|
|
assert_eq!(parsed["data"]["message"], "analyzing...");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_connection_tracking() {
|
|
let (addr, state, _agent_rx) = start_test_server().await;
|
|
let tracker = state.ws_tracker.as_ref().unwrap();
|
|
|
|
assert_eq!(tracker.connection_count(), 0);
|
|
|
|
// Connect first client
|
|
let ws1 = connect_ws(addr).await;
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
assert_eq!(tracker.connection_count(), 1);
|
|
|
|
// Connect second client
|
|
let ws2 = connect_ws(addr).await;
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
assert_eq!(tracker.connection_count(), 2);
|
|
|
|
// Disconnect first
|
|
drop(ws1);
|
|
tokio::time::sleep(Duration::from_millis(100)).await;
|
|
assert_eq!(tracker.connection_count(), 1);
|
|
|
|
// Disconnect second
|
|
drop(ws2);
|
|
tokio::time::sleep(Duration::from_millis(100)).await;
|
|
assert_eq!(tracker.connection_count(), 0);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_invalid_message_returns_error() {
|
|
let (addr, _state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
|
|
// Send invalid JSON
|
|
ws.send(Message::Text("not json".into())).await.unwrap();
|
|
|
|
// Should get an error message back
|
|
let text = recv_text(&mut ws).await;
|
|
let parsed: serde_json::Value = serde_json::from_str(&text).unwrap();
|
|
assert_eq!(parsed["type"], "error");
|
|
assert!(
|
|
parsed["message"]
|
|
.as_str()
|
|
.unwrap()
|
|
.contains("Invalid message")
|
|
);
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_unknown_type_returns_error() {
|
|
let (addr, _state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
|
|
// Send valid JSON but unknown message type
|
|
ws.send(Message::Text(r#"{"type":"foobar"}"#.into()))
|
|
.await
|
|
.unwrap();
|
|
|
|
let text = recv_text(&mut ws).await;
|
|
let parsed: serde_json::Value = serde_json::from_str(&text).unwrap();
|
|
assert_eq!(parsed["type"], "error");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_gateway_status_endpoint() {
|
|
let (addr, _state, _agent_rx) = start_test_server().await;
|
|
|
|
// Connect a WS client
|
|
let _ws = connect_ws(addr).await;
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
|
|
// Hit the status endpoint
|
|
let client = reqwest::Client::new();
|
|
let resp = client
|
|
.get(format!("http://{}/api/gateway/status", addr))
|
|
.header("Authorization", format!("Bearer {}", AUTH_TOKEN))
|
|
.send()
|
|
.await
|
|
.expect("Failed to fetch status");
|
|
|
|
assert_eq!(resp.status(), 200);
|
|
|
|
let body: serde_json::Value = resp.json().await.unwrap();
|
|
assert_eq!(body["ws_connections"], 1);
|
|
assert!(body["total_connections"].as_u64().unwrap() >= 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_no_auth_rejected() {
|
|
let (addr, _state, _agent_rx) = start_test_server().await;
|
|
|
|
// Try to connect without auth token
|
|
let url = format!("ws://{}/api/chat/ws", addr);
|
|
let request = url.into_client_request().unwrap();
|
|
let result = tokio_tungstenite::connect_async(request).await;
|
|
|
|
// Should fail (401 from auth middleware before WS upgrade)
|
|
assert!(result.is_err());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_ws_multiple_events_in_sequence() {
|
|
let (addr, state, _agent_rx) = start_test_server().await;
|
|
let mut ws = connect_ws(addr).await;
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
|
|
// Broadcast multiple events rapidly
|
|
state.sse.broadcast(SseEvent::Thinking {
|
|
message: "step 1".to_string(),
|
|
thread_id: None,
|
|
});
|
|
state.sse.broadcast(SseEvent::ToolStarted {
|
|
name: "shell".to_string(),
|
|
thread_id: None,
|
|
});
|
|
state.sse.broadcast(SseEvent::ToolCompleted {
|
|
name: "shell".to_string(),
|
|
success: true,
|
|
thread_id: None,
|
|
});
|
|
state.sse.broadcast(SseEvent::Response {
|
|
content: "done".to_string(),
|
|
thread_id: "t1".to_string(),
|
|
});
|
|
|
|
// Receive all 4 in order
|
|
let t1 = recv_text(&mut ws).await;
|
|
let t2 = recv_text(&mut ws).await;
|
|
let t3 = recv_text(&mut ws).await;
|
|
let t4 = recv_text(&mut ws).await;
|
|
|
|
let p1: serde_json::Value = serde_json::from_str(&t1).unwrap();
|
|
let p2: serde_json::Value = serde_json::from_str(&t2).unwrap();
|
|
let p3: serde_json::Value = serde_json::from_str(&t3).unwrap();
|
|
let p4: serde_json::Value = serde_json::from_str(&t4).unwrap();
|
|
|
|
assert_eq!(p1["event_type"], "thinking");
|
|
assert_eq!(p2["event_type"], "tool_started");
|
|
assert_eq!(p3["event_type"], "tool_completed");
|
|
assert_eq!(p4["event_type"], "response");
|
|
|
|
ws.close(None).await.unwrap();
|
|
}
|