diff --git a/.quixos-subtree-source.json b/.quixos-subtree-source.json index 40850d4..a072eef 100644 --- a/.quixos-subtree-source.json +++ b/.quixos-subtree-source.json @@ -1,7 +1,7 @@ { "version": 1, "sourceRepo": "https://gitea-external.egads.tutti.syntaxblitz.net/quixos/quixos.git", - "sourceCommit": "ebf8fe2843731ad538caabac1e2780fe240d35d7", + "sourceCommit": "ddccb886544038ee85e4f1fd4db1c35b24869199", "sourcePath": "quixos-instance/packages/camino-package-runtime", "exportName": "camino-package-runtime", "mirrorRemote": "https://gitea-external.egads.tutti.syntaxblitz.net/quixos/camino-package-runtime.git" diff --git a/src/index.ts b/src/index.ts index a0f08b1..9ede950 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,4 +1,5 @@ import http from "node:http"; +import { AsyncLocalStorage } from "node:async_hooks"; import { create } from "@bufbuild/protobuf"; import { createClient, type Client } from "@connectrpc/connect"; import { @@ -33,6 +34,79 @@ export type CaminoClient = Client; export type { InvokeRequest, Value }; +export type DerivedDependency = + | { + kind: "object"; + objectId: string; + } + | { + kind: "field"; + objectId: string; + fieldName: string; + } + | { + kind: "edges"; + objectId: string; + sourceField?: string; + }; + +type DerivedDependencyTracker = { + record(dependency: DerivedDependency): void; + dependencies(): DerivedDependency[]; +}; + +const dependencyKey = (dependency: DerivedDependency) => { + switch (dependency.kind) { + case "object": + return `object:${dependency.objectId}`; + case "field": + return `field:${dependency.objectId}:${dependency.fieldName}`; + case "edges": + return `edges:${dependency.objectId}:${dependency.sourceField ?? ""}`; + } +}; + +const createDependencyTracker = (): DerivedDependencyTracker => { + const dependencies = new Map(); + return { + record(dependency) { + dependencies.set(dependencyKey(dependency), dependency); + }, + dependencies() { + return [...dependencies.values()]; + }, + }; +}; + +const dependencyStorage = new AsyncLocalStorage(); + +export const recordDerivedDependency = (dependency: DerivedDependency) => { + dependencyStorage.getStore()?.record(dependency); +}; + +export const getActiveDerivedDependencies = (): DerivedDependency[] => + dependencyStorage.getStore()?.dependencies() ?? []; + +export const runWithDerivedDependencyTracking = async ( + fn: () => Promise | T, +): Promise<{ value: T; dependencies: DerivedDependency[] }> => { + const tracker = createDependencyTracker(); + const value = await dependencyStorage.run(tracker, fn); + return { + value, + dependencies: tracker.dependencies(), + }; +}; + +const logDerivedDependencies = (dependencies: DerivedDependency[]) => { + if (process.env.QUIXOS_DERIVED_DEPS_LOG !== "1") { + return; + } + process.stderr.write( + `[derived-deps] ${JSON.stringify({ dependencies })}\n`, + ); +}; + export type Field = { get(): Promise; set(value: T): Promise; @@ -55,7 +129,16 @@ export const fieldOps = ( export const derivedWithDeps = (operation: { get: RuntimeHandler; -}): FieldOperation => fieldOps(operation); +}): FieldOperation => + fieldOps({ + async get(context) { + const { value, dependencies } = await runWithDerivedDependencyTracking( + () => operation.get(context), + ); + logDerivedDependencies(dependencies); + return value; + }, + }); export type RuntimeFunctionSpec = { exportName: ExportName; @@ -221,6 +304,7 @@ export const createField = ( fieldName: string, ): Field => ({ async get() { + recordDerivedDependency({ kind: "field", objectId, fieldName }); const response = await camino.getObject({ objectId }); return protoValueToJs(response.object?.fields[fieldName]) as T | undefined; }, @@ -233,6 +317,51 @@ export const createField = ( }, }); +const createDependencyTrackingCaminoClient = ( + camino: CaminoClient, +): CaminoClient => + new Proxy(camino, { + get(target, property, receiver) { + if (property === "getObject") { + return async (request: { objectId: string }, ...args: unknown[]) => { + if (request.objectId) { + recordDerivedDependency({ + kind: "object", + objectId: request.objectId, + }); + } + return await ( + Reflect.get(target, property, receiver) as ( + request: { objectId: string }, + ...args: unknown[] + ) => Promise + )(request, ...args); + }; + } + if (property === "listEdges") { + return async ( + request: { objectId: string; sourceField?: string }, + ...args: unknown[] + ) => { + if (request.objectId) { + recordDerivedDependency({ + kind: "edges", + objectId: request.objectId, + ...(request.sourceField ? { sourceField: request.sourceField } : {}), + }); + } + return await ( + Reflect.get(target, property, receiver) as ( + request: { objectId: string; sourceField?: string }, + ...args: unknown[] + ) => Promise + )(request, ...args); + }; + } + return Reflect.get(target, property, receiver); + }, + }); + export const serveQuixosPackageRuntime = < Context, Spec extends readonly RuntimeFunctionSpec[], @@ -255,6 +384,7 @@ export const serveQuixosPackageRuntime = < CaminoService, createConnectTransport({ baseUrl: caminoUrl, httpVersion: "1.1" }), ); + const trackingCamino = createDependencyTrackingCaminoClient(camino); if (!Number.isInteger(port) || port <= 0) { throw new Error("QUIXOS_RUNTIME_PORT must be a positive integer"); @@ -325,7 +455,7 @@ export const serveQuixosPackageRuntime = < return create(InvokeResponseSchema, { ok: true, result: jsToProtoValue( - await runtimeFunction(params.createContext(camino, request)), + await runtimeFunction(params.createContext(trackingCamino, request)), ), }); } catch (error) {