Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions crates/control-plane-api/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ colored_json = { workspace = true }
derivative = { workspace = true }
enumset = { workspace = true }
futures = { workspace = true }
hex = { workspace = true }
humantime = { workspace = true }
humantime-serde = { workspace = true }
itertools = { workspace = true }
Expand All @@ -63,6 +64,7 @@ rustls = { workspace = true }
schemars = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
sha2 = { workspace = true }
sqlx = { workspace = true }
tempfile = { workspace = true }
thiserror = { workspace = true }
Expand All @@ -88,11 +90,9 @@ flow-client-next = { path = "../flow-client-next" }
runtime-local = { path = "../runtime-local" }
service-kit = { path = "../service-kit" }

hex = { workspace = true }
hmac = { workspace = true }
insta = { workspace = true }
md5 = { workspace = true }
sha2 = { workspace = true }
tokio = { workspace = true, features = ["test-util"] }
tracing-subscriber = { workspace = true }

Expand Down
10 changes: 9 additions & 1 deletion crates/control-plane-api/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,15 @@ Service-account API keys are minted by `server/public/graphql/service_accounts.r
Like refresh-token creation, API-key creation rejects access tokens carrying a
capability mask or prefix scope, because the new credential would not preserve
those restrictions. Tenant creation in `server/public/graphql/tenant.rs` also
rejects either restriction before provisioning a tenant.
rejects either restriction before provisioning a tenant. Its `tenant/storage.rs`
helper derives collection and recovery storage mappings. Tenant provisioning,
those mappings, and MSA consent commit together.
`dataPlane` is required and must name an open public plane with ready signing
keys, matching `publicDataPlanes`. Keyless planes are excluded from the tenant
storage mapping as well.
Public AWS planes always use their derived colocated trial bucket; local and
non-AWS planes use the GCS trial bucket. Provision the matching AWS bucket and
permissions before opening a plane for signup.

## Development

Expand Down
222 changes: 221 additions & 1 deletion crates/control-plane-api/src/server/public/graphql/tenant.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
use async_graphql::{Context, Result, SimpleObject};
use validator::Validate;

mod storage;

const TENANT_UNAVAILABLE_MESSAGE: &str = "The organization name is already in use, \
please choose a different one or contact support@estuary.dev.";

Expand Down Expand Up @@ -110,6 +112,11 @@ impl TenantMutation {
#[graphql(desc = "ID of the latest MSA terms the submitting user has read and accepts.")]
submitting_user_agrees_to_terms_id: models::Id,
survey: Option<async_graphql::Json<serde_json::Value>>,
#[graphql(
desc = "Full catalog name of an open public data plane to use by default.",
validator(max_length = 256)
)]
data_plane: String,
) -> async_graphql::Result<bool> {
let env = ctx.data::<crate::Envelope>()?;
let claims = env
Expand Down Expand Up @@ -148,8 +155,10 @@ impl TenantMutation {
}

let tenant_name = create_tenant(
env.snapshot(),
claims.sub,
&name,
&data_plane,
survey.map(|v| v.0).unwrap_or(serde_json::Value::Null),
&mut txn,
)
Expand Down Expand Up @@ -180,8 +189,10 @@ impl TenantMutation {
}

async fn create_tenant(
snapshot: &crate::Snapshot,
user_id: uuid::Uuid,
name: &str,
data_plane: &str,
survey: serde_json::Value,
txn: &mut sqlx::Transaction<'_, sqlx::Postgres>,
) -> async_graphql::Result<String> {
Expand Down Expand Up @@ -220,6 +231,22 @@ async fn create_tenant(
return Err(async_graphql::Error::new(TENANT_UNAVAILABLE_MESSAGE));
}

let mut public_planes: Vec<String> = sqlx::query_scalar(
"SELECT data_plane_name FROM data_planes
WHERE starts_with(data_plane_name, 'ops/dp/public/') AND NOT closed
ORDER BY id DESC",
)
.fetch_all(&mut **txn)
.await?;
// Match publicDataPlanes: managed planes are registered before their HMAC
// material is ready, and cannot receive tenants until they can sign claims.
public_planes.retain(|name| {
snapshot
.data_plane_by_catalog_name(name)
.is_some_and(|plane| plane.can_sign())
});
let (tenant_storage, recovery_storage) = storage::specs(public_planes, data_plane)?;

crate::directives::beta_onboard::provision_tenant(
// TODO: remove unused email param when retiring betaOnboard directive
"",
Expand All @@ -240,6 +267,19 @@ async fn create_tenant(
}
err.into()
})?;
// Keep the legacy directive's provisioning behavior unchanged. Its initial
// mappings and these replacements are committed atomically with the tenant.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We'll simplify some of this stuff when the UI is updated and we can retire the directive

sqlx::query(
"UPDATE storage_mappings SET spec = CASE WHEN catalog_prefix = $1 THEN $2::json ELSE $3::json END
WHERE catalog_prefix = $1 OR catalog_prefix = $4",
)
.bind(&tenant_name)
.bind(tenant_storage)
.bind(recovery_storage)
.bind(format!("recovery/{tenant_name}"))
.execute(&mut **txn)
.await?;

let metadata = if survey.is_null() {
serde_json::json!({})
} else {
Expand Down Expand Up @@ -285,9 +325,10 @@ mod test {

fn request(tenant: &str) -> serde_json::Value {
serde_json::json!({
"query": "mutation($name: String!, $submittingUserAgreesToTermsId: Id!, $survey: JSON) { tenantCreate(name: $name, submittingUserAgreesToTermsId: $submittingUserAgreesToTermsId, survey: $survey) }",
"query": "mutation($name: String!, $submittingUserAgreesToTermsId: Id!, $survey: JSON, $dataPlane: String!) { tenantCreate(name: $name, submittingUserAgreesToTermsId: $submittingUserAgreesToTermsId, survey: $survey, dataPlane: $dataPlane) }",
"variables": {
"name": tenant,
"dataPlane": "ops/dp/public/aws-us-west-2-c1",
"submittingUserAgreesToTermsId": TERMS_ID,
"survey": { "origin": "search", "details": "testing" },
}
Expand Down Expand Up @@ -353,6 +394,185 @@ mod test {
"#);
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
)]
async fn tenant_create_requires_data_plane(pool: sqlx::PgPool) {
let server = server(&pool).await;
let token = server.make_access_token(ALICE, None);
for null in [false, true] {
let mut req = request("acmeCo");
if null {
req["variables"]["dataPlane"] = serde_json::Value::Null;
} else {
req["variables"]
.as_object_mut()
.unwrap()
.remove("dataPlane");
}
let response: serde_json::Value = server.graphql(&req, Some(&token)).await;
let message = response["errors"][0]["message"].as_str().unwrap();
assert!(message.contains("dataPlane"), "{message}");
}
let count: i64 =
sqlx::query_scalar("SELECT count(*) FROM tenants WHERE tenant = 'acmeCo/'")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
)]
async fn tenant_create_data_plane_selection(pool: sqlx::PgPool) {
// Fixtures are inserted after migrations, so explicitly close this plane.
sqlx::query("UPDATE data_planes SET closed = true WHERE data_plane_name = 'ops/dp/public/gcp-us-central1-c2'")
.execute(&pool).await.unwrap();
sqlx::query("UPDATE data_planes SET data_plane_name = 'ops/dp/private/acmeCo/az-westeurope-c3' WHERE data_plane_name = 'ops/dp/public/az-westeurope-c3'")
.execute(&pool).await.unwrap();
let server = server(&pool).await;
let token = server.make_access_token(ALICE, None);
for plane in [
"ops/dp/public/missing",
"ops/dp/private/acmeCo/az-westeurope-c3",
"ops/dp/public/gcp-us-central1-c2",
] {
let mut req = request("acmeCo");
req["variables"]["dataPlane"] = plane.into();
let response: serde_json::Value = server.graphql(&req, Some(&token)).await;
assert_eq!(
response["errors"][0]["message"],
format!("{plane} is not a selectable public data-plane")
);
}
let count: i64 = sqlx::query_scalar("SELECT count(*) FROM user_grants WHERE user_id = $1")
.bind(ALICE)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);

let mut req = request("acmeCo");
req["variables"]["dataPlane"] = "ops/dp/public/aws-us-west-2-c1".into();
let response: serde_json::Value = server.graphql(&req, Some(&token)).await;
insta::assert_json_snapshot!(response, @r#"
{
"data": {
"tenantCreate": true
}
}
"#);
let planes: serde_json::Value = sqlx::query_scalar(
"SELECT spec->'data_planes' FROM storage_mappings WHERE catalog_prefix = 'acmeCo/'",
)
.fetch_one(&pool)
.await
.unwrap();
insta::assert_json_snapshot!(planes, @r#"
[
"ops/dp/public/aws-us-west-2-c1"
]
"#);
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
)]
async fn tenant_create_rejects_unprovisioned_plane(pool: sqlx::PgPool) {
let server = server(&pool).await;
let token = server.make_access_token(ALICE, None);
let mut req = request("acmeCo");
req["variables"]["dataPlane"] = "ops/dp/public/az-westeurope-c3".into();
let response: serde_json::Value = server.graphql(&req, Some(&token)).await;
insta::assert_json_snapshot!(response["errors"][0]["message"], @r#""ops/dp/public/az-westeurope-c3 is not a selectable public data-plane""#);
let counts: (i64, i64, i64, i64) = sqlx::query_as(
"SELECT (SELECT count(*) FROM tenants WHERE tenant = 'acmeCo/'),
(SELECT count(*) FROM storage_mappings WHERE catalog_prefix IN ('acmeCo/', 'recovery/acmeCo/')),
(SELECT count(*) FROM user_grants WHERE user_id = $1),
(SELECT count(*) FROM internal.tenant_consent WHERE user_id = $1)",
).bind(ALICE).fetch_one(&pool).await.unwrap();
assert_eq!(counts, (0, 0, 0, 0));

let response: serde_json::Value = server.graphql(&request("acmeCo"), Some(&token)).await;
insta::assert_json_snapshot!(response, @r#"
{
"data": {
"tenantCreate": true
}
}
"#);
let planes: serde_json::Value = sqlx::query_scalar(
"SELECT spec->'data_planes' FROM storage_mappings WHERE catalog_prefix = 'acmeCo/'",
)
.fetch_one(&pool)
.await
.unwrap();
insta::assert_json_snapshot!(planes, @r#"
[
"ops/dp/public/aws-us-west-2-c1",
"ops/dp/public/gcp-us-central1-c2"
]
"#);
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
)]
async fn tenant_create_no_open_planes(pool: sqlx::PgPool) {
sqlx::query("UPDATE data_planes SET closed = true")
.execute(&pool)
.await
.unwrap();
let server = server(&pool).await;
let token = server.make_access_token(ALICE, None);
let response: serde_json::Value = server.graphql(&request("acmeCo"), Some(&token)).await;
assert_eq!(
response["errors"][0]["message"],
"there are no open public data-planes to place a new tenant on"
);
let count: i64 =
sqlx::query_scalar("SELECT count(*) FROM tenants WHERE tenant = 'acmeCo/'")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(count, 0);
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
)]
async fn tenant_create_colocated_storage(pool: sqlx::PgPool) {
let server = server(&pool).await;
let token = server.make_access_token(ALICE, None);
let response: serde_json::Value = server.graphql(&request("acmeCo"), Some(&token)).await;
insta::assert_json_snapshot!(response, @r#"
{
"data": {
"tenantCreate": true
}
}
"#);
let specs: Vec<serde_json::Value> = sqlx::query_scalar("SELECT spec FROM storage_mappings WHERE catalog_prefix IN ('acmeCo/', 'recovery/acmeCo/') ORDER BY catalog_prefix")
.fetch_all(&pool).await.unwrap();
assert_eq!(specs.len(), 2);
assert_eq!(specs[0]["stores"][0]["provider"], "S3");
assert_eq!(specs[0]["stores"][0]["region"], "us-west-2");
assert_eq!(
specs[0]["stores"][0]["bucket"],
specs[1]["stores"][0]["bucket"]
);
assert_eq!(specs[0]["stores"][0]["prefix"], "collection-data/");
assert!(specs[1]["stores"][0].get("prefix").is_none());
assert_eq!(specs[0]["data_planes"][0], "ops/dp/public/aws-us-west-2-c1");
assert!(specs[1].get("data_planes").is_none());
}

#[sqlx::test(
migrations = "../../supabase/migrations",
fixtures(path = "../../../fixtures", scripts("data_planes"))
Expand Down
Loading
Loading