import { applyExtendedPatch, AssetUploadProgress, ClientToServerMessage, ExtendedPatch, ExtractRequestBody, ExtractResponseBody, } from "@noya-app/noya-multiplayer-react"; import { CreateResourceParametersWithoutAccessibleByFileId, DeleteResourceParameters, Resource, UploadableAsset, } from "@noya-app/noya-schemas"; import { isDeepEqual, uuid } from "@noya-app/noya-utils"; import { Observable } from "@noya-app/observable"; import type { ConnectionEvent } from "./ConnectionEventManager"; import { ResourceEditServerMessage, ResourceEditSession, } from "./ResourceEditSession"; import { RPCManager } from "./rpcManager"; export type ResourceManagerOptions = { initialResources?: Resource[]; // Send a message over the active WebSocket connection sendMessage?: (message: ClientToServerMessage) => void; // Subscribe to server messages coming over WebSocket subscribeConnectionEvents?: ( handler: (event: ConnectionEvent) => void ) => () => void; // Resolve an asset id or stableId to a real id resolveAssetId?: (idOrStableId: string) => string; }; export class ResourceManager { isInitialized$ = new Observable(false); resources$ = new Observable([], { isEqual: isDeepEqual }); patches$ = new Observable<{ id: string; patches: ExtendedPatch[] }[]>([], { isEqual: isDeepEqual, }); optimisticResources$ = Observable.combine( [this.resources$, this.patches$], ([resources, patches]) => { return patches.reduce((acc, p) => { try { return applyExtendedPatch(acc, p.patches); } catch (error) { return acc; } }, resources); }, { isEqual: isDeepEqual } ); uploads$ = new Observable>({ // example: { // id: "example", // name: "example.png", // percent: 50, // loaded: 1024 * 1024 * 2, // total: 1024 * 1024 * 5, // contentType: "image/png", // }, // example2: { // id: "example2", // name: "example2.png", // percent: 50, // loaded: 1024 * 1024 * 2, // total: 1024 * 1024 * 5, // contentType: "image/png", // }, }); constructor( public rpcManager: RPCManager, options: ResourceManagerOptions = {} ) { this.setResources(options.initialResources ?? []); this.isInitialized$.set(options.initialResources ? true : false); this._sendMessage = options.sendMessage; if (options.subscribeConnectionEvents) { this._unsubscribeConnectionEvents = options.subscribeConnectionEvents( this._handleConnectionEvent ); } this._resolveAssetId = options.resolveAssetId; } setResources(resources: Resource[]) { const resourcesWithUrls = resources.map((resource) => ({ ...resource, url: resource.url ?? `/api/transform/${resource.id}?format=image`, })); this.resources$.set(resourcesWithUrls); } async createResource< Result extends | CreateResourceParametersWithoutAccessibleByFileId | CreateResourceParametersWithoutAccessibleByFileId[], >( parameters: Result ): Promise { const payload: CreateResourceParametersWithoutAccessibleByFileId[] = Array.isArray(parameters) ? parameters : [parameters]; if (payload.length === 0) { throw new Error("No resources to create"); } const patchResult = await this.patchResources({ create: payload, }); const resources = patchResult.created; const result = Array.isArray(parameters) ? resources : resources[0]; return result as Result extends unknown[] ? Resource[] : Resource; } async _withPatch( patches: ExtendedPatch[], callback: () => Promise ): Promise { const patchId = uuid(); this.patches$.set([...this.patches$.get(), { id: patchId, patches }]); try { const result = await callback(); return result; } finally { this.patches$.set(this.patches$.get().filter((p) => p.id !== patchId)); } } async deleteResource< Result extends DeleteResourceParameters | DeleteResourceParameters[], >( parameters: Result ): Promise { const payload: DeleteResourceParameters[] = Array.isArray(parameters) ? parameters : [parameters]; const patchResult = await this.patchResources({ delete: payload, }); const deletedResources = patchResult.deleted; const finalResult = Array.isArray(parameters) ? deletedResources : deletedResources[0]; return finalResult as Result extends unknown[] ? Resource[] : Resource; } async patchResources({ create, update, delete: deleteParameters, }: ExtractRequestBody<"PATCH /api/resources">): Promise< ExtractResponseBody<"PATCH /api/resources"> > { const hasCreate = (create?.length ?? 0) > 0; const hasUpdate = (update?.length ?? 0) > 0; const hasDelete = (deleteParameters?.length ?? 0) > 0; if (!hasCreate && !hasUpdate && !hasDelete) { throw new Error( "PATCH body must include create, update, or delete arrays" ); } const resolveAssetId = (assetId: string) => this._resolveAssetId?.(assetId) ?? assetId; const resolvedCreate = hasCreate ? create!.map((payload) => payload.type === "asset" && payload.assetId ? { ...payload, assetId: resolveAssetId(payload.assetId) } : payload ) : undefined; const resolvedUpdate = hasUpdate ? update!.map((payload) => ({ ...payload, ...(payload.assetId ? { assetId: resolveAssetId(payload.assetId) } : {}), })) : undefined; const patches: ExtendedPatch[] = []; if (deleteParameters?.length) { patches.push( ...deleteParameters.map( ({ id }) => ({ op: "remove", path: [{ id }], }) satisfies ExtendedPatch ) ); } if (resolvedUpdate?.length) { for (const item of resolvedUpdate) { if (item.path !== undefined) { patches.push({ op: "replace", path: [{ id: item.id }, "path"], value: item.path, }); } if (item.type !== undefined) { patches.push({ op: "replace", path: [{ id: item.id }, "type"], value: item.type, }); } if (item.assetId !== undefined) { patches.push({ op: "replace", path: [{ id: item.id }, "assetId"], value: item.assetId, }); } if (item.fileId !== undefined) { patches.push({ op: "replace", path: [{ id: item.id }, "fileId"], value: item.fileId, }); } if (item.resourceId !== undefined) { patches.push({ op: "replace", path: [{ id: item.id }, "resourceId"], value: item.resourceId, }); } } } const executePatch = async () => { const body: ExtractRequestBody<"PATCH /api/resources"> = {}; if (resolvedCreate?.length) { body.create = resolvedCreate; } if (resolvedUpdate?.length) { body.update = resolvedUpdate; } if (deleteParameters?.length) { body.delete = deleteParameters; } const uploadEntries = [ ...(resolvedCreate?.flatMap((payload) => payload.type === "asset" && payload.asset ? [{ name: payload.path }] : [] ) ?? []), ...(resolvedUpdate?.flatMap((payload) => payload.asset ? [{ name: payload.path ?? payload.id }] : [] ) ?? []), ].map((entry) => ({ id: uuid(), name: entry.name })); try { const response = await this.rpcManager.requestRoute( "PATCH /api/resources", { body: JSON.stringify(body), headers: { "Content-Type": "application/json" }, ...(uploadEntries.length ? { onUploadProgress: (progress) => { const entries = uploadEntries.map((entry) => [ entry.id, { ...progress, name: entry.name, id: entry.id, } as AssetUploadProgress, ]) as [string, AssetUploadProgress][]; this.uploads$.set({ ...this.uploads$.get(), ...Object.fromEntries(entries), }); }, } : {}), } ); return this.rpcManager.getResponseBody( response ) as ExtractResponseBody<"PATCH /api/resources">; } finally { if (uploadEntries.length) { const uploads = { ...this.uploads$.get() }; for (const entry of uploadEntries) { delete uploads[entry.id]; } this.uploads$.set(uploads); } } }; if (patches.length > 0) { return this._withPatch(patches, executePatch); } return executePatch(); } async updateResource({ id, path, assetId, asset, type, fileId, fileVersionId, resourceId, }: { id: string; path?: string; assetId?: string; asset?: UploadableAsset; type?: Resource["type"]; fileId?: string; fileVersionId?: string; resourceId?: string; }) { const result = await this.patchResources({ update: [ { id, ...(path ? { path } : {}), ...(assetId ? { assetId } : {}), ...(asset ? { asset } : {}), ...(type ? { type } : {}), ...(fileId ? { fileId } : {}), ...(fileVersionId ? { fileVersionId } : {}), ...(resourceId ? { resourceId } : {}), }, ], }); return result.updated[0]; } async listResources() { const response = await this.rpcManager.requestRoute("GET /api/resources"); return this.rpcManager.getResponseBody(response); } // --- Editing handle API --- _sendMessage?: (message: ClientToServerMessage) => void; _unsubscribeConnectionEvents?: () => void; _resolveAssetId?: (idOrStableId: string) => string; // resourceId -> session _editSessions = new Map(); _handleConnectionEvent = (event: ConnectionEvent) => { if (event.type !== "receive") return; const message = event.message as ResourceEditServerMessage; switch (message.type) { case "resourceEdit.object": case "resourceEdit.acceptPatch": case "resourceEdit.rejectPatch": { const session = this._editSessions.get(message.resourceId); session?.handleServerMessage(message); break; } } }; openEditingHandle(resourceId: string) { let session = this._editSessions.get(resourceId); if (!session) { session = new ResourceEditSession({ resourceId, sendMessage: (m) => this._sendMessage?.(m), onClose: () => this._editSessions.delete(resourceId), }); this._editSessions.set(resourceId, session); } session.open(); return session.getHandle(); } }