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
40 changes: 40 additions & 0 deletions apps/server/src/auth/RpcAuthorization.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,15 @@ import {
WsRpcGroup,
} from "@t3tools/contracts";
import { describe, expect, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as RpcTest from "effect/unstable/rpc/RpcTest";

import {
RPC_REQUIRED_SCOPES,
requiredScopeForRpcMethod,
requiredScopeForDeviceList,
rpcScopeAuthorizationLayer,
} from "./RpcAuthorization.ts";

describe("RPC authorization scopes", () => {
Expand Down Expand Up @@ -116,3 +120,39 @@ it("requires operate permission for tool updates even alongside a read-only chec
);
expect(requiredScopeForDeviceList({ updateTool: "hub" })).toBe(AuthOrchestrationOperateScope);
});

describe("RPC scope middleware", () => {
const tested = [WS_METHODS.serverProbe, WS_METHODS.serverRetryResourceTelemetry] as const;
const group = WsRpcGroup.omit(
...[...WsRpcGroup.requests.keys()].filter(
(tag): tag is Exclude<keyof typeof RPC_REQUIRED_SCOPES, (typeof tested)[number]> =>
!(tested as ReadonlyArray<string>).includes(tag),
),
);

it.effect("checks each RPC's declared scope before its handler runs", () =>
Effect.gen(function* () {
const handled: Array<string> = [];
const client = yield* RpcTest.makeClient(group).pipe(
Effect.provide(
Layer.mergeAll(
group.toLayerHandler(WS_METHODS.serverProbe, () => Effect.succeed({})),
group.toLayerHandler(WS_METHODS.serverRetryResourceTelemetry, () =>
Effect.sync(() => handled.push("retry")).pipe(Effect.andThen(Effect.never)),
),
rpcScopeAuthorizationLayer([AuthOrchestrationReadScope]),
),
),
);

expect(yield* client[WS_METHODS.serverProbe]({})).toEqual({});
expect(
yield* client[WS_METHODS.serverRetryResourceTelemetry]({}).pipe(Effect.flip),
).toMatchObject({
_tag: "EnvironmentAuthorizationError",
requiredScope: AuthOrchestrationOperateScope,
});
expect(handled).toEqual([]);
}).pipe(Effect.scoped),
);
});
19 changes: 19 additions & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,13 @@ import {
AuthTerminalOperateScope,
ORCHESTRATION_V2_WS_METHODS,
type AuthEnvironmentScope,
EnvironmentAuthorizationError,
RpcScopeAuthorization,
WS_METHODS,
WsRpcGroup,
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import type * as RpcGroup from "effect/unstable/rpc/RpcGroup";

type WsRpcMethod = RpcGroup.Rpcs<typeof WsRpcGroup>["_tag"];
Expand Down Expand Up @@ -212,6 +216,21 @@ export function requiredScopeForRpcMethod(method: string): AuthEnvironmentScope
return requiredScope;
}

export const rpcAuthorizationError = (requiredScope: AuthEnvironmentScope) =>
new EnvironmentAuthorizationError({
message: `The authenticated token is missing required scope: ${requiredScope}.`,
requiredScope,
});

/** Authorizes every RPC on one connection against that connection's session scopes. */
export const rpcScopeAuthorizationLayer = (scopes: ReadonlyArray<AuthEnvironmentScope>) =>
Layer.succeed(RpcScopeAuthorization)((effect, { rpc }) => {
const requiredScope = requiredScopeForRpcMethod(rpc._tag);
return scopes.includes(requiredScope)
? effect
: Effect.fail(rpcAuthorizationError(requiredScope));
});

/** Retrying can install or restart tools even though ordinary listing is readable. */
export const requiredScopeForDeviceList = (input: DeviceListInput): AuthEnvironmentScope =>
input.retryHostId || input.updateTool
Expand Down
12 changes: 9 additions & 3 deletions apps/server/src/mcp/PreviewAutomationBroker.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { expect, it } from "@effect/vitest";
import {
AuthOrchestrationOperateScope,
EnvironmentId,
PreviewAutomationClientDisconnectedError,
PreviewAutomationInvalidSelectorError,
Expand All @@ -17,6 +18,7 @@ import {
type PreviewAutomationStreamEvent,
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Exit from "effect/Exit";
import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
Expand All @@ -27,6 +29,7 @@ import * as TestClock from "effect/testing/TestClock";
import * as RpcGroup from "effect/unstable/rpc/RpcGroup";
import * as RpcTest from "effect/unstable/rpc/RpcTest";

import { rpcScopeAuthorizationLayer } from "../auth/RpcAuthorization.ts";
import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";

const makeBroker = PreviewAutomationBroker.make.pipe(Effect.provide(NodeServices.layer));
Expand Down Expand Up @@ -1255,9 +1258,12 @@ it.effect("evicts an unanswered host and lets later calls use a healthy runtime"
);
const client = yield* RpcTest.makeClient(group).pipe(
Effect.provide(
group.toLayer({
[WS_METHODS.previewAutomationConnect]: (host) => Stream.unwrap(broker.connect(host)),
}),
Layer.merge(
group.toLayer({
[WS_METHODS.previewAutomationConnect]: (host) => Stream.unwrap(broker.connect(host)),
}),
rpcScopeAuthorizationLayer([AuthOrchestrationOperateScope]),
),
),
);
const events = client[WS_METHODS.previewAutomationConnect](makeHost());
Expand Down
63 changes: 12 additions & 51 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -157,9 +157,9 @@ import * as ThreadSearch from "./orchestration-v2/ThreadSearch.ts";
import * as OrchestrationEventStore from "./persistence/Services/OrchestrationEventStore.ts";
import { userFacingDispatchErrorMessage } from "./orchestration-v2/UserFacingErrors.ts";
import {
observeRpcEffect as instrumentRpcEffect,
observeRpcStream as instrumentRpcStream,
observeRpcStreamEffect as instrumentRpcStreamEffect,
observeRpcEffect,
observeRpcStream,
observeRpcStreamEffect,
} from "./observability/RpcInstrumentation.ts";
import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts";
import * as ProviderInstanceRegistry from "./provider/Services/ProviderInstanceRegistry.ts";
Expand Down Expand Up @@ -207,7 +207,11 @@ import * as ServerEnvironment from "./environment/ServerEnvironment.ts";
import * as RemoteOpenTargets from "./environment/RemoteOpenTargets.ts";
import * as BackgroundPolicy from "./background/BackgroundPolicy.ts";
import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts";
import { requiredScopeForRpcMethod, requiredScopeForDeviceList } from "./auth/RpcAuthorization.ts";
import {
requiredScopeForDeviceList,
rpcAuthorizationError,
rpcScopeAuthorizationLayer,
} from "./auth/RpcAuthorization.ts";
import * as ProcessDiagnostics from "./diagnostics/ProcessDiagnostics.ts";
import * as ProcessResourceMonitor from "./diagnostics/ProcessResourceMonitor.ts";
import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts";
Expand Down Expand Up @@ -1309,25 +1313,15 @@ const makeWsRpcLayer = (
const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor;
const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry;
const relayClient = yield* RelayClient.RelayClient;
const authorizationError = (requiredScope: AuthEnvironmentScope) =>
new EnvironmentAuthorizationError({
message: `The authenticated token is missing required scope: ${requiredScope}.`,
requiredScope,
});
// RpcScopeAuthorization checks each RPC's declared scope before its handler
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
// runs. This covers the one RPC whose scope depends on its input.
const authorizeEffect = <A, E, R>(
requiredScope: AuthEnvironmentScope,
effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, E | EnvironmentAuthorizationError, R> =>
currentSession.scopes.includes(requiredScope)
? effect
: Effect.fail(authorizationError(requiredScope));
const authorizeStream = <A, E, R>(
requiredScope: AuthEnvironmentScope,
stream: Stream.Stream<A, E, R>,
): Stream.Stream<A, E | EnvironmentAuthorizationError, R> =>
currentSession.scopes.includes(requiredScope)
? stream
: Stream.fail(authorizationError(requiredScope));
: Effect.fail(rpcAuthorizationError(requiredScope));

const acpRegistryProject = Effect.fn("ws.acpRegistry.project")(function* (
projectId: ProjectId,
Expand Down Expand Up @@ -1650,40 +1644,6 @@ const makeWsRpcLayer = (
yield* providerRegistry.refreshInstance(input.instanceId);
return { disabled: true } as const;
});
const observeRpcEffect = <A, E, R>(
method: string,
effect: Effect.Effect<A, E, R>,
traceAttributes?: Readonly<Record<string, unknown>>,
) =>
instrumentRpcEffect(
method,
authorizeEffect(requiredScopeForRpcMethod(method), effect),
traceAttributes,
);
const observeRpcStream = <A, E, R>(
method: string,
stream: Stream.Stream<A, E, R>,
traceAttributes?: Readonly<Record<string, unknown>>,
) =>
instrumentRpcStream(
method,
authorizeStream(requiredScopeForRpcMethod(method), stream),
traceAttributes,
);
const observeRpcStreamEffect = <A, StreamError, StreamContext, EffectError, EffectContext>(
method: string,
effect: Effect.Effect<
Stream.Stream<A, StreamError, StreamContext>,
EffectError,
EffectContext
>,
traceAttributes?: Readonly<Record<string, unknown>>,
) =>
instrumentRpcStreamEffect(
method,
authorizeEffect(requiredScopeForRpcMethod(method), effect),
traceAttributes,
);
const loadAuthAccessSnapshot = () =>
Effect.all({
pairingLinks: serverAuth.listPairingLinks(),
Expand Down Expand Up @@ -3842,6 +3802,7 @@ export const websocketRpcRouteLayer = Layer.unwrap(
const { protocol, httpEffect } = yield* RpcServer.makeProtocolWithHttpEffectWebsocket;
yield* RpcServer.make(ServerWsRpcGroup, { disableTracing: true }).pipe(
Effect.provideService(RpcServer.Protocol, withTerminalOutputWindow(protocol)),
Effect.provide(rpcScopeAuthorizationLayer(session.scopes)),
Effect.forkScoped,
);
// @effect-diagnostics-next-line returnEffectInGen:off
Expand Down
3 changes: 2 additions & 1 deletion docs/internals/environment-auth.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,8 @@ Bearer and DPoP clients obtain short-lived WebSocket tickets through authenticat
HTTP so long-lived tokens stay out of socket URLs. Browser sessions can
authenticate the upgrade with their cookie. A successful handshake grants no
extra authority: [every RPC declares a required
scope](../../apps/server/src/auth/RpcAuthorization.ts).
scope](../../apps/server/src/auth/RpcAuthorization.ts), and the WebSocket RPC
group's `RpcScopeAuthorization` middleware checks it before any handler runs.

Desktop restarts forget the previous local bearer token, so its reusable
bootstrap grant replaces earlier sessions for the same subject and method.
Expand Down
13 changes: 12 additions & 1 deletion packages/contracts/src/rpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
import * as Schema from "effect/Schema";
import * as Rpc from "effect/unstable/rpc/Rpc";
import * as RpcGroup from "effect/unstable/rpc/RpcGroup";
import * as RpcMiddleware from "effect/unstable/rpc/RpcMiddleware";
import { NonNegativeInt, TrimmedNonEmptyString } from "./baseSchemas.ts";
import {
CodexAuthCallbackInput,
Expand Down Expand Up @@ -1682,6 +1683,16 @@ const WsSubscribeResourceTelemetryRpc = Rpc.make(WS_METHODS.subscribeResourceTel
stream: true,
});

/**
* Checks the connection's scopes against the scope each RPC declares, before
* the handler runs. Every RPC in `WsRpcGroup` carries it, so a handler cannot
* be added without authorization.
*/
export class RpcScopeAuthorization extends RpcMiddleware.Service<RpcScopeAuthorization>()(
"t3/contracts/RpcScopeAuthorization",
{ error: EnvironmentAuthorizationError },
) {}

export const WsRpcGroup = RpcGroup.make(
WsServerProbeRpc,
WsServerGetConfigRpc,
Expand Down Expand Up @@ -1856,4 +1867,4 @@ export const WsRpcGroup = RpcGroup.make(
WsOrchestrationV2SubscribeArchivedShellRpc,
WsOrchestrationV2SubscribeShellRpc,
WsOrchestrationV2SubscribeThreadRpc,
);
).middleware(RpcScopeAuthorization);
Loading