Track dependencies for derived runtime reads
This commit is contained in:
+132
-2
@@ -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<typeof CaminoService>;
|
||||
|
||||
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<string, DerivedDependency>();
|
||||
return {
|
||||
record(dependency) {
|
||||
dependencies.set(dependencyKey(dependency), dependency);
|
||||
},
|
||||
dependencies() {
|
||||
return [...dependencies.values()];
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
const dependencyStorage = new AsyncLocalStorage<DerivedDependencyTracker>();
|
||||
|
||||
export const recordDerivedDependency = (dependency: DerivedDependency) => {
|
||||
dependencyStorage.getStore()?.record(dependency);
|
||||
};
|
||||
|
||||
export const getActiveDerivedDependencies = (): DerivedDependency[] =>
|
||||
dependencyStorage.getStore()?.dependencies() ?? [];
|
||||
|
||||
export const runWithDerivedDependencyTracking = async <T>(
|
||||
fn: () => Promise<T> | 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<T> = {
|
||||
get(): Promise<T | undefined>;
|
||||
set(value: T): Promise<void>;
|
||||
@@ -55,7 +129,16 @@ export const fieldOps = <Context>(
|
||||
|
||||
export const derivedWithDeps = <Context>(operation: {
|
||||
get: RuntimeHandler<Context>;
|
||||
}): FieldOperation<Context> => fieldOps(operation);
|
||||
}): FieldOperation<Context> =>
|
||||
fieldOps({
|
||||
async get(context) {
|
||||
const { value, dependencies } = await runWithDerivedDependencyTracking(
|
||||
() => operation.get(context),
|
||||
);
|
||||
logDerivedDependencies(dependencies);
|
||||
return value;
|
||||
},
|
||||
});
|
||||
|
||||
export type RuntimeFunctionSpec<ExportName extends string = string> = {
|
||||
exportName: ExportName;
|
||||
@@ -221,6 +304,7 @@ export const createField = <T>(
|
||||
fieldName: string,
|
||||
): Field<T> => ({
|
||||
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 = <T>(
|
||||
},
|
||||
});
|
||||
|
||||
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<unknown>
|
||||
)(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<unknown>
|
||||
)(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) {
|
||||
|
||||
Reference in New Issue
Block a user