mirror of
https://github.com/neondatabase/neon.git
synced 2026-08-09 08:09:15 +00:00
b655c7030f
## Problem Sometimes we have test data in the form of S3 contents that we would like to run live in a neon_local environment. ## Summary of changes - Add a storage controller API that imports an existing tenant. Currently this is equivalent to doing a create with a high generation number, but in future this would be something smarter to probe S3 to find the shards in a tenant and find generation numbers. - Add a `neon_local` command that invokes the import API, and then inspects timelines in the newly attached tenant to create matching branches.
217 lines
6.5 KiB
Rust
217 lines
6.5 KiB
Rust
use pageserver_api::{
|
|
models::{
|
|
LocationConfig, LocationConfigListResponse, PageserverUtilization, SecondaryProgress,
|
|
TenantScanRemoteStorageResponse, TenantShardSplitRequest, TenantShardSplitResponse,
|
|
TimelineCreateRequest, TimelineInfo,
|
|
},
|
|
shard::TenantShardId,
|
|
};
|
|
use pageserver_client::mgmt_api::{Client, Result};
|
|
use reqwest::StatusCode;
|
|
use utils::id::{NodeId, TenantId, TimelineId};
|
|
|
|
/// Thin wrapper around [`pageserver_client::mgmt_api::Client`]. It allows the storage
|
|
/// controller to collect metrics in a non-intrusive manner.
|
|
#[derive(Debug, Clone)]
|
|
pub(crate) struct PageserverClient {
|
|
inner: Client,
|
|
node_id_label: String,
|
|
}
|
|
|
|
macro_rules! measured_request {
|
|
($name:literal, $method:expr, $node_id: expr, $invoke:expr) => {{
|
|
let labels = crate::metrics::PageserverRequestLabelGroup {
|
|
pageserver_id: $node_id,
|
|
path: $name,
|
|
method: $method,
|
|
};
|
|
|
|
let latency = &crate::metrics::METRICS_REGISTRY
|
|
.metrics_group
|
|
.storage_controller_pageserver_request_latency;
|
|
let _timer_guard = latency.start_timer(labels.clone());
|
|
|
|
let res = $invoke;
|
|
|
|
if res.is_err() {
|
|
let error_counters = &crate::metrics::METRICS_REGISTRY
|
|
.metrics_group
|
|
.storage_controller_pageserver_request_error;
|
|
error_counters.inc(labels)
|
|
}
|
|
|
|
res
|
|
}};
|
|
}
|
|
|
|
impl PageserverClient {
|
|
pub(crate) fn new(node_id: NodeId, mgmt_api_endpoint: String, jwt: Option<&str>) -> Self {
|
|
Self {
|
|
inner: Client::from_client(reqwest::Client::new(), mgmt_api_endpoint, jwt),
|
|
node_id_label: node_id.0.to_string(),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn from_client(
|
|
node_id: NodeId,
|
|
raw_client: reqwest::Client,
|
|
mgmt_api_endpoint: String,
|
|
jwt: Option<&str>,
|
|
) -> Self {
|
|
Self {
|
|
inner: Client::from_client(raw_client, mgmt_api_endpoint, jwt),
|
|
node_id_label: node_id.0.to_string(),
|
|
}
|
|
}
|
|
|
|
pub(crate) async fn tenant_delete(&self, tenant_shard_id: TenantShardId) -> Result<StatusCode> {
|
|
measured_request!(
|
|
"tenant",
|
|
crate::metrics::Method::Delete,
|
|
&self.node_id_label,
|
|
self.inner.tenant_delete(tenant_shard_id).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn tenant_time_travel_remote_storage(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
timestamp: &str,
|
|
done_if_after: &str,
|
|
) -> Result<()> {
|
|
measured_request!(
|
|
"tenant_time_travel_remote_storage",
|
|
crate::metrics::Method::Put,
|
|
&self.node_id_label,
|
|
self.inner
|
|
.tenant_time_travel_remote_storage(tenant_shard_id, timestamp, done_if_after)
|
|
.await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn tenant_scan_remote_storage(
|
|
&self,
|
|
tenant_id: TenantId,
|
|
) -> Result<TenantScanRemoteStorageResponse> {
|
|
measured_request!(
|
|
"tenant_scan_remote_storage",
|
|
crate::metrics::Method::Get,
|
|
&self.node_id_label,
|
|
self.inner.tenant_scan_remote_storage(tenant_id).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn tenant_secondary_download(
|
|
&self,
|
|
tenant_id: TenantShardId,
|
|
wait: Option<std::time::Duration>,
|
|
) -> Result<(StatusCode, SecondaryProgress)> {
|
|
measured_request!(
|
|
"tenant_secondary_download",
|
|
crate::metrics::Method::Post,
|
|
&self.node_id_label,
|
|
self.inner.tenant_secondary_download(tenant_id, wait).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn location_config(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
config: LocationConfig,
|
|
flush_ms: Option<std::time::Duration>,
|
|
lazy: bool,
|
|
) -> Result<()> {
|
|
measured_request!(
|
|
"location_config",
|
|
crate::metrics::Method::Put,
|
|
&self.node_id_label,
|
|
self.inner
|
|
.location_config(tenant_shard_id, config, flush_ms, lazy)
|
|
.await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn list_location_config(&self) -> Result<LocationConfigListResponse> {
|
|
measured_request!(
|
|
"location_configs",
|
|
crate::metrics::Method::Get,
|
|
&self.node_id_label,
|
|
self.inner.list_location_config().await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn get_location_config(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
) -> Result<Option<LocationConfig>> {
|
|
measured_request!(
|
|
"location_config",
|
|
crate::metrics::Method::Get,
|
|
&self.node_id_label,
|
|
self.inner.get_location_config(tenant_shard_id).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn timeline_create(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
req: &TimelineCreateRequest,
|
|
) -> Result<TimelineInfo> {
|
|
measured_request!(
|
|
"timeline",
|
|
crate::metrics::Method::Post,
|
|
&self.node_id_label,
|
|
self.inner.timeline_create(tenant_shard_id, req).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn timeline_delete(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
timeline_id: TimelineId,
|
|
) -> Result<StatusCode> {
|
|
measured_request!(
|
|
"timeline",
|
|
crate::metrics::Method::Delete,
|
|
&self.node_id_label,
|
|
self.inner
|
|
.timeline_delete(tenant_shard_id, timeline_id)
|
|
.await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn tenant_shard_split(
|
|
&self,
|
|
tenant_shard_id: TenantShardId,
|
|
req: TenantShardSplitRequest,
|
|
) -> Result<TenantShardSplitResponse> {
|
|
measured_request!(
|
|
"tenant_shard_split",
|
|
crate::metrics::Method::Put,
|
|
&self.node_id_label,
|
|
self.inner.tenant_shard_split(tenant_shard_id, req).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn timeline_list(
|
|
&self,
|
|
tenant_shard_id: &TenantShardId,
|
|
) -> Result<Vec<TimelineInfo>> {
|
|
measured_request!(
|
|
"timelines",
|
|
crate::metrics::Method::Get,
|
|
&self.node_id_label,
|
|
self.inner.timeline_list(tenant_shard_id).await
|
|
)
|
|
}
|
|
|
|
pub(crate) async fn get_utilization(&self) -> Result<PageserverUtilization> {
|
|
measured_request!(
|
|
"utilization",
|
|
crate::metrics::Method::Get,
|
|
&self.node_id_label,
|
|
self.inner.get_utilization().await
|
|
)
|
|
}
|
|
}
|