Compare commits

..

2 Commits

Author SHA1 Message Date
Alex Auvolat be17e25bee CLI: garage json-api: allow payload to be omitted more often 2026-07-25 15:24:31 +02:00
Alex Auvolat 8d85301808 Admin API: add parameters details, offset and limit to ListKeys and ListBuckets 2026-07-25 15:12:37 +02:00
27 changed files with 420 additions and 460 deletions
Generated
-1
View File
@@ -1890,7 +1890,6 @@ dependencies = [
"serde",
"serde_json",
"sha2 0.10.9",
"subtle",
"thiserror 2.0.18",
"tokio",
"toml",
-1
View File
@@ -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"] }
+197 -9
View File
@@ -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"
+6 -11
View File
@@ -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:
+6 -5
View 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).
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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
View File
@@ -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")]
+2 -8
View File
@@ -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
View File
@@ -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
View File
@@ -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))
}
}
}
+3 -1
View File
@@ -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.
+10 -5
View File
@@ -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
View File
@@ -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(
+10 -3
View File
@@ -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,
+10 -4
View File
@@ -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);
+22 -9
View File
@@ -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
+1 -2
View File
@@ -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
View File
@@ -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)")?;
+5 -8
View File
@@ -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)]
-14
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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
-1
View File
@@ -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
View File
@@ -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,