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
16 changes: 0 additions & 16 deletions common/src/telemetry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -273,14 +273,6 @@ pub enum TelemetryEvent {
grace_secs: u64,
},

/// A member that was down during a password rotation adopted the
/// variables at boot: etcd accepted their password and refused the
/// pinned one.
CredentialsAdopted {
node: String,
variables: Vec<String>,
},

// === Generic Events ===
/// Component started
ComponentStarted { component: String, version: String },
Expand Down Expand Up @@ -334,7 +326,6 @@ impl TelemetryEvent {
Self::HaproxyStarted { .. } => "HAPROXY_STARTED",
Self::HaproxyConfigGenerating { .. } => "HAPROXY_CONFIG_GENERATING",
Self::StandaloneOrphanSlotsDropped { .. } => "STANDALONE_ORPHAN_SLOTS_DROPPED",
Self::CredentialsAdopted { .. } => "POSTGRES_HA_CREDENTIALS_ADOPTED",
Self::ComponentStarted { .. } => "COMPONENT_STARTED",
Self::ComponentError { .. } => "COMPONENT_ERROR",
}
Expand Down Expand Up @@ -603,13 +594,6 @@ impl TelemetryEvent {
grace_secs
)
}
Self::CredentialsAdopted { node, variables } => {
format!(
"{} adopted the rotated credential variables at boot ({})",
node,
variables.join(", ")
)
}
Self::ComponentStarted { component, version } => {
format!("{} v{} started", component, version)
}
Expand Down
92 changes: 12 additions & 80 deletions postgres-patroni/src/bin/patroni_runner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,19 +14,17 @@ use postgres_patroni::bootstrap::{reconcile_pg_stat_statements, refresh_collatio
use postgres_patroni::health_server::{self, HealthServerConfig};
use postgres_patroni::major_upgrade;
use postgres_patroni::patroni::etcd_preflight::{
etcd_password_source_variable, probe_etcd_credential, rejection_message, rotation_adoption,
EtcdAuthProbe, RotationAdoption, REJECTION_PREFIX,
etcd_password_source_variable, probe_etcd_credential, rejection_message, EtcdAuthProbe,
REJECTION_PREFIX,
};
use postgres_patroni::patroni::live_credentials;
use postgres_patroni::patroni::rest_preflight::{
divergence_message, rest_credential_diverges, DIVERGENCE_PREFIX,
};
use postgres_patroni::patroni::{
apply_credential_pin_with_proof, credential_drift, credentials_from_env_requested,
generate_patroni_config, read_credential_pin, reconcile_pgbackrest_archive_config,
run_monitoring_loop, spawn_backup_watcher, spawn_self_heal_watcher,
spawn_slot_recovery_watcher, update_pg_hba_for_replication, Config, Credential,
RestapiAddressSource,
apply_credential_pin, credential_drift, credentials_from_env_requested,
generate_patroni_config, reconcile_pgbackrest_archive_config, run_monitoring_loop,
spawn_backup_watcher, spawn_self_heal_watcher, spawn_slot_recovery_watcher,
update_pg_hba_for_replication, Config, RestapiAddressSource,
};
use postgres_patroni::pgbackrest::{derive_pgbackrest_repo_path, read_wal_level};
use postgres_patroni::wal_archive::{
Expand Down Expand Up @@ -520,23 +518,12 @@ fn drifted_summary(drifted: &[&str]) -> String {
/// starts; a blank `PATRONI_ETCD3_PASSWORD` replaces the etcd password the
/// file carries with spaces. Blank variables are removed from the child's
/// environment so both sides agree they are unset.
///
/// Since the live rotation route (`patroni::rotation`), the password variables
/// are WITHHELD from Patroni's environment altogether rather than mirrored
/// into it: patroni.yml already carries the post-pin values, and with no
/// environment copy to apply over the file, a rewrite of the file plus
/// `POST /reload` moves Patroni to a new password without a restart. The
/// usernames, hosts and every other `PATRONI_*` variable still pass through.
async fn start_patroni(config: &Config) -> Result<tokio::process::Child> {
info!(
node = %config.name,
"starting Patroni with patroni.yml as the only source of its passwords"
);
let mut command = Command::new("patroni");
command.arg("/etc/patroni/patroni.yml");
for var in PASSWORD_VARS_FILE_ONLY {
command.env_remove(var);
}
command
.arg("/etc/patroni/patroni.yml")
.env("PATRONI_REPLICATION_PASSWORD", &config.repl_pass)
.env("PATRONI_SUPERUSER_PASSWORD", &config.superuser_pass);
for var in blank_credential_vars(|name| env::var(name).ok()) {
warn!(
variable = var,
Expand All @@ -554,17 +541,6 @@ async fn start_patroni(config: &Config) -> Result<tokio::process::Child> {
Ok(child)
}

/// Password variables Patroni reads from patroni.yml only. Patroni's loader
/// applies its environment over the file, so any of these in the child's
/// environment would survive a rewrite of the file and defeat a live
/// rotation; the runner renders their post-pin values into the file instead.
const PASSWORD_VARS_FILE_ONLY: [&str; 4] = [
"PATRONI_SUPERUSER_PASSWORD",
"PATRONI_REPLICATION_PASSWORD",
"PATRONI_RESTAPI_PASSWORD",
"PATRONI_ETCD3_PASSWORD",
];

/// Control-plane credential variables the runner treats as unset when blank
/// (see `patroni::config::resolve_restapi_auth` / `resolve_etcd_auth`) and
/// Patroni would treat as set.
Expand Down Expand Up @@ -1394,11 +1370,7 @@ async fn async_main() -> Result<()> {
// reseed wipe removes along with the rest of pgdata); the pin itself is
// applied further down, once the post-reseed state of the data directory
// is known, and reports the same list.
let mut credential_drift = credential_drift(&config, credentials_from_env_requested());
// Set when etcd proves the variables carry the cluster's rotated
// password (see rotation_adoption); apply_credential_pin re-pins from
// the variables instead of keeping the stale pin.
let mut credentials_proven_by_etcd = false;
let credential_drift = credential_drift(&config, credentials_from_env_requested());

// etcd's root password is fixed when the etcd entrypoint first enables
// authentication; the credential this member presents is re-derived from
Expand Down Expand Up @@ -1427,41 +1399,6 @@ async fn async_main() -> Result<()> {
});
anyhow::bail!("{REJECTION_PREFIX}; see the message above for the variable to restore");
}

// A member that was down while the cluster rotated its password boots
// with the variables already moved and a pin that still holds the old
// password. etcd is the oracle: the rotation changed root's password
// to the new value, so etcd accepting the variables' password and
// refusing the pinned one proves the variables are the cluster's
// credentials now. Only that exact pair adopts; a dedicated
// PATRONI_ETCD3_PASSWORD says nothing about the superuser password.
if probe == EtcdAuthProbe::Accepted
&& !credential_drift.is_empty()
&& etcd_password_source_variable(env::var("PATRONI_ETCD3_PASSWORD").ok().as_deref())
== "PATRONI_SUPERUSER_PASSWORD"
{
if let Some(pinned) = read_credential_pin(&config.data_dir) {
let pinned_cred = Credential {
username: cred.username.clone(),
password: pinned.superuser_pass.clone(),
};
let pinned_probe =
probe_etcd_credential(&config.etcd_hosts, &pinned_cred, Duration::from_secs(3))
.await;
if rotation_adoption(probe, pinned_probe) == RotationAdoption::Adopt {
info!(
drifted = ?credential_drift,
"etcd accepts the variables' password and refuses the pinned one: the cluster rotated its password while this member was down; adopting the variables"
);
telemetry.send(TelemetryEvent::CredentialsAdopted {
node: config.name.clone(),
variables: credential_drift.iter().map(|v| v.to_string()).collect(),
});
credentials_proven_by_etcd = true;
credential_drift = Vec::new();
}
}
}
}

// The same edit seen from the REST API, for the cluster where etcd does
Expand Down Expand Up @@ -1532,18 +1469,13 @@ async fn async_main() -> Result<()> {
// patroni::credential_pin). Must run before generate_patroni_config and
// before anything else reads config.*_pass.
let has_cluster_data = has_pg_control && has_marker;
let pin_outcome = apply_credential_pin_with_proof(
let pin_outcome = apply_credential_pin(
&mut config,
has_cluster_data,
credentials_from_env_requested(),
credentials_proven_by_etcd,
&telemetry,
);
info!(outcome = ?pin_outcome, "credential pin reconciled");
// The passwords in force for the rest of this process: the rotation
// route compares against them and moves them; the REST client and the
// wrapper's psql calls read them at call time.
live_credentials::seed(&config);

// Recover the debris of an interrupted clone. A non-empty data directory
// with NO pg_control is what a pg_basebackup killed mid-stream leaves
Expand Down
2 changes: 1 addition & 1 deletion postgres-patroni/src/health_server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ async fn run(config: HealthServerConfig) -> Result<()> {
info!(
port,
patroni_fallback_port = patroni_port,
"Health server listening (endpoints: /primary, /replica, /health, POST /credentials/rotate)"
"Health server listening (endpoints: /primary, /replica, /health)"
);

axum::serve(listener, app)
Expand Down
168 changes: 2 additions & 166 deletions postgres-patroni/src/health_server/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,82 +2,19 @@

use super::config::HealthServerConfig;
use super::postgres::{check_patroni_role, is_in_recovery};
use crate::patroni::rotation::{self, MemberPorts, RotateRequest};
use axum::body::Bytes;
use axum::http::header::{AUTHORIZATION, WWW_AUTHENTICATE};
use axum::http::HeaderMap;
use axum::response::Response;
use axum::{
extract::State,
http::StatusCode,
response::IntoResponse,
routing::{get, post},
Json, Router,
};
use axum::{extract::State, http::StatusCode, response::IntoResponse, routing::get, Router};
use std::time::Duration;
use tracing::{debug, warn};
use tracing::debug;

/// Create the router with all health check endpoints
///
/// The GET routes stay open: HAProxy builds both backends from /primary and
/// /replica. `POST /credentials/rotate` is the one mutating route and
/// requires the member's REST API credential (see `patroni::rotation`).
pub fn create_router(config: HealthServerConfig) -> Router {
Router::new()
.route("/primary", get(primary_handler))
.route("/replica", get(replica_handler))
.route("/health", get(health_handler))
.route(rotation::ROUTE, post(rotate_handler))
.with_state(config)
}

/// Handler for `POST /credentials/rotate`
///
/// Basic auth with the member's current REST API credential, JSON body
/// `{"password": "<new>"}`. 200 with a summary when the member moved (or
/// already ran on that password), 401 without the credential, 400 on a bad
/// body, 409 when a role does not carry the password yet, 503 when the
/// verifier cannot be read, 500 when the switch failed and was undone.
async fn rotate_handler(
State(config): State<HealthServerConfig>,
headers: HeaderMap,
body: Bytes,
) -> Response {
let authorization = headers.get(AUTHORIZATION).and_then(|v| v.to_str().ok());
if !rotation::authorize(authorization) {
warn!("credentials/rotate refused: missing or wrong REST API credential");
return (
StatusCode::UNAUTHORIZED,
[(WWW_AUTHENTICATE, "Basic realm=\"postgres-ha\"")],
Json(serde_json::json!({ "error": rotation::RotateError::Unauthorized.message() })),
)
.into_response();
}
let request: RotateRequest = match serde_json::from_slice(&body) {
Ok(request) => request,
Err(e) => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({ "error": format!("body must be {{\"password\": \"...\"}}: {e}") })),
)
.into_response()
}
};
let ports = MemberPorts {
pg_port: config.pg_port,
patroni_port: config.patroni_port,
};
match rotation::rotate(&request.password, &ports).await {
Ok(summary) => (StatusCode::OK, Json(summary)).into_response(),
Err(e) => {
let status =
StatusCode::from_u16(e.status()).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
warn!(status = status.as_u16(), error = %e.message(), "credentials/rotate did not rotate");
(status, Json(serde_json::json!({ "error": e.message() }))).into_response()
}
}
}

/// Handler for /primary endpoint
///
/// Returns 200 if this node is the primary (pg_is_in_recovery() = false)
Expand Down Expand Up @@ -216,104 +153,3 @@ async fn health_check(config: &HealthServerConfig) -> (StatusCode, &'static str)
}
}
}

#[cfg(test)]
mod tests {
use super::*;
use crate::patroni::live_credentials;
use crate::patroni::rest::test_support::ENV_LOCK;
use base64::{engine::general_purpose::STANDARD as BASE64, Engine as _};

fn test_config() -> HealthServerConfig {
HealthServerConfig {
port: 0,
pg_port: 5432,
pg_user: "postgres".into(),
pg_password: String::new(),
pg_database: "postgres".into(),
patroni_port: 8008,
check_timeout_ms: 2000,
}
}

async fn serve() -> String {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
axum::serve(listener, create_router(test_config()))
.await
.unwrap();
});
format!("http://{addr}")
}

fn basic(user: &str, pass: &str) -> String {
format!("Basic {}", BASE64.encode(format!("{user}:{pass}")))
}

#[tokio::test]
async fn rotate_route_requires_the_rest_credential_and_reads_the_body() {
let _lock = ENV_LOCK.lock().unwrap();
live_credentials::reset();
std::env::set_var("PATRONI_SUPERUSER_USERNAME", "postgres");
std::env::set_var("PATRONI_SUPERUSER_PASSWORD", "su");
std::env::set_var("PATRONI_RESTAPI_PASSWORD", "rest");
live_credentials::set_rotated("current", true, true);
let base = serve().await;
let client = reqwest::Client::new();
let url = format!("{base}{}", rotation::ROUTE);

// No credential: 401 and a challenge.
let response = client
.post(&url)
.json(&serde_json::json!({ "password": "x" }))
.send()
.await
.unwrap();
assert_eq!(response.status().as_u16(), 401);
assert!(response.headers().contains_key(WWW_AUTHENTICATE));

// Wrong credential (the superuser password is not the REST one here).
let response = client
.post(&url)
.header(AUTHORIZATION, basic("postgres", "su"))
.json(&serde_json::json!({ "password": "x" }))
.send()
.await
.unwrap();
assert_eq!(response.status().as_u16(), 401);

// Right credential, unusable body: 400.
let response = client
.post(&url)
.header(AUTHORIZATION, basic("postgres", "current"))
.body("not json")
.send()
.await
.unwrap();
assert_eq!(response.status().as_u16(), 400);

// Right credential, the password already in force: 200 "already",
// nothing touched (no Postgres or Patroni needed).
let response = client
.post(&url)
.header(AUTHORIZATION, basic("postgres", "current"))
.json(&serde_json::json!({ "password": "current" }))
.send()
.await
.unwrap();
assert_eq!(response.status().as_u16(), 200);
let body: serde_json::Value = response.json().await.unwrap();
assert_eq!(body["status"], "already");
assert_eq!(body["reloaded"], false);

// The GET routes stay open.
let response = client.get(format!("{base}/health")).send().await.unwrap();
assert_ne!(response.status().as_u16(), 401);

live_credentials::reset();
std::env::remove_var("PATRONI_RESTAPI_PASSWORD");
std::env::remove_var("PATRONI_SUPERUSER_PASSWORD");
std::env::remove_var("PATRONI_SUPERUSER_USERNAME");
}
}
Loading
Loading