) -> &HttpState {
async fn handle_re_attach(mut req: Request) -> Result, ApiError> {
let reattach_req = json_request::(&mut req).await?;
let state = get_state(&req);
- json_response(
- StatusCode::OK,
- state
- .service
- .re_attach(reattach_req)
- .await
- .map_err(ApiError::InternalServerError)?,
- )
+ json_response(StatusCode::OK, state.service.re_attach(reattach_req).await?)
}
/// Pageserver calls into this before doing deletions, to confirm that it still
@@ -114,7 +107,10 @@ async fn handle_tenant_create(
mut req: Request,
) -> Result, ApiError> {
let create_req = json_request::(&mut req).await?;
- json_response(StatusCode::OK, service.tenant_create(create_req).await?)
+ json_response(
+ StatusCode::CREATED,
+ service.tenant_create(create_req).await?,
+ )
}
// For tenant and timeline deletions, which both implement an "initially return 202, then 404 once
@@ -177,6 +173,39 @@ async fn handle_tenant_location_config(
)
}
+async fn handle_tenant_time_travel_remote_storage(
+ service: Arc,
+ mut req: Request,
+) -> Result, ApiError> {
+ let tenant_id: TenantId = parse_request_param(&req, "tenant_id")?;
+ let time_travel_req = json_request::(&mut req).await?;
+
+ let timestamp_raw = must_get_query_param(&req, "travel_to")?;
+ let _timestamp = humantime::parse_rfc3339(×tamp_raw).map_err(|_e| {
+ ApiError::BadRequest(anyhow::anyhow!(
+ "Invalid time for travel_to: {timestamp_raw:?}"
+ ))
+ })?;
+
+ let done_if_after_raw = must_get_query_param(&req, "done_if_after")?;
+ let _done_if_after = humantime::parse_rfc3339(&done_if_after_raw).map_err(|_e| {
+ ApiError::BadRequest(anyhow::anyhow!(
+ "Invalid time for done_if_after: {done_if_after_raw:?}"
+ ))
+ })?;
+
+ service
+ .tenant_time_travel_remote_storage(
+ &time_travel_req,
+ tenant_id,
+ timestamp_raw,
+ done_if_after_raw,
+ )
+ .await?;
+
+ json_response(StatusCode::OK, ())
+}
+
async fn handle_tenant_delete(
service: Arc,
req: Request,
@@ -196,7 +225,7 @@ async fn handle_tenant_timeline_create(
let tenant_id: TenantId = parse_request_param(&req, "tenant_id")?;
let create_req = json_request::(&mut req).await?;
json_response(
- StatusCode::OK,
+ StatusCode::CREATED,
service
.tenant_timeline_create(tenant_id, create_req)
.await?,
@@ -296,7 +325,10 @@ async fn handle_node_configure(mut req: Request) -> Result,
}
let state = get_state(&req);
- json_response(StatusCode::OK, state.service.node_configure(config_req)?)
+ json_response(
+ StatusCode::OK,
+ state.service.node_configure(config_req).await?,
+ )
}
async fn handle_tenant_shard_split(
@@ -333,6 +365,22 @@ async fn handle_tenant_drop(req: Request) -> Result, ApiErr
json_response(StatusCode::OK, state.service.tenant_drop(tenant_id).await?)
}
+async fn handle_tenants_dump(req: Request) -> Result, ApiError> {
+ let state = get_state(&req);
+ state.service.tenants_dump()
+}
+
+async fn handle_scheduler_dump(req: Request) -> Result, ApiError> {
+ let state = get_state(&req);
+ state.service.scheduler_dump()
+}
+
+async fn handle_consistency_check(req: Request) -> Result, ApiError> {
+ let state = get_state(&req);
+
+ json_response(StatusCode::OK, state.service.consistency_check().await?)
+}
+
/// Status endpoint is just used for checking that our HTTP listener is up
async fn handle_status(_req: Request) -> Result, ApiError> {
json_response(StatusCode::OK, ())
@@ -421,6 +469,13 @@ pub fn make_router(
.post("/debug/v1/node/:node_id/drop", |r| {
request_span(r, handle_node_drop)
})
+ .get("/debug/v1/tenant", |r| request_span(r, handle_tenants_dump))
+ .get("/debug/v1/scheduler", |r| {
+ request_span(r, handle_scheduler_dump)
+ })
+ .post("/debug/v1/consistency_check", |r| {
+ request_span(r, handle_consistency_check)
+ })
.get("/control/v1/tenant/:tenant_id/locate", |r| {
tenant_service_handler(r, handle_tenant_locate)
})
@@ -451,6 +506,9 @@ pub fn make_router(
.put("/v1/tenant/:tenant_id/location_config", |r| {
tenant_service_handler(r, handle_tenant_location_config)
})
+ .put("/v1/tenant/:tenant_id/time_travel_remote_storage", |r| {
+ tenant_service_handler(r, handle_tenant_time_travel_remote_storage)
+ })
// Timeline operations
.delete("/v1/tenant/:tenant_id/timeline/:timeline_id", |r| {
tenant_service_handler(r, handle_tenant_timeline_delete)
diff --git a/control_plane/attachment_service/src/lib.rs b/control_plane/attachment_service/src/lib.rs
index 238efdf5a8..e950a57e57 100644
--- a/control_plane/attachment_service/src/lib.rs
+++ b/control_plane/attachment_service/src/lib.rs
@@ -3,6 +3,7 @@ use utils::seqwait::MonotonicCounter;
mod compute_hook;
pub mod http;
+pub mod metrics;
mod node;
pub mod persistence;
mod reconciler;
@@ -11,7 +12,7 @@ mod schema;
pub mod service;
mod tenant_state;
-#[derive(Clone, Serialize, Deserialize)]
+#[derive(Clone, Serialize, Deserialize, Debug)]
enum PlacementPolicy {
/// Cheapest way to attach a tenant: just one pageserver, no secondary
Single,
@@ -22,7 +23,7 @@ enum PlacementPolicy {
Detached,
}
-#[derive(Ord, PartialOrd, Eq, PartialEq, Copy, Clone)]
+#[derive(Ord, PartialOrd, Eq, PartialEq, Copy, Clone, Serialize)]
struct Sequence(u64);
impl Sequence {
diff --git a/control_plane/attachment_service/src/main.rs b/control_plane/attachment_service/src/main.rs
index b323ae8820..db4f00644f 100644
--- a/control_plane/attachment_service/src/main.rs
+++ b/control_plane/attachment_service/src/main.rs
@@ -6,6 +6,7 @@
///
use anyhow::{anyhow, Context};
use attachment_service::http::make_router;
+use attachment_service::metrics::preinitialize_metrics;
use attachment_service::persistence::Persistence;
use attachment_service::service::{Config, Service};
use aws_config::{self, BehaviorVersion, Region};
@@ -205,6 +206,8 @@ async fn async_main() -> anyhow::Result<()> {
logging::Output::Stdout,
)?;
+ preinitialize_metrics();
+
let args = Cli::parse();
tracing::info!(
"version: {}, launch_timestamp: {}, build_tag {}, state at {}, listening on {}",
diff --git a/control_plane/attachment_service/src/metrics.rs b/control_plane/attachment_service/src/metrics.rs
new file mode 100644
index 0000000000..ffe093b9c8
--- /dev/null
+++ b/control_plane/attachment_service/src/metrics.rs
@@ -0,0 +1,32 @@
+use metrics::{register_int_counter, register_int_counter_vec, IntCounter, IntCounterVec};
+use once_cell::sync::Lazy;
+
+pub(crate) struct ReconcilerMetrics {
+ pub(crate) spawned: IntCounter,
+ pub(crate) complete: IntCounterVec,
+}
+
+impl ReconcilerMetrics {
+ // Labels used on [`Self::complete`]
+ pub(crate) const SUCCESS: &'static str = "ok";
+ pub(crate) const ERROR: &'static str = "success";
+ pub(crate) const CANCEL: &'static str = "cancel";
+}
+
+pub(crate) static RECONCILER: Lazy = Lazy::new(|| ReconcilerMetrics {
+ spawned: register_int_counter!(
+ "storage_controller_reconcile_spawn",
+ "Count of how many times we spawn a reconcile task",
+ )
+ .expect("failed to define a metric"),
+ complete: register_int_counter_vec!(
+ "storage_controller_reconcile_complete",
+ "Reconciler tasks completed, broken down by success/failure/cancelled",
+ &["status"],
+ )
+ .expect("failed to define a metric"),
+});
+
+pub fn preinitialize_metrics() {
+ Lazy::force(&RECONCILER);
+}
diff --git a/control_plane/attachment_service/src/node.rs b/control_plane/attachment_service/src/node.rs
index 47f61702d8..09162701ac 100644
--- a/control_plane/attachment_service/src/node.rs
+++ b/control_plane/attachment_service/src/node.rs
@@ -1,9 +1,16 @@
use control_plane::attachment_service::{NodeAvailability, NodeSchedulingPolicy};
+use serde::Serialize;
use utils::id::NodeId;
use crate::persistence::NodePersistence;
-#[derive(Clone)]
+/// Represents the in-memory description of a Node.
+///
+/// Scheduling statistics are maintened separately in [`crate::scheduler`].
+///
+/// The persistent subset of the Node is defined in [`crate::persistence::NodePersistence`]: the
+/// implementation of serialization on this type is only for debug dumps.
+#[derive(Clone, Serialize)]
pub(crate) struct Node {
pub(crate) id: NodeId,
diff --git a/control_plane/attachment_service/src/persistence.rs b/control_plane/attachment_service/src/persistence.rs
index c5829cae88..4f336093cf 100644
--- a/control_plane/attachment_service/src/persistence.rs
+++ b/control_plane/attachment_service/src/persistence.rs
@@ -6,7 +6,7 @@ use std::time::Duration;
use self::split_state::SplitState;
use camino::Utf8Path;
use camino::Utf8PathBuf;
-use control_plane::attachment_service::{NodeAvailability, NodeSchedulingPolicy};
+use control_plane::attachment_service::NodeSchedulingPolicy;
use diesel::pg::PgConnection;
use diesel::prelude::*;
use diesel::Connection;
@@ -130,24 +130,10 @@ impl Persistence {
}
/// At startup, populate the list of nodes which our shards may be placed on
- pub(crate) async fn list_nodes(&self) -> DatabaseResult> {
- let nodes: Vec = self
+ pub(crate) async fn list_nodes(&self) -> DatabaseResult> {
+ let nodes: Vec = self
.with_conn(move |conn| -> DatabaseResult<_> {
- Ok(crate::schema::nodes::table
- .load::(conn)?
- .into_iter()
- .map(|n| Node {
- id: NodeId(n.node_id as u64),
- // At startup we consider a node offline until proven otherwise.
- availability: NodeAvailability::Offline,
- scheduling: NodeSchedulingPolicy::from_str(&n.scheduling_policy)
- .expect("Bad scheduling policy in DB"),
- listen_http_addr: n.listen_http_addr,
- listen_http_port: n.listen_http_port as u16,
- listen_pg_addr: n.listen_pg_addr,
- listen_pg_port: n.listen_pg_port as u16,
- })
- .collect::>())
+ Ok(crate::schema::nodes::table.load::(conn)?)
})
.await?;
@@ -156,6 +142,31 @@ impl Persistence {
Ok(nodes)
}
+ pub(crate) async fn update_node(
+ &self,
+ input_node_id: NodeId,
+ input_scheduling: NodeSchedulingPolicy,
+ ) -> DatabaseResult<()> {
+ use crate::schema::nodes::dsl::*;
+ let updated = self
+ .with_conn(move |conn| {
+ let updated = diesel::update(nodes)
+ .filter(node_id.eq(input_node_id.0 as i64))
+ .set((scheduling_policy.eq(String::from(input_scheduling)),))
+ .execute(conn)?;
+ Ok(updated)
+ })
+ .await?;
+
+ if updated != 1 {
+ Err(DatabaseError::Logical(format!(
+ "Node {node_id:?} not found for update",
+ )))
+ } else {
+ Ok(())
+ }
+ }
+
/// At startup, load the high level state for shards, such as their config + policy. This will
/// be enriched at runtime with state discovered on pageservers.
pub(crate) async fn list_tenant_shards(&self) -> DatabaseResult> {
@@ -477,7 +488,7 @@ impl Persistence {
}
/// Parts of [`crate::tenant_state::TenantState`] that are stored durably
-#[derive(Queryable, Selectable, Insertable, Serialize, Deserialize, Clone)]
+#[derive(Queryable, Selectable, Insertable, Serialize, Deserialize, Clone, Eq, PartialEq)]
#[diesel(table_name = crate::schema::tenant_shards)]
pub(crate) struct TenantShardPersistence {
#[serde(default)]
@@ -506,7 +517,7 @@ pub(crate) struct TenantShardPersistence {
}
/// Parts of [`crate::node::Node`] that are stored durably
-#[derive(Serialize, Deserialize, Queryable, Selectable, Insertable)]
+#[derive(Serialize, Deserialize, Queryable, Selectable, Insertable, Eq, PartialEq)]
#[diesel(table_name = crate::schema::nodes)]
pub(crate) struct NodePersistence {
pub(crate) node_id: i64,
diff --git a/control_plane/attachment_service/src/reconciler.rs b/control_plane/attachment_service/src/reconciler.rs
index a4fbd80dc3..751b06f93a 100644
--- a/control_plane/attachment_service/src/reconciler.rs
+++ b/control_plane/attachment_service/src/reconciler.rs
@@ -27,7 +27,7 @@ pub(super) struct Reconciler {
pub(super) tenant_shard_id: TenantShardId,
pub(crate) shard: ShardIdentity,
pub(crate) generation: Generation,
- pub(crate) intent: IntentState,
+ pub(crate) intent: TargetState,
pub(crate) config: TenantConfig,
pub(crate) observed: ObservedState,
@@ -62,10 +62,38 @@ pub(super) struct Reconciler {
pub(crate) persistence: Arc,
}
+/// This is a snapshot of [`crate::tenant_state::IntentState`], but it does not do any
+/// reference counting for Scheduler. The IntentState is what the scheduler works with,
+/// and the TargetState is just the instruction for a particular Reconciler run.
+#[derive(Debug)]
+pub(crate) struct TargetState {
+ pub(crate) attached: Option,
+ pub(crate) secondary: Vec,
+}
+
+impl TargetState {
+ pub(crate) fn from_intent(intent: &IntentState) -> Self {
+ Self {
+ attached: *intent.get_attached(),
+ secondary: intent.get_secondary().clone(),
+ }
+ }
+
+ fn all_pageservers(&self) -> Vec {
+ let mut result = self.secondary.clone();
+ if let Some(node_id) = &self.attached {
+ result.push(*node_id);
+ }
+ result
+ }
+}
+
#[derive(thiserror::Error, Debug)]
pub(crate) enum ReconcileError {
#[error(transparent)]
Notify(#[from] NotifyError),
+ #[error("Cancelled")]
+ Cancel,
#[error(transparent)]
Other(#[from] anyhow::Error),
}
@@ -410,7 +438,7 @@ impl Reconciler {
match self.observed.locations.get(&node_id) {
Some(conf) if conf.conf.as_ref() == Some(&wanted_conf) => {
// Nothing to do
- tracing::info!("Observed configuration already correct.")
+ tracing::info!(%node_id, "Observed configuration already correct.")
}
_ => {
// In all cases other than a matching observed configuration, we will
@@ -421,7 +449,7 @@ impl Reconciler {
.increment_generation(self.tenant_shard_id, node_id)
.await?;
wanted_conf.generation = self.generation.into();
- tracing::info!("Observed configuration requires update.");
+ tracing::info!(%node_id, "Observed configuration requires update.");
self.location_config(node_id, wanted_conf, None).await?;
self.compute_notify().await?;
}
@@ -471,6 +499,9 @@ impl Reconciler {
}
for (node_id, conf) in changes {
+ if self.cancel.is_cancelled() {
+ return Err(ReconcileError::Cancel);
+ }
self.location_config(node_id, conf, None).await?;
}
diff --git a/control_plane/attachment_service/src/scheduler.rs b/control_plane/attachment_service/src/scheduler.rs
index 3b4c9e3464..7059071bee 100644
--- a/control_plane/attachment_service/src/scheduler.rs
+++ b/control_plane/attachment_service/src/scheduler.rs
@@ -1,8 +1,7 @@
-use pageserver_api::shard::TenantShardId;
-use std::collections::{BTreeMap, HashMap};
-use utils::{http::error::ApiError, id::NodeId};
-
use crate::{node::Node, tenant_state::TenantState};
+use serde::Serialize;
+use std::collections::HashMap;
+use utils::{http::error::ApiError, id::NodeId};
/// Scenarios in which we cannot find a suitable location for a tenant shard
#[derive(thiserror::Error, Debug)]
@@ -19,52 +18,203 @@ impl From for ApiError {
}
}
+#[derive(Serialize, Eq, PartialEq)]
+struct SchedulerNode {
+ /// How many shards are currently scheduled on this node, via their [`crate::tenant_state::IntentState`].
+ shard_count: usize,
+
+ /// Whether this node is currently elegible to have new shards scheduled (this is derived
+ /// from a node's availability state and scheduling policy).
+ may_schedule: bool,
+}
+
+/// This type is responsible for selecting which node is used when a tenant shard needs to choose a pageserver
+/// on which to run.
+///
+/// The type has no persistent state of its own: this is all populated at startup. The Serialize
+/// impl is only for debug dumps.
+#[derive(Serialize)]
pub(crate) struct Scheduler {
- tenant_counts: HashMap,
+ nodes: HashMap,
}
impl Scheduler {
- pub(crate) fn new(
- tenants: &BTreeMap,
- nodes: &HashMap,
- ) -> Self {
- let mut tenant_counts = HashMap::new();
- for node_id in nodes.keys() {
- tenant_counts.insert(*node_id, 0);
+ pub(crate) fn new<'a>(nodes: impl Iterator) -> Self {
+ let mut scheduler_nodes = HashMap::new();
+ for node in nodes {
+ scheduler_nodes.insert(
+ node.id,
+ SchedulerNode {
+ shard_count: 0,
+ may_schedule: node.may_schedule(),
+ },
+ );
}
- for tenant in tenants.values() {
- if let Some(ps) = tenant.intent.attached {
- let entry = tenant_counts.entry(ps).or_insert(0);
- *entry += 1;
- }
+ Self {
+ nodes: scheduler_nodes,
}
-
- for (node_id, node) in nodes {
- if !node.may_schedule() {
- tenant_counts.remove(node_id);
- }
- }
-
- Self { tenant_counts }
}
- pub(crate) fn schedule_shard(
- &mut self,
- hard_exclude: &[NodeId],
- ) -> Result {
- if self.tenant_counts.is_empty() {
+ /// For debug/support: check that our internal statistics are in sync with the state of
+ /// the nodes & tenant shards.
+ ///
+ /// If anything is inconsistent, log details and return an error.
+ pub(crate) fn consistency_check<'a>(
+ &self,
+ nodes: impl Iterator,
+ shards: impl Iterator,
+ ) -> anyhow::Result<()> {
+ let mut expect_nodes: HashMap = HashMap::new();
+ for node in nodes {
+ expect_nodes.insert(
+ node.id,
+ SchedulerNode {
+ shard_count: 0,
+ may_schedule: node.may_schedule(),
+ },
+ );
+ }
+
+ for shard in shards {
+ if let Some(node_id) = shard.intent.get_attached() {
+ match expect_nodes.get_mut(node_id) {
+ Some(node) => node.shard_count += 1,
+ None => anyhow::bail!(
+ "Tenant {} references nonexistent node {}",
+ shard.tenant_shard_id,
+ node_id
+ ),
+ }
+ }
+
+ for node_id in shard.intent.get_secondary() {
+ match expect_nodes.get_mut(node_id) {
+ Some(node) => node.shard_count += 1,
+ None => anyhow::bail!(
+ "Tenant {} references nonexistent node {}",
+ shard.tenant_shard_id,
+ node_id
+ ),
+ }
+ }
+ }
+
+ for (node_id, expect_node) in &expect_nodes {
+ let Some(self_node) = self.nodes.get(node_id) else {
+ anyhow::bail!("Node {node_id} not found in Self")
+ };
+
+ if self_node != expect_node {
+ tracing::error!("Inconsistency detected in scheduling state for node {node_id}");
+ tracing::error!("Expected state: {}", serde_json::to_string(expect_node)?);
+ tracing::error!("Self state: {}", serde_json::to_string(self_node)?);
+
+ anyhow::bail!("Inconsistent state on {node_id}");
+ }
+ }
+
+ if expect_nodes.len() != self.nodes.len() {
+ // We just checked that all the expected nodes are present. If the lengths don't match,
+ // it means that we have nodes in Self that are unexpected.
+ for node_id in self.nodes.keys() {
+ if !expect_nodes.contains_key(node_id) {
+ anyhow::bail!("Node {node_id} found in Self but not in expected nodes");
+ }
+ }
+ }
+
+ Ok(())
+ }
+
+ /// Increment the reference count of a node. This reference count is used to guide scheduling
+ /// decisions, not for memory management: it represents one tenant shard whose IntentState targets
+ /// this node.
+ ///
+ /// It is an error to call this for a node that is not known to the scheduler (i.e. passed into
+ /// [`Self::new`] or [`Self::node_upsert`])
+ pub(crate) fn node_inc_ref(&mut self, node_id: NodeId) {
+ let Some(node) = self.nodes.get_mut(&node_id) else {
+ tracing::error!("Scheduler missing node {node_id}");
+ debug_assert!(false);
+ return;
+ };
+
+ node.shard_count += 1;
+ }
+
+ /// Decrement a node's reference count. Inverse of [`Self::node_inc_ref`].
+ pub(crate) fn node_dec_ref(&mut self, node_id: NodeId) {
+ let Some(node) = self.nodes.get_mut(&node_id) else {
+ debug_assert!(false);
+ tracing::error!("Scheduler missing node {node_id}");
+ return;
+ };
+
+ node.shard_count -= 1;
+ }
+
+ pub(crate) fn node_upsert(&mut self, node: &Node) {
+ use std::collections::hash_map::Entry::*;
+ match self.nodes.entry(node.id) {
+ Occupied(mut entry) => {
+ entry.get_mut().may_schedule = node.may_schedule();
+ }
+ Vacant(entry) => {
+ entry.insert(SchedulerNode {
+ shard_count: 0,
+ may_schedule: node.may_schedule(),
+ });
+ }
+ }
+ }
+
+ pub(crate) fn node_remove(&mut self, node_id: NodeId) {
+ if self.nodes.remove(&node_id).is_none() {
+ tracing::warn!(node_id=%node_id, "Removed non-existent node from scheduler");
+ }
+ }
+
+ /// Where we have several nodes to choose from, for example when picking a secondary location
+ /// to promote to an attached location, this method may be used to pick the best choice based
+ /// on the scheduler's knowledge of utilization and availability.
+ ///
+ /// If the input is empty, or all the nodes are not elegible for scheduling, return None: the
+ /// caller can pick a node some other way.
+ pub(crate) fn node_preferred(&self, nodes: &[NodeId]) -> Option {
+ if nodes.is_empty() {
+ return None;
+ }
+
+ let node = nodes
+ .iter()
+ .map(|node_id| {
+ let may_schedule = self
+ .nodes
+ .get(node_id)
+ .map(|n| n.may_schedule)
+ .unwrap_or(false);
+ (*node_id, may_schedule)
+ })
+ .max_by_key(|(_n, may_schedule)| *may_schedule);
+
+ // If even the preferred node has may_schedule==false, return None
+ node.and_then(|(node_id, may_schedule)| if may_schedule { Some(node_id) } else { None })
+ }
+
+ pub(crate) fn schedule_shard(&self, hard_exclude: &[NodeId]) -> Result {
+ if self.nodes.is_empty() {
return Err(ScheduleError::NoPageservers);
}
let mut tenant_counts: Vec<(NodeId, usize)> = self
- .tenant_counts
+ .nodes
.iter()
.filter_map(|(k, v)| {
- if hard_exclude.contains(k) {
+ if hard_exclude.contains(k) || !v.may_schedule {
None
} else {
- Some((*k, *v))
+ Some((*k, v.shard_count))
}
})
.collect();
@@ -73,7 +223,18 @@ impl Scheduler {
tenant_counts.sort_by_key(|i| (i.1, i.0));
if tenant_counts.is_empty() {
- // After applying constraints, no pageservers were left
+ // After applying constraints, no pageservers were left. We log some detail about
+ // the state of nodes to help understand why this happened. This is not logged as an error because
+ // it is legitimately possible for enough nodes to be Offline to prevent scheduling a shard.
+ tracing::info!("Scheduling failure, while excluding {hard_exclude:?}, node states:");
+ for (node_id, node) in &self.nodes {
+ tracing::info!(
+ "Node {node_id}: may_schedule={} shards={}",
+ node.may_schedule,
+ node.shard_count
+ );
+ }
+
return Err(ScheduleError::ImpossibleConstraint);
}
@@ -82,7 +243,89 @@ impl Scheduler {
"scheduler selected node {node_id} (elegible nodes {:?}, exclude: {hard_exclude:?})",
tenant_counts.iter().map(|i| i.0 .0).collect::>()
);
- *self.tenant_counts.get_mut(&node_id).unwrap() += 1;
+
+ // Note that we do not update shard count here to reflect the scheduling: that
+ // is IntentState's job when the scheduled location is used.
+
Ok(node_id)
}
}
+
+#[cfg(test)]
+pub(crate) mod test_utils {
+
+ use crate::node::Node;
+ use control_plane::attachment_service::{NodeAvailability, NodeSchedulingPolicy};
+ use std::collections::HashMap;
+ use utils::id::NodeId;
+ /// Test helper: synthesize the requested number of nodes, all in active state.
+ ///
+ /// Node IDs start at one.
+ pub(crate) fn make_test_nodes(n: u64) -> HashMap {
+ (1..n + 1)
+ .map(|i| {
+ (
+ NodeId(i),
+ Node {
+ id: NodeId(i),
+ availability: NodeAvailability::Active,
+ scheduling: NodeSchedulingPolicy::Active,
+ listen_http_addr: format!("httphost-{i}"),
+ listen_http_port: 80 + i as u16,
+ listen_pg_addr: format!("pghost-{i}"),
+ listen_pg_port: 5432 + i as u16,
+ },
+ )
+ })
+ .collect()
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use utils::id::NodeId;
+
+ use crate::tenant_state::IntentState;
+ #[test]
+ fn scheduler_basic() -> anyhow::Result<()> {
+ let nodes = test_utils::make_test_nodes(2);
+
+ let mut scheduler = Scheduler::new(nodes.values());
+ let mut t1_intent = IntentState::new();
+ let mut t2_intent = IntentState::new();
+
+ let scheduled = scheduler.schedule_shard(&[])?;
+ t1_intent.set_attached(&mut scheduler, Some(scheduled));
+ let scheduled = scheduler.schedule_shard(&[])?;
+ t2_intent.set_attached(&mut scheduler, Some(scheduled));
+
+ assert_eq!(scheduler.nodes.get(&NodeId(1)).unwrap().shard_count, 1);
+ assert_eq!(scheduler.nodes.get(&NodeId(2)).unwrap().shard_count, 1);
+
+ let scheduled = scheduler.schedule_shard(&t1_intent.all_pageservers())?;
+ t1_intent.push_secondary(&mut scheduler, scheduled);
+
+ assert_eq!(scheduler.nodes.get(&NodeId(1)).unwrap().shard_count, 1);
+ assert_eq!(scheduler.nodes.get(&NodeId(2)).unwrap().shard_count, 2);
+
+ t1_intent.clear(&mut scheduler);
+ assert_eq!(scheduler.nodes.get(&NodeId(1)).unwrap().shard_count, 0);
+ assert_eq!(scheduler.nodes.get(&NodeId(2)).unwrap().shard_count, 1);
+
+ if cfg!(debug_assertions) {
+ // Dropping an IntentState without clearing it causes a panic in debug mode,
+ // because we have failed to properly update scheduler shard counts.
+ let result = std::panic::catch_unwind(move || {
+ drop(t2_intent);
+ });
+ assert!(result.is_err());
+ } else {
+ t2_intent.clear(&mut scheduler);
+ assert_eq!(scheduler.nodes.get(&NodeId(1)).unwrap().shard_count, 0);
+ assert_eq!(scheduler.nodes.get(&NodeId(2)).unwrap().shard_count, 0);
+ }
+
+ Ok(())
+ }
+}
diff --git a/control_plane/attachment_service/src/service.rs b/control_plane/attachment_service/src/service.rs
index 149cb7f2ba..8a80d0c746 100644
--- a/control_plane/attachment_service/src/service.rs
+++ b/control_plane/attachment_service/src/service.rs
@@ -1,4 +1,5 @@
use std::{
+ borrow::Cow,
cmp::Ordering,
collections::{BTreeMap, HashMap, HashSet},
str::FromStr,
@@ -6,6 +7,7 @@ use std::{
time::{Duration, Instant},
};
+use anyhow::Context;
use control_plane::attachment_service::{
AttachHookRequest, AttachHookResponse, InspectRequest, InspectResponse, NodeAvailability,
NodeConfigureRequest, NodeRegisterRequest, NodeSchedulingPolicy, TenantCreateResponse,
@@ -13,18 +15,18 @@ use control_plane::attachment_service::{
TenantShardMigrateRequest, TenantShardMigrateResponse,
};
use diesel::result::DatabaseErrorKind;
-use futures::StreamExt;
+use futures::{stream::FuturesUnordered, StreamExt};
use hyper::StatusCode;
use pageserver_api::{
control_api::{
ReAttachRequest, ReAttachResponse, ReAttachResponseTenant, ValidateRequest,
ValidateResponse, ValidateResponseTenant,
},
- models,
models::{
- LocationConfig, LocationConfigMode, ShardParameters, TenantConfig, TenantCreateRequest,
- TenantLocationConfigRequest, TenantLocationConfigResponse, TenantShardLocation,
- TenantShardSplitRequest, TenantShardSplitResponse, TimelineCreateRequest, TimelineInfo,
+ self, LocationConfig, LocationConfigListResponse, LocationConfigMode, ShardParameters,
+ TenantConfig, TenantCreateRequest, TenantLocationConfigRequest,
+ TenantLocationConfigResponse, TenantShardLocation, TenantShardSplitRequest,
+ TenantShardSplitResponse, TenantTimeTravelRequest, TimelineCreateRequest, TimelineInfo,
},
shard::{ShardCount, ShardIdentity, ShardNumber, ShardStripeSize, TenantShardId},
};
@@ -44,10 +46,7 @@ use utils::{
use crate::{
compute_hook::{self, ComputeHook},
node::Node,
- persistence::{
- split_state::SplitState, DatabaseError, NodePersistence, Persistence,
- TenantShardPersistence,
- },
+ persistence::{split_state::SplitState, DatabaseError, Persistence, TenantShardPersistence},
reconciler::attached_location_conf,
scheduler::Scheduler,
tenant_state::{
@@ -57,6 +56,11 @@ use crate::{
PlacementPolicy, Sequence,
};
+// For operations that should be quick, like attaching a new tenant
+const SHORT_RECONCILE_TIMEOUT: Duration = Duration::from_secs(5);
+
+// For operations that might be slow, like migrating a tenant with
+// some data in it.
const RECONCILE_TIMEOUT: Duration = Duration::from_secs(30);
/// How long [`Service::startup_reconcile`] is allowed to take before it should give
@@ -69,6 +73,8 @@ struct ServiceState {
nodes: Arc>,
+ scheduler: Scheduler,
+
compute_hook: Arc,
result_tx: tokio::sync::mpsc::UnboundedSender,
@@ -80,14 +86,26 @@ impl ServiceState {
result_tx: tokio::sync::mpsc::UnboundedSender,
nodes: HashMap,
tenants: BTreeMap,
+ scheduler: Scheduler,
) -> Self {
Self {
tenants,
nodes: Arc::new(nodes),
+ scheduler,
compute_hook: Arc::new(ComputeHook::new(config)),
result_tx,
}
}
+
+ fn parts_mut(
+ &mut self,
+ ) -> (
+ &mut Arc>,
+ &mut BTreeMap,
+ &mut Scheduler,
+ ) {
+ (&mut self.nodes, &mut self.tenants, &mut self.scheduler)
+ }
}
#[derive(Clone)]
@@ -155,98 +173,68 @@ impl Service {
/// Called once on startup, this function attempts to contact all pageservers to build an up-to-date
/// view of the world, and determine which pageservers are responsive.
#[instrument(skip_all)]
- async fn startup_reconcile(&self) {
+ async fn startup_reconcile(self: &Arc) {
// For all tenant shards, a vector of observed states on nodes (where None means
// indeterminate, same as in [`ObservedStateLocation`])
let mut observed = HashMap::new();
let mut nodes_online = HashSet::new();
- // TODO: issue these requests concurrently
- {
- let nodes = {
- let locked = self.inner.read().unwrap();
- locked.nodes.clone()
- };
- for node in nodes.values() {
- let http_client = reqwest::ClientBuilder::new()
- .timeout(Duration::from_secs(5))
- .build()
- .expect("Failed to construct HTTP client");
- let client = mgmt_api::Client::from_client(
- http_client,
- node.base_url(),
- self.config.jwt_token.as_deref(),
- );
+ // Startup reconciliation does I/O to other services: whether they
+ // are responsive or not, we should aim to finish within our deadline, because:
+ // - If we don't, a k8s readiness hook watching /ready will kill us.
+ // - While we're waiting for startup reconciliation, we are not fully
+ // available for end user operations like creating/deleting tenants and timelines.
+ //
+ // We set multiple deadlines to break up the time available between the phases of work: this is
+ // arbitrary, but avoids a situation where the first phase could burn our entire timeout period.
+ let start_at = Instant::now();
+ let node_scan_deadline = start_at
+ .checked_add(STARTUP_RECONCILE_TIMEOUT / 2)
+ .expect("Reconcile timeout is a modest constant");
- fn is_fatal(e: &mgmt_api::Error) -> bool {
- use mgmt_api::Error::*;
- match e {
- ReceiveBody(_) | ReceiveErrorBody(_) => false,
- ApiError(StatusCode::SERVICE_UNAVAILABLE, _)
- | ApiError(StatusCode::GATEWAY_TIMEOUT, _)
- | ApiError(StatusCode::REQUEST_TIMEOUT, _) => false,
- ApiError(_, _) => true,
- }
- }
+ let compute_notify_deadline = start_at
+ .checked_add((STARTUP_RECONCILE_TIMEOUT / 4) * 3)
+ .expect("Reconcile timeout is a modest constant");
- let list_response = backoff::retry(
- || client.list_location_config(),
- is_fatal,
- 1,
- 5,
- "Location config listing",
- &self.cancel,
- )
- .await;
- let Some(list_response) = list_response else {
- tracing::info!("Shutdown during startup_reconcile");
- return;
- };
+ // Accumulate a list of any tenant locations that ought to be detached
+ let mut cleanup = Vec::new();
- tracing::info!("Scanning shards on node {}...", node.id);
- match list_response {
- Err(e) => {
- tracing::warn!("Could not contact pageserver {} ({e})", node.id);
- // TODO: be more tolerant, do some retries, in case
- // pageserver is being restarted at the same time as we are
- }
- Ok(listing) => {
- tracing::info!(
- "Received {} shard statuses from pageserver {}, setting it to Active",
- listing.tenant_shards.len(),
- node.id
- );
- nodes_online.insert(node.id);
+ let node_listings = self.scan_node_locations(node_scan_deadline).await;
+ for (node_id, list_response) in node_listings {
+ let tenant_shards = list_response.tenant_shards;
+ tracing::info!(
+ "Received {} shard statuses from pageserver {}, setting it to Active",
+ tenant_shards.len(),
+ node_id
+ );
+ nodes_online.insert(node_id);
- for (tenant_shard_id, conf_opt) in listing.tenant_shards {
- observed.insert(tenant_shard_id, (node.id, conf_opt));
- }
- }
- }
+ for (tenant_shard_id, conf_opt) in tenant_shards {
+ observed.insert(tenant_shard_id, (node_id, conf_opt));
}
}
- let mut cleanup = Vec::new();
-
+ // List of tenants for which we will attempt to notify compute of their location at startup
let mut compute_notifications = Vec::new();
// Populate intent and observed states for all tenants, based on reported state on pageservers
- let (shard_count, nodes) = {
+ let shard_count = {
let mut locked = self.inner.write().unwrap();
+ let (nodes, tenants, scheduler) = locked.parts_mut();
// Mark nodes online if they responded to us: nodes are offline by default after a restart.
- let mut nodes = (*locked.nodes).clone();
- for (node_id, node) in nodes.iter_mut() {
+ let mut new_nodes = (**nodes).clone();
+ for (node_id, node) in new_nodes.iter_mut() {
if nodes_online.contains(node_id) {
node.availability = NodeAvailability::Active;
+ scheduler.node_upsert(node);
}
}
- locked.nodes = Arc::new(nodes);
- let nodes = locked.nodes.clone();
+ *nodes = Arc::new(new_nodes);
for (tenant_shard_id, (node_id, observed_loc)) in observed {
- let Some(tenant_state) = locked.tenants.get_mut(&tenant_shard_id) else {
+ let Some(tenant_state) = tenants.get_mut(&tenant_shard_id) else {
cleanup.push((tenant_shard_id, node_id));
continue;
};
@@ -258,10 +246,9 @@ impl Service {
}
// Populate each tenant's intent state
- let mut scheduler = Scheduler::new(&locked.tenants, &nodes);
- for (tenant_shard_id, tenant_state) in locked.tenants.iter_mut() {
+ for (tenant_shard_id, tenant_state) in tenants.iter_mut() {
tenant_state.intent_from_observed();
- if let Err(e) = tenant_state.schedule(&mut scheduler) {
+ if let Err(e) = tenant_state.schedule(scheduler) {
// Non-fatal error: we are unable to properly schedule the tenant, perhaps because
// not enough pageservers are available. The tenant may well still be available
// to clients.
@@ -276,18 +263,171 @@ impl Service {
}
}
- (locked.tenants.len(), nodes)
+ tenants.len()
};
// TODO: if any tenant's intent now differs from its loaded generation_pageserver, we should clear that
// generation_pageserver in the database.
- // Clean up any tenants that were found on pageservers but are not known to us.
+ // Emit compute hook notifications for all tenants which are already stably attached. Other tenants
+ // will emit compute hook notifications when they reconcile.
+ //
+ // Ordering: we must complete these notification attempts before doing any other reconciliation for the
+ // tenants named here, because otherwise our calls to notify() might race with more recent values
+ // generated by reconciliation.
+ let notify_failures = self
+ .compute_notify_many(compute_notifications, compute_notify_deadline)
+ .await;
+
+ // Compute notify is fallible. If it fails here, do not delay overall startup: set the
+ // flag on these shards that they have a pending notification.
+ // Update tenant state for any that failed to do their initial compute notify, so that they'll retry later.
+ {
+ let mut locked = self.inner.write().unwrap();
+ for tenant_shard_id in notify_failures.into_iter() {
+ if let Some(shard) = locked.tenants.get_mut(&tenant_shard_id) {
+ shard.pending_compute_notification = true;
+ }
+ }
+ }
+
+ // Finally, now that the service is up and running, launch reconcile operations for any tenants
+ // which require it: under normal circumstances this should only include tenants that were in some
+ // transient state before we restarted, or any tenants whose compute hooks failed above.
+ let reconcile_tasks = self.reconcile_all();
+ // We will not wait for these reconciliation tasks to run here: we're now done with startup and
+ // normal operations may proceed.
+
+ // Clean up any tenants that were found on pageservers but are not known to us. Do this in the
+ // background because it does not need to complete in order to proceed with other work.
+ if !cleanup.is_empty() {
+ tracing::info!("Cleaning up {} locations in the background", cleanup.len());
+ tokio::task::spawn({
+ let cleanup_self = self.clone();
+ async move { cleanup_self.cleanup_locations(cleanup).await }
+ });
+ }
+
+ tracing::info!("Startup complete, spawned {reconcile_tasks} reconciliation tasks ({shard_count} shards total)");
+ }
+
+ /// Used during [`Self::startup_reconcile`]: issue GETs to all nodes concurrently, with a deadline.
+ ///
+ /// The result includes only nodes which responded within the deadline
+ async fn scan_node_locations(
+ &self,
+ deadline: Instant,
+ ) -> HashMap {
+ let nodes = {
+ let locked = self.inner.read().unwrap();
+ locked.nodes.clone()
+ };
+
+ let mut node_results = HashMap::new();
+
+ let mut node_list_futs = FuturesUnordered::new();
+
+ for node in nodes.values() {
+ node_list_futs.push({
+ async move {
+ let http_client = reqwest::ClientBuilder::new()
+ .timeout(Duration::from_secs(5))
+ .build()
+ .expect("Failed to construct HTTP client");
+ let client = mgmt_api::Client::from_client(
+ http_client,
+ node.base_url(),
+ self.config.jwt_token.as_deref(),
+ );
+
+ fn is_fatal(e: &mgmt_api::Error) -> bool {
+ use mgmt_api::Error::*;
+ match e {
+ ReceiveBody(_) | ReceiveErrorBody(_) => false,
+ ApiError(StatusCode::SERVICE_UNAVAILABLE, _)
+ | ApiError(StatusCode::GATEWAY_TIMEOUT, _)
+ | ApiError(StatusCode::REQUEST_TIMEOUT, _) => false,
+ ApiError(_, _) => true,
+ }
+ }
+
+ tracing::info!("Scanning shards on node {}...", node.id);
+ let description = format!("List locations on {}", node.id);
+ let response = backoff::retry(
+ || client.list_location_config(),
+ is_fatal,
+ 1,
+ 5,
+ &description,
+ &self.cancel,
+ )
+ .await;
+
+ (node.id, response)
+ }
+ });
+ }
+
+ loop {
+ let (node_id, result) = tokio::select! {
+ next = node_list_futs.next() => {
+ match next {
+ Some(result) => result,
+ None =>{
+ // We got results for all our nodes
+ break;
+ }
+
+ }
+ },
+ _ = tokio::time::sleep(deadline.duration_since(Instant::now())) => {
+ // Give up waiting for anyone who hasn't responded: we will yield the results that we have
+ tracing::info!("Reached deadline while waiting for nodes to respond to location listing requests");
+ break;
+ }
+ };
+
+ let Some(list_response) = result else {
+ tracing::info!("Shutdown during startup_reconcile");
+ break;
+ };
+
+ match list_response {
+ Err(e) => {
+ tracing::warn!("Could not scan node {} ({e})", node_id);
+ }
+ Ok(listing) => {
+ node_results.insert(node_id, listing);
+ }
+ }
+ }
+
+ node_results
+ }
+
+ /// Used during [`Self::startup_reconcile`]: detach a list of unknown-to-us tenants from pageservers.
+ ///
+ /// This is safe to run in the background, because if we don't have this TenantShardId in our map of
+ /// tenants, then it is probably something incompletely deleted before: we will not fight with any
+ /// other task trying to attach it.
+ #[instrument(skip_all)]
+ async fn cleanup_locations(&self, cleanup: Vec<(TenantShardId, NodeId)>) {
+ let nodes = self.inner.read().unwrap().nodes.clone();
+
for (tenant_shard_id, node_id) in cleanup {
// A node reported a tenant_shard_id which is unknown to us: detach it.
- let node = nodes
- .get(&node_id)
- .expect("Always exists: only known nodes are scanned");
+ let Some(node) = nodes.get(&node_id) else {
+ // This is legitimate; we run in the background and [`Self::startup_reconcile`] might have identified
+ // a location to clean up on a node that has since been removed.
+ tracing::info!(
+ "Not cleaning up location {node_id}/{tenant_shard_id}: node not found"
+ );
+ continue;
+ };
+
+ if self.cancel.is_cancelled() {
+ break;
+ }
let client = mgmt_api::Client::new(node.base_url(), self.config.jwt_token.as_deref());
match client
@@ -320,58 +460,71 @@ impl Service {
}
}
}
+ }
- // Emit compute hook notifications for all tenants which are already stably attached. Other tenants
- // will emit compute hook notifications when they reconcile.
- //
- // Ordering: we must complete these notification attempts before doing any other reconciliation for the
- // tenants named here, because otherwise our calls to notify() might race with more recent values
- // generated by reconciliation.
-
- // Compute notify is fallible. If it fails here, do not delay overall startup: set the
- // flag on these shards that they have a pending notification.
+ /// Used during [`Self::startup_reconcile`]: issue many concurrent compute notifications.
+ ///
+ /// Returns a set of any shards for which notifications where not acked within the deadline.
+ async fn compute_notify_many(
+ &self,
+ notifications: Vec<(TenantShardId, NodeId)>,
+ deadline: Instant,
+ ) -> HashSet {
let compute_hook = self.inner.read().unwrap().compute_hook.clone();
+ let attempt_shards = notifications.iter().map(|i| i.0).collect::>();
+ let mut success_shards = HashSet::new();
+
// Construct an async stream of futures to invoke the compute notify function: we do this
// in order to subsequently use .buffered() on the stream to execute with bounded parallelism.
- let stream = futures::stream::iter(compute_notifications.into_iter())
+ let mut stream = futures::stream::iter(notifications.into_iter())
.map(|(tenant_shard_id, node_id)| {
let compute_hook = compute_hook.clone();
let cancel = self.cancel.clone();
async move {
if let Err(e) = compute_hook.notify(tenant_shard_id, node_id, &cancel).await {
tracing::error!(
- tenant_shard_id=%tenant_shard_id,
- node_id=%node_id,
+ %tenant_shard_id,
+ %node_id,
"Failed to notify compute on startup for shard: {e}"
);
- Some(tenant_shard_id)
- } else {
None
+ } else {
+ Some(tenant_shard_id)
}
}
})
.buffered(compute_hook::API_CONCURRENCY);
- let notify_results = stream.collect::>().await;
- // Update tenant state for any that failed to do their initial compute notify, so that they'll retry later.
- {
- let mut locked = self.inner.write().unwrap();
- for tenant_shard_id in notify_results.into_iter().flatten() {
- if let Some(shard) = locked.tenants.get_mut(&tenant_shard_id) {
- shard.pending_compute_notification = true;
+ loop {
+ tokio::select! {
+ next = stream.next() => {
+ match next {
+ Some(Some(success_shard)) => {
+ // A notification succeeded
+ success_shards.insert(success_shard);
+ },
+ Some(None) => {
+ // A notification that failed
+ },
+ None => {
+ tracing::info!("Successfully sent all compute notifications");
+ break;
+ }
+ }
+ },
+ _ = tokio::time::sleep(deadline.duration_since(Instant::now())) => {
+ // Give up sending any that didn't succeed yet
+ tracing::info!("Reached deadline while sending compute notifications");
+ break;
}
- }
+ };
}
- // Finally, now that the service is up and running, launch reconcile operations for any tenants
- // which require it: under normal circumstances this should only include tenants that were in some
- // transient state before we restarted, or any tenants whose compute hooks failed above.
- let reconcile_tasks = self.reconcile_all();
- // We will not wait for these reconciliation tasks to run here: we're now done with startup and
- // normal operations may proceed.
-
- tracing::info!("Startup complete, spawned {reconcile_tasks} reconciliation tasks ({shard_count} shards total)");
+ attempt_shards
+ .difference(&success_shards)
+ .cloned()
+ .collect()
}
/// Long running background task that periodically wakes up and looks for shards that need
@@ -393,7 +546,56 @@ impl Service {
}
}
- #[instrument(skip_all)]
+ /// Apply the contents of a [`ReconcileResult`] to our in-memory state: if the reconciliation
+ /// was successful, this will update the observed state of the tenant such that subsequent
+ /// calls to [`TenantState::maybe_reconcile`] will do nothing.
+ #[instrument(skip_all, fields(
+ tenant_id=%result.tenant_shard_id.tenant_id, shard_id=%result.tenant_shard_id.shard_slug(),
+ sequence=%result.sequence
+ ))]
+ fn process_result(&self, result: ReconcileResult) {
+ let mut locked = self.inner.write().unwrap();
+ let Some(tenant) = locked.tenants.get_mut(&result.tenant_shard_id) else {
+ // A reconciliation result might race with removing a tenant: drop results for
+ // tenants that aren't in our map.
+ return;
+ };
+
+ // Usually generation should only be updated via this path, so the max() isn't
+ // needed, but it is used to handle out-of-band updates via. e.g. test hook.
+ tenant.generation = std::cmp::max(tenant.generation, result.generation);
+
+ // If the reconciler signals that it failed to notify compute, set this state on
+ // the shard so that a future [`TenantState::maybe_reconcile`] will try again.
+ tenant.pending_compute_notification = result.pending_compute_notification;
+
+ match result.result {
+ Ok(()) => {
+ for (node_id, loc) in &result.observed.locations {
+ if let Some(conf) = &loc.conf {
+ tracing::info!("Updating observed location {}: {:?}", node_id, conf);
+ } else {
+ tracing::info!("Setting observed location {} to None", node_id,)
+ }
+ }
+ tenant.observed = result.observed;
+ tenant.waiter.advance(result.sequence);
+ }
+ Err(e) => {
+ tracing::warn!("Reconcile error: {}", e);
+
+ // Ordering: populate last_error before advancing error_seq,
+ // so that waiters will see the correct error after waiting.
+ *(tenant.last_error.lock().unwrap()) = format!("{e}");
+ tenant.error_waiter.advance(result.sequence);
+
+ for (node_id, o) in result.observed.locations {
+ tenant.observed.locations.insert(node_id, o);
+ }
+ }
+ }
+ }
+
async fn process_results(
&self,
mut result_rx: tokio::sync::mpsc::UnboundedReceiver,
@@ -412,55 +614,7 @@ impl Service {
}
};
- tracing::info!(
- "Reconcile result for sequence {}, ok={}",
- result.sequence,
- result.result.is_ok()
- );
- let mut locked = self.inner.write().unwrap();
- let Some(tenant) = locked.tenants.get_mut(&result.tenant_shard_id) else {
- // A reconciliation result might race with removing a tenant: drop results for
- // tenants that aren't in our map.
- continue;
- };
-
- // Usually generation should only be updated via this path, so the max() isn't
- // needed, but it is used to handle out-of-band updates via. e.g. test hook.
- tenant.generation = std::cmp::max(tenant.generation, result.generation);
-
- // If the reconciler signals that it failed to notify compute, set this state on
- // the shard so that a future [`TenantState::maybe_reconcile`] will try again.
- tenant.pending_compute_notification = result.pending_compute_notification;
-
- match result.result {
- Ok(()) => {
- for (node_id, loc) in &result.observed.locations {
- if let Some(conf) = &loc.conf {
- tracing::info!("Updating observed location {}: {:?}", node_id, conf);
- } else {
- tracing::info!("Setting observed location {} to None", node_id,)
- }
- }
- tenant.observed = result.observed;
- tenant.waiter.advance(result.sequence);
- }
- Err(e) => {
- tracing::warn!(
- "Reconcile error on tenant {}: {}",
- tenant.tenant_shard_id,
- e
- );
-
- // Ordering: populate last_error before advancing error_seq,
- // so that waiters will see the correct error after waiting.
- *(tenant.last_error.lock().unwrap()) = format!("{e}");
- tenant.error_waiter.advance(result.sequence);
-
- for (node_id, o) in result.observed.locations {
- tenant.observed.locations.insert(node_id, o);
- }
- }
- }
+ self.process_result(result);
}
}
@@ -468,7 +622,22 @@ impl Service {
let (result_tx, result_rx) = tokio::sync::mpsc::unbounded_channel();
tracing::info!("Loading nodes from database...");
- let nodes = persistence.list_nodes().await?;
+ let nodes = persistence
+ .list_nodes()
+ .await?
+ .into_iter()
+ .map(|n| Node {
+ id: NodeId(n.node_id as u64),
+ // At startup we consider a node offline until proven otherwise.
+ availability: NodeAvailability::Offline,
+ scheduling: NodeSchedulingPolicy::from_str(&n.scheduling_policy)
+ .expect("Bad scheduling policy in DB"),
+ listen_http_addr: n.listen_http_addr,
+ listen_http_port: n.listen_http_port as u16,
+ listen_pg_addr: n.listen_pg_addr,
+ listen_pg_port: n.listen_pg_port as u16,
+ })
+ .collect::>();
let nodes: HashMap = nodes.into_iter().map(|n| (n.id, n)).collect();
tracing::info!("Loaded {} nodes from database.", nodes.len());
@@ -481,6 +650,34 @@ impl Service {
let mut tenants = BTreeMap::new();
+ let mut scheduler = Scheduler::new(nodes.values());
+
+ #[cfg(feature = "testing")]
+ {
+ // Hack: insert scheduler state for all nodes referenced by shards, as compatibility
+ // tests only store the shards, not the nodes. The nodes will be loaded shortly
+ // after when pageservers start up and register.
+ let mut node_ids = HashSet::new();
+ for tsp in &tenant_shard_persistence {
+ if tsp.generation_pageserver != i64::MAX {
+ node_ids.insert(tsp.generation_pageserver);
+ }
+ }
+ for node_id in node_ids {
+ tracing::info!("Creating node {} in scheduler for tests", node_id);
+ let node = Node {
+ id: NodeId(node_id as u64),
+ availability: NodeAvailability::Active,
+ scheduling: NodeSchedulingPolicy::Active,
+ listen_http_addr: "".to_string(),
+ listen_http_port: 123,
+ listen_pg_addr: "".to_string(),
+ listen_pg_port: 123,
+ };
+
+ scheduler.node_upsert(&node);
+ }
+ }
for tsp in tenant_shard_persistence {
let tenant_shard_id = TenantShardId {
tenant_id: TenantId::from_str(tsp.tenant_id.as_str())?,
@@ -501,7 +698,10 @@ impl Service {
// it with what we can infer: the node for which a generation was most recently issued.
let mut intent = IntentState::new();
if tsp.generation_pageserver != i64::MAX {
- intent.attached = Some(NodeId(tsp.generation_pageserver as u64))
+ intent.set_attached(
+ &mut scheduler,
+ Some(NodeId(tsp.generation_pageserver as u64)),
+ );
}
let new_tenant = TenantState {
@@ -532,6 +732,7 @@ impl Service {
result_tx,
nodes,
tenants,
+ scheduler,
))),
config,
persistence,
@@ -636,8 +837,9 @@ impl Service {
};
let mut locked = self.inner.write().unwrap();
- let tenant_state = locked
- .tenants
+ let (_nodes, tenants, scheduler) = locked.parts_mut();
+
+ let tenant_state = tenants
.get_mut(&attach_req.tenant_shard_id)
.expect("Checked for existence above");
@@ -657,7 +859,7 @@ impl Service {
generation = ?tenant_state.generation,
"issuing",
);
- } else if let Some(ps_id) = tenant_state.intent.attached {
+ } else if let Some(ps_id) = tenant_state.intent.get_attached() {
tracing::info!(
tenant_id = %attach_req.tenant_shard_id,
%ps_id,
@@ -669,7 +871,9 @@ impl Service {
tenant_id = %attach_req.tenant_shard_id,
"no-op: tenant already has no pageserver");
}
- tenant_state.intent.attached = attach_req.node_id;
+ tenant_state
+ .intent
+ .set_attached(scheduler, attach_req.node_id);
tracing::info!(
"attach_hook: tenant {} set generation {:?}, pageserver {}",
@@ -716,7 +920,7 @@ impl Service {
InspectResponse {
attachment: tenant_state.and_then(|s| {
s.intent
- .attached
+ .get_attached()
.map(|ps| (s.generation.into().unwrap(), ps))
}),
}
@@ -725,7 +929,16 @@ impl Service {
pub(crate) async fn re_attach(
&self,
reattach_req: ReAttachRequest,
- ) -> anyhow::Result {
+ ) -> Result {
+ // Take a re-attach as indication that the node is available: this is a precursor to proper
+ // heartbeating in https://github.com/neondatabase/neon/issues/6844
+ self.node_configure(NodeConfigureRequest {
+ node_id: reattach_req.node_id,
+ availability: Some(NodeAvailability::Active),
+ scheduling: None,
+ })
+ .await?;
+
// Ordering: we must persist generation number updates before making them visible in the in-memory state
let incremented_generations = self.persistence.re_attach(reattach_req.node_id).await?;
@@ -759,6 +972,15 @@ impl Service {
};
shard_state.generation = std::cmp::max(shard_state.generation, new_gen);
+ if let Some(observed) = shard_state
+ .observed
+ .locations
+ .get_mut(&reattach_req.node_id)
+ {
+ if let Some(conf) = observed.conf.as_mut() {
+ conf.generation = new_gen.into();
+ }
+ }
// TODO: cancel/restart any running reconciliation for this tenant, it might be trying
// to call location_conf API with an old generation. Wait for cancellation to complete
@@ -807,6 +1029,16 @@ impl Service {
&self,
create_req: TenantCreateRequest,
) -> Result {
+ let (response, waiters) = self.do_tenant_create(create_req).await?;
+
+ self.await_waiters(waiters, SHORT_RECONCILE_TIMEOUT).await?;
+ Ok(response)
+ }
+
+ pub(crate) async fn do_tenant_create(
+ &self,
+ create_req: TenantCreateRequest,
+ ) -> Result<(TenantCreateResponse, Vec), ApiError> {
// This service expects to handle sharding itself: it is an error to try and directly create
// a particular shard here.
let tenant_id = if !create_req.new_tenant_id.is_unsharded() {
@@ -862,16 +1094,15 @@ impl Service {
let (waiters, response_shards) = {
let mut locked = self.inner.write().unwrap();
+ let (_nodes, tenants, scheduler) = locked.parts_mut();
let mut response_shards = Vec::new();
- let mut scheduler = Scheduler::new(&locked.tenants, &locked.nodes);
-
for tenant_shard_id in create_ids {
tracing::info!("Creating shard {tenant_shard_id}...");
use std::collections::btree_map::Entry;
- match locked.tenants.entry(tenant_shard_id) {
+ match tenants.entry(tenant_shard_id) {
Entry::Occupied(mut entry) => {
tracing::info!(
"Tenant shard {tenant_shard_id} already exists while creating"
@@ -881,7 +1112,7 @@ impl Service {
// attached and secondary locations (independently) away frorm those
// pageservers also holding a shard for this tenant.
- entry.get_mut().schedule(&mut scheduler).map_err(|e| {
+ entry.get_mut().schedule(scheduler).map_err(|e| {
ApiError::Conflict(format!(
"Failed to schedule shard {tenant_shard_id}: {e}"
))
@@ -892,7 +1123,7 @@ impl Service {
node_id: entry
.get()
.intent
- .attached
+ .get_attached()
.expect("We just set pageserver if it was None"),
generation: entry.get().generation.into().unwrap(),
});
@@ -914,7 +1145,7 @@ impl Service {
}
state.config = create_req.config.clone();
- state.schedule(&mut scheduler).map_err(|e| {
+ state.schedule(scheduler).map_err(|e| {
ApiError::Conflict(format!(
"Failed to schedule shard {tenant_shard_id}: {e}"
))
@@ -924,7 +1155,7 @@ impl Service {
shard_id: tenant_shard_id,
node_id: state
.intent
- .attached
+ .get_attached()
.expect("We just set pageserver if it was None"),
generation: state.generation.into().unwrap(),
});
@@ -957,11 +1188,12 @@ impl Service {
(waiters, response_shards)
};
- self.await_waiters(waiters).await?;
-
- Ok(TenantCreateResponse {
- shards: response_shards,
- })
+ Ok((
+ TenantCreateResponse {
+ shards: response_shards,
+ },
+ waiters,
+ ))
}
/// Helper for functions that reconcile a number of shards, and would like to do a timeout-bounded
@@ -969,8 +1201,9 @@ impl Service {
async fn await_waiters(
&self,
waiters: Vec,
+ timeout: Duration,
) -> Result<(), ReconcileWaitError> {
- let deadline = Instant::now().checked_add(Duration::from_secs(30)).unwrap();
+ let deadline = Instant::now().checked_add(timeout).unwrap();
for waiter in waiters {
let timeout = deadline.duration_since(Instant::now());
waiter.wait_timeout(timeout).await?;
@@ -1002,16 +1235,11 @@ impl Service {
let mut locked = self.inner.write().unwrap();
let result_tx = locked.result_tx.clone();
let compute_hook = locked.compute_hook.clone();
- let pageservers = locked.nodes.clone();
-
- let mut scheduler = Scheduler::new(&locked.tenants, &locked.nodes);
+ let (nodes, tenants, scheduler) = locked.parts_mut();
// Maybe we have existing shards
let mut create = true;
- for (shard_id, shard) in locked
- .tenants
- .range_mut(TenantShardId::tenant_range(tenant_id))
- {
+ for (shard_id, shard) in tenants.range_mut(TenantShardId::tenant_range(tenant_id)) {
// Saw an existing shard: this is not a creation
create = false;
@@ -1035,7 +1263,7 @@ impl Service {
| LocationConfigMode::AttachedSingle
| LocationConfigMode::AttachedStale => {
// TODO: persistence for changes in policy
- if pageservers.len() > 1 {
+ if nodes.len() > 1 {
shard.policy = PlacementPolicy::Double(1)
} else {
// Convenience for dev/test: if we just have one pageserver, import
@@ -1045,11 +1273,11 @@ impl Service {
}
}
- shard.schedule(&mut scheduler)?;
+ shard.schedule(scheduler)?;
let maybe_waiter = shard.maybe_reconcile(
result_tx.clone(),
- &pageservers,
+ nodes,
&compute_hook,
&self.config,
&self.persistence,
@@ -1060,10 +1288,10 @@ impl Service {
waiters.push(waiter);
}
- if let Some(node_id) = shard.intent.attached {
+ if let Some(node_id) = shard.intent.get_attached() {
result.shards.push(TenantShardLocation {
shard_id: *shard_id,
- node_id,
+ node_id: *node_id,
})
}
}
@@ -1113,12 +1341,8 @@ impl Service {
}
};
- // TODO: if we timeout/fail on reconcile, we should still succeed this request,
- // because otherwise a broken compute hook causes a feedback loop where
- // location_config returns 500 and gets retried forever.
-
- if let Some(create_req) = maybe_create {
- let create_resp = self.tenant_create(create_req).await?;
+ let waiters = if let Some(create_req) = maybe_create {
+ let (create_resp, waiters) = self.do_tenant_create(create_req).await?;
result.shards = create_resp
.shards
.into_iter()
@@ -1127,20 +1351,115 @@ impl Service {
shard_id: s.shard_id,
})
.collect();
+ waiters
} else {
- // This was an update, wait for reconciliation
- if let Err(e) = self.await_waiters(waiters).await {
- // Do not treat a reconcile error as fatal: we have already applied any requested
- // Intent changes, and the reconcile can fail for external reasons like unavailable
- // compute notification API. In these cases, it is important that we do not
- // cause the cloud control plane to retry forever on this API.
- tracing::warn!(
- "Failed to reconcile after /location_config: {e}, returning success anyway"
- );
+ waiters
+ };
+
+ if let Err(e) = self.await_waiters(waiters, SHORT_RECONCILE_TIMEOUT).await {
+ // Do not treat a reconcile error as fatal: we have already applied any requested
+ // Intent changes, and the reconcile can fail for external reasons like unavailable
+ // compute notification API. In these cases, it is important that we do not
+ // cause the cloud control plane to retry forever on this API.
+ tracing::warn!(
+ "Failed to reconcile after /location_config: {e}, returning success anyway"
+ );
+ }
+
+ // Logging the full result is useful because it lets us cross-check what the cloud control
+ // plane's tenant_shards table should contain.
+ tracing::info!("Complete, returning {result:?}");
+
+ Ok(result)
+ }
+
+ pub(crate) async fn tenant_time_travel_remote_storage(
+ &self,
+ time_travel_req: &TenantTimeTravelRequest,
+ tenant_id: TenantId,
+ timestamp: Cow<'_, str>,
+ done_if_after: Cow<'_, str>,
+ ) -> Result<(), ApiError> {
+ let node = {
+ let locked = self.inner.read().unwrap();
+ // Just a sanity check to prevent misuse: the API expects that the tenant is fully
+ // detached everywhere, and nothing writes to S3 storage. Here, we verify that,
+ // but only at the start of the process, so it's really just to prevent operator
+ // mistakes.
+ for (shard_id, shard) in locked.tenants.range(TenantShardId::tenant_range(tenant_id)) {
+ if shard.intent.get_attached().is_some() || !shard.intent.get_secondary().is_empty()
+ {
+ return Err(ApiError::InternalServerError(anyhow::anyhow!(
+ "We want tenant to be attached in shard with tenant_shard_id={shard_id}"
+ )));
+ }
+ let maybe_attached = shard
+ .observed
+ .locations
+ .iter()
+ .filter_map(|(node_id, observed_location)| {
+ observed_location
+ .conf
+ .as_ref()
+ .map(|loc| (node_id, observed_location, loc.mode))
+ })
+ .find(|(_, _, mode)| *mode != LocationConfigMode::Detached);
+ if let Some((node_id, _observed_location, mode)) = maybe_attached {
+ return Err(ApiError::InternalServerError(anyhow::anyhow!("We observed attached={mode:?} tenant in node_id={node_id} shard with tenant_shard_id={shard_id}")));
+ }
+ }
+ let scheduler = &locked.scheduler;
+ // Right now we only perform the operation on a single node without parallelization
+ // TODO fan out the operation to multiple nodes for better performance
+ let node_id = scheduler.schedule_shard(&[])?;
+ let node = locked
+ .nodes
+ .get(&node_id)
+ .expect("Pageservers may not be deleted while lock is active");
+ node.clone()
+ };
+
+ // The shard count is encoded in the remote storage's URL, so we need to handle all historically used shard counts
+ let mut counts = time_travel_req
+ .shard_counts
+ .iter()
+ .copied()
+ .collect::>()
+ .into_iter()
+ .collect::>();
+ counts.sort_unstable();
+
+ for count in counts {
+ let shard_ids = (0..count.count())
+ .map(|i| TenantShardId {
+ tenant_id,
+ shard_number: ShardNumber(i),
+ shard_count: count,
+ })
+ .collect::>();
+ for tenant_shard_id in shard_ids {
+ let client =
+ mgmt_api::Client::new(node.base_url(), self.config.jwt_token.as_deref());
+
+ tracing::info!("Doing time travel recovery for shard {tenant_shard_id}",);
+
+ client
+ .tenant_time_travel_remote_storage(
+ tenant_shard_id,
+ ×tamp,
+ &done_if_after,
+ )
+ .await
+ .map_err(|e| {
+ ApiError::InternalServerError(anyhow::anyhow!(
+ "Error doing time travel recovery for shard {tenant_shard_id} on node {}: {e}",
+ node.id
+ ))
+ })?;
}
}
- Ok(result)
+ Ok(())
}
pub(crate) async fn tenant_delete(&self, tenant_id: TenantId) -> Result {
@@ -1154,7 +1473,7 @@ impl Service {
for (tenant_shard_id, shard) in
locked.tenants.range(TenantShardId::tenant_range(tenant_id))
{
- let node_id = shard.intent.attached.ok_or_else(|| {
+ let node_id = shard.intent.get_attached().ok_or_else(|| {
ApiError::InternalServerError(anyhow::anyhow!("Shard not scheduled"))
})?;
let node = locked
@@ -1211,9 +1530,16 @@ impl Service {
// Drop in-memory state
{
let mut locked = self.inner.write().unwrap();
- locked
- .tenants
- .retain(|tenant_shard_id, _shard| tenant_shard_id.tenant_id != tenant_id);
+ let (_nodes, tenants, scheduler) = locked.parts_mut();
+
+ // Dereference Scheduler from shards before dropping them
+ for (_tenant_shard_id, shard) in
+ tenants.range_mut(TenantShardId::tenant_range(tenant_id))
+ {
+ shard.intent.clear(scheduler);
+ }
+
+ tenants.retain(|tenant_shard_id, _shard| tenant_shard_id.tenant_id != tenant_id);
tracing::info!(
"Deleted tenant {tenant_id}, now have {} tenants",
locked.tenants.len()
@@ -1229,8 +1555,6 @@ impl Service {
tenant_id: TenantId,
mut create_req: TimelineCreateRequest,
) -> Result {
- let mut timeline_info = None;
-
tracing::info!(
"Creating timeline {}/{}",
tenant_id,
@@ -1241,14 +1565,14 @@ impl Service {
// TODO: refuse to do this if shard splitting is in progress
// (https://github.com/neondatabase/neon/issues/6676)
- let targets = {
+ let mut targets = {
let locked = self.inner.read().unwrap();
let mut targets = Vec::new();
for (tenant_shard_id, shard) in
locked.tenants.range(TenantShardId::tenant_range(tenant_id))
{
- let node_id = shard.intent.attached.ok_or_else(|| {
+ let node_id = shard.intent.get_attached().ok_or_else(|| {
ApiError::InternalServerError(anyhow::anyhow!("Shard not scheduled"))
})?;
let node = locked
@@ -1265,21 +1589,24 @@ impl Service {
return Err(ApiError::NotFound(
anyhow::anyhow!("Tenant not found").into(),
));
- }
-
- for (tenant_shard_id, node) in targets {
- // TODO: issue shard timeline creates in parallel, once the 0th is done.
-
- let client = mgmt_api::Client::new(node.base_url(), self.config.jwt_token.as_deref());
+ };
+ let shard_zero = targets.remove(0);
+ async fn create_one(
+ tenant_shard_id: TenantShardId,
+ node: Node,
+ jwt: Option,
+ create_req: TimelineCreateRequest,
+ ) -> Result {
tracing::info!(
"Creating timeline on shard {}/{}, attached to node {}",
tenant_shard_id,
create_req.new_timeline_id,
node.id
);
+ let client = mgmt_api::Client::new(node.base_url(), jwt.as_deref());
- let shard_timeline_info = client
+ client
.timeline_create(tenant_shard_id, &create_req)
.await
.map_err(|e| match e {
@@ -1292,23 +1619,66 @@ impl Service {
ApiError::InternalServerError(anyhow::anyhow!(msg))
}
_ => ApiError::Conflict(format!("Failed to create timeline: {e}")),
- })?;
-
- if timeline_info.is_none() {
- // If the caller specified an ancestor but no ancestor LSN, we are responsible for
- // propagating the LSN chosen by the first shard to the other shards: it is important
- // that all shards end up with the same ancestor_start_lsn.
- if create_req.ancestor_timeline_id.is_some()
- && create_req.ancestor_start_lsn.is_none()
- {
- create_req.ancestor_start_lsn = shard_timeline_info.ancestor_lsn;
- }
-
- // We will return the TimelineInfo from the first shard
- timeline_info = Some(shard_timeline_info);
- }
+ })
}
- Ok(timeline_info.expect("targets cannot be empty"))
+
+ // Because the caller might not provide an explicit LSN, we must do the creation first on a single shard, and then
+ // use whatever LSN that shard picked when creating on subsequent shards. We arbitrarily use shard zero as the shard
+ // that will get the first creation request, and propagate the LSN to all the >0 shards.
+ let timeline_info = create_one(
+ shard_zero.0,
+ shard_zero.1,
+ self.config.jwt_token.clone(),
+ create_req.clone(),
+ )
+ .await?;
+
+ // Propagate the LSN that shard zero picked, if caller didn't provide one
+ if create_req.ancestor_timeline_id.is_some() && create_req.ancestor_start_lsn.is_none() {
+ create_req.ancestor_start_lsn = timeline_info.ancestor_lsn;
+ }
+
+ // Create timeline on remaining shards with number >0
+ if !targets.is_empty() {
+ // If we had multiple shards, issue requests for the remainder now.
+ let jwt = self.config.jwt_token.clone();
+ self.tenant_for_shards(targets, |tenant_shard_id: TenantShardId, node: Node| {
+ let create_req = create_req.clone();
+ Box::pin(create_one(tenant_shard_id, node, jwt.clone(), create_req))
+ })
+ .await?;
+ }
+
+ Ok(timeline_info)
+ }
+
+ /// Helper for concurrently calling a pageserver API on a number of shards, such as timeline creation.
+ ///
+ /// On success, the returned vector contains exactly the same number of elements as the input `locations`.
+ async fn tenant_for_shards(
+ &self,
+ locations: Vec<(TenantShardId, Node)>,
+ mut req_fn: F,
+ ) -> Result, ApiError>
+ where
+ F: FnMut(
+ TenantShardId,
+ Node,
+ )
+ -> std::pin::Pin> + Send>>,
+ {
+ let mut futs = FuturesUnordered::new();
+ let mut results = Vec::with_capacity(locations.len());
+
+ for (tenant_shard_id, node) in locations {
+ futs.push(req_fn(tenant_shard_id, node));
+ }
+
+ while let Some(r) = futs.next().await {
+ results.push(r?);
+ }
+
+ Ok(results)
}
pub(crate) async fn tenant_timeline_delete(
@@ -1322,14 +1692,14 @@ impl Service {
// TODO: refuse to do this if shard splitting is in progress
// (https://github.com/neondatabase/neon/issues/6676)
- let targets = {
+ let mut targets = {
let locked = self.inner.read().unwrap();
let mut targets = Vec::new();
for (tenant_shard_id, shard) in
locked.tenants.range(TenantShardId::tenant_range(tenant_id))
{
- let node_id = shard.intent.attached.ok_or_else(|| {
+ let node_id = shard.intent.get_attached().ok_or_else(|| {
ApiError::InternalServerError(anyhow::anyhow!("Shard not scheduled"))
})?;
let node = locked
@@ -1347,12 +1717,14 @@ impl Service {
anyhow::anyhow!("Tenant not found").into(),
));
}
+ let shard_zero = targets.remove(0);
- // TODO: call into shards concurrently
- let mut any_pending = false;
- for (tenant_shard_id, node) in targets {
- let client = mgmt_api::Client::new(node.base_url(), self.config.jwt_token.as_deref());
-
+ async fn delete_one(
+ tenant_shard_id: TenantShardId,
+ timeline_id: TimelineId,
+ node: Node,
+ jwt: Option,
+ ) -> Result {
tracing::info!(
"Deleting timeline on shard {}/{}, attached to node {}",
tenant_shard_id,
@@ -1360,7 +1732,8 @@ impl Service {
node.id
);
- let status = client
+ let client = mgmt_api::Client::new(node.base_url(), jwt.as_deref());
+ client
.timeline_delete(tenant_shard_id, timeline_id)
.await
.map_err(|e| {
@@ -1368,18 +1741,36 @@ impl Service {
"Error deleting timeline {timeline_id} on {tenant_shard_id} on node {}: {e}",
node.id
))
- })?;
-
- if status == StatusCode::ACCEPTED {
- any_pending = true;
- }
+ })
}
- if any_pending {
- Ok(StatusCode::ACCEPTED)
- } else {
- Ok(StatusCode::NOT_FOUND)
+ let statuses = self
+ .tenant_for_shards(targets, |tenant_shard_id: TenantShardId, node: Node| {
+ Box::pin(delete_one(
+ tenant_shard_id,
+ timeline_id,
+ node,
+ self.config.jwt_token.clone(),
+ ))
+ })
+ .await?;
+
+ // If any shards >0 haven't finished deletion yet, don't start deletion on shard zero
+ if statuses.iter().any(|s| s != &StatusCode::NOT_FOUND) {
+ return Ok(StatusCode::ACCEPTED);
}
+
+ // Delete shard zero last: this is not strictly necessary, but since a caller's GET on a timeline will be routed
+ // to shard zero, it gives a more obvious behavior that a GET returns 404 once the deletion is done.
+ let shard_zero_status = delete_one(
+ shard_zero.0,
+ timeline_id,
+ shard_zero.1,
+ self.config.jwt_token.clone(),
+ )
+ .await?;
+
+ Ok(shard_zero_status)
}
/// When you need to send an HTTP request to the pageserver that holds shard0 of a tenant, this
@@ -1401,13 +1792,18 @@ impl Service {
// TODO: should use the ID last published to compute_hook, rather than the intent: the intent might
// point to somewhere we haven't attached yet.
- let Some(node_id) = shard.intent.attached else {
+ let Some(node_id) = shard.intent.get_attached() else {
+ tracing::warn!(
+ tenant_id=%tenant_shard_id.tenant_id, shard_id=%tenant_shard_id.shard_slug(),
+ "Shard not scheduled (policy {:?}), cannot generate pass-through URL",
+ shard.policy
+ );
return Err(ApiError::Conflict(
"Cannot call timeline API on non-attached tenant".to_string(),
));
};
- let Some(node) = locked.nodes.get(&node_id) else {
+ let Some(node) = locked.nodes.get(node_id) else {
// This should never happen
return Err(ApiError::InternalServerError(anyhow::anyhow!(
"Shard refers to nonexistent node"
@@ -1432,12 +1828,13 @@ impl Service {
for (tenant_shard_id, shard) in locked.tenants.range(TenantShardId::tenant_range(tenant_id))
{
- let node_id = shard
- .intent
- .attached
- .ok_or(ApiError::BadRequest(anyhow::anyhow!(
- "Cannot locate a tenant that is not attached"
- )))?;
+ let node_id =
+ shard
+ .intent
+ .get_attached()
+ .ok_or(ApiError::BadRequest(anyhow::anyhow!(
+ "Cannot locate a tenant that is not attached"
+ )))?;
let node = pageservers
.get(&node_id)
@@ -1510,106 +1907,104 @@ impl Service {
}
// Validate input, and calculate which shards we will create
- let (old_shard_count, targets, compute_hook) = {
- let locked = self.inner.read().unwrap();
-
- let pageservers = locked.nodes.clone();
-
- let mut targets = Vec::new();
-
- // In case this is a retry, count how many already-split shards we found
- let mut children_found = Vec::new();
- let mut old_shard_count = None;
-
- for (tenant_shard_id, shard) in
- locked.tenants.range(TenantShardId::tenant_range(tenant_id))
+ let (old_shard_count, targets, compute_hook) =
{
- match shard.shard.count.count().cmp(&split_req.new_shard_count) {
- Ordering::Equal => {
- // Already split this
- children_found.push(*tenant_shard_id);
- continue;
- }
- Ordering::Greater => {
- return Err(ApiError::BadRequest(anyhow::anyhow!(
- "Requested count {} but already have shards at count {}",
- split_req.new_shard_count,
- shard.shard.count.count()
- )));
- }
- Ordering::Less => {
- // Fall through: this shard has lower count than requested,
- // is a candidate for splitting.
- }
- }
+ let locked = self.inner.read().unwrap();
- match old_shard_count {
- None => old_shard_count = Some(shard.shard.count),
- Some(old_shard_count) => {
- if old_shard_count != shard.shard.count {
- // We may hit this case if a caller asked for two splits to
- // different sizes, before the first one is complete.
- // e.g. 1->2, 2->4, where the 4 call comes while we have a mixture
- // of shard_count=1 and shard_count=2 shards in the map.
- return Err(ApiError::Conflict(
- "Cannot split, currently mid-split".to_string(),
- ));
+ let pageservers = locked.nodes.clone();
+
+ let mut targets = Vec::new();
+
+ // In case this is a retry, count how many already-split shards we found
+ let mut children_found = Vec::new();
+ let mut old_shard_count = None;
+
+ for (tenant_shard_id, shard) in
+ locked.tenants.range(TenantShardId::tenant_range(tenant_id))
+ {
+ match shard.shard.count.count().cmp(&split_req.new_shard_count) {
+ Ordering::Equal => {
+ // Already split this
+ children_found.push(*tenant_shard_id);
+ continue;
+ }
+ Ordering::Greater => {
+ return Err(ApiError::BadRequest(anyhow::anyhow!(
+ "Requested count {} but already have shards at count {}",
+ split_req.new_shard_count,
+ shard.shard.count.count()
+ )));
+ }
+ Ordering::Less => {
+ // Fall through: this shard has lower count than requested,
+ // is a candidate for splitting.
}
}
- }
- if policy.is_none() {
- policy = Some(shard.policy.clone());
- }
- if shard_ident.is_none() {
- shard_ident = Some(shard.shard);
- }
- if tenant_shard_id.shard_count.count() == split_req.new_shard_count {
- tracing::info!(
- "Tenant shard {} already has shard count {}",
- tenant_shard_id,
- split_req.new_shard_count
- );
- continue;
- }
+ match old_shard_count {
+ None => old_shard_count = Some(shard.shard.count),
+ Some(old_shard_count) => {
+ if old_shard_count != shard.shard.count {
+ // We may hit this case if a caller asked for two splits to
+ // different sizes, before the first one is complete.
+ // e.g. 1->2, 2->4, where the 4 call comes while we have a mixture
+ // of shard_count=1 and shard_count=2 shards in the map.
+ return Err(ApiError::Conflict(
+ "Cannot split, currently mid-split".to_string(),
+ ));
+ }
+ }
+ }
+ if policy.is_none() {
+ policy = Some(shard.policy.clone());
+ }
+ if shard_ident.is_none() {
+ shard_ident = Some(shard.shard);
+ }
- let node_id =
- shard
- .intent
- .attached
- .ok_or(ApiError::BadRequest(anyhow::anyhow!(
- "Cannot split a tenant that is not attached"
- )))?;
+ if tenant_shard_id.shard_count.count() == split_req.new_shard_count {
+ tracing::info!(
+ "Tenant shard {} already has shard count {}",
+ tenant_shard_id,
+ split_req.new_shard_count
+ );
+ continue;
+ }
- let node = pageservers
- .get(&node_id)
- .expect("Pageservers may not be deleted while referenced");
+ let node_id = shard.intent.get_attached().ok_or(ApiError::BadRequest(
+ anyhow::anyhow!("Cannot split a tenant that is not attached"),
+ ))?;
- // TODO: if any reconciliation is currently in progress for this shard, wait for it.
+ let node = pageservers
+ .get(&node_id)
+ .expect("Pageservers may not be deleted while referenced");
- targets.push(SplitTarget {
- parent_id: *tenant_shard_id,
- node: node.clone(),
- child_ids: tenant_shard_id.split(ShardCount::new(split_req.new_shard_count)),
- });
- }
+ // TODO: if any reconciliation is currently in progress for this shard, wait for it.
- if targets.is_empty() {
- if children_found.len() == split_req.new_shard_count as usize {
- return Ok(TenantShardSplitResponse {
- new_shards: children_found,
+ targets.push(SplitTarget {
+ parent_id: *tenant_shard_id,
+ node: node.clone(),
+ child_ids: tenant_shard_id
+ .split(ShardCount::new(split_req.new_shard_count)),
});
- } else {
- // No shards found to split, and no existing children found: the
- // tenant doesn't exist at all.
- return Err(ApiError::NotFound(
- anyhow::anyhow!("Tenant {} not found", tenant_id).into(),
- ));
}
- }
- (old_shard_count, targets, locked.compute_hook.clone())
- };
+ if targets.is_empty() {
+ if children_found.len() == split_req.new_shard_count as usize {
+ return Ok(TenantShardSplitResponse {
+ new_shards: children_found,
+ });
+ } else {
+ // No shards found to split, and no existing children found: the
+ // tenant doesn't exist at all.
+ return Err(ApiError::NotFound(
+ anyhow::anyhow!("Tenant {} not found", tenant_id).into(),
+ ));
+ }
+ }
+
+ (old_shard_count, targets, locked.compute_hook.clone())
+ };
// unwrap safety: we would have returned above if we didn't find at least one shard to split
let old_shard_count = old_shard_count.unwrap();
@@ -1751,6 +2146,7 @@ impl Service {
let mut child_locations = Vec::new();
{
let mut locked = self.inner.write().unwrap();
+ let (_nodes, tenants, scheduler) = locked.parts_mut();
for target in targets {
let SplitTarget {
parent_id,
@@ -1758,19 +2154,14 @@ impl Service {
child_ids,
} = target;
let (pageserver, generation, config) = {
- let old_state = locked
- .tenants
+ let mut old_state = tenants
.remove(&parent_id)
.expect("It was present, we just split it");
- (
- old_state.intent.attached.unwrap(),
- old_state.generation,
- old_state.config.clone(),
- )
+ let old_attached = old_state.intent.get_attached().unwrap();
+ old_state.intent.clear(scheduler);
+ (old_attached, old_state.generation, old_state.config.clone())
};
- locked.tenants.remove(&parent_id);
-
for child in child_ids {
let mut child_shard = shard_ident;
child_shard.number = child.shard_number;
@@ -1785,7 +2176,7 @@ impl Service {
);
let mut child_state = TenantState::new(child, child_shard, policy.clone());
- child_state.intent = IntentState::single(Some(pageserver));
+ child_state.intent = IntentState::single(scheduler, Some(pageserver));
child_state.observed = ObservedState {
locations: child_observed,
};
@@ -1798,7 +2189,7 @@ impl Service {
child_locations.push((child, pageserver));
- locked.tenants.insert(child, child_state);
+ tenants.insert(child, child_state);
response.new_shards.push(child);
}
}
@@ -1834,35 +2225,34 @@ impl Service {
) -> Result {
let waiter = {
let mut locked = self.inner.write().unwrap();
-
let result_tx = locked.result_tx.clone();
- let pageservers = locked.nodes.clone();
let compute_hook = locked.compute_hook.clone();
+ let (nodes, tenants, scheduler) = locked.parts_mut();
- let Some(shard) = locked.tenants.get_mut(&tenant_shard_id) else {
+ let Some(shard) = tenants.get_mut(&tenant_shard_id) else {
return Err(ApiError::NotFound(
anyhow::anyhow!("Tenant shard not found").into(),
));
};
- if shard.intent.attached == Some(migrate_req.node_id) {
+ if shard.intent.get_attached() == &Some(migrate_req.node_id) {
// No-op case: we will still proceed to wait for reconciliation in case it is
// incomplete from an earlier update to the intent.
tracing::info!("Migrating: intent is unchanged {:?}", shard.intent);
} else {
- let old_attached = shard.intent.attached;
+ let old_attached = *shard.intent.get_attached();
match shard.policy {
PlacementPolicy::Single => {
- shard.intent.secondary.clear();
+ shard.intent.clear_secondary(scheduler);
}
PlacementPolicy::Double(_n) => {
// If our new attached node was a secondary, it no longer should be.
- shard.intent.secondary.retain(|s| s != &migrate_req.node_id);
+ shard.intent.remove_secondary(scheduler, migrate_req.node_id);
// If we were already attached to something, demote that to a secondary
if let Some(old_attached) = old_attached {
- shard.intent.secondary.push(old_attached);
+ shard.intent.push_secondary(scheduler, old_attached);
}
}
PlacementPolicy::Detached => {
@@ -1871,7 +2261,9 @@ impl Service {
)))
}
}
- shard.intent.attached = Some(migrate_req.node_id);
+ shard
+ .intent
+ .set_attached(scheduler, Some(migrate_req.node_id));
tracing::info!("Migrating: new intent {:?}", shard.intent);
shard.sequence = shard.sequence.next();
@@ -1879,7 +2271,7 @@ impl Service {
shard.maybe_reconcile(
result_tx,
- &pageservers,
+ nodes,
&compute_hook,
&self.config,
&self.persistence,
@@ -1903,18 +2295,129 @@ impl Service {
self.persistence.delete_tenant(tenant_id).await?;
let mut locked = self.inner.write().unwrap();
+ let (_nodes, tenants, scheduler) = locked.parts_mut();
let mut shards = Vec::new();
- for (tenant_shard_id, _) in locked.tenants.range(TenantShardId::tenant_range(tenant_id)) {
+ for (tenant_shard_id, _) in tenants.range(TenantShardId::tenant_range(tenant_id)) {
shards.push(*tenant_shard_id);
}
- for shard in shards {
- locked.tenants.remove(&shard);
+ for shard_id in shards {
+ if let Some(mut shard) = tenants.remove(&shard_id) {
+ shard.intent.clear(scheduler);
+ }
}
Ok(())
}
+ /// For debug/support: a full JSON dump of TenantStates. Returns a response so that
+ /// we don't have to make TenantState clonable in the return path.
+ pub(crate) fn tenants_dump(&self) -> Result, ApiError> {
+ let serialized = {
+ let locked = self.inner.read().unwrap();
+ let result = locked.tenants.values().collect::>();
+ serde_json::to_string(&result).map_err(|e| ApiError::InternalServerError(e.into()))?
+ };
+
+ hyper::Response::builder()
+ .status(hyper::StatusCode::OK)
+ .header(hyper::header::CONTENT_TYPE, "application/json")
+ .body(hyper::Body::from(serialized))
+ .map_err(|e| ApiError::InternalServerError(e.into()))
+ }
+
+ /// Check the consistency of in-memory state vs. persistent state, and check that the
+ /// scheduler's statistics are up to date.
+ ///
+ /// These consistency checks expect an **idle** system. If changes are going on while
+ /// we run, then we can falsely indicate a consistency issue. This is sufficient for end-of-test
+ /// checks, but not suitable for running continuously in the background in the field.
+ pub(crate) async fn consistency_check(&self) -> Result<(), ApiError> {
+ let (mut expect_nodes, mut expect_shards) = {
+ let locked = self.inner.read().unwrap();
+
+ locked
+ .scheduler
+ .consistency_check(locked.nodes.values(), locked.tenants.values())
+ .context("Scheduler checks")
+ .map_err(ApiError::InternalServerError)?;
+
+ let expect_nodes = locked
+ .nodes
+ .values()
+ .map(|n| n.to_persistent())
+ .collect::>();
+
+ let expect_shards = locked
+ .tenants
+ .values()
+ .map(|t| t.to_persistent())
+ .collect::>();
+
+ (expect_nodes, expect_shards)
+ };
+
+ let mut nodes = self.persistence.list_nodes().await?;
+ expect_nodes.sort_by_key(|n| n.node_id);
+ nodes.sort_by_key(|n| n.node_id);
+
+ if nodes != expect_nodes {
+ tracing::error!("Consistency check failed on nodes.");
+ tracing::error!(
+ "Nodes in memory: {}",
+ serde_json::to_string(&expect_nodes)
+ .map_err(|e| ApiError::InternalServerError(e.into()))?
+ );
+ tracing::error!(
+ "Nodes in database: {}",
+ serde_json::to_string(&nodes)
+ .map_err(|e| ApiError::InternalServerError(e.into()))?
+ );
+ return Err(ApiError::InternalServerError(anyhow::anyhow!(
+ "Node consistency failure"
+ )));
+ }
+
+ let mut shards = self.persistence.list_tenant_shards().await?;
+ shards.sort_by_key(|tsp| (tsp.tenant_id.clone(), tsp.shard_number, tsp.shard_count));
+ expect_shards.sort_by_key(|tsp| (tsp.tenant_id.clone(), tsp.shard_number, tsp.shard_count));
+
+ if shards != expect_shards {
+ tracing::error!("Consistency check failed on shards.");
+ tracing::error!(
+ "Shards in memory: {}",
+ serde_json::to_string(&expect_shards)
+ .map_err(|e| ApiError::InternalServerError(e.into()))?
+ );
+ tracing::error!(
+ "Shards in database: {}",
+ serde_json::to_string(&shards)
+ .map_err(|e| ApiError::InternalServerError(e.into()))?
+ );
+ return Err(ApiError::InternalServerError(anyhow::anyhow!(
+ "Shard consistency failure"
+ )));
+ }
+
+ Ok(())
+ }
+
+ /// For debug/support: a JSON dump of the [`Scheduler`]. Returns a response so that
+ /// we don't have to make TenantState clonable in the return path.
+ pub(crate) fn scheduler_dump(&self) -> Result, ApiError> {
+ let serialized = {
+ let locked = self.inner.read().unwrap();
+ serde_json::to_string(&locked.scheduler)
+ .map_err(|e| ApiError::InternalServerError(e.into()))?
+ };
+
+ hyper::Response::builder()
+ .status(hyper::StatusCode::OK)
+ .header(hyper::header::CONTENT_TYPE, "application/json")
+ .body(hyper::Body::from(serialized))
+ .map_err(|e| ApiError::InternalServerError(e.into()))
+ }
+
/// This is for debug/support only: we simply drop all state for a tenant, without
/// detaching or deleting it on pageservers. We do not try and re-schedule any
/// tenants that were on this node.
@@ -1933,19 +2436,21 @@ impl Service {
nodes.remove(&node_id);
locked.nodes = Arc::new(nodes);
+ locked.scheduler.node_remove(node_id);
+
Ok(())
}
- pub(crate) async fn node_list(&self) -> Result, ApiError> {
- // It is convenient to avoid taking the big lock and converting Node to a serializable
- // structure, by fetching from storage instead of reading in-memory state.
- let nodes = self
- .persistence
- .list_nodes()
- .await?
- .into_iter()
- .map(|n| n.to_persistent())
- .collect();
+ pub(crate) async fn node_list(&self) -> Result, ApiError> {
+ let nodes = {
+ self.inner
+ .read()
+ .unwrap()
+ .nodes
+ .values()
+ .cloned()
+ .collect::>()
+ };
Ok(nodes)
}
@@ -2004,6 +2509,7 @@ impl Service {
let mut locked = self.inner.write().unwrap();
let mut new_nodes = (*locked.nodes).clone();
+ locked.scheduler.node_upsert(&new_node);
new_nodes.insert(register_req.node_id, new_node);
locked.nodes = Arc::new(new_nodes);
@@ -2016,12 +2522,24 @@ impl Service {
Ok(())
}
- pub(crate) fn node_configure(&self, config_req: NodeConfigureRequest) -> Result<(), ApiError> {
+ pub(crate) async fn node_configure(
+ &self,
+ config_req: NodeConfigureRequest,
+ ) -> Result<(), ApiError> {
+ if let Some(scheduling) = config_req.scheduling {
+ // Scheduling is a persistent part of Node: we must write updates to the database before
+ // applying them in memory
+ self.persistence
+ .update_node(config_req.node_id, scheduling)
+ .await?;
+ }
+
let mut locked = self.inner.write().unwrap();
let result_tx = locked.result_tx.clone();
let compute_hook = locked.compute_hook.clone();
+ let (nodes, tenants, scheduler) = locked.parts_mut();
- let mut new_nodes = (*locked.nodes).clone();
+ let mut new_nodes = (**nodes).clone();
let Some(node) = new_nodes.get_mut(&config_req.node_id) else {
return Err(ApiError::NotFound(
@@ -2057,11 +2575,14 @@ impl Service {
// to wake up and start working.
}
+ // Update the scheduler, in case the elegibility of the node for new shards has changed
+ scheduler.node_upsert(node);
+
let new_nodes = Arc::new(new_nodes);
- let mut scheduler = Scheduler::new(&locked.tenants, &new_nodes);
if offline_transition {
- for (tenant_shard_id, tenant_state) in &mut locked.tenants {
+ let mut tenants_affected: usize = 0;
+ for (tenant_shard_id, tenant_state) in tenants {
if let Some(observed_loc) =
tenant_state.observed.locations.get_mut(&config_req.node_id)
{
@@ -2072,7 +2593,7 @@ impl Service {
if tenant_state.intent.notify_offline(config_req.node_id) {
tenant_state.sequence = tenant_state.sequence.next();
- match tenant_state.schedule(&mut scheduler) {
+ match tenant_state.schedule(scheduler) {
Err(e) => {
// It is possible that some tenants will become unschedulable when too many pageservers
// go offline: in this case there isn't much we can do other than make the issue observable.
@@ -2080,19 +2601,29 @@ impl Service {
tracing::warn!(%tenant_shard_id, "Scheduling error when marking pageserver {} offline: {e}", config_req.node_id);
}
Ok(()) => {
- tenant_state.maybe_reconcile(
- result_tx.clone(),
- &new_nodes,
- &compute_hook,
- &self.config,
- &self.persistence,
- &self.gate,
- &self.cancel,
- );
+ if tenant_state
+ .maybe_reconcile(
+ result_tx.clone(),
+ &new_nodes,
+ &compute_hook,
+ &self.config,
+ &self.persistence,
+ &self.gate,
+ &self.cancel,
+ )
+ .is_some()
+ {
+ tenants_affected += 1;
+ };
}
}
}
}
+ tracing::info!(
+ "Launched {} reconciler tasks for tenants affected by node {} going offline",
+ tenants_affected,
+ config_req.node_id
+ )
}
if active_transition {
@@ -2135,18 +2666,14 @@ impl Service {
let mut waiters = Vec::new();
let result_tx = locked.result_tx.clone();
let compute_hook = locked.compute_hook.clone();
- let mut scheduler = Scheduler::new(&locked.tenants, &locked.nodes);
- let pageservers = locked.nodes.clone();
+ let (nodes, tenants, scheduler) = locked.parts_mut();
- for (_tenant_shard_id, shard) in locked
- .tenants
- .range_mut(TenantShardId::tenant_range(tenant_id))
- {
- shard.schedule(&mut scheduler)?;
+ for (_tenant_shard_id, shard) in tenants.range_mut(TenantShardId::tenant_range(tenant_id)) {
+ shard.schedule(scheduler)?;
if let Some(waiter) = shard.maybe_reconcile(
result_tx.clone(),
- &pageservers,
+ nodes,
&compute_hook,
&self.config,
&self.persistence,
diff --git a/control_plane/attachment_service/src/tenant_state.rs b/control_plane/attachment_service/src/tenant_state.rs
index dd753ece3d..02f0171c29 100644
--- a/control_plane/attachment_service/src/tenant_state.rs
+++ b/control_plane/attachment_service/src/tenant_state.rs
@@ -1,10 +1,12 @@
use std::{collections::HashMap, sync::Arc, time::Duration};
+use crate::{metrics, persistence::TenantShardPersistence};
use control_plane::attachment_service::NodeAvailability;
use pageserver_api::{
models::{LocationConfig, LocationConfigMode, TenantConfig},
shard::{ShardIdentity, TenantShardId},
};
+use serde::Serialize;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::{instrument, Instrument};
@@ -19,11 +21,27 @@ use crate::{
compute_hook::ComputeHook,
node::Node,
persistence::{split_state::SplitState, Persistence},
- reconciler::{attached_location_conf, secondary_location_conf, ReconcileError, Reconciler},
+ reconciler::{
+ attached_location_conf, secondary_location_conf, ReconcileError, Reconciler, TargetState,
+ },
scheduler::{ScheduleError, Scheduler},
service, PlacementPolicy, Sequence,
};
+/// Serialization helper
+fn read_mutex_content(v: &std::sync::Mutex, serializer: S) -> Result
+where
+ S: serde::ser::Serializer,
+ T: Clone + std::fmt::Display,
+{
+ serializer.collect_str(&v.lock().unwrap())
+}
+
+/// In-memory state for a particular tenant shard.
+///
+/// This struct implement Serialize for debugging purposes, but is _not_ persisted
+/// itself: see [`crate::persistence`] for the subset of tenant shard state that is persisted.
+#[derive(Serialize)]
pub(crate) struct TenantState {
pub(crate) tenant_shard_id: TenantShardId,
@@ -58,6 +76,7 @@ pub(crate) struct TenantState {
/// If a reconcile task is currently in flight, it may be joined here (it is
/// only safe to join if either the result has been received or the reconciler's
/// cancellation token has been fired)
+ #[serde(skip)]
pub(crate) reconciler: Option,
/// If a tenant is being split, then all shards with that TenantId will have a
@@ -67,16 +86,19 @@ pub(crate) struct TenantState {
/// Optionally wait for reconciliation to complete up to a particular
/// sequence number.
+ #[serde(skip)]
pub(crate) waiter: std::sync::Arc>,
/// Indicates sequence number for which we have encountered an error reconciling. If
/// this advances ahead of [`Self::waiter`] then a reconciliation error has occurred,
/// and callers should stop waiting for `waiter` and propagate the error.
+ #[serde(skip)]
pub(crate) error_waiter: std::sync::Arc>,
/// The most recent error from a reconcile on this tenant
/// TODO: generalize to an array of recent events
/// TOOD: use a ArcSwap instead of mutex for faster reads?
+ #[serde(serialize_with = "read_mutex_content")]
pub(crate) last_error: std::sync::Arc>,
/// If we have a pending compute notification that for some reason we weren't able to send,
@@ -86,13 +108,131 @@ pub(crate) struct TenantState {
pub(crate) pending_compute_notification: bool,
}
-#[derive(Default, Clone, Debug)]
+#[derive(Default, Clone, Debug, Serialize)]
pub(crate) struct IntentState {
- pub(crate) attached: Option,
- pub(crate) secondary: Vec,
+ attached: Option,
+ secondary: Vec,
}
-#[derive(Default, Clone)]
+impl IntentState {
+ pub(crate) fn new() -> Self {
+ Self {
+ attached: None,
+ secondary: vec![],
+ }
+ }
+ pub(crate) fn single(scheduler: &mut Scheduler, node_id: Option) -> Self {
+ if let Some(node_id) = node_id {
+ scheduler.node_inc_ref(node_id);
+ }
+ Self {
+ attached: node_id,
+ secondary: vec![],
+ }
+ }
+
+ pub(crate) fn set_attached(&mut self, scheduler: &mut Scheduler, new_attached: Option) {
+ if self.attached != new_attached {
+ if let Some(old_attached) = self.attached.take() {
+ scheduler.node_dec_ref(old_attached);
+ }
+ if let Some(new_attached) = &new_attached {
+ scheduler.node_inc_ref(*new_attached);
+ }
+ self.attached = new_attached;
+ }
+ }
+
+ /// Like set_attached, but the node is from [`Self::secondary`]. This swaps the node from
+ /// secondary to attached while maintaining the scheduler's reference counts.
+ pub(crate) fn promote_attached(
+ &mut self,
+ _scheduler: &mut Scheduler,
+ promote_secondary: NodeId,
+ ) {
+ // If we call this with a node that isn't in secondary, it would cause incorrect
+ // scheduler reference counting, since we assume the node is already referenced as a secondary.
+ debug_assert!(self.secondary.contains(&promote_secondary));
+
+ // TODO: when scheduler starts tracking attached + secondary counts separately, we will
+ // need to call into it here.
+ self.secondary.retain(|n| n != &promote_secondary);
+ self.attached = Some(promote_secondary);
+ }
+
+ pub(crate) fn push_secondary(&mut self, scheduler: &mut Scheduler, new_secondary: NodeId) {
+ debug_assert!(!self.secondary.contains(&new_secondary));
+ scheduler.node_inc_ref(new_secondary);
+ self.secondary.push(new_secondary);
+ }
+
+ /// It is legal to call this with a node that is not currently a secondary: that is a no-op
+ pub(crate) fn remove_secondary(&mut self, scheduler: &mut Scheduler, node_id: NodeId) {
+ let index = self.secondary.iter().position(|n| *n == node_id);
+ if let Some(index) = index {
+ scheduler.node_dec_ref(node_id);
+ self.secondary.remove(index);
+ }
+ }
+
+ pub(crate) fn clear_secondary(&mut self, scheduler: &mut Scheduler) {
+ for secondary in self.secondary.drain(..) {
+ scheduler.node_dec_ref(secondary);
+ }
+ }
+
+ pub(crate) fn clear(&mut self, scheduler: &mut Scheduler) {
+ if let Some(old_attached) = self.attached.take() {
+ scheduler.node_dec_ref(old_attached);
+ }
+
+ self.clear_secondary(scheduler);
+ }
+
+ pub(crate) fn all_pageservers(&self) -> Vec {
+ let mut result = Vec::new();
+ if let Some(p) = self.attached {
+ result.push(p)
+ }
+
+ result.extend(self.secondary.iter().copied());
+
+ result
+ }
+
+ pub(crate) fn get_attached(&self) -> &Option {
+ &self.attached
+ }
+
+ pub(crate) fn get_secondary(&self) -> &Vec {
+ &self.secondary
+ }
+
+ /// When a node goes offline, we update intents to avoid using it
+ /// as their attached pageserver.
+ ///
+ /// Returns true if a change was made
+ pub(crate) fn notify_offline(&mut self, node_id: NodeId) -> bool {
+ if self.attached == Some(node_id) {
+ // TODO: when scheduler starts tracking attached + secondary counts separately, we will
+ // need to call into it here.
+ self.attached = None;
+ self.secondary.push(node_id);
+ true
+ } else {
+ false
+ }
+ }
+}
+
+impl Drop for IntentState {
+ fn drop(&mut self) {
+ // Must clear before dropping, to avoid leaving stale refcounts in the Scheduler
+ debug_assert!(self.attached.is_none() && self.secondary.is_empty());
+ }
+}
+
+#[derive(Default, Clone, Serialize)]
pub(crate) struct ObservedState {
pub(crate) locations: HashMap,
}
@@ -106,7 +246,7 @@ pub(crate) struct ObservedState {
/// what it is (e.g. we failed partway through configuring it)
/// * Instance exists with conf==Some: this tells us what we last successfully configured on this node,
/// and that configuration will still be present unless something external interfered.
-#[derive(Clone)]
+#[derive(Clone, Serialize)]
pub(crate) struct ObservedStateLocation {
/// If None, it means we do not know the status of this shard's location on this node, but
/// we know that we might have some state on this node.
@@ -182,46 +322,6 @@ pub(crate) struct ReconcileResult {
pub(crate) pending_compute_notification: bool,
}
-impl IntentState {
- pub(crate) fn new() -> Self {
- Self {
- attached: None,
- secondary: vec![],
- }
- }
- pub(crate) fn all_pageservers(&self) -> Vec {
- let mut result = Vec::new();
- if let Some(p) = self.attached {
- result.push(p)
- }
-
- result.extend(self.secondary.iter().copied());
-
- result
- }
-
- pub(crate) fn single(node_id: Option) -> Self {
- Self {
- attached: node_id,
- secondary: vec![],
- }
- }
-
- /// When a node goes offline, we update intents to avoid using it
- /// as their attached pageserver.
- ///
- /// Returns true if a change was made
- pub(crate) fn notify_offline(&mut self, node_id: NodeId) -> bool {
- if self.attached == Some(node_id) {
- self.attached = None;
- self.secondary.push(node_id);
- true
- } else {
- false
- }
- }
-}
-
impl ObservedState {
pub(crate) fn new() -> Self {
Self {
@@ -289,6 +389,9 @@ impl TenantState {
// All remaining observed locations generate secondary intents. This includes None
// observations, as these may well have some local content on disk that is usable (this
// is an edge case that might occur if we restarted during a migration or other change)
+ //
+ // We may leave intent.attached empty if we didn't find any attached locations: [`Self::schedule`]
+ // will take care of promoting one of these secondaries to be attached.
self.observed.locations.keys().for_each(|node_id| {
if Some(*node_id) != self.intent.attached {
self.intent.secondary.push(*node_id);
@@ -296,6 +399,33 @@ impl TenantState {
});
}
+ /// Part of [`Self::schedule`] that is used to choose exactly one node to act as the
+ /// attached pageserver for a shard.
+ ///
+ /// Returns whether we modified it, and the NodeId selected.
+ fn schedule_attached(
+ &mut self,
+ scheduler: &mut Scheduler,
+ ) -> Result<(bool, NodeId), ScheduleError> {
+ // No work to do if we already have an attached tenant
+ if let Some(node_id) = self.intent.attached {
+ return Ok((false, node_id));
+ }
+
+ if let Some(promote_secondary) = scheduler.node_preferred(&self.intent.secondary) {
+ // Promote a secondary
+ tracing::debug!("Promoted secondary {} to attached", promote_secondary);
+ self.intent.promote_attached(scheduler, promote_secondary);
+ Ok((true, promote_secondary))
+ } else {
+ // Pick a fresh node: either we had no secondaries or none were schedulable
+ let node_id = scheduler.schedule_shard(&self.intent.secondary)?;
+ tracing::debug!("Selected {} as attached", node_id);
+ self.intent.set_attached(scheduler, Some(node_id));
+ Ok((true, node_id))
+ }
+ }
+
pub(crate) fn schedule(&mut self, scheduler: &mut Scheduler) -> Result<(), ScheduleError> {
// TODO: before scheduling new nodes, check if any existing content in
// self.intent refers to pageservers that are offline, and pick other
@@ -306,36 +436,29 @@ impl TenantState {
// Build the set of pageservers already in use by this tenant, to avoid scheduling
// more work on the same pageservers we're already using.
- let mut used_pageservers = self.intent.all_pageservers();
let mut modified = false;
use PlacementPolicy::*;
match self.policy {
Single => {
// Should have exactly one attached, and zero secondaries
- if self.intent.attached.is_none() {
- let node_id = scheduler.schedule_shard(&used_pageservers)?;
- self.intent.attached = Some(node_id);
- used_pageservers.push(node_id);
- modified = true;
- }
+ let (modified_attached, _attached_node_id) = self.schedule_attached(scheduler)?;
+ modified |= modified_attached;
+
if !self.intent.secondary.is_empty() {
- self.intent.secondary.clear();
+ self.intent.clear_secondary(scheduler);
modified = true;
}
}
Double(secondary_count) => {
// Should have exactly one attached, and N secondaries
- if self.intent.attached.is_none() {
- let node_id = scheduler.schedule_shard(&used_pageservers)?;
- self.intent.attached = Some(node_id);
- used_pageservers.push(node_id);
- modified = true;
- }
+ let (modified_attached, attached_node_id) = self.schedule_attached(scheduler)?;
+ modified |= modified_attached;
+ let mut used_pageservers = vec![attached_node_id];
while self.intent.secondary.len() < secondary_count {
let node_id = scheduler.schedule_shard(&used_pageservers)?;
- self.intent.secondary.push(node_id);
+ self.intent.push_secondary(scheduler, node_id);
used_pageservers.push(node_id);
modified = true;
}
@@ -343,12 +466,12 @@ impl TenantState {
Detached => {
// Should have no attached or secondary pageservers
if self.intent.attached.is_some() {
- self.intent.attached = None;
+ self.intent.set_attached(scheduler, None);
modified = true;
}
if !self.intent.secondary.is_empty() {
- self.intent.secondary.clear();
+ self.intent.clear_secondary(scheduler);
modified = true;
}
}
@@ -414,6 +537,13 @@ impl TenantState {
}
}
+ for node_id in self.observed.locations.keys() {
+ if self.intent.attached != Some(*node_id) && !self.intent.secondary.contains(node_id) {
+ // We have observed state that isn't part of our intent: need to clean it up.
+ return true;
+ }
+ }
+
// Even if there is no pageserver work to be done, if we have a pending notification to computes,
// wake up a reconciler to send it.
if self.pending_compute_notification {
@@ -490,7 +620,7 @@ impl TenantState {
tenant_shard_id: self.tenant_shard_id,
shard: self.shard,
generation: self.generation,
- intent: self.intent.clone(),
+ intent: TargetState::from_intent(&self.intent),
config: self.config.clone(),
observed: self.observed.clone(),
pageservers: pageservers.clone(),
@@ -509,6 +639,7 @@ impl TenantState {
let reconciler_span = tracing::info_span!(parent: None, "reconciler", seq=%reconcile_seq,
tenant_id=%reconciler.tenant_shard_id.tenant_id,
shard_id=%reconciler.tenant_shard_id.shard_slug());
+ metrics::RECONCILER.spawned.inc();
let join_handle = tokio::task::spawn(
async move {
// Wait for any previous reconcile task to complete before we start
@@ -525,6 +656,10 @@ impl TenantState {
// TODO: wrap all remote API operations in cancellation check
// as well.
if reconciler.cancel.is_cancelled() {
+ metrics::RECONCILER
+ .complete
+ .with_label_values(&[metrics::ReconcilerMetrics::CANCEL])
+ .inc();
return;
}
@@ -538,6 +673,20 @@ impl TenantState {
reconciler.compute_notify().await.ok();
}
+ // Update result counter
+ match &result {
+ Ok(_) => metrics::RECONCILER
+ .complete
+ .with_label_values(&[metrics::ReconcilerMetrics::SUCCESS]),
+ Err(ReconcileError::Cancel) => metrics::RECONCILER
+ .complete
+ .with_label_values(&[metrics::ReconcilerMetrics::CANCEL]),
+ Err(_) => metrics::RECONCILER
+ .complete
+ .with_label_values(&[metrics::ReconcilerMetrics::ERROR]),
+ }
+ .inc();
+
result_tx
.send(ReconcileResult {
sequence: reconcile_seq,
@@ -580,4 +729,103 @@ impl TenantState {
debug_assert!(!self.intent.all_pageservers().contains(&node_id));
}
+
+ pub(crate) fn to_persistent(&self) -> TenantShardPersistence {
+ TenantShardPersistence {
+ tenant_id: self.tenant_shard_id.tenant_id.to_string(),
+ shard_number: self.tenant_shard_id.shard_number.0 as i32,
+ shard_count: self.tenant_shard_id.shard_count.literal() as i32,
+ shard_stripe_size: self.shard.stripe_size.0 as i32,
+ generation: self.generation.into().unwrap_or(0) as i32,
+ generation_pageserver: self
+ .intent
+ .get_attached()
+ .map(|n| n.0 as i64)
+ .unwrap_or(i64::MAX),
+
+ placement_policy: serde_json::to_string(&self.policy).unwrap(),
+ config: serde_json::to_string(&self.config).unwrap(),
+ splitting: SplitState::default(),
+ }
+ }
+}
+
+#[cfg(test)]
+pub(crate) mod tests {
+ use pageserver_api::shard::{ShardCount, ShardNumber};
+ use utils::id::TenantId;
+
+ use crate::scheduler::test_utils::make_test_nodes;
+
+ use super::*;
+
+ fn make_test_tenant_shard(policy: PlacementPolicy) -> TenantState {
+ let tenant_id = TenantId::generate();
+ let shard_number = ShardNumber(0);
+ let shard_count = ShardCount::new(1);
+
+ let tenant_shard_id = TenantShardId {
+ tenant_id,
+ shard_number,
+ shard_count,
+ };
+ TenantState::new(
+ tenant_shard_id,
+ ShardIdentity::new(
+ shard_number,
+ shard_count,
+ pageserver_api::shard::ShardStripeSize(32768),
+ )
+ .unwrap(),
+ policy,
+ )
+ }
+
+ /// Test the scheduling behaviors used when a tenant configured for HA is subject
+ /// to nodes being marked offline.
+ #[test]
+ fn tenant_ha_scheduling() -> anyhow::Result<()> {
+ // Start with three nodes. Our tenant will only use two. The third one is
+ // expected to remain unused.
+ let mut nodes = make_test_nodes(3);
+
+ let mut scheduler = Scheduler::new(nodes.values());
+
+ let mut tenant_state = make_test_tenant_shard(PlacementPolicy::Double(1));
+ tenant_state
+ .schedule(&mut scheduler)
+ .expect("we have enough nodes, scheduling should work");
+
+ // Expect to initially be schedule on to different nodes
+ assert_eq!(tenant_state.intent.secondary.len(), 1);
+ assert!(tenant_state.intent.attached.is_some());
+
+ let attached_node_id = tenant_state.intent.attached.unwrap();
+ let secondary_node_id = *tenant_state.intent.secondary.iter().last().unwrap();
+ assert_ne!(attached_node_id, secondary_node_id);
+
+ // Notifying the attached node is offline should demote it to a secondary
+ let changed = tenant_state.intent.notify_offline(attached_node_id);
+ assert!(changed);
+
+ // Update the scheduler state to indicate the node is offline
+ nodes.get_mut(&attached_node_id).unwrap().availability = NodeAvailability::Offline;
+ scheduler.node_upsert(nodes.get(&attached_node_id).unwrap());
+
+ // Scheduling the node should promote the still-available secondary node to attached
+ tenant_state
+ .schedule(&mut scheduler)
+ .expect("active nodes are available");
+ assert_eq!(tenant_state.intent.attached.unwrap(), secondary_node_id);
+
+ // The original attached node should have been retained as a secondary
+ assert_eq!(
+ *tenant_state.intent.secondary.iter().last().unwrap(),
+ attached_node_id
+ );
+
+ tenant_state.intent.clear(&mut scheduler);
+
+ Ok(())
+ }
}
diff --git a/control_plane/src/attachment_service.rs b/control_plane/src/attachment_service.rs
index 14bfda47c3..4a1d316fe7 100644
--- a/control_plane/src/attachment_service.rs
+++ b/control_plane/src/attachment_service.rs
@@ -113,7 +113,7 @@ pub struct TenantShardMigrateRequest {
pub node_id: NodeId,
}
-#[derive(Serialize, Deserialize, Clone, Copy)]
+#[derive(Serialize, Deserialize, Clone, Copy, Eq, PartialEq)]
pub enum NodeAvailability {
// Normal, happy state
Active,
@@ -137,7 +137,7 @@ impl FromStr for NodeAvailability {
/// FIXME: this is a duplicate of the type in the attachment_service crate, because the
/// type needs to be defined with diesel traits in there.
-#[derive(Serialize, Deserialize, Clone, Copy)]
+#[derive(Serialize, Deserialize, Clone, Copy, Eq, PartialEq)]
pub enum NodeSchedulingPolicy {
Active,
Filling,
diff --git a/control_plane/src/bin/neon_local.rs b/control_plane/src/bin/neon_local.rs
index a155e9ebb2..5c0d008943 100644
--- a/control_plane/src/bin/neon_local.rs
+++ b/control_plane/src/bin/neon_local.rs
@@ -652,6 +652,10 @@ async fn handle_timeline(timeline_match: &ArgMatches, env: &mut local_env::Local
let name = import_match
.get_one::("node-name")
.ok_or_else(|| anyhow!("No node name provided"))?;
+ let update_catalog = import_match
+ .get_one::("update-catalog")
+ .cloned()
+ .unwrap_or_default();
// Parse base inputs
let base_tarfile = import_match
@@ -694,6 +698,7 @@ async fn handle_timeline(timeline_match: &ArgMatches, env: &mut local_env::Local
None,
pg_version,
ComputeMode::Primary,
+ !update_catalog,
)?;
println!("Done");
}
@@ -831,6 +836,10 @@ async fn handle_endpoint(ep_match: &ArgMatches, env: &local_env::LocalEnv) -> Re
.get_one::("endpoint_id")
.map(String::to_string)
.unwrap_or_else(|| format!("ep-{branch_name}"));
+ let update_catalog = sub_args
+ .get_one::("update-catalog")
+ .cloned()
+ .unwrap_or_default();
let lsn = sub_args
.get_one::("lsn")
@@ -880,6 +889,7 @@ async fn handle_endpoint(ep_match: &ArgMatches, env: &local_env::LocalEnv) -> Re
http_port,
pg_version,
mode,
+ !update_catalog,
)?;
}
"start" => {
@@ -918,6 +928,11 @@ async fn handle_endpoint(ep_match: &ArgMatches, env: &local_env::LocalEnv) -> Re
.get(endpoint_id.as_str())
.ok_or_else(|| anyhow::anyhow!("endpoint {endpoint_id} not found"))?;
+ let create_test_user = sub_args
+ .get_one::("create-test-user")
+ .cloned()
+ .unwrap_or_default();
+
cplane.check_conflicting_endpoints(
endpoint.mode,
endpoint.tenant_id,
@@ -972,6 +987,7 @@ async fn handle_endpoint(ep_match: &ArgMatches, env: &local_env::LocalEnv) -> Re
pageservers,
remote_ext_config,
stripe_size.0 as usize,
+ create_test_user,
)
.await?;
}
@@ -1457,6 +1473,18 @@ fn cli() -> Command {
.required(false)
.default_value("1");
+ let update_catalog = Arg::new("update-catalog")
+ .value_parser(value_parser!(bool))
+ .long("update-catalog")
+ .help("If set, will set up the catalog for neon_superuser")
+ .required(false);
+
+ let create_test_user = Arg::new("create-test-user")
+ .value_parser(value_parser!(bool))
+ .long("create-test-user")
+ .help("If set, will create test user `user` and `neondb` database. Requires `update-catalog = true`")
+ .required(false);
+
Command::new("Neon CLI")
.arg_required_else_help(true)
.version(GIT_VERSION)
@@ -1517,6 +1545,7 @@ fn cli() -> Command {
.arg(Arg::new("end-lsn").long("end-lsn")
.help("Lsn the basebackup ends at"))
.arg(pg_version_arg.clone())
+ .arg(update_catalog.clone())
)
).subcommand(
Command::new("tenant")
@@ -1630,6 +1659,7 @@ fn cli() -> Command {
.required(false))
.arg(pg_version_arg.clone())
.arg(hot_standby_arg.clone())
+ .arg(update_catalog)
)
.subcommand(Command::new("start")
.about("Start postgres.\n If the endpoint doesn't exist yet, it is created.")
@@ -1637,6 +1667,7 @@ fn cli() -> Command {
.arg(endpoint_pageserver_id_arg.clone())
.arg(safekeepers_arg)
.arg(remote_ext_config_args)
+ .arg(create_test_user)
)
.subcommand(Command::new("reconfigure")
.about("Reconfigure the endpoint")
diff --git a/control_plane/src/endpoint.rs b/control_plane/src/endpoint.rs
index ce8f035dfc..de7eb797d6 100644
--- a/control_plane/src/endpoint.rs
+++ b/control_plane/src/endpoint.rs
@@ -41,11 +41,15 @@ use std::net::SocketAddr;
use std::net::TcpStream;
use std::path::PathBuf;
use std::process::Command;
+use std::str::FromStr;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{anyhow, bail, Context, Result};
+use compute_api::spec::Database;
+use compute_api::spec::PgIdent;
use compute_api::spec::RemoteExtSpec;
+use compute_api::spec::Role;
use nix::sys::signal::kill;
use nix::sys::signal::Signal;
use serde::{Deserialize, Serialize};
@@ -122,6 +126,7 @@ impl ComputeControlPlane {
http_port: Option,
pg_version: u32,
mode: ComputeMode,
+ skip_pg_catalog_updates: bool,
) -> Result> {
let pg_port = pg_port.unwrap_or_else(|| self.get_port());
let http_port = http_port.unwrap_or_else(|| self.get_port() + 1);
@@ -140,7 +145,7 @@ impl ComputeControlPlane {
// before and after start are the same. So, skip catalog updates,
// with this we basically test a case of waking up an idle compute, where
// we also skip catalog updates in the cloud.
- skip_pg_catalog_updates: true,
+ skip_pg_catalog_updates,
features: vec![],
});
@@ -155,7 +160,7 @@ impl ComputeControlPlane {
http_port,
pg_port,
pg_version,
- skip_pg_catalog_updates: true,
+ skip_pg_catalog_updates,
features: vec![],
})?,
)?;
@@ -500,6 +505,7 @@ impl Endpoint {
pageservers: Vec<(Host, u16)>,
remote_ext_config: Option<&String>,
shard_stripe_size: usize,
+ create_test_user: bool,
) -> Result<()> {
if self.status() == EndpointStatus::Running {
anyhow::bail!("The endpoint is already running");
@@ -551,8 +557,26 @@ impl Endpoint {
cluster_id: None, // project ID: not used
name: None, // project name: not used
state: None,
- roles: vec![],
- databases: vec![],
+ roles: if create_test_user {
+ vec![Role {
+ name: PgIdent::from_str("test").unwrap(),
+ encrypted_password: None,
+ options: None,
+ }]
+ } else {
+ Vec::new()
+ },
+ databases: if create_test_user {
+ vec![Database {
+ name: PgIdent::from_str("neondb").unwrap(),
+ owner: PgIdent::from_str("test").unwrap(),
+ options: None,
+ restrict_conn: false,
+ invalid: false,
+ }]
+ } else {
+ Vec::new()
+ },
settings: None,
postgresql_conf: Some(postgresql_conf),
},
@@ -566,6 +590,7 @@ impl Endpoint {
remote_extensions,
pgbouncer_settings: None,
shard_stripe_size: Some(shard_stripe_size),
+ primary_is_running: None,
};
let spec_path = self.endpoint_path().join("spec.json");
std::fs::write(spec_path, serde_json::to_string_pretty(&spec)?)?;
@@ -577,11 +602,16 @@ impl Endpoint {
.open(self.endpoint_path().join("compute.log"))?;
// Launch compute_ctl
- println!("Starting postgres node at '{}'", self.connstr());
+ let conn_str = self.connstr("cloud_admin", "postgres");
+ println!("Starting postgres node at '{}'", conn_str);
+ if create_test_user {
+ let conn_str = self.connstr("user", "neondb");
+ println!("Also at '{}'", conn_str);
+ }
let mut cmd = Command::new(self.env.neon_distrib_dir.join("compute_ctl"));
cmd.args(["--http-port", &self.http_address.port().to_string()])
.args(["--pgdata", self.pgdata().to_str().unwrap()])
- .args(["--connstr", &self.connstr()])
+ .args(["--connstr", &conn_str])
.args([
"--spec-path",
self.endpoint_path().join("spec.json").to_str().unwrap(),
@@ -785,13 +815,13 @@ impl Endpoint {
Ok(())
}
- pub fn connstr(&self) -> String {
+ pub fn connstr(&self, user: &str, db_name: &str) -> String {
format!(
"postgresql://{}@{}:{}/{}",
- "cloud_admin",
+ user,
self.pg_address.ip(),
self.pg_address.port(),
- "postgres"
+ db_name
)
}
}
diff --git a/control_plane/src/pageserver.rs b/control_plane/src/pageserver.rs
index 8dd86bad96..5909477586 100644
--- a/control_plane/src/pageserver.rs
+++ b/control_plane/src/pageserver.rs
@@ -210,6 +210,25 @@ impl PageServerNode {
update_config: bool,
register: bool,
) -> anyhow::Result<()> {
+ // Register the node with the storage controller before starting pageserver: pageserver must be registered to
+ // successfully call /re-attach and finish starting up.
+ if register {
+ let attachment_service = AttachmentService::from_env(&self.env);
+ let (pg_host, pg_port) =
+ parse_host_port(&self.conf.listen_pg_addr).expect("Unable to parse listen_pg_addr");
+ let (http_host, http_port) = parse_host_port(&self.conf.listen_http_addr)
+ .expect("Unable to parse listen_http_addr");
+ attachment_service
+ .node_register(NodeRegisterRequest {
+ node_id: self.conf.id,
+ listen_pg_addr: pg_host.to_string(),
+ listen_pg_port: pg_port.unwrap_or(5432),
+ listen_http_addr: http_host.to_string(),
+ listen_http_port: http_port.unwrap_or(80),
+ })
+ .await?;
+ }
+
// TODO: using a thread here because start_process() is not async but we need to call check_status()
let datadir = self.repo_path();
print!(
@@ -248,23 +267,6 @@ impl PageServerNode {
)
.await?;
- if register {
- let attachment_service = AttachmentService::from_env(&self.env);
- let (pg_host, pg_port) =
- parse_host_port(&self.conf.listen_pg_addr).expect("Unable to parse listen_pg_addr");
- let (http_host, http_port) = parse_host_port(&self.conf.listen_http_addr)
- .expect("Unable to parse listen_http_addr");
- attachment_service
- .node_register(NodeRegisterRequest {
- node_id: self.conf.id,
- listen_pg_addr: pg_host.to_string(),
- listen_pg_port: pg_port.unwrap_or(5432),
- listen_http_addr: http_host.to_string(),
- listen_http_port: http_port.unwrap_or(80),
- })
- .await?;
- }
-
Ok(())
}
diff --git a/libs/compute_api/src/spec.rs b/libs/compute_api/src/spec.rs
index 2f412b61a3..71ae66c45c 100644
--- a/libs/compute_api/src/spec.rs
+++ b/libs/compute_api/src/spec.rs
@@ -79,6 +79,12 @@ pub struct ComputeSpec {
// Stripe size for pageserver sharding, in pages
#[serde(default)]
pub shard_stripe_size: Option,
+
+ // When we are starting a new replica in hot standby mode,
+ // we need to know if the primary is running.
+ // This is used to determine if replica should wait for
+ // RUNNING_XACTS from primary or not.
+ pub primary_is_running: Option,
}
/// Feature flag to signal `compute_ctl` to enable certain experimental functionality.
diff --git a/libs/metrics/src/lib.rs b/libs/metrics/src/lib.rs
index 18786106d1..744fc18e61 100644
--- a/libs/metrics/src/lib.rs
+++ b/libs/metrics/src/lib.rs
@@ -201,6 +201,11 @@ impl GenericCounterPairVec