diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 64bfea1b9805..2f32b6524d7b 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -5473,6 +5473,43 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect("serves absolute host media without a local thread and rejects relative media", () => + Effect.gen(function* () { + yield* buildAppUnderTest(); + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const directory = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-host-media-" }); + const wsUrl = yield* getWsServerUrl("/ws"); + const threadId = ThreadId.make("thread-on-another-environment"); + + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + for (const [name, mimeType] of [ + ["screenshot.png", "image/png"], + ["recording.mp4", "video/mp4"], + ] as const) { + const filePath = path.join(directory, name); + yield* fileSystem.writeFileString(filePath, "host media bytes"); + const issued = yield* client[WS_METHODS.assetsCreateUrl]({ + resource: { _tag: "media-file", threadId, path: filePath }, + }); + const response = yield* HttpClient.get(issued.relativeUrl); + assert.equal(response.status, 200); + assert.equal(response.headers["content-type"], mimeType); + assert.equal(yield* response.text, "host media bytes"); + + const error = yield* client[WS_METHODS.assetsCreateUrl]({ + resource: { _tag: "media-file", threadId, path: name }, + }).pipe(Effect.flip); + assert.equal(error._tag, "AssetWorkspaceContextNotFoundError"); + } + }), + ), + ); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("uploads image bytes through a signed URL issued by websocket rpc", () => Effect.gen(function* () { const config = yield* buildAppUnderTest(); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 46396bf1e4a7..740be1330817 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -9,6 +9,7 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Path from "effect/Path"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; @@ -2390,9 +2391,12 @@ const makeWsRpcLayer = ( observeRpcEffect( WS_METHODS.assetsCreateUrl, Effect.gen(function* () { + const path = yield* Path.Path; + // An absolute media path can be linked from a thread on another environment. if ( input.resource._tag === "attachment" || - input.resource._tag === "native-app-icon" + input.resource._tag === "native-app-icon" || + (input.resource._tag === "media-file" && path.isAbsolute(input.resource.path)) ) { return yield* issueAssetUrl({ resource: input.resource }); } diff --git a/apps/web/src/state/assets.ts b/apps/web/src/state/assets.ts index d1ef71f8662d..672e2d88672d 100644 --- a/apps/web/src/state/assets.ts +++ b/apps/web/src/state/assets.ts @@ -2,12 +2,30 @@ import { createAssetEnvironmentAtoms, createProjectFaviconUrlAtomFamily, } from "@t3tools/client-runtime/state/assets"; +import { Atom } from "effect/unstable/reactivity"; import { connectionAtomRuntime } from "../connection/runtime"; import { projectFaviconCache } from "../assets/projectFaviconCache"; +import { isElectron } from "../env"; +import { primaryEnvironmentIdAtom } from "./primaryEnvironment"; import { environmentSession } from "./session"; -export const assetEnvironment = createAssetEnvironmentAtoms(connectionAtomRuntime); +const localMediaEnvironment = Atom.make((get) => { + if (!isElectron) return null; + const environmentId = get(primaryEnvironmentIdAtom); + if (environmentId === null) return null; + const connection = get(environmentSession.preparedConnectionValueAtom(environmentId)); + // The session's bootstrap config clears on disconnect and refreshes on reconnect. + const config = get(environmentSession.initialConfigValueAtom(environmentId)); + return connection._tag === "None" || config === null + ? null + : { environmentId, httpBaseUrl: connection.value.httpBaseUrl }; +}); + +export const assetEnvironment = createAssetEnvironmentAtoms( + connectionAtomRuntime, + localMediaEnvironment, +); export const projectFaviconUrlAtom = createProjectFaviconUrlAtomFamily({ imageCache: projectFaviconCache, diff --git a/packages/client-runtime/src/state/assets.test.ts b/packages/client-runtime/src/state/assets.test.ts index 1cbc970df928..bd266f261675 100644 --- a/packages/client-runtime/src/state/assets.test.ts +++ b/packages/client-runtime/src/state/assets.test.ts @@ -1,11 +1,32 @@ import { describe, expect, it } from "@effect/vitest"; -import { type AssetCreateUrlResult, EnvironmentId } from "@t3tools/contracts"; +import { + type AssetCreateUrlResult, + AssetWorkspaceAssetNotFoundError, + AssetWorkspaceAssetInspectionError, + AssetWorkspaceContextNotFoundError, + EnvironmentAuthorizationError, + EnvironmentId, + ThreadId, + WS_METHODS, +} from "@t3tools/contracts"; import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; import * as Option from "effect/Option"; import * as Layer from "effect/Layer"; import { AsyncResult, Atom, AtomRegistry } from "effect/unstable/reactivity"; -import type { EnvironmentRegistry } from "../connection/registry.ts"; +import { EnvironmentRegistry } from "../connection/registry.ts"; +import { + AVAILABLE_CONNECTION_STATE, + PrimaryConnectionTarget, + type PreparedConnection, + type SupervisorConnectionState, +} from "../connection/model.ts"; +import { EnvironmentSupervisor } from "../connection/supervisor.ts"; +import type { RpcSession } from "../rpc/session.ts"; +import type { WsRpcProtocolClient } from "../rpc/protocol.ts"; import { createProjectFaviconCache } from "../projectFaviconCache.ts"; import { createAssetEnvironmentAtoms, @@ -37,6 +58,116 @@ describe("asset collection keys", () => { }); describe("createAssetEnvironmentAtoms", () => { + for (const scenario of [ + { name: "missing video", path: "/tmp/clip.mp4", fallback: true }, + { name: "literal filename characters", path: "/tmp/frame#one?two.png", fallback: true }, + { name: "windows path", path: "C:\\Users\\demo\\clip.mp4", fallback: true }, + { name: "inspection failure", path: "/tmp/clip.mp4", error: "inspection", fallback: true }, + { name: "foreign thread", path: "/tmp/clip.mp4", error: "context", fallback: true }, + { name: "remote success", path: "/tmp/clip.mp4", success: true }, + { name: "relative path", path: "clip.mp4" }, + { name: "no primary", path: "/tmp/clip.mp4", primary: "none" }, + { name: "primary reconnect", path: "/tmp/frame.png", primary: "reconnecting" }, + { name: "same environment", path: "/tmp/clip.mp4", primary: "same" }, + { name: "non-media", path: "/tmp/report.html" }, + { name: "authorization failure", path: "/tmp/clip.mp4", error: "auth" }, + ]) { + it.effect(`uses the correct environment for ${scenario.name}`, () => + Effect.gen(function* () { + const remoteId = EnvironmentId.make("remote"); + const localId = EnvironmentId.make("local"); + const resource = { + _tag: "media-file" as const, + threadId: ThreadId.make("foreign-thread"), + path: scenario.path, + }; + const error = + scenario.error === "auth" + ? new EnvironmentAuthorizationError({ + message: "denied", + requiredScope: "orchestration:read", + }) + : scenario.error === "inspection" + ? new AssetWorkspaceAssetInspectionError({ resource, cause: new Error("unreadable") }) + : scenario.error === "context" + ? new AssetWorkspaceContextNotFoundError({ resource }) + : new AssetWorkspaceAssetNotFoundError({ resource }); + const calls: EnvironmentId[] = []; + const supervisors = new Map(); + for (const environmentId of [remoteId, localId]) { + const client = { + [WS_METHODS.assetsCreateUrl]: () => { + calls.push(environmentId); + return environmentId === remoteId && !scenario.success + ? Effect.fail(error) + : Effect.succeed({ + relativeUrl: `/api/assets/${environmentId}/media`, + expiresAt: 999999, + }); + }, + } as unknown as WsRpcProtocolClient; + const session = { client } as RpcSession; + supervisors.set( + environmentId, + EnvironmentSupervisor.of({ + target: new PrimaryConnectionTarget({ + environmentId, + label: environmentId, + httpBaseUrl: `https://${environmentId}.test`, + wsBaseUrl: `wss://${environmentId}.test`, + }), + state: yield* SubscriptionRef.make({ + ...AVAILABLE_CONNECTION_STATE, + phase: "connected" as const, + }), + session: yield* SubscriptionRef.make(Option.some(session)), + prepared: yield* SubscriptionRef.make(Option.none()), + connect: Effect.void, + disconnect: Effect.void, + retryNow: Effect.void, + }), + ); + } + const environments = EnvironmentRegistry.of({ + run: (id, effect) => + Effect.provideService(effect, EnvironmentSupervisor, supervisors.get(id)!), + followStream: (id, stream) => + Stream.provideService(stream, EnvironmentSupervisor, supervisors.get(id)!), + } as EnvironmentRegistry["Service"]); + const registry = AtomRegistry.make(); + yield* Effect.addFinalizer(() => Effect.sync(() => registry.dispose())); + const localTarget = { + environmentId: scenario.primary === "same" ? remoteId : localId, + httpBaseUrl: "https://local.test", + }; + const localEnvironment = Atom.make( + scenario.primary === "none" || scenario.primary === "reconnecting" ? null : localTarget, + ); + const assets = createAssetEnvironmentAtoms( + Atom.runtime(Layer.succeed(EnvironmentRegistry, environments)), + localEnvironment, + ); + const query = assets.createUrl({ environmentId: remoteId, input: { resource } }); + const result = AtomRegistry.getResult(registry, query, { suspendOnWaiting: true }); + if (scenario.fallback || scenario.success) { + expect((yield* result).relativeUrl).toBe( + scenario.fallback + ? "https://local.test/api/assets/local/media" + : "/api/assets/remote/media", + ); + } else { + expect(yield* Effect.flip(result)).toEqual(error); + } + expect(calls).toEqual(scenario.fallback ? [remoteId, localId] : [remoteId]); + if (scenario.primary === "reconnecting") { + registry.set(localEnvironment, localTarget); + expect((yield* result).relativeUrl).toBe("https://local.test/api/assets/local/media"); + expect(calls).toEqual([remoteId, remoteId, localId]); + } + }).pipe(Effect.scoped), + ); + } + it("keys asset URL queries by environment and resource", () => { const runtime = Atom.runtime(Layer.empty) as unknown as Atom.AtomRuntime< EnvironmentRegistry, diff --git a/packages/client-runtime/src/state/assets.ts b/packages/client-runtime/src/state/assets.ts index b8646d911cc2..bf857fd44daa 100644 --- a/packages/client-runtime/src/state/assets.ts +++ b/packages/client-runtime/src/state/assets.ts @@ -1,22 +1,28 @@ import { + type AssetCreateUrlInput, type AssetCreateUrlResult, type AssetImageDimensions, AssetResource, EnvironmentId, WS_METHODS, } from "@t3tools/contracts"; +import { mediaMimeTypeFromExtension } from "@t3tools/shared/filePreview"; +import { isWindowsAbsolutePath } from "@t3tools/shared/path"; import { getProjectFaviconResourceKey, isProjectFaviconFallbackUrl, } from "@t3tools/shared/projectFavicon"; import * as Effect from "effect/Effect"; import * as Option from "effect/Option"; +import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; -import { AsyncResult, Atom } from "effect/unstable/reactivity"; +import { AsyncResult, Atom, AtomRegistry } from "effect/unstable/reactivity"; -import type { EnvironmentRegistry } from "../connection/registry.ts"; +import * as EnvironmentRegistry from "../connection/registry.ts"; +import * as EnvironmentSupervisor from "../connection/supervisor.ts"; +import { request } from "../rpc/client.ts"; import type { ProjectFaviconCache, ProjectFaviconTarget } from "../projectFaviconCache.ts"; -import { createEnvironmentRpcQueryAtomFamily } from "./runtime.ts"; +import { createEnvironmentQueryAtomFamily } from "./runtime.ts"; const ASSET_URL_REFRESH_INTERVAL_MS = 30 * 60_000; const ASSET_URL_STALE_TIME_MS = 5 * 60_000; @@ -91,14 +97,50 @@ export function assetUrlStateFromResult( } export function createAssetEnvironmentAtoms( - runtime: Atom.AtomRuntime, + runtime: Atom.AtomRuntime, + localMediaEnvironment?: Atom.Atom<{ + readonly environmentId: EnvironmentId; + readonly httpBaseUrl: string; + } | null>, ) { - const createUrl = createEnvironmentRpcQueryAtomFamily(runtime, { + const execute = Effect.fn("assets.createUrl")(function* (input: AssetCreateUrlInput) { + const result = yield* request(WS_METHODS.assetsCreateUrl, input).pipe(Effect.result); + if (Result.isSuccess(result)) return result.success; + const error = result.failure; + const resource = input.resource; + const local = localMediaEnvironment + ? (yield* AtomRegistry.AtomRegistry).get(localMediaEnvironment) + : null; + const supervisor = yield* EnvironmentSupervisor.EnvironmentSupervisor; + if ( + !local || + local.environmentId === supervisor.target.environmentId || + !( + error._tag === "AssetWorkspaceAssetNotFoundError" || + error._tag === "AssetWorkspaceAssetInspectionError" || + error._tag === "AssetWorkspaceContextNotFoundError" + ) || + resource._tag !== "media-file" || + !(resource.path.startsWith("/") || isWindowsAbsolutePath(resource.path)) || + mediaMimeTypeFromExtension(resource.path.slice(resource.path.lastIndexOf("."))) === null + ) + return yield* error; + const registry = yield* EnvironmentRegistry.EnvironmentRegistry; + const asset = yield* registry.run( + local.environmentId, + request(WS_METHODS.assetsCreateUrl, input), + ); + // Callers resolve against the thread's server, so preserve the local server's origin. + return { ...asset, relativeUrl: new URL(asset.relativeUrl, local.httpBaseUrl).href }; + }); + const createUrl = createEnvironmentQueryAtomFamily(runtime, { label: "environment-data:assets:create-url", - tag: WS_METHODS.assetsCreateUrl, + execute, staleTimeMs: ASSET_URL_STALE_TIME_MS, idleTtlMs: ASSET_URL_IDLE_TTL_MS, refreshIntervalMs: ASSET_URL_REFRESH_INTERVAL_MS, + refreshTrigger: ({ input }) => + input.resource._tag === "media-file" ? localMediaEnvironment : undefined, }); const createUrlsFamily = Atom.family((key: string) => { const [environmentId, resources] = parseAssetCollectionKey(key); diff --git a/packages/client-runtime/src/state/runtime.ts b/packages/client-runtime/src/state/runtime.ts index 3e61909ee711..46cb9fa0f115 100644 --- a/packages/client-runtime/src/state/runtime.ts +++ b/packages/client-runtime/src/state/runtime.ts @@ -482,7 +482,12 @@ export function followStreamInEnvironment( export function createEnvironmentQueryAtomFamily( runtime: Atom.AtomRuntime, - options: EnvironmentQueryAtomOptions, + options: EnvironmentQueryAtomOptions< + Input, + A, + E, + EnvironmentSupervisor | EnvironmentRegistry | AtomRegistry.AtomRegistry | R + >, ): (target: { readonly environmentId: EnvironmentIdType; readonly input: Input;