chore: remove the unused restic-api reference implementation (#684)

This commit is contained in:
Antoine Lecompte
2026-09-15 13:08:31 +00:00
committed by GitHub
parent 813fe2f9fc
commit 79c914c375
44 changed files with 8 additions and 4193 deletions
-671
View File
@@ -1,671 +0,0 @@
#!/usr/bin/env bash
#MISE description="Benchmark restic backends (restic-api, michael, or both)"
#MISE depends=["docker:start", "restic-api:build", "michael:build"]
#MISE dir="{{config_root}}"
#
# Benchmark restic-compatible backends using the real restic CLI.
#
# Supports restic-api (Node.js/NestJS) and michael (Go), with an optional
# comparison mode that runs both sequentially and prints a side-by-side table.
#
# Prerequisites:
# - MinIO running (mise docker:start)
# - restic, node, perl installed (mise install)
# - For --compare: pnpm deps + go (mise install:deps)
#
# Usage:
# # Benchmark both backends and compare (default):
# mise bench
#
# # Benchmark a single backend:
# mise bench -- --backend restic-api
# mise bench -- --backend michael
#
# Options:
# --backend NAME Benchmark only one backend: restic-api | michael
# --small-files N Number of small files (default: 1000)
# --medium-files N Number of medium files (default: 1000)
# --large-files N Number of large files (default: 100)
# --large-size-mb N Size of each large file (default: 100)
# --no-cleanup Keep S3 buckets after run
#
set -euo pipefail
REPO_ROOT="$PWD"
# ─── Configuration ──────────────────────────────────────────────────────────────
RESTIC_ENDPOINT="${RESTIC_ENDPOINT:-http://localhost:3010}"
# bench spawns + port-manages the backend, so split the endpoint back into host/port
_endpoint_hostport="${RESTIC_ENDPOINT#*://}"
API_HOST="${_endpoint_hostport%%:*}"
API_PORT="${_endpoint_hostport##*:}"
API_PORT="${API_PORT%%/*}"
# the restic-api / michael server bench launches binds RESTIC_API_PORT
export RESTIC_API_PORT="$API_PORT"
export JWT_SECRET="${JWT_SECRET:-cca13c34b450a77c1d4b9ecd25dff6aebc6d7417afdb31864f5943c59abd03a1}"
BACKEND=""
COMPARE=true
SMALL_COUNT=50
SMALL_SIZE_KB=100
MEDIUM_COUNT=50
MEDIUM_SIZE_MB=20
LARGE_COUNT=20
LARGE_SIZE_MB=100
DO_CLEANUP=true
while [[ $# -gt 0 ]]; do
case "$1" in
--backend) BACKEND="$2"; COMPARE=false; shift 2 ;;
--small-files) SMALL_COUNT="$2"; shift 2 ;;
--medium-files) MEDIUM_COUNT="$2"; shift 2 ;;
--large-files) LARGE_COUNT="$2"; shift 2 ;;
--large-size-mb) LARGE_SIZE_MB="$2"; shift 2 ;;
--no-cleanup) DO_CLEANUP=false; shift ;;
*) echo "Unknown option: $1"; exit 1 ;;
esac
done
BENCH_DIR=$(mktemp -d)
MEM_PID=""
API_PID=""
MANAGED_PID=""
REPOS=()
# ─── Formatting ─────────────────────────────────────────────────────────────────
bold() { printf '\033[1m%s\033[0m' "$*"; }
dim() { printf '\033[2m%s\033[0m' "$*"; }
green() { printf '\033[32m%s\033[0m' "$*"; }
red() { printf '\033[31m%s\033[0m' "$*"; }
# ─── Helpers ────────────────────────────────────────────────────────────────────
on_exit() {
stop_mem_monitor
stop_backend
if [[ "$DO_CLEANUP" == true ]] && (( ${#REPOS[@]} > 0 )); then
cleanup_repos
fi
rm -rf "$BENCH_DIR"
}
trap on_exit EXIT
now_ms() {
perl -MTime::HiRes -e 'printf "%.0f\n", Time::HiRes::time() * 1000'
}
elapsed_s() {
perl -e "printf '%.3f', ($2 - $1) / 1000"
}
throughput_mbs() {
local bytes="$1" secs="$2"
perl -e "if ($secs > 0) { printf '%.1f', $bytes / 1048576 / $secs } else { print 'inf' }"
}
human_size() {
perl -e '
my $b = $ARGV[0];
if ($b >= 1073741824) { printf "%.1f GB", $b/1073741824 }
elsif ($b >= 1048576) { printf "%.1f MB", $b/1048576 }
elsif ($b >= 1024) { printf "%.1f KB", $b/1024 }
else { printf "%d B", $b }
' "$1"
}
make_jwt() {
local repo="$1"
node --no-warnings -e "
const c = require('node:crypto');
const h = Buffer.from(JSON.stringify({alg:'HS256',typ:'JWT'})).toString('base64url');
const p = Buffer.from(JSON.stringify({
user: c.randomUUID(),
repository: process.argv[1],
writeOnce: false,
iat: Math.floor(Date.now()/1000),
exp: Math.floor(Date.now()/1000) + 7200
})).toString('base64url');
const s = c.createHmac('sha256', process.env.JWT_SECRET)
.update(h+'.'+p).digest('base64url');
process.stdout.write(h+'.'+p+'.'+s);
" "$repo"
}
init_repo() {
local repo
repo=$(node --no-warnings -e "process.stdout.write(require('node:crypto').randomUUID())")
local token
token=$(make_jwt "$repo")
REPOS+=("$repo")
export RESTIC_REPOSITORY="rest:http://_:${token}@${API_HOST}:${API_PORT}/${repo}"
export RESTIC_PASSWORD="bench"
}
cleanup_repos() {
echo ""
echo " $(dim "Cleaning up ${#REPOS[@]} benchmark repositories...")"
(cd "$REPO_ROOT/packages/restic-api" && node --no-warnings -e "
const { S3Client, ListObjectsV2Command, DeleteObjectsCommand, DeleteBucketCommand }
= require('@aws-sdk/client-s3');
const client = new S3Client({
credentials: {
accessKeyId: process.env.S3_ACCESS_KEY_ID || 'minio',
secretAccessKey: process.env.S3_SECRET_ACCESS_KEY || 'miniominio',
},
region: process.env.S3_REGION || 'minio',
endpoint: process.env.S3_ENDPOINT || 'http://localhost:9000',
forcePathStyle: true,
});
async function nuke(Bucket) {
try {
const { Contents } = await client.send(new ListObjectsV2Command({ Bucket }));
if (Contents?.length) {
await client.send(new DeleteObjectsCommand({
Bucket, Delete: { Objects: Contents.map(({ Key }) => ({ Key })) }
}));
}
await client.send(new DeleteBucketCommand({ Bucket }));
} catch {}
}
Promise.all(process.argv.slice(1).map(nuke)).then(() => process.exit(0));
" "${REPOS[@]}") || true
}
# ─── Results I/O ────────────────────────────────────────────────────────────────
save_result() {
echo "$1=$2" >> "$BENCH_DIR/results/${CURRENT_BACKEND}.txt"
}
read_result() {
local key="$1" file="$2"
grep "^${key}=" "$file" 2>/dev/null | tail -1 | cut -d= -f2
}
# ─── Memory monitoring ──────────────────────────────────────────────────────────
start_mem_monitor() {
API_PID=$(lsof -ti ":$API_PORT" -sTCP:LISTEN 2>/dev/null | head -1 || true)
if [[ -z "$API_PID" ]]; then
echo " $(dim '(could not find API process — memory stats unavailable)')"
return
fi
local peak_file="$BENCH_DIR/.mem_peak"
echo "0" > "$peak_file"
(
while kill -0 "$API_PID" 2>/dev/null; do
rss_kb=$(ps -o rss= -p "$API_PID" 2>/dev/null | tr -d ' ' || echo 0)
peak=$(cat "$peak_file")
if (( rss_kb > peak )); then echo "$rss_kb" > "$peak_file"; fi
sleep 0.5
done
) &
MEM_PID=$!
}
stop_mem_monitor() {
[[ -n "${MEM_PID:-}" ]] && kill "$MEM_PID" 2>/dev/null || true
MEM_PID=""
API_PID=""
}
mem_rss_kb() {
[[ -z "${API_PID:-}" ]] && echo 0 && return
ps -o rss= -p "$API_PID" 2>/dev/null | tr -d ' ' || echo 0
}
mem_peak_kb() {
local peak_file="$BENCH_DIR/.mem_peak"
[[ -f "$peak_file" ]] && cat "$peak_file" || echo 0
}
fmt_mem_mb() {
perl -e "printf '%.1f MB', $1 / 1024"
}
# ─── Backend lifecycle ──────────────────────────────────────────────────────────
start_backend() {
local name="$1"
echo " $(dim "Starting ${name}...")"
# Ensure port is free before starting
lsof -ti ":$API_PORT" -sTCP:LISTEN 2>/dev/null | xargs kill 2>/dev/null || true
local elapsed=0
while nc -z "$API_HOST" "$API_PORT" 2>/dev/null; do
sleep 0.2
elapsed=$((elapsed + 1))
if (( elapsed > 50 )); then
echo "ERROR: port $API_PORT still in use after 10 seconds"
exit 1
fi
done
case "$name" in
restic-api)
(cd "$REPO_ROOT" && exec pnpm --filter restic-api start 2>&1) > "$BENCH_DIR/backend.log" 2>&1 &
MANAGED_PID=$!
;;
michael)
(cd "$REPO_ROOT/packages/michael" && \
OTLP_METRICS_ENDPOINT=localhost:8428 \
OTLP_METRICS_URL_PATH=/opentelemetry/api/v1/push \
exec go run . 2>&1) > "$BENCH_DIR/backend.log" 2>&1 &
MANAGED_PID=$!
;;
*)
echo "ERROR: unknown backend '$name' (expected restic-api or michael)"
exit 1
;;
esac
local elapsed=0
while ! nc -z "$API_HOST" "$API_PORT" 2>/dev/null; do
sleep 0.2
elapsed=$((elapsed + 1))
if (( elapsed > 150 )); then
echo "ERROR: $name did not start within 30 seconds"
echo "--- backend log ---"
cat "$BENCH_DIR/backend.log"
exit 1
fi
done
echo " $(dim "${name} ready on :${API_PORT}")"
}
stop_backend() {
if [[ -n "${MANAGED_PID:-}" ]]; then
kill "$MANAGED_PID" 2>/dev/null || true
wait "$MANAGED_PID" 2>/dev/null || true
MANAGED_PID=""
fi
# Also kill anything still on the port (child processes)
lsof -ti ":$API_PORT" -sTCP:LISTEN 2>/dev/null | xargs kill 2>/dev/null || true
# Wait for port to free up
local elapsed=0
while nc -z "$API_HOST" "$API_PORT" 2>/dev/null; do
sleep 0.2
elapsed=$((elapsed + 1))
if (( elapsed > 50 )); then break; fi
done
}
# ─── Prerequisites ──────────────────────────────────────────────────────────────
check_prereqs() {
local missing=()
command -v restic >/dev/null 2>&1 || missing+=("restic")
command -v node >/dev/null 2>&1 || missing+=("node")
command -v perl >/dev/null 2>&1 || missing+=("perl")
if [[ "$COMPARE" == true ]]; then
command -v pnpm >/dev/null 2>&1 || missing+=("pnpm")
command -v go >/dev/null 2>&1 || missing+=("go")
fi
if (( ${#missing[@]} > 0 )); then
echo "ERROR: missing required tools: ${missing[*]}"
exit 1
fi
}
check_api_reachable() {
local http_code
http_code=$(curl -s -o /dev/null -w "%{http_code}" --connect-timeout 3 \
"http://${API_HOST}:${API_PORT}/" 2>/dev/null || echo "000")
if [[ "$http_code" == "000" ]]; then
echo "ERROR: backend not reachable at http://${API_HOST}:${API_PORT}"
echo " Start it with: mise restic-api:dev or mise michael:dev"
exit 1
fi
}
# ─── Data generation ────────────────────────────────────────────────────────────
generate_data() {
echo "$(bold 'Generating test data...')"
local dir="$BENCH_DIR/data/small" && mkdir -p "$dir"
for i in $(seq 1 "$SMALL_COUNT"); do
dd if=/dev/urandom of="$dir/f_$i" bs=1024 count="$SMALL_SIZE_KB" 2>/dev/null
done
echo " small: ${SMALL_COUNT} files x ${SMALL_SIZE_KB} KB"
dir="$BENCH_DIR/data/medium" && mkdir -p "$dir"
for i in $(seq 1 "$MEDIUM_COUNT"); do
dd if=/dev/urandom of="$dir/f_$i" bs=1048576 count="$MEDIUM_SIZE_MB" 2>/dev/null
done
echo " medium: ${MEDIUM_COUNT} files x ${MEDIUM_SIZE_MB} MB"
dir="$BENCH_DIR/data/large" && mkdir -p "$dir"
for i in $(seq 1 "$LARGE_COUNT"); do
dd if=/dev/urandom of="$dir/f_$i" bs=1048576 count="$LARGE_SIZE_MB" 2>/dev/null
done
echo " large: ${LARGE_COUNT} files x ${LARGE_SIZE_MB} MB"
local total_bytes=$(( SMALL_COUNT * SMALL_SIZE_KB * 1024 + MEDIUM_COUNT * MEDIUM_SIZE_MB * 1048576 + LARGE_COUNT * LARGE_SIZE_MB * 1048576 ))
local peak_bytes=$(( total_bytes * 3 ))
echo ""
echo " $(bold "total: $(human_size $total_bytes)") $(dim "(~$(human_size $peak_bytes) peak disk + S3)")"
echo ""
}
# ─── Throughput scenario ────────────────────────────────────────────────────────
# run_throughput <tag> <label> <data_dir> <total_bytes> <file_count>
run_throughput() {
local tag="$1" label="$2" data_dir="$3" total_bytes="$4" file_count="$5"
local restic_log="$BENCH_DIR/restic_${tag}.log"
echo "$(bold " $label")"
init_repo
# init
local t0 t1
t0=$(now_ms)
if ! restic init -q 2>"$restic_log"; then
echo " $(red 'ERROR: restic init failed')"
cat "$restic_log" | sed 's/^/ /' >&2
return 1
fi
t1=$(now_ms)
echo " init: $(elapsed_s "$t0" "$t1")s"
# backup
t0=$(now_ms)
if ! restic backup -q "$data_dir" 2>"$restic_log"; then
echo " $(red 'ERROR: restic backup failed')"
cat "$restic_log" | sed 's/^/ /' >&2
return 1
fi
t1=$(now_ms)
local backup_s backup_tp
backup_s=$(elapsed_s "$t0" "$t1")
backup_tp=$(throughput_mbs "$total_bytes" "$backup_s")
echo " backup: ${backup_s}s @ ${backup_tp} MB/s"
save_result "${tag}_backup_s" "$backup_s"
save_result "${tag}_backup_mbs" "$backup_tp"
if (( file_count > 1 )); then
local ms_per_file
ms_per_file=$(perl -e "printf '%.1f', ($t1 - $t0) / $file_count")
echo " $(dim " ~${ms_per_file} ms/file")"
fi
# restore
local restore_dir="$BENCH_DIR/restore_${RANDOM}"
mkdir -p "$restore_dir"
t0=$(now_ms)
if ! restic restore latest --target "$restore_dir" -q 2>"$restic_log"; then
echo " $(red 'ERROR: restic restore failed')"
cat "$restic_log" | sed 's/^/ /' >&2
rm -rf "$restore_dir"
return 1
fi
t1=$(now_ms)
local restore_s restore_tp
restore_s=$(elapsed_s "$t0" "$t1")
restore_tp=$(throughput_mbs "$total_bytes" "$restore_s")
echo " restore: ${restore_s}s @ ${restore_tp} MB/s"
save_result "${tag}_restore_s" "$restore_s"
save_result "${tag}_restore_mbs" "$restore_tp"
rm -rf "$restore_dir"
# incremental (re-backup identical data — exercises index/dedup path)
t0=$(now_ms)
if ! restic backup -q "$data_dir" 2>"$restic_log"; then
echo " $(red 'ERROR: restic incremental backup failed')"
cat "$restic_log" | sed 's/^/ /' >&2
return 1
fi
t1=$(now_ms)
local incr_s
incr_s=$(elapsed_s "$t0" "$t1")
echo " incremental: ${incr_s}s (no new data)"
save_result "${tag}_incr_s" "$incr_s"
# memory
if [[ -n "${API_PID:-}" ]]; then
echo " memory: $(fmt_mem_mb "$(mem_rss_kb)") current, $(fmt_mem_mb "$(mem_peak_kb)") peak"
fi
echo ""
}
# ─── Run full benchmark suite for one backend ──────────────────────────────────
CURRENT_BACKEND=""
run_suite() {
local backend="$1"
CURRENT_BACKEND="$backend"
mkdir -p "$BENCH_DIR/results"
: > "$BENCH_DIR/results/${backend}.txt"
echo ""
echo "$(bold '════════════════════════════════════════════════════════════')"
echo "$(bold " ${backend}")"
echo "$(bold '════════════════════════════════════════════════════════════')"
echo ""
# memory baseline
stop_mem_monitor
echo "0" > "$BENCH_DIR/.mem_peak"
start_mem_monitor
local baseline_kb
baseline_kb=$(mem_rss_kb)
save_result "mem_baseline_kb" "$baseline_kb"
echo "$(bold 'Memory')"
echo " baseline: $(fmt_mem_mb "$baseline_kb") (PID ${API_PID:-?})"
echo ""
local small_bytes=$((SMALL_COUNT * SMALL_SIZE_KB * 1024))
local medium_bytes=$((MEDIUM_COUNT * MEDIUM_SIZE_MB * 1048576))
local large_bytes=$((LARGE_COUNT * LARGE_SIZE_MB * 1048576))
echo "$(bold 'Throughput')"
echo ""
run_throughput "small" \
"small (${SMALL_COUNT} x ${SMALL_SIZE_KB} KB = $(human_size $small_bytes))" \
"$BENCH_DIR/data/small" "$small_bytes" "$SMALL_COUNT"
run_throughput "medium" \
"medium (${MEDIUM_COUNT} x ${MEDIUM_SIZE_MB} MB = $(human_size $medium_bytes))" \
"$BENCH_DIR/data/medium" "$medium_bytes" "$MEDIUM_COUNT"
run_throughput "large" \
"large (${LARGE_COUNT} x ${LARGE_SIZE_MB} MB = $(human_size $large_bytes))" \
"$BENCH_DIR/data/large" "$large_bytes" "$LARGE_COUNT"
# save final memory
save_result "mem_peak_kb" "$(mem_peak_kb)"
}
# ─── Print single-backend summary ──────────────────────────────────────────────
print_summary() {
local f="$BENCH_DIR/results/$1.txt"
echo "$(bold '════════════════════════════════════════════════════════════')"
echo "$(bold "Summary: $1")"
echo ""
printf ' %-28s %10s %10s %10s %10s\n' "" "backup" "restore" "up MB/s" "down MB/s"
printf ' %-28s %10s %10s %10s %10s\n' "" "──────" "───────" "───────" "─────────"
for tag in small medium large; do
printf ' %-28s %9ss %9ss %9s %9s\n' \
"$tag" \
"$(read_result "${tag}_backup_s" "$f")" \
"$(read_result "${tag}_restore_s" "$f")" \
"$(read_result "${tag}_backup_mbs" "$f")" \
"$(read_result "${tag}_restore_mbs" "$f")"
done
echo ""
local peak_kb
peak_kb=$(read_result mem_peak_kb "$f")
local base_kb
base_kb=$(read_result mem_baseline_kb "$f")
echo " memory: baseline $(fmt_mem_mb "$base_kb"), peak $(fmt_mem_mb "$peak_kb")"
if [[ "$DO_CLEANUP" == false ]]; then
echo ""
echo " $(dim "Repos kept (--no-cleanup): ${REPOS[*]}")"
fi
echo ""
}
# ─── Print comparison table ────────────────────────────────────────────────────
# pct_delta <a> <b> → prints "(+12.3%)" or "(-12.3%)"
# positive means b > a
pct_delta() {
perl -e '
my ($a, $b) = @ARGV;
if ($a == 0) { print "n/a"; exit }
my $d = ($b - $a) / abs($a) * 100;
printf "%+.1f%%", $d;
' "$1" "$2"
}
# color_delta <a> <b> <higher_is_better>
# prints the delta in green if improvement, red if regression
color_delta() {
local a="$1" b="$2" higher_better="${3:-false}"
local raw
raw=$(pct_delta "$a" "$b")
[[ "$raw" == "n/a" ]] && printf '%s' "$raw" && return
local numeric
numeric=$(perl -e '
my ($a, $b) = @ARGV;
if ($a == 0) { print 0; exit }
printf "%.1f", ($b - $a) / abs($a) * 100;
' "$a" "$b")
local is_better
if [[ "$higher_better" == true ]]; then
is_better=$(perl -e "print $numeric > 0 ? 1 : 0")
else
is_better=$(perl -e "print $numeric < 0 ? 1 : 0")
fi
if (( is_better )); then
green "$raw"
else
# within 5% = dim, otherwise red
local abs_val
abs_val=$(perl -e "printf '%.1f', abs($numeric)")
if perl -e "exit( $abs_val < 5.0 ? 0 : 1 )"; then
dim "$raw"
else
red "$raw"
fi
fi
}
print_comparison() {
local fa="$BENCH_DIR/results/restic-api.txt"
local fb="$BENCH_DIR/results/michael.txt"
echo ""
echo "$(bold '════════════════════════════════════════════════════════════════════════════════')"
echo "$(bold 'Comparison: restic-api vs michael')"
echo "$(bold '════════════════════════════════════════════════════════════════════════════════')"
echo ""
# Header
printf ' %-24s %16s %16s %10s\n' "" "restic-api" "michael" "delta"
printf ' %-24s %16s %16s %10s\n' "" "──────────" "───────" "─────"
# Throughput rows — show MB/s (higher is better)
for tag in small medium large; do
for op in backup restore; do
local key="${tag}_${op}_mbs"
local va vb
va=$(read_result "$key" "$fa")
vb=$(read_result "$key" "$fb")
local delta
delta=$(color_delta "$va" "$vb" true)
printf ' %-24s %12s MB/s %12s MB/s %b\n' \
"${tag} ${op}" "$va" "$vb" "$delta"
done
local key="${tag}_incr_s"
local va vb
va=$(read_result "$key" "$fa")
vb=$(read_result "$key" "$fb")
local delta
delta=$(color_delta "$va" "$vb" false)
printf ' %-24s %15ss %15ss %b\n' \
"${tag} incremental" "$va" "$vb" "$delta"
done
echo ""
# Memory — lower is better
local va vb delta
va=$(read_result mem_baseline_kb "$fa")
vb=$(read_result mem_baseline_kb "$fb")
delta=$(color_delta "$va" "$vb" false)
printf ' %-24s %16s %16s %b\n' \
"memory baseline" "$(fmt_mem_mb "$va")" "$(fmt_mem_mb "$vb")" "$delta"
va=$(read_result mem_peak_kb "$fa")
vb=$(read_result mem_peak_kb "$fb")
delta=$(color_delta "$va" "$vb" false)
printf ' %-24s %16s %16s %b\n' \
"memory peak" "$(fmt_mem_mb "$va")" "$(fmt_mem_mb "$vb")" "$delta"
echo ""
}
# ─── Main ───────────────────────────────────────────────────────────────────────
main() {
echo ""
echo "$(bold 'restic backend benchmark')"
echo "$(dim '────────────────────────────────────────────────────────────')"
echo " restic $(restic version 2>&1 | head -1)"
echo " node $(node -v)"
if [[ "$COMPARE" == true ]]; then
echo " go $(go version 2>&1 | awk '{print $3}')"
echo " mode comparison (restic-api vs michael)"
else
echo " backend ${BACKEND}"
fi
echo " api http://${API_HOST}:${API_PORT}"
echo " tmpdir ${BENCH_DIR}"
echo ""
check_prereqs
generate_data
if [[ "$COMPARE" == true ]]; then
for backend in restic-api michael; do
start_backend "$backend"
run_suite "$backend"
stop_mem_monitor
stop_backend
done
print_comparison
else
if nc -z "$API_HOST" "$API_PORT" 2>/dev/null; then
echo " $(dim "Using already-running backend on :${API_PORT}")"
else
start_backend "$BACKEND"
fi
run_suite "$BACKEND"
stop_mem_monitor
stop_backend
print_summary "$BACKEND"
fi
}
main "$@"
+1 -1
View File
@@ -3,7 +3,7 @@
#MISE dir="{{config_root}}/packages/michael"
set -e
# Complements the fleet-level `mise bench` (real restic CLI against a running
# Complements `yuctl tools bench` (real restic CLI against a running
# backend): these measure per-request cost inside the process.
#
# The source-network benchmarks resolve against the committed three-network
+1 -1
View File
@@ -1,7 +1,7 @@
#!/usr/bin/env bash
#MISE description="Run michael in development mode"
set -e
source "$(dirname "$0")/../restic-api/env"
source "$(dirname "$0")/env"
# No load balancer in front of a local run, so there is nothing to drain away
# from — skip the fleet's 15s SIGTERM drain and let Ctrl-C exit immediately.
-6
View File
@@ -1,6 +0,0 @@
#!/usr/bin/env bash
#MISE description="Build restic-api"
#MISE depends=["common:build"]
set -e
NODE_ENV=production pnpm --filter restic-api build
-6
View File
@@ -1,6 +0,0 @@
#!/usr/bin/env bash
#MISE description="Run restic-api in development mode"
set -e
source "$(dirname "$0")/env"
pnpm --filter restic-api start:dev
-6
View File
@@ -1,6 +0,0 @@
#!/usr/bin/env bash
#MISE description="Run unit tests for restic-api"
set -e
source "$(dirname "$0")/../env"
pnpm --filter restic-api test "$@"
-6
View File
@@ -1,6 +0,0 @@
#!/usr/bin/env bash
#MISE description="Run integration tests for restic-api"
set -e
source "$(dirname "$0")/../env"
pnpm --filter restic-api test:integration "$@"
-6
View File
@@ -1,6 +0,0 @@
#!/usr/bin/env bash
#MISE description="Run unit tests in watch mode for restic-api"
set -e
source "$(dirname "$0")/../env"
pnpm --filter restic-api test --watch "$@"
+1 -1
View File
@@ -2,7 +2,7 @@
#MISE description="Run all end-to-end tests"
#MISE depends=["test:e2e:wait"]
set -e
source "$(dirname "$0")/../../restic-api/env"
source "$(dirname "$0")/../../michael/env"
source "$(dirname "$0")/../../yucca-api/env"
pnpm --filter e2e test "$@"
+2 -2
View File
@@ -3,7 +3,7 @@
# Targets the compose/`mise dev` flow (the k3d flow has its own readiness
# handling in packages/e2e/k3d/run.sh and does not use this task).
set -e
source "$(dirname "$0")/../../restic-api/env"
source "$(dirname "$0")/../../michael/env"
source "$(dirname "$0")/../../yucca-api/env"
echo Ensure dev environment is running before invoking tests.
@@ -24,7 +24,7 @@ wait_for_port() {
}
wait_for_port 8092 mock-oidc-provider
wait_for_port "$RESTIC_API_PORT" restic-api
wait_for_port "$RESTIC_API_PORT" michael
wait_for_port "$YUCCA_API_PORT" yucca-api
wait_for_port 22676 orchestration-api
wait_for_port "${WEB_PORT:-36033}" web
+1 -1
View File
@@ -62,4 +62,4 @@ export RADOS_ENDPOINT RADOS_ACCESS_KEY_ID RADOS_SECRET_ACCESS_KEY
echo "==> run the S3-backed integration suites"
# Concurrent: each scopes itself to freshly generated repository/bucket names.
mise run michael:test:integration ::: restic-api:test:integration ::: yucca-metrics-worker:test:integration
mise run michael:test:integration ::: yucca-metrics-worker:test:integration
-1
View File
@@ -109,7 +109,6 @@ Zod-validated `env.ts`, JWT auth guards via `@AuthRoute()`, OTel from `@common/s
| `michael` | Go | **Production** restic REST backend — S3 proxy with JWT verification, WORM enforcement, backend pooling. |
| `columbo` | Go | Ticket investigation agent: LLM loop (OpenRouter) over per-user-scoped o11y queries, answers only into staff threads. See `docs/columbo.md`. |
| `monk` | Go | Ceph scrub-backlog exporter: polls `pg ls`, serves measured per-pool scrub-age metrics on :9284. Deployed onto mon hosts by ansible, not K8s. |
| `restic-api` | NestJS | Earlier TS implementation of the restic backend, kept as **reference**; not deployed. |
| `yucca-metrics-worker` | NestJS | 5-min cron: RadosGW usage → meter tables → per-connection rollup (`connectionMetrics`, billing floor), OTel gauges. |
| `redis` (valkey) | | Shared platform cache (ephemeral; keys `yucca:<service>:<purpose>:*`). Primary-region only. |
| `mock-oidc-provider` | Node | Dev/test OIDC IdP (code + device flow). |
+2 -2
View File
@@ -112,9 +112,9 @@ wait_for_http yucca-api http://localhost:3020/api/meta 200
# call, and only answers once discovery has reached a live yucca-api.
wait_for_http orchestration-api http://localhost:22676/api/yucca/onboarding 200
echo "==> jest e2e (it-works, restic-api, yucca-api, orchestration-api)"
echo "==> jest e2e (it-works, restic, yucca-api, orchestration-api)"
# shellcheck disable=SC1091
source .mise/tasks/restic-api/env
source .mise/tasks/michael/env
# shellcheck disable=SC1091
source .mise/tasks/yucca-api/env
# Three, not one per suite: orchestration-api alone runs ~105s and gates the
-56
View File
@@ -1,56 +0,0 @@
# compiled output
/dist
/node_modules
/build
# Logs
logs
*.log
npm-debug.log*
pnpm-debug.log*
yarn-debug.log*
yarn-error.log*
lerna-debug.log*
# OS
.DS_Store
# Tests
/coverage
/.nyc_output
# IDEs and editors
/.idea
.project
.classpath
.c9/
*.launch
.settings/
*.sublime-workspace
# IDE - VSCode
.vscode/*
!.vscode/settings.json
!.vscode/tasks.json
!.vscode/launch.json
!.vscode/extensions.json
# dotenv environment variable files
.env
.env.development.local
.env.test.local
.env.production.local
.env.local
# temp directory
.temp
.tmp
# Runtime data
pids
*.pid
*.seed
*.pid.lock
# Diagnostic reports (https://nodejs.org/api/report.html)
report.[0-9]*.[0-9]*.[0-9]*.[0-9]*.json
-8
View File
@@ -1,8 +0,0 @@
{
"$schema": "https://json.schemastore.org/nest-cli",
"collection": "@nestjs/schematics",
"sourceRoot": "src",
"compilerOptions": {
"deleteOutDir": true
}
}
-77
View File
@@ -1,77 +0,0 @@
{
"name": "restic-api",
"version": "0.43.0",
"description": "",
"author": "",
"private": true,
"license": "UNLICENSED",
"files": [
"dist"
],
"scripts": {
"build": "nest build",
"start": "nest start",
"start:dev": "nest start --watch",
"start:debug": "nest start --debug --watch",
"start:prod": "node dist/main",
"test": "jest",
"test:watch": "jest --watch",
"test:cov": "jest --coverage",
"test:debug": "node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand",
"test:integration": "jest --config ./test/jest-integration.json"
},
"dependencies": {
"@aws-sdk/client-s3": "catalog:",
"@aws-sdk/lib-storage": "catalog:",
"@common/server": "workspace:^",
"@nestjs/common": "catalog:",
"@nestjs/core": "catalog:",
"@nestjs/jwt": "catalog:",
"@nestjs/platform-express": "catalog:",
"@opentelemetry/api": "catalog:",
"nestjs-zod": "catalog:",
"reflect-metadata": "catalog:",
"rxjs": "catalog:",
"zod": "catalog:"
},
"devDependencies": {
"@nestjs/cli": "catalog:",
"@nestjs/schematics": "catalog:",
"@nestjs/testing": "catalog:",
"@types/express": "catalog:",
"@types/jest": "catalog:",
"@types/node": "catalog:",
"@types/supertest": "catalog:",
"globals": "catalog:",
"jest": "catalog:",
"source-map-support": "catalog:",
"supertest": "catalog:",
"ts-jest": "catalog:",
"ts-loader": "catalog:",
"ts-node": "catalog:",
"tsconfig-paths": "catalog:",
"typescript": "catalog:",
"typescript-eslint": "catalog:",
"vitest": "catalog:"
},
"jest": {
"moduleFileExtensions": [
"js",
"json",
"ts"
],
"rootDir": "src",
"testRegex": ".*\\.spec\\.ts$",
"transform": {
"^.+\\.(t|j)s$": "ts-jest"
},
"moduleNameMapper": {
"^src/(.*)$": "<rootDir>/$1"
},
"collectCoverageFrom": [
"**/*.(t|j)s"
],
"coverageDirectory": "../coverage",
"testEnvironment": "node"
}
}
-52
View File
@@ -1,52 +0,0 @@
import {
LoggerRepository,
LoggingInterceptor,
OtelModule,
shutdownOtel,
WideContextRepository,
} from '@common/server/otel';
import { Module, type OnApplicationShutdown, Provider } from '@nestjs/common';
import { APP_GUARD, APP_INTERCEPTOR, APP_PIPE } from '@nestjs/core';
import { JwtModule } from '@nestjs/jwt';
import { ZodSerializerInterceptor, ZodValidationPipe } from 'nestjs-zod';
import { AppController } from './controllers/app.controller';
import { env } from './env';
import { AuthGuard } from './middleware/auth.guard';
import { ResticInterceptor } from './middleware/restic.interceptor';
import { StorageRepository } from './repositories/storage.repository';
import { AppService } from './services/app.service';
import { AuthService } from './services/auth.service';
export const imports = [
JwtModule.register({
global: true,
publicKey: env.JWT_PUBLIC_KEY,
verifyOptions: { algorithms: ['ES256'] },
}),
];
export const controllers = [AppController];
export const providers: Provider[] = [
WideContextRepository,
LoggerRepository,
StorageRepository,
AuthService,
AppService,
{ provide: APP_GUARD, useClass: AuthGuard },
{ provide: APP_INTERCEPTOR, useClass: ResticInterceptor },
{ provide: APP_INTERCEPTOR, useClass: LoggingInterceptor },
{ provide: APP_INTERCEPTOR, useClass: ZodSerializerInterceptor },
{ provide: APP_PIPE, useClass: ZodValidationPipe },
];
@Module({
imports: [OtelModule, ...imports],
controllers,
providers,
})
export class AppModule implements OnApplicationShutdown {
async onApplicationShutdown() {
await shutdownOtel();
}
}
@@ -1,121 +0,0 @@
import { Traceable } from '@common/server/otel';
import {
Controller,
Delete,
Get,
Head,
Headers,
HttpCode,
HttpStatus,
Param,
ParseBoolPipe,
Post,
Query,
Req,
Res,
} from '@nestjs/common';
import { type Request, type Response } from 'express';
import { BlobInfoResponseDto } from 'src/dto/app.dto';
import { type AuthDto } from 'src/dto/auth.dto';
import { Auth, AuthRoute } from 'src/middleware/auth.guard';
import { ResticRoute } from 'src/middleware/restic.interceptor';
import { AppService } from 'src/services/app.service';
import { respondWithObject } from 'src/utils/s3';
import { BlobParamsDto, BlobWithNameParamsDto } from 'src/validation';
@Traceable()
@Controller()
export class AppController {
constructor(private readonly service: AppService) {}
@Post(':path')
@AuthRoute()
@HttpCode(HttpStatus.OK)
async createRepository(@Auth() auth: AuthDto, @Query('create', ParseBoolPipe) isCreate: boolean): Promise<void> {
await this.service.createRepository(auth.repository, isCreate);
}
@Delete(':path')
@AuthRoute()
@HttpCode(HttpStatus.NOT_IMPLEMENTED)
deleteRepository(@Auth() _auth: AuthDto): void {
this.service.deleteRepository();
}
@Head(':path/config')
@AuthRoute()
async checkConfig(@Auth() auth: AuthDto, @Res() res: Response): Promise<void> {
const size = await this.service.checkConfig(auth);
res.set('Content-Length', String(size)).end();
}
@Get(':path/config')
@AuthRoute()
async getConfig(@Auth() auth: AuthDto, @Req() req: Request, @Res() res: Response): Promise<void> {
const config = await this.service.getConfig(auth);
respondWithObject(config, req, res);
}
@Post(':path/config')
@AuthRoute()
@HttpCode(HttpStatus.OK)
async saveConfig(@Auth() auth: AuthDto, @Req() req: Request): Promise<void> {
await this.service.saveConfig(auth, req);
}
@Delete(':path/config')
@AuthRoute()
@HttpCode(HttpStatus.OK)
async deleteConfig(@Auth() auth: AuthDto): Promise<void> {
await this.service.deleteConfig(auth);
}
@Get(':path/:type')
@AuthRoute()
@ResticRoute()
async listBlobs(@Auth() auth: AuthDto, @Param() { type }: BlobParamsDto): Promise<BlobInfoResponseDto[]> {
return this.service.listBlobs(auth, type);
}
@Head(':path/:type/:name')
@AuthRoute()
async checkBlob(
@Auth() auth: AuthDto,
@Param() { type, name }: BlobWithNameParamsDto,
@Res() res: Response,
): Promise<void> {
const size = await this.service.checkBlob(auth, type, name);
res.set('Content-Length', String(size)).end();
}
@Get(':path/:type/:name')
@AuthRoute()
async getBlob(
@Auth() auth: AuthDto,
@Param() { type, name }: BlobWithNameParamsDto,
@Headers('range') range: string | undefined,
@Req() req: Request,
@Res() res: Response,
): Promise<void> {
const blob = await this.service.getBlob(auth, type, name, range);
respondWithObject(blob, req, res);
}
@Post(':path/:type/:name')
@AuthRoute()
@HttpCode(HttpStatus.OK)
async saveBlob(
@Auth() auth: AuthDto,
@Param() { type, name }: BlobWithNameParamsDto,
@Req() req: Request,
): Promise<void> {
await this.service.saveBlob(auth, type, name, req);
}
@Delete(':path/:type/:name')
@AuthRoute()
@HttpCode(HttpStatus.OK)
async deleteBlob(@Auth() auth: AuthDto, @Param() { type, name }: BlobWithNameParamsDto): Promise<void> {
await this.service.deleteBlob(auth, type, name);
}
}
-11
View File
@@ -1,11 +0,0 @@
import { createZodDto } from 'nestjs-zod';
import { z } from 'zod';
const BlobInfoResponseSchema = z
.object({
name: z.string(),
size: z.number(),
})
.meta({ id: 'BlobInfoResponseDto' });
export class BlobInfoResponseDto extends createZodDto(BlobInfoResponseSchema) {}
-12
View File
@@ -1,12 +0,0 @@
import { createZodDto } from 'nestjs-zod';
import { z } from 'zod';
const AuthSchema = z
.object({
user: z.uuid(),
repository: z.uuid(),
writeOnce: z.boolean(),
})
.meta({ id: 'AuthDto' });
export class AuthDto extends createZodDto(AuthSchema) {}
-17
View File
@@ -1,17 +0,0 @@
export enum ContentType {
Binary = 'application/octet-stream',
ResticV2 = 'application/vnd.x.restic.rest.v2',
}
export enum MetadataKey {
Auth = 'AUTH',
Restic = 'RESTIC_V2',
}
export enum BlobType {
Data = 'data',
Index = 'index',
Keys = 'keys',
Locks = 'locks',
Snapshots = 'snapshots',
}
-24
View File
@@ -1,24 +0,0 @@
import { z } from 'zod';
const schema = z.object({
NODE_ENV: z.enum(['development', 'production', 'test', 'provision']).default('development'),
RESTIC_API_PORT: z.coerce.number().min(1000),
JWT_PUBLIC_KEY: z.string(),
S3_ACCESS_KEY_ID: z.string(),
S3_SECRET_ACCESS_KEY: z.string(),
S3_REGION: z.string(),
S3_ENDPOINT: z.string(),
S3_FORCE_PATH_STYLE: z.coerce.boolean().default(false),
OTEL_DEBUG: z.coerce.boolean().default(false),
OTEL_SAMPLE_RATE: z.number().min(0).max(1).default(1),
OTEL_METRICS_EXPORT_INTERVAL: z.number().default(10_000),
OTEL_METRICS: z.string().default('http://localhost:8428/opentelemetry/v1/metrics'),
OTEL_TRACING: z.string().default('http://localhost:10428/insert/opentelemetry/v1/traces'),
OTEL_LOGGING: z.string().default('http://localhost:9428/insert/opentelemetry/v1/logs'),
});
export const env = schema.parse(process.env);
-13
View File
@@ -1,13 +0,0 @@
import { BadRequestException, InternalServerErrorException } from '@nestjs/common';
export class S3Error extends InternalServerErrorException {
constructor() {
super('An error occurred with the storage server');
}
}
export class ChecksumMismatchError extends BadRequestException {
constructor() {
super('Content hash does not match blob name');
}
}
-13
View File
@@ -1,13 +0,0 @@
import '@common/server/otel';
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
import { env } from './env';
async function bootstrap() {
const app = await NestFactory.create(AppModule);
app.enableShutdownHooks();
await app.listen(env.RESTIC_API_PORT);
}
void bootstrap();
@@ -1,57 +0,0 @@
import {
BadRequestException,
CanActivate,
ExecutionContext,
Injectable,
Scope,
SetMetadata,
applyDecorators,
createParamDecorator,
} from '@nestjs/common';
import { Reflector } from '@nestjs/core';
import { Request } from 'express';
import { type AuthDto } from 'src/dto/auth.dto';
import { MetadataKey } from 'src/enum';
import { AuthService } from 'src/services/auth.service';
export const AuthRoute = (options = {}): MethodDecorator => {
return applyDecorators(SetMetadata(MetadataKey.Auth, options));
};
export interface AuthRequest extends Request {
auth?: AuthDto;
}
export interface AuthenticatedRequest extends Request {
auth: AuthDto;
}
export const Auth = createParamDecorator((_, context: ExecutionContext): AuthDto => {
return context.switchToHttp().getRequest<AuthenticatedRequest>().auth;
});
@Injectable({ scope: Scope.REQUEST })
export class AuthGuard implements CanActivate {
constructor(
private reflector: Reflector,
private service: AuthService,
) {}
async canActivate(context: ExecutionContext): Promise<boolean> {
const targets = [context.getHandler()];
const options = this.reflector.getAllAndOverride<{ _emptyObject: never } | undefined>(MetadataKey.Auth, targets);
if (!options) {
return true;
}
const request = context.switchToHttp().getRequest<AuthRequest>();
request.auth = await this.service.authenticate(request.headers);
const path = request.params.path;
if (path && path !== request.auth.repository) {
throw new BadRequestException('Repository mismatch');
}
return true;
}
}
@@ -1,44 +0,0 @@
import {
CallHandler,
ExecutionContext,
Injectable,
NestInterceptor,
NotImplementedException,
SetMetadata,
applyDecorators,
} from '@nestjs/common';
import { Reflector } from '@nestjs/core';
import { Request, Response } from 'express';
import { Observable, map } from 'rxjs';
import { ContentType, MetadataKey } from 'src/enum';
export const ResticRoute = (): MethodDecorator => {
return applyDecorators(SetMetadata(MetadataKey.Restic, true));
};
@Injectable()
export class ResticInterceptor implements NestInterceptor {
constructor(private reflector: Reflector) {}
intercept(context: ExecutionContext, next: CallHandler): Observable<unknown> {
const isResticV2 = this.reflector.get<boolean>(MetadataKey.Restic, context.getHandler());
if (!isResticV2) {
return next.handle();
}
const request = context.switchToHttp().getRequest<Request>();
const response = context.switchToHttp().getResponse<Response>();
if (request.headers.accept !== ContentType.ResticV2) {
throw new NotImplementedException();
}
return next.handle().pipe(
map((data) => {
response.setHeader('Content-Type', ContentType.ResticV2);
response.end(JSON.stringify(data));
return void 0;
}),
);
}
}
@@ -1,128 +0,0 @@
import {
CreateBucketCommand,
DeleteObjectCommand,
GetObjectCommand,
HeadBucketCommand,
HeadObjectCommand,
ListObjectsV2Command,
NotFound,
PutObjectCommandInput,
S3Client,
} from '@aws-sdk/client-s3';
import { Upload } from '@aws-sdk/lib-storage';
import { Traceable } from '@common/server/otel';
import { Injectable } from '@nestjs/common';
import { createHash } from 'node:crypto';
import { Readable, Transform } from 'node:stream';
import { env } from 'src/env';
import { ChecksumMismatchError } from 'src/errors';
@Traceable()
@Injectable()
export class StorageRepository {
private client: S3Client;
constructor() {
this.client = new S3Client({
credentials: {
accessKeyId: env.S3_ACCESS_KEY_ID,
secretAccessKey: env.S3_SECRET_ACCESS_KEY,
},
region: env.S3_REGION,
endpoint: env.S3_ENDPOINT,
forcePathStyle: env.S3_FORCE_PATH_STYLE,
});
}
async checkBucket(Bucket: string): Promise<boolean> {
try {
await this.client.send(new HeadBucketCommand({ Bucket }));
return true;
} catch (error) {
if (error instanceof NotFound) {
return false;
}
throw error;
}
}
createBucket(Bucket: string) {
return this.client.send(
new CreateBucketCommand({
Bucket,
}),
);
}
async putObject(Bucket: string, Key: string, Body: Readable, writeOnce: boolean, sha256Hex?: string) {
const hash = sha256Hex ? createHash('sha256') : undefined;
const params: PutObjectCommandInput = {
Bucket,
Key,
Body: hash
? Body.pipe(
new Transform({
transform(chunk, _, callback) {
hash.update(chunk);
callback(null, chunk);
},
}),
)
: Body,
};
if (writeOnce) {
params.IfNoneMatch = '*';
}
const upload = new Upload({
client: this.client,
params,
});
await upload.done();
if (hash && hash.digest('hex') !== sha256Hex) {
await this.deleteObject(Bucket, Key);
throw new ChecksumMismatchError();
}
}
headObject(Bucket: string, Key: string) {
return this.client.send(
new HeadObjectCommand({
Bucket,
Key,
}),
);
}
async listObjects(Bucket: string, Prefix: string) {
return await this.client.send(
new ListObjectsV2Command({
Bucket,
Prefix,
}),
);
}
async getObject(Bucket: string, Key: string, Range?: string) {
return await this.client.send(
new GetObjectCommand({
Bucket,
Key,
Range,
}),
);
}
deleteObject(Bucket: string, Key: string) {
return this.client.send(
new DeleteObjectCommand({
Bucket,
Key,
}),
);
}
}
@@ -1,404 +0,0 @@
import { S3ServiceException } from '@aws-sdk/client-s3';
import { Readable } from 'node:stream';
import { text } from 'node:stream/consumers';
import { AuthDto } from 'src/dto/auth.dto';
import { BlobType } from 'src/enum';
import { type Mocks, newMocks } from '../../test/mocks';
import { AppService } from './app.service';
const mockAuth = (writeOnce = false): AuthDto => ({
user: 'user',
repository: 'repository',
writeOnce,
});
describe(AppService.name, () => {
let mocks: Mocks;
let sut: AppService;
beforeEach(() => {
mocks = newMocks();
sut = new AppService(mocks.storage as never, mocks.metricService, mocks.wideContext);
});
it('should exist', () => {
expect(sut).toBeDefined();
});
describe('createRepository', () => {
it('should fail if isCreate is false', async () => {
await expect(sut.createRepository('repository', false)).rejects.toThrow();
expect(mocks.storage.checkBucket).toHaveBeenCalledTimes(0);
});
it('should fail if bucket exists', async () => {
mocks.storage.checkBucket.mockResolvedValue(true);
await expect(sut.createRepository('repository', true)).rejects.toThrow();
expect(mocks.storage.checkBucket).toHaveBeenCalled();
expect(mocks.storage.createBucket).toHaveBeenCalledTimes(0);
});
it('should succeed if bucket does not exist', async () => {
mocks.storage.checkBucket.mockResolvedValue(false);
await sut.createRepository('repository', true);
expect(mocks.storage.checkBucket).toHaveBeenCalled();
expect(mocks.storage.createBucket).toHaveBeenCalled();
});
it('should fail if S3 command throws', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.checkBucket.mockRejectedValueOnce(S3Error);
await expect(sut.createRepository('repository', true)).rejects.toThrow();
expect(mocks.storage.checkBucket).toHaveBeenCalled();
expect(mocks.storage.createBucket).toHaveBeenCalledTimes(0);
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
it('should fail if S3 command throws', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.createBucket.mockRejectedValueOnce(S3Error);
await expect(sut.createRepository('repository', true)).rejects.toThrow();
expect(mocks.storage.checkBucket).toHaveBeenCalled();
expect(mocks.storage.createBucket).toHaveBeenCalled();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('deleteRepository', () => {
it('should do nothing', () => {
sut.deleteRepository();
});
});
describe('checkConfig', () => {
it('should return content length', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 123, $metadata: void 0 as never });
const result = await sut.checkConfig(mockAuth());
expect(result).toBe(123);
expect(mocks.storage.headObject).toHaveBeenCalledWith('repository', 'config');
});
it('should return 0 if content length is undefined', async () => {
mocks.storage.headObject.mockResolvedValue({ $metadata: void 0 as never });
const result = await sut.checkConfig(mockAuth());
expect(result).toBe(0);
});
it('should throw if headObject fails', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.headObject.mockRejectedValue(S3Error);
await expect(sut.checkConfig(mockAuth())).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('getConfig', () => {
it('should return the stream', async () => {
const result = await sut.getConfig(mockAuth());
expect(result).toEqual(
expect.objectContaining({
object: expect.objectContaining({
ContentLength: expect.any(Number),
}),
}),
);
expect(mocks.storage.getObject).toHaveBeenCalledWith('repository', 'config');
const data = await text(result.stream()!);
expect(data).toHaveLength(result.object.ContentLength!);
expect(sut.blobsDownloadedBytes.add).toHaveBeenCalledWith(
result.object.ContentLength,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsDownloadedBytes.add).toHaveBeenCalledTimes(1);
expect(sut.blobsRequestedBytes.add).toHaveBeenCalledWith(
result.object.ContentLength,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsRequestedBytes.add).toHaveBeenCalledTimes(1);
});
it('should throw if getObject fails', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.getObject.mockImplementation(() => new Promise((_, reject) => reject(S3Error)));
await expect(sut.getConfig(mockAuth())).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('saveConfig', () => {
it('should save config', async () => {
const body = Readable.from('body');
await sut.saveConfig(mockAuth(), body as never);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'config', expect.anything(), false);
expect(sut.blobsUploadedBytes.add).toHaveBeenCalledWith(
4,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsUploadedBytes.add).toHaveBeenCalledTimes(1);
});
it('should pass writeOnce flag', async () => {
const body = Readable.from('body');
await sut.saveConfig(mockAuth(true), body as never);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'config', expect.anything(), true);
});
it('should throw on 412 error', async () => {
const body = Readable.from('body');
const error = new S3ServiceException({
name: 'PreconditionFailed',
$fault: 'client',
$metadata: { httpStatusCode: 412 },
});
mocks.storage.putObject.mockRejectedValue(error);
await expect(sut.saveConfig(mockAuth(true), body as never)).rejects.toThrow('Config already exists');
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(error);
});
it('should throw on other errors', async () => {
const body = Readable.from('body');
const S3Error = Symbol('S3Error');
mocks.storage.putObject.mockRejectedValue(S3Error);
await expect(sut.saveConfig(mockAuth(), body as never)).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('deleteConfig', () => {
it('should delete config', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 500, $metadata: void 0 as never });
mocks.storage.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteConfig(mockAuth());
expect(mocks.storage.headObject).toHaveBeenCalledWith('repository', 'config');
expect(mocks.storage.deleteObject).toHaveBeenCalledWith('repository', 'config');
});
it('should throw when writeOnce and not locks', async () => {
await expect(sut.deleteConfig(mockAuth(true))).rejects.toThrow();
expect(mocks.storage.deleteObject).not.toHaveBeenCalled();
});
it('should throw if deleteObject fails', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 500, $metadata: void 0 as never });
const S3Error = Symbol('S3Error');
mocks.storage.deleteObject.mockRejectedValue(S3Error);
await expect(sut.deleteConfig(mockAuth())).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('listBlobs', () => {
it('should return mapped blobs', async () => {
mocks.storage.listObjects.mockResolvedValue({
Contents: [
{ Key: 'data/abc123', Size: 100 },
{ Key: 'data/def456', Size: 200 },
],
KeyCount: 2,
$metadata: void 0 as never,
});
const result = await sut.listBlobs(mockAuth(), BlobType.Data);
expect(result).toEqual([
{ name: 'abc123', size: 100 },
{ name: 'def456', size: 200 },
]);
expect(mocks.storage.listObjects).toHaveBeenCalledWith('repository', 'data/');
});
it('should return empty array when KeyCount is 0', async () => {
mocks.storage.listObjects.mockResolvedValue({ KeyCount: 0, $metadata: void 0 as never });
const result = await sut.listBlobs(mockAuth(), BlobType.Data);
expect(result).toEqual([]);
});
it('should throw if Contents is undefined', async () => {
mocks.storage.listObjects.mockResolvedValue({ KeyCount: 1, $metadata: void 0 as never });
await expect(sut.listBlobs(mockAuth(), BlobType.Data)).rejects.toThrow();
});
it('should throw if Key or Size is missing', async () => {
const Contents = [{ Key: 'data/abc123' }];
mocks.storage.listObjects.mockResolvedValue({
Contents,
KeyCount: 1,
$metadata: void 0 as never,
});
await expect(sut.listBlobs(mockAuth(), BlobType.Data)).rejects.toThrow();
expect(mocks.wideContext.addContext).toHaveBeenCalledWith('contents', Contents);
});
it('should throw if listObjects fails', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.listObjects.mockRejectedValue(S3Error);
await expect(sut.listBlobs(mockAuth(), BlobType.Data)).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('checkBlob', () => {
it('should return content length', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 456, $metadata: void 0 as never });
const result = await sut.checkBlob(mockAuth(), BlobType.Data, 'abc123');
expect(result).toBe(456);
expect(mocks.storage.headObject).toHaveBeenCalledWith('repository', 'data/abc123');
});
it('should return 0 if content length is undefined', async () => {
mocks.storage.headObject.mockResolvedValue({ $metadata: void 0 as never });
const result = await sut.checkBlob(mockAuth(), BlobType.Data, 'abc123');
expect(result).toBe(0);
});
it('should throw if headObject fails', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.headObject.mockRejectedValue(S3Error);
await expect(sut.checkBlob(mockAuth(), BlobType.Data, 'abc123')).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('getBlob', () => {
it('should return the object', async () => {
const result = await sut.getBlob(mockAuth(), BlobType.Data, 'abc123');
expect(result).toEqual(
expect.objectContaining({
object: expect.objectContaining({
ContentLength: expect.any(Number),
}),
}),
);
expect(mocks.storage.getObject).toHaveBeenCalledWith('repository', 'data/abc123', undefined);
const data = await text(result.stream()!);
expect(data).toHaveLength(result.object.ContentLength!);
expect(sut.blobsDownloadedBytes.add).toHaveBeenCalledWith(
result.object.ContentLength,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsDownloadedBytes.add).toHaveBeenCalledTimes(1);
expect(sut.blobsRequestedBytes.add).toHaveBeenCalledWith(
result.object.ContentLength,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsRequestedBytes.add).toHaveBeenCalledTimes(1);
});
it('should pass range to getObjectStream', async () => {
await sut.getBlob(mockAuth(), BlobType.Data, 'abc123', 'bytes=0-100');
expect(mocks.storage.getObject).toHaveBeenCalledWith('repository', 'data/abc123', 'bytes=0-100');
});
it('should throw if getObject fails', async () => {
const S3Error = Symbol('S3Error');
mocks.storage.getObject.mockImplementation(() => new Promise((_, reject) => reject(S3Error)));
await expect(sut.getBlob(mockAuth(), BlobType.Data, 'abc123')).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('saveBlob', () => {
it('should save blob', async () => {
const body = Readable.from('body');
await sut.saveBlob(mockAuth(), BlobType.Data, 'abc123', body as never);
expect(mocks.storage.putObject).toHaveBeenCalledWith(
'repository',
'data/abc123',
expect.anything(),
true,
'abc123',
);
expect(sut.blobsUploadedBytes.add).toHaveBeenCalledWith(
4,
expect.objectContaining({
customerId: 'user',
repositoryId: 'repository',
}),
);
expect(sut.blobsUploadedBytes.add).toHaveBeenCalledTimes(1);
});
it('should pass writeOnce flag', async () => {
const body = Readable.from('body');
await sut.saveBlob(mockAuth(true), BlobType.Data, 'abc123', body as never);
expect(mocks.storage.putObject).toHaveBeenCalledWith(
'repository',
'data/abc123',
expect.anything(),
true,
'abc123',
);
});
it('should throw ConflictException on 412 error', async () => {
const body = Readable.from('body');
const error = new S3ServiceException({
name: 'PreconditionFailed',
$fault: 'client',
$metadata: { httpStatusCode: 412 },
});
mocks.storage.putObject.mockRejectedValue(error);
await expect(sut.saveBlob(mockAuth(true), BlobType.Data, 'abc123', body as never)).rejects.toThrow(
'Blob already exists',
);
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(error);
});
it('should throw on other errors', async () => {
const body = Readable.from('body');
const S3Error = Symbol('S3Error');
mocks.storage.putObject.mockRejectedValue(S3Error);
await expect(sut.saveBlob(mockAuth(), BlobType.Data, 'abc123', body as never)).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
describe('deleteBlob', () => {
it('should delete blob', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 500, $metadata: void 0 as never });
mocks.storage.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteBlob(mockAuth(), BlobType.Data, 'abc123');
expect(mocks.storage.headObject).toHaveBeenCalledWith('repository', 'data/abc123');
expect(mocks.storage.deleteObject).toHaveBeenCalledWith('repository', 'data/abc123');
});
it('should allow delete of locks with writeOnce', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 500, $metadata: void 0 as never });
mocks.storage.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteBlob(mockAuth(true), BlobType.Locks, 'abc123');
expect(mocks.storage.deleteObject).toHaveBeenCalledWith('repository', 'locks/abc123');
});
it('should throw when writeOnce and not locks', async () => {
await expect(sut.deleteBlob(mockAuth(true), BlobType.Data, 'abc123')).rejects.toThrow();
expect(mocks.storage.deleteObject).not.toHaveBeenCalled();
});
it('should throw if deleteObject fails', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 500, $metadata: void 0 as never });
const S3Error = Symbol('S3Error');
mocks.storage.deleteObject.mockRejectedValue(S3Error);
await expect(sut.deleteBlob(mockAuth(), BlobType.Data, 'abc123')).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
});
@@ -1,254 +0,0 @@
import { S3ServiceException } from '@aws-sdk/client-s3';
import { MetricService, Traceable, WideContextRepository } from '@common/server/otel';
import {
BadRequestException,
ConflictException,
ForbiddenException,
Injectable,
NotFoundException,
} from '@nestjs/common';
import { Counter, UpDownCounter } from '@opentelemetry/api';
import { Readable } from 'node:stream';
import { BlobInfoResponseDto } from 'src/dto/app.dto';
import { AuthDto } from 'src/dto/auth.dto';
import { BlobType } from 'src/enum';
import { ChecksumMismatchError, S3Error } from 'src/errors';
import { StorageRepository } from 'src/repositories/storage.repository';
import { attachMeterToStream, contextFromAuth } from 'src/utils/meters';
import { attachMeterToS3GetObject, S3RemoteObject } from 'src/utils/s3';
@Traceable()
@Injectable()
export class AppService {
blobsRequestedBytes: Counter;
blobsDownloadedBytes: Counter;
blobsUploadedBytes: Counter;
blobsStoredBytes: UpDownCounter;
constructor(
private readonly storage: StorageRepository,
private readonly metricService: MetricService,
private readonly wideContext: WideContextRepository,
) {
this.blobsRequestedBytes = this.metricService.getCounter('blobs.requested_bytes', {
description: 'Total no. of blob bytes requested for download',
});
this.blobsDownloadedBytes = this.metricService.getCounter('blobs.downloaded_bytes', {
description: 'Total no. of blob bytes download',
});
this.blobsUploadedBytes = this.metricService.getCounter('blobs.uploaded_bytes', {
description: 'Total no. of blob bytes uploaded',
});
this.blobsStoredBytes = this.metricService.getUpDownCounter('blobs.stored_bytes', {
description: 'Total no. of blob bytes stored',
});
}
async createRepository(repository: string, isCreate: boolean): Promise<void> {
if (!isCreate) {
throw new BadRequestException('isCreate must be true when creating repository');
}
let exists: boolean;
try {
exists = await this.storage.checkBucket(repository);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
if (exists) {
throw new ConflictException('Repository already exists');
}
try {
await this.storage.createBucket(repository);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
deleteRepository(): void {}
async checkConfig(auth: AuthDto): Promise<number> {
try {
const { ContentLength } = await this.storage.headObject(auth.repository, 'config');
return ContentLength || 0;
} catch (error) {
this.wideContext.setErrorCause(error);
throw new NotFoundException();
}
}
async getConfig(auth: AuthDto): Promise<S3RemoteObject> {
try {
return attachMeterToS3GetObject(
auth,
await this.storage.getObject(auth.repository, 'config'),
this.blobsRequestedBytes,
this.blobsDownloadedBytes,
);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async saveConfig(auth: AuthDto, body: Readable): Promise<void> {
try {
let stored = 0;
await this.storage.putObject(
auth.repository,
'config',
attachMeterToStream(body, this.blobsUploadedBytes, (bytes) => (stored += bytes), contextFromAuth(auth)),
auth.writeOnce,
);
this.blobsStoredBytes.add(stored, contextFromAuth(auth));
} catch (error) {
this.wideContext.setErrorCause(error);
if (error instanceof S3ServiceException && error.$metadata.httpStatusCode === 412) {
throw new ForbiddenException('Config already exists');
}
throw new S3Error();
}
}
async deleteConfig(auth: AuthDto): Promise<void> {
if (auth.writeOnce) {
throw new ForbiddenException('Not permitted to write to WORM repository');
}
try {
const { ContentLength } = await this.storage.headObject(auth.repository, 'config');
if (!ContentLength) {
throw 'Missing ContentLength from headObject output';
}
await this.storage.deleteObject(auth.repository, 'config');
this.blobsStoredBytes.add(-ContentLength, contextFromAuth(auth));
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async listBlobs(auth: AuthDto, type: BlobType): Promise<BlobInfoResponseDto[]> {
try {
const suffix = `${type}/`;
const { Contents, KeyCount } = await this.storage.listObjects(auth.repository, suffix);
if (KeyCount === 0) {
return [];
}
if (!Contents) {
this.wideContext.setErrorCause('Contents missing from ListObjects response');
throw void 0;
}
if (Contents.some(({ Key, Size }) => !Key || !Size)) {
this.wideContext.setErrorCause('Contents are malformed from ListObjects response');
this.wideContext.addContext('contents', Contents);
throw void 0;
}
return Contents!.map(({ Key, Size }) => ({
name: Key!.slice(suffix.length),
size: Size!,
}));
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async checkBlob(auth: AuthDto, type: BlobType, name: string): Promise<number> {
try {
const { ContentLength } = await this.storage.headObject(auth.repository, `${type}/${name}`);
return ContentLength || 0;
} catch (error) {
this.wideContext.setErrorCause(error);
throw new NotFoundException();
}
}
async getBlob(auth: AuthDto, type: BlobType, name: string, range?: string): Promise<S3RemoteObject> {
try {
return attachMeterToS3GetObject(
auth,
await this.storage.getObject(auth.repository, `${type}/${name}`, range),
this.blobsRequestedBytes,
this.blobsDownloadedBytes,
);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async saveBlob(auth: AuthDto, type: BlobType, name: string, body: Readable): Promise<void> {
try {
let stored = 0;
await this.storage.putObject(
auth.repository,
`${type}/${name}`,
attachMeterToStream(body, this.blobsUploadedBytes, (bytes) => (stored += bytes), contextFromAuth(auth)),
true,
name,
);
this.blobsStoredBytes.add(stored, contextFromAuth(auth));
} catch (error) {
this.wideContext.setErrorCause(error);
if (error instanceof ChecksumMismatchError) {
throw error;
}
if (error instanceof S3ServiceException) {
if (error.$metadata.httpStatusCode === 412) {
throw new ForbiddenException('Blob already exists');
}
if (error.name === 'XAmzContentChecksumMismatch') {
throw new BadRequestException('Content hash does not match blob name');
}
}
throw new S3Error();
}
}
async deleteBlob(auth: AuthDto, type: BlobType, name: string): Promise<void> {
if (auth.writeOnce && type !== BlobType.Locks) {
throw new ForbiddenException('Not permitted to write to WORM repository');
}
try {
const { ContentLength } = await this.storage.headObject(auth.repository, `${type}/${name}`);
if (!ContentLength) {
throw 'Missing ContentLength from headObject output';
}
await this.storage.deleteObject(auth.repository, `${type}/${name}`);
this.blobsStoredBytes.add(-ContentLength, contextFromAuth(auth));
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
}
@@ -1,55 +0,0 @@
import { randomUUID } from 'node:crypto';
import { newJwtMock, newWideContextMock } from '../../test/mocks';
import { AuthService } from './auth.service';
describe(AuthService.name, () => {
let jwt: ReturnType<typeof newJwtMock>;
let wideContext: ReturnType<typeof newWideContextMock>;
let sut: AuthService;
beforeEach(() => {
jwt = newJwtMock();
wideContext = newWideContextMock();
sut = new AuthService(jwt as never, wideContext as never);
});
it('should exist', () => {
expect(sut).toBeDefined();
});
describe('authenticate', () => {
it('should throw if Authorization header is missing', async () => {
await expect(sut.authenticate({})).rejects.toThrow('Missing Authorization header');
});
it('should throw if not Basic auth', async () => {
await expect(sut.authenticate({ authorization: 'Bearer token' })).rejects.toThrow('Expected Basic auth');
});
it('should throw if token is missing', async () => {
const auth = Buffer.from('username').toString('base64');
await expect(sut.authenticate({ authorization: `Basic ${auth}` })).rejects.toThrow('Expected Basic auth token');
});
it('should throw if JWT verification fails', async () => {
jwt.verifyAsync.mockRejectedValue(new Error('invalid'));
const auth = Buffer.from('username:token').toString('base64');
await expect(sut.authenticate({ authorization: `Basic ${auth}` })).rejects.toThrow('Invalid JWT Token');
});
it('should throw if JWT payload is invalid', async () => {
jwt.verifyAsync.mockResolvedValue({ invalid: 'payload' });
const auth = Buffer.from('username:token').toString('base64');
await expect(sut.authenticate({ authorization: `Basic ${auth}` })).rejects.toThrow('Bad Request Exception');
});
it('should return auth dto on success', async () => {
const payload = { user: randomUUID(), repository: randomUUID(), writeOnce: true };
jwt.verifyAsync.mockResolvedValue(payload);
const auth = Buffer.from('username:token').toString('base64');
const result = await sut.authenticate({ authorization: `Basic ${auth}` });
expect(result).toEqual(payload);
expect(jwt.verifyAsync).toHaveBeenCalledWith('token');
});
});
});
@@ -1,46 +0,0 @@
import { WideContextRepository } from '@common/server/otel';
import { BadRequestException, Injectable, UnauthorizedException } from '@nestjs/common';
import { JwtService } from '@nestjs/jwt';
import { type IncomingHttpHeaders } from 'node:http';
import { AuthDto } from 'src/dto/auth.dto';
import { contextFromAuth } from 'src/utils/meters';
@Injectable()
export class AuthService {
constructor(
private readonly jwt: JwtService,
private readonly wideContext: WideContextRepository,
) {}
async authenticate(headers: IncomingHttpHeaders): Promise<AuthDto> {
if (!headers.authorization) {
throw new UnauthorizedException('Missing Authorization header');
}
if (!headers.authorization.startsWith('Basic ')) {
throw new UnauthorizedException('Expected Basic auth');
}
const auth = Buffer.from(headers.authorization.split(' ').pop() || '', 'base64').toString();
const [_, token] = auth.split(':');
if (!token) {
throw new UnauthorizedException('Expected Basic auth token');
}
let jwt;
try {
jwt = await this.jwt.verifyAsync(token);
} catch {
throw new UnauthorizedException('Invalid JWT Token');
}
const result = AuthDto.schema.safeParse(jwt);
if (!result.success) {
throw new BadRequestException(result.error.issues.map((issue) => issue.message));
}
this.wideContext.assignContext(contextFromAuth(result.data));
return result.data;
}
}
-29
View File
@@ -1,29 +0,0 @@
import { Attributes, Counter } from '@opentelemetry/api';
import { PassThrough, Readable } from 'node:stream';
import { AuthDto } from 'src/dto/auth.dto';
export function attachMeterToStream(
stream: Readable,
meterBytes: Counter,
callbackBytes: (bytes: number) => void,
context: Attributes,
) {
const passthrough = new PassThrough({
transform(chunk, _, callback) {
meterBytes.add(chunk.length, context);
callbackBytes(chunk.length);
callback(null, chunk);
},
});
stream.pipe(passthrough);
return passthrough;
}
export function contextFromAuth(auth: AuthDto) {
return {
customerId: auth.user,
repositoryId: auth.repository,
};
}
-80
View File
@@ -1,80 +0,0 @@
import { GetObjectCommandOutput } from '@aws-sdk/client-s3';
import { HttpStatus } from '@nestjs/common';
import { Counter } from '@opentelemetry/api';
import { Request, Response } from 'express';
import { Readable } from 'node:stream';
import { ReadableStream } from 'node:stream/web';
import { AuthDto } from 'src/dto/auth.dto';
import { ContentType } from '../enum';
import { attachMeterToStream, contextFromAuth } from './meters';
export interface S3RemoteObject {
stream(): Readable | undefined;
object: GetObjectCommandOutput;
}
export function attachMeterToS3GetObject(
auth: AuthDto,
object: GetObjectCommandOutput,
requestedBytes: Counter,
downloadedBytes: Counter,
): S3RemoteObject {
const context = contextFromAuth(auth);
requestedBytes.add(object.ContentLength || 0, context);
return {
object,
stream() {
const webStream = object.Body?.transformToWebStream();
if (webStream) {
return attachMeterToStream(
Readable.fromWeb(webStream as ReadableStream),
downloadedBytes,
() => void 0,
context,
);
}
},
};
}
/**
* Process S3 object as web response
*
* References:
* http#ServeContent
* https://pkg.go.dev/net/http#ServeContent
*/
export function respondWithObject({ stream, object }: S3RemoteObject, request: Request, response: Response) {
if (request.headers['if-none-match'] === object.ETag) {
return response.send(HttpStatus.NOT_MODIFIED);
}
const range = request.headers.range;
if (range && range !== 'bytes=0-') {
response.status(HttpStatus.PARTIAL_CONTENT);
} else {
response.status(HttpStatus.OK);
}
if (object.ETag) {
response.header('ETag', object.ETag);
}
response.set('Content-Type', object.ContentType ?? ContentType.Binary);
if (object.ContentRange) {
response.set('Content-Range', object.ContentRange);
}
if (object.ContentLength) {
response.set('Content-Length', object.ContentLength.toString());
}
const readable = stream();
if (readable) {
readable.pipe(response);
} else {
return response.send(HttpStatus.INTERNAL_SERVER_ERROR);
}
}
-12
View File
@@ -1,12 +0,0 @@
import { createZodDto } from 'nestjs-zod';
import { z } from 'zod';
import { BlobType } from './enum';
const BlobParamsSchema = z.object({ type: z.enum(BlobType) }).meta({ id: 'BlobParamsDto' });
const BlobWithNameParamsSchema = BlobParamsSchema.extend({
name: z.string().regex(/^[a-f0-9]{64}$/),
}).meta({ id: 'BlobWithNameParamsDto' });
export class BlobParamsDto extends createZodDto(BlobParamsSchema) {}
export class BlobWithNameParamsDto extends createZodDto(BlobWithNameParamsSchema) {}
@@ -1,401 +0,0 @@
import { MetricService } from '@common/server/otel';
import { INestApplication } from '@nestjs/common';
import { JwtService } from '@nestjs/jwt';
import { Test, TestingModule } from '@nestjs/testing';
import { createHash, randomUUID } from 'node:crypto';
import request from 'supertest';
import { App } from 'supertest/types';
import { controllers, imports, providers } from '../src/app.module';
import { newMetricServiceMock } from './mocks';
const makeAuthHeader = (token: string) => 'Basic ' + Buffer.from(`_:${token}`).toString('base64');
describe('AppController (e2e)', () => {
let app: INestApplication<App>;
let repository: string;
let authHeader: string;
let wormAuthHeader: string;
beforeAll(async () => {
const moduleFixture: TestingModule = await Test.createTestingModule({
imports,
controllers,
providers: [MetricService, ...providers],
})
.overrideProvider(MetricService)
.useValue(newMetricServiceMock())
.compile();
app = moduleFixture.createNestApplication();
await app.init();
});
afterAll(async () => {
await app.close();
});
beforeEach(async () => {
const signer = new JwtService({
privateKey: process.env.JWT_PRIVATE_KEY,
signOptions: { algorithm: 'ES256' },
});
repository = randomUUID();
authHeader = makeAuthHeader(signer.sign({ user: randomUUID(), repository, writeOnce: false }));
wormAuthHeader = makeAuthHeader(signer.sign({ user: randomUUID(), repository, writeOnce: true }));
});
describe('POST /:path', () => {
it('fails if create is not true', async () => {
await request(app.getHttpServer())
.post(`/${repository}?create=false`)
.set('Authorization', authHeader)
.expect(400);
});
it('creates the repository', async () => {
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
});
it('fails if repository already exists', async () => {
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(409);
});
});
describe('DELETE /:path', () => {
it('returns not implemented', async () => {
await request(app.getHttpServer()).delete(`/${repository}`).set('Authorization', authHeader).expect(501);
});
});
describe('HEAD /:path/config', () => {
it('returns 404 if config does not exist', async () => {
await request(app.getHttpServer()).head(`/${repository}/config`).set('Authorization', authHeader).expect(404);
});
it('returns content-length if config exists', async () => {
const config = 'test-config-data';
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/config`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(config))
.expect(200);
await request(app.getHttpServer())
.head(`/${repository}/config`)
.set('Authorization', authHeader)
.expect(200)
.expect('Content-Length', String(config.length));
});
});
describe('GET, POST /:path/config', () => {
it('saves and returns config data', async () => {
const config = 'test-config-data';
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/config`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(config))
.expect(200);
const response = await request(app.getHttpServer())
.get(`/${repository}/config`)
.set('Authorization', authHeader)
.expect(200)
.expect('Content-Type', 'application/octet-stream');
expect(response.body.toString()).toBe(config);
});
it('fails to overwrite config with worm', async () => {
const config = 'test-config-data';
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', wormAuthHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/config`)
.set('Authorization', wormAuthHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(config))
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/config`)
.set('Authorization', wormAuthHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(config))
.expect(403);
});
});
describe('GET /:path/:type', () => {
it('returns 501 without restic v2 accept header', async () => {
await request(app.getHttpServer()).get(`/${repository}/data`).set('Authorization', authHeader).expect(501);
});
it('returns empty list when no blobs exist', async () => {
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
const response = await request(app.getHttpServer())
.get(`/${repository}/data`)
.set('Authorization', authHeader)
.set('Accept', 'application/vnd.x.restic.rest.v2')
.expect(200)
.expect('Content-Type', 'application/vnd.x.restic.rest.v2');
expect(JSON.parse(response.text)).toEqual([]);
});
it('returns list of blobs when they exist', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
const response = await request(app.getHttpServer())
.get(`/${repository}/data`)
.set('Authorization', authHeader)
.set('Accept', 'application/vnd.x.restic.rest.v2')
.expect(200)
.expect('Content-Type', 'application/vnd.x.restic.rest.v2');
expect(JSON.parse(response.text)).toEqual([{ name: blobName, size: blobData.length }]);
});
});
describe('HEAD /:path/:type/:name', () => {
it('returns 404 if blob does not exist', async () => {
await request(app.getHttpServer())
.head(`/${repository}/data/${'a'.repeat(64)}`)
.set('Authorization', authHeader)
.expect(404);
});
it('returns content-length if blob exists', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
await request(app.getHttpServer())
.head(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.expect(200)
.expect('Content-Length', String(blobData.length));
});
});
describe('GET, POST /:path/:type/:name', () => {
it('saves and returns blob data', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
const response = await request(app.getHttpServer())
.get(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.expect(200)
.expect('Content-Type', 'application/octet-stream');
expect(response.body.toString()).toBe(blobData);
});
it('returns partial content with range header', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
const response = await request(app.getHttpServer())
.get(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Range', 'bytes=0-3')
.expect(206)
.expect('Content-Type', 'application/octet-stream');
expect(response.body.toString()).toBe('test');
});
it('fails to overwrite blob with worm', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', wormAuthHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', wormAuthHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', wormAuthHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(403);
});
it('fails to write blob with non-hash for name', async () => {
const blobData = 'test-blob-data';
const blobName = 'invalid-hash';
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(400);
});
it('fails to write blob with non-matching hash for name', async () => {
const blobData = 'test-blob-data';
const blobName = 'a'.repeat(64);
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(400);
});
});
describe('DELETE /:path/:type/:name', () => {
it('deletes a blob', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
await request(app.getHttpServer())
.delete(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.expect(200);
await request(app.getHttpServer())
.head(`/${repository}/data/${blobName}`)
.set('Authorization', authHeader)
.expect(404);
});
it('fails to delete blob with worm', async () => {
const blobData = 'test-blob-data';
const blobName = createHash('sha256').update(blobData).digest('hex');
await request(app.getHttpServer())
.post(`/${repository}?create=true`)
.set('Authorization', wormAuthHeader)
.expect(200);
await request(app.getHttpServer())
.post(`/${repository}/data/${blobName}`)
.set('Authorization', wormAuthHeader)
.set('Content-Type', 'application/octet-stream')
.send(Buffer.from(blobData))
.expect(200);
await request(app.getHttpServer())
.delete(`/${repository}/data/${blobName}`)
.set('Authorization', wormAuthHeader)
.expect(403);
});
});
});
@@ -1,12 +0,0 @@
{
"moduleFileExtensions": ["js", "json", "ts"],
"rootDir": ".",
"testEnvironment": "node",
"testRegex": ".integration-spec.ts$",
"transform": {
"^.+\\.(t|j)s$": "ts-jest"
},
"moduleNameMapper": {
"^src/(.*)$": "<rootDir>/../src/$1"
}
}
-71
View File
@@ -1,71 +0,0 @@
import type { LoggerRepository, MetricService, WideContextRepository } from '@common/server/otel';
import { Readable } from 'node:stream';
import { text } from 'node:stream/consumers';
import { StorageRepository } from 'src/repositories/storage.repository';
export type RepositoryInterface<T extends object> = Pick<T, keyof T>;
export const newJwtMock = () => ({
verifyAsync: jest.fn(),
});
export const newWideContextMock = () => ({
assignContext: jest.fn(),
});
export const newLoggerRepositoryMock = (): jest.Mocked<RepositoryInterface<LoggerRepository>> => {
return {
debug: jest.fn(),
error: jest.fn(),
info: jest.fn(),
warn: jest.fn(),
};
};
export const newStorageRepositoryMock = (): jest.Mocked<RepositoryInterface<StorageRepository>> => {
return {
checkBucket: jest.fn(),
createBucket: jest.fn(),
deleteObject: jest.fn(),
getObject: jest.fn().mockResolvedValue({
ContentLength: 1000,
Body: {
transformToWebStream() {
return Readable.toWeb(Readable.from('_'.repeat(1000)));
},
},
}),
headObject: jest.fn(),
listObjects: jest.fn(),
putObject: jest.fn().mockImplementation((_1, _2, body: Readable) => text(body)),
};
};
export const newWideContextRepositoryMock = (): jest.Mocked<RepositoryInterface<WideContextRepository>> => ({
context: {},
addContext: jest.fn(),
applyContext: jest.fn(),
assignContext: jest.fn(),
setErrorCause: jest.fn(),
});
export const newMetricServiceMock = (): jest.Mocked<RepositoryInterface<MetricService>> => ({
getCounter: jest.fn().mockImplementation(() => ({ add: jest.fn() })),
getGauge: jest.fn(),
getHistogram: jest.fn(),
getObservableCounter: jest.fn(),
getObservableGauge: jest.fn(),
getObservableUpDownCounter: jest.fn(),
getUpDownCounter: jest.fn().mockImplementation(() => ({ add: jest.fn() })),
});
export const newMocks = () => {
return {
logger: newLoggerRepositoryMock(),
storage: newStorageRepositoryMock(),
metricService: newMetricServiceMock(),
wideContext: newWideContextRepositoryMock(),
};
};
export type Mocks = ReturnType<typeof newMocks>;
@@ -1,93 +0,0 @@
import { createHash, randomBytes, randomUUID } from 'node:crypto';
import { Readable } from 'node:stream';
import { text } from 'node:stream/consumers';
import { ReadableStream } from 'node:stream/web';
import { StorageRepository } from 'src/repositories/storage.repository';
describe('StorageRepository (e2e)', () => {
let sut: StorageRepository;
beforeEach(() => {
sut = new StorageRepository();
});
it('works', async () => {
const Bucket = randomUUID();
const Key = randomUUID();
const contents = 'test object';
await expect(sut.checkBucket(Bucket)).resolves.toBe(false);
await sut.createBucket(Bucket);
await expect(sut.checkBucket(Bucket)).resolves.toBe(true);
await expect(sut.headObject(Bucket, Key)).rejects.toThrow();
await sut.putObject(Bucket, Key, Readable.from(contents), false);
await sut.headObject(Bucket, Key);
await expect(sut.listObjects(Bucket, Key)).resolves.toEqual(
expect.objectContaining({
Contents: expect.arrayContaining([expect.objectContaining({ Key, Size: contents.length })]),
}),
);
const object = await sut.getObject(Bucket, Key);
const stream = Readable.fromWeb(object.Body!.transformToWebStream() as ReadableStream);
await expect(text(stream!)).resolves.toBe(contents);
await sut.deleteObject(Bucket, Key);
await expect(sut.headObject(Bucket, Key)).rejects.toThrow();
});
it('respects worm', async () => {
const Bucket = randomUUID();
const Key = randomUUID();
const contents = 'test object';
await sut.createBucket(Bucket);
await sut.putObject(Bucket, Key, Readable.from(contents), false);
await expect(sut.putObject(Bucket, Key, Readable.from(contents), true)).rejects.toThrow();
});
it('checks hash', async () => {
const Bucket = randomUUID();
const Key = randomUUID();
const contents = randomBytes(5_000_000);
const hash = createHash('sha256');
hash.update(contents);
await sut.createBucket(Bucket);
await sut.putObject(Bucket, Key, Readable.from(contents), false, hash.digest('hex'));
await sut.deleteObject(Bucket, Key);
});
it('rejects invalid hash', async () => {
const Bucket = randomUUID();
const Key = randomUUID();
const contents = randomBytes(5_000_000);
const hash = createHash('sha256');
hash.update('invalid');
await sut.createBucket(Bucket);
await expect(sut.putObject(Bucket, Key, Readable.from(contents), false, hash.digest('hex'))).rejects.toThrow();
});
it('checks hash where files are chunked', async () => {
const Bucket = randomUUID();
const Key = randomUUID();
const contents = randomBytes(20_000_000);
const hash = createHash('sha256');
hash.update(contents);
await sut.createBucket(Bucket);
await sut.putObject(Bucket, Key, Readable.from(contents), false, hash.digest('hex'));
await sut.deleteObject(Bucket, Key);
});
});
-4
View File
@@ -1,4 +0,0 @@
{
"extends": "./tsconfig.json",
"exclude": ["node_modules", "test", "dist", "**/*spec.ts"]
}
-29
View File
@@ -1,29 +0,0 @@
{
"compilerOptions": {
"module": "nodenext",
"moduleResolution": "nodenext",
"resolvePackageJsonExports": true,
"esModuleInterop": true,
"isolatedModules": true,
"declaration": true,
"removeComments": true,
"emitDecoratorMetadata": true,
"experimentalDecorators": true,
"allowSyntheticDefaultImports": true,
"target": "ES2023",
"sourceMap": true,
"outDir": "./dist",
"incremental": true,
"skipLibCheck": true,
"strictNullChecks": true,
"forceConsistentCasingInFileNames": true,
"noImplicitAny": true,
"strictBindCallApply": false,
"noFallthroughCasesInSwitch": false,
"paths": {
"src/*": ["./src/*"]
},
"types": ["jest"]
},
"include": ["src", "test"]
}
-1353
View File
File diff suppressed because it is too large Load Diff
-2
View File
@@ -11,8 +11,6 @@ onlyBuiltDependencies:
- unrs-resolver
catalog:
'@aws-sdk/client-s3': ^3.965.0
'@aws-sdk/lib-storage': ^3.965.0
'@better-svelte-email/cli': ^2.1.3
'@better-svelte-email/components': ^2.1.3
'@better-svelte-email/server': ^2.1.3
-5
View File
@@ -37,11 +37,6 @@
"path": "packages/mock-oidc-provider/package.json",
"jsonpath": "$.version"
},
{
"type": "json",
"path": "packages/restic-api/package.json",
"jsonpath": "$.version"
},
{
"type": "json",
"path": "packages/web/package.json",