/* eslint-disable turbo/no-undeclared-env-vars */ import { createEvaluateScriptLocally, createEvaluateScriptRemotely, createFetchFileVersionEnv, createPipelineTaskRunner, PipelineGraphAPI, PipelineState, RemoteEvaluateParameters, } from "@noya-app/noya-pipeline"; import { Observable } from "@noya-app/observable"; import { TaskSnapshot } from "@noya-app/task-runner"; import type { PipelineStatus } from "./multiplayer"; import { RPCManager } from "./rpcManager"; export type TaskArtifact = { name: string; size: number; url: string }; export type TaskMetadata = { artifacts: TaskArtifact[]; }; export class PipelineManager { constructor(public rpcManager: RPCManager) {} status$ = new Observable("idle"); tasks$ = new Observable[]>([]); result$ = new Observable(undefined); async run(parameters: RemoteEvaluateParameters) { const response = await this.rpcManager.requestRoute("POST /api/evaluate", { body: JSON.stringify(parameters), headers: { "Content-Type": "application/json", }, }); return JSON.parse(response.body).result; } async runScript( scriptUrl: string, parameters: Pick ) { const response = await fetch(scriptUrl); if (!response.ok) { throw new Error("Failed to fetch script"); } const scriptText = await response.text(); return this.run({ type: "code", code: scriptText, args: parameters.args, env: parameters.env, }); } async runScriptOnFileState(scriptUrl: string, state: unknown) { return await this.runScript(scriptUrl, { args: [{ state }], env: {}, }); } async runPipeline(pipelineState?: PipelineState) { const response = await this.rpcManager.requestRoute( "POST /api/runWorkflow", { body: JSON.stringify({ ...(pipelineState && { pipelineState }), }), headers: { "Content-Type": "application/json", }, } ); const { result, snapshot } = JSON.parse(response.body); return { result, snapshot }; } getPipelineGraphAPI( options: { getFileVersionData: PipelineGraphAPI["getFileVersionData"]; getFileVersionInputs: PipelineGraphAPI["getFileVersionInputs"]; getFileVersionTransforms: PipelineGraphAPI["getFileVersionTransforms"]; getFileVersionIds: PipelineGraphAPI["getFileVersionIds"]; getFileSecret: PipelineGraphAPI["getFileSecret"]; getFileVersionAssets: PipelineGraphAPI["getFileVersionAssets"]; getFileAssets: PipelineGraphAPI["getFileAssets"]; createFileVersion: PipelineGraphAPI["createFileVersion"]; } & ( | { executionEnvironment: "local"; } | { executionEnvironment: "remote"; baseUrl: string; } ) ): PipelineGraphAPI { return { evaluateScript: options.executionEnvironment === "local" ? createEvaluateScriptLocally() : createEvaluateScriptRemotely({ url: process.env.NEXT_PUBLIC_DEV_EVALUATE_SCRIPT_URL!, bundlerUrl: process.env.NEXT_PUBLIC_DEV_EVALUATE_SCRIPT_URL!, }), getFileVersionEnv: options.executionEnvironment === "local" ? () => Promise.resolve({}) : createFetchFileVersionEnv({ baseUrl: options.baseUrl, ENCRYPTION_KEY: process.env.NEXT_PUBLIC_DEV_ENCRYPTION_KEY!, MULTIPLAYER_SHARED_TOKEN: process.env.NEXT_PUBLIC_DEV_MULTIPLAYER_SHARED_TOKEN!, }), getFileVersionData: options.getFileVersionData, getFileVersionInputs: options.getFileVersionInputs, getFileVersionTransforms: options.getFileVersionTransforms, getFileVersionIds: options.getFileVersionIds, getFileSecret: options.getFileSecret, getFileVersionAssets: options.getFileVersionAssets, getFileAssets: options.getFileAssets, createArtifact: options.executionEnvironment === "local" ? async (uint8Array) => { const blob = new Blob([uint8Array]); return URL.createObjectURL(blob); } : undefined, createFileVersion: options.createFileVersion, }; } async runPipelineLocally( pipelineState: PipelineState, options: Parameters[0] ) { const { taskRunner, rootTask } = await createPipelineTaskRunner({ state: pipelineState, api: this.getPipelineGraphAPI(options), }); this.status$.set("running"); taskRunner.subscribe(() => { this.tasks$.set(taskRunner.getSnapshot()); }); try { const result = await taskRunner.runTask(rootTask); this.status$.set("success"); this.result$.set(result); return { result, snapshot: taskRunner.getSnapshot() }; } catch (error) { this.status$.set("error"); return { error, snapshot: taskRunner.getSnapshot() }; } } }