From f85678b68a2f011fea35ecfd21f3a9754bfde7db Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Tue, 29 Sep 2026 19:05:19 +0000 Subject: [PATCH 1/2] graphql: submit draft-backed discoveries Add createDiscover using an owned draft's capture or a readable live capture. Validate authorization, connector readiness, and placement before atomically inserting the discover row and scheduling discovery. Preserve serialized definitions and publication preconditions. Cover staged and live capture selection, byte-preserved endpoint configuration, the update_only policy, rejections, data plane selection, capability checks, and the discovery-to-publication workflow. --- ...1adf5b0be3f156ae09c944b32287155e6bfa8.json | 56 +++ ...9da15907549455162e92fda32dc017819ae24.json | 85 ++++ ...8898bd162a4db410131ddd946c2f47bb8ec13.json | 42 ++ ...a014b655e943518f735d3b3d0c5995ffa08c5.json | 29 ++ ...4ed5e56a8d66d7bd3e826324248ecea7bfdc.json} | 4 +- crates/agent/src/integration_tests/harness.rs | 4 +- ...s__graphql_discover_connector_failure.snap | 14 + ...aphql_discover_error_cleared_by_retry.snap | 10 + .../src/integration_tests/user_discovers.rs | 155 ++++++ crates/control-plane-api/src/discovers/mod.rs | 190 +++++++ .../src/server/public/graphql/discovers.rs | 103 ++-- ...s__test__discover_invalid_submissions.snap | 38 ++ ...test__discover_storage_mapping_planes.snap | 36 ++ ...rs__test__discover_token_capabilities.snap | 28 ++ ...ive_capture_submission_and_rejections.snap | 10 + ...ers__test__staged_discover_submission.snap | 18 + .../server/public/graphql/discovers/test.rs | 462 ++++++++++++++++++ .../src/server/public/graphql/mod.rs | 1 + crates/flow-client/control-plane-api.graphql | 10 + 19 files changed, 1261 insertions(+), 34 deletions(-) create mode 100644 .sqlx/query-19374f4457465c28356551df76c1adf5b0be3f156ae09c944b32287155e6bfa8.json create mode 100644 .sqlx/query-8dc9722aaeae926a290bef3d7f09da15907549455162e92fda32dc017819ae24.json create mode 100644 .sqlx/query-8eb20e235dd4ee616cdd0c9d2718898bd162a4db410131ddd946c2f47bb8ec13.json create mode 100644 .sqlx/query-cd2369734b5f0c5de149a3dbf75a014b655e943518f735d3b3d0c5995ffa08c5.json rename .sqlx/{query-5d53ab4cf24739fbc363d697cad9531f7e35b13e7cb35b38579efcd69d0eba07.json => query-e62d1a2bf12842fcd261a90a24434ed5e56a8d66d7bd3e826324248ecea7bfdc.json} (77%) create mode 100644 crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_connector_failure.snap create mode 100644 crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_error_cleared_by_retry.snap create mode 100644 crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_invalid_submissions.snap create mode 100644 crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_storage_mapping_planes.snap create mode 100644 crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_token_capabilities.snap create mode 100644 crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__live_capture_submission_and_rejections.snap create mode 100644 crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__staged_discover_submission.snap diff --git a/.sqlx/query-19374f4457465c28356551df76c1adf5b0be3f156ae09c944b32287155e6bfa8.json b/.sqlx/query-19374f4457465c28356551df76c1adf5b0be3f156ae09c944b32287155e6bfa8.json new file mode 100644 index 00000000000..d7e9829c85a --- /dev/null +++ b/.sqlx/query-19374f4457465c28356551df76c1adf5b0be3f156ae09c944b32287155e6bfa8.json @@ -0,0 +1,56 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT ds.catalog_name IS NOT NULL AS \"staged!\",\n ds.spec_type AS \"spec_type?: models::CatalogType\",\n ds.spec::text AS spec\n FROM drafts d\n LEFT JOIN draft_specs ds ON ds.draft_id = d.id AND ds.catalog_name = $3\n WHERE d.id = $1 AND d.user_id = $2\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "staged!", + "type_info": "Bool", + "origin": "Expression" + }, + { + "ordinal": 1, + "name": "spec_type?: models::CatalogType", + "type_info": { + "Custom": { + "name": "catalog_spec_type", + "kind": { + "Enum": [ + "capture", + "collection", + "materialization", + "test" + ] + } + } + }, + "origin": { + "Table": { + "table": "draft_specs", + "name": "spec_type" + } + } + }, + { + "ordinal": 2, + "name": "spec", + "type_info": "Text", + "origin": "Expression" + } + ], + "parameters": { + "Left": [ + "Macaddr8", + "Uuid", + "Text" + ] + }, + "nullable": [ + null, + true, + null + ] + }, + "hash": "19374f4457465c28356551df76c1adf5b0be3f156ae09c944b32287155e6bfa8" +} diff --git a/.sqlx/query-8dc9722aaeae926a290bef3d7f09da15907549455162e92fda32dc017819ae24.json b/.sqlx/query-8dc9722aaeae926a290bef3d7f09da15907549455162e92fda32dc017819ae24.json new file mode 100644 index 00000000000..842085a4e62 --- /dev/null +++ b/.sqlx/query-8dc9722aaeae926a290bef3d7f09da15907549455162e92fda32dc017819ae24.json @@ -0,0 +1,85 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO discovers (\n draft_id, capture_name, connector_tag_id, endpoint_config,\n update_only, data_plane_name\n ) VALUES ($1, $2, $3, $4, $5, $6)\n RETURNING id AS \"id!: models::Id\", created_at, updated_at\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: models::Id", + "type_info": "Macaddr8", + "origin": { + "Table": { + "table": "discovers", + "name": "id" + } + } + }, + { + "ordinal": 1, + "name": "created_at", + "type_info": "Timestamptz", + "origin": { + "Table": { + "table": "discovers", + "name": "created_at" + } + } + }, + { + "ordinal": 2, + "name": "updated_at", + "type_info": "Timestamptz", + "origin": { + "Table": { + "table": "discovers", + "name": "updated_at" + } + } + } + ], + "parameters": { + "Left": [ + { + "Custom": { + "name": "flowid", + "kind": { + "Domain": "Macaddr8" + } + } + }, + { + "Custom": { + "name": "catalog_name", + "kind": { + "Domain": "Text" + } + } + }, + { + "Custom": { + "name": "flowid", + "kind": { + "Domain": "Macaddr8" + } + } + }, + { + "Custom": { + "name": "json_obj", + "kind": { + "Domain": "Json" + } + } + }, + "Bool", + "Text" + ] + }, + "nullable": [ + false, + false, + false + ] + }, + "hash": "8dc9722aaeae926a290bef3d7f09da15907549455162e92fda32dc017819ae24" +} diff --git a/.sqlx/query-8eb20e235dd4ee616cdd0c9d2718898bd162a4db410131ddd946c2f47bb8ec13.json b/.sqlx/query-8eb20e235dd4ee616cdd0c9d2718898bd162a4db410131ddd946c2f47bb8ec13.json new file mode 100644 index 00000000000..c5ed3c4e2df --- /dev/null +++ b/.sqlx/query-8eb20e235dd4ee616cdd0c9d2718898bd162a4db410131ddd946c2f47bb8ec13.json @@ -0,0 +1,42 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT ls.spec::text AS spec, ls.spec_type::text AS spec_type,\n ls.data_plane_id AS \"data_plane_id: models::Id\"\n FROM live_specs ls\n WHERE ls.catalog_name = $1\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "spec", + "type_info": "Text", + "origin": "Expression" + }, + { + "ordinal": 1, + "name": "spec_type", + "type_info": "Text", + "origin": "Expression" + }, + { + "ordinal": 2, + "name": "data_plane_id: models::Id", + "type_info": "Macaddr8", + "origin": { + "Table": { + "table": "live_specs", + "name": "data_plane_id" + } + } + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + null, + null, + false + ] + }, + "hash": "8eb20e235dd4ee616cdd0c9d2718898bd162a4db410131ddd946c2f47bb8ec13" +} diff --git a/.sqlx/query-cd2369734b5f0c5de149a3dbf75a014b655e943518f735d3b3d0c5995ffa08c5.json b/.sqlx/query-cd2369734b5f0c5de149a3dbf75a014b655e943518f735d3b3d0c5995ffa08c5.json new file mode 100644 index 00000000000..2780913648a --- /dev/null +++ b/.sqlx/query-cd2369734b5f0c5de149a3dbf75a014b655e943518f735d3b3d0c5995ffa08c5.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT ct.id AS \"id!: models::Id\" FROM connector_tags ct\n JOIN connectors c ON c.id = ct.connector_id\n WHERE c.image_name = $1 AND ct.image_tag = $2\n AND ct.protocol = 'capture' AND ct.job_status->>'type' = 'success'\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: models::Id", + "type_info": "Macaddr8", + "origin": { + "Table": { + "table": "connector_tags", + "name": "id" + } + } + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "cd2369734b5f0c5de149a3dbf75a014b655e943518f735d3b3d0c5995ffa08c5" +} diff --git a/.sqlx/query-5d53ab4cf24739fbc363d697cad9531f7e35b13e7cb35b38579efcd69d0eba07.json b/.sqlx/query-e62d1a2bf12842fcd261a90a24434ed5e56a8d66d7bd3e826324248ecea7bfdc.json similarity index 77% rename from .sqlx/query-5d53ab4cf24739fbc363d697cad9531f7e35b13e7cb35b38579efcd69d0eba07.json rename to .sqlx/query-e62d1a2bf12842fcd261a90a24434ed5e56a8d66d7bd3e826324248ecea7bfdc.json index 666edf32557..ced90fee0d7 100644 --- a/.sqlx/query-5d53ab4cf24739fbc363d697cad9531f7e35b13e7cb35b38579efcd69d0eba07.json +++ b/.sqlx/query-e62d1a2bf12842fcd261a90a24434ed5e56a8d66d7bd3e826324248ecea7bfdc.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT di.id AS \"id!: models::Id\", di.draft_id AS \"draft_id!: models::Id\",\n di.capture_name AS \"capture_name!: models::Name\", di.data_plane_name,\n di.job_status AS \"status!: sqlx::types::Json\",\n di.created_at, di.updated_at\n FROM discovers di\n JOIN drafts d ON d.id = di.draft_id\n WHERE di.id = $1 AND d.user_id = $2\n ", + "query": "\n SELECT di.id AS \"id!: models::Id\", di.draft_id AS \"draft_id!: models::Id\",\n di.capture_name AS \"capture_name!: models::Name\", di.data_plane_name,\n di.job_status AS \"status!: sqlx::types::Json\",\n di.created_at, di.updated_at\n FROM discovers di\n JOIN drafts d ON d.id = di.draft_id\n WHERE di.id = $1 AND d.user_id = $2\n ", "describe": { "columns": [ { @@ -97,5 +97,5 @@ false ] }, - "hash": "5d53ab4cf24739fbc363d697cad9531f7e35b13e7cb35b38579efcd69d0eba07" + "hash": "e62d1a2bf12842fcd261a90a24434ed5e56a8d66d7bd3e826324248ecea7bfdc" } diff --git a/crates/agent/src/integration_tests/harness.rs b/crates/agent/src/integration_tests/harness.rs index a72cec1a6f3..9dba9f93add 100644 --- a/crates/agent/src/integration_tests/harness.rs +++ b/crates/agent/src/integration_tests/harness.rs @@ -414,9 +414,11 @@ impl TestHarness { self.add_data_plane( "ops/dp/public/test", "test.dp.estuary-data.com", - vec!["secret-key".to_string()], + vec!["dGVzdA==".to_string()], ) .await; + sqlx::query("UPDATE data_planes SET hmac_keys = ARRAY['dGVzdA=='] WHERE data_plane_name = 'ops/dp/public/test'") + .execute(&self.pool).await.unwrap(); } /// Ideally, we'd get a whole separate database for each integration test diff --git a/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_connector_failure.snap b/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_connector_failure.snap new file mode 100644 index 00000000000..7ff0ee9a788 --- /dev/null +++ b/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_connector_failure.snap @@ -0,0 +1,14 @@ +--- +source: crates/agent/src/integration_tests/user_discovers.rs +expression: failed_lookup +--- +{ + "discover": { + "errors": [ + { + "detail": "connector unavailable" + } + ], + "status": "DISCOVER_FAILED" + } +} diff --git a/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_error_cleared_by_retry.snap b/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_error_cleared_by_retry.snap new file mode 100644 index 00000000000..4a5288f8a66 --- /dev/null +++ b/crates/agent/src/integration_tests/snapshots/agent__integration_tests__user_discovers__graphql_discover_error_cleared_by_retry.snap @@ -0,0 +1,10 @@ +--- +source: crates/agent/src/integration_tests/user_discovers.rs +expression: earlier +--- +{ + "discover": { + "errors": [], + "status": "DISCOVER_FAILED" + } +} diff --git a/crates/agent/src/integration_tests/user_discovers.rs b/crates/agent/src/integration_tests/user_discovers.rs index 40821705c33..74e3f23ca16 100644 --- a/crates/agent/src/integration_tests/user_discovers.rs +++ b/crates/agent/src/integration_tests/user_discovers.rs @@ -5,6 +5,154 @@ use crate::{ }; use proto_flow::capture::response::{Discovered, discovered::Binding}; +const GRAPHQL_CAPTURE: &str = "acmeCo/capture"; +const GRAPHQL_LOOKUP: &str = "query($id: Id!) { discover(id:$id) { status errors { detail } } }"; + +async fn submit_graphql_discover( + harness: &mut TestHarness, + user: uuid::Uuid, + draft: models::Id, +) -> models::Id { + // Use the harness's publication plane regardless of persistent fixture ordering. + let submitted: serde_json::Value = harness + .execute_graphql_query( + user, + "mutation($draft: Id!, $name: Name!, $plane: String!) { createDiscover(draftId:$draft, captureName:$name, dataPlane:$plane) { id status } }", + &serde_json::json!({"draft": draft, "name": GRAPHQL_CAPTURE, "plane": "ops/dp/public/test"}), + ) + .await + .unwrap(); + serde_json::from_value(submitted["createDiscover"]["id"].clone()).unwrap() +} + +fn graphql_discovered(names: &[&str]) -> Discovered { + Discovered { + bindings: names + .iter() + .map(|name| Binding { + recommended_name: (*name).to_owned(), + document_schema_json: document_schema(1), + resource_config_json: serde_json::json!({"id": name}).to_string().into(), + key: vec!["/id".to_owned()], + ..Default::default() + }) + .collect(), + } +} + +/// Discover a new capture through GraphQL and publish it with one binding, +/// returning the publication's id. +async fn publish_graphql_capture(harness: &mut TestHarness, user: uuid::Uuid) -> models::Id { + let draft = harness + .create_draft( + user, + "initial capture fixture", + draft_catalog(serde_json::json!({ + "captures": { GRAPHQL_CAPTURE: { + "endpoint": {"connector": {"image": "source/test:test", "config": {}}}, + "bindings": [] + }} + })), + ) + .await; + let id = submit_graphql_discover(harness, user, draft).await; + harness.connectors.mock_discover( + GRAPHQL_CAPTURE, + Ok((spec_fixture(), graphql_discovered(&["widgets"]))), + ); + let discovered = harness.run_queued_discover(id).await; + assert!( + discovered.job_status.is_success(), + "{:?}", + discovered.errors + ); + let published = harness + .create_user_publication(user, draft, "publish capture fixture") + .await; + assert!(published.status.is_success(), "{published:?}"); + published.pub_id.unwrap() +} + +#[tokio::test] +async fn test_graphql_rediscover_connector_failure_and_retry() { + let mut harness = + TestHarness::init("test_graphql_rediscover_connector_failure_and_retry").await; + let user = harness.setup_tenant("acmeCo").await; + let pub_id = publish_graphql_capture(&mut harness, user).await; + let draft = harness + .create_draft(user, "rediscover live capture", Default::default()) + .await; + + // A failed discover of the live capture doesn't add it to the draft. + let failed_id = submit_graphql_discover(&mut harness, user, draft).await; + harness + .connectors + .mock_discover(GRAPHQL_CAPTURE, Err("connector unavailable".to_owned())); + let failed = harness.run_queued_discover(failed_id).await; + assert_eq!(failed.draft.spec_count(), 0); + let failed_lookup: serde_json::Value = harness + .execute_graphql_query(user, GRAPHQL_LOOKUP, &serde_json::json!({"id": failed_id})) + .await + .unwrap(); + insta::assert_json_snapshot!("graphql_discover_connector_failure", failed_lookup); + + // A retry merges new bindings into the live capture, expects its last + // publication, and replaces the draft errors that the failure recorded. + let retry_id = submit_graphql_discover(&mut harness, user, draft).await; + harness.connectors.mock_discover( + GRAPHQL_CAPTURE, + Ok((spec_fixture(), graphql_discovered(&["widgets", "gadgets"]))), + ); + let merged = harness.run_queued_discover(retry_id).await; + assert!(merged.job_status.is_success(), "{:?}", merged.errors); + assert_eq!(merged.draft.collections.len(), 2); + assert_eq!( + merged.draft.captures[0] + .model + .as_ref() + .unwrap() + .bindings + .len(), + 2 + ); + assert_eq!(merged.draft.captures[0].expect_pub_id, Some(pub_id)); + let earlier: serde_json::Value = harness + .execute_graphql_query(user, GRAPHQL_LOOKUP, &serde_json::json!({"id": failed_id})) + .await + .unwrap(); + insta::assert_json_snapshot!("graphql_discover_error_cleared_by_retry", earlier); +} + +#[tokio::test] +async fn test_graphql_discover_invalid_draft_is_atomic() { + let mut harness = TestHarness::init("test_graphql_discover_invalid_draft_is_atomic").await; + let user = harness.setup_tenant("acmeCo").await; + publish_graphql_capture(&mut harness, user).await; + let draft = harness + .create_draft( + user, + "rediscover with malformed collection", + Default::default(), + ) + .await; + let id = submit_graphql_discover(&mut harness, user, draft).await; + sqlx::query("INSERT INTO draft_specs (draft_id, catalog_name, spec_type, spec) VALUES ($1, 'acmeCo/malformed', 'collection', '{}'::json)") + .bind(draft).execute(&harness.pool).await.unwrap(); + harness.connectors.mock_discover( + GRAPHQL_CAPTURE, + Ok((spec_fixture(), graphql_discovered(&["widgets", "gadgets"]))), + ); + let failed = harness.run_queued_discover(id).await; + assert_eq!( + failed.job_status, + crate::discovers::JobStatus::DiscoverFailed + ); + assert_eq!(failed.errors[0].0, "flow://collection/acmeCo/malformed"); + // The reloaded draft holds only the malformed entry, which doesn't parse: + // the discover committed neither the capture nor its collections. + assert_eq!(failed.draft.spec_count(), 0); +} + #[tokio::test] async fn test_discover_rejects_non_discovered_responses() { let mut harness = TestHarness::init("test_discover_rejects_non_discovered_responses").await; @@ -618,6 +766,13 @@ async fn test_discover_no_data_plane() { Vec::new(), ) .await; + // Data-plane fixtures survive harness resets, so restore this case on every setup. + sqlx::query( + "UPDATE data_planes SET hmac_keys = '{}' WHERE data_plane_name = 'ops/dp/public/keyless'", + ) + .execute(&harness.pool) + .await + .unwrap(); for (case, data_plane_name) in [ ("unauthorized plane", "ops/dp/private/other"), diff --git a/crates/control-plane-api/src/discovers/mod.rs b/crates/control-plane-api/src/discovers/mod.rs index f6454948d84..c9a8d05daa9 100644 --- a/crates/control-plane-api/src/discovers/mod.rs +++ b/crates/control-plane-api/src/discovers/mod.rs @@ -12,6 +12,196 @@ use std::collections::HashSet; // Re-export key types and functions that executors will need pub use db::{Row, fetch_discover, resolve}; +/// Metadata of a discovery committed with its executor task. +pub struct CreatedDiscover { + pub id: models::Id, + pub data_plane_name: String, + pub created_at: chrono::DateTime, + pub updated_at: chrono::DateTime, +} + +/// Queues discovery using the capture's staged definition, or its live definition +/// when no draft entry exists and `subject` can read it. +/// The executor writes discovery results to the draft. +pub async fn create( + pool: &sqlx::PgPool, + snapshot: &crate::Snapshot, + subject: &models::authz::Subject, + draft_id: models::Id, + capture_name: &str, + requested_data_plane_name: Option<&str>, +) -> anyhow::Result { + // A draft owned by someone else reads as missing. `staged` distinguishes a + // draft without an entry for the capture from an entry staging a deletion. + let draft = sqlx::query!( + r#" + SELECT ds.catalog_name IS NOT NULL AS "staged!", + ds.spec_type AS "spec_type?: models::CatalogType", + ds.spec::text AS spec + FROM drafts d + LEFT JOIN draft_specs ds ON ds.draft_id = d.id AND ds.catalog_name = $3 + WHERE d.id = $1 AND d.user_id = $2 + "#, + draft_id as models::Id, + subject.user_id, + capture_name, + ) + .fetch_optional(pool) + .await? + .ok_or_else(|| anyhow::anyhow!("draft not found"))?; + + // Data plane selection is based on the live capture even when its definition + // cannot be disclosed to this caller or the draft contains edits. + let live_capture = sqlx::query!( + r#" + SELECT ls.spec::text AS spec, ls.spec_type::text AS spec_type, + ls.data_plane_id AS "data_plane_id: models::Id" + FROM live_specs ls + WHERE ls.catalog_name = $1 + "#, + capture_name, + ) + .fetch_optional(pool) + .await? + .filter(|row| row.spec_type.as_deref() == Some("capture") && row.spec.is_some()); + + let model = select_capture_model( + capture_name, + draft + .staged + .then_some((draft.spec_type, draft.spec.as_deref())), + live_capture.as_ref().and_then(|row| row.spec.as_deref()), + snapshot, + subject, + )?; + let (connector_config, image_name, image_tag) = extract_discovery_endpoint(&model)?; + + let connector_tag_id = sqlx::query_scalar!( + r#" + SELECT ct.id AS "id!: models::Id" FROM connector_tags ct + JOIN connectors c ON c.id = ct.connector_id + WHERE c.image_name = $1 AND ct.image_tag = $2 + AND ct.protocol = 'capture' AND ct.job_status->>'type' = 'success' + "#, + image_name, + image_tag, + ) + .fetch_optional(pool) + .await? + .ok_or_else(|| anyhow::anyhow!("capture connector tag is not ready"))?; + + let data_plane_error = || { + snapshot.request_refresh(); + anyhow::anyhow!("data plane not found or unauthorized") + }; + let data_plane_name = if let Some(live) = live_capture { + let current_data_plane_name = &snapshot + .data_plane_by_id(live.data_plane_id) + .ok_or_else(data_plane_error)? + .data_plane_name; + if requested_data_plane_name.is_some_and(|selected| selected != current_data_plane_name) { + anyhow::bail!("data plane differs from the live capture"); + } + current_data_plane_name.clone() + } else { + snapshot + .storage_mapping_for(capture_name) + .and_then(|mapping| { + let (prefix, data_planes) = + mapping.ok_or_else(|| anyhow::anyhow!("no storage mapping for capture"))?; + crate::storage_mappings::select_data_plane( + prefix.as_str(), + data_planes, + requested_data_plane_name, + ) + }) + .inspect_err(|_| snapshot.request_refresh())? + .to_owned() + }; + snapshot + .is_user_authorized(subject, &data_plane_name, models::Capability::Read) + .then(|| snapshot.data_plane_by_catalog_name(&data_plane_name)) + .flatten() + .filter(|data_plane| data_plane.connector_route().is_ok()) + .ok_or_else(data_plane_error)?; + + let update_only = model + .auto_discover + .as_ref() + .is_some_and(|policy| !policy.add_new_bindings); + // The `create_discover_task` trigger schedules the executor within this statement. + let row = sqlx::query!( + r#" + INSERT INTO discovers ( + draft_id, capture_name, connector_tag_id, endpoint_config, + update_only, data_plane_name + ) VALUES ($1, $2, $3, $4, $5, $6) + RETURNING id AS "id!: models::Id", created_at, updated_at + "#, + draft_id as models::Id, + capture_name, + connector_tag_id as models::Id, + crate::TextJson(connector_config.config.clone()) as crate::TextJson, + update_only, + data_plane_name, + ) + .fetch_one(pool) + .await?; + + Ok(CreatedDiscover { + id: row.id, + data_plane_name, + created_at: row.created_at, + updated_at: row.updated_at, + }) +} + +/// `staged` is the type and model of the draft's entry for the capture, if it has one. +fn select_capture_model( + capture_name: &str, + staged: Option<(Option, Option<&str>)>, + live_spec: Option<&str>, + snapshot: &crate::Snapshot, + subject: &models::authz::Subject, +) -> anyhow::Result { + let spec = if let Some((spec_type, spec)) = staged { + if spec_type != Some(models::CatalogType::Capture) { + anyhow::bail!("draft entry is not a capture"); + } + spec.ok_or_else(|| anyhow::anyhow!("draft entry is a deletion"))? + } else { + let spec = live_spec.ok_or_else(|| anyhow::anyhow!("capture not found"))?; + if !snapshot.is_user_authorized( + subject, + capture_name, + models::authz::Capability::CatalogRead, + ) { + anyhow::bail!("capture not found"); + } + spec + }; + serde_json::from_str(spec).map_err(|_| anyhow::anyhow!("invalid capture model")) +} + +fn extract_discovery_endpoint( + model: &models::CaptureDef, +) -> anyhow::Result<(&models::ConnectorConfig, String, String)> { + if model.delete { + anyhow::bail!("capture is staged for deletion"); + } + let models::CaptureEndpoint::Connector(connector_config) = &model.endpoint else { + anyhow::bail!("capture requires a connector endpoint"); + }; + let (image_name, image_tag) = models::split_image_tag(&connector_config.image); + if image_name.is_empty() || image_tag.is_empty() { + anyhow::bail!("capture requires a tagged connector image"); + } + if !serde_json::from_str::(connector_config.config.get())?.is_object() { + anyhow::bail!("endpoint configuration must be an inline JSON object"); + } + Ok((connector_config, image_name, image_tag)) +} + /// Represents the desire to discover an endpoint. The discovered bindings will be merged with /// those in the `base_model`. pub struct Discover<'a> { diff --git a/crates/control-plane-api/src/server/public/graphql/discovers.rs b/crates/control-plane-api/src/server/public/graphql/discovers.rs index d6e10a9af91..07626ddd6b0 100644 --- a/crates/control-plane-api/src/server/public/graphql/discovers.rs +++ b/crates/control-plane-api/src/server/public/graphql/discovers.rs @@ -121,40 +121,81 @@ impl DiscoversQuery { id: models::Id, ) -> async_graphql::Result> { let env = ctx.data::()?; - fetch_discover(id, env.claims()?.subject().user_id, &env.pg_pool).await + let user_id = env.claims()?.subject().user_id; + let row = sqlx::query!( + r#" + SELECT di.id AS "id!: models::Id", di.draft_id AS "draft_id!: models::Id", + di.capture_name AS "capture_name!: models::Name", di.data_plane_name, + di.job_status AS "status!: sqlx::types::Json", + di.created_at, di.updated_at + FROM discovers di + JOIN drafts d ON d.id = di.draft_id + WHERE di.id = $1 AND d.user_id = $2 + "#, + id as models::Id, + user_id, + ) + .fetch_optional(&env.pg_pool) + .await?; + + Ok(row.map(|row| Discover { + id: row.id, + draft_id: row.draft_id, + capture_name: row.capture_name, + data_plane_name: row.data_plane_name, + status: row.status.0, + created_at: row.created_at, + updated_at: row.updated_at, + })) } } -async fn fetch_discover( - id: models::Id, - user_id: uuid::Uuid, - pool: &sqlx::PgPool, -) -> async_graphql::Result> { - let row = sqlx::query!( - r#" - SELECT di.id AS "id!: models::Id", di.draft_id AS "draft_id!: models::Id", - di.capture_name AS "capture_name!: models::Name", di.data_plane_name, - di.job_status AS "status!: sqlx::types::Json", - di.created_at, di.updated_at - FROM discovers di - JOIN drafts d ON d.id = di.draft_id - WHERE di.id = $1 AND d.user_id = $2 - "#, - id as models::Id, - user_id, - ) - .fetch_optional(pool) - .await?; - - Ok(row.map(|row| Discover { - id: row.id, - draft_id: row.draft_id, - capture_name: row.capture_name, - data_plane_name: row.data_plane_name, - status: row.status.0, - created_at: row.created_at, - updated_at: row.updated_at, - })) +#[derive(Debug, Default)] +pub struct DiscoversMutation; + +#[async_graphql::Object] +impl DiscoversMutation { + /// Queue discovery for a capture using its staged or live definition. + /// Discovery updates the given draft with its results. + async fn create_discover( + &self, + ctx: &async_graphql::Context<'_>, + draft_id: models::Id, + capture_name: models::Name, + #[graphql( + desc = "Optional data plane name for discovery. Selected automatically when omitted." + )] + data_plane: Option, + ) -> async_graphql::Result { + let env = ctx.data::()?; + let subject = env.claims()?.subject(); + + super::verify_authorization( + env, + capture_name.as_str(), + models::authz::Capability::SpecEdit, + ) + .await?; + + let row = crate::discovers::create( + &env.pg_pool, + env.snapshot(), + &subject, + draft_id, + capture_name.as_str(), + data_plane.as_deref(), + ) + .await?; + Ok(Discover { + id: row.id, + draft_id, + capture_name, + data_plane_name: row.data_plane_name, + status: models::discovers::JobStatus::Queued, + created_at: row.created_at, + updated_at: row.updated_at, + }) + } } #[cfg(test)] diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_invalid_submissions.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_invalid_submissions.snap new file mode 100644 index 00000000000..35823bfc9a3 --- /dev/null +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_invalid_submissions.snap @@ -0,0 +1,38 @@ +--- +source: crates/control-plane-api/src/server/public/graphql/discovers/test.rs +expression: rejections +--- +[ + { + "capture": "aliceCo/deleted", + "message": "draft entry is a deletion" + }, + { + "capture": "aliceCo/other-type", + "message": "draft entry is not a capture" + }, + { + "capture": "aliceCo/malformed", + "message": "invalid capture model" + }, + { + "capture": "aliceCo/delete-flag", + "message": "capture is staged for deletion" + }, + { + "capture": "aliceCo/local", + "message": "capture requires a connector endpoint" + }, + { + "capture": "aliceCo/untagged", + "message": "capture requires a tagged connector image" + }, + { + "capture": "aliceCo/file-config", + "message": "endpoint configuration must be an inline JSON object" + }, + { + "capture": "aliceCo/failed-tag", + "message": "capture connector tag is not ready" + } +] diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_storage_mapping_planes.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_storage_mapping_planes.snap new file mode 100644 index 00000000000..a92276b8e4b --- /dev/null +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_storage_mapping_planes.snap @@ -0,0 +1,36 @@ +--- +source: crates/control-plane-api/src/server/public/graphql/discovers/test.rs +expression: outcomes +--- +[ + { + "capture": "aliceCo/parent-default", + "dataPlane": "ops/dp/public/aws-us-west-2-c1", + "error": null + }, + { + "capture": "aliceCo/private/nested-default", + "dataPlane": "ops/dp/public/gcp-us-central1-c2", + "error": null + }, + { + "capture": "aliceCo/private/nested-explicit", + "dataPlane": "ops/dp/public/gcp-us-central1-c2", + "error": null + }, + { + "capture": "aliceCo/private/parent-plane", + "dataPlane": null, + "error": "storage mapping aliceCo/private/ doesn't permit data plane ops/dp/public/aws-us-west-2-c1" + }, + { + "capture": "aliceCo/planeless/capture", + "dataPlane": null, + "error": "storage mapping aliceCo/planeless/ is missing associated data planes" + }, + { + "capture": "carolCo/unmapped", + "dataPlane": null, + "error": "no storage mapping for capture" + } +] diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_token_capabilities.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_token_capabilities.snap new file mode 100644 index 00000000000..b2020f55944 --- /dev/null +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_token_capabilities.snap @@ -0,0 +1,28 @@ +--- +source: crates/control-plane-api/src/server/public/graphql/discovers/test.rs +expression: outcomes +--- +[ + { + "error": "PermissionDenied: user is not authorized to access prefix or name 'aliceCo/in/capture-foo' with required capability SpecEdit", + "mask": [ + "viewer" + ], + "status": null + }, + { + "error": "data plane not found or unauthorized", + "mask": [ + "editor" + ], + "status": null + }, + { + "error": null, + "mask": [ + "editor", + "viewer" + ], + "status": "QUEUED" + } +] diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__live_capture_submission_and_rejections.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__live_capture_submission_and_rejections.snap new file mode 100644 index 00000000000..98c9f943020 --- /dev/null +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__live_capture_submission_and_rejections.snap @@ -0,0 +1,10 @@ +--- +source: crates/control-plane-api/src/server/public/graphql/discovers/test.rs +expression: "serde_json::json!({\n \"livePlane\": live[\"data\"][\"createDiscover\"][\"dataPlaneName\"],\n \"stagedPlane\": staged[\"data\"][\"createDiscover\"][\"dataPlaneName\"],\n \"missing\": missing[\"errors\"][0][\"message\"], \"wrongPlane\":\n wrong_plane[\"errors\"][0][\"message\"],\n})" +--- +{ + "livePlane": "ops/dp/public/aws-us-west-2-c1", + "missing": "capture not found", + "stagedPlane": "ops/dp/public/aws-us-west-2-c1", + "wrongPlane": "data plane differs from the live capture" +} diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__staged_discover_submission.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__staged_discover_submission.snap new file mode 100644 index 00000000000..f49c883e86c --- /dev/null +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__staged_discover_submission.snap @@ -0,0 +1,18 @@ +--- +source: crates/control-plane-api/src/server/public/graphql/discovers/test.rs +expression: response +--- +{ + "data": { + "createDiscover": { + "captureName": "aliceCo/new-capture", + "createdAt": "[ts]", + "dataPlaneName": "ops/dp/public/aws-us-west-2-c1", + "draftId": "[draft-id]", + "errors": [], + "id": "[id]", + "status": "QUEUED", + "updatedAt": "[ts]" + } + } +} diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/test.rs b/crates/control-plane-api/src/server/public/graphql/discovers/test.rs index 08d10f37e96..99bd6444f00 100644 --- a/crates/control-plane-api/src/server/public/graphql/discovers/test.rs +++ b/crates/control-plane-api/src/server/public/graphql/discovers/test.rs @@ -2,6 +2,13 @@ use crate::test_server; const ALICE: uuid::Uuid = uuid::Uuid::from_bytes([0x11; 16]); const BOB: uuid::Uuid = uuid::Uuid::from_bytes([0x22; 16]); +const CREATE: &str = r#" +mutation ($draftId: Id!, $captureName: Name!, $dataPlane: String) { + createDiscover(draftId: $draftId, captureName: $captureName, dataPlane: $dataPlane) { + id draftId captureName dataPlaneName status createdAt updatedAt + errors { catalogName scope detail } + } +}"#; const LOOKUP: &str = r#" query ($id: Id!, $after: String, $first: Int) { discover(id: $id) { @@ -13,6 +20,8 @@ query ($id: Id!, $after: String, $first: Int) { } } }"#; +const MODEL: &str = + r#"{"endpoint":{"connector":{"image":"source/test:test","config":{}}},"bindings":[]}"#; /// Insert a draft owned by Alice, then start a server over a Snapshot read /// from `pool`. The Snapshot is fixed, so fixtures it carries (grants, storage @@ -75,6 +84,46 @@ async fn insert_discover(pool: &sqlx::PgPool, draft_id: models::Id, name: &str) .unwrap() } +/// `MODEL` with `fields` replacing its top-level fields. +fn model_with(fields: serde_json::Value) -> String { + let serde_json::Value::Object(fields) = fields else { + panic!("fields must be an object"); + }; + let mut model: serde_json::Value = serde_json::from_str(MODEL).unwrap(); + model.as_object_mut().unwrap().extend(fields); + model.to_string() +} + +async fn stage_capture(pool: &sqlx::PgPool, draft_id: models::Id, name: &str, model: &str) { + sqlx::query( + "INSERT INTO draft_specs (draft_id, catalog_name, spec_type, spec) VALUES ($1, $2, 'capture', $3::json)", + ) + .bind(draft_id) + .bind(name) + .bind(model) + .execute(pool) + .await + .unwrap(); +} + +async fn submit( + server: &test_server::TestServer, + token: Option<&str>, + draft_id: models::Id, + name: &str, + plane: Option<&str>, +) -> serde_json::Value { + server + .graphql( + &serde_json::json!({ + "query": CREATE, + "variables": { "draftId": draft_id, "captureName": name, "dataPlane": plane } + }), + token, + ) + .await +} + async fn lookup( server: &test_server::TestServer, token: Option<&str>, @@ -88,6 +137,31 @@ async fn lookup( .await } +/// The `endpoint_config` and `update_only` which a successful submission queued. +async fn queued(pool: &sqlx::PgPool, response: &serde_json::Value) -> (String, bool) { + let id: models::Id = serde_json::from_value(response["data"]["createDiscover"]["id"].clone()) + .unwrap_or_else(|err| panic!("{err}: {response}")); + let (config, update_only): (String, bool) = + sqlx::query_as("SELECT endpoint_config::text, update_only FROM discovers WHERE id = $1") + .bind(id) + .fetch_one(pool) + .await + .unwrap(); + (config.trim().to_owned(), update_only) +} + +async fn draft_state(pool: &sqlx::PgPool, draft_id: models::Id) -> serde_json::Value { + sqlx::query_scalar(r#" + SELECT jsonb_build_object( + 'draft', (SELECT to_jsonb(d) FROM drafts d WHERE id = $1), + 'specs', (SELECT jsonb_agg(to_jsonb(s) ORDER BY catalog_name) FROM draft_specs s WHERE draft_id = $1), + 'errors', (SELECT jsonb_agg(to_jsonb(e) ORDER BY scope, detail) FROM draft_errors e WHERE draft_id = $1), + 'jobs', (SELECT jsonb_agg(to_jsonb(j) ORDER BY id) FROM discovers j WHERE draft_id = $1), + 'tasks', (SELECT jsonb_agg(to_jsonb(t) ORDER BY task_id) FROM internal.tasks t) + ) + "#).bind(draft_id).fetch_one(pool).await.unwrap() +} + #[sqlx::test( migrations = "../../supabase/migrations", fixtures( @@ -111,6 +185,125 @@ async fn lookup_is_private_to_the_draft_owner(pool: sqlx::PgPool) { }); } +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn staged_submission_and_ownership(pool: sqlx::PgPool) { + let _guard = test_server::init(); + let (server, draft_id, alice, _) = setup(&pool).await; + stage_capture(&pool, draft_id, "aliceCo/new-capture", MODEL).await; + + let response = submit(&server, Some(&alice), draft_id, "aliceCo/new-capture", None).await; + insta::assert_json_snapshot!("staged_discover_submission", response, { + ".data.createDiscover.id" => "[id]", + ".data.createDiscover.draftId" => "[draft-id]", + ".data.createDiscover.createdAt" => "[ts]", + ".data.createDiscover.updatedAt" => "[ts]", + }); + let tasks: i64 = sqlx::query_scalar( + "SELECT count(*) FROM internal.tasks WHERE task_id IN (SELECT id FROM discovers WHERE draft_id = $1)", + ) + .bind(draft_id) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(tasks, 1, "submission schedules the executor"); + + // A draft owned by someone else is indistinguishable from a missing one. + let foreign_draft_id = insert_draft(&pool, BOB).await; + let missing_draft_id = models::Id::new(u64::MAX.to_be_bytes()); + for draft_id in [foreign_draft_id, missing_draft_id] { + let response = submit(&server, Some(&alice), draft_id, "aliceCo/new-capture", None).await; + assert_eq!( + response["errors"][0]["message"], "draft not found", + "{response}" + ); + } +} + +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn submission_selects_staged_or_live_capture(pool: sqlx::PgPool) { + let _guard = test_server::init(); + // Keys are deliberately unsorted: SOPS authenticates an encrypted config by + // walking it in order, so the queued config must preserve its bytes. + const LIVE: &str = r#"{"zebra":"live","alpha":{"zulu":1,"bravo":2}}"#; + const STAGED: &str = r#"{"zebra":"staged","alpha":{"zulu":1,"bravo":2}}"#; + let model = |config: &str| { + format!( + r#"{{"endpoint":{{"connector":{{"image":"source/test:test","config":{config}}}}},"bindings":[]}}"# + ) + }; + sqlx::query( + "UPDATE live_specs SET spec = $1::json WHERE catalog_name = 'aliceCo/in/capture-foo'", + ) + .bind(model(LIVE)) + .execute(&pool) + .await + .unwrap(); + let (server, draft_id, alice, _) = setup(&pool).await; + sqlx::query("INSERT INTO draft_specs (draft_id, catalog_name, spec_type, spec) VALUES ($1, 'aliceCo/unrelated', 'collection', '{}'::json)") + .bind(draft_id).execute(&pool).await.unwrap(); + sqlx::query("INSERT INTO draft_errors (draft_id, scope, detail) VALUES ($1, 'flow://collection/aliceCo/unrelated', 'existing error')") + .bind(draft_id).execute(&pool).await.unwrap(); + + // Without a staged entry, the readable live capture is used, and the draft + // keeps its entries and errors. + let before = draft_state(&pool, draft_id).await; + let live = submit( + &server, + Some(&alice), + draft_id, + "aliceCo/in/capture-foo", + None, + ) + .await; + let after = draft_state(&pool, draft_id).await; + for field in ["draft", "specs", "errors"] { + assert_eq!(after[field], before[field], "{field}"); + } + // A staged entry takes precedence over the live capture. + stage_capture(&pool, draft_id, "aliceCo/in/capture-foo", &model(STAGED)).await; + let staged = submit( + &server, + Some(&alice), + draft_id, + "aliceCo/in/capture-foo", + None, + ) + .await; + assert_eq!(queued(&pool, &live).await.0, LIVE); + assert_eq!(queued(&pool, &staged).await.0, STAGED); + + let missing = submit(&server, Some(&alice), draft_id, "aliceCo/missing", None).await; + let wrong_plane = submit( + &server, + Some(&alice), + draft_id, + "aliceCo/in/capture-foo", + Some("ops/dp/public/gcp-us-central1-c2"), + ) + .await; + insta::assert_json_snapshot!( + "live_capture_submission_and_rejections", + serde_json::json!({ + "livePlane": live["data"]["createDiscover"]["dataPlaneName"], + "stagedPlane": staged["data"]["createDiscover"]["dataPlaneName"], + "missing": missing["errors"][0]["message"], + "wrongPlane": wrong_plane["errors"][0]["message"], + }) + ); +} + #[sqlx::test( migrations = "../../supabase/migrations", fixtures( @@ -201,6 +394,207 @@ async fn logs_paginate_and_arrive_after_completion(pool: sqlx::PgPool) { ); } +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn update_only_follows_auto_discover(pool: sqlx::PgPool) { + let _guard = test_server::init(); + let (server, draft_id, alice, _) = setup(&pool).await; + for (name, fields, update_only) in [ + ("aliceCo/absent-policy", serde_json::json!({}), false), + ( + "aliceCo/null-policy", + serde_json::json!({ "autoDiscover": null }), + false, + ), + ( + "aliceCo/empty-policy", + serde_json::json!({ "autoDiscover": {} }), + true, + ), + ( + "aliceCo/enabled-policy", + serde_json::json!({ "autoDiscover": { "addNewBindings": true } }), + false, + ), + ] { + stage_capture(&pool, draft_id, name, &model_with(fields)).await; + let response = submit(&server, Some(&alice), draft_id, name, None).await; + assert_eq!(queued(&pool, &response).await.1, update_only, "{name}"); + } +} + +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn invalid_submissions_change_nothing(pool: sqlx::PgPool) { + let _guard = test_server::init(); + let (server, draft_id, alice, _) = setup(&pool).await; + let connector = |image: &str, config: serde_json::Value| { + Some(model_with(serde_json::json!({ + "endpoint": { "connector": { "image": image, "config": config } } + }))) + }; + let cases = [ + ("aliceCo/deleted", "capture", None), + ("aliceCo/other-type", "collection", Some("{}".to_owned())), + ("aliceCo/malformed", "capture", Some("{}".to_owned())), + ( + "aliceCo/delete-flag", + "capture", + Some(model_with(serde_json::json!({ "delete": true }))), + ), + ( + "aliceCo/local", + "capture", + Some(model_with(serde_json::json!({ + "endpoint": { "local": { "command": ["true"], "config": {} } } + }))), + ), + ( + "aliceCo/untagged", + "capture", + connector("source/test", serde_json::json!({})), + ), + ( + "aliceCo/file-config", + "capture", + connector("source/test:test", serde_json::json!("config.json")), + ), + ( + "aliceCo/failed-tag", + "capture", + connector("source/multi-tag-test:v2", serde_json::json!({})), + ), + ]; + for (name, spec_type, model) in &cases { + sqlx::query("INSERT INTO draft_specs (draft_id, catalog_name, spec_type, spec) VALUES ($1, $2, $3::catalog_spec_type, $4::json)") + .bind(draft_id) + .bind(name) + .bind(spec_type) + .bind(model) + .execute(&pool) + .await + .unwrap(); + } + + let before = draft_state(&pool, draft_id).await; + let mut rejections = Vec::new(); + for (name, _, _) in cases { + let response = submit(&server, Some(&alice), draft_id, name, None).await; + rejections.push(serde_json::json!({ + "capture": name, + "message": response["errors"][0]["message"], + })); + } + insta::assert_json_snapshot!("discover_invalid_submissions", rejections); + assert_eq!(draft_state(&pool, draft_id).await, before); +} + +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn new_captures_use_their_storage_mapping(pool: sqlx::PgPool) { + let _guard = test_server::init(); + sqlx::query( + "INSERT INTO storage_mappings (catalog_prefix, spec) VALUES ('aliceCo/planeless/', '{}')", + ) + .execute(&pool) + .await + .unwrap(); + // `carolCo/` has no storage mapping at all. + sqlx::query("INSERT INTO user_grants (user_id, object_role, capability) VALUES ($1, 'carolCo/', 'admin')") + .bind(ALICE).execute(&pool).await.unwrap(); + let (server, draft_id, alice, revoke) = setup(&pool).await; + + let mut outcomes = Vec::new(); + for (name, plane) in [ + ("aliceCo/parent-default", None), + ("aliceCo/private/nested-default", None), + ( + "aliceCo/private/nested-explicit", + Some("ops/dp/public/gcp-us-central1-c2"), + ), + // The most specific mapping decides alone, though its parent admits this plane. + ( + "aliceCo/private/parent-plane", + Some("ops/dp/public/aws-us-west-2-c1"), + ), + ("aliceCo/planeless/capture", None), + ("carolCo/unmapped", None), + ] { + stage_capture(&pool, draft_id, name, MODEL).await; + let response = submit(&server, Some(&alice), draft_id, name, plane).await; + outcomes.push(serde_json::json!({ + "capture": name, + "dataPlane": response["data"]["createDiscover"]["dataPlaneName"], + "error": response["errors"][0]["message"], + })); + } + insta::assert_json_snapshot!("discover_storage_mapping_planes", outcomes); + assert!( + revoke.is_cancelled(), + "mapping rejections request a Snapshot refresh" + ); +} + +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn submission_enforces_token_capabilities(pool: sqlx::PgPool) { + let _guard = test_server::init(); + sqlx::query( + "UPDATE live_specs SET spec = $1::json WHERE catalog_name = 'aliceCo/in/capture-foo'", + ) + .bind(MODEL) + .execute(&pool) + .await + .unwrap(); + let (server, draft_id, _, _) = setup(&pool).await; + + // Discovery needs SpecEdit on the capture and legacy `read` of its data + // plane, which the editor bundle alone doesn't convey. + let mut outcomes = Vec::new(); + for mask in [&["viewer"][..], &["editor"], &["editor", "viewer"]] { + let token = server.make_restricted_access_token( + ALICE, + None, + Some(mask.iter().map(|bundle| bundle.to_string()).collect()), + None, + ); + let response = submit( + &server, + Some(&token), + draft_id, + "aliceCo/in/capture-foo", + None, + ) + .await; + outcomes.push(serde_json::json!({ + "mask": mask, + "status": response["data"]["createDiscover"]["status"], + "error": response["errors"][0]["message"], + })); + } + insta::assert_json_snapshot!("discover_token_capabilities", outcomes); +} + #[sqlx::test( migrations = "../../supabase/migrations", fixtures( @@ -256,3 +650,71 @@ async fn log_page_arguments(pool: sqlx::PgPool) { } insta::assert_json_snapshot!("discover_log_arguments", rejected); } + +#[sqlx::test( + migrations = "../../supabase/migrations", + fixtures( + path = "../../../../fixtures", + scripts("data_planes", "alice", "drafts", "connectors", "storage_mappings") + ) +)] +async fn plane_gate_requests_refresh_without_waiting(pool: sqlx::PgPool) { + let _guard = test_server::init(); + sqlx::query( + "UPDATE live_specs SET spec = $1::json WHERE catalog_name = 'aliceCo/in/capture-foo'", + ) + .bind(MODEL) + .execute(&pool) + .await + .unwrap(); + let draft_id = insert_draft(&pool, ALICE).await; + let mut no_access = crate::snapshot::try_fetch(&pool, &mut Default::default()) + .await + .unwrap(); + let mut unlisted = no_access.clone(); + no_access + .role_grants + .retain(|grant| grant.object_role.as_str() != "ops/dp/public/"); + unlisted + .data_planes + .retain(|plane| plane.data_plane_name != "ops/dp/public/aws-us-west-2-c1"); + sqlx::query("UPDATE data_planes SET hmac_keys = ARRAY['invalid-base64%%%'], encrypted_hmac_keys = '{}'::json WHERE data_plane_name = 'ops/dp/public/aws-us-west-2-c1'") + .execute(&pool) + .await + .unwrap(); + let unsigned = crate::snapshot::try_fetch(&pool, &mut Default::default()) + .await + .unwrap(); + + for (case, data) in [ + ("access", no_access), + ("snapshot lookup", unlisted), + ("signing readiness", unsigned), + ] { + // A Snapshot taken before the request starts would make a provisional + // authorization failure await a refresh, which a fixed watch never serves. + let (server, revoke) = + start(&pool, data, tokens::now() - chrono::TimeDelta::minutes(1)).await; + let alice = server.make_access_token(ALICE, Some("alice@example.com")); + let response = tokio::time::timeout( + std::time::Duration::from_secs(5), + submit( + &server, + Some(&alice), + draft_id, + "aliceCo/in/capture-foo", + None, + ), + ) + .await + .expect("plane rejection must not await a new snapshot"); + assert_eq!( + response["errors"][0]["message"], "data plane not found or unauthorized", + "{case}" + ); + assert!( + revoke.is_cancelled(), + "{case} must request a background refresh" + ); + } +} diff --git a/crates/control-plane-api/src/server/public/graphql/mod.rs b/crates/control-plane-api/src/server/public/graphql/mod.rs index 0c56ef0eb24..370bd33785f 100644 --- a/crates/control-plane-api/src/server/public/graphql/mod.rs +++ b/crates/control-plane-api/src/server/public/graphql/mod.rs @@ -154,6 +154,7 @@ pub struct MutationRoot( service_accounts::ServiceAccountsMutation, secrets::SecretsMutation, drafts::DraftsMutation, + discovers::DiscoversMutation, ); pub fn create_schema(alert_config_defaults: models::AlertConfig) -> GraphQLSchema { diff --git a/crates/flow-client/control-plane-api.graphql b/crates/flow-client/control-plane-api.graphql index 8767e2bd366..d277285f0c7 100644 --- a/crates/flow-client/control-plane-api.graphql +++ b/crates/flow-client/control-plane-api.graphql @@ -2181,6 +2181,16 @@ type MutationRoot { Returns the deleted ID. Nothing published is affected. """ deleteDraft(id: Id!): Id! + """ + Queue discovery for a capture using its staged or live definition. + Discovery updates the given draft with its results. + """ + createDiscover( draftId: Id!, captureName: Name!, + """ + Optional data plane name for discovery. Selected automatically when omitted. + """ + dataPlane: String + ): Discover! } """ From b3dbc25ae4cbba20f191efbdf172630ce9103a45 Mon Sep 17 00:00:00 2001 From: Joseph Shearer Date: Tue, 29 Sep 2026 16:14:33 +0000 Subject: [PATCH 2/2] logs: preserve microsecond timestamp order across writer batches Assign each line a timestamp at PostgreSQL's microsecond precision, advancing past the writer's previous timestamp even when its clock moves backward. Retain that timestamp across batches so timestamp cursors can resume within and between batches without skipping lines from this writer. Exercise the real writer through the discover log connection, including page boundaries within a batch, across batches, and late writes. Keep a focused unit test for clock recession. Ordering remains local to one writer's lifetime. --- crates/control-plane-api/src/logs.rs | 42 +++++++- ...__discovers__test__discover_log_pages.snap | 32 +++--- .../server/public/graphql/discovers/test.rs | 98 +++++++++++++++---- 3 files changed, 132 insertions(+), 40 deletions(-) diff --git a/crates/control-plane-api/src/logs.rs b/crates/control-plane-api/src/logs.rs index 33810ce8275..ae6014a2a59 100644 --- a/crates/control-plane-api/src/logs.rs +++ b/crates/control-plane-api/src/logs.rs @@ -71,6 +71,8 @@ pub async fn serve_sink( let mut tokens = Vec::new(); let mut streams = Vec::new(); let mut lines = Vec::new(); + let mut logged_at = Vec::new(); + let mut previous_logged_at = None; let mut held_conn = None; let mut interval = tokio::time::interval(std::time::Duration::from_secs(15)); @@ -118,6 +120,13 @@ pub async fn serve_sink( lines.push(sanitize_null_bytes(line)); } + let mut next = next_logged_at(previous_logged_at, chrono::Utc::now()); + for _ in &lines { + logged_at.push(next); + next += chrono::Duration::microseconds(1); + } + previous_logged_at = logged_at.last().copied(); + if let None = held_conn { held_conn = Some(pg_pool.acquire().await?); debug!("acquired new pg_conn"); @@ -127,13 +136,14 @@ pub async fn serve_sink( // Dispatch the vector of lines to the table. let r = sqlx::query( r#" - insert into internal.log_lines (token, stream, log_line) - select * from unnest($1, $2, $3) + insert into internal.log_lines (token, stream, log_line, logged_at) + select * from unnest($1, $2, $3, $4) "#, ) .bind(&tokens) .bind(&streams) .bind(&lines) + .bind(&logged_at) .execute(held_conn.as_deref_mut().unwrap()) .await?; @@ -142,9 +152,24 @@ pub async fn serve_sink( tokens.clear(); streams.clear(); lines.clear(); + logged_at.clear(); } } +/// PostgreSQL stores timestamps at microsecond precision. Keep the cursor +/// strictly increasing across batches, including when the local clock recedes. +fn next_logged_at( + previous: Option>, + now: chrono::DateTime, +) -> chrono::DateTime { + let now = chrono::DateTime::from_timestamp_micros(now.timestamp_micros()) + .expect("current time is representable"); + previous + .map(|previous| previous + chrono::Duration::microseconds(1)) + .filter(|next| *next > now) + .unwrap_or(now) +} + #[derive(Debug, Clone)] pub struct OpsHandler { tx: Tx, @@ -238,9 +263,20 @@ fn render_ops_log_for_ui(log: &ops::Log) -> String { #[cfg(test)] mod test { - use super::render_ops_log_for_ui; + use super::{next_logged_at, render_ops_log_for_ui}; use proto_flow::ops; + #[test] + fn timestamps_advance_across_batches_and_clock_recession() { + let now = "2026-01-01T00:00:00.123456789Z".parse().unwrap(); + let first = next_logged_at(None, now); + assert_eq!(first.to_rfc3339(), "2026-01-01T00:00:00.123456+00:00"); + let second = next_logged_at(Some(first), now); + assert_eq!(second - first, chrono::Duration::microseconds(1)); + let backwards = next_logged_at(Some(second), now - chrono::Duration::seconds(1)); + assert_eq!(backwards - second, chrono::Duration::microseconds(1)); + } + #[test] fn test_log_rendering() { let fixture = ops::Log { diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_log_pages.snap b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_log_pages.snap index 6e2b53f5aec..69de6c0183b 100644 --- a/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_log_pages.snap +++ b/crates/control-plane-api/src/server/public/graphql/discovers/snapshots/control_plane_api__server__public__graphql__discovers__test__discover_log_pages.snap @@ -7,16 +7,16 @@ expression: "serde_json::json!({\n \"pages\": pages, \"afterCompletion\":\n "logs": { "edges": [ { - "cursor": "2026-01-01T00:00:01+00:00", + "cursor": "[cursor]", "node": { "line": "late line", - "loggedAt": "2026-01-01T00:00:01+00:00", + "loggedAt": "[ts]", "stream": "test" } } ], "pageInfo": { - "endCursor": "2026-01-01T00:00:01+00:00", + "endCursor": "[cursor]", "hasNextPage": false } }, @@ -26,64 +26,64 @@ expression: "serde_json::json!({\n \"pages\": pages, \"afterCompletion\":\n { "edges": [ { - "cursor": "2026-01-01T00:00:00.000001+00:00", + "cursor": "[cursor]", "node": { "line": "first", - "loggedAt": "2026-01-01T00:00:00.000001+00:00", + "loggedAt": "[ts]", "stream": "test" } }, { - "cursor": "2026-01-01T00:00:00.000002+00:00", + "cursor": "[cursor]", "node": { "line": "second", - "loggedAt": "2026-01-01T00:00:00.000002+00:00", + "loggedAt": "[ts]", "stream": "test" } } ], "pageInfo": { - "endCursor": "2026-01-01T00:00:00.000002+00:00", + "endCursor": "[cursor]", "hasNextPage": true } }, { "edges": [ { - "cursor": "2026-01-01T00:00:00.000003+00:00", + "cursor": "[cursor]", "node": { "line": "third", - "loggedAt": "2026-01-01T00:00:00.000003+00:00", + "loggedAt": "[ts]", "stream": "test" } }, { - "cursor": "2026-01-01T00:00:00.000004+00:00", + "cursor": "[cursor]", "node": { "line": "fourth", - "loggedAt": "2026-01-01T00:00:00.000004+00:00", + "loggedAt": "[ts]", "stream": "test" } } ], "pageInfo": { - "endCursor": "2026-01-01T00:00:00.000004+00:00", + "endCursor": "[cursor]", "hasNextPage": true } }, { "edges": [ { - "cursor": "2026-01-01T00:00:00.000005+00:00", + "cursor": "[cursor]", "node": { "line": "fifth", - "loggedAt": "2026-01-01T00:00:00.000005+00:00", + "loggedAt": "[ts]", "stream": "test" } } ], "pageInfo": { - "endCursor": "2026-01-01T00:00:00.000005+00:00", + "endCursor": "[cursor]", "hasNextPage": false } } diff --git a/crates/control-plane-api/src/server/public/graphql/discovers/test.rs b/crates/control-plane-api/src/server/public/graphql/discovers/test.rs index 99bd6444f00..b58e2b35dac 100644 --- a/crates/control-plane-api/src/server/public/graphql/discovers/test.rs +++ b/crates/control-plane-api/src/server/public/graphql/discovers/test.rs @@ -315,22 +315,46 @@ async fn logs_paginate_and_arrive_after_completion(pool: sqlx::PgPool) { let _guard = test_server::init(); let (server, draft_id, alice, _) = setup(&pool).await; let id = insert_discover(&pool, draft_id, "aliceCo/logged").await; - sqlx::query( - r#" - INSERT INTO internal.log_lines (token, stream, log_line, logged_at) - SELECT logs_token, 'test', line, ts::timestamptz FROM discovers, - (VALUES ('first', '2026-01-01T00:00:00.000001Z'), - ('second', '2026-01-01T00:00:00.000002Z'), - ('third', '2026-01-01T00:00:00.000003Z'), - ('fourth', '2026-01-01T00:00:00.000004Z'), - ('fifth', '2026-01-01T00:00:00.000005Z')) AS lines(line, ts) - WHERE id = $1 - "#, + let token: uuid::Uuid = sqlx::query_scalar("SELECT logs_token FROM discovers WHERE id = $1") + .bind(id) + .fetch_one(&pool) + .await + .unwrap(); + let (tx, rx) = tokio::sync::mpsc::channel(8); + // Queue the first batch before starting the writer so a page splits it. + crate::logs::capture_lines( + tx.clone(), + "test".into(), + token, + &b"first\nsecond\nthird"[..], ) - .bind(id) - .execute(&pool) .await .unwrap(); + let writer = tokio::spawn(crate::logs::serve_sink(pool.clone(), rx)); + let wait_pool = &pool; + let wait_for_lines = |expected| async move { + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let count: i64 = + sqlx::query_scalar("SELECT count(*) FROM internal.log_lines WHERE token = $1") + .bind(token) + .fetch_one(wait_pool) + .await + .unwrap(); + if count == expected { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + }; + wait_for_lines(3).await; + crate::logs::capture_lines(tx.clone(), "test".into(), token, &b"fourth\nfifth"[..]) + .await + .unwrap(); + wait_for_lines(5).await; // Lines of another user's discover are never returned. sqlx::query(r#" WITH foreign_draft AS ( @@ -369,19 +393,43 @@ async fn logs_paginate_and_arrive_after_completion(pool: sqlx::PgPool) { .execute(&pool) .await .unwrap(); - sqlx::query( - "INSERT INTO internal.log_lines (token, stream, log_line, logged_at) SELECT logs_token, 'test', 'late line', '2026-01-01T00:00:01Z' FROM discovers WHERE id = $1", - ) - .bind(id) - .execute(&pool) - .await - .unwrap(); + crate::logs::capture_lines(tx.clone(), "test".into(), token, &b"late line"[..]) + .await + .unwrap(); + wait_for_lines(6).await; + drop(tx); + writer.await.unwrap().unwrap(); let later = lookup( &server, Some(&alice), serde_json::json!({ "id": id, "after": after }), ) .await; + + // The writer stamps lines from its own clock, so check their order here and + // redact them from the snapshot. + let logged_at = pages + .iter() + .chain([&later["data"]["discover"]["logs"]]) + .flat_map(|logs| logs["edges"].as_array().unwrap()) + .map(|edge| { + edge["node"]["loggedAt"] + .as_str() + .unwrap() + .parse::>() + .unwrap() + }) + .collect::>(); + assert!(logged_at.windows(2).all(|pair| pair[0] < pair[1])); + assert!( + logged_at + .iter() + .all(|ts| ts.timestamp_subsec_nanos() % 1_000 == 0) + ); + assert_eq!( + logged_at[2] - logged_at[0], + chrono::Duration::microseconds(2) + ); insta::assert_json_snapshot!( "discover_log_pages", serde_json::json!({ @@ -390,7 +438,15 @@ async fn logs_paginate_and_arrive_after_completion(pool: sqlx::PgPool) { "status": later["data"]["discover"]["status"], "logs": later["data"]["discover"]["logs"], }, - }) + }), + { + ".pages[].edges[].cursor" => "[cursor]", + ".pages[].edges[].node.loggedAt" => "[ts]", + ".pages[].pageInfo.endCursor" => "[cursor]", + ".afterCompletion.logs.edges[].cursor" => "[cursor]", + ".afterCompletion.logs.edges[].node.loggedAt" => "[ts]", + ".afterCompletion.logs.pageInfo.endCursor" => "[cursor]", + } ); }