Skip to content
Open
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
104 changes: 104 additions & 0 deletions packages/client-runtime/src/state/session.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
import { describe, expect, it } from "@effect/vitest";
import { EnvironmentId, type AuthSessionState } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Stream from "effect/Stream";
import * as SubscriptionRef from "effect/SubscriptionRef";
import { AsyncResult, Atom, AtomRegistry } from "effect/reactivity";

import { RelayConnectionTarget, type PreparedConnection } from "../connection/model.ts";
import { EnvironmentRegistry } from "../connection/registry.ts";
import { EnvironmentSupervisor } from "../connection/supervisor.ts";
import * as RpcHttp from "../rpc/http.ts";
import { createEnvironmentSessionAtoms } from "./session.ts";

const TARGET = new RelayConnectionTarget({
environmentId: EnvironmentId.make("environment-1"),
label: "Environment",
});
const PREPARED: PreparedConnection = {
environmentId: TARGET.environmentId,
label: TARGET.label,
httpBaseUrl: "https://environment.example.test",
socketUrl: "wss://environment.example.test/ws",
httpAuthorization: { _tag: "Bearer", token: "bearer-token" },
target: TARGET,
};
const SESSION = {
authenticated: true,
auth: {
policy: "remote-reachable",
bootstrapMethods: ["desktop-bootstrap"],
sessionMethods: ["bearer-access-token"],
sessionCookieName: "t3_session",
},
scopes: ["orchestration:read", "orchestration:operate"],
} satisfies AuthSessionState;

function makeSessionAtom(reply: (requestNumber: number) => Promise<Response>) {
let calls = 0;
const fetchFn: typeof fetch = () => reply(++calls);
const layer = Effect.gen(function* () {
const supervisor = {
target: TARGET,
prepared: yield* SubscriptionRef.make(Option.some(PREPARED)),
} as unknown as EnvironmentSupervisor["Service"];
const registry = {
followStream: <A, E, R>(_environmentId: EnvironmentId, stream: Stream.Stream<A, E, R>) =>
stream.pipe(Stream.provideService(EnvironmentSupervisor, supervisor)),
} as unknown as EnvironmentRegistry["Service"];
return Layer.merge(
Layer.succeed(EnvironmentRegistry, registry),
RpcHttp.layerRemoteHttpClient(fetchFn),
);
}).pipe(Layer.unwrap);
const atom = createEnvironmentSessionAtoms(Atom.runtime(layer)).sessionStateAtom(
TARGET.environmentId,
);
return { atom, calls: () => calls };
}

const settle = (
registry: AtomRegistry.AtomRegistry,
atom: ReturnType<typeof makeSessionAtom>["atom"],
) => AtomRegistry.getResult(registry, atom, { suspendOnWaiting: true }).pipe(Effect.result);

describe("environment session state", () => {
it.effect("keeps a confirmed grant through a failed refresh and retries it", () =>
Effect.gen(function* () {
const session = makeSessionAtom(async (requestNumber) => {
if (requestNumber === 2) throw new TypeError("Failed to fetch");
return Response.json(SESSION);
});
const registry = AtomRegistry.make();
yield* Effect.addFinalizer(() => Effect.sync(() => registry.dispose()));
const observed: Array<string> = [];
registry.mount(session.atom);
registry.subscribe(session.atom, (result) => observed.push(result._tag));

expect(yield* settle(registry, session.atom)).toMatchObject({ success: SESSION });
registry.refresh(session.atom);
expect(yield* settle(registry, session.atom)).toMatchObject({ success: SESSION });

expect(session.calls()).toBe(3);
expect(observed).not.toContain("Failure");
}).pipe(Effect.scoped),
);

it.effect("reports a failed first load without retrying", () =>
Effect.gen(function* () {
const session = makeSessionAtom(async () => {
throw new TypeError("Failed to fetch");
});
const registry = AtomRegistry.make();
yield* Effect.addFinalizer(() => Effect.sync(() => registry.dispose()));
registry.mount(session.atom);

const result = yield* settle(registry, session.atom);
expect(result._tag).toBe("Failure");
expect(AsyncResult.isFailure(registry.get(session.atom))).toBe(true);
expect(session.calls()).toBe(1);
}).pipe(Effect.scoped),
);
});
38 changes: 37 additions & 1 deletion packages/client-runtime/src/state/session.ts
Original file line number Diff line number Diff line change
@@ -1,18 +1,22 @@
import type { AuthSessionState, EnvironmentId, ServerConfig } from "@t3tools/contracts";
import * as Data from "effect/Data";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Schedule from "effect/Schedule";
import * as Stream from "effect/Stream";
import * as SubscriptionRef from "effect/SubscriptionRef";
import { HttpClient } from "effect/http";
import { AsyncResult, Atom } from "effect/reactivity";

import * as RemoteEnvironmentAuthorization from "../authorization/service.ts";
import { mapRemoteEnvironmentError } from "../connection/errors.ts";
import * as EnvironmentRegistry from "../connection/registry.ts";
import type { PreparedConnection } from "../connection/model.ts";
import * as EnvironmentSupervisor from "../connection/supervisor.ts";
import { environmentEndpointUrl } from "../environment/endpoint.ts";
import * as ManagedRelay from "../relay/managedRelay.ts";
import type { RemoteEnvironmentRequestError } from "../rpc/http.ts";
import { safeErrorLogAttributes } from "../errors/safeLog.ts";
import { executeAuthenticatedEnvironmentHttpRequest } from "./environmentHttpAuth.ts";
import { followStreamInEnvironment } from "./environmentStreams.ts";
Expand All @@ -39,6 +43,20 @@ function initialConfigOption<E>(
// Bounded so a wedged environment cannot pin the permissions check (and with it
// the settings UI) in a loading state for long.
const DEFAULT_SESSION_STATE_TIMEOUT_MS = 6_000;
const SESSION_STATE_RETRY_SCHEDULE = Schedule.exponential("1 second").pipe(
Schedule.modifyDelay(({ duration }) =>
Effect.succeed(Duration.min(duration, Duration.seconds(30))),
),
);

function isTransientSessionError(
error: SessionHttpClientUnavailable | RemoteEnvironmentRequestError,
): boolean {
return (
error._tag !== "SessionHttpClientUnavailable" &&
mapRemoteEnvironmentError(error)._tag === "ConnectionTransientError"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High state/session.ts:57

Revoked relay access is treated as transient, so a confirmed session keeps its old scopes and retries rejected credentials instead of surfacing the authorization failure; thread controls can remain enabled while the check waits. mapRemoteEnvironmentError classifies the RemoteEnvironmentAuthFetchError wrapper without inspecting its cause, so preserve or classify permanent authorization failures before retrying.

🤖 Copy this AI Prompt to have your agent fix this:
In file @packages/client-runtime/src/state/session.ts around line 57:

Revoked relay access is treated as transient, so a confirmed session keeps its old scopes and retries rejected credentials instead of surfacing the authorization failure; thread controls can remain enabled while the check waits. `mapRemoteEnvironmentError` classifies the `RemoteEnvironmentAuthFetchError` wrapper without inspecting its cause, so preserve or classify permanent authorization failures before retrying.

);
}

/**
* Read the granted scopes of this client's session on one environment via its
Expand Down Expand Up @@ -138,7 +156,17 @@ function makeEnvironmentSessionAtoms<R, E>(
if (prepared === null) {
return Effect.never;
}
return Effect.gen(function* () {
// Scope checks reject a Failure outright, so one slow response would
// revoke a grant this connection already confirmed until the client
// reloads. With a confirmed grant, keep it (the result stays waiting)
// and retry transient failures; rejected credentials still fail.
const hasConfirmedGrant = Option.isSome(
Option.flatMap(
get.self<AsyncResult.AsyncResult<AuthSessionState, unknown>>(),
AsyncResult.value,
),
);
const fetchSession = Effect.gen(function* () {
const signer = yield* Effect.serviceOption(ManagedRelay.ManagedRelayDpopSigner);
const remoteAuthorization = yield* Effect.serviceOption(
RemoteEnvironmentAuthorization.RemoteEnvironmentAuthorization,
Expand All @@ -151,6 +179,14 @@ function makeEnvironmentSessionAtoms<R, E>(
remoteAuthorization,
}).pipe(Effect.provideService(HttpClient.HttpClient, client.value));
});
return hasConfirmedGrant
? fetchSession.pipe(
Effect.retry({
while: isTransientSessionError,
schedule: SESSION_STATE_RETRY_SCHEDULE,
}),
)
: fetchSession;
})
.pipe(
Atom.swr({ staleTime: 30_000, revalidateOnMount: true }),
Expand Down
Loading