mirror of
https://github.com/deuxfleurs-org/garage.git
synced 2026-09-06 03:59:15 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| be17e25bee | |||
| 8d85301808 |
Generated
-1
@@ -1890,7 +1890,6 @@ dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
"sha2 0.10.9",
|
||||
"subtle",
|
||||
"thiserror 2.0.18",
|
||||
"tokio",
|
||||
"toml",
|
||||
|
||||
@@ -76,7 +76,6 @@ pnet_datalink = "0.35"
|
||||
rand = "0.9"
|
||||
sha1 = "0.10"
|
||||
sha2 = "0.10"
|
||||
subtle = "2.6.1"
|
||||
timeago = { version = "0.5", default-features = false }
|
||||
xxhash-rust = { version = "0.8", default-features = false, features = ["xxh3"] }
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@
|
||||
"name": "AGPL-3.0",
|
||||
"identifier": "AGPL-3.0"
|
||||
},
|
||||
"version": "v2.3.0"
|
||||
"version": "v2.4.0"
|
||||
},
|
||||
"servers": [
|
||||
{
|
||||
@@ -1243,6 +1243,36 @@
|
||||
],
|
||||
"description": "List all the buckets on the cluster with their UUID and their global and local aliases.",
|
||||
"operationId": "ListBuckets",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "details",
|
||||
"in": "query",
|
||||
"description": "Returned detailed informations in the same format as GetBucketInfo for each bucket",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "boolean"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "offset",
|
||||
"in": "query",
|
||||
"description": "Bucket ID of the first bucket to return",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "limit",
|
||||
"in": "query",
|
||||
"description": "Maximum number of buckets to return in a single call",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "integer",
|
||||
"minimum": 0
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "Returns the UUID of all the buckets and all their aliases",
|
||||
@@ -1267,6 +1297,36 @@
|
||||
],
|
||||
"description": "Returns all API access keys in the cluster.",
|
||||
"operationId": "ListKeys",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "details",
|
||||
"in": "query",
|
||||
"description": "Returned detailed informations in the same format as GetKeyInfo for each bucket",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "boolean"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "offset",
|
||||
"in": "query",
|
||||
"description": "Key ID of the first key to return",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "limit",
|
||||
"in": "query",
|
||||
"description": "Maximum number of keys to return in a single call",
|
||||
"required": false,
|
||||
"schema": {
|
||||
"type": "integer",
|
||||
"minimum": 0
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "Returns the key identifier (aka `AWS_ACCESS_KEY_ID`) and its associated, human friendly, name if any (otherwise return an empty string)",
|
||||
@@ -3200,10 +3260,20 @@
|
||||
}
|
||||
},
|
||||
"ListBucketsResponse": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/ListBucketsResponseItem"
|
||||
}
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/ListBucketsResponseItem"
|
||||
}
|
||||
},
|
||||
{
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/GetBucketInfoResponse"
|
||||
}
|
||||
}
|
||||
]
|
||||
},
|
||||
"ListBucketsResponseItem": {
|
||||
"type": "object",
|
||||
@@ -3236,10 +3306,20 @@
|
||||
}
|
||||
},
|
||||
"ListKeysResponse": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/ListKeysResponseItem"
|
||||
}
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/ListKeysResponseItem"
|
||||
}
|
||||
},
|
||||
{
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/GetKeyInfoResponse"
|
||||
}
|
||||
}
|
||||
]
|
||||
},
|
||||
"ListKeysResponseItem": {
|
||||
"type": "object",
|
||||
@@ -3321,10 +3401,35 @@
|
||||
"dbEngine"
|
||||
],
|
||||
"properties": {
|
||||
"addr": {
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
],
|
||||
"description": "Socket address used by other nodes to connect to this node for RPC"
|
||||
},
|
||||
"dataPartition": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/FreeSpaceResp",
|
||||
"description": "Total and available space on the disk partition(s) containing the data\ndirectory(ies)"
|
||||
}
|
||||
]
|
||||
},
|
||||
"dbEngine": {
|
||||
"type": "string",
|
||||
"description": "database engine used for metadata"
|
||||
},
|
||||
"draining": {
|
||||
"type": [
|
||||
"boolean",
|
||||
"null"
|
||||
],
|
||||
"description": "Whether this node is part of an older layout version and is draining data."
|
||||
},
|
||||
"garageFeatures": {
|
||||
"type": [
|
||||
"array",
|
||||
@@ -3346,9 +3451,38 @@
|
||||
],
|
||||
"description": "hostname of this node"
|
||||
},
|
||||
"isUp": {
|
||||
"type": [
|
||||
"boolean",
|
||||
"null"
|
||||
],
|
||||
"description": "Whether this node is connected in the cluster"
|
||||
},
|
||||
"metadataPartition": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/FreeSpaceResp",
|
||||
"description": "Total and available space on the disk partition containing the\nmetadata directory"
|
||||
}
|
||||
]
|
||||
},
|
||||
"nodeId": {
|
||||
"type": "string"
|
||||
},
|
||||
"role": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/NodeAssignedRole",
|
||||
"description": "Role assigned to this node in the current cluster layout"
|
||||
}
|
||||
]
|
||||
},
|
||||
"rustVersion": {
|
||||
"type": "string",
|
||||
"description": "rustc version with which this garage release was compiled"
|
||||
@@ -3684,10 +3818,35 @@
|
||||
"dbEngine"
|
||||
],
|
||||
"properties": {
|
||||
"addr": {
|
||||
"type": [
|
||||
"string",
|
||||
"null"
|
||||
],
|
||||
"description": "Socket address used by other nodes to connect to this node for RPC"
|
||||
},
|
||||
"dataPartition": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/FreeSpaceResp",
|
||||
"description": "Total and available space on the disk partition(s) containing the data\ndirectory(ies)"
|
||||
}
|
||||
]
|
||||
},
|
||||
"dbEngine": {
|
||||
"type": "string",
|
||||
"description": "database engine used for metadata"
|
||||
},
|
||||
"draining": {
|
||||
"type": [
|
||||
"boolean",
|
||||
"null"
|
||||
],
|
||||
"description": "Whether this node is part of an older layout version and is draining data."
|
||||
},
|
||||
"garageFeatures": {
|
||||
"type": [
|
||||
"array",
|
||||
@@ -3709,9 +3868,38 @@
|
||||
],
|
||||
"description": "hostname of this node"
|
||||
},
|
||||
"isUp": {
|
||||
"type": [
|
||||
"boolean",
|
||||
"null"
|
||||
],
|
||||
"description": "Whether this node is connected in the cluster"
|
||||
},
|
||||
"metadataPartition": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/FreeSpaceResp",
|
||||
"description": "Total and available space on the disk partition containing the\nmetadata directory"
|
||||
}
|
||||
]
|
||||
},
|
||||
"nodeId": {
|
||||
"type": "string"
|
||||
},
|
||||
"role": {
|
||||
"oneOf": [
|
||||
{
|
||||
"type": "null"
|
||||
},
|
||||
{
|
||||
"$ref": "#/components/schemas/NodeAssignedRole",
|
||||
"description": "Role assigned to this node in the current cluster layout"
|
||||
}
|
||||
]
|
||||
},
|
||||
"rustVersion": {
|
||||
"type": "string",
|
||||
"description": "rustc version with which this garage release was compiled"
|
||||
|
||||
@@ -133,17 +133,12 @@ Use the following command to launch the Garage server:
|
||||
garage server --single-node --default-bucket
|
||||
```
|
||||
|
||||
- the `--single-node` flag instructs Garage to automatically configure a
|
||||
single-node cluster without data replication;
|
||||
- the `--default-bucket` flag instructs Garage to create a default access key
|
||||
and a default bucket using the environment variables we defined above (it
|
||||
implies `--default-access-key`).
|
||||
The `--single-node` flag instructs Garage to automatically configure a single-node cluster without data replication.
|
||||
The `--default-bucket` flag instructs Garage to create a default access key and a default bucket using the environment variables we defined above.
|
||||
Both flags are optional and can be omitted, in which case you will have to follow manual configuration steps described below.
|
||||
|
||||
> You can refer to the [manual configuration
|
||||
> steps](#manual-configuration) if:
|
||||
>
|
||||
> - you decide to no use these optional flags;
|
||||
> - you are running an **older version of Garage (before v2.3.0)**.
|
||||
**For older versions of Garage (before v2.3.0):** automatic configuration using `--single-node` and `--default-bucket` is not available,
|
||||
you must follow the manual configuration steps.
|
||||
|
||||
Alternatively, if you cannot or do not wish to run the Garage binary directly,
|
||||
you may use Docker to run Garage in a container using the following command:
|
||||
@@ -297,7 +292,7 @@ An exhaustive list is maintained in the ["Integrations" > "Browsing tools" secti
|
||||
|
||||
|
||||
|
||||
## Manual configuration {#manual-configuration}
|
||||
## Manual configuration
|
||||
|
||||
This section provides instructions that are equivalent to using the
|
||||
`--single-node` and `--default-bucket` flags for automatic configuration. If
|
||||
|
||||
@@ -175,9 +175,6 @@ they do not exist in the configuration file:
|
||||
Garage daemon send its logs to `journald` (using the native protocol of `systemd-journald`)
|
||||
instead of printing to stderr.
|
||||
|
||||
- `NO_COLOR` (since `v2.4.0`): set this to `0` or `false` to disable
|
||||
ANSI color codes in Garage's logs.
|
||||
|
||||
The following environment variables can be used to override the corresponding
|
||||
values in the configuration file:
|
||||
|
||||
|
||||
@@ -8,11 +8,12 @@ which is an alternative storage API designed to help efficiently store
|
||||
many small values in buckets (in opposition to S3 which is more designed
|
||||
to store large blobs).
|
||||
|
||||
K2V is included in release builds since version 0.8.0. Precompiled builds
|
||||
of earlier versions including `k2v` can be found in our download page under
|
||||
"Extra builds": they can be easily identified as their tag name ends with
|
||||
`-k2v` (example: `v0.7.2-k2v`). Otherwise, when compiling Garage, the Cargo
|
||||
feature flag `k2v` must be activated.
|
||||
K2V is currently disabled at compile time in all builds, as the
|
||||
specification is still subject to changes. To build a Garage version with
|
||||
K2V, the Cargo feature flag `k2v` must be activated. Special builds with
|
||||
the `k2v` feature flag enabled can be obtained from our download page under
|
||||
"Extra builds": such builds can be identified easily as their tag name ends
|
||||
with `-k2v` (example: `v0.7.2-k2v`).
|
||||
|
||||
The specification of the K2V API can be found
|
||||
[here](https://git.deuxfleurs.fr/Deuxfleurs/garage/src/commit/f8be15c37db857e177d543de7be863692628d567/doc/drafts/k2v-spec.md).
|
||||
|
||||
@@ -18,7 +18,7 @@ fi
|
||||
|
||||
$GARAGE_BIN -c /tmp/config.1.toml bucket create eprouvette
|
||||
if [ "$GARAGE_OLDVER" = "v08" ]; then
|
||||
KEY_INFO=$($GARAGE_BIN -c /tmp/config.1.toml key new --name opérateur)
|
||||
KEY_INFO=$($GARAGE_BIN -c /tmp/config.1.toml key create opérateur)
|
||||
ACCESS_KEY=`echo $KEY_INFO|grep -Po 'GK[a-f0-9]+'`
|
||||
SECRET_KEY=`echo $KEY_INFO|grep -Po 'Secret key: [a-f0-9]+'|grep -Po '[a-f0-9]+$'`
|
||||
elif [ "$GARAGE_OLDVER" = "v1" ]; then
|
||||
|
||||
@@ -191,7 +191,7 @@ impl RequestHandler for GetCurrentAdminTokenInfoRequest {
|
||||
.admin
|
||||
.metrics_token
|
||||
.as_ref()
|
||||
.is_some_and(|s| s.eq_ct(&self.admin_token))
|
||||
.is_some_and(|s| s == &self.admin_token)
|
||||
{
|
||||
return Ok(GetCurrentAdminTokenInfoResponse(
|
||||
GetAdminTokenInfoResponse {
|
||||
@@ -210,7 +210,7 @@ impl RequestHandler for GetCurrentAdminTokenInfoRequest {
|
||||
.admin
|
||||
.admin_token
|
||||
.as_ref()
|
||||
.is_some_and(|s| s.eq_ct(&self.admin_token))
|
||||
.is_some_and(|s| s == &self.admin_token)
|
||||
{
|
||||
return Ok(GetCurrentAdminTokenInfoResponse(
|
||||
GetAdminTokenInfoResponse {
|
||||
|
||||
+36
-6
@@ -688,11 +688,26 @@ pub struct ClusterLayoutSkipDeadNodesResponse {
|
||||
|
||||
// ---- ListKeys ----
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ListKeysRequest;
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, Default, IntoParams)]
|
||||
#[into_params(parameter_in = Query)]
|
||||
pub struct ListKeysRequest {
|
||||
/// Returned detailed informations in the same format as GetKeyInfo for each bucket
|
||||
#[serde(default)]
|
||||
pub details: bool,
|
||||
/// Key ID of the first key to return
|
||||
#[serde(default)]
|
||||
pub offset: Option<String>,
|
||||
/// Maximum number of keys to return in a single call
|
||||
#[serde(default)]
|
||||
pub limit: Option<usize>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
|
||||
pub struct ListKeysResponse(pub Vec<ListKeysResponseItem>);
|
||||
#[serde(untagged)]
|
||||
pub enum ListKeysResponse {
|
||||
WithoutDetails(Vec<ListKeysResponseItem>),
|
||||
WithDetails(Vec<GetKeyInfoResponse>),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
@@ -830,11 +845,26 @@ pub struct DeleteKeyResponse;
|
||||
|
||||
// ---- ListBuckets ----
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ListBucketsRequest;
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, Default, IntoParams)]
|
||||
#[into_params(parameter_in = Query)]
|
||||
pub struct ListBucketsRequest {
|
||||
/// Returned detailed informations in the same format as GetBucketInfo for each bucket
|
||||
#[serde(default)]
|
||||
pub details: bool,
|
||||
/// Bucket ID of the first bucket to return
|
||||
#[serde(default)]
|
||||
pub offset: Option<String>,
|
||||
/// Maximum number of buckets to return in a single call
|
||||
#[serde(default)]
|
||||
pub limit: Option<usize>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
|
||||
pub struct ListBucketsResponse(pub Vec<ListBucketsResponseItem>);
|
||||
#[serde(untagged)]
|
||||
pub enum ListBucketsResponse {
|
||||
WithoutDetails(Vec<ListBucketsResponseItem>),
|
||||
WithDetails(Vec<GetBucketInfoResponse>),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
|
||||
@@ -117,14 +117,8 @@ impl AdminApiServer {
|
||||
#[cfg(feature = "metrics")] exporter: PrometheusExporter,
|
||||
) -> Arc<Self> {
|
||||
let cfg = &garage.config.admin;
|
||||
let metrics_token = cfg
|
||||
.metrics_token
|
||||
.as_ref()
|
||||
.map(|token| hash_bearer_token(token.extract_secret()));
|
||||
let admin_token = cfg
|
||||
.admin_token
|
||||
.as_ref()
|
||||
.map(|token| hash_bearer_token(token.extract_secret()));
|
||||
let metrics_token = cfg.metrics_token.as_deref().map(hash_bearer_token);
|
||||
let admin_token = cfg.admin_token.as_deref().map(hash_bearer_token);
|
||||
let metrics_require_token = cfg.metrics_require_token;
|
||||
|
||||
let endpoint = garage.system.netapp.endpoint(ADMIN_RPC_PATH.into());
|
||||
|
||||
+55
-31
@@ -3,6 +3,7 @@ use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use chrono::DateTime;
|
||||
use futures::StreamExt;
|
||||
|
||||
use garage_util::crdt::*;
|
||||
use garage_util::data::*;
|
||||
@@ -32,47 +33,70 @@ impl RequestHandler for ListBucketsRequest {
|
||||
garage: &Arc<Garage>,
|
||||
_admin: &Admin,
|
||||
) -> Result<ListBucketsResponse, Error> {
|
||||
let limit = self
|
||||
.limit
|
||||
.unwrap_or_else(|| if self.details { 1000 } else { 10_000 });
|
||||
|
||||
let offset = match self.offset {
|
||||
Some(id) => Some(parse_bucket_id(&id)?),
|
||||
None => None,
|
||||
};
|
||||
|
||||
let buckets = garage
|
||||
.bucket_table
|
||||
.get_range(
|
||||
&EmptyKey,
|
||||
None,
|
||||
offset,
|
||||
Some(DeletedFilter::NotDeleted),
|
||||
1_000_000,
|
||||
limit,
|
||||
EnumerationOrder::Forward,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let res = buckets
|
||||
.into_iter()
|
||||
.map(|b| {
|
||||
let state = b.state.as_option().unwrap();
|
||||
ListBucketsResponseItem {
|
||||
id: hex::encode(b.id),
|
||||
created: DateTime::from_timestamp_millis(state.creation_date as i64)
|
||||
.expect("invalid timestamp stored in db"),
|
||||
global_aliases: state
|
||||
.aliases
|
||||
.items()
|
||||
.iter()
|
||||
.filter(|(_, _, a)| *a)
|
||||
.map(|(n, _, _)| n.to_string())
|
||||
.collect::<Vec<_>>(),
|
||||
local_aliases: state
|
||||
.local_aliases
|
||||
.items()
|
||||
.iter()
|
||||
.filter(|(_, _, a)| *a)
|
||||
.map(|((k, n), _, _)| BucketLocalAlias {
|
||||
access_key_id: k.to_string(),
|
||||
alias: n.to_string(),
|
||||
})
|
||||
.collect::<Vec<_>>(),
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
if self.details {
|
||||
let mut stream = buckets
|
||||
.into_iter()
|
||||
.map(|b| bucket_info_results(garage, b.id))
|
||||
.collect::<futures::stream::FuturesOrdered<_>>();
|
||||
|
||||
Ok(ListBucketsResponse(res))
|
||||
let mut res = vec![];
|
||||
while let Some(next) = stream.next().await {
|
||||
res.push(next?);
|
||||
}
|
||||
|
||||
Ok(ListBucketsResponse::WithDetails(res))
|
||||
} else {
|
||||
let res = buckets
|
||||
.into_iter()
|
||||
.map(|b| {
|
||||
let state = b.state.as_option().unwrap();
|
||||
ListBucketsResponseItem {
|
||||
id: hex::encode(b.id),
|
||||
created: DateTime::from_timestamp_millis(state.creation_date as i64)
|
||||
.expect("invalid timestamp stored in db"),
|
||||
global_aliases: state
|
||||
.aliases
|
||||
.items()
|
||||
.iter()
|
||||
.filter(|(_, _, a)| *a)
|
||||
.map(|(n, _, _)| n.to_string())
|
||||
.collect::<Vec<_>>(),
|
||||
local_aliases: state
|
||||
.local_aliases
|
||||
.items()
|
||||
.iter()
|
||||
.filter(|(_, _, a)| *a)
|
||||
.map(|((k, n), _, _)| BucketLocalAlias {
|
||||
access_key_id: k.to_string(),
|
||||
alias: n.to_string(),
|
||||
})
|
||||
.collect::<Vec<_>>(),
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
Ok(ListBucketsResponse::WithoutDetails(res))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+44
-23
@@ -2,6 +2,7 @@ use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use chrono::DateTime;
|
||||
use futures::StreamExt;
|
||||
|
||||
use garage_table::*;
|
||||
use garage_util::time::now_msec;
|
||||
@@ -20,37 +21,57 @@ impl RequestHandler for ListKeysRequest {
|
||||
async fn handle(self, garage: &Arc<Garage>, _admin: &Admin) -> Result<ListKeysResponse, Error> {
|
||||
let now = now_msec();
|
||||
|
||||
let res = garage
|
||||
let limit = self
|
||||
.limit
|
||||
.unwrap_or_else(|| if self.details { 1000 } else { 10_000 });
|
||||
|
||||
let keys = garage
|
||||
.key_table
|
||||
.get_range(
|
||||
&EmptyKey,
|
||||
None,
|
||||
self.offset,
|
||||
Some(KeyFilter::Deleted(DeletedFilter::NotDeleted)),
|
||||
10000,
|
||||
limit,
|
||||
EnumerationOrder::Forward,
|
||||
)
|
||||
.await?
|
||||
.iter()
|
||||
.map(|k| {
|
||||
let p = k.params().unwrap();
|
||||
.await?;
|
||||
|
||||
ListKeysResponseItem {
|
||||
id: k.key_id.to_string(),
|
||||
name: p.name.get().clone(),
|
||||
created: p.created.map(|x| {
|
||||
DateTime::from_timestamp_millis(x as i64)
|
||||
.expect("invalid timestamp stored in db")
|
||||
}),
|
||||
expiration: p.expiration.get().inner().map(|x| {
|
||||
DateTime::from_timestamp_millis(x.0 as i64)
|
||||
.expect("invalid timestamp stored in db")
|
||||
}),
|
||||
expired: p.is_expired(now),
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
if self.details {
|
||||
let mut stream = keys
|
||||
.into_iter()
|
||||
.map(|k| key_info_results(garage, k, false))
|
||||
.collect::<futures::stream::FuturesOrdered<_>>();
|
||||
|
||||
Ok(ListKeysResponse(res))
|
||||
let mut res = vec![];
|
||||
while let Some(next) = stream.next().await {
|
||||
res.push(next?);
|
||||
}
|
||||
|
||||
Ok(ListKeysResponse::WithDetails(res))
|
||||
} else {
|
||||
let res = keys
|
||||
.iter()
|
||||
.map(|k| {
|
||||
let p = k.params().unwrap();
|
||||
|
||||
ListKeysResponseItem {
|
||||
id: k.key_id.to_string(),
|
||||
name: p.name.get().clone(),
|
||||
created: p.created.map(|x| {
|
||||
DateTime::from_timestamp_millis(x as i64)
|
||||
.expect("invalid timestamp stored in db")
|
||||
}),
|
||||
expiration: p.expiration.get().inner().map(|x| {
|
||||
DateTime::from_timestamp_millis(x.0 as i64)
|
||||
.expect("invalid timestamp stored in db")
|
||||
}),
|
||||
expired: p.is_expired(now),
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
Ok(ListKeysResponse::WithoutDetails(res))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -364,6 +364,7 @@ fn ClusterLayoutSkipDeadNodes() {}
|
||||
path = "/v2/ListKeys",
|
||||
tag = "Access key",
|
||||
description = "Returns all API access keys in the cluster.",
|
||||
params(ListKeysRequest),
|
||||
responses(
|
||||
(status = 200, description = "Returns the key identifier (aka `AWS_ACCESS_KEY_ID`) and its associated, human friendly, name if any (otherwise return an empty string)", body = ListKeysResponse),
|
||||
(status = 500, description = "Internal server error")
|
||||
@@ -453,6 +454,7 @@ fn DeleteKey() {}
|
||||
path = "/v2/ListBuckets",
|
||||
tag = "Bucket",
|
||||
description = "List all the buckets on the cluster with their UUID and their global and local aliases.",
|
||||
params(ListBucketsRequest),
|
||||
responses(
|
||||
(status = 200, description = "Returns the UUID of all the buckets and all their aliases", body = ListBucketsResponse),
|
||||
(status = 500, description = "Internal server error")
|
||||
@@ -876,7 +878,7 @@ impl Modify for SecurityAddon {
|
||||
#[derive(OpenApi)]
|
||||
#[openapi(
|
||||
info(
|
||||
version = "v2.3.0",
|
||||
version = "v2.4.0",
|
||||
title = "Garage administration API",
|
||||
description = "Administrate your Garage cluster programmatically, including status, layout, keys, buckets, and maintenance tasks.
|
||||
|
||||
|
||||
@@ -55,10 +55,10 @@ impl AdminApiRequest {
|
||||
POST CreateKey (body),
|
||||
POST ImportKey (body),
|
||||
POST DeleteKey (query::id),
|
||||
GET ListKeys (),
|
||||
GET ListKeys (parse_default(false)::details, query_opt::offset, opt_parse::limit),
|
||||
// Bucket endpoints
|
||||
GET GetBucketInfo (query_opt::id, query_opt::global_alias, query_opt::search),
|
||||
GET ListBuckets (),
|
||||
GET ListBuckets (parse_default(false)::details, query_opt::offset, opt_parse::limit),
|
||||
POST CreateBucket (body),
|
||||
POST DeleteBucket (query::id),
|
||||
POST UpdateBucket (body_field, query::id),
|
||||
@@ -129,7 +129,7 @@ impl AdminApiRequest {
|
||||
)),
|
||||
|
||||
// Keys
|
||||
Endpoint::ListKeys => Ok(AdminApiRequest::ListKeys(ListKeysRequest)),
|
||||
Endpoint::ListKeys => Ok(AdminApiRequest::ListKeys(ListKeysRequest::default())),
|
||||
Endpoint::GetKeyInfo {
|
||||
id,
|
||||
search,
|
||||
@@ -161,7 +161,9 @@ impl AdminApiRequest {
|
||||
// Endpoint::DeleteKey { id } => Ok(AdminApiRequest::DeleteKey(DeleteKeyRequest { id })),
|
||||
|
||||
// Buckets
|
||||
Endpoint::ListBuckets => Ok(AdminApiRequest::ListBuckets(ListBucketsRequest)),
|
||||
Endpoint::ListBuckets => {
|
||||
Ok(AdminApiRequest::ListBuckets(ListBucketsRequest::default()))
|
||||
}
|
||||
Endpoint::GetBucketInfo { id, global_alias } => {
|
||||
Ok(AdminApiRequest::GetBucketInfo(GetBucketInfoRequest {
|
||||
id,
|
||||
@@ -271,6 +273,9 @@ generateQueryParameters! {
|
||||
"accessKeyId" => access_key_id,
|
||||
"showSecretKey" => show_secret_key,
|
||||
"bucketId" => bucket_id,
|
||||
"key" => key
|
||||
"key" => key,
|
||||
"details" => details,
|
||||
"offset" => offset,
|
||||
"limit" => limit
|
||||
]
|
||||
}
|
||||
|
||||
+2
-228
@@ -47,26 +47,15 @@ where
|
||||
HI: Iterator<Item = S>,
|
||||
S: AsRef<str>,
|
||||
{
|
||||
rule.allow_origins.iter().any(|x| wildcard_match(x, origin))
|
||||
rule.allow_origins.iter().any(|x| x == "*" || x == origin)
|
||||
&& rule.allow_methods.iter().any(|x| x == "*" || x == method)
|
||||
&& request_headers.all(|h| {
|
||||
rule.allow_headers
|
||||
.iter()
|
||||
.any(|x| wildcard_match(x, h.as_ref()))
|
||||
.any(|x| x == "*" || x == h.as_ref())
|
||||
})
|
||||
}
|
||||
|
||||
/// Checks whether `candidate` matches the pattern `allowed_wildcard`.
|
||||
#[inline]
|
||||
fn wildcard_match(allowed_wildcard: &String, candidate: &str) -> bool {
|
||||
if allowed_wildcard.contains("*") {
|
||||
let parts = allowed_wildcard.split("*").collect::<Vec<&str>>();
|
||||
parts.len() == 2 && candidate.starts_with(parts[0]) && candidate.ends_with(parts[1])
|
||||
} else {
|
||||
candidate == allowed_wildcard
|
||||
}
|
||||
}
|
||||
|
||||
pub fn add_cors_headers(
|
||||
resp: &mut Response<impl Body>,
|
||||
rule: &GarageCorsRule,
|
||||
@@ -201,221 +190,6 @@ pub fn handle_options_for_bucket<B>(
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn cors_rule(
|
||||
allow_origins: &[&str],
|
||||
allow_methods: &[&str],
|
||||
allow_headers: &[&str],
|
||||
) -> GarageCorsRule {
|
||||
GarageCorsRule {
|
||||
id: None,
|
||||
max_age_seconds: None,
|
||||
allow_origins: allow_origins.iter().map(|s| s.to_string()).collect(),
|
||||
allow_methods: allow_methods.iter().map(|s| s.to_string()).collect(),
|
||||
allow_headers: allow_headers.iter().map(|s| s.to_string()).collect(),
|
||||
expose_headers: vec![],
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn matches_when_origin_method_and_headers_are_explicitly_allowed() {
|
||||
let rule = cors_rule(
|
||||
&["https://app.example.test"],
|
||||
&["GET", "PUT"],
|
||||
&["content-type", "x-custom"],
|
||||
);
|
||||
let headers = vec!["content-type", "x-custom"];
|
||||
|
||||
assert!(cors_rule_matches(
|
||||
&rule,
|
||||
"https://app.example.test",
|
||||
"PUT",
|
||||
headers.iter(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn does_not_match_when_origin_is_not_allowed() {
|
||||
let rule = cors_rule(&["https://app.example.test"], &["GET"], &["*"]);
|
||||
|
||||
assert!(!cors_rule_matches(
|
||||
&rule,
|
||||
"https://evil.example.test",
|
||||
"GET",
|
||||
std::iter::empty::<&str>(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn does_not_match_when_method_is_not_allowed() {
|
||||
let rule = cors_rule(&["*"], &["GET"], &["*"]);
|
||||
|
||||
assert!(!cors_rule_matches(
|
||||
&rule,
|
||||
"https://app.example.test",
|
||||
"DELETE",
|
||||
std::iter::empty::<&str>(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn does_not_match_when_a_requested_header_is_not_allowed() {
|
||||
let rule = cors_rule(&["*"], &["GET"], &["content-type"]);
|
||||
let headers = vec!["content-type", "x-not-allowed"];
|
||||
|
||||
assert!(!cors_rule_matches(
|
||||
&rule,
|
||||
"https://app.example.test",
|
||||
"GET",
|
||||
headers.iter(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wildcard_origin_method_and_headers_match_anything() {
|
||||
let rule = cors_rule(&["*"], &["*"], &["*"]);
|
||||
let headers = vec!["x-anything"];
|
||||
|
||||
assert!(cors_rule_matches(
|
||||
&rule,
|
||||
"https://app.example.test",
|
||||
"DELETE",
|
||||
headers.iter(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wildcard_origin_regex() {
|
||||
let rule = cors_rule(&["https://*.localhost.com"], &["*"], &["*"]);
|
||||
let headers = vec!["x-anything"];
|
||||
|
||||
assert!(cors_rule_matches(
|
||||
&rule,
|
||||
"https://s3.localhost.com",
|
||||
"DELETE",
|
||||
headers.iter(),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn origin_matching_cases() {
|
||||
// (allow_origins, origin, expect_match)
|
||||
let cases: &[(&[&str], &str, bool)] = &[
|
||||
// exact match
|
||||
(
|
||||
&["https://app.example.test"],
|
||||
"https://app.example.test",
|
||||
true,
|
||||
),
|
||||
(
|
||||
&["https://app.example.test"],
|
||||
"https://other.example.test",
|
||||
false,
|
||||
),
|
||||
// full wildcard
|
||||
(&["*"], "https://anything.example.test", true),
|
||||
// subdomain glob
|
||||
(
|
||||
&["https://*.example.test"],
|
||||
"https://foo.example.test",
|
||||
true,
|
||||
),
|
||||
(&["https://*.example.test"], "https://example.test", false),
|
||||
(
|
||||
&["https://*.example.test"],
|
||||
"http://foo.example.test",
|
||||
false,
|
||||
),
|
||||
// multiple allowed origins, at least one should match
|
||||
(
|
||||
&["https://a.example.test", "https://b.example.test"],
|
||||
"https://b.example.test",
|
||||
true,
|
||||
),
|
||||
// match multiple origins
|
||||
(
|
||||
&["https://a*.example.test", "https://ab*.example.test"],
|
||||
"https://abc.example.test",
|
||||
true,
|
||||
),
|
||||
(
|
||||
&["https://a.example.test", "https://b.example.test"],
|
||||
"https://c.example.test",
|
||||
false,
|
||||
),
|
||||
// at most one '*' in a pattern is allowed
|
||||
(&["https://*.example.*"], "https://a.example.test", false),
|
||||
// domain changed with wildcard
|
||||
(
|
||||
&["https://*example.test"],
|
||||
"https://garageexample.test",
|
||||
true,
|
||||
),
|
||||
// trailing '*' matches any suffix, including the empty string,
|
||||
// so this also matches origins with anything (or nothing) after
|
||||
// "example."
|
||||
(&["https://example.*"], "https://example.test", true),
|
||||
(&["https://*example.test"], "https://example.test", true),
|
||||
(&["https://example.*"], "https://example.", true),
|
||||
];
|
||||
|
||||
for (allow_origins, origin, expect_match) in cases {
|
||||
let rule = cors_rule(allow_origins, &["GET"], &["*"]);
|
||||
let got = cors_rule_matches(&rule, origin, "GET", std::iter::empty::<&str>());
|
||||
assert_eq!(
|
||||
got, *expect_match,
|
||||
"allow_origins={allow_origins:?}, origin={origin:?}: expected match={expect_match}, got {got}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn header_matching_cases() {
|
||||
// (allow_headers, requested_headers, expect_match)
|
||||
let cases: &[(&[&str], &[&str], bool)] = &[
|
||||
// exact match
|
||||
(&["content-type"], &["content-type"], true),
|
||||
(&["content-type"], &["x-custom"], false),
|
||||
// full wildcard
|
||||
(&["*"], &["x-anything"], true),
|
||||
// no headers requested always matches, regardless of allow_headers
|
||||
(&["content-type"], &[], true),
|
||||
(&[], &[], true),
|
||||
// prefix glob
|
||||
(&["x-amz-*"], &["x-amz-meta-foo"], true),
|
||||
(&["x-amz-*"], &["x-amz-"], true),
|
||||
(&["x-amz-*"], &["x-other"], false),
|
||||
// suffix glob
|
||||
(&["*-meta"], &["foo-meta"], true),
|
||||
(&["*-meta"], &["-meta"], true),
|
||||
(&["*-meta"], &["foo-meta-bar"], false),
|
||||
// multiple allowed headers, at least one should match per requested header
|
||||
(
|
||||
&["content-type", "x-amz-*"],
|
||||
&["content-type", "x-amz-meta-foo"],
|
||||
true,
|
||||
),
|
||||
(&["content-type", "x-amz-*"], &["x-other"], false),
|
||||
// all requested headers must be covered
|
||||
(&["content-type"], &["content-type", "x-custom"], false),
|
||||
// at most one '*' in a pattern is allowed
|
||||
(&["x-*-*"], &["x-a-b"], false),
|
||||
];
|
||||
|
||||
for (allow_headers, requested_headers, expect_match) in cases {
|
||||
let rule = cors_rule(&["*"], &["GET"], allow_headers);
|
||||
let got = cors_rule_matches(
|
||||
&rule,
|
||||
"https://app.example.test",
|
||||
"GET",
|
||||
requested_headers.iter(),
|
||||
);
|
||||
assert_eq!(
|
||||
got, *expect_match,
|
||||
"allow_headers={allow_headers:?}, requested_headers={requested_headers:?}: expected match={expect_match}, got {got}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn bucket_params_with_rule(allow_origins: Vec<&str>) -> BucketParams {
|
||||
let mut bucket_params = BucketParams::default();
|
||||
bucket_params.cors_config.update(
|
||||
|
||||
@@ -31,12 +31,19 @@ impl Cli {
|
||||
}
|
||||
|
||||
pub async fn cmd_list_buckets(&self) -> Result<(), Error> {
|
||||
let mut buckets = self.api_request(ListBucketsRequest).await?;
|
||||
let mut buckets = match self.api_request(ListBucketsRequest::default()).await? {
|
||||
ListBucketsResponse::WithoutDetails(list) => list,
|
||||
_ => {
|
||||
return Err(Error::Message(
|
||||
"Unexpected ListBuckets response format".into(),
|
||||
))
|
||||
}
|
||||
};
|
||||
|
||||
buckets.0.sort_by_key(|x| x.created);
|
||||
buckets.sort_by_key(|x| x.created);
|
||||
|
||||
let mut table = vec!["ID\tCreated\tGlobal aliases\tLocal aliases".to_string()];
|
||||
for bucket in buckets.0.iter() {
|
||||
for bucket in buckets.iter() {
|
||||
table.push(format!(
|
||||
"{:.16}\t{}\t{}\t{}",
|
||||
bucket.id,
|
||||
|
||||
@@ -28,12 +28,15 @@ impl Cli {
|
||||
}
|
||||
|
||||
pub async fn cmd_list_keys(&self) -> Result<(), Error> {
|
||||
let mut keys = self.api_request(ListKeysRequest).await?;
|
||||
let mut keys = match self.api_request(ListKeysRequest::default()).await? {
|
||||
ListKeysResponse::WithoutDetails(list) => list,
|
||||
_ => return Err(Error::Message("Unexpected ListKeys response format".into())),
|
||||
};
|
||||
|
||||
keys.0.sort_by_key(|x| x.created);
|
||||
keys.sort_by_key(|x| x.created);
|
||||
|
||||
let mut table = vec!["ID\tCreated\tName\tExpiration".to_string()];
|
||||
for key in keys.0.iter() {
|
||||
for key in keys.iter() {
|
||||
let exp = if key.expired {
|
||||
Cow::from("expired")
|
||||
} else {
|
||||
@@ -243,7 +246,10 @@ impl Cli {
|
||||
}
|
||||
|
||||
pub async fn cmd_delete_expired_keys(&self, yes: bool) -> Result<(), Error> {
|
||||
let mut list = self.api_request(ListKeysRequest).await?.0;
|
||||
let mut list = match self.api_request(ListKeysRequest::default()).await? {
|
||||
ListKeysResponse::WithoutDetails(list) => list,
|
||||
_ => return Err(Error::Message("Unexpected ListKeys response format".into())),
|
||||
};
|
||||
|
||||
list.retain(|key| key.expired);
|
||||
|
||||
|
||||
@@ -110,16 +110,29 @@ impl Cli {
|
||||
Ok(resp.success.into_iter().next().unwrap().1)
|
||||
}
|
||||
|
||||
pub async fn cmd_json_api(&self, endpoint: String, payload: String) -> Result<(), Error> {
|
||||
let payload: serde_json::Value = if payload == "-" {
|
||||
serde_json::from_reader(&std::io::stdin())?
|
||||
} else {
|
||||
serde_json::from_str(&payload)?
|
||||
};
|
||||
pub async fn cmd_json_api(
|
||||
&self,
|
||||
endpoint: String,
|
||||
payload: Option<String>,
|
||||
) -> Result<(), Error> {
|
||||
let request: AdminApiRequest = if let Some(payload) = payload {
|
||||
let payload: serde_json::Value = if payload == "-" {
|
||||
serde_json::from_reader(&std::io::stdin())?
|
||||
} else {
|
||||
serde_json::from_str(&payload)?
|
||||
};
|
||||
|
||||
let request: AdminApiRequest = serde_json::from_value(serde_json::json!({
|
||||
endpoint.clone(): payload,
|
||||
}))?;
|
||||
serde_json::from_value(serde_json::json!({
|
||||
endpoint.clone(): payload,
|
||||
}))?
|
||||
} else {
|
||||
serde_json::from_value(serde_json::json!({
|
||||
endpoint.clone(): null,
|
||||
}))
|
||||
.or(serde_json::from_value(serde_json::json!({
|
||||
endpoint.clone(): {},
|
||||
})))?
|
||||
};
|
||||
|
||||
let resp = match self
|
||||
.proxy_rpc_endpoint
|
||||
|
||||
@@ -78,8 +78,7 @@ pub enum Command {
|
||||
/// The admin API endpoint to invoke, e.g. `GetClusterStatus`
|
||||
endpoint: String,
|
||||
/// The JSON payload, or `-` to read from `stdin`
|
||||
#[structopt(default_value = "null")]
|
||||
payload: String,
|
||||
payload: Option<String>,
|
||||
},
|
||||
|
||||
/// Generate completions for a shell
|
||||
|
||||
+1
-7
@@ -276,11 +276,6 @@ fn init_logging(opt: &Opt) {
|
||||
tracing_subscriber::fmt()
|
||||
.with_writer(std::io::stderr)
|
||||
.with_env_filter(env_filter)
|
||||
.with_ansi(
|
||||
std::env::var("NO_COLOR")
|
||||
.map(|x| x != "0" && !x.eq_ignore_ascii_case("false"))
|
||||
.unwrap_or(true),
|
||||
)
|
||||
.init();
|
||||
}
|
||||
|
||||
@@ -307,8 +302,7 @@ async fn cli_command(opt: Opt) -> Result<(), Error> {
|
||||
|
||||
let net_key_hex_str = rpc_secret.ok_or("No RPC secret provided")?;
|
||||
let network_key = NetworkKey::from_slice(
|
||||
&hex::decode(net_key_hex_str.extract_secret())
|
||||
.err_context("Invalid RPC secret key (bad hex)")?[..],
|
||||
&hex::decode(&net_key_hex_str).err_context("Invalid RPC secret key (bad hex)")?[..],
|
||||
)
|
||||
.ok_or("Invalid RPC secret provided (wrong length)")?;
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@ use std::path::PathBuf;
|
||||
|
||||
use structopt::StructOpt;
|
||||
|
||||
use garage_util::config::{Config, Secret};
|
||||
use garage_util::config::Config;
|
||||
use garage_util::error::Error;
|
||||
|
||||
/// Structure for secret values or paths that are passed as CLI arguments or environment
|
||||
@@ -99,7 +99,7 @@ pub fn fill_secrets(mut config: Config, secrets: Secrets) -> Result<Config, Erro
|
||||
}
|
||||
|
||||
pub(crate) fn fill_secret(
|
||||
config_secret: &mut Option<Secret<String>>,
|
||||
config_secret: &mut Option<String>,
|
||||
config_secret_file: &Option<PathBuf>,
|
||||
cli_secret: &Option<String>,
|
||||
cli_secret_file: &Option<PathBuf>,
|
||||
@@ -110,7 +110,7 @@ pub(crate) fn fill_secret(
|
||||
(Some(_), Some(_)) => {
|
||||
return Err(format!("only one of `{}` and `{}_file` can be set", name, name).into());
|
||||
}
|
||||
(Some(secret), None) => Some(Secret::new(secret.to_string())),
|
||||
(Some(secret), None) => Some(secret.to_string()),
|
||||
(None, Some(file)) => Some(read_secret_file(file, allow_world_readable)?),
|
||||
(None, None) => None,
|
||||
};
|
||||
@@ -132,10 +132,7 @@ pub(crate) fn fill_secret(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn read_secret_file(
|
||||
file_path: &PathBuf,
|
||||
allow_world_readable: bool,
|
||||
) -> Result<Secret<String>, Error> {
|
||||
fn read_secret_file(file_path: &PathBuf, allow_world_readable: bool) -> Result<String, Error> {
|
||||
if !allow_world_readable {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
@@ -155,7 +152,7 @@ fn read_secret_file(
|
||||
|
||||
// trim_end: allows for use case such as `echo "$(openssl rand -hex 32)" > somefile`.
|
||||
// also editors sometimes add a trailing newline
|
||||
Ok(Secret::new(String::from(secret_buf.trim_end())))
|
||||
Ok(String::from(secret_buf.trim_end()))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -152,14 +152,6 @@ pub async fn run_server(
|
||||
}
|
||||
}
|
||||
|
||||
// Deregister from Consul (if enabled) in the background, in parallel with the
|
||||
// rest of the shutdown sequence, so that it doesn't add to shutdown latency.
|
||||
#[cfg(feature = "consul-discovery")]
|
||||
let deregister_consul_task = tokio::spawn({
|
||||
let system = garage.system.clone();
|
||||
async move { system.deregister_from_discovery().await }
|
||||
});
|
||||
|
||||
// Remove RPC handlers for system to break reference cycles
|
||||
info!("Deregistering RPC handlers for shutdown...");
|
||||
garage.system.netapp.drop_all_handlers();
|
||||
@@ -176,12 +168,6 @@ pub async fn run_server(
|
||||
// Await for all background tasks to end
|
||||
await_background_done.await?;
|
||||
|
||||
// Await for Consul deregistration to end, if it hasn't already
|
||||
#[cfg(feature = "consul-discovery")]
|
||||
if let Err(e) = deregister_consul_task.await {
|
||||
error!("Error while joining Consul deregistration task: {}", e);
|
||||
}
|
||||
|
||||
info!("Cleaning up...");
|
||||
|
||||
Ok(())
|
||||
|
||||
+1
-1
@@ -137,7 +137,7 @@ impl Garage {
|
||||
info!("Initializing RPC...");
|
||||
let network_key = hex::decode(config.rpc_secret.as_ref().ok_or_message(
|
||||
"rpc_secret value is missing, not present in config file or in environment",
|
||||
)?.extract_secret())
|
||||
)?)
|
||||
.ok()
|
||||
.and_then(|x| NetworkKey::from_slice(&x))
|
||||
.ok_or_message("Invalid RPC secret key: expected 32 bytes of random hex, please check the documentation for requirements")?;
|
||||
|
||||
+1
-31
@@ -115,7 +115,7 @@ impl ConsulDiscovery {
|
||||
let mut headers = reqwest::header::HeaderMap::new();
|
||||
headers.insert(
|
||||
"x-consul-token",
|
||||
reqwest::header::HeaderValue::from_str(token.extract_secret())?,
|
||||
reqwest::header::HeaderValue::from_str(token)?,
|
||||
);
|
||||
builder = builder.default_headers(headers);
|
||||
}
|
||||
@@ -183,36 +183,6 @@ impl ConsulDiscovery {
|
||||
}
|
||||
// ---- PUBLISHING TO CONSUL CATALOG ----
|
||||
|
||||
#[cfg(feature = "consul-discovery")]
|
||||
pub async fn deregister_consul_service(&self, node_id: NodeID) -> Result<(), ConsulError> {
|
||||
let node = format!("garage:{}", hex::encode(&node_id[..8]));
|
||||
let url = format!(
|
||||
"{}/v1/{}",
|
||||
self.config.consul_http_addr,
|
||||
(match &self.config.api {
|
||||
ConsulDiscoveryAPI::Catalog => format!("catalog/deregister"),
|
||||
ConsulDiscoveryAPI::Agent => format!("agent/service/deregister/{}", node),
|
||||
})
|
||||
);
|
||||
|
||||
let req = self.client.put(&url);
|
||||
|
||||
let http = if matches!(&self.config.api, ConsulDiscoveryAPI::Catalog) {
|
||||
let deregister_request = serde_json::json!({
|
||||
"Node": node,
|
||||
"ServiceID": node,
|
||||
});
|
||||
let req = req.json(&deregister_request);
|
||||
req.send().await?
|
||||
} else {
|
||||
req.send().await?
|
||||
};
|
||||
http.error_for_status()?;
|
||||
|
||||
debug!("Deregistered service {} from Consul", node);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn publish_consul_service(
|
||||
&self,
|
||||
node_id: NodeID,
|
||||
|
||||
+1
-10
@@ -358,7 +358,7 @@ impl System {
|
||||
);
|
||||
}
|
||||
|
||||
pub fn cleanup(self: &Arc<Self>) {
|
||||
pub fn cleanup(&self) {
|
||||
// Break reference cycle
|
||||
self.metrics.store(None);
|
||||
}
|
||||
@@ -650,15 +650,6 @@ impl System {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "consul-discovery")]
|
||||
pub async fn deregister_from_discovery(self: &Arc<Self>) {
|
||||
if let Some(c) = &self.consul_discovery {
|
||||
if let Err(e) = c.deregister_consul_service(self.netapp.id).await {
|
||||
error!("Error while deregistering from Consul: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn discovery_loop(self: &Arc<Self>, mut stop_signal: watch::Receiver<bool>) {
|
||||
while !*stop_signal.borrow() {
|
||||
let peers_up = self
|
||||
|
||||
@@ -32,7 +32,6 @@ lazy_static.workspace = true
|
||||
tracing.workspace = true
|
||||
rand.workspace = true
|
||||
sha2.workspace = true
|
||||
subtle.workspace = true
|
||||
|
||||
chrono.workspace = true
|
||||
rmp-serde.workspace = true
|
||||
|
||||
+4
-35
@@ -90,7 +90,7 @@ pub struct Config {
|
||||
pub allow_world_readable_secrets: bool,
|
||||
|
||||
/// RPC secret key: 32 bytes hex encoded
|
||||
pub rpc_secret: Option<Secret<String>>,
|
||||
pub rpc_secret: Option<String>,
|
||||
/// Optional file where RPC secret key is read from
|
||||
pub rpc_secret_file: Option<PathBuf>,
|
||||
/// Address to bind for RPC
|
||||
@@ -205,37 +205,6 @@ pub struct WebConfig {
|
||||
pub add_host_to_metrics: bool,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Clone)]
|
||||
#[serde(transparent)]
|
||||
pub struct Secret<T>(T);
|
||||
|
||||
impl<T> Secret<T> {
|
||||
pub fn new(secret: T) -> Self {
|
||||
Secret(secret)
|
||||
}
|
||||
|
||||
pub fn extract_secret(&self) -> &T {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: std::ops::Deref<Target = str>> Secret<T> {
|
||||
pub fn eq_ct(&self, other: &T) -> bool {
|
||||
use subtle::ConstantTimeEq;
|
||||
self.0
|
||||
.deref()
|
||||
.as_bytes()
|
||||
.ct_eq(other.deref().as_bytes())
|
||||
.into()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> std::fmt::Debug for Secret<T> {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("Secret").finish_non_exhaustive()
|
||||
}
|
||||
}
|
||||
|
||||
/// Configuration for the admin and monitoring HTTP API
|
||||
#[derive(Deserialize, Debug, Clone, Default)]
|
||||
pub struct AdminConfig {
|
||||
@@ -243,7 +212,7 @@ pub struct AdminConfig {
|
||||
pub api_bind_addr: Option<UnixOrTCPSocketAddress>,
|
||||
|
||||
/// Bearer token to use to scrape metrics
|
||||
pub metrics_token: Option<Secret<String>>,
|
||||
pub metrics_token: Option<String>,
|
||||
/// File to read metrics token from
|
||||
pub metrics_token_file: Option<PathBuf>,
|
||||
/// Whether to require an access token for accessing the metrics endpoint
|
||||
@@ -251,7 +220,7 @@ pub struct AdminConfig {
|
||||
pub metrics_require_token: bool,
|
||||
|
||||
/// Bearer token to use to access Admin API endpoints
|
||||
pub admin_token: Option<Secret<String>>,
|
||||
pub admin_token: Option<String>,
|
||||
/// File to read admin token from
|
||||
pub admin_token_file: Option<PathBuf>,
|
||||
|
||||
@@ -283,7 +252,7 @@ pub struct ConsulDiscoveryConfig {
|
||||
/// Client TLS key to use when connecting to Consul
|
||||
pub client_key: Option<String>,
|
||||
/// /// Token to use for connecting to consul
|
||||
pub token: Option<Secret<String>>,
|
||||
pub token: Option<String>,
|
||||
/// Skip TLS hostname verification
|
||||
#[serde(default)]
|
||||
pub tls_skip_verify: bool,
|
||||
|
||||
Reference in New Issue
Block a user