diff --git a/apps/web/__tests__/unit/content-transfer.test.ts b/apps/web/__tests__/unit/content-transfer.test.ts index 126f9a7f907..43e0958b95d 100644 --- a/apps/web/__tests__/unit/content-transfer.test.ts +++ b/apps/web/__tests__/unit/content-transfer.test.ts @@ -175,6 +175,24 @@ describe("content transfer planning", () => { }; expect(getContentTransferStorageBlockReason(input)).toBeNull(); + expect( + getContentTransferStorageBlockReason({ + ...input, + bucketId: "cap-tokyo", + bucketOwnerId: null, + bucketOrganizationId: null, + }), + ).toBeNull(); + expect( + getContentTransferStorageBlockReason({ ...input, bucketOwnerId: null }), + ).toBe("The storage bucket is missing"); + expect( + getContentTransferStorageBlockReason({ + ...input, + bucketId: "cap-tokyo", + storageIntegrationId: "drive", + }), + ).toBe("The Cap has conflicting storage assignments"); expect( getContentTransferStorageBlockReason({ ...input, diff --git a/apps/web/__tests__/unit/desktop-video-create.test.ts b/apps/web/__tests__/unit/desktop-video-create.test.ts index b3ae9980899..603ca7c94ce 100644 --- a/apps/web/__tests__/unit/desktop-video-create.test.ts +++ b/apps/web/__tests__/unit/desktop-video-create.test.ts @@ -5,8 +5,10 @@ import { Storage as StorageDomain, Video, } from "@cap/web-domain"; -import { Effect, Option } from "effect"; -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { Effect, Layer, Option } from "effect"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const regionalStorage = vi.hoisted(() => ({ select: vi.fn(), s3: vi.fn() })); const deletion = vi.hoisted(() => ({ deleteVideo: vi.fn(), @@ -86,19 +88,23 @@ vi.mock("@cap/web-backend", async () => { class Videos extends Effect.Service()("Videos", { sync: () => ({ delete: deletion.deleteVideo }), }) {} + class S3Buckets extends Effect.Service()("S3Buckets", { + sync: () => ({ getRegionalUploadBucketId: regionalStorage.select }), + }) {} return { makeCurrentUserLayer, Videos, + S3Buckets, Storage: { getOrganizationWritableAccess: vi.fn(), - getS3WritableAccessForUser: vi.fn(), + getS3WritableAccessForUser: regionalStorage.s3, }, }; }); vi.mock("@/lib/server", async () => { const { Effect } = await import("effect"); - const { Videos } = await import("@cap/web-backend"); + const { Videos, S3Buckets } = await import("@cap/web-backend"); return { runPromise: vi.fn(async (value: unknown) => Effect.isEffect(value) @@ -107,9 +113,11 @@ vi.mock("@/lib/server", async () => { value as Effect.Effect< unknown, unknown, - InstanceType + InstanceType | InstanceType > - ).pipe(Effect.provide(Videos.Default)), + ).pipe( + Effect.provide(Layer.mergeAll(Videos.Default, S3Buckets.Default)), + ), ) : value, ), @@ -126,6 +134,8 @@ vi.mock("@/lib/google-drive-storage-quota", () => ({ // The live-transcription stack drags in the whole workflow graph // (server-only modules included); these tests only care that create works. +vi.mock("next/server", () => ({ after: vi.fn() })); + vi.mock("@/lib/live-transcribe", () => ({ maybeStartLiveTranscription: vi.fn(async () => "skipped"), })); @@ -667,3 +677,156 @@ describe("GET /create", () => { expect(mockDb.insert).not.toHaveBeenCalled(); }); }); + +describe("new Instant recording regions", () => { + let app: typeof import("@/app/api/desktop/[...route]/video")["app"]; + beforeEach(async () => { + vi.clearAllMocks(); + vi.stubEnv("VERCEL", "1"); + resetMockDb(); + stubStorage(); + regionalStorage.s3.mockReturnValue( + Effect.succeed({ + bucketId: Option.none(), + storageIntegrationId: Option.none(), + }), + ); + regionalStorage.select.mockImplementation( + (latitude: string | undefined, longitude: string | undefined) => + latitude === "35.68" && longitude === "139.69" + ? Option.some("cap-tokyo") + : Option.none(), + ); + defaultSharing.getNewVideoPublic.mockResolvedValue(true); + mockGetCurrentUser.mockResolvedValue({ + id: "user-1", + defaultOrgId: "org-1", + activeOrganizationId: "org-1", + }); + mockDb.where + .mockResolvedValueOnce([ + { id: "org-1", name: "Org", createdAt: new Date() }, + ]) + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([{ count: 5 }]); + app = (await import("@/app/api/desktop/[...route]/video")).app; + }); + afterEach(() => vi.unstubAllEnvs()); + + it.each(["desktopMP4", "desktopSegments"])( + "persists the selected bucket for %s", + async (mode) => { + const response = await app.request( + `https://cap.test/create?recordingMode=${mode}`, + { + headers: { + "x-vercel-ip-latitude": "35.68", + "x-vercel-ip-longitude": "139.69", + }, + }, + ); + expect(response.status).toBe(200); + expect(insertedValues(schema.videos)?.bucket).toBe("cap-tokyo"); + expect(regionalStorage.select).toHaveBeenCalledWith("35.68", "139.69"); + }, + ); + + it.each>([ + { "x-vercel-ip-latitude": "40.71", "x-vercel-ip-longitude": "-74.01" }, + {}, + { "x-vercel-ip-latitude": "bad", "x-vercel-ip-longitude": "139.69" }, + ])( + "keeps Virginia when no regional bucket is selected (%#)", + async (headers) => { + const response = await app.request( + "https://cap.test/create?recordingMode=desktopMP4", + { headers }, + ); + expect(response.status).toBe(200); + expect(insertedValues(schema.videos)?.bucket).toBeNull(); + }, + ); + + it.each([ + { + bucket: "custom-bucket", + integration: null, + vercel: "1", + query: "recordingMode=desktopMP4", + }, + { + bucket: null, + integration: "drive-id", + vercel: "1", + query: "recordingMode=desktopMP4", + }, + { + bucket: null, + integration: null, + vercel: "", + query: "recordingMode=desktopMP4", + }, + { + bucket: null, + integration: null, + vercel: "1", + query: "isScreenshot=true&recordingMode=desktopMP4", + }, + { + bucket: null, + integration: null, + vercel: "1", + query: "recordingMode=hls", + }, + ])( + "preserves storage outside eligible new Instant recordings (%#)", + async ({ bucket, integration, vercel, query }) => { + vi.stubEnv("VERCEL", vercel); + regionalStorage.s3.mockReturnValue( + Effect.succeed({ + bucketId: Option.fromNullable(bucket), + storageIntegrationId: Option.fromNullable(integration), + }), + ); + const response = await app.request(`https://cap.test/create?${query}`, { + headers: { + "x-vercel-ip-latitude": "35.68", + "x-vercel-ip-longitude": "139.69", + }, + }); + expect(response.status).toBe(200); + expect(insertedValues(schema.videos)).toMatchObject({ + bucket, + storageIntegrationId: integration, + }); + expect(regionalStorage.select).not.toHaveBeenCalled(); + }, + ); + + it.each([null, "cap-tokyo", "custom-bucket"])( + "does not reroute an existing recording after travel: %s", + async (bucket) => { + mockDb.where.mockReset().mockResolvedValue([ + { + id: "existing", + ownerId: "user-1", + bucket, + source: { type: "desktopMP4" }, + }, + ]); + const response = await app.request( + "https://cap.test/create?videoId=existing&recordingMode=desktopMP4", + { + headers: { + "x-vercel-ip-latitude": "35.68", + "x-vercel-ip-longitude": "139.69", + }, + }, + ); + expect(response.status).toBe(200); + expect(mockDb.insert).not.toHaveBeenCalled(); + expect(mockDb.update).not.toHaveBeenCalled(); + expect(regionalStorage.select).not.toHaveBeenCalled(); + }, + ); +}); diff --git a/apps/web/__tests__/unit/regional-organization-cleanup.test.ts b/apps/web/__tests__/unit/regional-organization-cleanup.test.ts new file mode 100644 index 00000000000..0c395603f8f --- /dev/null +++ b/apps/web/__tests__/unit/regional-organization-cleanup.test.ts @@ -0,0 +1,152 @@ +import { CurrentUser, Organisation, S3Bucket, User } from "@cap/web-domain"; +import { Effect, Option } from "effect"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + database: vi.fn(), + bucket: vi.fn(), + deleted: vi.fn(), +})); +vi.mock("@cap/web-backend/src/Database", async () => { + const { Effect } = await import("effect"); + class Database extends Effect.Service()("Database", { + sync: () => ({ use: mocks.database }), + }) {} + return { Database }; +}); +vi.mock("@cap/web-backend/src/S3Buckets", async () => { + const { Effect } = await import("effect"); + class S3Buckets extends Effect.Service()("S3Buckets", { + sync: () => ({ getBucketAccess: mocks.bucket }), + }) {} + return { S3Buckets }; +}); +vi.mock("@cap/web-backend/src/ImageUploads", async () => { + const { Effect } = await import("effect"); + class ImageUploads extends Effect.Service()("ImageUploads", { + sync: () => ({}), + }) {} + return { ImageUploads }; +}); +vi.mock("@cap/web-backend/src/Tinybird", async () => { + const { Effect } = await import("effect"); + class Tinybird extends Effect.Service()("Tinybird", { + sync: () => ({ deleteData: () => Effect.void }), + }) {} + return { Tinybird }; +}); +vi.mock("@cap/web-backend/src/Organisations/OrganisationsPolicy", async () => { + const { Effect } = await import("effect"); + const { Policy } = await import("@cap/web-domain"); + class OrganisationsPolicy extends Effect.Service()( + "OrganisationsPolicy", + { + sync: () => ({ + isOwner: () => Policy.policy(() => Effect.succeed(true)), + }), + }, + ) {} + return { OrganisationsPolicy }; +}); + +import { Organisations } from "@cap/web-backend/src/Organisations"; + +const cleanup = () => + Effect.runPromise( + Effect.flatMap(Organisations, (organizations) => + organizations.softDelete(Organisation.OrganisationId.make("org")), + ).pipe( + Effect.provide(Organisations.Default), + Effect.provideService(CurrentUser, { + id: User.UserId.make("owner"), + email: "owner@cap.test", + activeOrganizationId: Organisation.OrganisationId.make("org"), + iconUrlOrKey: Option.none(), + }), + ), + ); + +beforeEach(() => { + vi.resetAllMocks(); + mocks.database + .mockReturnValueOnce(Effect.succeed([{ id: "org", ownerId: "owner" }])) + .mockReturnValueOnce( + Effect.succeed([ + { + id: "virginia", + ownerId: "owner", + bucket: null, + storageIntegrationId: null, + }, + { + id: "tokyo", + ownerId: "owner", + bucket: S3Bucket.S3BucketId.make("cap-tokyo"), + storageIntegrationId: null, + }, + { + id: "custom", + ownerId: "owner", + bucket: "custom-bucket", + storageIntegrationId: null, + }, + { + id: "drive", + ownerId: "owner", + bucket: null, + storageIntegrationId: "drive-id", + }, + ]), + ) + .mockReturnValue(Effect.void); + mocks.bucket.mockImplementation((bucket: Option.Option) => + Effect.succeed([ + { + listObjects: ({ + prefix, + continuationToken, + }: { + prefix: string; + continuationToken?: string; + }) => + Effect.succeed({ + Contents: [{ Key: `${prefix}${continuationToken ?? "first"}` }], + IsTruncated: !continuationToken, + NextContinuationToken: "second", + }), + deleteObjects: (objects: Array<{ Key: string }>) => + Effect.sync(() => mocks.deleted(Option.getOrNull(bucket), objects)), + }, + Option.none(), + ]), + ); +}); + +describe("organization regional media cleanup", () => { + it("deletes both pages in each managed bucket and leaves customer storage alone", async () => { + await cleanup(); + expect(mocks.deleted.mock.calls).toEqual( + expect.arrayContaining([ + [null, [{ Key: "owner/virginia/first" }]], + [null, [{ Key: "owner/virginia/second" }]], + ["cap-tokyo", [{ Key: "owner/tokyo/first" }]], + ["cap-tokyo", [{ Key: "owner/tokyo/second" }]], + [null, [{ Key: "organizations/org/first" }]], + [null, [{ Key: "organizations/org/second" }]], + ]), + ); + expect(mocks.deleted).toHaveBeenCalledTimes(6); + expect(mocks.database).toHaveBeenCalledTimes(3); + }); + + it("keeps the database records when the regional bucket cannot be opened", async () => { + const original = mocks.bucket.getMockImplementation(); + mocks.bucket.mockImplementation((bucket: Option.Option) => + Option.getOrNull(bucket) === S3Bucket.S3BucketId.make("cap-tokyo") + ? Effect.fail(new Error("Tokyo unavailable")) + : original?.(bucket), + ); + await expect(cleanup()).rejects.toThrow("Tokyo unavailable"); + expect(mocks.database).toHaveBeenCalledTimes(2); + }); +}); diff --git a/apps/web/__tests__/unit/regional-upload-selection.test.ts b/apps/web/__tests__/unit/regional-upload-selection.test.ts new file mode 100644 index 00000000000..86a9e9e8dd5 --- /dev/null +++ b/apps/web/__tests__/unit/regional-upload-selection.test.ts @@ -0,0 +1,147 @@ +import { + getNearestRegionalBucket, + parseRegionalBuckets, +} from "@cap/web-backend/src/S3Buckets/RegionalBuckets"; +import { S3Bucket } from "@cap/web-domain"; +import { Option } from "effect"; +import { describe, expect, it } from "vitest"; + +const regions = S3Bucket.RegionalBuckets; + +describe("upload region selection", () => { + it.each([ + ["New York", "40.71", "-74.01", null], + ["San Francisco", "37.77", "-122.42", "cap-oregon"], + ["London", "51.51", "-0.13", "cap-london"], + ["Berlin", "52.52", "13.41", "cap-frankfurt"], + ["Buenos Aires", "-34.60", "-58.38", null], + ["Johannesburg", "-26.20", "28.04", "cap-uae"], + ["Delhi", "28.61", "77.21", "cap-mumbai"], + ["Bangkok", "13.76", "100.50", "cap-malaysia"], + ["Tokyo", "35.68", "139.69", "cap-tokyo"], + ["Melbourne", "-37.81", "144.96", "cap-sydney"], + ["Columbus", "40.10", "-83.00", null], + ["Montreal", "45.50", "-73.57", "cap-montreal"], + ["Calgary", "51.04", "-114.07", "cap-oregon"], + ["Paris", "48.86", "2.35", "cap-paris"], + ["Stockholm", "59.33", "18.07", "cap-stockholm"], + ["Milan", "45.46", "9.19", "cap-frankfurt"], + ["Madrid", "40.42", "-3.70", "cap-paris"], + ["Manama", "26.22", "50.59", "cap-uae"], + ["Abu Dhabi", "24.45", "54.38", "cap-uae"], + ["Hyderabad", "17.39", "78.49", "cap-mumbai"], + ["Singapore", "1.35", "103.82", "cap-singapore"], + ["Kuala Lumpur", "3.14", "101.69", "cap-malaysia"], + ["Osaka", "34.69", "135.50", "cap-tokyo"], + ["Sydney", "-33.87", "151.21", "cap-sydney"], + ["Fiji", "-18.14", "178.44", "cap-sydney"], + ["Samoa", "-13.85", "-171.75", "cap-sydney"], + ])( + "selects the nearest configured region for %s", + (_, latitude, longitude, expected) => { + expect( + Option.getOrNull( + getNearestRegionalBucket(latitude, longitude, regions), + ), + ).toBe(expected); + }, + ); + + it.each([ + [undefined, undefined], + ["35.68", undefined], + [undefined, "139.69"], + ["", ""], + [" ", "139.69"], + ["35.68", " "], + ["NaN", "139.69"], + ["Infinity", "139.69"], + ["91", "0"], + ["-91", "0"], + ["0", "181"], + ["0", "-181"], + ["35.68,1", "139.69"], + ])( + "defaults to Virginia for missing or invalid coordinates (%#)", + (latitude, longitude) => { + expect( + Option.isNone(getNearestRegionalBucket(latitude, longitude, regions)), + ).toBe(true); + }, + ); + + it("uses only configured destinations and always considers Virginia", () => { + const tokyo = regions.filter( + (region) => region.region === "ap-northeast-1", + ); + expect(Option.isNone(getNearestRegionalBucket("35.68", "139.69", []))).toBe( + true, + ); + expect( + Option.isNone(getNearestRegionalBucket("51.51", "-0.13", tokyo)), + ).toBe(true); + expect( + Option.getOrNull(getNearestRegionalBucket("35.68", "139.69", tokyo)), + ).toBe("cap-tokyo"); + }); + + it("keeps managed IDs distinct from customer IDs and within the database limit", () => { + for (const region of regions) { + expect(region.id).toContain("-"); + expect(region.id.length).toBeLessThanOrEqual(15); + expect(S3Bucket.isCapManagedBucket(region.id)).toBe(true); + } + expect(new Set(regions.map((region) => region.id)).size).toBe( + regions.length, + ); + expect(S3Bucket.isCapManagedBucket(null)).toBe(true); + expect(S3Bucket.isCapManagedBucket("customer1234567")).toBe(false); + }); + + it("ignores unsupported regions and rejects malformed configuration", () => { + const bucket = { + bucket: "cap-test-tokyo", + bucketUrl: "https://tokyo.cap.test", + distributionId: "ETOKYO", + }; + expect( + parseRegionalBuckets(JSON.stringify({ "not-a-region": bucket })), + ).toEqual([]); + expect( + parseRegionalBuckets(JSON.stringify({ "ap-northeast-1": bucket })), + ).toMatchObject([{ id: "cap-tokyo", ...bucket }]); + expect( + parseRegionalBuckets( + JSON.stringify({ + "ap-northeast-1": { ...bucket, distributionId: " " }, + }), + ), + ).toEqual([]); + }); + it("does not admit the excluded high-cost regions", () => { + const bucket = { + bucket: "unused-bucket", + bucketUrl: "https://unused.cap.test", + distributionId: "EUNUSED", + }; + expect( + parseRegionalBuckets( + JSON.stringify({ "sa-east-1": bucket, "af-south-1": bucket }), + ), + ).toEqual([]); + }); + + it("preserves valid regions when another region is misconfigured", () => { + const config = parseRegionalBuckets( + JSON.stringify({ + "ap-northeast-1": { + bucket: "cap-test-tokyo", + bucketUrl: "https://tokyo.cap.test", + distributionId: "ETOKYO", + }, + "eu-central-1": { bucket: "cap-test-frankfurt" }, + }), + ); + expect(config.map((bucket) => bucket.id)).toEqual(["cap-tokyo"]); + }); +}); diff --git a/apps/web/__tests__/unit/s3-bucket-connections.test.ts b/apps/web/__tests__/unit/s3-bucket-connections.test.ts index 07df81095f0..80e7a442cc3 100644 --- a/apps/web/__tests__/unit/s3-bucket-connections.test.ts +++ b/apps/web/__tests__/unit/s3-bucket-connections.test.ts @@ -1,5 +1,7 @@ +import { generateKeyPairSync } from "node:crypto"; import { createServer, request as httpRequest } from "node:http"; import type { Socket } from "node:net"; +import * as S3 from "@aws-sdk/client-s3"; import { S3Bucket } from "@cap/web-domain"; import { ConfigProvider, Effect, Layer, ManagedRuntime, Option } from "effect"; import { describe, expect, it, vi } from "vitest"; @@ -42,7 +44,7 @@ vi.mock("@cap/web-backend/src/S3Buckets/S3BucketsRepo.ts", async () => { import { S3Buckets } from "@cap/web-backend/src/S3Buckets"; import { s3ConnectionPool } from "@cap/web-backend/src/S3Buckets/S3ConnectionPool"; -async function storageFixture() { +async function storageFixture(config: Record = {}) { let connections = 0; const sockets = new Set(); const authorizations: string[] = []; @@ -73,6 +75,7 @@ async function storageFixture() { ["CAP_AWS_BUCKET", "capso"], ["S3_INTERNAL_ENDPOINT", endpoint], ["S3_PUBLIC_ENDPOINT", endpoint], + ...Object.entries(config), ]), ), ), @@ -228,3 +231,241 @@ describe("S3 connection reuse", () => { } }); }); + +const tokyoConfig = { + bucket: "cap-test-tokyo", + bucketUrl: "https://tokyo-cdn.cap.test", + distributionId: "ETOKYO", +}; +const regionalConfig = { + CAP_REGIONAL_UPLOADS_ENABLED: "true", + CAP_REGIONAL_UPLOAD_BUCKETS: JSON.stringify({ + "ap-northeast-1": tokyoConfig, + }), + CAP_CLOUDFRONT_DISTRIBUTION_ID: "EVIRGINIA", + CAP_AWS_BUCKET_URL: "https://cdn.cap.test", + CLOUDFRONT_KEYPAIR_ID: "KTEST", + CLOUDFRONT_KEYPAIR_PRIVATE_KEY: generateKeyPairSync("rsa", { + modulusLength: 2048, + }) + .privateKey.export({ type: "pkcs8", format: "pem" }) + .toString(), +}; + +describe("regional storage", () => { + it("selects per upload without I/O and keeps regional reads after routing is disabled", async () => { + const fixture = await storageFixture(regionalConfig); + const repoCalls = mocks.getById.mock.calls.length; + const send = vi.spyOn(S3.S3Client.prototype, "send"); + try { + for (const [latitude, longitude, expected] of [ + ["35.68", "139.69", "cap-tokyo"], + ["40.71", "-74.01", null], + ["35.68", "139.69", "cap-tokyo"], + [undefined, undefined, null], + ["bad", "139.69", null], + ] as const) { + expect( + Option.getOrNull( + fixture.service.getRegionalUploadBucketId(latitude, longitude), + ), + ).toBe(expected); + } + const [regional] = await fixture.runtime.runPromise( + fixture.service.getBucketAccess( + Option.some(S3Bucket.S3BucketId.make("cap-tokyo")), + ), + ); + const key = "owner/video/segments/segment_000001.m4s"; + for (const effect of [ + regional.getPresignedPutUrl(key), + regional.getInternalSignedObjectUrl(key), + regional.getInternalPresignedPutUrl(key), + ]) { + const url = new URL(await fixture.runtime.runPromise(effect)); + expect(url.hostname).toBe( + "cap-test-tokyo.s3.ap-northeast-1.amazonaws.com", + ); + expect(url.pathname).toBe(`/${key}`); + expect(url.searchParams.get("X-Amz-Credential")).toContain( + "/ap-northeast-1/s3/", + ); + } + const playback = new URL( + await fixture.runtime.runPromise(regional.getSignedObjectUrl(key)), + ); + expect(playback.origin).toBe(tokyoConfig.bucketUrl); + expect(playback.searchParams.get("Key-Pair-Id")).toBe("KTEST"); + const [original] = await fixture.runtime.runPromise( + fixture.service.getBucketAccess(Option.none()), + ); + expect( + new URL( + await fixture.runtime.runPromise(original.getSignedObjectUrl(key)), + ).origin, + ).toBe(regionalConfig.CAP_AWS_BUCKET_URL); + expect(send).not.toHaveBeenCalled(); + expect(mocks.getById.mock.calls).toHaveLength(repoCalls); + } finally { + send.mockRestore(); + await fixture.close(); + } + const disabled = await storageFixture({ + ...regionalConfig, + CAP_REGIONAL_UPLOADS_ENABLED: "false", + }); + try { + expect( + Option.isNone( + disabled.service.getRegionalUploadBucketId("35.68", "139.69"), + ), + ).toBe(true); + const [regional] = await disabled.runtime.runPromise( + disabled.service.getBucketAccess( + Option.some(S3Bucket.S3BucketId.make("cap-tokyo")), + ), + ); + expect(regional.bucketName).toBe("cap-test-tokyo"); + } finally { + await disabled.close(); + } + }); + + it.each([ + {}, + { CAP_REGIONAL_UPLOADS_ENABLED: "true" }, + { ...regionalConfig, CAP_REGIONAL_UPLOADS_ENABLED: "invalid" }, + { ...regionalConfig, CAP_AWS_REGION: "eu-west-1" }, + ...[ + "not json", + "null", + "[]", + JSON.stringify({ + "ap-northeast-1": { ...tokyoConfig, bucket: "INVALID" }, + }), + JSON.stringify({ + "ap-northeast-1": { + ...tokyoConfig, + bucketUrl: "http://tokyo-cdn.cap.test", + }, + }), + JSON.stringify({ + "ap-northeast-1": { + ...tokyoConfig, + bucketUrl: "https://tokyo-cdn.cap.test/path", + }, + }), + JSON.stringify({ + "ap-northeast-1": { ...tokyoConfig, distributionId: "" }, + }), + ].map((value) => ({ + ...regionalConfig, + CAP_REGIONAL_UPLOAD_BUCKETS: value, + })), + Object.fromEntries( + Object.entries(regionalConfig).filter( + ([key]) => key !== "CLOUDFRONT_KEYPAIR_ID", + ), + ), + ])( + "keeps the original upload endpoint when routing is unavailable (%#)", + async (config) => { + const fixture = await storageFixture({ + ...config, + S3_PUBLIC_ENDPOINT: "https://s3-accelerate.amazonaws.com", + S3_PATH_STYLE: "false", + }); + try { + expect( + Option.isNone( + fixture.service.getRegionalUploadBucketId("35.68", "139.69"), + ), + ).toBe(true); + const [original] = await fixture.runtime.runPromise( + fixture.service.getBucketAccess(Option.none()), + ); + const url = new URL( + await fixture.runtime.runPromise( + original.getPresignedPutUrl("owner/video/result.mp4"), + ), + ); + expect(url.hostname).toBe("capso.s3-accelerate.amazonaws.com"); + expect(url.searchParams.get("X-Amz-Credential")).toContain( + `/${"CAP_AWS_REGION" in config ? config.CAP_AWS_REGION : "us-east-1"}/s3/`, + ); + } finally { + await fixture.close(); + } + }, + ); + + it("fails a persisted regional recording safely instead of writing to Virginia when config is lost", async () => { + const fixture = await storageFixture(); + const repoCalls = mocks.getById.mock.calls.length; + try { + await expect( + fixture.runtime.runPromise( + fixture.service.getBucketAccess( + Option.some(S3Bucket.S3BucketId.make("cap-tokyo")), + ), + ), + ).rejects.toThrow("Regional storage configuration"); + expect(mocks.getById.mock.calls).toHaveLength(repoCalls); + } finally { + await fixture.close(); + } + }); + + it.each(S3Bucket.RegionalBuckets)( + "keeps signing, processing, copying, and cleanup in $region", + async (region) => { + const bucket = `${region.id}-test`; + const fixture = await storageFixture({ + ...regionalConfig, + CAP_REGIONAL_UPLOAD_BUCKETS: JSON.stringify({ + [region.region]: { ...tokyoConfig, bucket }, + }), + }); + const send = vi + .spyOn(S3.S3Client.prototype, "send") + .mockImplementation(async () => ({})); + try { + const [regional] = await fixture.runtime.runPromise( + fixture.service.getBucketAccess(Option.some(region.id)), + ); + const upload = new URL( + await fixture.runtime.runPromise( + regional.getPresignedPutUrl("owner/video/result.mp4"), + ), + ); + expect(upload.hostname).toBe( + `${bucket}.s3.${region.region}.amazonaws.com`, + ); + expect(upload.searchParams.get("X-Amz-Credential")).toContain( + `/${region.region}/s3/`, + ); + await fixture.runtime.runPromise( + regional.headObject("owner/video/result.mp4"), + ); + await fixture.runtime.runPromise( + regional.copyObject( + `${bucket}/owner/video/result.mp4`, + "new-owner/video/result.mp4", + ), + ); + await fixture.runtime.runPromise( + regional.listObjects({ prefix: "owner/video/" }), + ); + await fixture.runtime.runPromise( + regional.deleteObjects([{ Key: "owner/video/result.mp4" }]), + ); + expect(send).toHaveBeenCalledTimes(4); + for (const [command] of send.mock.calls) + expect(command.input).toHaveProperty("Bucket", bucket); + } finally { + send.mockRestore(); + await fixture.close(); + } + }, + ); +}); diff --git a/apps/web/__tests__/unit/video-cloudfront.test.ts b/apps/web/__tests__/unit/video-cloudfront.test.ts new file mode 100644 index 00000000000..c74255c8f54 --- /dev/null +++ b/apps/web/__tests__/unit/video-cloudfront.test.ts @@ -0,0 +1,37 @@ +import { describe, expect, it, vi } from "vitest"; + +vi.mock("@cap/env", () => ({ + serverEnv: () => ({ + CAP_CLOUDFRONT_DISTRIBUTION_ID: "EVIRGINIA", + CAP_REGIONAL_UPLOAD_BUCKETS: JSON.stringify({ + "ap-northeast-1": { + bucket: "cap-test-tokyo", + bucketUrl: "https://tokyo.cap.test", + distributionId: "ETOKYO", + }, + "eu-central-1": { + bucket: "cap-test-frankfurt", + bucketUrl: "https://frankfurt.cap.test", + distributionId: "EFRANKFURT", + }, + }), + CAP_REGIONAL_UPLOADS_ENABLED: false, + }), +})); + +import { getVideoCloudFrontDistributionId } from "@/lib/video-cloudfront"; + +describe("video cache invalidation destination", () => { + it.each([ + [null, "EVIRGINIA"], + ["cap-tokyo", "ETOKYO"], + ["cap-frankfurt", "EFRANKFURT"], + ["cap-sydney", undefined], + ["custom-bucket", undefined], + ])( + "uses the persisted bucket even with routing disabled: %s", + (bucket, expected) => { + expect(getVideoCloudFrontDistributionId(bucket)).toBe(expected); + }, + ); +}); diff --git a/apps/web/actions/admin/replace-video.ts b/apps/web/actions/admin/replace-video.ts index 0a6879c9ca9..883d766cdff 100644 --- a/apps/web/actions/admin/replace-video.ts +++ b/apps/web/actions/admin/replace-video.ts @@ -16,6 +16,7 @@ import { Effect } from "effect"; import { retireDesktopRecordingJobForOutputReplacement } from "@/lib/desktop-recording-jobs"; import { MESSENGER_ADMIN_EMAIL } from "@/lib/messenger/constants"; import { runPromise } from "@/lib/server"; +import { getVideoCloudFrontDistributionId } from "@/lib/video-cloudfront"; import { decodeStorageVideo } from "@/lib/video-storage"; async function requireAdmin() { @@ -101,11 +102,7 @@ export async function invalidateVideoCache(videoId: string) { await tx.delete(videoUploads).where(eq(videoUploads.videoId, video.id)); }); - if (video.bucket) { - return; - } - - const distributionId = serverEnv().CAP_CLOUDFRONT_DISTRIBUTION_ID; + const distributionId = getVideoCloudFrontDistributionId(video.bucket); if (!distributionId) { return; } diff --git a/apps/web/app/api/desktop/[...route]/video.ts b/apps/web/app/api/desktop/[...route]/video.ts index 173c1a933b3..cc4790fcf0d 100644 --- a/apps/web/app/api/desktop/[...route]/video.ts +++ b/apps/web/app/api/desktop/[...route]/video.ts @@ -13,7 +13,12 @@ import type { VideoMetadata } from "@cap/database/types"; import { getNewVideoPublic } from "@cap/database/video-sharing-default"; import { serverEnv } from "@cap/env"; import { userIsPro } from "@cap/utils"; -import { makeCurrentUserLayer, Storage, Videos } from "@cap/web-backend"; +import { + makeCurrentUserLayer, + S3Buckets, + Storage, + Videos, +} from "@cap/web-backend"; import { Organisation, Video } from "@cap/web-domain"; import { zValidator } from "@hono/zod-validator"; import { and, count, eq, lte } from "drizzle-orm"; @@ -299,10 +304,29 @@ app.get( ); } - const writable = await (clientSupportsGoogleDriveUpload - ? Storage.getWritableAccessForUser(user.id, videoOrgId) - : Storage.getS3WritableAccessForUser(user.id, videoOrgId) - ).pipe(runPromise); + const writable = await Effect.gen(function* () { + const writable = yield* clientSupportsGoogleDriveUpload + ? Storage.getWritableAccessForUser(user.id, videoOrgId) + : Storage.getS3WritableAccessForUser(user.id, videoOrgId); + if ( + Option.isNone(writable.bucketId) && + Option.isNone(writable.storageIntegrationId) && + !isScreenshot && + (recordingMode === "desktopSegments" || + recordingMode === "desktopMP4") && + process.env.VERCEL === "1" + ) { + const buckets = yield* S3Buckets; + return { + ...writable, + bucketId: buckets.getRegionalUploadBucketId( + c.req.header("x-vercel-ip-latitude"), + c.req.header("x-vercel-ip-longitude"), + ), + }; + } + return writable; + }).pipe(runPromise); await db() .insert(videos) diff --git a/apps/web/lib/content-transfer.ts b/apps/web/lib/content-transfer.ts index 42e2588b575..17b682d9590 100644 --- a/apps/web/lib/content-transfer.ts +++ b/apps/web/lib/content-transfer.ts @@ -1,3 +1,5 @@ +import { S3Bucket } from "@cap/web-domain"; + export const CONTENT_TRANSFER_KIND = "transfer_org_content" as const; export const MAX_CONTENT_TRANSFER_VIDEOS = 10_000; export const MAX_CONTENT_TRANSFER_FOLDERS = 2_000; @@ -290,7 +292,7 @@ export function getContentTransferStorageBlockReason({ } return "The Cap uses a personal storage integration owned by another user"; } - if (bucketId) { + if (!S3Bucket.isCapManagedBucket(bucketId)) { if (!bucketOwnerId) return "The storage bucket is missing"; if ( bucketOrganizationId === organizationId || diff --git a/apps/web/lib/desktop-reupload.ts b/apps/web/lib/desktop-reupload.ts index 7ce70d31452..f328548d9e5 100644 --- a/apps/web/lib/desktop-reupload.ts +++ b/apps/web/lib/desktop-reupload.ts @@ -19,6 +19,7 @@ import { type DesktopReuploadToken, } from "@/lib/desktop-reupload-token"; import { runPromise } from "@/lib/server"; +import { getVideoCloudFrontDistributionId } from "@/lib/video-cloudfront"; type ReuploadedVideo = Pick< Video.Video, @@ -90,12 +91,10 @@ export async function prepareDesktopReupload( } export async function invalidateReuploadedVideo(video: ReuploadedVideo) { - if ( - Option.isSome(video.bucketId) || - Option.isSome(video.storageIntegrationId) - ) - return; - const distributionId = serverEnv().CAP_CLOUDFRONT_DISTRIBUTION_ID; + if (Option.isSome(video.storageIntegrationId)) return; + const distributionId = getVideoCloudFrontDistributionId( + Option.getOrNull(video.bucketId), + ); if (!distributionId) return; const client = new CloudFrontClient({ region: serverEnv().CAP_AWS_REGION || "us-east-1", diff --git a/apps/web/lib/video-cloudfront.ts b/apps/web/lib/video-cloudfront.ts new file mode 100644 index 00000000000..42d8f332492 --- /dev/null +++ b/apps/web/lib/video-cloudfront.ts @@ -0,0 +1,9 @@ +import { serverEnv } from "@cap/env"; +import { parseRegionalBuckets } from "@cap/web-backend/src/S3Buckets/RegionalBuckets"; + +export function getVideoCloudFrontDistributionId(bucketId: string | null) { + if (!bucketId) return serverEnv().CAP_CLOUDFRONT_DISTRIBUTION_ID; + return parseRegionalBuckets(serverEnv().CAP_REGIONAL_UPLOAD_BUCKETS).find( + (bucket) => bucket.id === bucketId, + )?.distributionId; +} diff --git a/apps/web/workflows/admin-reprocess-video.ts b/apps/web/workflows/admin-reprocess-video.ts index 8baa0d6d0fc..7da8cebaad9 100644 --- a/apps/web/workflows/admin-reprocess-video.ts +++ b/apps/web/workflows/admin-reprocess-video.ts @@ -16,6 +16,7 @@ import { createMediaServerCapacityError, isMediaServerCapacityError, } from "@/lib/media-server-backpressure"; +import { getVideoCloudFrontDistributionId } from "@/lib/video-cloudfront"; import { decodeStorageVideo } from "@/lib/video-storage"; import { runWorkflowPromise } from "@/lib/workflow-runtime"; @@ -443,9 +444,7 @@ async function invalidateResultCache( ): Promise { "use step"; - if (bucketId) return; - - const distributionId = serverEnv().CAP_CLOUDFRONT_DISTRIBUTION_ID; + const distributionId = getVideoCloudFrontDistributionId(bucketId); if (!distributionId) return; const basePath = `/${ownerId}/${videoId}`; diff --git a/apps/web/workflows/edit-video.ts b/apps/web/workflows/edit-video.ts index 1f182eb5059..7ace4a35c1e 100644 --- a/apps/web/workflows/edit-video.ts +++ b/apps/web/workflows/edit-video.ts @@ -33,6 +33,7 @@ import { import { decryptEditTranscriptObject } from "@/lib/edit-transcript-storage"; import { startAiGeneration } from "@/lib/generate-ai"; import { transcribeVideo } from "@/lib/transcribe"; +import { getVideoCloudFrontDistributionId } from "@/lib/video-cloudfront"; import { clearFailedEdit, type EditOperation, @@ -723,9 +724,6 @@ async function invalidateEditedVideoCache( ): Promise { "use step"; - const distributionId = serverEnv().CAP_CLOUDFRONT_DISTRIBUTION_ID; - if (!distributionId) return; - const [video] = await db() .select({ ownerId: videos.ownerId, @@ -734,7 +732,9 @@ async function invalidateEditedVideoCache( .from(videos) .where(eq(videos.id, videoId as Video.VideoId)); - if (!video || video.bucket) return; + if (!video) return; + const distributionId = getVideoCloudFrontDistributionId(video.bucket); + if (!distributionId) return; const basePath = `/${video.ownerId}/${videoId}`; const paths = [ diff --git a/packages/env/server.ts b/packages/env/server.ts index 78d649a442f..ace6e634db0 100644 --- a/packages/env/server.ts +++ b/packages/env/server.ts @@ -59,6 +59,8 @@ function createServerEnv() { .optional() .describe("Public URL of the S3 bucket"), CAP_CLOUDFRONT_DISTRIBUTION_ID: z.string().optional(), + CAP_REGIONAL_UPLOADS_ENABLED: boolString(), + CAP_REGIONAL_UPLOAD_BUCKETS: z.string().optional(), CLOUDFRONT_KEYPAIR_ID: z.string().optional(), CLOUDFRONT_KEYPAIR_PRIVATE_KEY: z.string().optional(), diff --git a/packages/web-backend/src/Organisations/index.ts b/packages/web-backend/src/Organisations/index.ts index 55bf96d6665..044d1a67cb0 100644 --- a/packages/web-backend/src/Organisations/index.ts +++ b/packages/web-backend/src/Organisations/index.ts @@ -1,5 +1,5 @@ import * as Db from "@cap/database/schema"; -import { CurrentUser, Organisation, Policy } from "@cap/web-domain"; +import { CurrentUser, Organisation, Policy, S3Bucket } from "@cap/web-domain"; import * as Dz from "drizzle-orm"; import { Effect, Array as EffectArray, Option } from "effect"; import { Database } from "../Database"; @@ -103,16 +103,18 @@ export class Organisations extends Effect.Service()( .where(Dz.eq(Db.videos.orgId, id)), ); const capManagedVideos = videos.filter( - (video) => !video.bucket && !video.storageIntegrationId, + (video) => + S3Bucket.isCapManagedBucket(video.bucket) && + !video.storageIntegrationId, ); const [defaultBucket] = yield* s3Buckets.getBucketAccess(Option.none()); - const deleteS3Prefix = (prefix: string) => + const deleteS3Prefix = (prefix: string, bucket = defaultBucket) => Effect.gen(function* () { let continuationToken: string | undefined; do { - const listedObjects = yield* defaultBucket.listObjects({ + const listedObjects = yield* bucket.listObjects({ prefix, continuationToken, }); @@ -126,7 +128,7 @@ export class Organisations extends Effect.Service()( index < objects.length; index += s3DeleteBatchSize ) { - yield* defaultBucket.deleteObjects( + yield* bucket.deleteObjects( objects.slice(index, index + s3DeleteBatchSize), ); } @@ -139,7 +141,13 @@ export class Organisations extends Effect.Service()( yield* Effect.forEach( capManagedVideos, - (video) => deleteS3Prefix(`${video.ownerId}/${video.id}/`), + (video) => + Effect.gen(function* () { + const [bucket] = yield* s3Buckets.getBucketAccess( + Option.fromNullable(video.bucket), + ); + yield* deleteS3Prefix(`${video.ownerId}/${video.id}/`, bucket); + }), { concurrency: 3 }, ); yield* deleteS3Prefix(`organizations/${id}/`); diff --git a/packages/web-backend/src/S3Buckets/README.md b/packages/web-backend/src/S3Buckets/README.md new file mode 100644 index 00000000000..15c265d6e8f --- /dev/null +++ b/packages/web-backend/src/S3Buckets/README.md @@ -0,0 +1,57 @@ +# Regional Instant uploads + +Disabled by default. On Vercel, new desktop Instant recordings use the request's +`x-vercel-ip-latitude` and `x-vercel-ip-longitude` headers to choose the geographically +nearest configured region, including the existing Virginia destination. Missing or +invalid location, disabled routing, or no usable regional configuration keeps the +original path. Selection is local arithmetic: no geolocation request or extra database +query. Custom storage and Google Drive retain priority. + +Supported destinations are Virginia (existing default) plus 12 new regions: +Oregon, Montreal, London, Paris, Frankfurt, Stockholm, UAE, Mumbai, Malaysia, +Singapore, Tokyo, and Sydney. São Paulo and Cape Town are excluded for cost. +Only provisioned and configured regions participate. Geographic proximity does not +guarantee the fastest network route; validate each region before enabling it. + +The destination is stored on the recording, not the user. Travel affects the next new +recording; retries, processing, playback, edits, transfers, and deletion keep the +original destination. Reserved `cap-*` IDs fit `videos.bucket` and cannot collide with +generated customer IDs. No schema migration or desktop change is needed. + +Provision a private S3 bucket and CloudFront distribution for each enabled region. +Match Virginia's CloudFront cache, CORS, and signed ZIP-download behavior; +use the existing signing key group and S3 origin access control, and grant the +server/worker AWS identity bucket access and distribution invalidation. Keep direct S3 access private. Playback URLs +and relative HLS segments must retain the existing viewer-access behavior. +Enable opt-in AWS regions (such as UAE and Malaysia) in the account first. Configure every +web and workflow deployment with `CAP_REGIONAL_UPLOAD_BUCKETS`, a JSON object keyed +by supported AWS region. For example: + +```json +{ + "ap-northeast-1": { + "bucket": "your-tokyo-bucket", + "bucketUrl": "https://your-tokyo-cdn.example.com", + "distributionId": "YOUR_TOKYO_DISTRIBUTION_ID" + }, + "eu-central-1": { + "bucket": "your-frankfurt-bucket", + "bucketUrl": "https://your-frankfurt-cdn.example.com", + "distributionId": "YOUR_FRANKFURT_DISTRIBUTION_ID" + } +} +``` + +Use HTTPS CDN origins without a trailing slash or path. Retain the existing default +AWS/CloudFront configuration (`CAP_AWS_REGION=us-east-1`). After validating the +configured destinations, set `CAP_REGIONAL_UPLOADS_ENABLED=true`. Roll back selection +by setting it to `false`; retain regional bucket configuration while recordings refer +to it, and never repoint a region's bucket name. Missing configuration for an existing +regional recording fails explicitly rather than splitting its objects across regions. + +The Japan benchmark improved upload completion but used direct S3 playback and was +slower to first playback/final MP4 than accelerated Virginia with CDN. Before enabling, +repeat the GPUI benchmark with the configured CDN path and verify preparation, +playback, finalization, replacement, and deletion. Additional buckets do not replicate +recordings, but regional storage rates and cross-region processing transfers affect +cost. This change does not provision infrastructure or enable production routing. diff --git a/packages/web-backend/src/S3Buckets/RegionalBuckets.ts b/packages/web-backend/src/S3Buckets/RegionalBuckets.ts new file mode 100644 index 00000000000..1a498827638 --- /dev/null +++ b/packages/web-backend/src/S3Buckets/RegionalBuckets.ts @@ -0,0 +1,69 @@ +import { S3Bucket } from "@cap/web-domain"; +import { Option, Schema } from "effect"; + +const BucketConfig = Schema.Struct({ + bucket: Schema.String.pipe( + Schema.pattern(/^[a-z0-9][a-z0-9-]{1,61}[a-z0-9]$/), + ), + bucketUrl: Schema.String.pipe( + Schema.filter((value) => { + try { + const url = new URL(value); + return url.protocol === "https:" && url.origin === value; + } catch { + return false; + } + }), + ), + distributionId: Schema.NonEmptyTrimmedString, +}); +const decodeConfig = Schema.decodeUnknownOption( + Schema.parseJson( + Schema.Record({ key: Schema.String, value: Schema.Unknown }), + ), +); + +const decodeBucket = Schema.decodeUnknownOption(BucketConfig); + +export function parseRegionalBuckets(value: string | undefined) { + const config = Option.getOrUndefined(decodeConfig(value ?? "")); + return S3Bucket.RegionalBuckets.flatMap((region) => { + const bucket = Option.getOrUndefined(decodeBucket(config?.[region.region])); + return bucket ? [{ ...region, ...bucket }] : []; + }); +} + +export function getNearestRegionalBucket( + latitude: string | undefined, + longitude: string | undefined, + buckets: ReadonlyArray<(typeof S3Bucket.RegionalBuckets)[number]>, +) { + const lat = Number(latitude); + const lon = Number(longitude); + if ( + !latitude?.trim() || + !longitude?.trim() || + !Number.isFinite(lat) || + !Number.isFinite(lon) || + Math.abs(lat) > 90 || + Math.abs(lon) > 180 + ) + return Option.none(); + + const radians = Math.PI / 180; + const distance = (targetLat: number, targetLon: number) => + Math.sin(((targetLat - lat) * radians) / 2) ** 2 + + Math.cos(lat * radians) * + Math.cos(targetLat * radians) * + Math.sin(((targetLon - lon) * radians) / 2) ** 2; + let nearest = Option.none(); + let nearestDistance = distance(38.13, -78.45); + for (const bucket of buckets) { + const candidate = distance(bucket.latitude, bucket.longitude); + if (candidate < nearestDistance) { + nearest = Option.some(bucket.id); + nearestDistance = candidate; + } + } + return nearest; +} diff --git a/packages/web-backend/src/S3Buckets/index.ts b/packages/web-backend/src/S3Buckets/index.ts index 49ecb912ecc..d14c45de5e9 100644 --- a/packages/web-backend/src/S3Buckets/index.ts +++ b/packages/web-backend/src/S3Buckets/index.ts @@ -1,12 +1,16 @@ import * as S3 from "@aws-sdk/client-s3"; import * as CloudFrontPresigner from "@aws-sdk/cloudfront-signer"; import { decrypt } from "@cap/database/crypto"; -import type { Organisation, S3Bucket, User } from "@cap/web-domain"; +import { type Organisation, S3Bucket, type User } from "@cap/web-domain"; import type { RequestPresigningArguments } from "@smithy/types"; import { Config, Effect, Layer, Option } from "effect"; import { AwsCredentials } from "../Aws.ts"; import { Database } from "../Database.ts"; +import { + getNearestRegionalBucket, + parseRegionalBuckets, +} from "./RegionalBuckets.ts"; import { createS3BucketAccess } from "./S3BucketAccess.ts"; import { S3BucketClientProvider } from "./S3BucketClientProvider.ts"; import { S3BucketsRepo } from "./S3BucketsRepo.ts"; @@ -100,47 +104,89 @@ export class S3Buckets extends Effect.Service()("S3Buckets", { Effect.map(Option.fromNullable), ); - const cloudfrontBucketAccess = cloudfrontEnvs.pipe( - Option.map((cloudfrontEnvs) => - Effect.flatMap(createS3BucketAccess, (s3) => { - const getCloudFrontSignedUrl = ( - key: string, - signingArgs?: RequestPresigningArguments, - ) => { - const url = `${cloudfrontEnvs.bucketUrl}/${key}`; - const expiresIn = signingArgs?.expiresIn ?? 3600; - const expires = Math.floor((Date.now() + expiresIn * 1000) / 1000); - - const policy = { - Statement: [ - { - Resource: url, - Condition: { - DateLessThan: { - "AWS:EpochTime": Math.floor(expires), + const cloudfrontBucketAccess = (bucketUrl?: string) => + cloudfrontEnvs.pipe( + Option.map((cloudfrontEnvs) => + Effect.flatMap(createS3BucketAccess, (s3) => { + const getCloudFrontSignedUrl = ( + key: string, + signingArgs?: RequestPresigningArguments, + ) => { + const url = `${bucketUrl ?? cloudfrontEnvs.bucketUrl}/${key}`; + const expiresIn = signingArgs?.expiresIn ?? 3600; + const expires = Math.floor( + (Date.now() + expiresIn * 1000) / 1000, + ); + + const policy = { + Statement: [ + { + Resource: url, + Condition: { + DateLessThan: { + "AWS:EpochTime": Math.floor(expires), + }, }, }, - }, - ], + ], + }; + + return Effect.succeed( + CloudFrontPresigner.getSignedUrl({ + url, + keyPairId: cloudfrontEnvs.keypairId, + privateKey: cloudfrontEnvs.privateKey, + policy: JSON.stringify(policy), + }), + ); }; - return Effect.succeed( - CloudFrontPresigner.getSignedUrl({ - url, - keyPairId: cloudfrontEnvs.keypairId, - privateKey: cloudfrontEnvs.privateKey, - policy: JSON.stringify(policy), - }), - ); - }; + return Effect.succeed({ + ...s3, + getSignedObjectUrl: getCloudFrontSignedUrl, + }); + }), + ), + ); - return Effect.succeed({ - ...s3, - getSignedObjectUrl: getCloudFrontSignedUrl, - }); - }), + const regionalConfigs = parseRegionalBuckets( + Option.getOrUndefined( + yield* Config.string("CAP_REGIONAL_UPLOAD_BUCKETS").pipe(Config.option), ), ); + const regionalUploadsEnabled = yield* Config.boolean( + "CAP_REGIONAL_UPLOADS_ENABLED", + ).pipe(Effect.orElseSucceed(() => false)); + const regionalBucketAccess = new Map( + regionalConfigs.flatMap((config) => { + const access = cloudfrontBucketAccess(config.bucketUrl); + if (Option.isNone(access)) return []; + const client = new S3.S3Client({ + region: config.region, + credentials, + forcePathStyle: false, + requestHandler, + }); + return [ + [ + config.id, + access.value.pipe( + Effect.provide( + Layer.succeed(S3BucketClientProvider, { + getInternal: Effect.succeed(client), + getPublic: Effect.succeed(client), + bucket: config.bucket, + isPathStyle: false, + }), + ), + ), + ] as const, + ]; + }), + ); + const uploadRegions = regionalConfigs.filter((config) => + regionalBucketAccess.has(config.id), + ); const getBucketAccess = Effect.fn("S3Buckets.getProviderLayer")(function* ( customBucket: Option.Option, @@ -154,7 +200,7 @@ export class S3Buckets extends Effect.Service()("S3Buckets", { isPathStyle: defaultConfigs.forcePathStyle, }); - return Option.match(cloudfrontBucketAccess, { + return Option.match(cloudfrontBucketAccess(), { onSome: (access) => access, onNone: () => createS3BucketAccess, }).pipe(Effect.provide(provider)); @@ -183,9 +229,31 @@ export class S3Buckets extends Effect.Service()("S3Buckets", { }); return { + getRegionalUploadBucketId: ( + latitude: string | undefined, + longitude: string | undefined, + ) => + regionalUploadsEnabled && defaultConfigs.region === "us-east-1" + ? getNearestRegionalBucket(latitude, longitude, uploadRegions) + : Option.none(), getBucketAccess: Effect.fn("S3Buckets.getBucketAccess")(function* ( bucketId?: Option.Option, ) { + const regionalBucket = S3Bucket.getRegionalBucket( + Option.getOrNull(bucketId ?? Option.none()), + ); + if (regionalBucket) { + const access = regionalBucketAccess.get(regionalBucket.id); + if (!access) + return yield* Effect.fail( + new S3Bucket.S3Error({ + cause: new Error( + `Regional storage configuration is missing or invalid: ${regionalBucket.region}`, + ), + }), + ); + return [yield* access, Option.none()] as const; + } const customBucket = yield* (bucketId ?? Option.none()).pipe( Option.map(repo.getById), Effect.transposeOption, diff --git a/packages/web-domain/src/S3Bucket.ts b/packages/web-domain/src/S3Bucket.ts index 60f264f3ce2..f13a241ca45 100644 --- a/packages/web-domain/src/S3Bucket.ts +++ b/packages/web-domain/src/S3Bucket.ts @@ -4,6 +4,88 @@ import { UserId } from "./User.ts"; export const S3BucketId = Schema.String.pipe(Schema.brand("S3BucketId")); export type S3BucketId = typeof S3BucketId.Type; +// Reserved IDs fit videos.bucket and cannot collide with customer IDs (no hyphens). +export const RegionalBuckets = [ + { + region: "us-west-2", + id: S3BucketId.make("cap-oregon"), + latitude: 45.84, + longitude: -119.7, + }, + { + region: "ca-central-1", + id: S3BucketId.make("cap-montreal"), + latitude: 45.5, + longitude: -73.57, + }, + { + region: "eu-west-2", + id: S3BucketId.make("cap-london"), + latitude: 51.51, + longitude: -0.13, + }, + { + region: "eu-west-3", + id: S3BucketId.make("cap-paris"), + latitude: 48.86, + longitude: 2.35, + }, + { + region: "eu-central-1", + id: S3BucketId.make("cap-frankfurt"), + latitude: 50.11, + longitude: 8.68, + }, + { + region: "eu-north-1", + id: S3BucketId.make("cap-stockholm"), + latitude: 59.33, + longitude: 18.07, + }, + { + region: "me-central-1", + id: S3BucketId.make("cap-uae"), + latitude: 24.45, + longitude: 54.38, + }, + { + region: "ap-south-1", + id: S3BucketId.make("cap-mumbai"), + latitude: 19.08, + longitude: 72.88, + }, + { + region: "ap-southeast-1", + id: S3BucketId.make("cap-singapore"), + latitude: 1.35, + longitude: 103.82, + }, + { + region: "ap-southeast-5", + id: S3BucketId.make("cap-malaysia"), + latitude: 3.14, + longitude: 101.69, + }, + { + region: "ap-northeast-1", + id: S3BucketId.make("cap-tokyo"), + latitude: 35.68, + longitude: 139.69, + }, + { + region: "ap-southeast-2", + id: S3BucketId.make("cap-sydney"), + latitude: -33.87, + longitude: 151.21, + }, +] as const; + +export const getRegionalBucket = (id: string | null | undefined) => + RegionalBuckets.find((bucket) => bucket.id === id); + +export const isCapManagedBucket = (id: string | null | undefined) => + !id || getRegionalBucket(id) !== undefined; + export class S3Bucket extends Schema.Class("S3Bucket")({ id: S3BucketId, ownerId: UserId,