Skip to main content

weft_core/api/
mod.rs

1pub mod app_detail;
2pub mod apps;
3pub mod capabilities;
4pub mod core_capabilities;
5pub mod generations;
6pub mod health;
7pub mod openai_compat;
8pub mod package_webhook;
9pub mod package_ws;
10pub mod packages;
11pub mod packages_runtime;
12pub mod plans;
13pub mod profile;
14pub mod providers;
15pub mod rpc;
16pub mod scenes;
17pub mod services;
18
19use axum::extract::{Request, State};
20use axum::http::{header::AUTHORIZATION, HeaderValue, StatusCode};
21use axum::middleware::{self, Next};
22use axum::response::{IntoResponse, Response};
23use axum::routing::{delete, get, post};
24use axum::{Json, Router};
25use openai_compat::AppState;
26use tower_http::cors::{Any, CorsLayer};
27use tower_http::services::ServeDir;
28
29fn route_is_unprotected(path: &str) -> bool {
30    path == "/health"
31        || path == "/api/health"
32        || path == "/v1/models"
33        // 生成的媒体(图/视频)经 HTTP 提供给 web UI 的 <img>/<video> 加载,
34        // 这些标签无法携带 Bearer token,故放行;loopback 绑定已限制访问面。
35        || path.starts_with("/media/")
36        // web UI 静态资源(HTML/JS/CSS),webview 首次加载不带 token header。
37        || path.starts_with("/app-ui")
38        // 包 web UI 静态资源(各包的 ui/dist/),同理 webview 子资源无法带 header。
39        || (path.starts_with("/packages/") && path.contains("/ui/"))
40        // 页面代理(iframe 嵌入原文),同理无法带 header。
41        || path.starts_with("/proxy/page")
42}
43
44fn matches_bearer_token(value: &HeaderValue, token: &str) -> bool {
45    value
46        .to_str()
47        .ok()
48        .and_then(|header| header.strip_prefix("Bearer "))
49        .map(|candidate| candidate == token)
50        .unwrap_or(false)
51}
52
53fn unauthorized_response() -> Response {
54    (
55        StatusCode::UNAUTHORIZED,
56        Json(serde_json::json!({
57            "error": "missing or invalid loopback bearer token"
58        })),
59    )
60        .into_response()
61}
62
63async fn require_loopback_token(
64    State(state): State<AppState>,
65    request: Request,
66    next: Next,
67) -> Response {
68    let path = request.uri().path();
69    // CORS 预检(OPTIONS)永不携带 token,必须放行交给 CorsLayer 处理,
70    // 否则 401 且无 CORS 头 → 浏览器侧 web UI 跨源请求全被拦截。
71    if request.method() == axum::http::Method::OPTIONS {
72        return next.run(request).await;
73    }
74    if route_is_unprotected(path) {
75        return next.run(request).await;
76    }
77
78    if path.starts_with("/ws/packages/") {
79        let token = state.runtime_token.as_deref().unwrap_or_default();
80        let authorized = request
81            .uri()
82            .query()
83            .and_then(|query| {
84                query.split('&').find_map(|entry| {
85                    let (key, value) = entry.split_once('=')?;
86                    if key == "token" { Some(value) } else { None }
87                })
88            })
89            .map(|value| value == token)
90            .unwrap_or(false);
91        return if authorized {
92            next.run(request).await
93        } else {
94            unauthorized_response()
95        };
96    }
97
98    // No runtime token configured → loopback auth is not armed; pass through.
99    // The token file's presence is what enables auth, so dev/test and
100    // not-yet-provisioned deployments keep working.
101    let Some(token) = state.runtime_token.as_deref() else {
102        return next.run(request).await;
103    };
104
105    if let Some(header) = request.headers().get(AUTHORIZATION) {
106        if matches_bearer_token(header, token) {
107            return next.run(request).await;
108        }
109    }
110
111    unauthorized_response()
112}
113
114/// GET /api/config/registry — get registry config
115async fn get_registry_config(State(state): State<AppState>) -> Json<serde_json::Value> {
116    let config = state.config.read().await;
117    Json(serde_json::json!({
118        "gitea_url": config.registry.gitea_url,
119        "gitea_token": config.registry.gitea_token.as_ref().map(|t| {
120            if t.len() > 4 {
121                format!("****{}", &t[t.len()-4..])
122            } else {
123                "****".to_string()
124            }
125        })
126    }))
127}
128
129/// PUT /api/config/registry — update registry config
130async fn update_registry_config(
131    State(state): State<AppState>,
132    Json(body): Json<serde_json::Value>,
133) -> Result<Json<serde_json::Value>, (StatusCode, Json<serde_json::Value>)> {
134    let gitea_url = body["gitea_url"].as_str().map(|s| s.to_string());
135    let gitea_token = body["gitea_token"].as_str().map(|s| s.to_string());
136
137    // Update in-memory config
138    {
139        let mut config = state.config.write().await;
140        if let Some(url) = gitea_url {
141            config.registry.gitea_url = url;
142        }
143        if let Some(token) = gitea_token {
144            config.registry.gitea_token = if token.is_empty() { None } else { Some(token) };
145        }
146    }
147
148    // Persist to disk
149    {
150        let config = state.config.read().await;
151        if let Err(e) = crate::config::store::save_config(&state.config_path, &config) {
152            tracing::error!("Failed to save config: {}", e);
153            return Err((
154                StatusCode::INTERNAL_SERVER_ERROR,
155                Json(serde_json::json!({"error": format!("Failed to save config: {}", e)})),
156            ));
157        }
158    }
159
160    Ok(Json(
161        serde_json::json!({"status": "ok", "message": "Registry config updated"}),
162    ))
163}
164
165async fn shutdown_core(State(state): State<AppState>) -> Json<serde_json::Value> {
166    if let Some(tx) = state.shutdown_tx.lock().unwrap().take() {
167        let _ = tx.send(());
168    }
169
170    Json(serde_json::json!({
171        "status": "ok",
172        "message": "shutdown requested",
173    }))
174}
175
176/// GET /api/stream/tokens?session_id=xxx
177/// Returns pending stream tokens for a session and clears them from the buffer.
178/// This endpoint is lock-free relative to WASM execution.
179async fn stream_tokens(
180    State(state): State<AppState>,
181    axum::extract::Query(params): axum::extract::Query<std::collections::HashMap<String, String>>,
182) -> Json<serde_json::Value> {
183    let session_id = params.get("session_id").cloned().unwrap_or_default();
184    if session_id.is_empty() {
185        return Json(serde_json::json!({"tokens": [], "error": "missing session_id"}));
186    }
187    let tokens = state
188        .stream_buffer
189        .lock()
190        .map(|mut buf| buf.remove(&session_id).unwrap_or_default())
191        .unwrap_or_default();
192    Json(serde_json::json!({"tokens": tokens}))
193}
194
195/// GET /api/stream/events?session_id=xxx&after_seq=N
196/// Reads session events directly from SQLite, bypassing the WASM lock.
197/// This allows real-time polling while send_message is executing.
198async fn stream_events(
199    axum::extract::Query(params): axum::extract::Query<std::collections::HashMap<String, String>>,
200) -> Json<serde_json::Value> {
201    let session_id = params.get("session_id").cloned().unwrap_or_default();
202    if session_id.is_empty() {
203        return Json(serde_json::json!({"events": [], "latest_seq": 0, "error": "missing session_id"}));
204    }
205    let after_seq: i64 = params.get("after_seq").and_then(|v| v.parse().ok()).unwrap_or(0);
206
207    let db_path = std::env::var("WEFT_SESSION_EVENTS_DB")
208        .unwrap_or_else(|_| "./data/session-events/session-events.sqlite".to_string());
209
210    let result = tokio::task::spawn_blocking(move || {
211        let conn = rusqlite::Connection::open(&db_path)?;
212        let mut stmt = conn.prepare(
213            "SELECT seq, event_id, event_type, payload_json, created_at \
214             FROM session_events WHERE session_id = ?1 AND seq > ?2 \
215             ORDER BY seq ASC LIMIT 200",
216        )?;
217        let rows: Vec<serde_json::Value> = stmt
218            .query_map(rusqlite::params![session_id, after_seq], |row| {
219                Ok((
220                    row.get::<_, i64>(0)?,
221                    row.get::<_, String>(1)?,
222                    row.get::<_, String>(2)?,
223                    row.get::<_, String>(3)?,
224                    row.get::<_, i64>(4)?,
225                ))
226            })?
227            .filter_map(|r| r.ok())
228            .map(|(seq, event_id, event_type, payload_json, created_at)| {
229                let payload = serde_json::from_str(&payload_json).unwrap_or(serde_json::Value::Null);
230                serde_json::json!({
231                    "seq": seq,
232                    "event_id": event_id,
233                    "type": event_type,
234                    "payload": payload,
235                    "created_at": created_at,
236                })
237            })
238            .collect();
239
240        let latest_seq: i64 = conn
241            .query_row(
242                "SELECT COALESCE(MAX(seq), 0) FROM session_events WHERE session_id = ?1",
243                rusqlite::params![session_id],
244                |row| row.get(0),
245            )
246            .unwrap_or(after_seq);
247
248        Ok::<_, rusqlite::Error>((rows, latest_seq))
249    })
250    .await;
251
252    match result {
253        Ok(Ok((events, latest_seq))) => Json(serde_json::json!({
254            "events": events,
255            "latest_seq": latest_seq,
256        })),
257        _ => Json(serde_json::json!({"events": [], "latest_seq": after_seq})),
258    }
259}
260
261/// 包 web UI 静态文件服务。从 packages/official/{name}/ui/dist/ 或
262/// packages/installed/{name}/ui/dist/ 读取文件,按扩展名推断 content-type。
263/// 任何包只要有 ui/dist/ 目录就自动被托管,无需硬编码。
264async fn serve_package_ui(
265    axum::extract::Path((name, path)): axum::extract::Path<(String, String)>,
266) -> Response {
267    if name.contains("..") || path.contains("..") {
268        return (StatusCode::BAD_REQUEST, "invalid path").into_response();
269    }
270    // 查找顺序:installed 优先(用户覆盖),然后 official
271    let candidates = [
272        format!("./packages/installed/{}/ui/dist/{}", name, path),
273        format!("./packages/official/{}/ui/dist/{}", name, path),
274        format!("./packages/installed/{}/ui/{}", name, path),
275        format!("./packages/official/{}/ui/{}", name, path),
276    ];
277    for candidate in &candidates {
278        let full = std::path::Path::new(candidate);
279        if let Ok(bytes) = tokio::fs::read(full).await {
280            let ct = match full.extension().and_then(|e| e.to_str()) {
281                Some("html") => "text/html; charset=utf-8",
282                Some("js") | Some("mjs") => "application/javascript; charset=utf-8",
283                Some("css") => "text/css; charset=utf-8",
284                Some("json") => "application/json",
285                Some("svg") => "image/svg+xml",
286                Some("png") => "image/png",
287                Some("ico") => "image/x-icon",
288                Some("woff2") => "font/woff2",
289                Some("woff") => "font/woff",
290                _ => "application/octet-stream",
291            };
292            return ([(axum::http::header::CONTENT_TYPE, ct)], bytes).into_response();
293        }
294    }
295    (StatusCode::NOT_FOUND, "package ui resource not found").into_response()
296}
297
298/// 提供生成媒体文件给 web UI。路径形如 /media/image-gen/img-1.png,
299/// 映射到 Core 工作目录下的 ./workspace/<path>。仅读,禁路径穿越。
300async fn serve_media(
301    axum::extract::Path(path): axum::extract::Path<String>,
302) -> Response {
303    if path.contains("..") {
304        return (StatusCode::BAD_REQUEST, "invalid path").into_response();
305    }
306    let full = std::path::Path::new("./workspace").join(&path);
307    match tokio::fs::read(&full).await {
308        Ok(bytes) => {
309            let ct = match full.extension().and_then(|e| e.to_str()) {
310                Some("png") => "image/png",
311                Some("jpg") | Some("jpeg") => "image/jpeg",
312                Some("webp") => "image/webp",
313                Some("gif") => "image/gif",
314                Some("mp4") => "video/mp4",
315                Some("webm") => "video/webm",
316                _ => "application/octet-stream",
317            };
318            ([(axum::http::header::CONTENT_TYPE, ct)], bytes).into_response()
319        }
320        Err(_) => (StatusCode::NOT_FOUND, "not found").into_response(),
321    }
322}
323
324pub fn build_router(state: AppState) -> Router {
325    let cors = CorsLayer::new()
326        .allow_origin(Any)
327        .allow_methods(Any)
328        .allow_headers(Any);
329
330    Router::new()
331        // web UI 静态资源(clients/web-canvas/dist),webview 加载 /app-ui/ 即可使用。
332        .nest_service("/app-ui", ServeDir::new("./clients/web-canvas/dist"))
333        // 包 web UI 静态资源:/packages/{name}/ui/{path} → packages/official/{name}/ui/dist/{path}
334        // 任何包只要有 ui/dist/ 目录就自动被托管,无需写死。
335        .route("/packages/{name}/ui/{*path}", get(serve_package_ui))
336        // 生成媒体(图/视频)的 HTTP 访问:web UI 用 /media/<workspace下相对路径> 加载。
337        .route("/media/{*path}", get(serve_media))
338        // OpenAI-compatible API
339        .route(
340            "/v1/chat/completions",
341            post(openai_compat::chat_completions),
342        )
343        .route("/v1/models", get(openai_compat::list_models))
344        // Management API
345        .route("/health", get(health::health))
346        .route("/api/health", get(health::health))
347        .route("/api/apps", get(apps::list_apps))
348        .route("/api/apps/{name}", get(app_detail::get_app))
349        .route("/api/apps/{name}/scenes", get(scenes::list_scenes))
350        .route("/api/apps/{name}/scenes", post(scenes::create_scene))
351        .route("/api/apps/{name}/scenes/{scene}", get(scenes::get_scene))
352        .route(
353            "/api/apps/{name}/scenes/{scene}",
354            axum::routing::delete(scenes::delete_scene),
355        )
356        .route(
357            "/api/apps/{name}/scenes/{scene}/bind",
358            post(scenes::bind_scene),
359        )
360        .route("/api/apps/{name}/health", get(app_detail::app_health))
361        .route("/api/apps/{name}/run", post(app_detail::app_run))
362        .route("/api/capabilities", get(capabilities::list_capabilities))
363        .route(
364            "/api/capabilities/{name}",
365            get(capabilities::get_capability),
366        )
367        .route(
368            "/api/capabilities/{name}/call",
369            post(capabilities::capability_call),
370        )
371        .route("/api/packages", get(packages::list_packages))
372        .route("/api/plans/activate", post(plans::activation_plan))
373        .route("/api/profile", get(profile::get_profile))
374        .route("/api/profile", axum::routing::put(profile::set_profile))
375        .route("/api/policy", get(profile::get_policy))
376        .route(
377            "/api/apps/{name}/generations",
378            get(generations::get_generation),
379        )
380        .route(
381            "/api/apps/{name}/generations/diagnostics",
382            get(generations::get_generation_diagnostics),
383        )
384        .route(
385            "/api/apps/{name}/generations/{from}/diff/{to}",
386            get(generations::get_generation_diff),
387        )
388        .route(
389            "/api/apps/{name}/propose",
390            post(generations::propose_generation),
391        )
392        .route(
393            "/api/apps/{name}/verify",
394            post(generations::verify_generation),
395        )
396        .route(
397            "/api/apps/{name}/activate",
398            post(generations::activate_generation),
399        )
400        .route(
401            "/api/apps/{name}/generations/{id}/activate",
402            post(generations::activate_existing_generation),
403        )
404        .route(
405            "/api/apps/{name}/rollback",
406            post(generations::rollback_generation),
407        )
408        .route("/api/shutdown", post(shutdown_core))
409        .route("/rpc", post(rpc::rpc_endpoint))
410        .route("/api/providers", get(providers::list_providers))
411        .route("/api/providers", post(providers::create_provider))
412        .route(
413            "/api/providers/fetch-models",
414            post(providers::fetch_models),
415        )
416        .route("/api/providers/{name}", get(providers::get_provider))
417        .route(
418            "/api/providers/{name}",
419            axum::routing::put(providers::update_provider),
420        )
421        .route("/api/providers/{name}", delete(providers::delete_provider))
422        .route("/api/providers/{name}/upstream-models", get(providers::list_upstream_models))
423        .route("/api/routing", axum::routing::put(providers::update_routing))
424        .route("/api/services", get(services::list_services))
425        .route("/api/services/{name}/start", post(services::start_service))
426        .route("/api/services/{name}/stop", post(services::stop_service))
427        .route(
428            "/api/services/{name}/restart",
429            post(services::restart_service),
430        )
431        .route(
432            "/api/services/{name}/webhook",
433            post(services::proxy_webhook),
434        )
435        // Package runtime management
436        .route(
437            "/api/packages/runtime",
438            get(packages_runtime::list_packages),
439        )
440        .route(
441            "/api/packages/install",
442            post(packages_runtime::install_package),
443        )
444        .route(
445            "/api/packages/{name}/reload",
446            post(packages_runtime::reload_package),
447        )
448        .route(
449            "/api/packages/{name}/toggle",
450            post(packages_runtime::toggle_package),
451        )
452        .route(
453            "/api/packages/{name}",
454            delete(packages_runtime::uninstall_package),
455        )
456        .route(
457            "/api/packages/{name}/dependencies",
458            get(packages_runtime::get_package_dependencies),
459        )
460        .route(
461            "/api/packages/{name}/config/schema",
462            get(packages_runtime::get_package_config_schema),
463        )
464        .route(
465            "/api/packages/{name}/config",
466            get(packages_runtime::get_package_config),
467        )
468        .route(
469            "/api/packages/{name}/config",
470            axum::routing::put(packages_runtime::save_package_config),
471        )
472        // Package Webhook
473        .route(
474            "/api/packages/{package_name}/webhook",
475            post(package_webhook::package_webhook_no_channel),
476        )
477        .route(
478            "/api/packages/{package_name}/webhook/{channel_type}",
479            post(package_webhook::package_webhook_with_channel),
480        )
481        .route(
482            "/api/packages/{package_name}/call",
483            post(package_ws::package_call),
484        )
485        // Package WebSocket
486        .route(
487            "/ws/packages/{package_name}",
488            get(package_ws::package_websocket),
489        )
490        // Config management
491        .route("/api/config/registry", get(get_registry_config))
492        .route(
493            "/api/config/registry",
494            axum::routing::put(update_registry_config),
495        )
496        .route("/api/stream/tokens", get(stream_tokens))
497        .route("/api/stream/events", get(stream_events))
498        .layer(cors)
499        .layer(middleware::from_fn_with_state(
500            state.clone(),
501            require_loopback_token,
502        ))
503        .with_state(state)
504}
505
506#[cfg(test)]
507mod tests {
508    use super::build_router;
509    use crate::api::openai_compat::AppState;
510    use crate::app::{
511        AppProfile, CapabilityRegistry, CorePolicy, GenerationStoreMap, PackageIndex,
512        PackageSource, ResolvedAppMap,
513    };
514    use crate::config::{
515        AppConfig, CoreConfig, FallbackConfig, KeyStrategyConfig, RegistryConfig, RoutingConfig,
516    };
517    use crate::defaults::{
518        error_handler::DefaultErrorHandler, key_selectors::FailoverSelector, router::DefaultRouter,
519    };
520    use crate::pipeline::Pipeline;
521    use crate::process::ProcessManager;
522    use crate::vkeys::VirtualKeyStore;
523    use axum::body::{to_bytes, Body};
524    use axum::http::{Request, StatusCode as HttpStatusCode};
525    use std::collections::HashMap;
526    use std::sync::{Arc, Mutex as StdMutex};
527    use tempfile::TempDir;
528    use tokio::sync::RwLock;
529    use tower::util::ServiceExt;
530
531    fn test_state(repo_root: std::path::PathBuf) -> AppState {
532        test_state_with_profile(repo_root, AppProfile::Developer)
533    }
534
535    fn test_state_with_profile(
536        repo_root: std::path::PathBuf,
537        active_profile: AppProfile,
538    ) -> AppState {
539        AppState {
540            config: Arc::new(RwLock::new(AppConfig {
541                core: CoreConfig::default(),
542                providers: vec![],
543                routing: RoutingConfig::default(),
544                key_strategy: KeyStrategyConfig::default(),
545                fallback: FallbackConfig::default(),
546                virtual_keys: vec![],
547                services: vec![],
548                packages: vec![],
549                registry: RegistryConfig::default(),
550                package_aliases: HashMap::new(),
551                web_search: Default::default(),
552                team: Default::default(),
553            })),
554            config_path: repo_root.join("config").join("config.toml"),
555            pipeline: Arc::new(Pipeline {
556                router: Arc::new(DefaultRouter {
557                    default_provider: "".into(),
558                }),
559                key_selector: Arc::new(FailoverSelector),
560                transforms: Arc::new(crate::defaults::transforms::TransformRegistry::with_defaults()),
561                error_handler: Arc::new(DefaultErrorHandler { max_retries: 0 }),
562                http_client: reqwest::Client::new(),
563            }),
564            process_manager: Arc::new(ProcessManager::new()),
565            vkey_store: Arc::new(VirtualKeyStore::new()),
566            package_manager: Arc::new(RwLock::new(crate::package::PackageManager::new())),
567            wasm_handle: Arc::new(RwLock::new(None)),
568            native_handle: Arc::new(RwLock::new(None)),
569            resolved_apps: Arc::new(RwLock::new(ResolvedAppMap::new())),
570            capability_registry: Arc::new(RwLock::new(CapabilityRegistry::new())),
571            active_profile: Arc::new(RwLock::new(active_profile)),
572            core_policy: Arc::new(CorePolicy::default_policy()),
573            generation_store: Arc::new(RwLock::new(GenerationStoreMap::new())),
574            package_index: Arc::new(PackageIndex {
575                version: 1,
576                revision: "test".into(),
577                source_url: "local://packages".into(),
578                package_sources: vec![PackageSource {
579                    name: "weft-claw-ui".into(),
580                    kind: "embedded".into(),
581                    package_kind: "provider".into(),
582                    runtime_provider: "weft-claw-ui".into(),
583                    current_source: "packages/installed/weft-claw".into(),
584                    trusted: true,
585                    signature: "builtin:test".into(),
586                    source_authority: "test".into(),
587                    source_public_keys: vec![],
588                    provides: vec!["ui.surface".into()],
589                    requires: vec![],
590                }],
591            }),
592            repo_root,
593            data_dir: std::path::PathBuf::from("data"),
594            runtime_token: None,
595            runtime_token_path: None,
596            chat_providers: Arc::new(RwLock::new(vec![])),
597            shutdown_tx: Arc::new(StdMutex::new(None)),
598            stream_buffer: Arc::new(StdMutex::new(std::collections::HashMap::new())),
599        }
600    }
601
602    fn create_package_ui_fixture() -> TempDir {
603        let dir = TempDir::new().expect("temp dir");
604        let package_dir = dir
605            .path()
606            .join("packages")
607            .join("installed")
608            .join("weft-claw")
609            .join("ui");
610        std::fs::create_dir_all(&package_dir).expect("create package ui dir");
611        std::fs::write(
612            package_dir.join("index.html"),
613            "<html><body>weft claw ui</body></html>",
614        )
615        .expect("write package ui html");
616        std::fs::write(
617            dir.path()
618                .join("packages")
619                .join("installed")
620                .join("weft-claw")
621                .join("package.toml"),
622            "[identity]\nname = \"weft-claw-ui\"\nversion = \"0.1.0\"\ndescription = \"test\"\n\n[package]\nentry = \"ui/index.html\"\nruntime = \"embedded\"\napi_version = \"v1\"\n",
623        )
624        .expect("write package manifest");
625        dir
626    }
627
628    #[tokio::test]
629    async fn package_ui_route_serves_installed_package_assets_from_repo_root() {
630        let fixture = create_package_ui_fixture();
631        let app = build_router(test_state(fixture.path().to_path_buf()));
632
633        let response = app
634            .oneshot(
635                Request::builder()
636                    .uri("/packages/weft-claw-ui/ui/index.html")
637                    .body(Body::empty())
638                    .unwrap(),
639            )
640            .await
641            .unwrap();
642
643        assert_eq!(response.status(), HttpStatusCode::OK);
644    }
645
646    #[tokio::test]
647    async fn activation_plan_route_reports_metadata_without_mutating_runtime_state() {
648        let fixture = TempDir::new().expect("temp dir");
649        let package_dir = fixture.path().join("materialized-package");
650        std::fs::create_dir_all(&package_dir).expect("package dir");
651        std::fs::write(
652            package_dir.join("package.toml"),
653            r#"[package_info]
654name = "plan-only-package"
655version = "0.1.0"
656description = "plan only"
657entry = "package.wasm"
658"#,
659        )
660        .expect("manifest");
661        std::fs::write(package_dir.join("package.wasm"), b"\0asm").expect("entry");
662
663        let state = test_state(fixture.path().to_path_buf());
664        let app = build_router(state.clone());
665        let response = app
666            .oneshot(
667                Request::builder()
668                    .method("POST")
669                    .uri("/api/plans/activate")
670                    .header("content-type", "application/json")
671                    .body(Body::from(
672                        serde_json::json!({
673                            "materialized_path": package_dir.display().to_string()
674                        })
675                        .to_string(),
676                    ))
677                    .expect("request"),
678            )
679            .await
680            .expect("response");
681
682        assert_eq!(response.status(), HttpStatusCode::OK);
683        let body = to_bytes(response.into_body(), usize::MAX)
684            .await
685            .expect("body bytes");
686        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
687        assert_eq!(
688            payload["status"],
689            serde_json::json!("activation_plan_ready")
690        );
691        assert_eq!(payload["plan_only"], serde_json::json!(true));
692        assert_eq!(payload["metadata_only"], serde_json::json!(true));
693        assert_eq!(payload["activation_performed"], serde_json::json!(false));
694        assert_eq!(payload["mutation_performed"], serde_json::json!(false));
695        assert_eq!(payload["lock_mutation_performed"], serde_json::json!(false));
696        assert_eq!(payload["activation_required"], serde_json::json!(true));
697        assert_eq!(payload["ready_for_activation"], serde_json::json!(true));
698        assert_eq!(
699            payload["package"]["name"],
700            serde_json::json!("plan-only-package")
701        );
702        assert!(payload["checks"]
703            .as_array()
704            .expect("checks")
705            .iter()
706            .all(|check| { check["ok"].as_bool().unwrap_or(false) }));
707        assert!(state.package_manager.read().await.list().is_empty());
708    }
709
710    #[tokio::test]
711    async fn activation_plan_route_blocks_missing_manifest_without_mutating_runtime_state() {
712        let fixture = TempDir::new().expect("temp dir");
713        let package_dir = fixture.path().join("materialized-package");
714        std::fs::create_dir_all(&package_dir).expect("package dir");
715
716        let state = test_state(fixture.path().to_path_buf());
717        let app = build_router(state.clone());
718        let response = app
719            .oneshot(
720                Request::builder()
721                    .method("POST")
722                    .uri("/api/plans/activate")
723                    .header("content-type", "application/json")
724                    .body(Body::from(
725                        serde_json::json!({
726                            "materialized_path": package_dir.display().to_string()
727                        })
728                        .to_string(),
729                    ))
730                    .expect("request"),
731            )
732            .await
733            .expect("response");
734
735        assert_eq!(response.status(), HttpStatusCode::OK);
736        let body = to_bytes(response.into_body(), usize::MAX)
737            .await
738            .expect("body bytes");
739        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
740        assert_eq!(
741            payload["status"],
742            serde_json::json!("activation_plan_blocked")
743        );
744        assert_eq!(payload["plan_only"], serde_json::json!(true));
745        assert_eq!(payload["activation_performed"], serde_json::json!(false));
746        assert_eq!(payload["mutation_performed"], serde_json::json!(false));
747        assert_eq!(payload["manifest_found"], serde_json::json!(false));
748        assert_eq!(payload["activation_required"], serde_json::json!(false));
749        assert_eq!(payload["ready_for_activation"], serde_json::json!(false));
750        assert!(state.package_manager.read().await.list().is_empty());
751    }
752
753    #[tokio::test]
754    async fn activation_plan_route_apply_true_registers_service_metadata_from_temp_root() {
755        let fixture = TempDir::new().expect("temp dir");
756        let package_dir = fixture.path().join("materialized-service-package");
757        std::fs::create_dir_all(&package_dir).expect("package dir");
758        std::fs::write(
759            package_dir.join("package.toml"),
760            r#"[package_info]
761name = "controlled-service-package"
762version = "0.1.0"
763description = "controlled service"
764entry = "server.js"
765provides = ["chat_channel"]
766chat_endpoint = "/chat"
767
768[package]
769runtime = "service"
770entry = "server.js"
771"#,
772        )
773        .expect("manifest");
774        std::fs::write(package_dir.join("server.js"), "console.log('service');\n").expect("entry");
775
776        let state = test_state(fixture.path().to_path_buf());
777        let app = build_router(state.clone());
778        let response = app
779            .oneshot(
780                Request::builder()
781                    .method("POST")
782                    .uri("/api/plans/activate")
783                    .header("content-type", "application/json")
784                    .body(Body::from(
785                        serde_json::json!({
786                            "materialized_path": package_dir.display().to_string(),
787                            "apply": true,
788                            "confirm": true
789                        })
790                        .to_string(),
791                    ))
792                    .expect("request"),
793            )
794            .await
795            .expect("response");
796
797        assert_eq!(response.status(), HttpStatusCode::OK);
798        let body = to_bytes(response.into_body(), usize::MAX)
799            .await
800            .expect("body bytes");
801        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
802        assert_eq!(
803            payload["status"],
804            serde_json::json!("activation_metadata_registered")
805        );
806        assert_eq!(payload["plan_only"], serde_json::json!(false));
807        assert_eq!(payload["activation_performed"], serde_json::json!(true));
808        assert_eq!(payload["mutation_performed"], serde_json::json!(true));
809        assert_eq!(payload["lock_mutation_performed"], serde_json::json!(false));
810        assert_eq!(payload["service_registered"], serde_json::json!(true));
811        assert_eq!(payload["runtime_started"], serde_json::json!(false));
812        assert!(state
813            .package_manager
814            .read()
815            .await
816            .get("controlled-service-package")
817            .is_some());
818        assert_eq!(state.chat_providers.read().await.len(), 1);
819    }
820
821    #[tokio::test]
822    async fn activation_plan_route_apply_true_requires_confirm_and_does_not_mutate() {
823        let fixture = TempDir::new().expect("temp dir");
824        let package_dir = fixture.path().join("materialized-service-package");
825        std::fs::create_dir_all(&package_dir).expect("package dir");
826        std::fs::write(
827            package_dir.join("package.toml"),
828            r#"[package_info]
829name = "unconfirmed-plugin"
830version = "0.1.0"
831description = "unconfirmed"
832entry = "server.js"
833
834[package]
835runtime = "service"
836entry = "server.js"
837"#,
838        )
839        .expect("manifest");
840        std::fs::write(package_dir.join("server.js"), "console.log('service');\n").expect("entry");
841
842        let state = test_state(fixture.path().to_path_buf());
843        let app = build_router(state.clone());
844        let response = app
845            .oneshot(
846                Request::builder()
847                    .method("POST")
848                    .uri("/api/plans/activate")
849                    .header("content-type", "application/json")
850                    .body(Body::from(
851                        serde_json::json!({
852                            "materialized_path": package_dir.display().to_string(),
853                            "apply": true
854                        })
855                        .to_string(),
856                    ))
857                    .expect("request"),
858            )
859            .await
860            .expect("response");
861
862        assert_eq!(response.status(), HttpStatusCode::OK);
863        let body = to_bytes(response.into_body(), usize::MAX)
864            .await
865            .expect("body bytes");
866        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
867        assert_eq!(
868            payload["status"],
869            serde_json::json!("activation_apply_blocked")
870        );
871        assert_eq!(payload["activation_performed"], serde_json::json!(false));
872        assert_eq!(payload["mutation_performed"], serde_json::json!(false));
873        assert!(payload["validation_issues"]
874            .as_array()
875            .expect("validation issues")
876            .iter()
877            .any(|issue| issue
878                .as_str()
879                .unwrap_or_default()
880                .contains("explicit_confirmation_required_for_apply")));
881        assert!(state.package_manager.read().await.list().is_empty());
882    }
883
884    #[tokio::test]
885    async fn activation_plan_route_apply_true_blocks_native_runtime_without_mutating() {
886        let fixture = TempDir::new().expect("temp dir");
887        let package_dir = fixture.path().join("materialized-native-package");
888        std::fs::create_dir_all(&package_dir).expect("package dir");
889        std::fs::write(
890            package_dir.join("package.toml"),
891            r#"[package_info]
892name = "native-plugin"
893version = "0.1.0"
894description = "native"
895entry = "native.dll"
896
897[package]
898runtime = "native"
899entry = "native.dll"
900"#,
901        )
902        .expect("manifest");
903
904        let state = test_state(fixture.path().to_path_buf());
905        let app = build_router(state.clone());
906        let response = app
907            .oneshot(
908                Request::builder()
909                    .method("POST")
910                    .uri("/api/plans/activate")
911                    .header("content-type", "application/json")
912                    .body(Body::from(
913                        serde_json::json!({
914                            "materialized_path": package_dir.display().to_string(),
915                            "apply": true,
916                            "confirm": true
917                        })
918                        .to_string(),
919                    ))
920                    .expect("request"),
921            )
922            .await
923            .expect("response");
924
925        assert_eq!(response.status(), HttpStatusCode::OK);
926        let body = to_bytes(response.into_body(), usize::MAX)
927            .await
928            .expect("body bytes");
929        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
930        assert_eq!(
931            payload["status"],
932            serde_json::json!("activation_apply_blocked")
933        );
934        assert_eq!(payload["activation_performed"], serde_json::json!(false));
935        assert!(payload["validation_issues"]
936            .as_array()
937            .expect("validation issues")
938            .iter()
939            .any(|issue| issue
940                .as_str()
941                .unwrap_or_default()
942                .contains("runtime_safe_for_controlled_activation")));
943        assert!(state.package_manager.read().await.list().is_empty());
944    }
945
946    #[tokio::test]
947    async fn activation_plan_route_developer_blocks_trusted_native_plan() {
948        let fixture = TempDir::new().expect("temp dir");
949        let package_dir = fixture.path().join("developer-native-package");
950        std::fs::create_dir_all(&package_dir).expect("package dir");
951        std::fs::write(
952            package_dir.join("package.toml"),
953            r#"[package_info]
954name = "developer-native-plugin"
955version = "0.1.0"
956description = "native"
957entry = "native.dll"
958
959[package]
960runtime = "native"
961entry = "native.dll"
962native_allowed = true
963expected_digest = "sha256:0123456789abcdef"
964"#,
965        )
966        .expect("manifest");
967        std::fs::write(package_dir.join("native.dll"), b"fake native library")
968            .expect("native file");
969
970        let state = test_state_with_profile(fixture.path().to_path_buf(), AppProfile::Developer);
971        let app = build_router(state.clone());
972        let response = app
973            .oneshot(
974                Request::builder()
975                    .method("POST")
976                    .uri("/api/plans/activate")
977                    .header("content-type", "application/json")
978                    .body(Body::from(
979                        serde_json::json!({
980                            "materialized_path": package_dir.display().to_string(),
981                            "confirm": true
982                        })
983                        .to_string(),
984                    ))
985                    .expect("request"),
986            )
987            .await
988            .expect("response");
989
990        assert_eq!(response.status(), HttpStatusCode::OK);
991        let body = to_bytes(response.into_body(), usize::MAX)
992            .await
993            .expect("body bytes");
994        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
995        assert_eq!(
996            payload["status"],
997            serde_json::json!("activation_plan_blocked")
998        );
999        assert_eq!(
1000            payload["ready_for_trusted_native_load"],
1001            serde_json::json!(false)
1002        );
1003        assert_eq!(payload["native_load_performed"], serde_json::json!(false));
1004        assert!(payload["validation_issues"]
1005            .as_array()
1006            .expect("validation issues")
1007            .iter()
1008            .any(|issue| issue
1009                .as_str()
1010                .unwrap_or_default()
1011                .contains("trusted_native_profile_required")));
1012        assert!(state.package_manager.read().await.list().is_empty());
1013        assert!(state.native_handle.read().await.is_none());
1014    }
1015
1016    #[tokio::test]
1017    async fn activation_plan_route_trusted_native_missing_digest_blocks_plan() {
1018        let fixture = TempDir::new().expect("temp dir");
1019        let package_dir = fixture.path().join("trusted-native-missing-digest-package");
1020        std::fs::create_dir_all(&package_dir).expect("package dir");
1021        std::fs::write(
1022            package_dir.join("package.toml"),
1023            r#"[package_info]
1024name = "trusted-native-missing-digest-plugin"
1025version = "0.1.0"
1026description = "native"
1027entry = "native.dll"
1028
1029[package]
1030runtime = "native"
1031entry = "native.dll"
1032native_allowed = true
1033"#,
1034        )
1035        .expect("manifest");
1036        std::fs::write(package_dir.join("native.dll"), b"fake native library")
1037            .expect("native file");
1038
1039        let state = test_state_with_profile(fixture.path().to_path_buf(), AppProfile::Trusted);
1040        let app = build_router(state.clone());
1041        let response = app
1042            .oneshot(
1043                Request::builder()
1044                    .method("POST")
1045                    .uri("/api/plans/activate")
1046                    .header("content-type", "application/json")
1047                    .body(Body::from(
1048                        serde_json::json!({
1049                            "materialized_path": package_dir.display().to_string(),
1050                            "confirm": true
1051                        })
1052                        .to_string(),
1053                    ))
1054                    .expect("request"),
1055            )
1056            .await
1057            .expect("response");
1058
1059        assert_eq!(response.status(), HttpStatusCode::OK);
1060        let body = to_bytes(response.into_body(), usize::MAX)
1061            .await
1062            .expect("body bytes");
1063        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
1064        assert_eq!(
1065            payload["status"],
1066            serde_json::json!("activation_plan_blocked")
1067        );
1068        assert_eq!(
1069            payload["ready_for_trusted_native_load"],
1070            serde_json::json!(false)
1071        );
1072        assert!(payload["validation_issues"]
1073            .as_array()
1074            .expect("validation issues")
1075            .iter()
1076            .any(|issue| issue
1077                .as_str()
1078                .unwrap_or_default()
1079                .contains("trusted_native_expected_digest_present")));
1080        assert!(state.native_handle.read().await.is_none());
1081    }
1082
1083    #[tokio::test]
1084    async fn activation_plan_route_trusted_native_ready_returns_plan_only() {
1085        let fixture = TempDir::new().expect("temp dir");
1086        let package_dir = fixture.path().join("trusted-native-ready-package");
1087        std::fs::create_dir_all(&package_dir).expect("package dir");
1088        std::fs::write(
1089            package_dir.join("package.toml"),
1090            r#"[package_info]
1091name = "trusted-native-ready-plugin"
1092version = "0.1.0"
1093description = "native"
1094entry = "native.dll"
1095
1096[package]
1097runtime = "native"
1098entry = "native.dll"
1099native_allowed = true
1100expected_digest = "sha256:0123456789abcdef"
1101"#,
1102        )
1103        .expect("manifest");
1104        std::fs::write(package_dir.join("native.dll"), b"fake native library")
1105            .expect("native file");
1106
1107        let state = test_state_with_profile(fixture.path().to_path_buf(), AppProfile::Trusted);
1108        let app = build_router(state.clone());
1109        let response = app
1110            .oneshot(
1111                Request::builder()
1112                    .method("POST")
1113                    .uri("/api/plans/activate")
1114                    .header("content-type", "application/json")
1115                    .body(Body::from(
1116                        serde_json::json!({
1117                            "materialized_path": package_dir.display().to_string(),
1118                            "confirm": true
1119                        })
1120                        .to_string(),
1121                    ))
1122                    .expect("request"),
1123            )
1124            .await
1125            .expect("response");
1126
1127        assert_eq!(response.status(), HttpStatusCode::OK);
1128        let body = to_bytes(response.into_body(), usize::MAX)
1129            .await
1130            .expect("body bytes");
1131        let payload: serde_json::Value = serde_json::from_slice(&body).expect("json payload");
1132        assert_eq!(
1133            payload["status"],
1134            serde_json::json!("ready_for_trusted_native_load")
1135        );
1136        assert_eq!(payload["plan_only"], serde_json::json!(true));
1137        assert_eq!(payload["activation_performed"], serde_json::json!(false));
1138        assert_eq!(payload["native_load_performed"], serde_json::json!(false));
1139        assert_eq!(
1140            payload["ready_for_trusted_native_load"],
1141            serde_json::json!(true)
1142        );
1143        assert!(payload["validation_issues"].as_array().unwrap().is_empty());
1144        assert!(payload["trusted_native"]["library_candidate"]
1145            .as_str()
1146            .unwrap_or_default()
1147            .ends_with("native.dll"));
1148        assert!(state.package_manager.read().await.list().is_empty());
1149        assert!(state.native_handle.read().await.is_none());
1150    }
1151}