diff --git a/apps/cli/docs/stack-commands.md b/apps/cli/docs/stack-commands.md index f6d4e42511..05933fc4c0 100644 --- a/apps/cli/docs/stack-commands.md +++ b/apps/cli/docs/stack-commands.md @@ -237,13 +237,14 @@ lint transaction (always rolled back). It does not launch a client binary. ## Reading stack logs -`supabase stack logs` prints the retained stdout/stderr of composition members and -exits; `-f/--follow` then streams new lines until interrupted. Select `--stack ` -or `--stack-id `; the repeatable `--service ` can include -standalone services too. History is read from the persisted log files, so it works -while the stack is stopped; `--follow` requires a running owner and fails before -printing anything without one. Neither mode starts an owner or service, and Ctrl-C -leaves services running. +`supabase stack logs` prints the retained stdout/stderr of composition members, plus +the `gateway` request lines of the shared API port, and exits; `-f/--follow` then +streams new lines until interrupted. Select `--stack ` or `--stack-id `; the +repeatable `--service ` can include standalone services too, and +`--service gateway` reads only the request lines. History is read from the persisted +log files, so it works while the stack is stopped; `--follow` requires a running owner +and fails before printing anything without one. Neither mode starts an owner or +service, and Ctrl-C leaves services running. `--tail N` (default 200) keeps the newest lines across the selected services and `--since` takes a duration (`10m`, `1h30m`), an ISO-8601 time, or `start` for each diff --git a/apps/cli/src/command-internal/db-config.integration.test.ts b/apps/cli/src/command-internal/db-config.integration.test.ts index 1002a2b5d7..8ce8f853f6 100644 --- a/apps/cli/src/command-internal/db-config.integration.test.ts +++ b/apps/cli/src/command-internal/db-config.integration.test.ts @@ -33,6 +33,7 @@ import { mockTty, } from "../../tests/helpers/mocks.ts"; import { VALID_TOKEN, mockCommandSettings } from "../../tests/helpers/command-mocks.ts"; +import { unusedGateway } from "../../tests/helpers/unused-stack.ts"; import { DebugFlag, DnsResolverFlag, @@ -489,6 +490,7 @@ describe("dbConfigResolver (db-url under the stack backend)", () => { }, stop: unused, destroy: unused, + gateway: unusedGateway, commands: { run: () => unused }, }; return Layer.succeed(StackApi, { diff --git a/apps/cli/src/commands/db/dump/dump.integration.test.ts b/apps/cli/src/commands/db/dump/dump.integration.test.ts index e020809e10..76a7e9518e 100644 --- a/apps/cli/src/commands/db/dump/dump.integration.test.ts +++ b/apps/cli/src/commands/db/dump/dump.integration.test.ts @@ -16,6 +16,7 @@ import { import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { mockOutput, mockTty, processEnvLayer } from "../../../../tests/helpers/mocks.ts"; +import { unusedGateway } from "../../../../tests/helpers/unused-stack.ts"; import { VALID_REF, mockCommandSettings, @@ -153,6 +154,7 @@ const managedDumpStackApi = (runtime: "native" | "docker") => { }, stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), + gateway: unusedGateway, commands: { run: runCommand }, } satisfies Stack; return Layer.succeed(StackApi, { diff --git a/apps/cli/src/commands/db/reset/reset.integration.test.ts b/apps/cli/src/commands/db/reset/reset.integration.test.ts index 1f73deb1b8..d3fed41fb4 100644 --- a/apps/cli/src/commands/db/reset/reset.integration.test.ts +++ b/apps/cli/src/commands/db/reset/reset.integration.test.ts @@ -39,6 +39,7 @@ import { sequentialExecBatch, transportFailure, } from "../../../../tests/helpers/command-mocks.ts"; +import { unusedGateway } from "../../../../tests/helpers/unused-stack.ts"; import { CommandPlatformApi } from "../../../auth/command-platform-api.service.ts"; import { CommandPlatformApiFactory } from "../../../auth/command-platform-api-factory.service.ts"; import { ProjectRefNotLinkedError } from "../../../config/project-ref.errors.ts"; @@ -808,6 +809,7 @@ function mockResetStackApi(opts: { }, stop: Effect.die("unused"), destroy: Effect.die("unused"), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, }; return { diff --git a/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts b/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts index f6838085f6..74e8cee717 100644 --- a/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts +++ b/apps/cli/src/commands/db/schema/declarative/generate/generate.integration.test.ts @@ -14,6 +14,7 @@ import { } from "effect"; import { StackError, type DatabaseInstance, type Stack } from "@supabase/stack/effect"; import { stripAnsi } from "../../../../../../tests/helpers/ansi.ts"; +import { unusedGateway } from "../../../../../../tests/helpers/unused-stack.ts"; import { alwaysReadyHttpClientLayer, @@ -164,6 +165,7 @@ function generateStackApi(workdir: string) { }, stop: unusedStack, destroy: unusedStack, + gateway: unusedGateway, commands: { run: unusedStackFn }, }; const identity = { projectRoot: workdir, branchContext: "main", stackName: "default" }; diff --git a/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts b/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts index c7035c2c4f..4d0321f07d 100644 --- a/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts +++ b/apps/cli/src/commands/db/schema/declarative/sync/sync.integration.test.ts @@ -14,6 +14,7 @@ import { } from "effect"; import { stripAnsi } from "../../../../../../tests/helpers/ansi.ts"; +import { unusedGateway } from "../../../../../../tests/helpers/unused-stack.ts"; import { alwaysReadyHttpClientLayer, defaultLocalResetRoute, @@ -166,6 +167,7 @@ function syncStackApi(workdir: string, port: number) { }, stop: unusedSync, destroy: unusedSync, + gateway: unusedGateway, commands: { run: unusedSyncFn }, }; const identity = { projectRoot: workdir, branchContext: "main", stackName: "default" }; diff --git a/apps/cli/src/commands/db/start/start.integration.test.ts b/apps/cli/src/commands/db/start/start.integration.test.ts index 48f591f4dc..eeaaefa772 100644 --- a/apps/cli/src/commands/db/start/start.integration.test.ts +++ b/apps/cli/src/commands/db/start/start.integration.test.ts @@ -23,6 +23,7 @@ import { mockProcessControl, mockRuntimeInfo, } from "../../../../tests/helpers/mocks.ts"; +import { unusedGateway } from "../../../../tests/helpers/unused-stack.ts"; import { mockCommandSettings, mockLocalDockerEngineUnavailableLayer, @@ -1733,6 +1734,7 @@ describe("db start stack backend", () => { }, stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, }; return { stack, state }; diff --git a/apps/cli/src/commands/experimental/stack/logs/SIDE_EFFECTS.md b/apps/cli/src/commands/experimental/stack/logs/SIDE_EFFECTS.md index 7684224b65..00804f94f6 100644 --- a/apps/cli/src/commands/experimental/stack/logs/SIDE_EFFECTS.md +++ b/apps/cli/src/commands/experimental/stack/logs/SIDE_EFFECTS.md @@ -8,9 +8,11 @@ starts/stops a service. ## Selection and files Select the current project/branch/name, `--stack `, or `--stack-id `. -The selectors are mutually exclusive. By default, only composition members are -included. `--service ` is repeatable and can also select -standalone instances; a value that matches no saved instance fails with status 1. +The selectors are mutually exclusive. By default, composition members are +included, plus `gateway` (the shared API port's request lines) once the stack has +claimed that port. `--service ` is repeatable and can also select +standalone instances or `gateway`; a value that matches no saved instance, or +`gateway` without the shared API port, fails with status 1. A missing stack fails with status 1. Reads saved definitions under `/stacks//` and diff --git a/apps/cli/src/commands/experimental/stack/logs/logs.command.ts b/apps/cli/src/commands/experimental/stack/logs/logs.command.ts index 0e3d972b62..96ae1d2583 100644 --- a/apps/cli/src/commands/experimental/stack/logs/logs.command.ts +++ b/apps/cli/src/commands/experimental/stack/logs/logs.command.ts @@ -17,7 +17,7 @@ const config = { ), service: Flag.string("service").pipe( Flag.withDescription( - "Read one service kind or instance ID; repeat to select several. Defaults to composition members.", + "Read one service kind, instance ID, or gateway (shared API requests); repeat to select several. Defaults to composition members and gateway.", ), Flag.atLeast(0), ), diff --git a/apps/cli/src/commands/experimental/stack/logs/logs.handler.ts b/apps/cli/src/commands/experimental/stack/logs/logs.handler.ts index 779791e276..1b836ec14e 100644 --- a/apps/cli/src/commands/experimental/stack/logs/logs.handler.ts +++ b/apps/cli/src/commands/experimental/stack/logs/logs.handler.ts @@ -1,5 +1,10 @@ import { Clock, Effect, Option, Path, Stream } from "effect"; -import { streamStackLogs, type SavedStack, type StackLogRecord } from "@supabase/stack/effect"; +import { + gatewayLog, + streamStackLogs, + type SavedStack, + type StackLogRecord, +} from "@supabase/stack/effect"; import { Output } from "../../../../shared/output/output.service.ts"; import { OutputFlag } from "../../../../command-internal/global-flags.ts"; import { dim } from "../../../../command-internal/colors.ts"; @@ -50,9 +55,16 @@ const select = ( service: creation.service, launchId, }); + // The owner records the shared API listener's requests once the stack has claimed its port. + const gateway: ReadonlyArray = definition.ports.some(({ key }) => key === "api") + ? [{ id: gatewayLog.instanceId, service: gatewayLog.service, launchId: undefined }] + : []; if (requested.length === 0) { const members = new Set(definition.composition.members.map(({ id }) => id)); - const selected = definition.instances.filter(({ id }) => members.has(id)).map(subject); + const selected = [ + ...definition.instances.filter(({ id }) => members.has(id)).map(subject), + ...gateway, + ]; return selected.length === 0 ? Effect.fail( new StackCommandLogsError({ @@ -65,17 +77,20 @@ const select = ( } const unmatched = requested.find( (value) => - !definition.instances.some(({ id, creation }) => id === value || creation.service === value), + !definition.instances.some( + ({ id, creation }) => id === value || creation.service === value, + ) && !gateway.some(({ service }) => service === value), ); if (unmatched !== undefined) return Effect.fail( new StackCommandLogsError({ reason: "flags", message: `No service matches ${unmatched}.` }), ); - return Effect.succeed( - definition.instances + return Effect.succeed([ + ...definition.instances .filter(({ id, creation }) => requested.includes(id) || requested.includes(creation.service)) .map(subject), - ); + ...gateway.filter(({ service }) => requested.includes(service)), + ]); }; export const stackLogs = Effect.fn("experimental.stack.logs")(function* (flags: StackLogsFlags) { @@ -189,7 +204,7 @@ export const stackLogs = Effect.fn("experimental.stack.logs")(function* (flags: ]), ); const missing = selected - .filter(({ id }) => !handles.has(id)) + .filter(({ id }) => id !== gatewayLog.instanceId && !handles.has(id)) .map(({ id, service }) => `${service} (${id})`); if (missing.length === selected.length) return yield* new StackCommandLogsError({ @@ -203,13 +218,14 @@ export const stackLogs = Effect.fn("experimental.stack.logs")(function* (flags: `Not following ${missing.join(", ")}, which the running stack does not serve.`, ); const streams = selected.flatMap(({ id, service }) => { - const handle = handles.get(id); - if (handle === undefined) return []; + const readLogs = + id === gatewayLog.instanceId ? stack.gateway.readLogs : handles.get(id)?.readLogs; + if (readLogs === undefined) return []; const from = printed.get(id); // Without history, the owner pins the start at its end under the writer's lock. const start = flags.tail === 0 ? { tail: 0 } : from === undefined ? {} : { from }; return [ - handle.readLogs({ follow: true, ...start, ...sinceTime }).pipe( + readLogs({ follow: true, ...start, ...sinceTime }).pipe( Stream.filter( (record) => isAfter(record, from) && isFromLaunch(record, launches.get(id)), ), diff --git a/apps/cli/src/commands/experimental/stack/logs/logs.integration.test.ts b/apps/cli/src/commands/experimental/stack/logs/logs.integration.test.ts index 35d5e069ef..6e88b68ba5 100644 --- a/apps/cli/src/commands/experimental/stack/logs/logs.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/logs/logs.integration.test.ts @@ -96,6 +96,8 @@ const fixture = Effect.fn("StackLogsTest.fixture")(function* (options: { readonly stateRoot: string; readonly stackId: string; }) => ReadonlyArray>; + readonly ports?: SavedStack["ports"]; + readonly gatewayLogs?: (options?: ReadLogsOptions) => Stream.Stream; }) { const fs = yield* FileSystem.FileSystem; const path = yield* Path.Path; @@ -113,6 +115,7 @@ const fixture = Effect.fn("StackLogsTest.fixture")(function* (options: { const definition: SavedStack = { ...found.value.definition, instances: options.instances, + ports: options.ports ?? [], composition: { members: options.members.map((id) => ({ id, activation: "eager" as const })), dependencies: [], @@ -147,6 +150,7 @@ const fixture = Effect.fn("StackLogsTest.fixture")(function* (options: { ...stack.services, list: Effect.succeed(options.handles?.({ stateRoot, stackId: stack.id }) ?? []), }, + gateway: { readLogs: options.gatewayLogs ?? stack.gateway.readLogs }, }), ), }), @@ -446,6 +450,111 @@ describe("stack logs", () => { }).pipe(Effect.scoped, Effect.provide(live)), ); + const apiPort = [{ key: "api", host: "127.0.0.1", port: 54321 }]; + const request = + '127.0.0.1 - - [29/Sep/2026:10:00:00 +0000] "GET /rest/v1/ HTTP/1.1" 200 2 "-" "curl/8.7.1" 3ms'; + + it.live("includes the gateway's requests by default and selects them with --service", () => + Effect.gen(function* () { + const f = yield* fixture({ + instances: [mail("mail-a")], + members: ["mail-a"], + ports: apiPort, + }); + yield* f.writeSegment("mail", "mail-a", [launch(t0), line(t0 + 1, "stdout", "mail up")]); + yield* f.writeSegment("gateway", "gateway", [launch(t0), line(t0 + 2, "stdout", request)]); + + const all = yield* f.run({}); + const gateway = yield* f.run({ service: ["gateway"] }, "stream-json"); + + expect(lines(all.output.stdoutText).map((value) => value.replace(/\d{2}:\S+ /u, ""))).toEqual( + [ + "gateway | --- launch 1 ---", + "mail | --- launch 1 ---", + "mail | mail up", + `gateway | ${request}`, + ], + ); + expect(eventLines(gateway.output.events)).toEqual([request]); + expect(gateway.output.events.at(-1)).toMatchObject({ + service: "gateway", + instance_id: "gateway", + }); + }).pipe(Effect.scoped, Effect.provide(live)), + ); + + it.live("rejects --service gateway when the stack has no shared API listener", () => + Effect.gen(function* () { + const f = yield* fixture({ instances: [mail("mail-a")], members: ["mail-a"] }); + yield* f.writeSegment("mail", "mail-a", [launch(t0), line(t0 + 1, "stdout", "mail up")]); + + const selected = yield* f.run({ service: ["gateway"] }); + const all = yield* f.run({}, "stream-json"); + + expect(failure(selected.exit)).toMatchObject({ + reason: "flags", + message: "No service matches gateway.", + }); + expect(eventLines(all.output.events)).toEqual(["mail up"]); + }).pipe(Effect.scoped, Effect.provide(live)), + ); + + it.live("follows the gateway's new requests through the owner", () => + Effect.gen(function* () { + const f = yield* fixture({ + instances: [], + members: [], + ports: apiPort, + running: true, + gatewayLogs: () => + Stream.make({ + kind: "stdout" as const, + timestamp: iso(t0 + 100), + launchId: 1, + text: request, + position: { generation: 1, byteOffset: 1_000 }, + }), + }); + + const { exit, output } = yield* f.run({ follow: true, tail: 0 }, "stream-json"); + + expect(Exit.isSuccess(exit)).toBe(true); + expect(output.events).toEqual([ + expect.objectContaining({ service: "gateway", line: request, source: "live" }), + ]); + }).pipe(Effect.scoped, Effect.provide(live)), + ); + + it.live("keeps every owner run's gateway requests with --since start", () => + Effect.gen(function* () { + const f = yield* fixture({ + instances: [{ ...mail("mail-a"), launchId: 2 }], + members: ["mail-a"], + ports: apiPort, + }); + yield* f.writeSegment("mail", "mail-a", [ + launch(t0, 1), + line(t0 + 1, "stdout", "mail old", 1), + launch(t0 + 10, 2), + line(t0 + 11, "stdout", "mail new", 2), + ]); + yield* f.writeSegment("gateway", "gateway", [ + launch(t0), + line(t0 + 2, "stdout", "first run request"), + launch(t0 + 10), + line(t0 + 12, "stdout", "second run request"), + ]); + + const { output } = yield* f.run({ since: Option.some("start") }, "stream-json"); + + expect(eventLines(output.events)).toEqual([ + "first run request", + "mail new", + "second run request", + ]); + }).pipe(Effect.scoped, Effect.provide(live)), + ); + it.live("follows --since start without delayed records of an earlier launch", () => Effect.gen(function* () { const record = (offset: number, launchId: number, text: string): LogRecord => ({ diff --git a/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts b/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts index 352518aadb..bce7ca9e63 100644 --- a/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/prepare/prepare.integration.test.ts @@ -17,6 +17,7 @@ import { containerEngineSpawner, type ContainerEngineState, } from "../../../../../tests/helpers/child-process-spawner.ts"; +import { unusedGateway } from "../../../../../tests/helpers/unused-stack.ts"; import { mockCommandSettings, mockTelemetryStateTracked, @@ -157,6 +158,7 @@ const makeFixture = (root: string, options: FixtureOptions = {}) => { }, stop: Effect.void, destroy: Effect.succeed({ runtimeCleanup: "complete" as const }), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, }; const output = mockOutput(); diff --git a/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md b/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md index 63a51546cf..e7e990fec7 100644 --- a/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md +++ b/apps/cli/src/commands/experimental/stack/start/SIDE_EFFECTS.md @@ -41,7 +41,13 @@ The owner persists each service's output under `$SUPABASE_HOME/stacks/ most about 10 MiB (plus the segment being written) per service instance. Destroying an instance or the stack deletes those logs; stopping the stack and resetting database data keep them. PostgREST runs with `PGRST_LOG_LEVEL=info`, so every request line, query string included, is -persisted and shipped to Analytics. +persisted and shipped to Analytics. The shared API port also records one nginx combined line plus +duration (`... "GET /rest/v1/todos?select=* HTTP/1.1" 200 126 "-" "curl/8.7.1" 12ms`) per request +and WebSocket upgrade under `logs/gateway/gateway/`, with the same retention; credential query and +fragment values (`apikey`, `token`, `token_hash`, `code`, access, refresh, ID, and provider +tokens, and the `X-Amz-Signature`, `X-Amz-Credential`, and `X-Amz-Security-Token` of S3 +presigned URLs), in the request target and the Referer, are written as `redacted`, as are values +that themselves carry such a pair (a `redirect_to` URL with a token) and URL userinfo. When an owner starts a stack saved with a Vector instance, it removes that instance, its composition members, dependencies and port claims from `state.json`, and its stack-owned Vector config files under `data//runtime/vector/`; its containers go with the stack's container sweep. @@ -89,11 +95,12 @@ Studio requires REST; excluding REST while keeping Studio fails before stopping ## Service logs in Analytics -The owner ships the persisted Auth, REST, Realtime, Storage, Functions, and database output lines -(not launch or lost markers) to Analytics' `POST /api/logs` ingest endpoint on its direct backend, -using the Analytics API key and the legacy Logflare source names (`gotrue.logs.prod`, -`postgREST.logs.prod`, `realtime.logs.prod`, `storage.logs.prod.2`, `deno-relay-logs`, -`postgres.logs`) with the legacy per-service field remaps. This applies to the Docker, Podman, and +The owner ships the persisted Auth, REST, Realtime, Storage, Functions, database, and gateway +output lines (not launch or lost markers) to Analytics' `POST /api/logs` ingest endpoint on its +direct backend, using the Analytics API key and the legacy Logflare source names +(`gotrue.logs.prod`, `postgREST.logs.prod`, `realtime.logs.prod`, `storage.logs.prod.2`, +`deno-relay-logs`, `postgres.logs`, and `cloudflare.logs.prod` for Studio's API Gateway page) with +the legacy per-service field remaps. This applies to the Docker, Podman, and native runtimes. Shipping runs only while the composed Analytics service is running and healthy; each instance keeps its position in `logs///cursor.json`, so lines written while Analytics is stopped, starting, or unhealthy are shipped with their original timestamps once diff --git a/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts b/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts index 022a52ae47..dc425005fb 100644 --- a/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/start/start-export-pointer.integration.test.ts @@ -19,6 +19,7 @@ import { mockTty, } from "../../../../../tests/helpers/mocks.ts"; import { mockTelemetryStateTracked } from "../../../../../tests/helpers/command-mocks.ts"; +import { unusedGateway } from "../../../../../tests/helpers/unused-stack.ts"; import { commandRuntimeLayer } from "../../../../shared/runtime/command-runtime.layer.ts"; import { CliArgs } from "../../../../shared/cli/cli-args.service.ts"; import { textCliOutputFormatter } from "../../../../shared/output/text-formatter.ts"; @@ -112,6 +113,7 @@ function makeDatabaseStack(sqlPort: number, credentials: StackCredentials): Stac }, stop: Effect.die("unused"), destroy: Effect.die("unused"), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, } satisfies Stack; } diff --git a/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts b/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts index eb5e59d608..507942e573 100644 --- a/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/start/start.integration.test.ts @@ -41,6 +41,7 @@ import { mockRuntimeInfo, mockTty, } from "../../../../../tests/helpers/mocks.ts"; +import { unusedGateway } from "../../../../../tests/helpers/unused-stack.ts"; import { containerEngineSpawner } from "../../../../../tests/helpers/child-process-spawner.ts"; import { DbConnection, @@ -371,6 +372,7 @@ const fakeStack = (compositionStart?: Stack["composition"]["start"]) => { hostDestroyed += 1; return { runtimeCleanup: "complete" as const }; }), + gateway: unusedGateway, commands: { run: () => Effect.die("command not used") }, }; return { diff --git a/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts b/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts index 3d2584df86..c41c585456 100644 --- a/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts +++ b/apps/cli/src/commands/experimental/stack/status/status.integration.test.ts @@ -9,6 +9,7 @@ import { StackError, } from "@supabase/stack/effect"; import { mockOutput } from "../../../../../tests/helpers/mocks.ts"; +import { unusedGateway } from "../../../../../tests/helpers/unused-stack.ts"; import { mockCommandSettings, mockTelemetryStateTracked, @@ -176,6 +177,7 @@ const makeStack = ( }, stop: Effect.die("unused"), destroy: Effect.die("unused"), + gateway: unusedGateway, commands: { run: (_tool, _options) => Effect.die("unused"), }, diff --git a/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts b/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts index 10cc300593..0dc3a08250 100644 --- a/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts +++ b/apps/cli/src/commands/functions/serve/serve.stack.integration.test.ts @@ -26,6 +26,7 @@ import { import { StackApi } from "../../../command-internal/stack-api.ts"; import { OutputFlag } from "../../../command-internal/global-flags.ts"; import { mockCommandSettings } from "../../../../tests/helpers/command-mocks.ts"; +import { unusedGateway } from "../../../../tests/helpers/unused-stack.ts"; import { mockOutput, mockProcessControl, @@ -287,6 +288,7 @@ const fixture = ( }, stop: Effect.die("unused"), destroy: Effect.die("unused"), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, } satisfies Stack; const identity = { projectRoot: "/project", branchContext: "main", stackName: "default" }; diff --git a/apps/cli/tests/helpers/storage.ts b/apps/cli/tests/helpers/storage.ts index 95cb7d9854..33430edc67 100644 --- a/apps/cli/tests/helpers/storage.ts +++ b/apps/cli/tests/helpers/storage.ts @@ -18,6 +18,7 @@ import { CliArgs } from "../../src/shared/cli/cli-args.service.ts"; import { CommandPlatformApi } from "../../src/auth/command-platform-api.service.ts"; import { CommandPlatformApiFactory } from "../../src/auth/command-platform-api-factory.service.ts"; import { ProjectRefNotLinkedError } from "../../src/config/project-ref.errors.ts"; +import { unusedGateway } from "./unused-stack.ts"; import { ProjectRefResolver } from "../../src/config/project-ref.service.ts"; import { generateGoJwt } from "../../src/command-internal/go-jwt.ts"; import { YesFlag } from "../../src/command-internal/global-flags.ts"; @@ -263,6 +264,7 @@ export function buildStorageStackApi( }, stop: Effect.die("unused"), destroy: Effect.die("unused"), + gateway: unusedGateway, commands: { run: () => Effect.die("unused") }, }; const definition = { diff --git a/apps/cli/tests/helpers/unused-stack.ts b/apps/cli/tests/helpers/unused-stack.ts index ea1039f398..381dc84a8a 100644 --- a/apps/cli/tests/helpers/unused-stack.ts +++ b/apps/cli/tests/helpers/unused-stack.ts @@ -1,4 +1,5 @@ -import { Effect, Layer } from "effect"; +import { Effect, Layer, Stream } from "effect"; +import type { Stack } from "@supabase/stack/effect"; import { StackApi } from "../../src/command-internal/stack-api.ts"; import { StackCatalogSetup } from "../../src/command-internal/stack-catalog-setup.ts"; @@ -13,3 +14,6 @@ export const unusedStackServices = Layer.mergeAll( }), Layer.succeed(StackCatalogSetup, { apply: unused }), ); + +/** Fills a fake `Stack`'s gateway log stream for tests that never read it. */ +export const unusedGateway: Stack["gateway"] = { readLogs: () => Stream.die("unused") }; diff --git a/packages/stack/ARCHITECTURE.md b/packages/stack/ARCHITECTURE.md index fc411ae733..5cc956aa0f 100644 --- a/packages/stack/ARCHITECTURE.md +++ b/packages/stack/ARCHITECTURE.md @@ -670,12 +670,24 @@ The owner is the only subscriber of each instance's output and persists it as re of instances no longer saved, so a failed deletion is retried. Destroying the stack removes `logs/`; resetting database data keeps them. +#### Gateway access logs + +The shared API listener is the stack's gateway. Each completed request, and each WebSocket upgrade +once its handshake status is known, becomes one nginx combined line plus duration, with credential +query and fragment values (API keys, tokens, token hashes, PKCE codes), values that nest one or a +URL with userinfo, and URL userinfo redacted in the target and the Referer. The proxy hands it to +a bounded sliding buffer after the response settles, so logging never delays a response or holds +a target's activity. The owner persists these lines as the `gateway` stream, +`logs/gateway/gateway/`, with one launch per owner run; it is not a service instance, so orphan +cleanup keeps it. + #### Shipping logs to Analytics While the composed Analytics instance is running and healthy, the owner ships the persisted stdout/stderr records of the Auth, REST, Realtime, Storage, Functions and database instances of the -composition to its direct backend, never the proxy, so shipping neither wakes it nor counts as -activity; standalone instances such as shadow databases are not shipped. The owner logs +composition, and of the gateway stream as Studio's API Gateway source, to its direct backend, never +the proxy, so shipping neither wakes it nor counts as activity; standalone instances such as shadow +databases are not shipped. The owner logs each target change, and the target is re-selected when the composition, Analytics' health or its launch changes. diff --git a/packages/stack/src/HttpProxy.integration.test.ts b/packages/stack/src/HttpProxy.integration.test.ts index 8e6398d412..2d95d4cb1a 100644 --- a/packages/stack/src/HttpProxy.integration.test.ts +++ b/packages/stack/src/HttpProxy.integration.test.ts @@ -1,6 +1,6 @@ import { NodeHttpClient, NodeServices } from "@effect/platform-node"; import { expect, it } from "@effect/vitest"; -import { Data, Deferred, Effect, Fiber, Layer } from "effect"; +import { Data, Deferred, Effect, Fiber, Layer, Queue } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { createServer, type Server, type ServerResponse } from "node:http"; // oxlint-disable-line effecttsgo/node-builtin-import -- raw server fixture. import { createServer as createTcpServer, Socket, type Server as NetServer } from "node:net"; // oxlint-disable-line effecttsgo/node-builtin-import -- raw disconnect fixture. @@ -8,7 +8,7 @@ import { createServer as createTcpServer, Socket, type Server as NetServer } fro import { WebSocket, WebSocketServer } from "ws"; import { captureLogs } from "../tests/logs.ts"; import { ProxyError } from "./Proxy.ts"; -import { makeHttpProxy, type HttpRoute } from "./HttpProxy.ts"; +import { makeHttpProxy, type HttpAccess, type HttpRoute } from "./HttpProxy.ts"; const listen = (server: Server | NetServer, options?: { readonly beforeClose?: () => void }) => Effect.acquireRelease( @@ -1027,3 +1027,366 @@ it.live("releases a waiting WebSocket target quietly when its client resets", () }), ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); }); + +/** Opens a raw client connection that has written `text`. */ +const rawClient = (port: number, text: string) => + Effect.acquireRelease( + Effect.callback((resume) => { + const socket = new Socket(); + const onConnectError = (cause: Error) => { + socket.destroy(); + resume(Effect.fail(new HttpProxyTestError({ message: cause.message, cause }))); + }; + socket.once("error", onConnectError); + socket.connect(port, "127.0.0.1", () => { + socket.off("error", onConnectError); + // Tests reset these connections on purpose. + socket.on("error", () => undefined); + socket.write(text); + resume(Effect.succeed(socket)); + }); + return Effect.sync(() => socket.destroy()); + }), + (socket) => Effect.sync(() => socket.destroy()), + ); + +const upgradeRequest = (path: string) => + `GET ${path} HTTP/1.1\r\nHost: localhost\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n`; + +/** Sends a raw upgrade request and waits until the proxy closes the connection. */ +const rawUpgrade = (port: number, path: string) => + rawClient(port, upgradeRequest(path)).pipe( + Effect.flatMap((socket) => + Effect.callback((resume) => { + socket.once("close", () => resume(Effect.void)); + // Drains any answer so an upstream's graceful end reaches the client as a close. + socket.resume(); + if (socket.destroyed) resume(Effect.void); + }), + ), + Effect.timeout("5 seconds"), + ); + +it.live( + "records each request once with its sent status, body bytes and redacted credentials", + () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const backend = createServer((request, response) => { + request.resume(); + if (request.url?.startsWith("/ok")) response.end("hello"); + else { + response.statusCode = 404; + response.end("nope"); + } + }); + const backendAddress = yield* listen(backend); + const accesses = yield* Queue.unbounded(); + const proxy = yield* makeHttpProxy({ + host: "127.0.0.1", + port: 0, + onAccess: (access) => Queue.offer(accesses, access), + }); + yield* proxy.setRoutes([ + { + id: "api", + prefix: "/api", + upstreamPrefix: "/", + target: Effect.succeed(backendAddress), + }, + { + id: "down", + prefix: "/down", + target: Effect.fail(new ProxyError({ message: "wake failed" })), + }, + ]); + const get = (path: string, headers: Readonly> = {}) => + request(proxy.port, path, new Uint8Array(), headers, "GET").pipe( + Effect.map(({ status }) => status), + ); + + const statuses = [ + yield* get("/api/ok?select=*&apikey=sb_secret_x", { + "user-agent": "proxy-test/1", + referer: "http://127.0.0.1:54321/x?token=t1&select=*", + }), + yield* get("/api/missing"), + yield* get("/elsewhere"), + yield* get("/down/thing"), + ]; + const recorded = yield* Queue.takeN(accesses, 4); + + expect(statuses).toEqual([200, 404, 404, 502]); + expect(recorded.toSorted((left, right) => left.time - right.time)).toEqual([ + { + time: expect.any(Number), + client: "127.0.0.1", + method: "GET", + target: "/api/ok?select=*&apikey=redacted", + protocol: "HTTP/1.1", + status: 200, + bytes: 5, + referer: "http://127.0.0.1:54321/x?token=redacted&select=*", + userAgent: "proxy-test/1", + durationMillis: expect.any(Number), + }, + expect.objectContaining({ target: "/api/missing", status: 404, bytes: 4 }), + expect.objectContaining({ target: "/elsewhere", status: 404, bytes: 9 }), + expect.objectContaining({ target: "/down/thing", status: 502, bytes: 11 }), + ]); + + // A later sentinel request proves none of the four requests above recorded twice. + const sentinelStatus = yield* get("/api/ok?select=sentinel"); + const sentinel = yield* Queue.take(accesses); + expect(sentinelStatus).toBe(200); + expect(sentinel).toMatchObject({ target: "/api/ok?select=sentinel" }); + expect(yield* Queue.size(accesses)).toBe(0); + }), + ).pipe( + Effect.provide( + Layer.mergeAll(NodeHttpClient.layerNodeHttp, NodeServices.layer, captureErrors(logs)), + ), + ); + }, +); + +it.live("records WebSocket upgrades at the handshake with the status sent to the client", () => { + const logs: Array = []; + return Effect.scoped( + Effect.gen(function* () { + const backend = createServer(); + const sockets = new WebSocketServer({ server: backend }); + const backendAddress = yield* listen(backend, { beforeClose: () => sockets.close() }); + const upgradeReceived = yield* Deferred.make(); + const silent = createTcpServer((connection) => { + connection.on("error", () => undefined); + connection.once("data", () => Deferred.doneUnsafe(upgradeReceived, Effect.void)); + }); + const silentAddress = yield* listen(silent); + // Answers with an interim 1xx before the final handshake status, in separate writes. + const interim = createTcpServer((connection) => { + connection.on("error", () => undefined); + connection.once("data", (chunk: Buffer) => { + const upgrading = chunk.toString("latin1").startsWith("GET /interim/upgrade "); + connection.write( + upgrading + ? "HTTP/1.1 103 Early Hints\r\nLink: ; rel=preload\r\n\r\n" + : "HTTP/1.1 100 Continue\r\n\r\n", + ); + connection.end( + upgrading + ? "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n" + : "HTTP/1.1 403 Forbidden\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", + ); + }); + }); + const interimAddress = yield* listen(interim); + const accesses = yield* Queue.unbounded(); + const proxy = yield* makeHttpProxy({ + host: "127.0.0.1", + port: 0, + onAccess: (access) => Queue.offer(accesses, access), + }); + yield* proxy.setRoutes([ + { id: "ws", prefix: "/socket", target: Effect.succeed(backendAddress) }, + { id: "silent", prefix: "/silent", target: Effect.succeed(silentAddress) }, + { id: "interim", prefix: "/interim", target: Effect.succeed(interimAddress) }, + { + id: "down", + prefix: "/down", + target: Effect.fail(new ProxyError({ message: "wake failed" })), + }, + ]); + const client = yield* Effect.acquireRelease( + Effect.callback((resume) => { + const socket = new WebSocket( + `ws://127.0.0.1:${proxy.port}/socket?apikey=sb_publishable_x&vsn=2.0.0`, + ); + socket.once("open", () => resume(Effect.succeed(socket))); + socket.once("error", (cause) => + resume(Effect.fail(new HttpProxyTestError({ message: cause.message, cause }))), + ); + }), + (socket) => Effect.sync(() => socket.terminate()), + ).pipe(Effect.timeout("5 seconds")); + + const opened = yield* Queue.take(accesses); + + expect(client.readyState).toBe(WebSocket.OPEN); + expect(opened).toMatchObject({ + method: "GET", + target: "/socket?apikey=redacted&vsn=2.0.0", + status: 101, + }); + expect(opened.bytes).toBeUndefined(); + + yield* rawUpgrade(proxy.port, "/elsewhere"); + expect(yield* Queue.take(accesses)).toMatchObject({ target: "/elsewhere", status: 404 }); + yield* rawUpgrade(proxy.port, "/down"); + expect(yield* Queue.take(accesses)).toMatchObject({ target: "/down", status: 502 }); + yield* rawUpgrade(proxy.port, "/interim/forbidden"); + expect(yield* Queue.take(accesses)).toMatchObject({ + target: "/interim/forbidden", + status: 403, + }); + yield* rawUpgrade(proxy.port, "/interim/upgrade"); + expect(yield* Queue.take(accesses)).toMatchObject({ + target: "/interim/upgrade", + status: 101, + }); + + const leaving = yield* rawClient(proxy.port, upgradeRequest("/silent")); + yield* Deferred.await(upgradeReceived); + leaving.destroy(); + expect(yield* Queue.take(accesses)).toMatchObject({ target: "/silent", status: 499 }); + }), + ).pipe(Effect.provide(Layer.merge(NodeServices.layer, captureErrors(logs)))); +}); + +it.live("records no body bytes for a client that left before its response completed", () => + Effect.scoped( + Effect.gen(function* () { + const backend = createServer((request, response) => { + request.resume(); + response.writeHead(200); + response.write("first chunk"); + }); + const backendAddress = yield* listen(backend, { + beforeClose: () => backend.closeAllConnections(), + }); + const acquiring = yield* Deferred.make(); + const accesses = yield* Queue.unbounded(); + const proxy = yield* makeHttpProxy({ + host: "127.0.0.1", + port: 0, + onAccess: (access) => Queue.offer(accesses, access), + }); + yield* proxy.setRoutes([ + { + id: "waking", + prefix: "/waking", + target: Deferred.succeed(acquiring, undefined).pipe(Effect.andThen(Effect.never)), + }, + { id: "streaming", prefix: "/streaming", target: Effect.succeed(backendAddress) }, + ]); + + const waiting = yield* rawClient( + proxy.port, + "GET /waking HTTP/1.1\r\nHost: localhost\r\n\r\n", + ); + yield* Deferred.await(acquiring); + waiting.destroy(); + const beforeHeaders = yield* Queue.take(accesses).pipe(Effect.timeout("5 seconds")); + const reading = yield* rawClient( + proxy.port, + "GET /streaming HTTP/1.1\r\nHost: localhost\r\n\r\n", + ); + yield* Effect.callback((resume) => { + reading.once("data", () => resume(Effect.void)); + }); + reading.destroy(); + const midBody = yield* Queue.take(accesses).pipe(Effect.timeout("5 seconds")); + + expect(beforeHeaders).toMatchObject({ target: "/waking", status: 499 }); + expect(beforeHeaders.bytes).toBeUndefined(); + expect(midBody).toMatchObject({ target: "/streaming", status: 200 }); + expect(midBody.bytes).toBeUndefined(); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); + +it.live("answers and releases targets while the access sink is stalled", () => + Effect.scoped( + Effect.gen(function* () { + const backend = createServer((request, response) => { + request.resume(); + response.end("ok"); + }); + const backendAddress = yield* listen(backend); + const released = yield* Queue.unbounded(); + const proxy = yield* makeHttpProxy({ + host: "127.0.0.1", + port: 0, + onAccess: () => Effect.never, + }); + yield* proxy.setRoutes([ + { + id: "api", + prefix: "/api", + target: Effect.acquireRelease(Effect.succeed(backendAddress), () => + Queue.offer(released, undefined), + ), + }, + ]); + + const statuses = yield* Effect.forEach( + Array.from({ length: 20 }, (_, index) => `/api/${index}`), + (path) => + request(proxy.port, path, new Uint8Array(), {}, "GET").pipe( + Effect.map(({ status }) => status), + ), + { concurrency: 5 }, + ).pipe(Effect.timeout("10 seconds")); + + expect(statuses).toEqual(Array.from({ length: 20 }, () => 200)); + expect(yield* Queue.takeN(released, 20).pipe(Effect.timeout("5 seconds"))).toHaveLength(20); + }), + ).pipe(Effect.provide(Layer.merge(NodeHttpClient.layerNodeHttp, NodeServices.layer))), +); + +it.live("records and releases a request whose client resets after the full body was sent", () => + Effect.scoped( + Effect.gen(function* () { + const bodySent = yield* Deferred.make(); + const body = new Uint8Array(256 * 1024).fill(65); + const backend = createServer((request, response) => { + request.resume(); + response.once("finish", () => Deferred.doneUnsafe(bodySent, Effect.void)); + response.end(body); + }); + const backendAddress = yield* listen(backend, { + beforeClose: () => backend.closeAllConnections(), + }); + const released = yield* Deferred.make(); + const accesses = yield* Queue.unbounded(); + const proxy = yield* makeHttpProxy({ + host: "127.0.0.1", + port: 0, + onAccess: (access) => Queue.offer(accesses, access), + }); + yield* proxy.setRoutes([ + { + id: "api", + prefix: "/api", + target: Effect.acquireRelease(Effect.succeed(backendAddress), () => + Deferred.succeed(released, undefined), + ), + }, + ]); + + // The client stops reading after its first bytes; socket buffers decide how much of the + // body the proxy flushed before the reset, so the record holds the whole body or no count. + const client = yield* rawClient( + proxy.port, + "GET /api/full HTTP/1.1\r\nHost: localhost\r\n\r\n", + ); + const firstBytes = Effect.callback((resume) => { + client.once("data", () => { + client.pause(); + resume(Effect.void); + }); + }); + yield* Effect.all([firstBytes, Deferred.await(bodySent)], { concurrency: "unbounded" }).pipe( + Effect.timeout("5 seconds"), + ); + client.resetAndDestroy(); + const recorded = yield* Queue.take(accesses).pipe(Effect.timeout("5 seconds")); + yield* Deferred.await(released).pipe(Effect.timeout("5 seconds")); + + expect(recorded).toMatchObject({ target: "/api/full", status: 200 }); + expect([undefined, body.length]).toContain(recorded.bytes); + expect(yield* Queue.size(accesses)).toBe(0); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/packages/stack/src/HttpProxy.ts b/packages/stack/src/HttpProxy.ts index d2ae57d119..c899dedd43 100644 --- a/packages/stack/src/HttpProxy.ts +++ b/packages/stack/src/HttpProxy.ts @@ -1,4 +1,5 @@ -import { Data, Effect, FiberSet, Ref, Scope } from "effect"; +import { Clock, Data, Deferred, Effect, Fiber, FiberSet, Ref, Scope } from "effect"; +import { decodeQuery, redactCredentials } from "./internal/redact-credentials.ts"; import { PortError } from "./Ports.ts"; import type { BackendAddress, ProxyError } from "./Proxy.ts"; import { @@ -41,6 +42,33 @@ interface HttpRouteKeyRewrite { }; } +/** One request or WebSocket upgrade the proxy completed. */ +export interface HttpAccess { + /** Epoch milliseconds when the request arrived. */ + readonly time: number; + readonly client: string; + readonly method: string; + /** The request path and query, with credential query and fragment values redacted. */ + readonly target: string; + readonly protocol: string; + /** + * The recorded outcome: the status sent, or 499 when the client left before a response. An + * upgrade records the upstream handshake status; without one, 404 for no route and 502 for a + * failure, while the client sees its socket reset. + */ + readonly status: number; + /** Body bytes of a response that finished; absent when it was cut short, and for upgrades. */ + readonly bytes?: number; + /** The Referer header, with credential query and fragment values redacted. */ + readonly referer?: string; + readonly userAgent?: string; + /** Until the response finished; for an upgrade, until the upstream handshake answered. */ + readonly durationMillis: number; +} + +/** Receives each completed request after its response; it must not block. */ +export type HttpAccessSink = (access: HttpAccess) => Effect.Effect; + export interface HttpProxy { readonly host: string; readonly port: number; @@ -117,14 +145,53 @@ const upstreamHeadersFor = (headers: IncomingMessage["headers"], route: HttpRout return result; }; -const decodeQuery = (value: string) => { - try { - return decodeURIComponent(value.replace(/\+/gu, " ")); - } catch { - return value; - } +/** Captures a request's access fields while its socket is open; completes them once it settles. */ +const accessFor = (request: IncomingMessage, time: number) => { + const referer = headerValue(request.headers.referer); + const userAgent = headerValue(request.headers["user-agent"]); + const fields = { + time, + client: request.socket.remoteAddress ?? "-", + method: request.method ?? "GET", + target: redactCredentials(request.url ?? "/"), + protocol: `HTTP/${request.httpVersion}`, + ...(referer === undefined ? {} : { referer: redactCredentials(referer) }), + ...(userAgent === undefined ? {} : { userAgent }), + }; + return (ended: number, status: number, bytes?: number): HttpAccess => ({ + ...fields, + status, + ...(bytes === undefined ? {} : { bytes }), + durationMillis: Math.max(0, ended - time), + }); }; +/** What one request's client was sent. */ +interface Sent { + /** Body bytes handed to the response. */ + bytes: number; + /** The client went away before the response completed. */ + clientLeft: boolean; +} + +const respond = (response: ServerResponse, sent: Sent, status: number, body?: string) => { + response.statusCode = status; + sent.bytes = body === undefined ? 0 : Buffer.byteLength(body); + response.end(body); +}; + +const responseSettled = (response: ServerResponse) => + Effect.callback((resume) => { + const onDone = () => resume(Effect.void); + response.once("finish", onDone); + response.once("close", onDone); + if (response.writableFinished || response.destroyed) onDone(); + return Effect.sync(() => { + response.off("finish", onDone); + response.off("close", onDone); + }); + }); + const queryValueFor = (value: string, keys: HttpRouteKeyRewrite["keys"]) => value === keys.secretKey ? keys.serviceRoleKey @@ -247,6 +314,7 @@ const forward = Effect.fn("HttpProxy.forward")( route: HttpRoute, backend: BackendAddress, agent: Agent | false, + sent: Sent, ) => Effect.callback((resume) => { let outgoing: ReturnType | undefined; @@ -278,8 +346,9 @@ const forward = Effect.fn("HttpProxy.forward")( abandon(Effect.fail(errorFor(cause, incoming !== undefined))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); const onFinish = () => finish(Effect.void); + // A close after `end()` but before `finish` means the client reset with writes still queued. const onResponseClose = () => { - if (!response.writableEnded) onClientGone(); + if (!response.writableFinished) onClientGone(); }; outgoing = upstreamRequest( { @@ -296,6 +365,9 @@ const forward = Effect.fn("HttpProxy.forward")( incoming = value; value.once("aborted", onError); value.on("error", onError); + value.on("data", (chunk: Buffer) => { + sent.bytes += chunk.length; + }); response.once("finish", onFinish); setCors(response, request); response.statusCode = value.statusCode ?? 502; @@ -323,23 +395,56 @@ const forward = Effect.fn("HttpProxy.forward")( ); const proxyRequest = Effect.fn("HttpProxy.proxyRequest")( - (request: IncomingMessage, response: ServerResponse, route: HttpRoute, agent: Agent) => + ( + request: IncomingMessage, + response: ServerResponse, + route: HttpRoute, + agent: Agent, + sent: Sent, + ) => Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ route_id: route.id }); const backend = yield* Effect.raceFirst(route.target, disconnected(request, response)); - yield* forward(request, response, route, backend, isReplayable(request) ? agent : false).pipe( + yield* forward( + request, + response, + route, + backend, + isReplayable(request) ? agent : false, + sent, + ).pipe( Effect.catchIf(isRetryable(request, response), (error) => Effect.logWarning( `Route ${route.id} ${request.method ?? "GET"} upstream failed before responding, retrying`, error, - ).pipe(Effect.andThen(forward(request, response, route, backend, false))), + ).pipe(Effect.andThen(forward(request, response, route, backend, false, sent))), ), ); - }), + }).pipe( + Effect.ensuring( + Effect.suspend(() => + response.headersSent + ? Effect.annotateCurrentSpan({ "http.response.status_code": response.statusCode }) + : Effect.void, + ), + ), + ), ); +const statusLine = /^HTTP\/\d(?:\.\d)? (\d{3})\b/u; +/** Upstream bytes read for the handshake status before giving up on finding one. */ +const answerLimit = 8 * 1024; + const upgrade = Effect.fn("HttpProxy.upgrade")( - (request: IncomingMessage, client: Duplex, head: Buffer, route: HttpRoute) => + ( + request: IncomingMessage, + client: Duplex, + head: Buffer, + route: HttpRoute, + handshake: Deferred.Deferred, + ) => Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ route_id: route.id }); const backend = yield* Effect.raceFirst( route.target, Effect.callback((resume) => { @@ -352,12 +457,14 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( const upstream = yield* connectInterruptibly(backend); yield* Effect.callback((resume) => { let settled = false; + let answer = ""; const cleanup = () => { - client.off("close", onClose); + client.off("close", onClientClose); upstream.off("close", onClose); client.off("end", onClientEnd); upstream.off("end", onUpstreamEnd); + upstream.off("data", onAnswer); }; const finish = (result: Effect.Effect) => { if (settled) return; @@ -373,10 +480,33 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( const onError = (cause: Error) => abandon(Effect.fail(errorFor(cause))); const onClientGone = () => abandon(Effect.fail(new HttpProxyDisconnected())); const onClose = () => abandon(Effect.void); + // A client leaving before the upstream answered never received a response. + const onClientClose = () => (Deferred.isDoneUnsafe(handshake) ? onClose() : onClientGone()); const onClientEnd = () => upstream.end(); const onUpstreamEnd = () => client.end(); + // Reads the final handshake status off the bytes relayed to the client, skipping interim + // 1xx responses such as 100 Continue (RFC 9110 section 15.2). + const onAnswer = (chunk: Buffer) => { + answer += chunk.toString("latin1"); + for ( + let headEnd = answer.indexOf("\r\n\r\n"); + headEnd >= 0; + headEnd = answer.indexOf("\r\n\r\n") + ) { + const status = Number(statusLine.exec(answer)?.[1] ?? Number.NaN); + if (status >= 100 && status < 200 && status !== 101) { + answer = answer.slice(headEnd + 4); + continue; + } + upstream.off("data", onAnswer); + if (!Number.isNaN(status)) Deferred.doneUnsafe(handshake, Effect.succeed(status)); + return; + } + if (answer.length >= answerLimit) upstream.off("data", onAnswer); + }; + upstream.on("data", onAnswer); client.on("error", onClientGone); - client.once("close", onClose); + client.once("close", onClientClose); upstream.on("error", onError); upstream.once("close", onClose); client.once("end", onClientEnd); @@ -398,12 +528,25 @@ const upgrade = Effect.fn("HttpProxy.upgrade")( upstream.destroy(); }); }); - }), + }).pipe( + Effect.ensuring( + Effect.suspend(() => + Deferred.isDoneUnsafe(handshake) + ? Deferred.await(handshake).pipe( + Effect.flatMap((status) => + Effect.annotateCurrentSpan({ "http.response.status_code": status }), + ), + ) + : Effect.void, + ), + ), + ), ); export const makeHttpProxy = (options: { readonly host: string; readonly port: number; + readonly onAccess?: HttpAccessSink | undefined; }): Effect.Effect => Effect.gen(function* () { const routes = yield* Ref.make>([]); @@ -415,42 +558,65 @@ export const makeHttpProxy = (options: { ); const runRequest = yield* FiberSet.makeRuntime(); const sockets = new Set(); + const onAccess = options.onAccess; const server = createServer((request, response) => { runRequest( - Effect.scoped( - Effect.gen(function* () { - const route = yield* Ref.get(routes).pipe( - Effect.map((current) => routeFor(request.url ?? "/", current)), - ); - setCors(response, request); - if (request.method === "OPTIONS") { - response.statusCode = 204; - response.end(); - } else if (route === undefined) { - response.statusCode = 404; - response.end("Not Found"); - } else { - yield* proxyRequest(request, response, route, agent).pipe( - Effect.tapError((cause) => - cause._tag === "HttpProxyDisconnected" - ? Effect.void - : Effect.logError(`Route ${route.id} request failed`, cause), - ), - Effect.catch(() => - Effect.sync(() => { - if (response.destroyed) return; - if (response.headersSent) response.destroy(); - else { - setCors(response, request); - response.statusCode = 502; - response.end("Bad Gateway"); - } - }), - ), + Effect.gen(function* () { + const sent: Sent = { bytes: 0, clientLeft: false }; + const access = + onAccess === undefined + ? undefined + : { + complete: accessFor(request, yield* Clock.currentTimeMillis), + // Timed when the response settles, before the target's release runs. + settled: yield* responseSettled(response).pipe( + Effect.andThen(Clock.currentTimeMillis), + Effect.forkChild({ startImmediately: true }), + ), + record: onAccess, + }; + // The access record waits outside this scope, so it never holds the target's activity. + yield* Effect.scoped( + Effect.gen(function* () { + const route = yield* Ref.get(routes).pipe( + Effect.map((current) => routeFor(request.url ?? "/", current)), ); - } - }), - ), + setCors(response, request); + if (request.method === "OPTIONS") respond(response, sent, 204); + else if (route === undefined) respond(response, sent, 404, "Not Found"); + else { + yield* proxyRequest(request, response, route, agent, sent).pipe( + Effect.tapError((cause) => + cause._tag === "HttpProxyDisconnected" + ? Effect.void + : Effect.logError(`Route ${route.id} request failed`, cause), + ), + Effect.catch((cause) => + Effect.sync(() => { + if (cause._tag === "HttpProxyDisconnected") sent.clientLeft = true; + if (response.destroyed) return; + if (response.headersSent || sent.clientLeft) response.destroy(); + else { + setCors(response, request); + respond(response, sent, 502, "Bad Gateway"); + } + }), + ), + ); + } + }), + ); + if (access === undefined) return; + const ended = yield* Fiber.join(access.settled); + const delivered = response.writableFinished && !sent.clientLeft; + yield* access.record( + access.complete( + ended, + response.headersSent ? response.statusCode : 499, + delivered ? sent.bytes : undefined, + ), + ); + }), ); }); server.on("connection", (socket) => { @@ -460,23 +626,49 @@ export const makeHttpProxy = (options: { server.on("upgrade", (request, socket, head) => { socket.on("error", () => socket.destroy()); runRequest( - Effect.scoped( - Effect.gen(function* () { - const route = yield* Ref.get(routes).pipe( - Effect.map((current) => routeFor(request.url ?? "/", current)), - ); - if (route === undefined) socket.destroy(); - else - yield* upgrade(request, socket, head, route).pipe( - Effect.tapError((cause) => - cause._tag === "HttpProxyDisconnected" - ? Effect.void - : Effect.logError(`Route ${route.id} upgrade failed`, cause), - ), - Effect.catch(() => Effect.sync(() => socket.destroy())), - ); - }), - ), + Effect.gen(function* () { + // Completed by the upstream's handshake status, otherwise by the upgrade's outcome. + const handshake = yield* Deferred.make(); + const recorded = + onAccess === undefined + ? undefined + : yield* Effect.gen(function* () { + const complete = accessFor(request, yield* Clock.currentTimeMillis); + return yield* Deferred.await(handshake).pipe( + Effect.flatMap((status) => + Clock.currentTimeMillis.pipe( + Effect.flatMap((ended) => onAccess(complete(ended, status))), + ), + ), + Effect.forkChild({ startImmediately: true }), + ); + }); + const route = yield* Ref.get(routes).pipe( + Effect.map((current) => routeFor(request.url ?? "/", current)), + ); + const outcome = + route === undefined + ? yield* Effect.sync(() => { + socket.destroy(); + return 404; + }) + : yield* Effect.scoped(upgrade(request, socket, head, route, handshake)).pipe( + Effect.as(502), + Effect.tapError((cause) => + cause._tag === "HttpProxyDisconnected" + ? Effect.void + : Effect.logError(`Route ${route.id} upgrade failed`, cause), + ), + Effect.catch((cause) => + Effect.sync(() => { + socket.destroy(); + return cause._tag === "HttpProxyDisconnected" ? 499 : 502; + }), + ), + ); + yield* Deferred.succeed(handshake, outcome); + if (recorded !== undefined) yield* Fiber.join(recorded); + }), ); }); yield* Effect.acquireRelease( diff --git a/packages/stack/src/Network.ts b/packages/stack/src/Network.ts index bb80763509..d5b2ca44ad 100644 --- a/packages/stack/src/Network.ts +++ b/packages/stack/src/Network.ts @@ -3,7 +3,7 @@ import { DOCKER_HOST_ALIAS } from "./runtime/Container.ts"; import * as State from "./State.ts"; import { makePorts, PortError } from "./Ports.ts"; import { bindTcp, serveTcp, type BackendAddress, type ProxyError } from "./Proxy.ts"; -import { makeHttpProxy, type HttpProxy, type HttpRoute } from "./HttpProxy.ts"; +import { makeHttpProxy, type HttpAccessSink, type HttpProxy, type HttpRoute } from "./HttpProxy.ts"; export type NetworkRuntime = "native" | "docker" | "podman"; @@ -70,6 +70,7 @@ const makeNetwork = (options: { readonly stackId: string; readonly runtime: NetworkRuntime; readonly state: State.Interface; + readonly onAccess?: HttpAccessSink; }) => Effect.gen(function* () { const ports = yield* makePorts(options.state).pipe( @@ -174,7 +175,13 @@ const makeNetwork = (options: { }); return { proxy: current.proxy }; } - return { proxy: yield* makeHttpProxy({ host, port }) }; + return { + proxy: yield* makeHttpProxy({ + host, + port, + onAccess: options.onAccess, + }), + }; } if (endpoint.protocol === "http") { const proxy = yield* makeHttpProxy({ host, port }); @@ -341,7 +348,12 @@ const makeNetwork = (options: { return { register, release: release() } satisfies Interface; }); -export const layer = (options: { readonly stackId: string; readonly runtime: NetworkRuntime }) => +export const layer = (options: { + readonly stackId: string; + readonly runtime: NetworkRuntime; + /** Receives the shared API listener's completed requests. */ + readonly onAccess?: HttpAccessSink; +}) => Layer.effect( Service, Effect.gen(function* () { diff --git a/packages/stack/src/Owner.analytics.integration.test.ts b/packages/stack/src/Owner.analytics.integration.test.ts index 578c863a83..20e44eeb87 100644 --- a/packages/stack/src/Owner.analytics.integration.test.ts +++ b/packages/stack/src/Owner.analytics.integration.test.ts @@ -261,6 +261,16 @@ it.live( timestamp: expect.any(String), }, ]); + expect( + yield* awaitStored( + "cloudflare.logs.prod", + `body->'metadata'->'request'->>'path' LIKE '%${restPath}' AND body->'metadata'->'response'->>'status_code' = '404'`, + ), + ).toEqual([ + expect.objectContaining({ + message: expect.stringContaining(`${restPath} HTTP/1.1" 404 `), + }), + ]); yield* Fiber.interrupt(keeper); const noise = yield* Effect.forkScoped( diff --git a/packages/stack/src/Owner.logs.integration.test.ts b/packages/stack/src/Owner.logs.integration.test.ts index 172f21b136..53c219aa2b 100644 --- a/packages/stack/src/Owner.logs.integration.test.ts +++ b/packages/stack/src/Owner.logs.integration.test.ts @@ -13,6 +13,7 @@ import { Scope, Stream, } from "effect"; +import { HttpClient } from "effect/unstable/http"; import { tmpdir } from "node:os"; import { ownerFor } from "../tests/owner-rpc.ts"; import type { LogRecord } from "./host/LogRecord.ts"; @@ -181,6 +182,61 @@ describe("owner persisted logs", () => { ).pipe(Effect.provide(services)), ); + it.live("keeps shared API requests as gateway logs across owner restarts until destroy", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const client = yield* HttpClient.HttpClient; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "owner-logs-gateway-" }); + const firstScope = yield* Scope.make(); + // Closes the first owner if the test fails before the restart closes it. + yield* Effect.addFinalizer(() => Scope.close(firstScope, Exit.void)); + const first = yield* openOwner("owner-logs-gateway-", "native", { root }).pipe( + Scope.provide(firstScope), + ); + const rest = yield* first.rpc.createService({ + service: "rest", + config: {}, + endpoints: { http: { port: "auto" } }, + }); + yield* first.rpc.configureComposition({ + members: [{ id: rest.id, activation: "lazy" }], + dependencies: [], + }); + yield* first.rpc.startComposition(); + const api = (yield* first.state.read(first.stack.id))?.ports.find( + ({ key }) => key === "api", + ); + if (api === undefined) return yield* Effect.die("the shared API port was not claimed"); + + const response = yield* client.get(`http://127.0.0.1:${api.port}/unrouted`); + const [recorded] = yield* firstOutput(first.rpc.readLogs({ id: "gateway", follow: true })); + + expect(response.status).toBe(404); + expect(recorded?.text).toContain("/unrouted"); + const directory = path.join(first.logsRoot, "gateway", "gateway"); + expect(yield* fs.readDirectory(directory)).toEqual(["0000000001.log"]); + yield* Scope.close(firstScope, Exit.void); + + const state = Context.get( + yield* Layer.build(State.layer({ root: first.stateRoot })), + State.Service, + ); + const saved = yield* state.read(first.stack.id); + if (saved === undefined) return yield* Effect.die("the stopped stack was not saved"); + const second = yield* ownerFor({ saved, state, root: `${root}/data`, cacheRoot }); + const history = Array.from( + yield* second.rpc.readLogs({ id: "gateway", follow: false }).pipe(Stream.runCollect), + ); + + expect(history).toContainEqual(recorded); + yield* second.namespace.destroy; + expect(yield* fs.exists(first.logsRoot)).toBe(false); + }), + ).pipe(Effect.provide(services)), + ); + it.live("continues after the launch ids in its logs when its saved state has none", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/Owner.ts b/packages/stack/src/Owner.ts index d8549c3956..6f131ac708 100644 --- a/packages/stack/src/Owner.ts +++ b/packages/stack/src/Owner.ts @@ -67,6 +67,7 @@ import { stackError, type OwnerRpc } from "./Rpc.ts"; import * as State from "./State.ts"; import type { SavedStack, StackCredentials, StackKeysInput } from "./State.ts"; import { makeDockerHelperRegistry } from "./storage/DockerHelperRegistry.ts"; +import * as GatewayLog from "./host/GatewayLog.ts"; import * as LogForwarder from "./host/LogForwarder.ts"; import * as LogflareStorage from "./host/LogflareStorage.ts"; import * as LogStore from "./host/LogStore.ts"; @@ -172,7 +173,10 @@ const withoutInstance = (current: SavedStack, id: string): SavedStack => const drainingBlocks: ReadonlyArray = ["start", "arm", "restart", "storage"]; -const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { +const makeOwner = Effect.fn("Owner.make")(function* ( + options: OwnerOptions, + gateway: GatewayLog.GatewayLog, +) { const services = yield* Effect.context< | FileSystem.FileSystem | Path.Path @@ -216,6 +220,15 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { logs: logStore, storedEvents: (analytics) => LogflareStorage.make(analyticsDatabase(analytics.id)), }).pipe(Effect.provideContext(services)); + yield* logStore.attach({ + ...GatewayLog.gatewayLog, + logs: gateway.logs, + observation: gateway.observation, + }); + yield* forwarder.attach({ + id: GatewayLog.gatewayLog.instanceId, + service: GatewayLog.gatewayLog.service, + }); const definitionGate = yield* Semaphore.make(1); const draining = yield* Ref.make(false); const { id: stackId, runtime } = options.saved; @@ -799,17 +812,24 @@ const makeOwner = Effect.fn("Owner.make")(function* (options: OwnerOptions) { }); export const layer = (options: Omit) => - Layer.effect( - Service, - Effect.gen(function* () { - const state = yield* State.Service; - return Service.of(yield* makeOwner({ ...options, state })); - }), - ).pipe( - Layer.provide( - Network.layer({ - stackId: options.saved.id, - runtime: options.saved.runtime, - }), + Layer.unwrap( + GatewayLog.make.pipe( + Effect.map((gateway) => + Layer.effect( + Service, + Effect.gen(function* () { + const state = yield* State.Service; + return Service.of(yield* makeOwner({ ...options, state }, gateway)); + }), + ).pipe( + Layer.provide( + Network.layer({ + stackId: options.saved.id, + runtime: options.saved.runtime, + onAccess: gateway.record, + }), + ), + ), + ), ), ); diff --git a/packages/stack/src/effect.ts b/packages/stack/src/effect.ts index 958f569abf..ba0cf8ec16 100644 --- a/packages/stack/src/effect.ts +++ b/packages/stack/src/effect.ts @@ -46,6 +46,7 @@ import * as State from "./State.ts"; import type { SavedStack, StackCredentials, StackKeysInput } from "./State.ts"; import { StackError, type Definition, type Observation } from "./Rpc.ts"; import { sinceMillis, streamStackLogs as streamPersistedLogs } from "./host/LogStore.ts"; +import { gatewayLog } from "./host/GatewayLog.ts"; import type { LogPosition, LogRecord, StackLogRecord } from "./host/LogRecord.ts"; import { reclaimStack } from "./Sweep.ts"; import { @@ -65,6 +66,7 @@ import type { export { initialization, postgres } from "./Commands.ts"; export { resolveNativePostgresUser } from "./runtime/postgres-user.ts"; export { apiRoute } from "./host/Endpoints.ts"; +export { gatewayLog }; export { StackError } from "./Rpc.ts"; export type { ServiceCreation } from "./services/Catalog.ts"; /** A service creation as `services.create` accepts it, before stack credentials fill its inputs. */ @@ -275,6 +277,10 @@ export interface Stack { readonly stop: Effect.Effect, StackError>; readonly restart: Effect.Effect, StackError>; }; + /** The shared API listener's access log, read through the live owner like an instance's. */ + readonly gateway: { + readonly readLogs: (options?: ReadLogsOptions) => Stream.Stream; + }; readonly stop: Effect.Effect; readonly destroy: Effect.Effect; readonly commands: { @@ -636,6 +642,18 @@ const makeHandle = Effect.fn("Stack.makeHandle")(function* ( const snapshotScope = (options: DatabaseSnapshotOptions | undefined) => options?.scope === undefined ? {} : { scope: options.scope }; + const readLogs = + (id: string) => + (options?: ReadLogsOptions): Stream.Stream => + stream("readLogs", (rpc) => + rpc.readLogs({ + id, + follow: options?.follow ?? false, + ...(options?.from === undefined ? {} : { from: options.from }), + ...(options?.since === undefined ? {} : { since: options.since }), + ...(options?.tail === undefined ? {} : { tail: options.tail }), + }), + ); const common = (id: string, service: K): ServiceInstance => ({ id, service, @@ -657,16 +675,7 @@ const makeHandle = Effect.fn("Stack.makeHandle")(function* ( prepare: call("prepare", (rpc) => rpc.prepareService({ id })), status: call("status", (rpc) => rpc.status({ id }), "attach"), followStatus: stream("followStatus", (rpc) => rpc.followStatus({ id })), - readLogs: (options) => - stream("readLogs", (rpc) => - rpc.readLogs({ - id, - follow: options?.follow ?? false, - ...(options?.from === undefined ? {} : { from: options.from }), - ...(options?.since === undefined ? {} : { since: options.since }), - ...(options?.tail === undefined ? {} : { tail: options.tail }), - }), - ), + readLogs: readLogs(id), credentials: (options) => call("credentials", (rpc) => rpc.credentials({ id, from: options?.from ?? "host" })), }); @@ -911,6 +920,7 @@ const makeHandle = Effect.fn("Stack.makeHandle")(function* ( stop: whileRunning("stopComposition", (rpc) => rpc.stopComposition(), []), restart: call("restartComposition", (rpc) => rpc.restartComposition()), }, + gateway: { readLogs: readLogs(gatewayLog.instanceId) }, stop: shutdown(false).pipe(Effect.asVoid), destroy: shutdown(true), commands: { run }, diff --git a/packages/stack/src/host/GatewayLog.ts b/packages/stack/src/host/GatewayLog.ts new file mode 100644 index 0000000000..b833be862a --- /dev/null +++ b/packages/stack/src/host/GatewayLog.ts @@ -0,0 +1,50 @@ +import { DateTime, Effect, PubSub, Stream } from "effect"; +import type { HttpAccess, HttpAccessSink } from "../HttpProxy.ts"; +import { launchOutputPublisher, type LaunchOutput } from "../runtime/Session.ts"; +import type { CatalogLogs } from "../services/Recipe.ts"; +import { monthNames } from "./LogflareEvents.ts"; + +/** The log stream of the shared API listener's access records, one per owner and stack. */ +export const gatewayLog = { service: "gateway", instanceId: "gateway" } as const; + +/** Bounds records not yet persisted; the oldest are dropped and reported as lost. */ +const bufferedRecords = 4096; + +const two = (value: number) => String(value).padStart(2, "0"); + +/** Formats epoch milliseconds as nginx's `$time_local` in UTC. */ +const nginxTime = (millis: number) => { + const parts = DateTime.toPartsUtc(DateTime.makeUnsafe(millis)); + return `${two(parts.day)}/${monthNames[parts.month - 1]}/${parts.year}:${two(parts.hour)}:${two(parts.minute)}:${two(parts.second)} +0000`; +}; + +/** Escapes quotes, backslashes and control characters like nginx's default log escaping. */ +const escapeLogValue = (value: string) => + value.replace( + // oxlint-disable-next-line no-control-regex -- control characters are what this escapes. + /["\\\u0000-\u001f\u007f]/gu, + (character) => `\\x${character.charCodeAt(0).toString(16).padStart(2, "0")}`, + ); + +/** Formats an access record as an nginx combined log line followed by its duration. */ +export const formatAccess = (access: HttpAccess) => + `${access.client} - - [${nginxTime(access.time)}] "${escapeLogValue(`${access.method} ${access.target} ${access.protocol}`)}" ${access.status} ${access.bytes ?? "-"} "${escapeLogValue(access.referer ?? "-")}" "${escapeLogValue(access.userAgent ?? "-")}" ${access.durationMillis}ms`; + +const encoder = new TextEncoder(); + +/** The gateway stream's output, its access sink, and an observation of one launch per owner run. */ +export interface GatewayLog { + readonly logs: CatalogLogs; + readonly record: HttpAccessSink; + readonly observation: Stream.Stream<{ readonly launchId: number }>; +} + +export const make = Effect.gen(function* () { + const output = yield* PubSub.sliding(bufferedRecords); + const publish = yield* (yield* launchOutputPublisher(output, 1)).part; + return { + logs: PubSub.subscribe(output), + record: (access) => publish("stdout", encoder.encode(`${formatAccess(access)}\n`)), + observation: Stream.make({ launchId: 1 }).pipe(Stream.concat(Stream.never)), + } satisfies GatewayLog; +}); diff --git a/packages/stack/src/host/LogForwarder.integration.test.ts b/packages/stack/src/host/LogForwarder.integration.test.ts index 4c7d74b6da..757c321eb1 100644 --- a/packages/stack/src/host/LogForwarder.integration.test.ts +++ b/packages/stack/src/host/LogForwarder.integration.test.ts @@ -29,6 +29,7 @@ import { makeManualClock, type ManualClock } from "../../tests/manual-clock.ts"; import type { ServiceObservation } from "../Service.ts"; import type { LaunchOutput } from "../runtime/Session.ts"; import { CatalogError } from "../services/Recipe.ts"; +import * as GatewayLog from "./GatewayLog.ts"; import type { LogRecord } from "./LogRecord.ts"; import * as LogForwarder from "./LogForwarder.ts"; import * as LogStore from "./LogStore.ts"; @@ -190,7 +191,7 @@ const composition = Effect.succeed({ const startForwarder = ( store: LogStore.Interface, logflare: FakeLogflare, - instances: ReadonlyArray, + instances: ReadonlyArray, clock?: Clock.Clock, stored?: LogForwarder.StoredEvents, ) => @@ -1158,4 +1159,40 @@ describe("LogForwarder", () => { expect(shipped).toEqual(retained.map((record) => record.text)); }).pipe(Effect.scoped, Effect.provide(layer)), ); + + it.live("ships gateway lines to cloudflare.logs.prod", () => + Effect.gen(function* () { + const { store, logflare, analytics } = yield* fixture({ logflare: {} }); + const gateway = yield* GatewayLog.make; + yield* store.attach({ + ...GatewayLog.gatewayLog, + logs: gateway.logs, + observation: gateway.observation, + }); + yield* analytics.set(true); + yield* startForwarder(store, logflare, [ + analytics.instance, + { id: GatewayLog.gatewayLog.instanceId, service: GatewayLog.gatewayLog.service }, + ]); + + yield* gateway.record({ + time: Date.parse("2026-10-01T09:25:23.000Z"), + client: "127.0.0.1", + method: "POST", + target: "/auth/v1/token?grant_type=password", + protocol: "HTTP/1.1", + status: 400, + bytes: 60, + durationMillis: 3, + }); + const shipped = yield* logflare.next; + yield* logflare.apply(); + + expect(shipped.url).toBe("/api/logs?source_name=cloudflare.logs.prod"); + expect(shipped.events).toEqual([expect.objectContaining({ appname: "gateway" })]); + expect(yield* logflare.storedIds("cloudflare.logs.prod", ids(shipped))).toEqual( + new Set(ids(shipped)), + ); + }).pipe(Effect.scoped, Effect.provide(layer)), + ); }); diff --git a/packages/stack/src/host/LogForwarder.ts b/packages/stack/src/host/LogForwarder.ts index 7fe0c513e8..395d11a85d 100644 --- a/packages/stack/src/host/LogForwarder.ts +++ b/packages/stack/src/host/LogForwarder.ts @@ -32,6 +32,7 @@ import { type LogflareEvent, type ShippedService, } from "./LogflareEvents.ts"; +import type { gatewayLog } from "./GatewayLog.ts"; import { LogPosition, type LogRecord } from "./LogRecord.ts"; import type * as LogStore from "./LogStore.ts"; @@ -44,6 +45,12 @@ export interface ForwardedInstance { readonly observation: Stream.Stream>; } +/** An owner log stream without a service instance; it ships while the owner runs. */ +export interface ForwardedStream { + readonly id: string; + readonly service: typeof gatewayLog.service; +} + export class StoredEventsError extends Schema.TaggedError()( "StoredEventsError", { message: Schema.String, cause: Schema.optionalKey(Schema.Defect()) }, @@ -59,8 +66,11 @@ export interface StoredEvents { } interface Interface { - /** Ships a shipped service's persisted logs, or tracks an Analytics instance as the target. */ - readonly attach: (instance: ForwardedInstance) => Effect.Effect; + /** + * Ships the persisted logs of a shipped service or an owner stream, or tracks an Analytics + * instance as the target. + */ + readonly attach: (instance: ForwardedInstance | ForwardedStream) => Effect.Effect; /** Stops following an instance; returns once no cursor write of it can still land. */ readonly detach: (instanceId: string) => Effect.Effect; /** Re-selects the shipping target after the composition changes. */ @@ -859,21 +869,30 @@ export const make = Effect.fn("LogForwarder.make")(function* (options: LogForwar Stream.runDrain, ); - /** Ships until the instance unregisters or the store detaches its logs, which ends a session. */ - const forward = (instance: ForwardedInstance, service: ShippedService) => + /** + * Ships until the instance unregisters (`until`) or the store detaches its logs, which ends a + * session. A service instance ships only while it is a composition member; the gateway stream + * belongs to the owner. + */ + const forward = ( + id: string, + service: ShippedService, + until: Effect.Effect, + composed: boolean, + ) => Effect.gen(function* () { while (true) { - yield* membership(instance.id, true); + if (composed) yield* membership(id, true); const current = yield* serving; const failing = yield* Ref.make(false); const posting = yield* Semaphore.make(1); // A stopped or refused target pauses shipping until it changes; a failed log read resumes // from the cursor after a backoff. - const detached = yield* session(instance.id, service, current, failing, posting).pipe( + const detached = yield* session(id, service, current, failing, posting).pipe( Effect.tapError((error) => error._tag === "StaleTarget" ? Effect.void - : warnOnce(failing, `Reading ${instance.id} logs to ship failed; retrying`, error), + : warnOnce(failing, `Reading ${id} logs to ship failed; retrying`, error), ), Effect.retry({ schedule: retrySchedule, @@ -884,19 +903,20 @@ export const make = Effect.fn("LogForwarder.make")(function* (options: LogForwar // A retarget or leaving the composition lets a started post finish, so its answer decides // the pending body. Effect.raceFirst( - Effect.raceFirst(retargeted(current), membership(instance.id, false)).pipe( - Effect.andThen(posting.take(1)), - Effect.as(false), - ), + Effect.raceFirst( + retargeted(current), + composed ? membership(id, false) : Effect.never, + ).pipe(Effect.andThen(posting.take(1)), Effect.as(false)), ), ); - if (detached) - return yield* Effect.logDebug(`Log shipping of ${instance.id} stopped with its logs`); + if (detached) return yield* Effect.logDebug(`Log shipping of ${id} stopped with its logs`); } - }).pipe(Effect.raceFirst(unregistered(instance))); + }).pipe(Effect.raceFirst(until)); const followers = yield* FiberMap.make(); - const attach = Effect.fn("LogForwarder.attach")(function* (instance: ForwardedInstance) { + const attach = Effect.fn("LogForwarder.attach")(function* ( + instance: ForwardedInstance | ForwardedStream, + ) { if (instance.service === "analytics") yield* FiberMap.run(followers, instance.id, trackTarget(instance)); else if (isShippedService(instance.service)) { @@ -905,7 +925,16 @@ export const make = Effect.fn("LogForwarder.make")(function* (options: LogForwar Effect.flatMap((directory) => reapStaleWrites(fs, path, directory)), Effect.ignore, ); - yield* FiberMap.run(followers, instance.id, forward(instance, instance.service)); + yield* FiberMap.run( + followers, + instance.id, + forward( + instance.id, + instance.service, + "observation" in instance ? unregistered(instance) : Effect.never, + "observation" in instance, + ), + ); } }); diff --git a/packages/stack/src/host/LogflareEvents.ts b/packages/stack/src/host/LogflareEvents.ts index 1b507f1df7..3b4fedcf72 100644 --- a/packages/stack/src/host/LogflareEvents.ts +++ b/packages/stack/src/host/LogflareEvents.ts @@ -1,7 +1,11 @@ import { DateTime, Option, Predicate } from "effect"; import type { ServiceCreation } from "../services/Catalog.ts"; +import type { gatewayLog } from "./GatewayLog.ts"; -/** Logflare source per shipped service kind; Studio's Logs pages query these names. */ +/** A service kind, or the owner's `gateway` stream of shared API listener requests. */ +type LogService = ServiceCreation["service"] | typeof gatewayLog.service; + +/** Logflare source per shipped log service; Studio's Logs pages query these names. */ export const logflareSources = { auth: "gotrue.logs.prod", rest: "postgREST.logs.prod", @@ -9,11 +13,12 @@ export const logflareSources = { storage: "storage.logs.prod.2", functions: "deno-relay-logs", database: "postgres.logs", -} as const satisfies Partial>; + gateway: "cloudflare.logs.prod", +} as const satisfies Partial>; export type ShippedService = keyof typeof logflareSources; -export const isShippedService = (service: ServiceCreation["service"]): service is ShippedService => +export const isShippedService = (service: LogService): service is ShippedService => Object.hasOwn(logflareSources, service); /** One ingest event in the shape Studio's local log queries expect. */ @@ -35,29 +40,30 @@ const parseJsonObject = (text: string): Record | undefined => { } }; -const months: Record = { - jan: 0, - feb: 1, - mar: 2, - apr: 3, - may: 4, - jun: 5, - jul: 6, - aug: 7, - sep: 8, - oct: 9, - nov: 10, - dec: 11, -}; +/** The `%b` month abbreviations of PostgREST and nginx log times, January first. */ +export const monthNames = [ + "Jan", + "Feb", + "Mar", + "Apr", + "May", + "Jun", + "Jul", + "Aug", + "Sep", + "Oct", + "Nov", + "Dec", +] as const; -/** Parses PostgREST's `%d/%b/%Y:%H:%M:%S %z` prefix into an ISO-8601 timestamp. */ -const parsePostgrestTime = (text: string): string | undefined => { +/** Parses the `%d/%b/%Y:%H:%M:%S %z` time of PostgREST and nginx logs into ISO-8601. */ +const parseLogTime = (text: string): string | undefined => { const match = /^(\d{2})\/([A-Za-z]{3})\/(\d{4}):(\d{2}):(\d{2}):(\d{2}) ([+-])(\d{2})(\d{2})$/u.exec(text); if (match === null) return undefined; const [, day, month, year, hour, minute, second, sign, zoneHours, zoneMinutes] = match; - const monthIndex = months[month?.toLowerCase() ?? ""]; - if (monthIndex === undefined) return undefined; + const monthIndex = monthNames.findIndex((name) => name.toLowerCase() === month?.toLowerCase()); + if (monthIndex < 0) return undefined; const utc = Date.UTC( Number(year), monthIndex, @@ -75,6 +81,19 @@ const parsePostgrestTime = (text: string): string | undefined => { /** The time and request of PostgREST's Apache combined request line. */ const postgrestRequest = /^\S+ \S+ \S+ \[([^\]]+)\] "([A-Z]+) (\S+) ([^"\s]+)" (\d{3}) /u; +/** The gateway's nginx combined line with its trailing duration. */ +const gatewayRequest = + /^(\S+) \S+ \S+ \[([^\]]+)\] "(\S+) (\S+) (\S+)" (\d{3}) (?:\d+|-) "([^"]*)" "([^"]*)" \d+ms$/u; + +const unescapeLogValue = (value: string) => + value.replace(/\\x([0-9a-f]{2})/giu, (_, hex: string) => + String.fromCharCode(Number.parseInt(hex, 16)), + ); + +/** A quoted combined-log field, absent when nginx wrote `-`. */ +const logValue = (value: string | undefined) => + value === undefined || value === "-" ? undefined : unescapeLogValue(value); + const withoutProject = ({ project: _project, ...event }: LogflareEvent): LogflareEvent => event; const remaps: Record LogflareEvent> = { @@ -86,7 +105,7 @@ const remaps: Record LogflareEvent> = }, rest: (event) => { const request = postgrestRequest.exec(event.event_message); - const requestTime = request === null ? undefined : parsePostgrestTime(request[1] ?? ""); + const requestTime = request === null ? undefined : parseLogTime(request[1] ?? ""); if (request !== null && requestTime !== undefined) return { ...event, @@ -101,7 +120,7 @@ const remaps: Record LogflareEvent> = }, }; const match = /^(.*?): (.*)$/u.exec(event.event_message); - const timestamp = match === null ? undefined : parsePostgrestTime(match[1] ?? ""); + const timestamp = match === null ? undefined : parseLogTime(match[1] ?? ""); return match === null || timestamp === undefined ? event : { @@ -151,6 +170,35 @@ const remaps: Record LogflareEvent> = }, }; }, + gateway: (event) => { + const match = gatewayRequest.exec(event.event_message); + const timestamp = match === null ? undefined : parseLogTime(match[2] ?? ""); + if (match === null || timestamp === undefined) return event; + const [, client, , method, target = "", protocol, status, referer, userAgent] = match; + const decoded = unescapeLogValue(target); + const queryAt = decoded.indexOf("?"); + const refererHeader = logValue(referer); + const userAgentHeader = logValue(userAgent); + return { + ...event, + timestamp, + metadata: { + ...event.metadata, + request: { + method, + path: queryAt < 0 ? decoded : decoded.slice(0, queryAt), + ...(queryAt < 0 ? {} : { search: decoded.slice(queryAt) }), + protocol, + headers: { + cf_connecting_ip: client, + ...(refererHeader === undefined ? {} : { referer: refererHeader }), + ...(userAgentHeader === undefined ? {} : { user_agent: userAgentHeader }), + }, + }, + response: { status_code: Number(status) }, + }, + }; + }, }; /** diff --git a/packages/stack/src/host/LogflareEvents.unit.test.ts b/packages/stack/src/host/LogflareEvents.unit.test.ts index ab6e8b9c2f..b77c0a0415 100644 --- a/packages/stack/src/host/LogflareEvents.unit.test.ts +++ b/packages/stack/src/host/LogflareEvents.unit.test.ts @@ -1,4 +1,5 @@ import { expect, it } from "@effect/vitest"; +import { formatAccess } from "./GatewayLog.ts"; import { logflareEvent } from "./LogflareEvents.ts"; const received = "2026-09-28T10:00:00.000Z"; @@ -123,6 +124,50 @@ it("derives the Postgres severity from the last level marker and defaults to LOG }); }); +it("ships a gateway access line as an API Gateway request stamped with its request time", () => { + const line = formatAccess({ + time: Date.parse("2026-10-01T09:25:23.456Z"), + client: "127.0.0.1", + method: "GET", + target: '/rest/v1/todos?select=*&q="x"', + protocol: "HTTP/1.1", + status: 200, + bytes: 126, + userAgent: "curl/8.7.1", + durationMillis: 12, + }); + + expect(line).toBe( + '127.0.0.1 - - [01/Oct/2026:09:25:23 +0000] "GET /rest/v1/todos?select=*&q=\\x22x\\x22 HTTP/1.1" 200 126 "-" "curl/8.7.1" 12ms', + ); + expect(logflareEvent("gateway", received, line)).toEqual({ + project: "default", + appname: "gateway", + event_message: line, + timestamp: "2026-10-01T09:25:23.000Z", + metadata: { + request: { + method: "GET", + path: "/rest/v1/todos", + search: '?select=*&q="x"', + protocol: "HTTP/1.1", + headers: { cf_connecting_ip: "127.0.0.1", user_agent: "curl/8.7.1" }, + }, + response: { status_code: 200 }, + }, + }); +}); + +it("passes a malformed gateway line through without request metadata", () => { + expect(logflareEvent("gateway", received, "not an access line")).toEqual({ + project: "default", + appname: "gateway", + event_message: "not an access line", + timestamp: received, + metadata: {}, + }); +}); + it("replaces NUL and unpaired surrogates, which Postgres jsonb rejects, in the message and metadata", () => { const line = String.raw`{"msg":"a\u0000b","detail":["\ud800"],"k\u0000":"\udc00 x"}`; const event = logflareEvent("auth", received, line); diff --git a/packages/stack/src/internal/redact-credentials.ts b/packages/stack/src/internal/redact-credentials.ts new file mode 100644 index 0000000000..4d9f257925 --- /dev/null +++ b/packages/stack/src/internal/redact-credentials.ts @@ -0,0 +1,77 @@ +/** Decodes a query component, keeping it as is when it is not valid percent-encoding. */ +export const decodeQuery = (value: string) => { + try { + return decodeURIComponent(value.replace(/\+/gu, " ")); + } catch { + return value; + } +}; + +/** Decodes repeatedly so a nested URL encoded more than once still exposes its delimiters. */ +const decodeNested = (value: string) => { + let current = value; + for (let pass = 0; pass < 4; pass++) { + const next = decodeQuery(current); + if (next === current) break; + current = next; + } + return current; +}; + +// Auth puts PKCE codes, OTP token hashes and OAuth tokens in redirect and verify URLs; S3 +// presigned URLs carry their SigV4 signature, credential scope and session token. +const credentialParameters = [ + "apikey", + "jwt", + "access_token", + "token", + "token_hash", + "code", + "refresh_token", + "id_token", + "provider_token", + "provider_refresh_token", + "x-amz-signature", + "x-amz-credential", + "x-amz-security-token", +]; +const scheme = String.raw`[a-z][a-z\d+.-]*`; + +/** + * A decoded parameter that is a credential pair, nests one (a `redirect_to` URL carrying a + * token), or nests a URL with userinfo. + */ +const sensitive = new RegExp( + String.raw`(?:^|[?&#=])(?:${credentialParameters.join("|")})=|${scheme}:\/\/[^/?#@]*@`, + "iu", +); +const userinfo = new RegExp(String.raw`^(${scheme}:\/\/)[^/?#@]*@`, "iu"); + +const redactPairs = (pairs: string) => + pairs + .split("&") + .map((parameter) => { + const separator = parameter.indexOf("="); + return separator >= 0 && sensitive.test(decodeNested(parameter)) + ? `${parameter.slice(0, separator)}=redacted` + : parameter; + }) + .join("&"); + +/** + * Redacts an absolute URL's userinfo and credential values in its query and fragment, where + * OAuth implicit grants put `access_token` for non-browser clients. Edits the text in place, so + * relative and unparsable URLs work and the rest of the URL keeps its original encoding. + */ +export const redactCredentials = (url: string) => { + const hashAt = url.indexOf("#"); + const beforeHash = hashAt < 0 ? url : url.slice(0, hashAt); + const queryAt = beforeHash.indexOf("?"); + const base = (queryAt < 0 ? beforeHash : beforeHash.slice(0, queryAt)).replace( + userinfo, + "$1redacted@", + ); + const query = queryAt < 0 ? "" : `?${redactPairs(beforeHash.slice(queryAt + 1))}`; + const fragment = hashAt < 0 ? "" : `#${redactPairs(url.slice(hashAt + 1))}`; + return `${base}${query}${fragment}`; +}; diff --git a/packages/stack/src/internal/redact-credentials.unit.test.ts b/packages/stack/src/internal/redact-credentials.unit.test.ts new file mode 100644 index 0000000000..93d44c6a49 --- /dev/null +++ b/packages/stack/src/internal/redact-credentials.unit.test.ts @@ -0,0 +1,82 @@ +import { describe, expect, it } from "@effect/vitest"; +import { redactCredentials } from "./redact-credentials.ts"; + +describe("redactCredentials", () => { + it.each([ + { + case: "redacts credential query values, any case", + url: "/api/ok?select=*&apikey=sb_secret_x&Access_Token=jwt", + redacted: "/api/ok?select=*&apikey=redacted&Access_Token=redacted", + }, + { + case: "redacts an Edge Function websocket JWT", + url: "/functions/v1/realtime-chat?jwt=eyJ.p.s", + redacted: "/functions/v1/realtime-chat?jwt=redacted", + }, + { + case: "redacts Auth codes and OTP token hashes", + url: "/auth/v1/verify?code=pkce&token_hash=h&type=signup", + redacted: "/auth/v1/verify?code=redacted&token_hash=redacted&type=signup", + }, + { + case: "redacts S3 presigned URL signatures and session tokens", + url: "/storage/v1/s3/b/o.txt?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=AK%2F20261001%2Flocal%2Fs3%2Faws4_request&X-Amz-Expires=60&X-Amz-Security-Token=st&X-Amz-Signature=abc", + redacted: + "/storage/v1/s3/b/o.txt?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=redacted&X-Amz-Expires=60&X-Amz-Security-Token=redacted&X-Amz-Signature=redacted", + }, + { + case: "redacts fragment credentials of an absolute URL", + url: "http://127.0.0.1:54321/x?token=t1&select=*#access_token=frag&type=bearer", + redacted: + "http://127.0.0.1:54321/x?token=redacted&select=*#access_token=redacted&type=bearer", + }, + { + case: "redacts refresh and provider tokens in a fragment", + url: "http://127.0.0.1:54321/y?apikey=k#refresh_token=r&provider_token=p&provider_refresh_token=q", + redacted: + "http://127.0.0.1:54321/y?apikey=redacted#refresh_token=redacted&provider_token=redacted&provider_refresh_token=redacted", + }, + { + case: "redacts an encoded nested URL carrying a token", + url: "/elsewhere?redirect_to=https%3A%2F%2Fclient%2Fcb%3Faccess_token%3DJWT", + redacted: "/elsewhere?redirect_to=redacted", + }, + { + case: "redacts a raw nested URL carrying a token", + url: "/down?redirect_to=https://client/cb?access_token=JWT", + redacted: "/down?redirect_to=redacted", + }, + { + case: "redacts a doubly encoded nested URL carrying a token", + url: "/down?back=https%253A%252F%252Fclient%252Fcb%253Faccess_token%253DJWT", + redacted: "/down?back=redacted", + }, + { + case: "redacts a nested URL with userinfo", + url: "/elsewhere?return_to=https%3A%2F%2Fuser%3Asecret%40client%2Fcb", + redacted: "/elsewhere?return_to=redacted", + }, + { + case: "redacts userinfo of an absolute URL", + url: "https://user:password@studio.test/relative?Access_Token=jwt", + redacted: "https://redacted@studio.test/relative?Access_Token=redacted", + }, + { + case: "redacts a value that embeds a credential pair", + url: "/api?data=a=token=b", + redacted: "/api?data=redacted", + }, + { + case: "keeps a nested URL without credentials, with its encoding", + url: "/elsewhere?next=http%3A%2F%2Flocalhost%3A3000%2F&country_code=FR", + redacted: "/elsewhere?next=http%3A%2F%2Flocalhost%3A3000%2F&country_code=FR", + }, + { + case: "keeps a path without a query", + url: "/rest/v1/todos", + redacted: "/rest/v1/todos", + }, + ])("$case", ({ url, redacted }) => { + expect(redactCredentials(url)).toBe(redacted); + }); +});