mirror of
https://github.com/immich-app/yucca.git
synced 2026-09-30 13:33:00 +08:00
feat(orchestrator): meta discovery and placement (#385)
This commit is contained in:
@@ -2,5 +2,4 @@
|
||||
#MISE description="Test yucca-sdk/orchestration-api"
|
||||
set -e
|
||||
|
||||
# pnpm --filter @futo-org/backups-orchestrator-api test "$@"
|
||||
echo No tests provided.
|
||||
pnpm --filter @futo-org/backups-orchestrator-api test "$@"
|
||||
|
||||
@@ -10,7 +10,8 @@
|
||||
# (hostmap.cjs); the playwright browser via --host-resolver-rules.
|
||||
# - yucca-api's topology is patched (e2e-only) to a localhost rest_url,
|
||||
# because restic runs on this host (the cluster keeps the in-cluster name
|
||||
# for real use).
|
||||
# for real use); the orchestrator discovers everything through the meta
|
||||
# pod's well-known pointer, exercising the real discovery chain.
|
||||
#
|
||||
# Prereq: mise k3d:up && mise tilt:up (cluster healthy, context k3d-yucca)
|
||||
set -euo pipefail
|
||||
@@ -48,6 +49,7 @@ mise run yucca-sdk:orchestration-ui:build >/dev/null
|
||||
echo "==> port-forward k3d services to the e2e host ports"
|
||||
kubectl port-forward -n yucca svc/yucca-michael 3010:3010 >/tmp/yucca-e2e-pf.log 2>&1 & PF_PIDS+=($!)
|
||||
kubectl port-forward -n yucca svc/yucca-mock-oidc 8092:8092 >>/tmp/yucca-e2e-pf.log 2>&1 & PF_PIDS+=($!)
|
||||
kubectl port-forward -n yucca svc/yucca-meta 8080:8080 >>/tmp/yucca-e2e-pf.log 2>&1 & PF_PIDS+=($!)
|
||||
# web on :36033 (orchestration-api + the web e2e target) AND :5173 (yucca-api's
|
||||
# OIDC redirect_uri, which the OIDC callback follows).
|
||||
kubectl port-forward -n yucca svc/yucca-web 36033:5173 >>/tmp/yucca-e2e-pf.log 2>&1 & PF_PIDS+=($!)
|
||||
@@ -67,6 +69,8 @@ kubectl port-forward -n yucca svc/yucca-api 3020:3020 >>/tmp/yucca-e2e-pf.log 2>
|
||||
sleep 5
|
||||
|
||||
echo "==> launch orchestration-api as a separate local process (:22676)"
|
||||
# Discovery runs the real chain: meta pod well-known -> /api/meta -> api_root.
|
||||
FUTO_BACKUPS_WELL_KNOWN_URL="http://localhost:8080/.well-known/yucca.json" \
|
||||
NODE_OPTIONS="--require $HERE/hostmap.cjs" \
|
||||
mise run yucca-sdk:orchestration-api:dev >/tmp/yucca-e2e-orch.log 2>&1 & ORCH_PID=$!
|
||||
for _ in $(seq 1 40); do
|
||||
|
||||
@@ -6,7 +6,7 @@ async function bootstrap() {
|
||||
const app = await NestFactory.create(
|
||||
OrchestrationApiModule.forRootAsync({
|
||||
useFactory: () => ({
|
||||
yuccaProductionApi: 'http://localhost:36033',
|
||||
wellKnownUrl: process.env.FUTO_BACKUPS_WELL_KNOWN_URL,
|
||||
developmentMode: true,
|
||||
immichIntegration: {
|
||||
dataFolders: ['upload', 'profile', 'library', 'backups', 'thumbs', 'encoded-video'],
|
||||
|
||||
@@ -7,9 +7,7 @@ import { ORCHESTRATION_PORT, OrchestrationApiModule } from '../src';
|
||||
async function main() {
|
||||
const app = await NestFactory.create<NestApplication>(
|
||||
OrchestrationApiModule.forRootAsync({
|
||||
useFactory: () => ({
|
||||
yuccaProductionApi: 'http://localhost',
|
||||
}),
|
||||
useFactory: () => ({}),
|
||||
}),
|
||||
);
|
||||
|
||||
|
||||
@@ -1540,6 +1540,10 @@
|
||||
"worm": {
|
||||
"type": "boolean"
|
||||
},
|
||||
"site": {
|
||||
"type": "string",
|
||||
"description": "Internal site code from environment metadata"
|
||||
},
|
||||
"paths": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
@@ -1665,6 +1669,14 @@
|
||||
"name": {
|
||||
"type": "string"
|
||||
},
|
||||
"siteCode": {
|
||||
"type": "string",
|
||||
"nullable": true
|
||||
},
|
||||
"storageClusterCode": {
|
||||
"type": "string",
|
||||
"nullable": true
|
||||
},
|
||||
"metrics": {
|
||||
"$ref": "#/components/schemas/RepositoryMetricsDto"
|
||||
},
|
||||
@@ -1682,6 +1694,8 @@
|
||||
"id",
|
||||
"worm",
|
||||
"name",
|
||||
"siteCode",
|
||||
"storageClusterCode",
|
||||
"metrics"
|
||||
]
|
||||
},
|
||||
@@ -1756,6 +1770,12 @@
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
"tags": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
"summary": {
|
||||
"$ref": "#/components/schemas/SnapshotSummaryDto"
|
||||
}
|
||||
@@ -1778,6 +1798,14 @@
|
||||
"name": {
|
||||
"type": "string"
|
||||
},
|
||||
"siteCode": {
|
||||
"type": "string",
|
||||
"nullable": true
|
||||
},
|
||||
"storageClusterCode": {
|
||||
"type": "string",
|
||||
"nullable": true
|
||||
},
|
||||
"metrics": {
|
||||
"$ref": "#/components/schemas/RepositoryMetricsDto"
|
||||
},
|
||||
@@ -1801,6 +1829,8 @@
|
||||
"id",
|
||||
"worm",
|
||||
"name",
|
||||
"siteCode",
|
||||
"storageClusterCode",
|
||||
"metrics",
|
||||
"snapshots"
|
||||
]
|
||||
|
||||
@@ -40,6 +40,7 @@
|
||||
"rxjs": "catalog:",
|
||||
"socket.io": "catalog:",
|
||||
"tail": "catalog:",
|
||||
"zod": "catalog:",
|
||||
"@futo-org/backups-api-client": "workspace:^"
|
||||
},
|
||||
"devDependencies": {
|
||||
|
||||
@@ -42,6 +42,8 @@ export class LocalBackend extends Backend {
|
||||
id,
|
||||
name: dto.name,
|
||||
worm: dto.worm,
|
||||
siteCode: null,
|
||||
storageClusterCode: null,
|
||||
metrics: {
|
||||
sizeBytes: 0,
|
||||
},
|
||||
@@ -55,6 +57,8 @@ export class LocalBackend extends Backend {
|
||||
id,
|
||||
name: dto.name ?? 'new name',
|
||||
worm: false,
|
||||
siteCode: null,
|
||||
storageClusterCode: null,
|
||||
metrics: {
|
||||
sizeBytes: 0,
|
||||
},
|
||||
@@ -83,6 +87,8 @@ export class LocalBackend extends Backend {
|
||||
id: files[index],
|
||||
name: 'Unknown',
|
||||
worm: false,
|
||||
siteCode: null,
|
||||
storageClusterCode: null,
|
||||
metrics: {
|
||||
sizeBytes: 0, // in local cache
|
||||
},
|
||||
|
||||
@@ -13,41 +13,12 @@ import {
|
||||
submitStructuredLog,
|
||||
updateRepository,
|
||||
} from '@futo-org/backups-api-client';
|
||||
import { YUCCA_WELL_KNOWN } from '../const';
|
||||
import { BackendType, CookieName } from '../enum';
|
||||
import { LoggingRepository } from '../repositories/logging.repository';
|
||||
import { BackendConfiguration } from '../schema/tables/backend.table';
|
||||
import { yuccaWellKnown } from '../wellKnown';
|
||||
import { Backend } from './backend';
|
||||
|
||||
type WellKnown = {
|
||||
backends: Record<string, { displayName: string; api: string }>;
|
||||
defaultBackend: string;
|
||||
};
|
||||
|
||||
class YuccaWellKnown {
|
||||
private data?: WellKnown;
|
||||
|
||||
async get() {
|
||||
this.data ??= await fetch(YUCCA_WELL_KNOWN, { signal: AbortSignal.timeout(8000) }).then((response) =>
|
||||
response.json(),
|
||||
);
|
||||
|
||||
return this.data! as WellKnown;
|
||||
}
|
||||
|
||||
async getBaseUrlById(id: string) {
|
||||
const { backends } = await this.get();
|
||||
return backends[id].api;
|
||||
}
|
||||
|
||||
async getBaseUrl() {
|
||||
const { backends, defaultBackend } = await this.get();
|
||||
return backends[defaultBackend].api;
|
||||
}
|
||||
}
|
||||
|
||||
export const yuccaWellKnown = new YuccaWellKnown();
|
||||
|
||||
export class YuccaBackend extends Backend {
|
||||
private readonly logger = LoggingRepository.create(YuccaBackend.name);
|
||||
|
||||
@@ -57,11 +28,7 @@ export class YuccaBackend extends Backend {
|
||||
|
||||
private async getRequestOptions() {
|
||||
return {
|
||||
baseUrl:
|
||||
this.configuration.url ??
|
||||
(this.configuration.uuid
|
||||
? await yuccaWellKnown.getBaseUrlById(this.configuration.uuid)
|
||||
: await yuccaWellKnown.getBaseUrl()),
|
||||
baseUrl: this.configuration.url ?? (await yuccaWellKnown.getBaseUrl()),
|
||||
headers: {
|
||||
cookie: `${CookieName.YuccaAccessToken}=${this.configuration.accessToken}`,
|
||||
},
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
export const ORCHESTRATION_PORT = 22_676;
|
||||
export const REPOSITORY_DEFAULT_CLOUD_UUID = 'd0368cdd-39ae-40e1-91d6-d81815c65c7e';
|
||||
export const YUCCA_WELL_KNOWN = 'https://futo.cloud/.well-known/yucca.json';
|
||||
export const YUCCA_WELL_KNOWN = 'https://meta.futo.cloud/.well-known/yucca.json';
|
||||
|
||||
@@ -35,6 +35,12 @@ export class RepositoryDto {
|
||||
|
||||
@ApiProperty({ type: String })
|
||||
name!: string;
|
||||
|
||||
@ApiProperty({ type: String, nullable: true })
|
||||
siteCode!: string | null;
|
||||
|
||||
@ApiProperty({ type: String, nullable: true })
|
||||
storageClusterCode!: string | null;
|
||||
}
|
||||
|
||||
export class RepositoryMetricsDto {
|
||||
@@ -117,6 +123,11 @@ export class RepositoryCreateRequestDto {
|
||||
@IsBoolean()
|
||||
worm!: boolean;
|
||||
|
||||
@ApiProperty({ type: String, required: false, description: 'Internal site code from environment metadata' })
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
site?: string;
|
||||
|
||||
@ApiProperty({ type: [String], required: false })
|
||||
@IsOptional()
|
||||
@IsArray()
|
||||
|
||||
@@ -4,3 +4,4 @@ export * from './orchestrationApi.module';
|
||||
export { LoggingRepository } from './repositories/logging.repository';
|
||||
export { ModuleConfigRepository } from './repositories/moduleConfig.repository';
|
||||
export { YuccaService } from './services/yucca.service';
|
||||
export { yuccaWellKnown, type Meta, type MetaConfig, type MetaSite } from './wellKnown';
|
||||
|
||||
@@ -27,7 +27,7 @@ export type ImmichIntegration = {
|
||||
|
||||
export type ModuleConfig = {
|
||||
statePath: string;
|
||||
yuccaProductionApi?: string;
|
||||
wellKnownUrl?: string;
|
||||
externalBaseUrl?: string;
|
||||
requireWsAuth?: boolean;
|
||||
requireLock?: boolean;
|
||||
|
||||
@@ -49,6 +49,7 @@ import { RunningTasksService } from './services/runningTasks.service';
|
||||
import { ScheduleService } from './services/schedule.service';
|
||||
import { TelemetryService } from './services/telemetry.service';
|
||||
import { YuccaService } from './services/yucca.service';
|
||||
import { yuccaWellKnown } from './wellKnown';
|
||||
|
||||
export const controllers = [
|
||||
AuthController,
|
||||
@@ -115,6 +116,7 @@ class OrchestrationConfigModule {
|
||||
useFactory: async (...args: any[]): Promise<ModuleConfig> => {
|
||||
const config = await options.useFactory(...args);
|
||||
config.statePath ??= resolve(homedir(), '.yucca');
|
||||
yuccaWellKnown.configure(config.wellKnownUrl);
|
||||
return config as ModuleConfig;
|
||||
},
|
||||
},
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
import { availableParallelism } from 'node:os';
|
||||
import { yuccaWellKnown } from '../wellKnown';
|
||||
import { ConfigRepository } from './config.repository';
|
||||
|
||||
describe(ConfigRepository.name, () => {
|
||||
const db = {
|
||||
selectFrom: jest.fn(() => ({
|
||||
where: jest.fn(() => ({
|
||||
select: jest.fn(() => ({
|
||||
executeTakeFirst: jest.fn().mockResolvedValue(),
|
||||
})),
|
||||
})),
|
||||
})),
|
||||
};
|
||||
|
||||
afterEach(() => {
|
||||
jest.restoreAllMocks();
|
||||
});
|
||||
|
||||
it('resolves options for a placed repository', async () => {
|
||||
const connections = jest.spyOn(yuccaWellKnown, 'getConnections').mockResolvedValue(7);
|
||||
const packSize = jest.spyOn(yuccaWellKnown, 'getPackSizeMib').mockResolvedValue(64);
|
||||
const repository = new ConfigRepository(db as never);
|
||||
|
||||
await expect(
|
||||
repository.getResticOptions({ siteCode: 'father', storageClusterCode: 'father-spice' }),
|
||||
).resolves.toEqual({ connections: 7, packSizeMib: 64 });
|
||||
|
||||
expect(connections).toHaveBeenCalledWith(availableParallelism(), 'father', 'father-spice');
|
||||
expect(packSize).toHaveBeenCalledWith('father', 'father-spice');
|
||||
});
|
||||
|
||||
it('falls back to global config and core count without placement', async () => {
|
||||
const connections = jest.spyOn(yuccaWellKnown, 'getConnections').mockResolvedValue();
|
||||
const packSize = jest.spyOn(yuccaWellKnown, 'getPackSizeMib').mockResolvedValue();
|
||||
const repository = new ConfigRepository(db as never);
|
||||
|
||||
await expect(repository.getResticOptions({ siteCode: null, storageClusterCode: null })).resolves.toEqual({
|
||||
connections: availableParallelism(),
|
||||
packSizeMib: undefined,
|
||||
});
|
||||
|
||||
expect(connections).toHaveBeenCalledWith(availableParallelism(), undefined, undefined);
|
||||
expect(packSize).toHaveBeenCalledWith(undefined, undefined);
|
||||
});
|
||||
});
|
||||
@@ -5,6 +5,9 @@ import { randomBytes } from 'node:crypto';
|
||||
import { availableParallelism } from 'node:os';
|
||||
import { ConfigurationKey } from '../enum';
|
||||
import { DB } from '../schema';
|
||||
import { yuccaWellKnown } from '../wellKnown';
|
||||
|
||||
export type ResticPlacement = { siteCode: string | null; storageClusterCode: string | null };
|
||||
|
||||
@Injectable()
|
||||
export class ConfigRepository {
|
||||
@@ -111,8 +114,18 @@ export class ConfigRepository {
|
||||
return this.set(ConfigurationKey.SkippedOnboardingExtraConfig, '1');
|
||||
}
|
||||
|
||||
async getResticOptionRestConnections() {
|
||||
const concurrency = await this.getOptional(ConfigurationKey.ResticOptionRestConnections);
|
||||
return concurrency ? Number.parseInt(concurrency) : availableParallelism();
|
||||
async getResticOptions(
|
||||
placement: ResticPlacement,
|
||||
): Promise<{ connections: number; packSizeMib: number | undefined }> {
|
||||
const siteCode = placement.siteCode ?? undefined;
|
||||
const clusterCode = placement.storageClusterCode ?? undefined;
|
||||
const cores = availableParallelism();
|
||||
|
||||
const override = await this.getOptional(ConfigurationKey.ResticOptionRestConnections);
|
||||
const connections = override
|
||||
? Number.parseInt(override)
|
||||
: ((await yuccaWellKnown.getConnections(cores, siteCode, clusterCode)) ?? cores);
|
||||
|
||||
return { connections, packSizeMib: await yuccaWellKnown.getPackSizeMib(siteCode, clusterCode) };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,8 @@ type RepositoryRow = {
|
||||
remoteId: string;
|
||||
backendId: string;
|
||||
retentionPolicy: RetentionPolicy | null;
|
||||
siteCode: string | null;
|
||||
storageClusterCode: string | null;
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
@@ -24,6 +26,8 @@ export class RepositoryRepository {
|
||||
remoteId: repository.remoteId,
|
||||
backendId: repository.backendId,
|
||||
retentionPolicy: repository.retentionPolicy === null ? null : JSON.stringify(repository.retentionPolicy),
|
||||
siteCode: repository.siteCode,
|
||||
storageClusterCode: repository.storageClusterCode,
|
||||
})
|
||||
.returningAll()
|
||||
.executeTakeFirstOrThrow();
|
||||
@@ -54,6 +58,8 @@ export class RepositoryRepository {
|
||||
remoteId: row.remoteId,
|
||||
backendId: row.backendId,
|
||||
retentionPolicy: row.retentionPolicy === null ? null : (JSON.parse(row.retentionPolicy) as RetentionPolicy),
|
||||
siteCode: row.siteCode,
|
||||
storageClusterCode: row.storageClusterCode,
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -65,6 +71,8 @@ export class RepositoryRepository {
|
||||
remoteId: row.remoteId,
|
||||
backendId: row.backendId,
|
||||
retentionPolicy: row.retentionPolicy === null ? null : (JSON.parse(row.retentionPolicy) as RetentionPolicy),
|
||||
siteCode: row.siteCode,
|
||||
storageClusterCode: row.storageClusterCode,
|
||||
}));
|
||||
}
|
||||
|
||||
|
||||
@@ -3,15 +3,16 @@ import { Injectable } from '@nestjs/common';
|
||||
import { Writable } from 'node:stream';
|
||||
import { RepositorySnapshotRestoreRequestDto } from '../dto/repository.dto';
|
||||
import { createSampledLogWriter, RetentionPolicy } from '../utils/restic';
|
||||
import { ConfigRepository } from './config.repository';
|
||||
import { ConfigRepository, ResticPlacement } from './config.repository';
|
||||
|
||||
@Injectable()
|
||||
export class ResticRepository {
|
||||
constructor(private readonly config: ConfigRepository) {}
|
||||
|
||||
async init(repository: string, key: Uint8Array) {
|
||||
async init(repository: string, key: Uint8Array, placement: ResticPlacement) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
await init()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.run();
|
||||
@@ -20,15 +21,18 @@ export class ResticRepository {
|
||||
async backup(
|
||||
repository: string,
|
||||
key: Uint8Array,
|
||||
placement: ResticPlacement,
|
||||
paths: string[],
|
||||
logStream?: Writable,
|
||||
signal?: AbortSignal,
|
||||
tags: string[] = [],
|
||||
) {
|
||||
const write = createSampledLogWriter(logStream);
|
||||
const { connections, packSizeMib } = await this.config.getResticOptions(placement);
|
||||
|
||||
return await backup()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.packSize(packSizeMib)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.tag(...tags)
|
||||
@@ -41,15 +45,17 @@ export class ResticRepository {
|
||||
async restore(
|
||||
repository: string,
|
||||
key: Uint8Array,
|
||||
placement: ResticPlacement,
|
||||
snapshotId: string,
|
||||
{ include, target }: RepositorySnapshotRestoreRequestDto,
|
||||
logStream?: Writable,
|
||||
signal?: AbortSignal,
|
||||
) {
|
||||
const write = createSampledLogWriter(logStream);
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
|
||||
let command = restore()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.snapshot(snapshotId)
|
||||
@@ -64,9 +70,10 @@ export class ResticRepository {
|
||||
return await command.run();
|
||||
}
|
||||
|
||||
async ls(repository: string, key: Uint8Array, snapshotId: string, path: string) {
|
||||
async ls(repository: string, key: Uint8Array, placement: ResticPlacement, snapshotId: string, path: string) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await ls()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.snapshot(snapshotId)
|
||||
@@ -74,35 +81,46 @@ export class ResticRepository {
|
||||
.run();
|
||||
}
|
||||
|
||||
async stats(repository: string, key: Uint8Array) {
|
||||
async stats(repository: string, key: Uint8Array, placement: ResticPlacement) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await stats()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.modeRawData()
|
||||
.run();
|
||||
}
|
||||
|
||||
async snapshots(repository: string, key: Uint8Array) {
|
||||
async snapshots(repository: string, key: Uint8Array, placement: ResticPlacement) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await snapshots()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.run();
|
||||
}
|
||||
|
||||
async snapshot(repository: string, key: Uint8Array, snapshotId: string) {
|
||||
async snapshot(repository: string, key: Uint8Array, placement: ResticPlacement, snapshotId: string) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await snapshots()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.snapshot(snapshotId)
|
||||
.run();
|
||||
}
|
||||
|
||||
async forget(repository: string, key: Uint8Array, snapshotId: string, prune = true, signal?: AbortSignal) {
|
||||
async forget(
|
||||
repository: string,
|
||||
key: Uint8Array,
|
||||
placement: ResticPlacement,
|
||||
snapshotId: string,
|
||||
prune = true,
|
||||
signal?: AbortSignal,
|
||||
) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await forget()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.snapshot(snapshotId)
|
||||
@@ -111,9 +129,16 @@ export class ResticRepository {
|
||||
.run();
|
||||
}
|
||||
|
||||
async forgetByPolicy(repository: string, key: Uint8Array, policy: RetentionPolicy, signal?: AbortSignal) {
|
||||
async forgetByPolicy(
|
||||
repository: string,
|
||||
key: Uint8Array,
|
||||
placement: ResticPlacement,
|
||||
policy: RetentionPolicy,
|
||||
signal?: AbortSignal,
|
||||
) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await forget()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.signal(signal)
|
||||
@@ -127,26 +152,30 @@ export class ResticRepository {
|
||||
.run();
|
||||
}
|
||||
|
||||
async prune(repository: string, key: Uint8Array, signal?: AbortSignal) {
|
||||
async prune(repository: string, key: Uint8Array, placement: ResticPlacement, signal?: AbortSignal) {
|
||||
const { connections, packSizeMib } = await this.config.getResticOptions(placement);
|
||||
return await prune()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.packSize(packSizeMib)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.signal(signal)
|
||||
.run();
|
||||
}
|
||||
|
||||
async keyList(repository: string, key: Uint8Array) {
|
||||
async keyList(repository: string, key: Uint8Array, placement: ResticPlacement) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await keyList()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
.run();
|
||||
}
|
||||
|
||||
async unlockAll(repository: string, key: Uint8Array) {
|
||||
async unlockAll(repository: string, key: Uint8Array, placement: ResticPlacement) {
|
||||
const { connections } = await this.config.getResticOptions(placement);
|
||||
return await unlock()
|
||||
.option(`rest.connections=${await this.config.getResticOptionRestConnections()}`)
|
||||
.option(`rest.connections=${connections}`)
|
||||
.removeAll()
|
||||
.repository(repository)
|
||||
.password(Buffer.from(key).toString('hex'))
|
||||
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
import { Kysely } from 'kysely';
|
||||
|
||||
export async function up(db: Kysely<any>): Promise<void> {
|
||||
await db.schema.alterTable('repositories').addColumn('siteCode', 'text').execute();
|
||||
await db.schema.alterTable('repositories').addColumn('storageClusterCode', 'text').execute();
|
||||
}
|
||||
|
||||
export async function down(db: Kysely<any>): Promise<void> {
|
||||
await db.schema.alterTable('repositories').dropColumn('storageClusterCode').execute();
|
||||
await db.schema.alterTable('repositories').dropColumn('siteCode').execute();
|
||||
}
|
||||
@@ -11,7 +11,6 @@ export type BackendConfiguration =
|
||||
*/
|
||||
type: BackendType.Yucca;
|
||||
url?: string;
|
||||
uuid?: string;
|
||||
accessToken: string;
|
||||
}
|
||||
| {
|
||||
|
||||
@@ -6,4 +6,8 @@ export class RepositoryTable {
|
||||
backendId!: string;
|
||||
|
||||
retentionPolicy!: string | null;
|
||||
|
||||
siteCode!: string | null;
|
||||
|
||||
storageClusterCode!: string | null;
|
||||
}
|
||||
|
||||
@@ -1,12 +1,11 @@
|
||||
import { Injectable, InternalServerErrorException } from '@nestjs/common';
|
||||
import { createEventSource, EventSourceClient } from 'eventsource-client';
|
||||
import { yuccaWellKnown } from '../backends/yucca.backend';
|
||||
import { REPOSITORY_DEFAULT_CLOUD_UUID } from '../const';
|
||||
import { BackendType } from '../enum';
|
||||
import { EventsGateway } from '../events/events.gateway';
|
||||
import { BackendRepository } from '../repositories/backend.repository';
|
||||
import { ConfigRepository } from '../repositories/config.repository';
|
||||
import { ModuleConfigRepository } from '../repositories/moduleConfig.repository';
|
||||
import { yuccaWellKnown } from '../wellKnown';
|
||||
import { TelemetryService } from './telemetry.service';
|
||||
|
||||
@Injectable()
|
||||
@@ -14,12 +13,11 @@ export class AuthService {
|
||||
constructor(
|
||||
readonly config: ConfigRepository,
|
||||
readonly backend: BackendRepository,
|
||||
readonly moduleConfig: ModuleConfigRepository,
|
||||
readonly events: EventsGateway,
|
||||
readonly telemetry: TelemetryService,
|
||||
) {}
|
||||
|
||||
private async waitForDeviceFlow(events: EventSourceClient, url?: string) {
|
||||
private async waitForDeviceFlow(events: EventSourceClient) {
|
||||
for await (const { data } of events) {
|
||||
const { type, accessToken } = JSON.parse(data);
|
||||
|
||||
@@ -28,7 +26,6 @@ export class AuthService {
|
||||
await this.backend.updateBackend(REPOSITORY_DEFAULT_CLOUD_UUID, {
|
||||
type: BackendType.Yucca,
|
||||
accessToken,
|
||||
url,
|
||||
});
|
||||
|
||||
this.telemetry.submitStructuredLog('Connected FUTO Backups backend', {
|
||||
@@ -63,8 +60,7 @@ export class AuthService {
|
||||
}
|
||||
|
||||
async oidcDeviceFlow(): Promise<{ userCode: string; verificationUri: string }> {
|
||||
const overrideEndpoint = this.moduleConfig.get().yuccaProductionApi;
|
||||
const endpoint = overrideEndpoint ?? (await yuccaWellKnown.getBaseUrl());
|
||||
const endpoint = await yuccaWellKnown.getBaseUrl();
|
||||
|
||||
const events: EventSourceClient = createEventSource({
|
||||
url: new URL('/api/auth/oidc/device', endpoint),
|
||||
@@ -77,7 +73,7 @@ export class AuthService {
|
||||
clearTimeout(connectTimeout);
|
||||
const { userCode, verificationUri } = JSON.parse(data);
|
||||
|
||||
void this.waitForDeviceFlow(events, overrideEndpoint).catch((error) => {
|
||||
void this.waitForDeviceFlow(events).catch((error) => {
|
||||
this.telemetry.submitStructuredLog('Device flow authentication errored', { error });
|
||||
this.events.publish({ type: 'DeviceFlowFailure' });
|
||||
});
|
||||
|
||||
@@ -24,10 +24,10 @@ import {
|
||||
RepositoryWithMetricsDto,
|
||||
RunHistoryResponseDto,
|
||||
} from '../dto/repository.dto';
|
||||
import { ResticTagPrefix, TaskType } from '../enum';
|
||||
import { BackendType, ResticTagPrefix, TaskType } from '../enum';
|
||||
import { EventsGateway } from '../events/events.gateway';
|
||||
import { BackendRepository } from '../repositories/backend.repository';
|
||||
import { ConfigRepository } from '../repositories/config.repository';
|
||||
import { ConfigRepository, ResticPlacement } from '../repositories/config.repository';
|
||||
import { DatabaseRepository } from '../repositories/database.repository';
|
||||
import { LoggingRepository } from '../repositories/logging.repository';
|
||||
import { ModuleConfigRepository } from '../repositories/moduleConfig.repository';
|
||||
@@ -110,14 +110,17 @@ export class RepositoryService {
|
||||
const id = randomUUID();
|
||||
|
||||
const endpoint = await backend.getResticEndpoint(remote.id);
|
||||
const placement = { siteCode: remote.siteCode, storageClusterCode: remote.storageClusterCode };
|
||||
const key = await this.config.deriveEncryptionKey(`repository-${remote.id}`);
|
||||
await this.restic.init(endpoint, key);
|
||||
await this.restic.init(endpoint, key, placement);
|
||||
|
||||
await this.repository.create({
|
||||
id,
|
||||
remoteId: remote.id,
|
||||
backendId,
|
||||
retentionPolicy: DEFAULT_RETENTION_POLICY,
|
||||
siteCode: remote.siteCode,
|
||||
storageClusterCode: remote.storageClusterCode,
|
||||
});
|
||||
|
||||
const paths = dto.paths ?? [];
|
||||
@@ -179,7 +182,7 @@ export class RepositoryService {
|
||||
const localPaths = await this.repositoryPath.getAll();
|
||||
const localMetrics = await this.repositoryLocalMetrics.getAll();
|
||||
|
||||
for (const { id, remoteId, backendId, retentionPolicy } of localRepositories) {
|
||||
for (const { id, remoteId, backendId, retentionPolicy, siteCode, storageClusterCode } of localRepositories) {
|
||||
const remoteRepository = remoteRepositories[backendId][remoteId];
|
||||
|
||||
const configuration: RepositoryConfigurationDto = {
|
||||
@@ -190,6 +193,12 @@ export class RepositoryService {
|
||||
const metrics = localMetrics.find((entry) => entry.id === id);
|
||||
|
||||
if (remoteRepository) {
|
||||
if (remoteRepository.siteCode !== siteCode || remoteRepository.storageClusterCode !== storageClusterCode) {
|
||||
await this.repository.update(id, {
|
||||
siteCode: remoteRepository.siteCode,
|
||||
storageClusterCode: remoteRepository.storageClusterCode,
|
||||
});
|
||||
}
|
||||
repositories.push({
|
||||
...remoteRepository,
|
||||
id,
|
||||
@@ -210,6 +219,8 @@ export class RepositoryService {
|
||||
id,
|
||||
name: 'Unknown',
|
||||
worm: false,
|
||||
siteCode,
|
||||
storageClusterCode,
|
||||
...(await this.getLocalRepository(id, configuration, metrics)),
|
||||
backends: {
|
||||
primary: {
|
||||
@@ -255,11 +266,11 @@ export class RepositoryService {
|
||||
|
||||
const snapshots = await Promise.allSettled(
|
||||
list.map(async (repository) => {
|
||||
const { endpoint, key } = await this.getResticParameters({
|
||||
const { endpoint, key, placement } = await this.getResticParameters({
|
||||
backendId: repository.backends!.primary.id,
|
||||
remoteId: repository.id,
|
||||
});
|
||||
return this.restic.snapshots(endpoint, key);
|
||||
return this.restic.snapshots(endpoint, key, placement);
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -385,27 +396,40 @@ export class RepositoryService {
|
||||
|
||||
private async getResticParameters(
|
||||
repository: string | { backendId: string; remoteId: string },
|
||||
): Promise<{ endpoint: string; key: Uint8Array }> {
|
||||
): Promise<{ endpoint: string; key: Uint8Array; placement: ResticPlacement }> {
|
||||
let backendId: string;
|
||||
let remoteId: string;
|
||||
let localId: string | undefined;
|
||||
let siteCode: string | null = null;
|
||||
let storageClusterCode: string | null = null;
|
||||
if (typeof repository === 'string') {
|
||||
const localRepository = await this.repository.get(repository);
|
||||
if (!localRepository) {
|
||||
throw new NotFoundException('Repository not found locally');
|
||||
}
|
||||
|
||||
localId = localRepository.id;
|
||||
backendId = localRepository.backendId;
|
||||
remoteId = localRepository.remoteId;
|
||||
siteCode = localRepository.siteCode;
|
||||
storageClusterCode = localRepository.storageClusterCode;
|
||||
} else {
|
||||
({ backendId, remoteId } = repository);
|
||||
}
|
||||
|
||||
const { backend } = await this.getBackendOrThrow(backendId);
|
||||
const { backend, configuration } = await this.getBackendOrThrow(backendId);
|
||||
if ((!siteCode || !storageClusterCode) && configuration.type === BackendType.Yucca) {
|
||||
const { repository: remote } = await backend.getRepository(remoteId);
|
||||
siteCode = remote.siteCode;
|
||||
storageClusterCode = remote.storageClusterCode;
|
||||
if (localId) {
|
||||
await this.repository.update(localId, { siteCode, storageClusterCode });
|
||||
}
|
||||
}
|
||||
const endpoint = await backend.getResticEndpoint(remoteId);
|
||||
|
||||
const key = await this.config.deriveEncryptionKey(`repository-${remoteId}`);
|
||||
|
||||
return { endpoint, key };
|
||||
return { endpoint, key, placement: { siteCode, storageClusterCode } };
|
||||
}
|
||||
|
||||
private async updateLocalMetrics(
|
||||
@@ -414,6 +438,7 @@ export class RepositoryService {
|
||||
resticParameters?: {
|
||||
endpoint: string;
|
||||
key: Uint8Array;
|
||||
placement: ResticPlacement;
|
||||
};
|
||||
additionalMetrics?: Updateable<RepositoryLocalMetricsTable>;
|
||||
},
|
||||
@@ -424,8 +449,8 @@ export class RepositoryService {
|
||||
|
||||
try {
|
||||
if (options.resticParameters) {
|
||||
const { endpoint, key } = options.resticParameters;
|
||||
const { total_size } = await this.restic.stats(endpoint, key);
|
||||
const { endpoint, key, placement } = options.resticParameters;
|
||||
const { total_size } = await this.restic.stats(endpoint, key, placement);
|
||||
metrics.sizeBytes = total_size;
|
||||
}
|
||||
|
||||
@@ -488,7 +513,7 @@ export class RepositoryService {
|
||||
|
||||
const { backendId, remoteId, retentionPolicy } = repository;
|
||||
const { backend } = await this.getBackendOrThrow(backendId);
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
|
||||
const paths = await this.repositoryPath.get(id);
|
||||
if (paths.length === 0) {
|
||||
@@ -541,8 +566,8 @@ export class RepositoryService {
|
||||
}
|
||||
}
|
||||
|
||||
await this.restic.unlockAll(endpoint, key);
|
||||
const summary = await this.restic.backup(endpoint, key, paths, log, taskSignal, tags);
|
||||
await this.restic.unlockAll(endpoint, key, placement);
|
||||
const summary = await this.restic.backup(endpoint, key, placement, paths, log, taskSignal, tags);
|
||||
|
||||
this.telemetry.submitStructuredLog('Finished backup to primary backend', {
|
||||
repositoryId: id,
|
||||
@@ -550,7 +575,7 @@ export class RepositoryService {
|
||||
});
|
||||
|
||||
if (retentionPolicy) {
|
||||
await this.runForgetAndPrune(endpoint, key, retentionPolicy, log, taskSignal);
|
||||
await this.runForgetAndPrune(endpoint, key, placement, retentionPolicy, log, taskSignal);
|
||||
|
||||
this.telemetry.submitStructuredLog('Finished prune on primary backend', {
|
||||
repositoryId: id,
|
||||
@@ -566,7 +591,7 @@ export class RepositoryService {
|
||||
const lastBackupDuration = Date.now() - startTime;
|
||||
|
||||
void this.updateLocalMetrics(id, {
|
||||
resticParameters: { endpoint, key },
|
||||
resticParameters: { endpoint, key, placement },
|
||||
additionalMetrics: {
|
||||
lastBackup,
|
||||
lastSuccessfulBackup,
|
||||
@@ -599,11 +624,12 @@ export class RepositoryService {
|
||||
private async runForgetAndPrune(
|
||||
endpoint: string,
|
||||
key: Uint8Array,
|
||||
placement: ResticPlacement,
|
||||
policy: RetentionPolicy,
|
||||
log: WriteStream,
|
||||
signal?: AbortSignal,
|
||||
): Promise<void> {
|
||||
const events = await this.restic.forgetByPolicy(endpoint, key, policy, signal);
|
||||
const events = await this.restic.forgetByPolicy(endpoint, key, placement, policy, signal);
|
||||
|
||||
for (const { keep, remove, reasons } of events) {
|
||||
if (keep) {
|
||||
@@ -626,7 +652,7 @@ export class RepositoryService {
|
||||
}
|
||||
}
|
||||
|
||||
await this.restic.prune(endpoint, key, signal);
|
||||
await this.restic.prune(endpoint, key, placement, signal);
|
||||
}
|
||||
|
||||
async pruneRepository(
|
||||
@@ -659,7 +685,7 @@ export class RepositoryService {
|
||||
retentionPolicy,
|
||||
});
|
||||
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
|
||||
return new Promise((resolve) => {
|
||||
const task = new Promise<void>(
|
||||
@@ -672,15 +698,15 @@ export class RepositoryService {
|
||||
|
||||
try {
|
||||
const taskSignal = this.tasks.startTask(id, TaskType.Forget, logId, signal);
|
||||
await this.restic.unlockAll(endpoint, key);
|
||||
await this.runForgetAndPrune(endpoint, key, retentionPolicy, log, taskSignal);
|
||||
await this.restic.unlockAll(endpoint, key, placement);
|
||||
await this.runForgetAndPrune(endpoint, key, placement, retentionPolicy, log, taskSignal);
|
||||
} finally {
|
||||
this.tasks.endTask(id);
|
||||
}
|
||||
},
|
||||
(error) => {
|
||||
void this.updateLocalMetrics(id, {
|
||||
resticParameters: { endpoint, key },
|
||||
resticParameters: { endpoint, key, placement },
|
||||
});
|
||||
|
||||
this.telemetry.submitStructuredLog('Finished repository prune', {
|
||||
@@ -702,10 +728,10 @@ export class RepositoryService {
|
||||
}
|
||||
|
||||
async checkImportRepository(id: string, backendId: string): Promise<RepositoryCheckImportResponseDto> {
|
||||
const { endpoint, key } = await this.getResticParameters({ backendId, remoteId: id });
|
||||
const { endpoint, key, placement } = await this.getResticParameters({ backendId, remoteId: id });
|
||||
|
||||
try {
|
||||
await this.restic.snapshots(endpoint, key);
|
||||
await this.restic.snapshots(endpoint, key, placement);
|
||||
|
||||
return {
|
||||
readable: true,
|
||||
@@ -729,12 +755,13 @@ export class RepositoryService {
|
||||
const localId = randomUUID();
|
||||
|
||||
const endpoint = await backend.getResticEndpoint(remote.id);
|
||||
const placement = { siteCode: remote.siteCode, storageClusterCode: remote.storageClusterCode };
|
||||
const key = await this.config.deriveEncryptionKey(`repository-${remote.id}`);
|
||||
await this.restic.keyList(endpoint, key);
|
||||
await this.restic.keyList(endpoint, key, placement);
|
||||
|
||||
let paths: string[] = [];
|
||||
try {
|
||||
const snapshots = await this.restic.snapshots(endpoint, key);
|
||||
const snapshots = await this.restic.snapshots(endpoint, key, placement);
|
||||
snapshots.sort((a, b) => +b.time - +a.time);
|
||||
paths = snapshots[0].paths;
|
||||
} catch (error) {
|
||||
@@ -750,6 +777,8 @@ export class RepositoryService {
|
||||
remoteId: remote.id,
|
||||
backendId,
|
||||
retentionPolicy: DEFAULT_RETENTION_POLICY,
|
||||
siteCode: remote.siteCode,
|
||||
storageClusterCode: remote.storageClusterCode,
|
||||
});
|
||||
|
||||
const repository: LocalRepositoryDto = {
|
||||
@@ -802,12 +831,15 @@ export class RepositoryService {
|
||||
});
|
||||
|
||||
const endpoint = await backend.getResticEndpoint(remote.id);
|
||||
const placement = { siteCode: remote.siteCode, storageClusterCode: remote.storageClusterCode };
|
||||
const key = await this.config.deriveEncryptionKey(`repository-${remote.id}`);
|
||||
await this.restic.init(endpoint, key);
|
||||
await this.restic.init(endpoint, key, placement);
|
||||
|
||||
await this.repository.update(id, {
|
||||
remoteId: remote.id,
|
||||
backendId: dto.backendId,
|
||||
siteCode: remote.siteCode,
|
||||
storageClusterCode: remote.storageClusterCode,
|
||||
});
|
||||
|
||||
const { id: _, ...repository }: LocalRepositoryDto = {
|
||||
@@ -838,8 +870,8 @@ export class RepositoryService {
|
||||
}
|
||||
|
||||
async getSnapshots(id: string): Promise<ListSnapshotsResponseDto> {
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const snapshots = await this.restic.snapshots(endpoint, key);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
const snapshots = await this.restic.snapshots(endpoint, key, placement);
|
||||
|
||||
return {
|
||||
snapshots: snapshots.map((snapshot) => this.mapSnapshot(snapshot)),
|
||||
@@ -847,8 +879,8 @@ export class RepositoryService {
|
||||
}
|
||||
|
||||
async getSnapshot(repositoryId: string, snapshotId: string): Promise<GetSnapshotResponseDto> {
|
||||
const { endpoint, key } = await this.getResticParameters(repositoryId);
|
||||
const snapshots = await this.restic.snapshot(endpoint, key, snapshotId);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(repositoryId);
|
||||
const snapshots = await this.restic.snapshot(endpoint, key, placement, snapshotId);
|
||||
|
||||
return { snapshot: this.mapSnapshot(snapshots[0]) };
|
||||
}
|
||||
@@ -880,11 +912,11 @@ export class RepositoryService {
|
||||
logId,
|
||||
});
|
||||
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
|
||||
try {
|
||||
const signal = this.tasks.startTask(id, TaskType.Restore, logId);
|
||||
summary = await this.restic.restore(endpoint, key, snapshotId, dto, log, signal);
|
||||
summary = await this.restic.restore(endpoint, key, placement, snapshotId, dto, log, signal);
|
||||
} finally {
|
||||
this.tasks.endTask(id);
|
||||
}
|
||||
@@ -936,11 +968,19 @@ export class RepositoryService {
|
||||
logId,
|
||||
});
|
||||
|
||||
const { endpoint, key } = await this.getResticParameters({ backendId, remoteId: id });
|
||||
const { endpoint, key, placement } = await this.getResticParameters({ backendId, remoteId: id });
|
||||
|
||||
try {
|
||||
const signal = this.tasks.startTask(id, TaskType.Restore, logId);
|
||||
summary = await this.restic.restore(endpoint, key, snapshotId, { include: dto.include }, log, signal);
|
||||
summary = await this.restic.restore(
|
||||
endpoint,
|
||||
key,
|
||||
placement,
|
||||
snapshotId,
|
||||
{ include: dto.include },
|
||||
log,
|
||||
signal,
|
||||
);
|
||||
|
||||
if (dto.yuccaConfig) {
|
||||
const target = await this.storage.tempdir();
|
||||
@@ -948,6 +988,7 @@ export class RepositoryService {
|
||||
await this.restic.restore(
|
||||
endpoint,
|
||||
key,
|
||||
placement,
|
||||
snapshotId,
|
||||
{ include: [dto.yuccaConfig], target },
|
||||
log,
|
||||
@@ -1006,13 +1047,13 @@ export class RepositoryService {
|
||||
snapshotId,
|
||||
});
|
||||
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
|
||||
let error;
|
||||
try {
|
||||
const signal = this.tasks.startTask(id, TaskType.Forget);
|
||||
await this.restic.unlockAll(endpoint, key);
|
||||
await this.restic.forget(endpoint, key, snapshotId, true, signal);
|
||||
await this.restic.unlockAll(endpoint, key, placement);
|
||||
await this.restic.forget(endpoint, key, placement, snapshotId, true, signal);
|
||||
} catch (error_) {
|
||||
error = error_;
|
||||
} finally {
|
||||
@@ -1030,7 +1071,7 @@ export class RepositoryService {
|
||||
}
|
||||
|
||||
await this.updateLocalMetrics(id, {
|
||||
resticParameters: { endpoint, key },
|
||||
resticParameters: { endpoint, key, placement },
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1042,8 +1083,8 @@ export class RepositoryService {
|
||||
const path = dto.path ?? '/';
|
||||
|
||||
try {
|
||||
const { endpoint, key } = await this.getResticParameters(id);
|
||||
const files = await this.restic.ls(endpoint, key, snapshotId, path);
|
||||
const { endpoint, key, placement } = await this.getResticParameters(id);
|
||||
const files = await this.restic.ls(endpoint, key, placement, snapshotId, path);
|
||||
|
||||
return {
|
||||
parent: dirname(path),
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
// Grammar for the client-evaluated `connections_math` expression:
|
||||
// expr := term (('+' | '-') term)*
|
||||
// term := factor (('*' | '/') factor)*
|
||||
// factor := integer | 'cores' | ('min' | 'max') '(' expr ',' expr ')' | '(' expr ')'
|
||||
// Clients substitute their local CPU count for `cores` and fall back to a
|
||||
// default when an expression fails to parse.
|
||||
|
||||
type Token = { kind: 'number'; value: number } | { kind: 'ident'; value: string } | { kind: 'symbol'; value: string };
|
||||
|
||||
const tokenize = (input: string): Token[] => {
|
||||
const tokens: Token[] = [];
|
||||
let i = 0;
|
||||
while (i < input.length) {
|
||||
const char = input[i];
|
||||
if (/\s/.test(char)) {
|
||||
i++;
|
||||
} else if (/\d/.test(char)) {
|
||||
let j = i;
|
||||
while (j < input.length && /\d/.test(input[j])) {
|
||||
j++;
|
||||
}
|
||||
tokens.push({ kind: 'number', value: Number.parseInt(input.slice(i, j), 10) });
|
||||
i = j;
|
||||
} else if (/[a-z]/.test(char)) {
|
||||
let j = i;
|
||||
while (j < input.length && /[a-z]/.test(input[j])) {
|
||||
j++;
|
||||
}
|
||||
tokens.push({ kind: 'ident', value: input.slice(i, j) });
|
||||
i = j;
|
||||
} else if ('+-*/(),'.includes(char)) {
|
||||
tokens.push({ kind: 'symbol', value: char });
|
||||
i++;
|
||||
} else {
|
||||
throw new Error(`Unexpected character '${char}' in connections_math expression`);
|
||||
}
|
||||
}
|
||||
return tokens;
|
||||
};
|
||||
|
||||
class Parser {
|
||||
private position = 0;
|
||||
|
||||
constructor(
|
||||
private readonly tokens: Token[],
|
||||
private readonly cores: number,
|
||||
) {}
|
||||
|
||||
parse(): number {
|
||||
const value = this.expr();
|
||||
if (this.position < this.tokens.length) {
|
||||
throw new Error('Unexpected trailing tokens in connections_math expression');
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
private expr(): number {
|
||||
let value = this.term();
|
||||
while (this.peekSymbol('+') || this.peekSymbol('-')) {
|
||||
const op = this.next();
|
||||
value = op.value === '+' ? value + this.term() : value - this.term();
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
private term(): number {
|
||||
let value = this.factor();
|
||||
while (this.peekSymbol('*') || this.peekSymbol('/')) {
|
||||
const op = this.next();
|
||||
value = op.value === '*' ? value * this.factor() : value / this.factor();
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
private factor(): number {
|
||||
const token = this.next();
|
||||
if (token.kind === 'number') {
|
||||
return token.value;
|
||||
}
|
||||
if (token.kind === 'ident') {
|
||||
if (token.value === 'cores') {
|
||||
return this.cores;
|
||||
}
|
||||
if (token.value === 'min' || token.value === 'max') {
|
||||
this.expectSymbol('(');
|
||||
const left = this.expr();
|
||||
this.expectSymbol(',');
|
||||
const right = this.expr();
|
||||
this.expectSymbol(')');
|
||||
return token.value === 'min' ? Math.min(left, right) : Math.max(left, right);
|
||||
}
|
||||
throw new Error(`Unknown identifier '${token.value}' in connections_math expression`);
|
||||
}
|
||||
if (token.kind === 'symbol' && token.value === '(') {
|
||||
const value = this.expr();
|
||||
this.expectSymbol(')');
|
||||
return value;
|
||||
}
|
||||
throw new Error(`Unexpected token '${token.value}' in connections_math expression`);
|
||||
}
|
||||
|
||||
private next(): Token {
|
||||
const token = this.tokens[this.position++];
|
||||
if (!token) {
|
||||
throw new Error('Unexpected end of connections_math expression');
|
||||
}
|
||||
return token;
|
||||
}
|
||||
|
||||
private peekSymbol(symbol: string): boolean {
|
||||
const token = this.tokens[this.position];
|
||||
return token?.kind === 'symbol' && token.value === symbol;
|
||||
}
|
||||
|
||||
private expectSymbol(symbol: string): void {
|
||||
const token = this.next();
|
||||
if (token.kind !== 'symbol' || token.value !== symbol) {
|
||||
throw new Error(`Expected '${symbol}' in connections_math expression`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export const evaluateConnectionsMath = (expression: string, cores: number): number => {
|
||||
const result = new Parser(tokenize(expression), cores).parse();
|
||||
if (!Number.isFinite(result)) {
|
||||
throw new TypeError('connections_math expression did not evaluate to a finite number');
|
||||
}
|
||||
return result;
|
||||
};
|
||||
|
||||
export const isValidConnectionsMath = (expression: string): boolean => {
|
||||
try {
|
||||
evaluateConnectionsMath(expression, 1);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
};
|
||||
@@ -0,0 +1,181 @@
|
||||
import { z } from 'zod';
|
||||
import { YUCCA_WELL_KNOWN } from './const';
|
||||
import { evaluateConnectionsMath } from './utils/connectionsMath';
|
||||
|
||||
const metaConfigSchema = z
|
||||
.object({
|
||||
restic_pack_size_mib: z.number().int().min(1).max(128).optional(),
|
||||
connections_math: z.string().optional(),
|
||||
})
|
||||
.strict();
|
||||
|
||||
const metaClusterSchema = z
|
||||
.object({
|
||||
code: z.string().min(1),
|
||||
display_name: z.string().min(1),
|
||||
cluster_config: metaConfigSchema,
|
||||
})
|
||||
.strict();
|
||||
|
||||
const metaSiteSchema = z
|
||||
.object({
|
||||
code: z.string().min(1),
|
||||
display_name: z.string().min(1),
|
||||
description: z.string(),
|
||||
rest_url: z.url(),
|
||||
default_cluster: z.string().min(1),
|
||||
site_config: metaConfigSchema,
|
||||
clusters: z.array(metaClusterSchema).min(1),
|
||||
})
|
||||
.strict();
|
||||
|
||||
const metaSchema = z
|
||||
.object({
|
||||
api_root: z.url(),
|
||||
config: metaConfigSchema,
|
||||
default_site: z.string().min(1),
|
||||
sites: z.array(metaSiteSchema).min(1),
|
||||
})
|
||||
.strict()
|
||||
.superRefine((meta, ctx) => {
|
||||
if (!meta.sites.some((site) => site.code === meta.default_site)) {
|
||||
ctx.addIssue({ code: 'custom', message: `default_site '${meta.default_site}' is not declared` });
|
||||
}
|
||||
for (const site of meta.sites) {
|
||||
if (!site.clusters.some((cluster) => cluster.code === site.default_cluster)) {
|
||||
ctx.addIssue({
|
||||
code: 'custom',
|
||||
message: `Site '${site.code}' default_cluster '${site.default_cluster}' is not declared`,
|
||||
});
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
const pointerSchema = z.object({ meta_url: z.url() }).strict();
|
||||
|
||||
export type MetaConfig = z.infer<typeof metaConfigSchema>;
|
||||
export type MetaCluster = z.infer<typeof metaClusterSchema>;
|
||||
export type MetaSite = z.infer<typeof metaSiteSchema>;
|
||||
export type Meta = z.infer<typeof metaSchema>;
|
||||
|
||||
const META_CACHE_TTL = 5 * 60 * 1000;
|
||||
const FAILURE_CACHE_TTL = 60 * 1000;
|
||||
const FETCH_TIMEOUT = 8000;
|
||||
|
||||
class UnknownTopologyCodeError extends Error {}
|
||||
|
||||
const fetchValidated = async <T>(url: string, schema: z.ZodType<T>): Promise<T> => {
|
||||
const response = await fetch(url, { signal: AbortSignal.timeout(FETCH_TIMEOUT) });
|
||||
if (!response.ok) {
|
||||
throw new Error(`GET ${url} failed with HTTP ${response.status}`);
|
||||
}
|
||||
return schema.parse(await response.json());
|
||||
};
|
||||
|
||||
export class YuccaWellKnown {
|
||||
private url = YUCCA_WELL_KNOWN;
|
||||
private data?: Meta;
|
||||
private fetchedAt = 0;
|
||||
private lastError?: unknown;
|
||||
private failedAt = 0;
|
||||
|
||||
configure(url?: string) {
|
||||
this.url = url ?? YUCCA_WELL_KNOWN;
|
||||
this.data = undefined;
|
||||
this.fetchedAt = 0;
|
||||
this.lastError = undefined;
|
||||
this.failedAt = 0;
|
||||
}
|
||||
|
||||
async get(): Promise<Meta> {
|
||||
const fresh = Date.now() - this.fetchedAt < META_CACHE_TTL;
|
||||
const failedRecently = this.lastError !== undefined && Date.now() - this.failedAt < FAILURE_CACHE_TTL;
|
||||
if (this.data && (fresh || failedRecently)) {
|
||||
return this.data;
|
||||
}
|
||||
if (!this.data && failedRecently) {
|
||||
throw this.lastError;
|
||||
}
|
||||
|
||||
try {
|
||||
const pointer = await fetchValidated(this.url, pointerSchema);
|
||||
const next = await fetchValidated(pointer.meta_url, metaSchema);
|
||||
this.data = next;
|
||||
this.fetchedAt = Date.now();
|
||||
this.lastError = undefined;
|
||||
return next;
|
||||
} catch (error) {
|
||||
this.lastError = error;
|
||||
this.failedAt = Date.now();
|
||||
if (this.data) {
|
||||
return this.data;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
// The api-client's operation paths already carry the /api prefix, so the
|
||||
// request base is api_root without it.
|
||||
async getBaseUrl(): Promise<string> {
|
||||
const meta = await this.get();
|
||||
return meta.api_root.replace(/\/api\/?$/, '');
|
||||
}
|
||||
|
||||
async getSites(): Promise<MetaSite[]> {
|
||||
const meta = await this.get();
|
||||
return meta.sites;
|
||||
}
|
||||
|
||||
async getConfig(siteCode?: string, clusterCode?: string): Promise<MetaConfig> {
|
||||
const meta = await this.get();
|
||||
if (!siteCode) {
|
||||
if (clusterCode) {
|
||||
throw new UnknownTopologyCodeError(`Cluster '${clusterCode}' was provided without a site`);
|
||||
}
|
||||
return { ...meta.config };
|
||||
}
|
||||
|
||||
const site = meta.sites.find((candidate) => candidate.code === siteCode);
|
||||
if (!site) {
|
||||
throw new UnknownTopologyCodeError(`Unknown site '${siteCode}'`);
|
||||
}
|
||||
if (!clusterCode) {
|
||||
return { ...meta.config, ...site.site_config };
|
||||
}
|
||||
|
||||
const cluster = site.clusters.find((candidate) => candidate.code === clusterCode);
|
||||
if (!cluster) {
|
||||
throw new UnknownTopologyCodeError(`Unknown cluster '${clusterCode}' for site '${siteCode}'`);
|
||||
}
|
||||
return { ...meta.config, ...site.site_config, ...cluster.cluster_config };
|
||||
}
|
||||
|
||||
async getConnections(cores: number, siteCode?: string, clusterCode?: string): Promise<number | undefined> {
|
||||
try {
|
||||
const { connections_math } = await this.getConfig(siteCode, clusterCode);
|
||||
if (!connections_math) {
|
||||
return undefined;
|
||||
}
|
||||
return Math.max(1, Math.floor(evaluateConnectionsMath(connections_math, cores)));
|
||||
} catch (error) {
|
||||
if (error instanceof UnknownTopologyCodeError) {
|
||||
throw error;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
async getPackSizeMib(siteCode?: string, clusterCode?: string): Promise<number | undefined> {
|
||||
try {
|
||||
const config = await this.getConfig(siteCode, clusterCode);
|
||||
return config.restic_pack_size_mib;
|
||||
} catch (error) {
|
||||
if (error instanceof UnknownTopologyCodeError) {
|
||||
throw error;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export const yuccaWellKnown = new YuccaWellKnown();
|
||||
@@ -0,0 +1,24 @@
|
||||
import { evaluateConnectionsMath, isValidConnectionsMath } from '../src/utils/connectionsMath';
|
||||
|
||||
describe('evaluateConnectionsMath', () => {
|
||||
it.each([
|
||||
['16', 4, 16],
|
||||
['cores', 4, 4],
|
||||
['cores * 2', 4, 8],
|
||||
['min(16, cores * 2)', 4, 8],
|
||||
['min(16, cores * 2)', 32, 16],
|
||||
['max(2, cores / 2)', 1, 2],
|
||||
['(cores + 2) * 2 - 1', 3, 9],
|
||||
['min(max(1, cores), 8)', 64, 8],
|
||||
])('evaluates %s with cores=%d to %d', (expression, cores, expected) => {
|
||||
expect(evaluateConnectionsMath(expression, cores)).toBe(expected);
|
||||
});
|
||||
|
||||
it.each(['', 'foo', 'min(1)', 'min(1, 2', '1 +', '2 ** 3', 'cores cores', '1; process.exit(1)'])(
|
||||
'rejects %s',
|
||||
(expression) => {
|
||||
expect(() => evaluateConnectionsMath(expression, 4)).toThrow();
|
||||
expect(isValidConnectionsMath(expression)).toBe(false);
|
||||
},
|
||||
);
|
||||
});
|
||||
@@ -54,7 +54,6 @@ export async function createTestingModule(): Promise<TestContext> {
|
||||
provide: ModuleConfigProvider,
|
||||
useValue: {
|
||||
statePath,
|
||||
yuccaProductionApi: 'http://test:3000',
|
||||
requireLock: false,
|
||||
requireWsAuth: false,
|
||||
},
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
import { YuccaWellKnown } from '../src/wellKnown';
|
||||
|
||||
const meta = {
|
||||
api_root: 'https://backups.futo.cloud/api',
|
||||
config: { restic_pack_size_mib: 16, connections_math: 'min(16, cores * 2)' },
|
||||
default_site: 'father',
|
||||
sites: [
|
||||
{
|
||||
code: 'father',
|
||||
display_name: 'Hetzner Falkenstein-1',
|
||||
description: 'Located in Falkenstein, Germany',
|
||||
rest_url: 'https://rest.htz-fsn1.backups.futo.cloud',
|
||||
default_cluster: 'father-spice',
|
||||
site_config: { restic_pack_size_mib: 64 },
|
||||
clusters: [
|
||||
{
|
||||
code: 'father-spice',
|
||||
display_name: 'Spice',
|
||||
cluster_config: { connections_math: '4' },
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
};
|
||||
|
||||
const jsonResponse = (body: unknown, status = 200) =>
|
||||
({ ok: status >= 200 && status < 300, status, json: () => Promise.resolve(body) }) as Response;
|
||||
|
||||
describe(YuccaWellKnown.name, () => {
|
||||
let sut: YuccaWellKnown;
|
||||
let fetchMock: jest.SpyInstance;
|
||||
|
||||
beforeEach(() => {
|
||||
sut = new YuccaWellKnown();
|
||||
sut.configure('https://meta.example.test/.well-known/yucca.json');
|
||||
fetchMock = jest.spyOn(globalThis, 'fetch').mockImplementation((url) => {
|
||||
if (String(url) === 'https://meta.example.test/.well-known/yucca.json') {
|
||||
return Promise.resolve(jsonResponse({ meta_url: 'https://backups.example.test/api/meta' }));
|
||||
}
|
||||
if (String(url) === 'https://backups.example.test/api/meta') {
|
||||
return Promise.resolve(jsonResponse(meta));
|
||||
}
|
||||
return Promise.reject(new Error(`Unexpected fetch: ${String(url)}`));
|
||||
});
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
fetchMock.mockRestore();
|
||||
});
|
||||
|
||||
it('follows the pointer to /meta and caches the result', async () => {
|
||||
await expect(sut.getBaseUrl()).resolves.toBe('https://backups.futo.cloud');
|
||||
await expect(sut.getSites()).resolves.toHaveLength(1);
|
||||
expect(fetchMock).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it('merges site overrides over the global config', async () => {
|
||||
await expect(sut.getConfig()).resolves.toEqual({
|
||||
restic_pack_size_mib: 16,
|
||||
connections_math: 'min(16, cores * 2)',
|
||||
});
|
||||
await expect(sut.getConfig('father')).resolves.toEqual({
|
||||
restic_pack_size_mib: 64,
|
||||
connections_math: 'min(16, cores * 2)',
|
||||
});
|
||||
await expect(sut.getConfig('father', 'father-spice')).resolves.toEqual({
|
||||
restic_pack_size_mib: 64,
|
||||
connections_math: '4',
|
||||
});
|
||||
});
|
||||
|
||||
it('evaluates connections_math with the given core count', async () => {
|
||||
await expect(sut.getConnections(4)).resolves.toBe(8);
|
||||
await expect(sut.getConnections(32)).resolves.toBe(16);
|
||||
await expect(sut.getPackSizeMib('father', 'father-spice')).resolves.toBe(64);
|
||||
});
|
||||
|
||||
it('re-resolves after configure()', async () => {
|
||||
await sut.getBaseUrl();
|
||||
expect(fetchMock).toHaveBeenCalledTimes(2);
|
||||
|
||||
sut.configure('https://meta.example.test/.well-known/yucca.json');
|
||||
await sut.getBaseUrl();
|
||||
expect(fetchMock).toHaveBeenCalledTimes(4);
|
||||
});
|
||||
|
||||
it('returns undefined helpers and throws get() when discovery is unreachable', async () => {
|
||||
fetchMock.mockRejectedValue(new Error('offline'));
|
||||
|
||||
await expect(sut.get()).rejects.toThrow('offline');
|
||||
await expect(sut.getConnections(4)).resolves.toBeUndefined();
|
||||
await expect(sut.getPackSizeMib()).resolves.toBeUndefined();
|
||||
|
||||
// The failure is memoized: no new fetch attempts within the failure TTL.
|
||||
const callsAfterFailure = fetchMock.mock.calls.length;
|
||||
await expect(sut.get()).rejects.toThrow('offline');
|
||||
expect(fetchMock.mock.calls.length).toBe(callsAfterFailure);
|
||||
});
|
||||
|
||||
it('serves stale data when a refresh fails', async () => {
|
||||
await sut.getBaseUrl();
|
||||
|
||||
sut['fetchedAt'] = 0; // expire the cache
|
||||
fetchMock.mockRejectedValue(new Error('offline'));
|
||||
|
||||
await expect(sut.getBaseUrl()).resolves.toBe('https://backups.futo.cloud');
|
||||
});
|
||||
|
||||
it('rejects a pointer without meta_url', async () => {
|
||||
fetchMock.mockResolvedValue(jsonResponse({ nope: true }));
|
||||
await expect(sut.get()).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('rejects JSON error responses and malformed metadata without replacing stale data', async () => {
|
||||
await sut.getBaseUrl();
|
||||
sut['fetchedAt'] = 0;
|
||||
|
||||
fetchMock.mockResolvedValue(jsonResponse({ error: 'unavailable' }, 503));
|
||||
await expect(sut.getBaseUrl()).resolves.toBe('https://backups.futo.cloud');
|
||||
|
||||
sut['failedAt'] = 0;
|
||||
sut['lastError'] = undefined;
|
||||
fetchMock
|
||||
.mockResolvedValueOnce(jsonResponse({ meta_url: 'https://backups.example.test/api/meta' }))
|
||||
.mockResolvedValueOnce(jsonResponse({ api_root: 123 }));
|
||||
await expect(sut.getBaseUrl()).resolves.toBe('https://backups.futo.cloud');
|
||||
});
|
||||
});
|
||||
@@ -113,6 +113,8 @@ export type ImportRecoveryKeyRequest = {
|
||||
export type RepositoryCreateRequestDto = {
|
||||
name: string;
|
||||
worm: boolean;
|
||||
/** Internal site code from environment metadata */
|
||||
site?: string;
|
||||
paths?: string[];
|
||||
};
|
||||
export type RepositoryMetricsDto = {
|
||||
@@ -143,6 +145,8 @@ export type LocalRepositoryDto = {
|
||||
id: string;
|
||||
worm: boolean;
|
||||
name: string;
|
||||
siteCode: string | null;
|
||||
storageClusterCode: string | null;
|
||||
metrics: RepositoryMetricsDto;
|
||||
meter?: RepositoryMeterDto;
|
||||
backends?: RepositoryBackendsDto;
|
||||
@@ -166,12 +170,15 @@ export type SnapshotDto = {
|
||||
id: string;
|
||||
time: string;
|
||||
paths: string[];
|
||||
tags?: string[];
|
||||
summary?: SnapshotSummaryDto;
|
||||
};
|
||||
export type InspectedLocalRepositoryDto = {
|
||||
id: string;
|
||||
worm: boolean;
|
||||
name: string;
|
||||
siteCode: string | null;
|
||||
storageClusterCode: string | null;
|
||||
metrics: RepositoryMetricsDto;
|
||||
meter?: RepositoryMeterDto;
|
||||
backends?: RepositoryBackendsDto;
|
||||
|
||||
@@ -18,6 +18,8 @@ export class MockProvider extends BaseProvider {
|
||||
id: 'repo1',
|
||||
name: 'My Repository',
|
||||
worm: false,
|
||||
siteCode: 'local',
|
||||
storageClusterCode: 'local-dev',
|
||||
metrics: {
|
||||
sizeBytes: 1337,
|
||||
},
|
||||
|
||||
Generated
+3
@@ -1102,6 +1102,9 @@ importers:
|
||||
tail:
|
||||
specifier: 'catalog:'
|
||||
version: 2.2.6
|
||||
zod:
|
||||
specifier: 'catalog:'
|
||||
version: 4.3.5
|
||||
devDependencies:
|
||||
'@nestjs/cli':
|
||||
specifier: 'catalog:'
|
||||
|
||||
Reference in New Issue
Block a user