Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
86 changes: 85 additions & 1 deletion apps/server/src/auth/http.test.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,20 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { EnvironmentHttpApi } from "@t3tools/contracts";
import {
AuthSessionId,
EnvironmentAuthenticatedAuth,
EnvironmentHttpApi,
} from "@t3tools/contracts";
import { RelayClientTracer } from "@t3tools/shared/relayTracing";
import { expect, it } from "@effect/vitest";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Redacted from "effect/Redacted";
import * as Schema from "effect/Schema";
import * as Tracer from "effect/Tracer";
import { HttpServerRequest } from "effect/http";
import * as Etag from "effect/http/Etag";
import * as HttpPlatform from "effect/http/HttpPlatform";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
Expand Down Expand Up @@ -134,3 +142,79 @@ it.effect("sets the selected browser session cookies through the HTTP route", ()
);
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("exports only verified T3 Connect requests", () =>
Effect.gen(function* () {
const productSpans: Array<string> = [];
const localSpans: Array<string> = [];
const collect = (into: Array<string>) =>
Tracer.make({
span: (options) => {
into.push(options.name);
return new Tracer.NativeSpan(options);
},
});
// "DPoP connect" is a T3 Connect session; any other DPoP token is rejected.
const environmentAuth = {
authenticateHttpRequest: (request: HttpServerRequest.HttpServerRequest) =>
(request.headers.authorization === "DPoP forged"
? Effect.fail(
new EnvironmentAuth.ServerAuthInvalidCredentialError({ diagnostic: "forged" }),
)
: Effect.succeed({
sessionId: AuthSessionId.make("session-1"),
subject:
request.headers.authorization === "DPoP connect"
? "cloud-connect"
: "cli-issued-session",
method: "bearer-access-token" as const,
scopes: ["orchestration:read" as const],
})
).pipe(Effect.withSpan("EnvironmentAuth.authenticateHttpRequest")),
} as unknown as EnvironmentAuth.EnvironmentAuth["Service"];
const middleware = yield* Layer.build(AuthHttp.layerAuthenticatedAuth).pipe(
Effect.provideService(EnvironmentAuth.EnvironmentAuth, environmentAuth),
Effect.map(Context.get(EnvironmentAuthenticatedAuth)),
);
const handle = (authorization: string) =>
(
middleware as unknown as (
effect: Effect.Effect<void>,
) => Effect.Effect<void, Error, HttpServerRequest.HttpServerRequest>
)(Effect.void.pipe(Effect.withSpan("environment.handler"))).pipe(
Effect.ignore,
Effect.provideService(
HttpServerRequest.HttpServerRequest,
HttpServerRequest.fromWeb(
new Request("https://environment.example.test/api/orchestration/shell", {
headers: {
authorization,
traceparent: "00-0123456789abcdef0123456789abcdef-0123456789abcdef-01",
},
}),
),
),
Effect.provideService(RelayClientTracer, Option.some(collect(productSpans))),
Effect.withTracer(collect(localSpans)),
);

yield* handle("DPoP connect");
expect(productSpans).toEqual([
"environment.relay.request",
"EnvironmentAuth.authenticateHttpRequest",
"environment.handler",
]);
expect(localSpans).toEqual(["EnvironmentAuth.authenticateHttpRequest"]);

productSpans.length = 0;
localSpans.length = 0;
yield* handle("DPoP forged");
expect(productSpans).toEqual([]);
expect(localSpans).toEqual(["EnvironmentAuth.authenticateHttpRequest"]);

localSpans.length = 0;
yield* handle("Bearer access-token");
expect(productSpans).toEqual([]);
expect(localSpans).toEqual(["EnvironmentAuth.authenticateHttpRequest", "environment.handler"]);
}).pipe(Effect.scoped),
);
10 changes: 7 additions & 3 deletions apps/server/src/auth/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,9 @@ import {
import type { AuthEnvironmentScope, DpopFailureReason } from "@t3tools/contracts";
import { parseAllowedOAuthScope } from "@t3tools/shared/oauthScope";
import { causeErrorTag } from "@t3tools/shared/observability";
import * as Clock from "effect/Clock";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import { identity } from "effect/Function";
import * as Layer from "effect/Layer";
import * as Cookies from "effect/http/Cookies";
import * as HttpEffect from "effect/http/HttpEffect";
Expand Down Expand Up @@ -205,6 +205,7 @@ export const layerAuthenticatedAuth = Layer.effect(
return (httpEffect) =>
Effect.gen(function* () {
const request = yield* HttpServerRequest.HttpServerRequest;
const startTime = yield* Clock.currentTimeNanos;
const session = yield* serverAuth.authenticateHttpRequest(request).pipe(
Effect.catchIf(EnvironmentAuth.isServerAuthCredentialError, (error) =>
failEnvironmentAuthInvalid(
Expand All @@ -216,13 +217,16 @@ export const layerAuthenticatedAuth = Layer.effect(
failEnvironmentInternal("internal_error", error),
),
);
return yield* httpEffect.pipe(
const endTime = yield* Clock.currentTimeNanos;
const handler = httpEffect.pipe(
Effect.provideService(EnvironmentAuthenticatedPrincipal, {
...session,
scopes: new Set(session.scopes),
}),
session.subject === "cloud-connect" ? traceAuthenticatedRelayRequest : identity,
);
return yield* session.subject === "cloud-connect"
? traceAuthenticatedRelayRequest(handler, { startTime, endTime })
: handler;
}).pipe(Effect.catchTags({ EnvironmentAuthInvalidError: appendDpopChallengeOnUnauthorized }));
}),
);
Expand Down
65 changes: 60 additions & 5 deletions apps/server/src/cloud/traceRelayRequest.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,13 @@ import * as Tracer from "effect/Tracer";
import { HttpServerRequest } from "effect/http";
import { RelayClientTracer } from "@t3tools/shared/relayTracing";

import { traceAuthenticatedRelayRequest, traceRelayRequest } from "./traceRelayRequest.ts";
import {
traceAuthenticatedRelayRequest,
traceLocalHandlerWork,
traceRelayRequest,
} from "./traceRelayRequest.ts";

const AUTHENTICATION = { startTime: 1_000n, endTime: 2_000n };

describe("relay request tracing", () => {
it.effect("does not accept an unauthenticated request trace parent", () =>
Expand Down Expand Up @@ -58,15 +64,64 @@ describe("relay request tracing", () => {

yield* traceAuthenticatedRelayRequest(
Effect.void.pipe(Effect.withSpan("relay.mint.handler")),
AUTHENTICATION,
).pipe(
Effect.provideService(HttpServerRequest.HttpServerRequest, request),
Effect.provideService(RelayClientTracer, Option.some(productTracer)),
);

expect(spans).toHaveLength(1);
const span = spans[0]!;
expect(span.traceId).toBe("0123456789abcdef0123456789abcdef");
expect(Option.getOrUndefined(span.parent)?.spanId).toBe("0123456789abcdef");
expect(spans.map((span) => span.name)).toEqual([
"environment.relay.request",
"EnvironmentAuth.authenticateHttpRequest",
"relay.mint.handler",
]);
const [relaySpan, authentication, handler] = spans;
expect(Option.getOrUndefined(authentication!.parent)?.spanId).toBe(relaySpan!.spanId);
expect(relaySpan!.traceId).toBe("0123456789abcdef0123456789abcdef");
expect(Option.getOrUndefined(relaySpan!.parent)?.spanId).toBe("0123456789abcdef");
expect(Option.getOrUndefined(handler!.parent)?.spanId).toBe(relaySpan!.spanId);
}),
);
});

describe("relay request tracing boundary", () => {
it.effect("exports a T3 Connect handler span but not its local work", () =>
Effect.gen(function* () {
const productSpans: Array<string> = [];
const localSpans: Array<string> = [];
const collect = (into: Array<string>) =>
Tracer.make({
span: (options) => {
into.push(options.name);
return new Tracer.NativeSpan(options);
},
});
const request = HttpServerRequest.fromWeb(
new Request("https://environment.example.test/api/orchestration/threads/thread-1"),
);

// Shaped like a real handler: an Effect.fn span whose service call runs
// on the local tracer.
const handler = Effect.fn("environment.orchestration.threadSnapshot")(function* () {
yield* Effect.void.pipe(
Effect.withSpan("sql.execute"),
Effect.withSpan("ServerSecretStore.get"),
traceLocalHandlerWork,
);
});

yield* traceAuthenticatedRelayRequest(handler(), AUTHENTICATION).pipe(
Effect.provideService(HttpServerRequest.HttpServerRequest, request),
Effect.provideService(RelayClientTracer, Option.some(collect(productSpans))),
Effect.withTracer(collect(localSpans)),
);

expect(productSpans).toEqual([
"environment.relay.request",
"EnvironmentAuth.authenticateHttpRequest",
"environment.orchestration.threadSnapshot",
]);
expect(localSpans).toEqual(["ServerSecretStore.get", "sql.execute"]);
}),
);
});
74 changes: 65 additions & 9 deletions apps/server/src/cloud/traceRelayRequest.ts
Original file line number Diff line number Diff line change
@@ -1,21 +1,77 @@
import { withRelayClientTracing } from "@t3tools/shared/relayTracing";
import {
RelayClientTracer,
withLocalTracing,
withRelayClientTracing,
} from "@t3tools/shared/relayTracing";
import * as Clock from "effect/Clock";
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Option from "effect/Option";
import * as Tracer from "effect/Tracer";
import { HttpServerRequest, HttpTraceContext } from "effect/http";

/**
* Exports every span of a handler that is itself T3 Connect work (the token
* exchange, the environment descriptor, credential minting), so the whole
* connection path can be measured.
*/
export const traceRelayRequest = <A, E, R>(
effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, E, R> => effect.pipe(withRelayClientTracing);

/**
* Traces a request that arrived over T3 Connect, continuing the client's
* trace. The request's span and its authentication timing are exported, so
* their latency and errors are visible; what the handler then does on the
* user's machine (database reads, project indexing, processes) is not, unless
* the handler is connection work and opts back in with {@link traceRelayRequest}.
*
* Call it only after the session is verified as T3 Connect: an unverified
* request must not add spans to a trace it names. Authentication therefore
* runs on the local tracer, and its timing is recorded here afterwards.
*/
export const traceAuthenticatedRelayRequest = <A, E, R>(
effect: Effect.Effect<A, E, R>,
authentication: { readonly startTime: bigint; readonly endTime: bigint },
): Effect.Effect<A, E, R | HttpServerRequest.HttpServerRequest> =>
HttpServerRequest.HttpServerRequest.pipe(
Effect.flatMap((request) =>
Option.match(HttpTraceContext.fromHeaders(request.headers), {
onNone: () => effect,
onSome: (parent) => effect.pipe(Effect.withParentSpan(parent)),
}),
),
withRelayClientTracing,
Effect.all([RelayClientTracer, HttpServerRequest.HttpServerRequest]).pipe(
Effect.flatMap(([tracer, request]): Effect.Effect<A, E, R> => {
if (Option.isNone(tracer)) return effect;
const parent = HttpTraceContext.fromHeaders(request.headers);
const sampled = Option.isNone(parent) || parent.value.sampled;
const startSpan = (
name: string,
kind: Tracer.SpanKind,
parent: Option.Option<Tracer.AnySpan>,
) =>
tracer.value.span({
name,
parent,
annotations: Context.empty(),
links: [],
startTime: authentication.startTime,
kind,
root: Option.isNone(parent),
sampled,
});
const span = startSpan("environment.relay.request", "server", parent);
startSpan("EnvironmentAuth.authenticateHttpRequest", "internal", Option.some(span)).end(
authentication.endTime,
Exit.void,
);
return effect.pipe(
Effect.withParentSpan(span),
Effect.onExit((exit) =>
Clock.currentTimeNanos.pipe(Effect.map((endTime) => span.end(endTime, exit))),
),
withRelayClientTracing,
);
}),
);

/**
* Keeps the work inside a relay-traced handler on the local tracer. The
* handler's span still records the timing and outcome.
*/
export const traceLocalHandlerWork = withLocalTracing;
10 changes: 7 additions & 3 deletions apps/server/src/orchestration-v2/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import {
failEnvironmentNotFound,
requireEnvironmentScope,
} from "../auth/http.ts";
import { traceLocalHandlerWork } from "../cloud/traceRelayRequest.ts";
import * as OrchestrationEventStore from "../persistence/OrchestrationEventStore.ts";
import * as ProjectEnrichmentService from "../project/ProjectEnrichmentService.ts";
import {
Expand Down Expand Up @@ -174,6 +175,7 @@ export const layer = HttpApiBuilder.group(
yield* annotateEnvironmentRequest(args.endpoint.name);
yield* requireEnvironmentScope(AuthOrchestrationReadScope);
return yield* loadShellSnapshot().pipe(
traceLocalHandlerWork,
Effect.catch((cause) =>
failEnvironmentInternal("orchestration_snapshot_failed", cause),
),
Expand All @@ -188,7 +190,7 @@ export const layer = HttpApiBuilder.group(
const snapshot = yield* loadThreadSnapshot(
args.params.threadId,
"orchestration_thread_snapshot_failed",
);
).pipe(traceLocalHandlerWork);
return {
snapshotSequence: snapshot.snapshotSequence,
projection: snapshot.projection,
Expand All @@ -200,7 +202,9 @@ export const layer = HttpApiBuilder.group(
Effect.fn("environment.orchestration.threadBoundedSnapshot")(function* (args) {
yield* annotateEnvironmentRequest(args.endpoint.name);
yield* requireEnvironmentScope(AuthOrchestrationReadScope);
const snapshot = yield* loadThreadSnapshotWindow(args.params.threadId);
const snapshot = yield* loadThreadSnapshotWindow(args.params.threadId).pipe(
traceLocalHandlerWork,
);
const bounded = buildBoundedThreadProjection({
projection: snapshot.projection,
snapshotSequence: snapshot.snapshotSequence,
Expand Down Expand Up @@ -234,7 +238,7 @@ export const layer = HttpApiBuilder.group(
args.params.threadId,
anchorItemId,
ThreadId.make(decodedCursor.st),
);
).pipe(traceLocalHandlerWork);
const pageOrError = selectHistoryPageFromCursorOrError({
items: snapshot.projection.visibleTurnItems,
cursor: args.query.cursor,
Expand Down
6 changes: 5 additions & 1 deletion apps/server/src/project/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
failEnvironmentInvalidRequest,
requireEnvironmentScope,
} from "../auth/http.ts";
import { traceLocalHandlerWork } from "../cloud/traceRelayRequest.ts";
import * as ServerRuntimeStartup from "../serverRuntimeStartup.ts";
import * as ProjectService from "./ProjectService.ts";
import { projectMutationOperation } from "./ProjectMutation.ts";
Expand Down Expand Up @@ -43,6 +44,7 @@ export const layer = HttpApiBuilder.group(
yield* annotateEnvironmentRequest(args.endpoint.name);
yield* requireEnvironmentScope(AuthOrchestrationReadScope);
return yield* projects.snapshot.pipe(
traceLocalHandlerWork,
Effect.catch((cause) => failEnvironmentInternal("project_snapshot_failed", cause)),
);
}),
Expand All @@ -53,7 +55,9 @@ export const layer = HttpApiBuilder.group(
yield* annotateEnvironmentRequest(args.endpoint.name);
yield* requireEnvironmentScope(AuthOrchestrationOperateScope);
const operation = projectMutationOperation(projects, args.payload);
return yield* startup.enqueueCommand(operation).pipe(Effect.catch(failProjectMutation));
return yield* startup
.enqueueCommand(operation)
.pipe(traceLocalHandlerWork, Effect.catch(failProjectMutation));
}),
);
}),
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/pullRequest/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import * as Effect from "effect/Effect";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";

import { annotateEnvironmentRequest, requireEnvironmentScope } from "../auth/http.ts";
import { traceLocalHandlerWork } from "../cloud/traceRelayRequest.ts";
import * as PullRequestService from "./PullRequestService.ts";

/** The patch is often the largest PR payload and benefits from HTTP compression and flow control. */
Expand All @@ -16,7 +17,7 @@ export const layer = HttpApiBuilder.group(
Effect.fn("environment.pullRequests.diff")(function* (args) {
yield* annotateEnvironmentRequest(args.endpoint.name);
yield* requireEnvironmentScope(AuthOrchestrationReadScope);
return yield* pullRequests.diff(args.payload);
return yield* pullRequests.diff(args.payload).pipe(traceLocalHandlerWork);
}),
);
}),
Expand Down
Loading
Loading