All files / apps/api-gateway/src/gateway/http downstream.client.ts

100% Statements 20/20
100% Branches 13/13
100% Functions 5/5
100% Lines 20/20

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 859x         9x             9x                     9x   27x 27x       18x       3x             3x             24x       24x   24x 24x                       3x 1x       2x       21x   21x 1x   20x      
import {
  BadGatewayException,
  GatewayTimeoutException,
  HttpException,
} from '@nestjs/common';
import {
  currentRequestId,
  currentUserId,
  REQUEST_ID_HEADER,
  USER_ID_HEADER,
} from '@app/observability';
 
export const REQUEST_TIMEOUT_MS = 3_000;
 
/**
 * Thin wrapper around fetch encoding the gateway's downstream policy:
 *
 * - every call has a hard timeout (a slow dependency must not exhaust the
 *   gateway's own connections — fail fast and let the client retry);
 * - downstream HTTP errors pass through untouched (status + body), so
 *   validation messages and 404s reach the caller exactly as produced;
 * - transport failures are translated: timeout → 504, unreachable → 502.
 */
export class DownstreamClient {
  constructor(
    private readonly baseUrl: string,
    private readonly serviceName: string,
  ) {}
 
  protected get<T>(path: string): Promise<T> {
    return this.request<T>(path);
  }
 
  protected post<T>(path: string, body: unknown): Promise<T> {
    return this.request<T>(path, {
      method: 'POST',
      body: JSON.stringify(body),
    });
  }
 
  protected patch<T>(path: string, body: unknown): Promise<T> {
    return this.request<T>(path, {
      method: 'PATCH',
      body: JSON.stringify(body),
    });
  }
 
  private async request<T>(path: string, init: RequestInit = {}): Promise<T> {
    const requestId = currentRequestId();
    // Identity crosses the sync boundary the same way the correlation id
    // does; services trust the header because only the gateway can reach
    // them (k8s NetworkPolicy).
    const userId = currentUserId();
    let response: Response;
    try {
      response = await fetch(`${this.baseUrl}${path}`, {
        ...init,
        headers: {
          'content-type': 'application/json',
          // Propagate the caller's correlation id so one user action is
          // traceable across every service it touches.
          ...(requestId ? { [REQUEST_ID_HEADER]: requestId } : {}),
          ...(userId ? { [USER_ID_HEADER]: userId } : {}),
        },
        signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS),
      });
    } catch (error) {
      if (error instanceof DOMException && error.name === 'TimeoutError') {
        throw new GatewayTimeoutException(
          `${this.serviceName} timed out after ${REQUEST_TIMEOUT_MS}ms`,
        );
      }
      throw new BadGatewayException(`${this.serviceName} is unreachable`);
    }
 
    const body: unknown =
      response.status === 204 ? undefined : await response.json();
 
    if (!response.ok) {
      throw new HttpException(body as Record<string, unknown>, response.status);
    }
    return body as T;
  }
}