feat: configure OpenTelemetry (#29)

This commit is contained in:
Paul Makles
2026-02-18 14:49:31 +00:00
committed by GitHub
parent c06ae1a2e9
commit c568a50687
45 changed files with 1722 additions and 318 deletions
-1
View File
@@ -2,7 +2,6 @@ name: ci
on:
pull_request:
branches: [main]
push:
branches: [main]
+1
View File
@@ -1,4 +1,5 @@
node_modules
*.tsbuildinfo
.env.local
.env
+1 -1
View File
@@ -9,7 +9,7 @@ cleanup() {
trap cleanup EXIT INT TERM
# start services
mise restic-api:dev &
pnpm --filter restic-api start:dev &
while ! nc -z localhost "$RESTIC_API_PORT" 2>/dev/null; do
sleep 0.1
+1
View File
@@ -0,0 +1 @@
node-linker=isolated
+20
View File
@@ -8,10 +8,26 @@
"license": "ISC",
"packageManager": "pnpm@10.27.0",
"dependencies": {
"@nestjs/common": "11.0.1",
"@opentelemetry/api": "^1.9.0",
"@opentelemetry/context-async-hooks": "^2.5.0",
"@opentelemetry/core": "^2.5.0",
"@opentelemetry/exporter-logs-otlp-proto": "^0.211.0",
"@opentelemetry/exporter-metrics-otlp-proto": "^0.211.0",
"@opentelemetry/exporter-trace-otlp-proto": "^0.211.0",
"@opentelemetry/instrumentation-pino": "^0.57.0",
"@opentelemetry/propagator-b3": "^2.5.0",
"@opentelemetry/propagator-jaeger": "^2.5.0",
"@opentelemetry/sdk-node": "^0.211.0",
"nestjs-otel": "^8.0.2",
"pino": "^10.3.0",
"pino-pretty": "^13.1.3",
"rxjs": "7.8.1",
"typescript": "^5.9.3",
"zod": "^4.3.5"
},
"devDependencies": {
"@types/express": "^5.0.6",
"@types/node": "^25.0.10"
},
"scripts": {
@@ -25,6 +41,10 @@
"./env": {
"types": "./dist/env.d.ts",
"default": "./dist/env.js"
},
"./otel": {
"types": "./dist/otel/index.d.ts",
"default": "./dist/otel/index.js"
}
}
}
+7
View File
@@ -20,6 +20,13 @@ const schema = z.object({
POSTGRES_PASSWORD: z.string(),
POSTGRES_DATABASE: z.string(),
POSTGRES_SSL: z.union([z.enum(['require', 'allow', 'prefer', 'verify-full']), z.boolean()]).default(false),
OTEL_DEBUG: z.coerce.boolean(),
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);
+7
View File
@@ -0,0 +1,7 @@
export * from './init.js';
export { LoggerRepository } from './logger.repository.js';
export { LoggingInterceptor } from './logging.interceptor.js';
export { WideContextRepository } from './wideContext.repository.js';
export { MetricService, TraceService, Traceable } from 'nestjs-otel';
+72
View File
@@ -0,0 +1,72 @@
import { AsyncLocalStorageContextManager } from '@opentelemetry/context-async-hooks';
import { CompositePropagator, W3CBaggagePropagator, W3CTraceContextPropagator } from '@opentelemetry/core';
import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-proto';
import { OTLPMetricExporter } from '@opentelemetry/exporter-metrics-otlp-proto';
import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-proto';
import { PinoInstrumentation } from '@opentelemetry/instrumentation-pino';
import { B3Propagator } from '@opentelemetry/propagator-b3';
import { JaegerPropagator } from '@opentelemetry/propagator-jaeger';
import { logs, metrics, NodeSDK, tracing } from '@opentelemetry/sdk-node';
import { env } from '../env.js';
const SpanProcessor = env.NODE_ENV === 'development' ? tracing.SimpleSpanProcessor : tracing.BatchSpanProcessor;
const LogProcessor = env.NODE_ENV === 'development' ? logs.SimpleLogRecordProcessor : logs.BatchLogRecordProcessor;
const otelSDK = new NodeSDK({
// metrics
metricReader: new metrics.PeriodicExportingMetricReader({
exporter: new OTLPMetricExporter({
url: env.OTEL_METRICS,
temporalityPreference: metrics.AggregationTemporality.CUMULATIVE,
}),
exportIntervalMillis: 1000,
}),
// tracing
sampler:
env.NODE_ENV === 'development' || env.OTEL_SAMPLE_RATE == 1
? new tracing.AlwaysOnSampler()
: new tracing.TraceIdRatioBasedSampler(env.OTEL_SAMPLE_RATE),
contextManager: new AsyncLocalStorageContextManager(),
textMapPropagator: new CompositePropagator({
propagators: [
new JaegerPropagator(),
new W3CTraceContextPropagator(),
new W3CBaggagePropagator(),
new B3Propagator(),
],
}),
spanProcessor: new SpanProcessor(
new OTLPTraceExporter({
url: env.OTEL_TRACING,
}),
),
// logging
logRecordProcessors: [
new LogProcessor(
new OTLPLogExporter({
url: env.OTEL_LOGGING,
}),
),
],
instrumentations: [new PinoInstrumentation()],
});
import { diag, DiagConsoleLogger, DiagLogLevel } from '@opentelemetry/api';
import { OpenTelemetryModule } from 'nestjs-otel';
if (env.OTEL_DEBUG) {
diag.setLogger(new DiagConsoleLogger(), DiagLogLevel.DEBUG);
}
otelSDK.start();
export const shutdownOtel = () =>
otelSDK.shutdown().then(
() => console.log('SDK shut down successfully'),
(error) => console.log('Error shutting down SDK', error),
);
export default otelSDK;
export const OtelModule = OpenTelemetryModule.forRoot();
@@ -0,0 +1,36 @@
import { Injectable, Scope } from '@nestjs/common';
import pino from 'pino';
import { env } from '../env.js';
@Injectable({ scope: Scope.DEFAULT })
export class LoggerRepository {
private logger: ReturnType<typeof pino>;
constructor() {
this.logger = pino(
env.NODE_ENV === 'development'
? {
transport: {
target: 'pino-pretty',
},
}
: {},
);
}
debug(...args: Parameters<typeof this.logger.info>): void {
this.logger.debug(...args);
}
info(...args: Parameters<typeof this.logger.info>): void {
this.logger.info(...args);
}
warn(...args: Parameters<typeof this.logger.warn>): void {
this.logger.warn(...args);
}
error(...args: Parameters<typeof this.logger.error>): void {
this.logger.error(...args);
}
}
@@ -0,0 +1,67 @@
import { type CallHandler, type ExecutionContext, Injectable, type NestInterceptor, Scope } from '@nestjs/common';
import { type Request } from 'express';
import { randomUUID } from 'node:crypto';
import { Observable, catchError, tap } from 'rxjs';
import { env } from '../env.js';
import { LoggerRepository } from './logger.repository.js';
import { WideContextRepository } from './wideContext.repository.js';
@Injectable({ scope: Scope.REQUEST })
export class LoggingInterceptor implements NestInterceptor {
constructor(
private readonly logger: LoggerRepository,
private readonly wideContext: WideContextRepository,
) {}
intercept(context: ExecutionContext, next: CallHandler): Observable<unknown> {
const startTime = Date.now();
const httpCtx = context.switchToHttp();
const request = httpCtx.getRequest<Request>();
const response = httpCtx.getResponse<Request>();
const event: Record<string, unknown> = {
request_id: request.headers['x-request-id'] ?? randomUUID(),
timestamp: new Date().toISOString(),
method: request.method,
path: request.path,
_msg: `${request.method} ${context.getClass().name}.${context.getHandler().name}`,
};
return next.handle().pipe(
tap(() => {
event.status_code = response.statusCode || -1;
event.outcome = 'success';
event.duration_ms = Date.now() - startTime;
event._msg += ' (OK)';
this.wideContext.applyContext(event);
if ((event.duration_ms as number) > 500) {
event._msg = '[SLOW] ' + event._msg;
this.logger.warn(event);
} else if (env.NODE_ENV === 'development' || Math.random() < env.OTEL_SAMPLE_RATE) {
this.logger.info(event);
}
}),
catchError((error) => {
event.status_code = error.status ?? 500;
event.outcome = 'error';
event.error ??= {
type: error.name,
message: error.message,
code: error.code,
cause: error.cause,
retriable: error.retriable ?? false,
} as never;
event.duration_ms = Date.now() - startTime;
event._msg += ` (ERROR ${error.name})`;
this.wideContext.applyContext(event);
this.logger.error(event);
throw error;
}),
);
}
}
@@ -0,0 +1,22 @@
import { Injectable, Scope } from '@nestjs/common';
@Injectable({ scope: Scope.REQUEST })
export class WideContextRepository {
context: Record<string, unknown> = {};
addContext(key: string, object: unknown) {
this.context[key] = object;
}
assignContext(object: unknown) {
Object.assign(this.context, object);
}
setErrorCause(cause: any) {
this.context['error.cause'] = cause;
}
applyContext(event: any) {
Object.assign(event, this.context);
}
}
+5 -1
View File
@@ -7,7 +7,7 @@
// Environment Settings
// See also https://aka.ms/tsconfig/module
"module": "nodenext",
"module": "commonjs",
"target": "esnext",
// For nodejs:
"lib": ["esnext"],
@@ -31,7 +31,11 @@
// "noFallthroughCasesInSwitch": true,
// "noPropertyAccessFromIndexSignature": true,
"experimentalDecorators": true,
"emitDecoratorMetadata": true,
// Recommended Options
"esModuleInterop": true,
"strict": true,
"jsx": "react-jsx",
"isolatedModules": true,
+16
View File
@@ -28,3 +28,19 @@ services:
timeout: 5s
retries: 30
start_period: 10s
victoria-metrics:
image: victoriametrics/victoria-metrics:latest
ports:
- 8428:8428
victoria-logs:
image: docker.io/victoriametrics/victoria-logs:v1.43.1
ports:
- 9428:9428
command: -storageDataPath=victoria-logs-data
victoria-traces:
image: docker.io/victoriametrics/victoria-traces:latest
ports:
- 10428:10428
+6 -4
View File
@@ -41,11 +41,13 @@ in pkgs.mkShell {
export PLAYWRIGHT_BROWSERS_PATH=${unstablePkgs.playwright-driver.browsers}
export PLAYWRIGHT_SKIP_VALIDATE_HOST_REQUIREMENTS=true
playwrightNpmVersion="$(npm show @playwright/test version)"
echo "❄️ Playwright nix version: ${unstablePkgs.playwright.version}"
echo "📦 Playwright npm version: $playwrightNpmVersion"
playwrightPnpmVersion=($(pnpm list -r @playwright/test | grep playwright))
playwrightPnpmVersion=''${playwrightPnpmVersion[1]}
if [ "${unstablePkgs.playwright.version}" != "$playwrightNpmVersion" ]; then
echo "❄️ Playwright nix version: ${unstablePkgs.playwright.version}"
echo "📦 Playwright npm version: $playwrightPnpmVersion"
if [ "${unstablePkgs.playwright.version}" != "$playwrightPnpmVersion" ]; then
echo "❌ Playwright versions in nix and npm are not the same!"
else
echo "✅ Playwright versions in nix and npm are the same"
-1
View File
@@ -3,7 +3,6 @@
"version": "0.0.1",
"description": "Monorepo for yucca",
"private": true,
"packageManager": "pnpm@10.27.0+sha512.72d699da16b1179c14ba9e64dc71c9a40988cbdc65c264cb0e489db7de917f20dcf4d64d8723625f2969ba52d4b7e2a1170682d9ac2a5dcaeaab732b7e16f04a",
"devDependencies": {
"@eslint/js": "^9.39.2",
"eslint": "^9.39.2",
+949 -57
View File
File diff suppressed because it is too large Load Diff
+1
View File
@@ -10,4 +10,5 @@ onlyBuiltDependencies:
- '@nestjs/core'
- '@scarf/scarf'
- esbuild
- protobufjs
- unrs-resolver
+1
View File
@@ -25,6 +25,7 @@
"@nestjs/core": "^11.0.1",
"@nestjs/jwt": "^11.0.2",
"@nestjs/platform-express": "^11.0.1",
"@opentelemetry/api": "^1.9.0",
"class-transformer": "^0.5.1",
"class-validator": "^0.14.3",
"reflect-metadata": "^0.2.2",
+36 -18
View File
@@ -1,30 +1,48 @@
import { env } from '@common/server/env';
import { Module } from '@nestjs/common';
import {
LoggerRepository,
LoggingInterceptor,
OtelModule,
shutdownOtel,
WideContextRepository,
} from '@common/server/otel';
import { Module, type OnApplicationShutdown } from '@nestjs/common';
import { APP_GUARD, APP_INTERCEPTOR } from '@nestjs/core';
import { JwtModule } from '@nestjs/jwt';
import { AppController } from './controllers/app.controller';
import { AuthGuard } from './middleware/auth.guard';
import { ResticInterceptor } from './middleware/restic.interceptor';
import { LoggerRepository } from './repositories/logger.repository';
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,
secret: env.JWT_SECRET,
}),
];
export const controllers = [AppController];
export const providers = [
WideContextRepository,
LoggerRepository,
StorageRepository,
AuthService,
AppService,
{ provide: APP_GUARD, useClass: AuthGuard },
{ provide: APP_INTERCEPTOR, useClass: ResticInterceptor },
{ provide: APP_INTERCEPTOR, useClass: LoggingInterceptor },
];
@Module({
imports: [
JwtModule.register({
global: true,
secret: env.JWT_SECRET,
}),
],
controllers: [AppController],
providers: [
LoggerRepository,
StorageRepository,
AuthService,
AppService,
{ provide: APP_GUARD, useClass: AuthGuard },
{ provide: APP_INTERCEPTOR, useClass: ResticInterceptor },
],
imports: [OtelModule, ...imports],
controllers,
providers,
})
export class AppModule {}
export class AppModule implements OnApplicationShutdown {
async onApplicationShutdown() {
await shutdownOtel();
}
}
+11 -9
View File
@@ -1,3 +1,4 @@
import { Traceable } from '@common/server/otel';
import {
Controller,
Delete,
@@ -22,6 +23,7 @@ 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) {}
@@ -43,14 +45,14 @@ export class AppController {
@Head(':path/config')
@AuthRoute()
async checkConfig(@Auth() auth: AuthDto, @Res() res: Response): Promise<void> {
const size = await this.service.checkConfig(auth.repository);
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.repository);
const config = await this.service.getConfig(auth);
respondWithObject(config, req, res);
}
@@ -58,21 +60,21 @@ export class AppController {
@AuthRoute()
@HttpCode(HttpStatus.OK)
async saveConfig(@Auth() auth: AuthDto, @Req() req: Request): Promise<void> {
await this.service.saveConfig(auth.repository, req, auth.writeOnce);
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.repository, auth.writeOnce);
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.repository, type);
return this.service.listBlobs(auth, type);
}
@Head(':path/:type/:name')
@@ -82,7 +84,7 @@ export class AppController {
@Param() { type, name }: BlobWithNameParamsDto,
@Res() res: Response,
): Promise<void> {
const size = await this.service.checkBlob(auth.repository, type, name);
const size = await this.service.checkBlob(auth, type, name);
res.set('Content-Length', String(size)).end();
}
@@ -95,7 +97,7 @@ export class AppController {
@Req() req: Request,
@Res() res: Response,
): Promise<void> {
const blob = await this.service.getBlob(auth.repository, type, name, range);
const blob = await this.service.getBlob(auth, type, name, range);
respondWithObject(blob, req, res);
}
@@ -107,13 +109,13 @@ export class AppController {
@Param() { type, name }: BlobWithNameParamsDto,
@Req() req: Request,
): Promise<void> {
await this.service.saveBlob(auth.repository, type, name, req, auth.writeOnce);
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.repository, type, name, auth.writeOnce);
await this.service.deleteBlob(auth, type, name);
}
}
+4 -3
View File
@@ -1,13 +1,14 @@
import '@common/server/otel';
import { env } from '@common/server/env';
import { ValidationPipe } from '@nestjs/common';
import { NestFactory } from '@nestjs/core';
// import { raw } from 'express';
import { env } from '@common/server/env';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.create(AppModule);
app.enableShutdownHooks();
app.useGlobalPipes(new ValidationPipe());
// TODO? app.use(raw({ type: ContentType.Binary, limit: '100mb' }));
await app.listen(env.RESTIC_API_PORT);
}
+2 -1
View File
@@ -3,6 +3,7 @@ import {
CanActivate,
ExecutionContext,
Injectable,
Scope,
SetMetadata,
applyDecorators,
createParamDecorator,
@@ -29,7 +30,7 @@ export const Auth = createParamDecorator((_, context: ExecutionContext): AuthDto
return context.switchToHttp().getRequest<AuthenticatedRequest>().auth;
});
@Injectable()
@Injectable({ scope: Scope.REQUEST })
export class AuthGuard implements CanActivate {
constructor(
private reflector: Reflector,
@@ -1,14 +0,0 @@
import { Injectable, Scope } from '@nestjs/common';
@Injectable({ scope: Scope.TRANSIENT })
export class LoggerRepository {
private context: string;
setContext(context: string) {
this.context = context;
}
debug(...args: any[]): void {
console.debug(`[${this.context}]`, ...args);
}
}
@@ -11,9 +11,11 @@ import {
} from '@aws-sdk/client-s3';
import { Upload } from '@aws-sdk/lib-storage';
import { env } from '@common/server/env';
import { Traceable } from '@common/server/otel';
import { Injectable } from '@nestjs/common';
import { Readable } from 'node:stream';
@Traceable()
@Injectable()
export class StorageRepository {
private client: S3Client;
+176 -65
View File
@@ -1,15 +1,24 @@
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.logger as never, mocks.storage as never);
sut = new AppService(mocks.storage as never, mocks.metricService, mocks.wideContext);
});
it('should exist', () => {
@@ -37,15 +46,21 @@ describe(AppService.name, () => {
});
it('should fail if S3 command throws', async () => {
mocks.storage.checkBucket.mockRejectedValueOnce(void 0);
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);
});
mocks.storage.createBucket.mockRejectedValueOnce(void 0);
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);
});
});
@@ -58,84 +73,125 @@ describe(AppService.name, () => {
describe('checkConfig', () => {
it('should return content length', async () => {
mocks.storage.headObject.mockResolvedValue({ ContentLength: 123, $metadata: void 0 as never });
const result = await sut.checkConfig('repository');
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('repository');
const result = await sut.checkConfig(mockAuth());
expect(result).toBe(0);
});
it('should throw if headObject fails', async () => {
mocks.storage.headObject.mockRejectedValue(void 0);
await expect(sut.checkConfig('repository')).rejects.toThrow();
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 stream = Symbol('Stream');
mocks.storage.getObject.mockResolvedValue(stream as never);
const result = await sut.getConfig('repository');
expect(result).toBe(stream);
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 () => {
mocks.storage.getObject.mockImplementation(() => new Promise((_, reject) => reject()));
await expect(sut.getConfig('repository')).rejects.toThrow();
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 = Symbol('Body');
mocks.storage.putObject.mockResolvedValue(void 0 as never);
await sut.saveConfig('repository', body as never, false);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'config', body, false);
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 = Symbol('Body');
mocks.storage.putObject.mockResolvedValue(void 0 as never);
await sut.saveConfig('repository', body as never, true);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'config', body, true);
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('repository', null as never, true)).rejects.toThrow('Config already exists');
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 () => {
mocks.storage.putObject.mockRejectedValue(new Error('other'));
await expect(sut.saveConfig('repository', null as never, false)).rejects.toThrow();
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.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteConfig('repository', false);
await sut.deleteConfig(mockAuth());
expect(mocks.storage.deleteObject).toHaveBeenCalledWith('repository', 'config');
});
it('should throw when writeOnce and not locks', async () => {
await expect(sut.deleteConfig('repository', true)).rejects.toThrow();
await expect(sut.deleteConfig(mockAuth(true))).rejects.toThrow();
expect(mocks.storage.deleteObject).not.toHaveBeenCalled();
});
it('should throw if deleteObject fails', async () => {
mocks.storage.deleteObject.mockRejectedValue(void 0);
await expect(sut.deleteConfig('repository', false)).rejects.toThrow();
const S3Error = Symbol('S3Error');
mocks.storage.deleteObject.mockRejectedValue(S3Error);
await expect(sut.deleteConfig(mockAuth())).rejects.toThrow();
expect(mocks.wideContext.setErrorCause).toHaveBeenCalledWith(S3Error);
});
});
@@ -150,7 +206,7 @@ describe(AppService.name, () => {
$metadata: void 0 as never,
});
const result = await sut.listBlobs('repository', BlobType.Data);
const result = await sut.listBlobs(mockAuth(), BlobType.Data);
expect(result).toEqual([
{ name: 'abc123', size: 100 },
{ name: 'def456', size: 200 },
@@ -160,127 +216,182 @@ describe(AppService.name, () => {
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('repository', BlobType.Data);
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('repository', BlobType.Data)).rejects.toThrow();
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: [{ Key: 'data/abc123' }],
Contents,
KeyCount: 1,
$metadata: void 0 as never,
});
await expect(sut.listBlobs('repository', BlobType.Data)).rejects.toThrow();
await expect(sut.listBlobs(mockAuth(), BlobType.Data)).rejects.toThrow();
expect(mocks.wideContext.addContext).toHaveBeenCalledWith('contents', Contents);
});
it('should throw if listObjects fails', async () => {
mocks.storage.listObjects.mockRejectedValue(void 0);
await expect(sut.listBlobs('repository', BlobType.Data)).rejects.toThrow();
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('repository', BlobType.Data, 'abc123');
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('repository', BlobType.Data, 'abc123');
const result = await sut.checkBlob(mockAuth(), BlobType.Data, 'abc123');
expect(result).toBe(0);
});
it('should throw if headObject fails', async () => {
mocks.storage.headObject.mockRejectedValue(void 0);
await expect(sut.checkBlob('repository', BlobType.Data, 'abc123')).rejects.toThrow();
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 stream', async () => {
const stream = Symbol('Stream');
mocks.storage.getObject.mockResolvedValue(stream as never);
const result = await sut.getBlob('repository', BlobType.Data, 'abc123');
expect(result).toBe(stream);
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 () => {
const stream = Symbol('Stream');
mocks.storage.getObject.mockResolvedValue(stream as never);
await sut.getBlob('repository', BlobType.Data, 'abc123', 'bytes=0-100');
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 () => {
mocks.storage.getObject.mockImplementation(() => new Promise((_, reject) => reject()));
await expect(sut.getBlob('repository', BlobType.Data, 'abc123')).rejects.toThrow();
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 = Symbol('Body');
mocks.storage.putObject.mockResolvedValue(void 0 as never);
await sut.saveBlob('repository', BlobType.Data, 'abc123', body as never, false);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'data/abc123', body, false, 'abc123');
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(),
false,
'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 = Symbol('Body');
mocks.storage.putObject.mockResolvedValue(void 0 as never);
await sut.saveBlob('repository', BlobType.Data, 'abc123', body as never, true);
expect(mocks.storage.putObject).toHaveBeenCalledWith('repository', 'data/abc123', body, true, 'abc123');
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('repository', BlobType.Data, 'abc123', null as never, true)).rejects.toThrow(
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 () => {
mocks.storage.putObject.mockRejectedValue(new Error('other'));
await expect(sut.saveBlob('repository', BlobType.Data, 'abc123', null as never, false)).rejects.toThrow();
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.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteBlob('repository', BlobType.Data, 'abc123', false);
await sut.deleteBlob(mockAuth(), BlobType.Data, 'abc123');
expect(mocks.storage.deleteObject).toHaveBeenCalledWith('repository', 'data/abc123');
});
it('should allow delete of locks with writeOnce', async () => {
mocks.storage.deleteObject.mockResolvedValue(void 0 as never);
await sut.deleteBlob('repository', BlobType.Locks, 'abc123', true);
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('repository', BlobType.Data, 'abc123', true)).rejects.toThrow();
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.deleteObject.mockRejectedValue(void 0);
await expect(sut.deleteBlob('repository', BlobType.Data, 'abc123', false)).rejects.toThrow();
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);
});
});
});
+101 -61
View File
@@ -1,4 +1,5 @@
import { GetObjectCommandOutput, S3ServiceException } from '@aws-sdk/client-s3';
import { S3ServiceException } from '@aws-sdk/client-s3';
import { MetricService, Traceable, WideContextRepository } from '@common/server/otel';
import {
BadRequestException,
ConflictException,
@@ -6,78 +7,103 @@ import {
Injectable,
NotFoundException,
} from '@nestjs/common';
import { Counter } 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 { S3Error } from 'src/errors';
import { LoggerRepository } from 'src/repositories/logger.repository';
import { StorageRepository } from 'src/repositories/storage.repository';
import { attachMeterToStream, contextFromAuth } from 'src/utils/meters';
import { attachMeterToS3Object, S3RemoteObject } from 'src/utils/s3';
@Traceable()
@Injectable()
export class AppService {
blobsRequestedBytes: Counter;
blobsDownloadedBytes: Counter;
blobsUploadedBytes: Counter;
constructor(
private readonly logger: LoggerRepository,
private readonly storage: StorageRepository,
private readonly metricService: MetricService,
private readonly wideContext: WideContextRepository,
) {
logger.setContext('AppService');
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',
});
}
async createRepository(repository: string, isCreate: boolean): Promise<void> {
if (!isCreate) {
throw new BadRequestException();
throw new BadRequestException('isCreate must be true when creating repository');
}
this.logger.debug(`Creating a new repository at ${repository}`);
let exists: boolean;
try {
exists = await this.storage.checkBucket(repository);
} catch {
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
if (exists) {
throw new ConflictException();
throw new ConflictException('Repository already exists');
}
try {
await this.storage.createBucket(repository);
} catch {
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
deleteRepository(): void {
this.logger.debug('Ignoring repository delete request');
}
async checkConfig(path: string): Promise<number> {
this.logger.debug(`Checking config at ${path}`);
deleteRepository(): void {}
async checkConfig(auth: AuthDto): Promise<number> {
try {
const { ContentLength } = await this.storage.headObject(path, 'config');
const { ContentLength } = await this.storage.headObject(auth.repository, 'config');
return ContentLength || 0;
} catch {
} catch (error) {
this.wideContext.setErrorCause(error);
throw new NotFoundException();
}
}
async getConfig(path: string): Promise<GetObjectCommandOutput> {
this.logger.debug(`Reading repository config at ${path}`);
async getConfig(auth: AuthDto): Promise<S3RemoteObject> {
try {
return await this.storage.getObject(path, 'config');
} catch {
return attachMeterToS3Object(
auth,
await this.storage.getObject(auth.repository, 'config'),
this.blobsRequestedBytes,
this.blobsDownloadedBytes,
);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async saveConfig(path: string, body: Readable, writeOnce: boolean): Promise<void> {
this.logger.debug(`Writing config to repository at ${path}`);
async saveConfig(auth: AuthDto, body: Readable): Promise<void> {
try {
await this.storage.putObject(path, 'config', body, writeOnce);
await this.storage.putObject(
auth.repository,
'config',
attachMeterToStream(body, this.blobsUploadedBytes, contextFromAuth(auth)),
auth.writeOnce,
);
} catch (error) {
this.wideContext.setErrorCause(error);
if (error instanceof S3ServiceException && error.$metadata.httpStatusCode === 412) {
throw new ForbiddenException('Config already exists');
}
@@ -86,32 +112,36 @@ export class AppService {
}
}
async deleteConfig(path: string, writeOnce: boolean): Promise<void> {
this.logger.debug(`Deleting repository config at ${path}`);
if (writeOnce) {
throw new ForbiddenException();
async deleteConfig(auth: AuthDto): Promise<void> {
if (auth.writeOnce) {
throw new ForbiddenException('Not permitted to write to WORM repository');
}
try {
await this.storage.deleteObject(path, 'config');
} catch {
await this.storage.deleteObject(auth.repository, 'config');
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async listBlobs(path: string, type: BlobType): Promise<BlobInfoResponseDto[]> {
this.logger.debug(`Listing repository blobs at ${path} for ${type}`);
async listBlobs(auth: AuthDto, type: BlobType): Promise<BlobInfoResponseDto[]> {
try {
const suffix = `${type}/`;
const { Contents, KeyCount } = await this.storage.listObjects(path, suffix);
const { Contents, KeyCount } = await this.storage.listObjects(auth.repository, suffix);
if (KeyCount === 0) {
return [];
}
if (!Contents || Contents.some(({ Key, Size }) => !Key || !Size)) {
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;
}
@@ -119,38 +149,48 @@ export class AppService {
name: Key!.slice(suffix.length),
size: Size!,
}));
} catch {
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
async checkBlob(path: string, type: BlobType, name: string): Promise<number> {
this.logger.debug(`Checking repository blob at ${path} for ${type}/${name}`);
async checkBlob(auth: AuthDto, type: BlobType, name: string): Promise<number> {
try {
const { ContentLength } = await this.storage.headObject(path, `${type}/${name}`);
const { ContentLength } = await this.storage.headObject(auth.repository, `${type}/${name}`);
return ContentLength || 0;
} catch {
} catch (error) {
this.wideContext.setErrorCause(error);
throw new NotFoundException();
}
}
async getBlob(path: string, type: BlobType, name: string, range?: string): Promise<GetObjectCommandOutput> {
this.logger.debug(`Downloading repository blob at ${path} for ${type}/${name} (range = ${range})`);
async getBlob(auth: AuthDto, type: BlobType, name: string, range?: string): Promise<S3RemoteObject> {
try {
return await this.storage.getObject(path, `${type}/${name}`, range);
} catch {
return attachMeterToS3Object(
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(path: string, type: BlobType, name: string, body: Readable, writeOnce: boolean): Promise<void> {
this.logger.debug(`Uploading repository blob at ${path} for ${type}/${name}`);
async saveBlob(auth: AuthDto, type: BlobType, name: string, body: Readable): Promise<void> {
try {
await this.storage.putObject(path, `${type}/${name}`, body, writeOnce, name);
await this.storage.putObject(
auth.repository,
`${type}/${name}`,
attachMeterToStream(body, this.blobsUploadedBytes, contextFromAuth(auth)),
auth.writeOnce,
name,
);
} catch (error) {
this.wideContext.setErrorCause(error);
if (error instanceof S3ServiceException) {
if (error.$metadata.httpStatusCode === 412) {
throw new ForbiddenException('Blob already exists');
@@ -165,16 +205,16 @@ export class AppService {
}
}
async deleteBlob(path: string, type: BlobType, name: string, writeOnce: boolean): Promise<void> {
this.logger.debug(`Deleting repository blob at ${path} for ${type}/${name}`);
if (writeOnce && type !== BlobType.Locks) {
throw new ForbiddenException();
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 {
await this.storage.deleteObject(path, `${type}/${name}`);
} catch {
await this.storage.deleteObject(auth.repository, `${type}/${name}`);
} catch (error) {
this.wideContext.setErrorCause(error);
throw new S3Error();
}
}
+4 -2
View File
@@ -1,14 +1,16 @@
import { randomUUID } from 'node:crypto';
import { newJwtMock } from '../../test/mocks';
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();
sut = new AuthService(jwt as never);
wideContext = newWideContextMock();
sut = new AuthService(jwt as never, wideContext as never);
});
it('should exist', () => {
+7 -1
View File
@@ -1,13 +1,18 @@
import { WideContextRepository } from '@common/server/otel';
import { BadRequestException, Injectable, UnauthorizedException } from '@nestjs/common';
import { JwtService } from '@nestjs/jwt';
import { plainToInstance } from 'class-transformer';
import { validate } from 'class-validator';
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) {}
constructor(
private readonly jwt: JwtService,
private readonly wideContext: WideContextRepository,
) {}
async authenticate(headers: IncomingHttpHeaders): Promise<AuthDto> {
if (!headers.authorization) {
@@ -38,6 +43,7 @@ export class AuthService {
throw new BadRequestException(errors.flatMap((err) => Object.values(err.constraints ?? {})));
}
this.wideContext.assignContext(contextFromAuth(instance));
return instance;
}
}
+23
View File
@@ -0,0 +1,23 @@
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, context: Attributes) {
const passthrough = new PassThrough({
transform(chunk, _, callback) {
meterBytes.add(chunk.length, context);
callback(null, chunk);
},
});
stream.pipe(passthrough);
return passthrough;
}
export function contextFromAuth(auth: AuthDto) {
return {
customerId: auth.user,
repositoryId: auth.repository,
};
}
+32 -4
View File
@@ -1,9 +1,37 @@
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 attachMeterToS3Object(
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, context);
}
},
};
}
/**
* Process S3 object as web response
@@ -12,7 +40,7 @@ import { ContentType } from '../enum';
* http#ServeContent
* https://pkg.go.dev/net/http#ServeContent
*/
export function respondWithObject(object: GetObjectCommandOutput, request: Request, response: Response) {
export function respondWithObject({ stream, object }: S3RemoteObject, request: Request, response: Response) {
if (request.headers['if-none-match'] === object.ETag) {
return response.send(HttpStatus.NOT_MODIFIED);
}
@@ -38,9 +66,9 @@ export function respondWithObject(object: GetObjectCommandOutput, request: Reque
response.set('Content-Length', object.ContentLength.toString());
}
const webStream = object.Body?.transformToWebStream();
if (webStream) {
Readable.fromWeb(webStream as ReadableStream).pipe(response);
const readable = stream();
if (readable) {
readable.pipe(response);
} else {
return response.send(HttpStatus.INTERNAL_SERVER_ERROR);
}
+10 -3
View File
@@ -1,10 +1,12 @@
import { MetricService } from '@common/server/otel';
import { INestApplication, ValidationPipe } 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 { AppModule } from './../src/app.module';
import { controllers, imports, providers } from './../src/app.module';
import { newMetricServiceMock } from './mocks';
const makeAuthHeader = (token: string) => 'Basic ' + Buffer.from(`_:${token}`).toString('base64');
@@ -17,8 +19,13 @@ describe('AppController (e2e)', () => {
beforeEach(async () => {
const moduleFixture: TestingModule = await Test.createTestingModule({
imports: [AppModule],
}).compile();
imports,
controllers,
providers: [MetricService, ...providers],
})
.overrideProvider(MetricService)
.useValue(newMetricServiceMock())
.compile();
app = moduleFixture.createNestApplication();
app.useGlobalPipes(new ValidationPipe());
+39 -4
View File
@@ -1,4 +1,6 @@
import { LoggerRepository } from 'src/repositories/logger.repository';
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>;
@@ -7,10 +9,16 @@ export const newJwtMock = () => ({
verifyAsync: jest.fn(),
});
export const newWideContextMock = () => ({
assignContext: jest.fn(),
});
export const newLoggerRepositoryMock = (): jest.Mocked<RepositoryInterface<LoggerRepository>> => {
return {
setContext: jest.fn(),
debug: jest.fn(),
error: jest.fn(),
info: jest.fn(),
warn: jest.fn(),
};
};
@@ -19,17 +27,44 @@ export const newStorageRepositoryMock = (): jest.Mocked<RepositoryInterface<Stor
checkBucket: jest.fn(),
createBucket: jest.fn(),
deleteObject: jest.fn(),
getObject: 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(),
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(),
});
export const newMocks = () => {
return {
logger: newLoggerRepositoryMock(),
storage: newStorageRepositoryMock(),
metricService: newMetricServiceMock(),
wideContext: newWideContextRepositoryMock(),
};
};
+6
View File
@@ -6,6 +6,12 @@ msgstr ""
"Content-Transfer-Encoding: 8bit\n"
"X-Generator: @lingui/cli\n"
"Language: test\n"
"Project-Id-Version: \n"
"Report-Msgid-Bugs-To: \n"
"PO-Revision-Date: \n"
"Last-Translator: \n"
"Language-Team: \n"
"Plural-Forms: \n"
#: src/routes/+page.svelte:40
msgid "{num, plural, one {# item} other {# items}}"
+1 -1
View File
@@ -14,7 +14,7 @@
let value = $state("Loading from API...");
onMount(() => {
hello().then((v) => (value = v.data!));
hello().then((v) => (value = v));
});
import { locale } from "svelte-i18n-lingui";
-13
View File
@@ -1,13 +0,0 @@
import { page } from 'vitest/browser';
import { describe, expect, it } from 'vitest';
import { render } from 'vitest-browser-svelte';
import Page from './+page.svelte';
describe('/+page.svelte', () => {
it('should render h1', async () => {
render(Page);
const heading = page.getByRole('heading', { level: 1 });
await expect.element(heading).toBeInTheDocument();
});
});
+1
View File
@@ -24,6 +24,7 @@
"@nestjs/jwt": "^11.0.2",
"@nestjs/platform-express": "^11.0.1",
"@nestjs/swagger": "^11.2.5",
"@opentelemetry/api": "^1.9.0",
"class-transformer": "^0.5.1",
"class-validator": "^0.14.3",
"kysely": "0.28.2",
+25 -10
View File
@@ -1,24 +1,39 @@
import { env } from '@common/server/env';
import { LoggerRepository, LoggingInterceptor, OtelModule, WideContextRepository } from '@common/server/otel';
import { Module } from '@nestjs/common';
import { APP_INTERCEPTOR } from '@nestjs/core';
import { JwtModule } from '@nestjs/jwt';
import { KyselyModule } from 'nestjs-kysely';
import { AppController } from './controllers/app.controller';
import { DatabaseRepository } from './repositories/database.repository';
import { DummyRepository } from './repositories/dummy.repository';
import { LoggerRepository } from './repositories/logger.repository';
import { AppService } from './services/app.service';
import { DatabaseService } from './services/database.service';
import { getKyselyConfig } from './utils/database';
export const imports = [
JwtModule.register({
global: true,
secret: env.JWT_SECRET,
}),
KyselyModule.forRoot(getKyselyConfig()),
];
export const controllers = [AppController];
export const providers = [
WideContextRepository,
LoggerRepository,
DatabaseRepository,
DummyRepository,
DatabaseService,
AppService,
{ provide: APP_INTERCEPTOR, useClass: LoggingInterceptor },
];
@Module({
imports: [
JwtModule.register({
global: true,
secret: env.JWT_SECRET,
}),
KyselyModule.forRoot(getKyselyConfig()),
],
controllers: [AppController],
providers: [LoggerRepository, DatabaseRepository, DummyRepository, DatabaseService, AppService],
imports: [OtelModule, ...imports],
controllers,
providers,
})
export class AppModule {}
+2
View File
@@ -1,3 +1,5 @@
import '@common/server/otel';
import { env } from '@common/server/env';
import { ValidationPipe } from '@nestjs/common';
import { NestFactory } from '@nestjs/core';
@@ -1,27 +1,25 @@
import { env } from '@common/server/env';
import { LoggerRepository } from '@common/server/otel';
import { Injectable } from '@nestjs/common';
import { FileMigrationProvider, Kysely, Migrator } from 'kysely';
import { InjectKysely } from 'nestjs-kysely';
import { readdir } from 'node:fs/promises';
import { join } from 'node:path';
import { DB } from 'src/schema';
import { LoggerRepository } from './logger.repository';
@Injectable()
export class DatabaseRepository {
constructor(
@InjectKysely() private db: Kysely<DB>,
private logger: LoggerRepository,
) {
this.logger.setContext(DatabaseRepository.name);
}
) {}
async shutdown() {
await this.db.destroy();
}
async runMigrations(): Promise<void> {
this.logger.log('Running migrations');
this.logger.debug('Running migrations');
const migrator = this.createMigrator();
const { error, results } = await migrator.migrateToLatest();
@@ -44,7 +42,7 @@ export class DatabaseRepository {
}
}
this.logger.log('Finished running migrations');
this.logger.info('Finished running migrations');
}
private createMigrator(): Migrator {
@@ -1,26 +0,0 @@
import { Injectable, Scope } from '@nestjs/common';
@Injectable({ scope: Scope.TRANSIENT })
export class LoggerRepository {
private context: string;
setContext(context: string) {
this.context = context;
}
log(...args: any[]): void {
console.log(`[${this.context}]`, ...args);
}
debug(...args: any[]): void {
console.debug(`[${this.context}]`, ...args);
}
warn(...args: any[]): void {
console.warn(`[${this.context}]`, ...args);
}
error(...args: any[]): void {
console.error(`[${this.context}]`, ...args);
}
}
+2 -4
View File
@@ -1,15 +1,13 @@
import { LoggerRepository } from '@common/server/otel';
import { Injectable } from '@nestjs/common';
import { DummyRepository } from 'src/repositories/dummy.repository';
import { LoggerRepository } from 'src/repositories/logger.repository';
@Injectable()
export class AppService {
constructor(
private readonly logger: LoggerRepository,
private readonly dummy: DummyRepository,
) {
logger.setContext('AppService');
}
) {}
hello(): Promise<string> {
this.logger.debug('Hello, World!');
+10 -3
View File
@@ -1,17 +1,24 @@
import { MetricService } from '@common/server/otel';
import { INestApplication, ValidationPipe } from '@nestjs/common';
import { Test, TestingModule } from '@nestjs/testing';
import { DatabaseRepository } from 'src/repositories/database.repository';
import request from 'supertest';
import { App } from 'supertest/types';
import { AppModule } from './../src/app.module';
import { controllers, imports, providers } from './../src/app.module';
import { newMetricServiceMock } from './mocks';
describe('AppController (e2e)', () => {
let app: INestApplication<App>;
beforeEach(async () => {
const moduleFixture: TestingModule = await Test.createTestingModule({
imports: [AppModule],
}).compile();
imports,
controllers,
providers: [MetricService, ...providers],
})
.overrideProvider(MetricService)
.useValue(newMetricServiceMock())
.compile();
app = moduleFixture.createNestApplication();
app.useGlobalPipes(new ValidationPipe());
+9 -2
View File
@@ -1,5 +1,5 @@
import { LoggerRepository } from '@common/server/otel';
import { DummyRepository } from 'src/repositories/dummy.repository';
import { LoggerRepository } from 'src/repositories/logger.repository';
export type RepositoryInterface<T extends object> = Pick<T, keyof T>;
@@ -9,8 +9,10 @@ export const newJwtMock = () => ({
export const newLoggerRepositoryMock = (): jest.Mocked<RepositoryInterface<LoggerRepository>> => {
return {
setContext: jest.fn(),
debug: jest.fn(),
error: jest.fn(),
info: jest.fn(),
warn: jest.fn(),
};
};
@@ -20,10 +22,15 @@ export const newDummyRepositoryMock = (): jest.Mocked<RepositoryInterface<DummyR
};
};
export const newMetricServiceMock = () => ({
getCounter: jest.fn().mockReturnValue({ add: jest.fn() }),
});
export const newMocks = () => {
return {
logger: newLoggerRepositoryMock(),
dummy: newDummyRepositoryMock(),
metrics: newMetricServiceMock(),
};
};
+1 -1
View File
@@ -6,7 +6,7 @@
"main": "dist/index.js",
"types": "dist/index.d.ts",
"scripts": {
"generate": "oazapfts ./openapi-specs.json src/fetch-client.ts",
"generate": "oazapfts --optimistic ./openapi-specs.json src/fetch-client.ts",
"build": "tsc"
},
"keywords": [],
+2 -2
View File
@@ -15,7 +15,7 @@ export const servers = {
server1: "/api"
};
export function hello(opts?: Oazapfts.RequestOpts) {
return oazapfts.fetchText("/", {
return oazapfts.ok(oazapfts.fetchText("/", {
...opts
});
}));
}