use super::super::common::test_service_with_external_references; use super::*; use axum::{Router, routing::get}; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, }; #[tokio::test] #[serial] async fn external_relative_chain_survives_reversed_builder_order_and_apply_never_refetches() { let fetches = Arc::new(AtomicUsize::new(0)); let root_fetches = Arc::clone(&fetches); let child_fetches = Arc::clone(&fetches); let app = Router::new() .route( "/root.yaml", get(move || { let fetches = Arc::clone(&root_fetches); async move { fetches.fetch_add(1, Ordering::SeqCst); "Item: { $ref: './child.yaml#/Item' }" } }), ) .route( "/child.yaml", get(move || { let fetches = Arc::clone(&child_fetches); async move { fetches.fetch_add(1, Ordering::SeqCst); "Item: { type: object, required: [id], properties: { id: { type: string } } }" } }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let origin = format!("http://{address}"); let registry = test_registry().await; let service = test_service_with_external_references( registry.clone(), test_storage_root("openapi_import_external_snapshot"), vec![format!("{origin}/")], ); let workspace_id = WorkspaceId::new("ws_default"); let upload = OpenApiUpload { bytes: format!( r#" openapi: 3.1.0 info: {{ title: External }} servers: [{{ url: https://api.example.test }}] paths: /items: get: operationId: listItems responses: '200': description: ok content: application/json: schema: {{ $ref: '{origin}/root.yaml#/Item' }} "# ) .into_bytes(), mime_type: "application/yaml".to_owned(), locale: OpenApiUploadLocale::En, }; let preview = service .preview_openapi_import(&workspace_id, upload) .await .unwrap(); assert_eq!(fetches.load(Ordering::SeqCst), 2); assert_eq!(preview.preview.groups[0].operations[0].output_fields, 1); let job = registry .get_import_job(&workspace_id, &preview.job_id.as_str().into()) .await .unwrap() .unwrap(); assert_eq!( job.preview_payload["dependencies"] .as_array() .unwrap() .len(), 2 ); assert_eq!( job.preview_payload["dependency_snapshots"] .as_array() .unwrap() .len(), 2 ); let applied = service .create_openapi_import( &workspace_id, &preview.job_id.as_str().into(), OpenApiImportCreatePayload { selected_operation_keys: vec!["GET /items".to_owned()], server_url: None, conflict_mode: "skip".to_owned(), }, ) .await .unwrap(); assert_eq!(applied.created.len(), 1); assert_eq!(fetches.load(Ordering::SeqCst), 2); let active_sources: i64 = sqlx::query_scalar( "select count(*) from artifact_sources where source_id like 'src_openapi_%' and lifecycle = 'active'", ) .fetch_one(registry.pool()) .await .unwrap(); assert_eq!(active_sources, 0); } #[tokio::test] #[serial] async fn missing_external_snapshot_fails_closed_before_draft_mutation() { let app = Router::new().route( "/schemas.yaml", get(|| async { "Item: { type: object, properties: { id: { type: string } } }" }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let origin = format!("http://{address}"); let registry = test_registry().await; let storage_root = test_storage_root("openapi_import_missing_external_snapshot"); let service = test_service_with_external_references( registry.clone(), storage_root.clone(), vec![format!("{origin}/")], ); let workspace_id = WorkspaceId::new("ws_default"); let preview = service .preview_openapi_import( &workspace_id, OpenApiUpload { bytes: format!( r#" openapi: 3.1.0 info: {{ title: External integrity }} servers: [{{ url: https://api.example.test }}] paths: /items: get: operationId: listItems responses: '200': description: ok content: {{ application/json: {{ schema: {{ $ref: '{origin}/schemas.yaml#/Item' }} }} }} "# ) .into_bytes(), mime_type: "application/yaml".to_owned(), locale: OpenApiUploadLocale::En, }, ) .await .unwrap(); let job = registry .get_import_job(&workspace_id, &preview.job_id.as_str().into()) .await .unwrap() .unwrap(); let dependency_source_id = job.preview_payload["dependencies"][0]["source_id"] .as_str() .unwrap(); sqlx::query( "update artifact_sources set lifecycle = 'detached', updated_at = now(), detached_at = now() where workspace_id = $1 and source_id = $2", ) .bind(workspace_id.as_str()) .bind(dependency_source_id) .execute(registry.pool()) .await .unwrap(); assert!( service .create_openapi_import( &workspace_id, &preview.job_id.as_str().into(), OpenApiImportCreatePayload { selected_operation_keys: vec!["GET /items".to_owned()], server_url: None, conflict_mode: "skip".to_owned(), }, ) .await .is_err() ); assert!( service .list_operations(&workspace_id) .await .unwrap() .is_empty() ); assert_failed_job_and_detached_sources( ®istry, &workspace_id, &preview.job_id.as_str().into(), "import_dependency_verification_failed", ) .await; let corrupt_preview = service .preview_openapi_import( &workspace_id, OpenApiUpload { bytes: format!( r#" openapi: 3.1.0 info: {{ title: Corrupt external integrity }} servers: [{{ url: https://api.example.test }}] paths: /items: get: operationId: listItemsCorrupt responses: '200': description: ok content: {{ application/json: {{ schema: {{ $ref: '{origin}/schemas.yaml#/Item' }} }} }} "# ) .into_bytes(), mime_type: "application/yaml".to_owned(), locale: OpenApiUploadLocale::En, }, ) .await .unwrap(); let corrupt_job = registry .get_import_job(&workspace_id, &corrupt_preview.job_id.as_str().into()) .await .unwrap() .unwrap(); let dependency_ref: crank_artifacts::ArtifactRef = corrupt_job.preview_payload["dependencies"] [0]["digest"] .as_str() .unwrap() .parse() .unwrap(); let dependency_path = storage_root .join("sha256") .join(&dependency_ref.digest_hex()[..2]) .join(dependency_ref.digest_hex()); std::fs::set_permissions( &dependency_path, std::os::unix::fs::PermissionsExt::from_mode(0o600), ) .unwrap(); std::fs::write(&dependency_path, b"corrupt external dependency").unwrap(); std::fs::set_permissions( &dependency_path, std::os::unix::fs::PermissionsExt::from_mode(0o400), ) .unwrap(); assert!( service .create_openapi_import( &workspace_id, &corrupt_preview.job_id.as_str().into(), OpenApiImportCreatePayload { selected_operation_keys: vec!["GET /items".to_owned()], server_url: None, conflict_mode: "skip".to_owned(), }, ) .await .is_err() ); assert_failed_job_and_detached_sources( ®istry, &workspace_id, &corrupt_preview.job_id.as_str().into(), "import_dependency_verification_failed", ) .await; } #[tokio::test] #[serial] async fn external_materialization_uses_one_wall_clock_timeout_and_cleans_cancelled_sources() { let fetches = Arc::new(AtomicUsize::new(0)); let root_fetches = Arc::clone(&fetches); let child_fetches = Arc::clone(&fetches); let app = Router::new() .route( "/root.yaml", get(move || { let fetches = Arc::clone(&root_fetches); async move { fetches.fetch_add(1, Ordering::SeqCst); "Item: { $ref: './slow-child.yaml#/Item' }" } }), ) .route( "/slow-child.yaml", get(move || { let fetches = Arc::clone(&child_fetches); async move { fetches.fetch_add(1, Ordering::SeqCst); tokio::time::sleep(std::time::Duration::from_secs(1)).await; "Item: { type: object, properties: { id: { type: string } } }" } }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let origin = format!("http://{address}"); let registry = test_registry().await; let outbound_policy = crank_runtime::OutboundHttpPolicy::allowing_hosts(["127.0.0.1"]); let runtime = crank_runtime::community_with_outbound_policy(outbound_policy.clone()).build(); let service = AdminServiceBuilder::new( registry.clone(), test_storage_root("openapi_import_chain_timeout"), test_auth_settings(), test_secret_crypto(), runtime, ) .with_external_reference_import(&crank_config::ExternalReferenceSettings { allowed_url_prefixes: vec![format!("{origin}/")], max_depth: 8, max_documents: 32, max_fetch_bytes: 64 * 1024, fetch_timeout_ms: 200, max_expanded_nodes: 10_000, }) .unwrap() .with_outbound_http_policy(outbound_policy) .build(); let workspace_id = WorkspaceId::new("ws_default"); let started = tokio::time::Instant::now(); let preview = service .preview_openapi_import( &workspace_id, OpenApiUpload { bytes: format!( r#" openapi: 3.1.0 info: {{ title: Timed chain }} servers: [{{ url: https://api.example.test }}] paths: /items: get: operationId: timedItems responses: '200': description: ok content: {{ application/json: {{ schema: {{ $ref: '{origin}/root.yaml#/Item' }} }} }} "# ) .into_bytes(), mime_type: "application/yaml".to_owned(), locale: OpenApiUploadLocale::En, }, ) .await .unwrap(); assert!(started.elapsed() < std::time::Duration::from_millis(800)); assert_eq!(fetches.load(Ordering::SeqCst), 2); let job = registry .get_import_job(&workspace_id, &preview.job_id.as_str().into()) .await .unwrap() .unwrap(); assert_eq!( job.preview_payload["dependencies"] .as_array() .unwrap() .len(), 0 ); for _ in 0..100 { let active_dependencies: i64 = sqlx::query_scalar( "select count(*) from artifact_sources where source_id like 'src_openapi_dep_%' and lifecycle = 'active'", ) .fetch_one(registry.pool()) .await .unwrap(); if active_dependencies == 0 { break; } tokio::time::sleep(std::time::Duration::from_millis(10)).await; } let detached_dependencies: i64 = sqlx::query_scalar( "select count(*) from artifact_sources where source_id like 'src_openapi_dep_%' and lifecycle = 'detached'", ) .fetch_one(registry.pool()) .await .unwrap(); assert_eq!(detached_dependencies, 1); } #[tokio::test] #[serial] async fn expired_job_detaches_primary_and_external_dependency_without_orphan() { let app = Router::new().route( "/schemas.yaml", get(|| async { "Item: { type: object, properties: { id: { type: string } } }" }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let origin = format!("http://{address}"); let registry = test_registry().await; let service = test_service_with_external_references( registry.clone(), test_storage_root("openapi_import_external_expiry"), vec![format!("{origin}/")], ); let workspace_id = WorkspaceId::new("ws_default"); let preview = service .preview_openapi_import( &workspace_id, OpenApiUpload { bytes: format!( r#" openapi: 3.1.0 info: {{ title: Expiring external }} servers: [{{ url: https://api.example.test }}] paths: /items: get: operationId: listItems responses: '200': description: ok content: {{ application/json: {{ schema: {{ $ref: '{origin}/schemas.yaml#/Item' }} }} }} "# ) .into_bytes(), mime_type: "application/yaml".to_owned(), locale: OpenApiUploadLocale::En, }, ) .await .unwrap(); sqlx::query("update import_jobs set expires_at = now() - interval '1 second' where id = $1") .bind(preview.job_id.as_str()) .execute(registry.pool()) .await .unwrap(); let report = registry.cleanup_expired_import_jobs(16).await.unwrap(); assert_eq!(report.deleted_jobs, 1); assert_eq!(report.detached_sources, 2); let active_sources: i64 = sqlx::query_scalar( "select count(*) from artifact_sources where source_id like 'src_openapi_%' and lifecycle = 'active'", ) .fetch_one(registry.pool()) .await .unwrap(); assert_eq!(active_sources, 0); } #[tokio::test] #[serial] async fn exact_v2_contract_replays_persisted_preview_without_v3_or_dependency_reads() { let registry = test_registry().await; let service = test_service( registry.clone(), test_storage_root("openapi_import_v2_compatibility"), test_auth_settings(), test_secret_crypto(), ); let workspace_id = WorkspaceId::new("ws_default"); let preview = service .preview_openapi_import(&workspace_id, openapi_upload()) .await .unwrap(); let job_id: crank_registry::ImportJobId = preview.job_id.as_str().into(); let missing_dependency = serde_json::json!([{ "source_id": "src_openapi_dep_v2_must_not_be_read", "digest": format!("sha256:{}", "0".repeat(64)), "canonical_uri": "https://schemas.example.test/v2.yaml" }]); sqlx::query( "update import_jobs set preview_payload = jsonb_set( jsonb_set( jsonb_set( jsonb_set( preview_payload, '{normalization,normalizer_version}', to_jsonb('normalized-ir-v2'::text) ), '{normalization,projection_version}', to_jsonb('preview-v2'::text) ), '{normalization,ir_fingerprint}', to_jsonb($1::text) ), '{dependency_snapshots}', $2::jsonb ) where id = $3", ) .bind("a".repeat(64)) .bind(missing_dependency) .bind(job_id.as_str()) .execute(registry.pool()) .await .unwrap(); let result = service .create_openapi_import( &workspace_id, &job_id, OpenApiImportCreatePayload { selected_operation_keys: vec!["GET /v2/latest".to_owned()], server_url: Some("https://api.frankfurter.dev".to_owned()), conflict_mode: "skip".to_owned(), }, ) .await .unwrap(); assert_eq!(result.created.len(), 1); let completed = registry .get_import_job(&workspace_id, &job_id) .await .unwrap() .unwrap(); assert_eq!(completed.status, ImportJobStatus::Completed); }