diff --git a/Cargo.lock b/Cargo.lock index 5835a20..1a6ac6a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1365,6 +1365,32 @@ dependencies = [ "slab", ] +[[package]] +name = "gcp_auth" +version = "0.12.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2b3d0b409a042a380111af38136310839af8ac1a0917fb6e84515ed1e4bf3ee" +dependencies = [ + "async-trait", + "base64", + "bytes", + "chrono", + "http 1.3.1", + "http-body-util", + "hyper 1.7.0", + "hyper-rustls 0.27.7", + "hyper-util", + "ring", + "rustls-pki-types", + "serde", + "serde_json", + "thiserror 2.0.17", + "tokio", + "tracing", + "tracing-futures", + "url", +] + [[package]] name = "generic-array" version = "0.14.7" @@ -2345,6 +2371,19 @@ dependencies = [ "tokio", ] +[[package]] +name = "mlf-dns-google" +version = "0.1.0" +dependencies = [ + "gcp_auth", + "mlf-plugin-host", + "reqwest", + "serde", + "serde_json", + "thiserror 2.0.17", + "tokio", +] + [[package]] name = "mlf-dns-namecheap" version = "0.1.0" @@ -3114,6 +3153,7 @@ checksum = "cd3c25631629d034ce7cd9940adc9d45762d46de2b0f57193c4443b92c6d4d40" dependencies = [ "aws-lc-rs", "once_cell", + "ring", "rustls-pki-types", "rustls-webpki 0.103.7", "subtle", @@ -4005,6 +4045,16 @@ dependencies = [ "valuable", ] +[[package]] +name = "tracing-futures" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97d095ae15e245a057c8e8451bab9b3ee1e1f68e9ba2b4fbc18d0ac5237835f2" +dependencies = [ + "pin-project", + "tracing", +] + [[package]] name = "tracing-log" version = "0.2.0" diff --git a/Cargo.toml b/Cargo.toml index 059cf9c..ee0cf89 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,6 +8,7 @@ members = [ "mlf-atproto", "mlf-cli", "dns-plugins/mlf-dns-godaddy", + "dns-plugins/mlf-dns-google", "dns-plugins/mlf-dns-namecheap", "dns-plugins/mlf-dns-porkbun", "dns-plugins/mlf-dns-route53", diff --git a/dns-plugins/mlf-dns-google/Cargo.toml b/dns-plugins/mlf-dns-google/Cargo.toml new file mode 100644 index 0000000..e72cd6e --- /dev/null +++ b/dns-plugins/mlf-dns-google/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "mlf-dns-google" +version = "0.1.0" +edition = "2024" +license = "MIT" +description = "Official MLF DNS provider plugin for Google Cloud DNS" + +[[bin]] +name = "mlf-dns-google" +path = "src/main.rs" + +[dependencies] +mlf-plugin-host = { path = "../../mlf-plugin-host" } +gcp_auth = "0.12" +reqwest = { version = "0.12", features = ["json"] } +serde = { version = "1", features = ["derive"] } +serde_json = "1" +thiserror = "2" +tokio = { version = "1", features = ["io-util", "macros", "rt"] } diff --git a/dns-plugins/mlf-dns-google/src/api.rs b/dns-plugins/mlf-dns-google/src/api.rs new file mode 100644 index 0000000..c8b2ee0 --- /dev/null +++ b/dns-plugins/mlf-dns-google/src/api.rs @@ -0,0 +1,368 @@ +//! Thin Google Cloud DNS wrapper. +//! +//! Auth: a service-account JSON key (the content, not a path) or +//! Application Default Credentials via `gcp_auth`. Every API call +//! sends a freshly-fetched bearer token. +//! +//! Cloud DNS addresses records by `(managedZone, name, type)`. +//! Resource-record-sets are atomically modified via the `changes` +//! endpoint, which takes `additions` and `deletions` lists in one +//! request — handy for upsert (delete old + add new together). + +use crate::Credentials; +use gcp_auth::{CustomServiceAccount, TokenProvider}; +use serde::{Deserialize, Serialize}; +use std::sync::Arc; +use thiserror::Error; + +const API_BASE: &str = "https://dns.googleapis.com/dns/v1"; +const SCOPE: &[&str] = &["https://www.googleapis.com/auth/ndev.clouddns.readwrite"]; + +#[derive(Error, Debug)] +pub enum GoogleDnsError { + #[error("Google DNS HTTP error: {0}")] + Http(String), + #[error("Google DNS API error: {0}")] + Api(String), + #[error("Auth error: {0}")] + Auth(String), + #[error("JSON decode error: {0}")] + Decode(String), +} + +impl From for GoogleDnsError { + fn from(e: reqwest::Error) -> Self { + GoogleDnsError::Http(e.to_string()) + } +} + +pub struct GoogleDnsClient { + client: reqwest::Client, + provider: Arc, + project: String, +} + +impl GoogleDnsClient { + pub async fn new(creds: &Credentials) -> Result { + let provider: Arc = + if let Some(json) = creds.service_account_json.as_deref() { + let sa = CustomServiceAccount::from_json(json) + .map_err(|e| GoogleDnsError::Auth(e.to_string()))?; + Arc::new(sa) + } else { + gcp_auth::provider() + .await + .map_err(|e| GoogleDnsError::Auth(e.to_string()))? + }; + Ok(Self { + client: reqwest::Client::new(), + provider, + project: creds.project_id.clone(), + }) + } + + async fn token(&self) -> Result { + let token = self + .provider + .token(SCOPE) + .await + .map_err(|e| GoogleDnsError::Auth(e.to_string()))?; + Ok(token.as_str().to_string()) + } + + /// Validate auth by listing the zones in the project. + pub async fn verify(&self) -> Result { + let zones = self.list_zones().await?; + Ok(format!( + "google ({} — {} zone(s))", + self.project, + zones.len() + )) + } + + /// All managed zones under the project. + pub async fn list_zones(&self) -> Result, GoogleDnsError> { + #[derive(Deserialize)] + struct Resp { + #[serde(default)] + #[serde(rename = "managedZones")] + managed_zones: Vec, + } + let token = self.token().await?; + let resp = self + .client + .get(format!("{API_BASE}/projects/{}/managedZones", self.project)) + .bearer_auth(token) + .send() + .await?; + if !resp.status().is_success() { + let status = resp.status(); + let body = resp.text().await.unwrap_or_default(); + return Err(GoogleDnsError::Api(format!("HTTP {status}: {body}"))); + } + let parsed: Resp = resp + .json() + .await + .map_err(|e| GoogleDnsError::Decode(e.to_string()))?; + Ok(parsed.managed_zones) + } + + pub async fn find_zone_for( + &self, + dns_name: &str, + ) -> Result, GoogleDnsError> { + let stripped = dns_name.strip_prefix("_lexicon.").unwrap_or(dns_name); + let zones = self.list_zones().await?; + // Zones in Cloud DNS end with a dot; compare without. + for candidate in parent_domains(stripped) { + if let Some(zone) = zones + .iter() + .find(|z| z.dns_name.trim_end_matches('.') == candidate) + { + return Ok(Some(zone.clone())); + } + } + Ok(None) + } + + pub async fn list_txt( + &self, + zone_name: &str, + dns_name: &str, + ) -> Result, GoogleDnsError> { + let name = with_trailing_dot(dns_name); + let token = self.token().await?; + #[derive(Deserialize)] + struct Resp { + #[serde(default)] + rrsets: Vec, + } + let resp = self + .client + .get(format!( + "{API_BASE}/projects/{}/managedZones/{}/rrsets", + self.project, zone_name + )) + .bearer_auth(token) + .query(&[("name", name.as_str()), ("type", "TXT")]) + .send() + .await?; + if !resp.status().is_success() { + let status = resp.status(); + let body = resp.text().await.unwrap_or_default(); + return Err(GoogleDnsError::Api(format!("HTTP {status}: {body}"))); + } + let parsed: Resp = resp + .json() + .await + .map_err(|e| GoogleDnsError::Decode(e.to_string()))?; + let mut out = Vec::new(); + for rrset in parsed.rrsets { + if rrset.r#type != "TXT" { + continue; + } + for value in rrset.rrdatas { + out.push(TxtRecord { + id: format!("{}/TXT", rrset.name.trim_end_matches('.')), + value: unquote(&value), + }); + } + } + Ok(out) + } + + /// Atomically replace the TXT rrset at `dns_name` with a single value. + pub async fn upsert_txt( + &self, + zone_name: &str, + dns_name: &str, + value: &str, + ttl: u32, + ) -> Result { + let name = with_trailing_dot(dns_name); + let quoted = format!("\"{}\"", escape_for_txt(value)); + let existing = self.get_rrset(zone_name, &name, "TXT").await?; + let change = Change { + additions: vec![Rrset { + name: name.clone(), + r#type: "TXT".into(), + ttl: ttl as i64, + rrdatas: vec![quoted], + }], + deletions: existing.into_iter().collect(), + }; + self.submit_change(zone_name, &change).await?; + Ok(name.trim_end_matches('.').to_string()) + } + + pub async fn delete_txt(&self, zone_name: &str, dns_name: &str) -> Result<(), GoogleDnsError> { + let name = with_trailing_dot(dns_name); + let existing = self.get_rrset(zone_name, &name, "TXT").await?; + if existing.is_empty() { + return Ok(()); + } + let change = Change { + additions: vec![], + deletions: existing, + }; + self.submit_change(zone_name, &change).await + } + + async fn get_rrset( + &self, + zone_name: &str, + dns_name: &str, + record_type: &str, + ) -> Result, GoogleDnsError> { + let token = self.token().await?; + #[derive(Deserialize)] + struct Resp { + #[serde(default)] + rrsets: Vec, + } + let resp = self + .client + .get(format!( + "{API_BASE}/projects/{}/managedZones/{}/rrsets", + self.project, zone_name + )) + .bearer_auth(token) + .query(&[("name", dns_name), ("type", record_type)]) + .send() + .await?; + if !resp.status().is_success() { + let status = resp.status(); + let body = resp.text().await.unwrap_or_default(); + return Err(GoogleDnsError::Api(format!("HTTP {status}: {body}"))); + } + let parsed: Resp = resp + .json() + .await + .map_err(|e| GoogleDnsError::Decode(e.to_string()))?; + Ok(parsed.rrsets) + } + + async fn submit_change(&self, zone_name: &str, change: &Change) -> Result<(), GoogleDnsError> { + let token = self.token().await?; + let resp = self + .client + .post(format!( + "{API_BASE}/projects/{}/managedZones/{}/changes", + self.project, zone_name + )) + .bearer_auth(token) + .json(change) + .send() + .await?; + if !resp.status().is_success() { + let status = resp.status(); + let body = resp.text().await.unwrap_or_default(); + return Err(GoogleDnsError::Api(format!("HTTP {status}: {body}"))); + } + Ok(()) + } +} + +#[derive(Debug, Clone, Deserialize)] +pub struct ManagedZone { + pub name: String, + #[serde(rename = "dnsName")] + pub dns_name: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct Rrset { + name: String, + r#type: String, + ttl: i64, + #[serde(default)] + rrdatas: Vec, +} + +#[derive(Debug, Clone, Serialize)] +struct Change { + additions: Vec, + deletions: Vec, +} + +#[derive(Debug, Clone)] +pub struct TxtRecord { + pub id: String, + pub value: String, +} + +fn parent_domains(name: &str) -> Vec { + let parts: Vec<&str> = name.split('.').collect(); + (0..parts.len()).map(|i| parts[i..].join(".")).collect() +} + +fn with_trailing_dot(s: &str) -> String { + if s.ends_with('.') { + s.to_string() + } else { + format!("{s}.") + } +} + +fn escape_for_txt(s: &str) -> String { + s.replace('\\', "\\\\").replace('"', "\\\"") +} + +/// Cloud DNS returns rrdatas surrounded with double quotes + backslash +/// escapes. Strip one set of outer quotes and unescape `\\` / `\"`. +fn unquote(s: &str) -> String { + let inner = s + .strip_prefix('"') + .and_then(|t| t.strip_suffix('"')) + .unwrap_or(s); + let mut out = String::with_capacity(inner.len()); + let mut chars = inner.chars(); + while let Some(c) = chars.next() { + if c == '\\' { + match chars.next() { + Some('\\') => out.push('\\'), + Some('"') => out.push('"'), + Some(other) => { + out.push('\\'); + out.push(other); + } + None => out.push('\\'), + } + } else { + out.push(c); + } + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parent_domains_walks_up() { + assert_eq!( + parent_domains("_lexicon.forum.example.com"), + vec![ + "_lexicon.forum.example.com", + "forum.example.com", + "example.com", + "com", + ] + ); + } + + #[test] + fn with_trailing_dot_is_idempotent() { + assert_eq!(with_trailing_dot("example.com"), "example.com."); + assert_eq!(with_trailing_dot("example.com."), "example.com."); + } + + #[test] + fn txt_escape_round_trip() { + let raw = r#"did=did:plc:"hello""#; + let escaped = escape_for_txt(raw); + let wrapped = format!("\"{escaped}\""); + assert_eq!(unquote(&wrapped), raw); + } +} diff --git a/dns-plugins/mlf-dns-google/src/main.rs b/dns-plugins/mlf-dns-google/src/main.rs new file mode 100644 index 0000000..838417d --- /dev/null +++ b/dns-plugins/mlf-dns-google/src/main.rs @@ -0,0 +1,278 @@ +//! Official MLF DNS provider plugin for Google Cloud DNS. +//! +//! Options schema: +//! - `project_id` (non-secret, required) — GCP project that owns the zones +//! - `service_account_json` (secret, optional) — JSON key content. If +//! omitted, the plugin falls back to Application Default Credentials +//! (useful on GCE / Cloud Run where the runtime provides them). +//! +//! The plugin holds a `gcp_auth::TokenProvider` alive for the session +//! so repeated ops share token caching. + +mod api; + +use api::{GoogleDnsClient, GoogleDnsError}; +use mlf_plugin_host::plugin::{Server, empty_data, params_as}; +use mlf_plugin_host::protocol::{HelloData, OptionField, PROTOCOL_VERSION, Request}; +use serde::{Deserialize, Serialize}; +use serde_json::{Value, json}; + +#[tokio::main(flavor = "current_thread")] +async fn main() -> std::io::Result<()> { + let mut server = Server::stdio(); + + let identity = HelloData { + name: "google".into(), + protocol_version: PROTOCOL_VERSION, + kind: Some("dns".into()), + capabilities: vec![ + "login".into(), + "list_txt".into(), + "upsert_txt".into(), + "delete_txt".into(), + "resolve_zone".into(), + ], + options_schema: vec![ + OptionField { + name: "project_id".into(), + label: "GCP project ID".into(), + help: Some( + "The project that owns the Cloud DNS managed zones you'll publish under." + .into(), + ), + secret: false, + required: true, + default: None, + }, + OptionField { + name: "service_account_json".into(), + label: "Service account JSON key (content, not path)".into(), + help: Some( + "Paste the full JSON body of a service-account key with \ + roles/dns.admin. If omitted, the plugin uses Application \ + Default Credentials (gcloud auth, metadata server, etc.)." + .into(), + ), + secret: true, + required: false, + default: None, + }, + ], + }; + + if server.handshake(identity).await.is_err() { + return Ok(()); + } + + let mut creds: Option = None; + + while let Ok(Some(req)) = server.next_request().await { + if let Err(e) = dispatch(&mut server, &req, &mut creds).await { + let _ = server.reply_err("internal", &e.to_string(), false).await; + } + } + + Ok(()) +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Credentials { + pub project_id: String, + #[serde(default)] + pub service_account_json: Option, +} + +#[derive(Debug, Deserialize)] +struct InitParams { + #[serde(default)] + credentials: Option, +} + +#[derive(Debug, Deserialize)] +struct ResolveZoneParams { + domain: String, +} + +#[derive(Debug, Deserialize)] +struct ListTxtParams { + name: String, +} + +#[derive(Debug, Deserialize)] +struct UpsertTxtParams { + name: String, + value: String, + #[serde(default)] + ttl: Option, +} + +#[derive(Debug, Deserialize)] +struct DeleteTxtParams { + name: String, + #[allow(dead_code)] + record_id: String, +} + +#[derive(thiserror::Error, Debug)] +enum DispatchError { + #[error("{0}")] + Plugin(#[from] mlf_plugin_host::plugin::PluginError), + #[error("{0}")] + Google(#[from] GoogleDnsError), +} + +async fn dispatch( + server: &mut Server, + req: &Request, + creds: &mut Option, +) -> Result<(), DispatchError> +where + W: tokio::io::AsyncWrite + Unpin, + R: tokio::io::AsyncBufReadExt + Unpin, +{ + match req.op.as_str() { + "init" => { + let InitParams { credentials } = params_as(req)?; + *creds = credentials; + server.reply_ok(empty_data()).await?; + } + "login" => { + let Some(c) = creds.as_ref() else { + server + .reply_err( + "no_credentials", + "login called before init set credentials", + false, + ) + .await?; + return Ok(()); + }; + match GoogleDnsClient::new(c).await { + Ok(client) => match client.verify().await { + Ok(name) => { + server + .reply_ok(json!({"credentials": c, "display_name": name})) + .await?; + } + Err(e) => { + server + .reply_err("invalid_credentials", &e.to_string(), false) + .await?; + } + }, + Err(e) => { + server + .reply_err("invalid_credentials", &e.to_string(), false) + .await?; + } + } + } + "logout" => { + *creds = None; + server.reply_ok(empty_data()).await?; + } + "resolve_zone" => { + let ResolveZoneParams { domain } = params_as(req)?; + let c = require_creds(server, creds).await?; + let client = GoogleDnsClient::new(&c).await?; + match client.find_zone_for(&domain).await? { + Some(zone) => { + server + .reply_ok(json!({"zone_id": zone.name, "covered": true})) + .await?; + } + None => { + server + .reply_ok(json!({"zone_id": Value::Null, "covered": false})) + .await?; + } + } + } + "list_txt" => { + let ListTxtParams { name } = params_as(req)?; + let c = require_creds(server, creds).await?; + let client = GoogleDnsClient::new(&c).await?; + let zone = match client.find_zone_for(&name).await? { + Some(z) => z, + None => { + server + .reply_err("unknown_zone", &format!("no zone covers {name}"), false) + .await?; + return Ok(()); + } + }; + let records = client.list_txt(&zone.name, &name).await?; + server + .reply_ok(json!({ + "records": records.into_iter().map(|r| json!({ + "id": r.id, + "value": r.value, + })).collect::>(), + })) + .await?; + } + "upsert_txt" => { + let UpsertTxtParams { name, value, ttl } = params_as(req)?; + let c = require_creds(server, creds).await?; + let client = GoogleDnsClient::new(&c).await?; + let zone = match client.find_zone_for(&name).await? { + Some(z) => z, + None => { + server + .reply_err("unknown_zone", &format!("no zone covers {name}"), false) + .await?; + return Ok(()); + } + }; + let id = client + .upsert_txt(&zone.name, &name, &value, ttl.unwrap_or(300)) + .await?; + server.reply_ok(json!({ "record_id": id })).await?; + } + "delete_txt" => { + let DeleteTxtParams { name, record_id: _ } = params_as(req)?; + let c = require_creds(server, creds).await?; + let client = GoogleDnsClient::new(&c).await?; + let zone = match client.find_zone_for(&name).await? { + Some(z) => z, + None => { + server + .reply_err("unknown_zone", &format!("no zone covers {name}"), false) + .await?; + return Ok(()); + } + }; + client.delete_txt(&zone.name, &name).await?; + server.reply_ok(empty_data()).await?; + } + other => { + server + .reply_err("unknown_op", &format!("unsupported op `{other}`"), false) + .await?; + } + } + Ok(()) +} + +async fn require_creds( + server: &mut Server, + creds: &Option, +) -> Result +where + W: tokio::io::AsyncWrite + Unpin, + R: tokio::io::AsyncBufReadExt + Unpin, +{ + if let Some(c) = creds.clone() { + return Ok(c); + } + server + .reply_err( + "no_credentials", + "host hasn't called init with credentials yet", + false, + ) + .await?; + Err(DispatchError::Plugin( + mlf_plugin_host::plugin::PluginError::Unexpected("missing credentials".into()), + )) +}