diff --git a/dist/bindings.d.ts b/dist/bindings.d.ts new file mode 100644 index 0000000..15a8f91 --- /dev/null +++ b/dist/bindings.d.ts @@ -0,0 +1,76 @@ +import { type Value } from "./camino/api_pb.js"; +import { type RuntimeHandler, type DerivedHandler } from "./index.js"; +declare const referenceBrand: unique symbol; +export type QxObjectRef = string & { + readonly [referenceBrand]: { + readonly [K in Identity]: true; + }; +}; +declare const watchBrand: unique symbol; +export type QxWatchHandle = string & { + readonly [watchBrand]: true; +}; +export type MessageBinding = { + encode(value: T): Value; + decode(value: Value): T; +}; +export type BindingValue = B extends MessageBinding ? T : never; +export type QxHandler = (context: C) => O | Promise; +export type QxDerived = { + kind: "derived"; + get: QxHandler; +}; +export declare const qxDerived: (get: QxHandler) => QxDerived; +/** Versioned binding ABI. This mirrors the language-neutral value IR. */ +export type QxValueType = { + kind: "builtin"; + name: "unit" | "watch-handle"; +} | { + kind: "scalar"; + name: string; +} | { + kind: "message"; + descriptorId: string; +} | { + kind: "object-ref"; + expectation: unknown; +} | { + kind: "optional" | "list"; + value: QxValueType; +}; +export type QxOperationSpec = { + id: string; + inputType: QxValueType; + outputType: QxValueType; +}; +export type QxPortSpec = { + kind: "state"; + id: string; + valueType: QxValueType; + primitives: string[]; +} | { + kind: "edge"; + id: string; + primitives: string[]; +} | { + kind: "interface"; + id: string; + operations: Record; +} | { + kind: "constructor"; + id: string; + inputType: QxValueType; +}; +export type QxHandlerSpec = { + inputType: QxValueType; + outputType: QxValueType; + eventType?: QxValueType; + ports: Record; +}; +export type QxMessages = Record>; +export declare const decodeQxValue: (type: QxValueType, value: Value | undefined, messages: QxMessages) => any; +export declare const encodeQxValue: (type: QxValueType, value: any, messages: QxMessages) => Value; +/** The sole unchecked cast connects generated contracts to the dynamic RPC runtime. */ +export declare const bindQxHandler: (spec: QxHandlerSpec, handler: QxHandler | QxDerived, messages: QxMessages) => RuntimeHandler | DerivedHandler; +export {}; +//# sourceMappingURL=bindings.d.ts.map \ No newline at end of file diff --git a/dist/bindings.d.ts.map b/dist/bindings.d.ts.map new file mode 100644 index 0000000..77c3341 --- /dev/null +++ b/dist/bindings.d.ts.map @@ -0,0 +1 @@ +{"version":3,"file":"bindings.d.ts","sourceRoot":"","sources":["../src/bindings.ts"],"names":[],"mappings":"AACA,OAAO,EAAkC,KAAK,KAAK,EAAE,MAAM,oBAAoB,CAAC;AAChF,OAAO,EACgB,KAAK,cAAc,EAAE,KAAK,cAAc,EAAE,MAAM,YAAY,CAAC;AAEpF,OAAO,CAAC,MAAM,cAAc,EAAE,OAAO,MAAM,CAAC;AAC5C,MAAM,MAAM,WAAW,CAAC,QAAQ,SAAS,MAAM,IAAI,MAAM,GAAG;IAAE,QAAQ,CAAC,CAAC,cAAc,CAAC,EAAE;QAAE,QAAQ,EAAE,CAAC,IAAI,QAAQ,GAAG,IAAI;KAAE,CAAA;CAAE,CAAC;AAC9H,OAAO,CAAC,MAAM,UAAU,EAAE,OAAO,MAAM,CAAC;AACxC,MAAM,MAAM,aAAa,GAAG,MAAM,GAAG;IAAE,QAAQ,CAAC,CAAC,UAAU,CAAC,EAAE,IAAI,CAAA;CAAE,CAAC;AACrE,MAAM,MAAM,cAAc,CAAC,CAAC,IAAI;IAAE,MAAM,CAAC,KAAK,EAAE,CAAC,GAAG,KAAK,CAAC;IAAC,MAAM,CAAC,KAAK,EAAE,KAAK,GAAG,CAAC,CAAA;CAAE,CAAC;AACrF,MAAM,MAAM,YAAY,CAAC,CAAC,IAAI,CAAC,SAAS,cAAc,CAAC,MAAM,CAAC,CAAC,GAAG,CAAC,GAAG,KAAK,CAAC;AAC5E,MAAM,MAAM,SAAS,CAAC,CAAC,EAAE,CAAC,IAAI,CAAC,OAAO,EAAE,CAAC,KAAK,CAAC,GAAG,OAAO,CAAC,CAAC,CAAC,CAAC;AAC7D,MAAM,MAAM,SAAS,CAAC,CAAC,EAAE,CAAC,IAAI;IAAE,IAAI,EAAE,SAAS,CAAC;IAAC,GAAG,EAAE,SAAS,CAAC,CAAC,EAAE,CAAC,CAAC,CAAA;CAAE,CAAC;AACxE,eAAO,MAAM,SAAS,GAAI,CAAC,EAAE,CAAC,OAAO,SAAS,CAAC,CAAC,EAAE,CAAC,CAAC,KAAG,SAAS,CAAC,CAAC,EAAE,CAAC,CAA+B,CAAC;AAErG,yEAAyE;AACzE,MAAM,MAAM,WAAW,GACnB;IAAE,IAAI,EAAE,SAAS,CAAC;IAAC,IAAI,EAAE,MAAM,GAAG,cAAc,CAAA;CAAE,GAClD;IAAE,IAAI,EAAE,QAAQ,CAAC;IAAC,IAAI,EAAE,MAAM,CAAA;CAAE,GAChC;IAAE,IAAI,EAAE,SAAS,CAAC;IAAC,YAAY,EAAE,MAAM,CAAA;CAAE,GACzC;IAAE,IAAI,EAAE,YAAY,CAAC;IAAC,WAAW,EAAE,OAAO,CAAA;CAAE,GAC5C;IAAE,IAAI,EAAE,UAAU,GAAG,MAAM,CAAC;IAAC,KAAK,EAAE,WAAW,CAAA;CAAE,CAAC;AACtD,MAAM,MAAM,eAAe,GAAG;IAAE,EAAE,EAAE,MAAM,CAAC;IAAC,SAAS,EAAE,WAAW,CAAC;IAAC,UAAU,EAAE,WAAW,CAAA;CAAE,CAAC;AAC9F,MAAM,MAAM,UAAU,GAClB;IAAE,IAAI,EAAE,OAAO,CAAC;IAAC,EAAE,EAAE,MAAM,CAAC;IAAC,SAAS,EAAE,WAAW,CAAC;IAAC,UAAU,EAAE,MAAM,EAAE,CAAA;CAAE,GAC3E;IAAE,IAAI,EAAE,MAAM,CAAC;IAAC,EAAE,EAAE,MAAM,CAAC;IAAC,UAAU,EAAE,MAAM,EAAE,CAAA;CAAE,GAClD;IAAE,IAAI,EAAE,WAAW,CAAC;IAAC,EAAE,EAAE,MAAM,CAAC;IAAC,UAAU,EAAE,MAAM,CAAC,MAAM,EAAE,eAAe,CAAC,CAAA;CAAE,GAC9E;IAAE,IAAI,EAAE,aAAa,CAAC;IAAC,EAAE,EAAE,MAAM,CAAC;IAAC,SAAS,EAAE,WAAW,CAAA;CAAE,CAAC;AAChE,MAAM,MAAM,aAAa,GAAG;IAC1B,SAAS,EAAE,WAAW,CAAC;IAAC,UAAU,EAAE,WAAW,CAAC;IAAC,SAAS,CAAC,EAAE,WAAW,CAAC;IACzE,KAAK,EAAE,MAAM,CAAC,MAAM,EAAE,UAAU,CAAC,CAAC;CACnC,CAAC;AACF,MAAM,MAAM,UAAU,GAAG,MAAM,CAAC,MAAM,EAAE,cAAc,CAAC,GAAG,CAAC,CAAC,CAAC;AAG7D,eAAO,MAAM,aAAa,SAAU,WAAW,SAAS,KAAK,GAAG,SAAS,YAAY,UAAU,KAAG,GAoBjG,CAAC;AAOF,eAAO,MAAM,aAAa,SAAU,WAAW,SAAS,GAAG,YAAY,UAAU,KAAG,KAOnF,CAAC;AAiBF,uFAAuF;AACvF,eAAO,MAAM,aAAa,GAAI,CAAC,EAAE,CAAC,QAC1B,aAAa,WAAW,SAAS,CAAC,CAAC,EAAE,CAAC,CAAC,GAAG,SAAS,CAAC,CAAC,EAAE,CAAC,CAAC,YAAY,UAAU,KACpF,cAAc,GAAG,cA+BnB,CAAC"} \ No newline at end of file diff --git a/dist/bindings.js b/dist/bindings.js new file mode 100644 index 0000000..6277cdc --- /dev/null +++ b/dist/bindings.js @@ -0,0 +1,101 @@ +import { create } from "@bufbuild/protobuf"; +import { ValueSchema, ObjectValueSchema } from "./camino/api_pb.js"; +import { derived, jsToProtoValue, liveValue, protoValueToJs } from "./index.js"; +export const qxDerived = (get) => ({ kind: "derived", get }); +// Conversion belongs at the binding boundary. It does not add orchestrator validation. +export const decodeQxValue = (type, value, messages) => { + if (type.kind === "builtin" && type.name === "unit") + return null; + if (!value) + throw new Error("Missing QX wire value"); + if (type.kind === "optional") + return value.kind.case === "nullValue" ? null : decodeQxValue(type.value, value, messages); + if (type.kind === "list") { + if (value.kind.case !== "listValue") + throw new Error("Expected QX list"); + return value.kind.value.values.map((entry) => decodeQxValue(type.value, entry, messages)); + } + if (type.kind === "message") + return requireMessage(messages, type.descriptorId).decode(value); + if (type.kind === "scalar") { + if (type.name === "int64" || type.name === "uint64") { + if (value.kind.case !== "integerValue") + throw new Error("Expected QX integer"); + return BigInt(value.kind.value); + } + if (type.name === "bytes") { + if (value.kind.case !== "bytesValue") + throw new Error("Expected QX bytes"); + return value.kind.value; + } + } + return protoValueToJs(value); +}; +const requireMessage = (messages, id) => { + const binding = messages[id]; + if (!binding) + throw new Error(`Missing message binding ${id}`); + return binding; +}; +export const encodeQxValue = (type, value, messages) => { + if (type.kind === "builtin" && type.name === "unit") + return jsToProtoValue(null); + if (type.kind === "optional") + return value === null ? jsToProtoValue(null) : encodeQxValue(type.value, value, messages); + if (type.kind === "list") + return jsToProtoValue(value.map((entry) => liveValue(encodeQxValue(type.value, entry, messages)))); + if (type.kind === "message") + return requireMessage(messages, type.descriptorId).encode(value); + if (type.kind === "object-ref") + return jsToProtoValue({ $quixosRef: value }); + return jsToProtoValue(value); +}; +const inputValue = (context, type) => { + if (type.kind === "message") + return create(ValueSchema, { kind: { case: "objectValue", + value: create(ObjectValueSchema, { fields: context.inputProto }) } }); + return context.inputProto.value; +}; +const inputFields = (type, value, messages) => { + if (type.kind === "builtin" && type.name === "unit") + return {}; + const encoded = encodeQxValue(type, value, messages); + if (type.kind === "message") { + if (encoded.kind.case !== "objectValue") + throw new Error("Message inputs must encode an object value"); + return Object.fromEntries(Object.entries(encoded.kind.value.fields).map(([key, entry]) => [key, liveValue(entry)])); + } + return { value: liveValue(encoded) }; +}; +/** The sole unchecked cast connects generated contracts to the dynamic RPC runtime. */ +export const bindQxHandler = (spec, handler, messages) => { + const execute = async (raw) => { + const ports = Object.fromEntries(Object.entries(spec.ports).map(([name, port]) => { + switch (port.kind) { + case "state": { + const state = raw.state(port.id); + return [name, { + ...(port.primitives.includes("read") ? { get: async () => decodeQxValue(port.valueType, (await state.live()).$quixosValue, messages) } : {}), + ...(port.primitives.includes("write") ? { set: async (value) => state.set(liveValue(encodeQxValue(port.valueType, value, messages))) } : {}), + }]; + } + case "edge": { + const edge = raw.edge(port.id); + return [name, Object.fromEntries(port.primitives.map((primitive) => [primitive, edge[primitive]]))]; + } + case "interface": { + const target = raw.interface(port.id); + return [name, Object.fromEntries(Object.entries(port.operations).map(([name, operation]) => [name, + async (input) => decodeQxValue(operation.outputType, (await target.live(operation.id, inputFields(operation.inputType, input, messages))).$quixosValue, messages), + ]))]; + } + case "constructor": return [name, { construct: (input) => raw.constructor(port.id).construct(inputFields(port.inputType, input, messages)) }]; + } + })); + const context = { objectId: raw.objectId, + input: decodeQxValue(spec.inputType, inputValue(raw, spec.inputType), messages), ports }; + const value = await (typeof handler === "function" ? handler(context) : handler.get(context)); + return liveValue(encodeQxValue(spec.eventType ?? spec.outputType, value, messages)); + }; + return typeof handler === "function" ? execute : derived(execute); +}; diff --git a/dist/index.d.ts b/dist/index.d.ts index d876040..91bb2a1 100644 --- a/dist/index.d.ts +++ b/dist/index.d.ts @@ -1,4 +1,5 @@ import http from "node:http"; +export * from "./bindings.js"; import { type Client, type ConnectRouter } from "@connectrpc/connect"; import { CaminoService, type Value } from "./camino/api_pb.js"; import { OrchestratorRuntime } from "./quixos/orch_pb.js"; @@ -36,6 +37,7 @@ export type EdgePort = { projectionId: string; resolve(): Promise; connect(targetObjectId: string): Promise; + disconnect(targetObjectId: string): Promise; }; export type InterfacePort = { objectId: string; @@ -80,5 +82,4 @@ type RuntimeRequest = { input: Record; dependencies: import("./quixos/refs_pb.js").InjectedDependency[]; }; -export {}; //# sourceMappingURL=index.d.ts.map \ No newline at end of file diff --git a/dist/index.d.ts.map b/dist/index.d.ts.map index 4259b5e..585d1f4 100644 --- a/dist/index.d.ts.map +++ b/dist/index.d.ts.map @@ -1 +1 @@ -{"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../src/index.ts"],"names":[],"mappings":"AAAA,OAAO,IAAI,MAAM,WAAW,CAAC;AAI7B,OAAO,EAAoC,KAAK,MAAM,EAAE,KAAK,aAAa,EAAE,MAAM,qBAAqB,CAAC;AAExG,OAAO,EACL,aAAa,EAOb,KAAK,KAAK,EACX,MAAM,oBAAoB,CAAC;AAC5B,OAAO,EAAE,mBAAmB,EAAE,MAAM,qBAAqB,CAAC;AAU1D,MAAM,MAAM,YAAY,GAAG,MAAM,CAAC,OAAO,aAAa,CAAC,CAAC;AACxD,MAAM,MAAM,UAAU,GAAG,MAAM,CAAC,OAAO,mBAAmB,CAAC,CAAC;AAE5D,MAAM,MAAM,iBAAiB,GACzB;IAAE,IAAI,EAAE,OAAO,CAAC;IAAC,QAAQ,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAA;CAAE,GACzD;IAAE,IAAI,EAAE,MAAM,CAAC;IAAC,QAAQ,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAA;CAAE,CAAC;AA+BnF,eAAO,MAAM,SAAS,aAAc,MAAM;IAAQ,UAAU;CAAa,CAAC;AAC1E,eAAO,MAAM,SAAS,UAAW,KAAK;IAAQ,YAAY;CAAU,CAAC;AAErE,eAAO,MAAM,cAAc,UAAW,OAAO,KAAG,KAoC/C,CAAC;AAEF,eAAO,MAAM,cAAc,UAAW,KAAK,GAAG,SAAS,KAAG,OAoBzD,CAAC;AAEF,eAAO,MAAM,eAAe,WAAY,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC;;CACmC,CAAC;AAEjG,MAAM,MAAM,SAAS,CAAC,CAAC,GAAG,OAAO,IAAI;IACnC,MAAM,EAAE,MAAM,CAAC;IACf,GAAG,IAAI,OAAO,CAAC,CAAC,CAAC,CAAC;IAClB,IAAI,IAAI,OAAO,CAAC,UAAU,CAAC,OAAO,SAAS,CAAC,CAAC,CAAC;IAC9C,GAAG,CAAC,KAAK,EAAE,CAAC,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;CAC9B,CAAC;AACF,MAAM,MAAM,QAAQ,GAAG;IACrB,UAAU,EAAE,MAAM,CAAC;IACnB,YAAY,EAAE,MAAM,CAAC;IACrB,OAAO,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC,CAAC;IAC7B,OAAO,CAAC,cAAc,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;CAChD,CAAC;AACF,MAAM,MAAM,aAAa,GAAG;IAC1B,QAAQ,EAAE,MAAM,CAAC;IACjB,mBAAmB,EAAE,MAAM,CAAC;IAC5B,MAAM,CAAC,WAAW,EAAE,MAAM,EAAE,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,OAAO,CAAC,CAAC;IAC/E,IAAI,CAAC,WAAW,EAAE,MAAM,EAAE,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,UAAU,CAAC,OAAO,SAAS,CAAC,CAAC,CAAC;CACnG,CAAC;AACF,MAAM,MAAM,eAAe,GAAG;IAC5B,MAAM,EAAE,MAAM,CAAC;IACf,SAAS,CAAC,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,MAAM,CAAC,CAAC;CAC7D,CAAC;AACF,MAAM,MAAM,WAAW,GAAG,SAAS,GAAG,QAAQ,GAAG,aAAa,GAAG,eAAe,CAAC;AAEjF,MAAM,MAAM,cAAc,GAAG;IAC3B,QAAQ,EAAE,MAAM,CAAC;IACjB,KAAK,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,CAAC;IAC/B,UAAU,EAAE,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC,CAAC;IAClC,KAAK,EAAE,WAAW,CAAC,MAAM,EAAE,WAAW,CAAC,CAAC;IACxC,KAAK,CAAC,CAAC,GAAG,OAAO,EAAE,MAAM,EAAE,MAAM,GAAG,SAAS,CAAC,CAAC,CAAC,CAAC;IACjD,IAAI,CAAC,MAAM,EAAE,MAAM,GAAG,QAAQ,CAAC;IAC/B,SAAS,CAAC,MAAM,EAAE,MAAM,GAAG,aAAa,CAAC;IACzC,WAAW,CAAC,MAAM,EAAE,MAAM,GAAG,eAAe,CAAC;CAC9C,CAAC;AAOF,eAAO,MAAM,oBAAoB,WACvB,YAAY,QACd,UAAU,WACP,cAAc,KACtB,cAuHF,CAAC;AAEF,MAAM,MAAM,cAAc,GAAG,CAAC,OAAO,EAAE,cAAc,KAAK,OAAO,GAAG,OAAO,CAAC,OAAO,CAAC,CAAC;AACrF,MAAM,MAAM,cAAc,GAAG;IAAE,IAAI,EAAE,SAAS,CAAC;IAAC,GAAG,EAAE,cAAc,CAAA;CAAE,CAAC;AACtE,eAAO,MAAM,OAAO,QAAS,cAAc,KAAG,cAA4C,CAAC;AAyB3F,eAAO,MAAM,0BAA0B,WAAY;IACjD,iBAAiB,EAAE,MAAM,CAAC;IAC1B,OAAO,EAAE,MAAM,CAAC,MAAM,EAAE,cAAc,GAAG,cAAc,CAAC,CAAC;IACzD,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,OAAO,CAAC,EAAE,MAAM,CAAC;CAClB,cAsBiB,aAAa,kBA8J9B,CAAC;AAEF,eAAO,MAAM,mBAAmB,WAAY;IAC1C,iBAAiB,EAAE,MAAM,CAAC;IAC1B,OAAO,EAAE,MAAM,CAAC,MAAM,EAAE,cAAc,GAAG,cAAc,CAAC,CAAC;CAC1D,yEAaA,CAAC;AACF,KAAK,cAAc,GAAG;IACpB,QAAQ,EAAE,MAAM,CAAC;IACjB,KAAK,EAAE,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC,CAAC;IAC7B,YAAY,EAAE,OAAO,qBAAqB,EAAE,kBAAkB,EAAE,CAAC;CAClE,CAAC"} \ No newline at end of file +{"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../src/index.ts"],"names":[],"mappings":"AAAA,OAAO,IAAI,MAAM,WAAW,CAAC;AAC7B,cAAc,eAAe,CAAC;AAI9B,OAAO,EAAoC,KAAK,MAAM,EAAE,KAAK,aAAa,EAAE,MAAM,qBAAqB,CAAC;AAExG,OAAO,EACL,aAAa,EAOb,KAAK,KAAK,EACX,MAAM,oBAAoB,CAAC;AAC5B,OAAO,EAAE,mBAAmB,EAAE,MAAM,qBAAqB,CAAC;AAU1D,MAAM,MAAM,YAAY,GAAG,MAAM,CAAC,OAAO,aAAa,CAAC,CAAC;AACxD,MAAM,MAAM,UAAU,GAAG,MAAM,CAAC,OAAO,mBAAmB,CAAC,CAAC;AAE5D,MAAM,MAAM,iBAAiB,GACzB;IAAE,IAAI,EAAE,OAAO,CAAC;IAAC,QAAQ,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAA;CAAE,GACzD;IAAE,IAAI,EAAE,MAAM,CAAC;IAAC,QAAQ,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAC;IAAC,YAAY,EAAE,MAAM,CAAA;CAAE,CAAC;AA+BnF,eAAO,MAAM,SAAS,aAAc,MAAM;IAAQ,UAAU;CAAa,CAAC;AAC1E,eAAO,MAAM,SAAS,UAAW,KAAK;IAAQ,YAAY;CAAU,CAAC;AAErE,eAAO,MAAM,cAAc,UAAW,OAAO,KAAG,KAoC/C,CAAC;AAEF,eAAO,MAAM,cAAc,UAAW,KAAK,GAAG,SAAS,KAAG,OAoBzD,CAAC;AAEF,eAAO,MAAM,eAAe,WAAY,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC;;CACmC,CAAC;AAEjG,MAAM,MAAM,SAAS,CAAC,CAAC,GAAG,OAAO,IAAI;IACnC,MAAM,EAAE,MAAM,CAAC;IACf,GAAG,IAAI,OAAO,CAAC,CAAC,CAAC,CAAC;IAClB,IAAI,IAAI,OAAO,CAAC,UAAU,CAAC,OAAO,SAAS,CAAC,CAAC,CAAC;IAC9C,GAAG,CAAC,KAAK,EAAE,CAAC,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;CAC9B,CAAC;AACF,MAAM,MAAM,QAAQ,GAAG;IACrB,UAAU,EAAE,MAAM,CAAC;IACnB,YAAY,EAAE,MAAM,CAAC;IACrB,OAAO,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC,CAAC;IAC7B,OAAO,CAAC,cAAc,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;IAC/C,UAAU,CAAC,cAAc,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC,CAAC;CACnD,CAAC;AACF,MAAM,MAAM,aAAa,GAAG;IAC1B,QAAQ,EAAE,MAAM,CAAC;IACjB,mBAAmB,EAAE,MAAM,CAAC;IAC5B,MAAM,CAAC,WAAW,EAAE,MAAM,EAAE,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,OAAO,CAAC,CAAC;IAC/E,IAAI,CAAC,WAAW,EAAE,MAAM,EAAE,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,UAAU,CAAC,OAAO,SAAS,CAAC,CAAC,CAAC;CACnG,CAAC;AACF,MAAM,MAAM,eAAe,GAAG;IAC5B,MAAM,EAAE,MAAM,CAAC;IACf,SAAS,CAAC,KAAK,CAAC,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,GAAG,OAAO,CAAC,MAAM,CAAC,CAAC;CAC7D,CAAC;AACF,MAAM,MAAM,WAAW,GAAG,SAAS,GAAG,QAAQ,GAAG,aAAa,GAAG,eAAe,CAAC;AAEjF,MAAM,MAAM,cAAc,GAAG;IAC3B,QAAQ,EAAE,MAAM,CAAC;IACjB,KAAK,EAAE,MAAM,CAAC,MAAM,EAAE,OAAO,CAAC,CAAC;IAC/B,UAAU,EAAE,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC,CAAC;IAClC,KAAK,EAAE,WAAW,CAAC,MAAM,EAAE,WAAW,CAAC,CAAC;IACxC,KAAK,CAAC,CAAC,GAAG,OAAO,EAAE,MAAM,EAAE,MAAM,GAAG,SAAS,CAAC,CAAC,CAAC,CAAC;IACjD,IAAI,CAAC,MAAM,EAAE,MAAM,GAAG,QAAQ,CAAC;IAC/B,SAAS,CAAC,MAAM,EAAE,MAAM,GAAG,aAAa,CAAC;IACzC,WAAW,CAAC,MAAM,EAAE,MAAM,GAAG,eAAe,CAAC;CAC9C,CAAC;AAOF,eAAO,MAAM,oBAAoB,WACvB,YAAY,QACd,UAAU,WACP,cAAc,KACtB,cA6HF,CAAC;AAEF,MAAM,MAAM,cAAc,GAAG,CAAC,OAAO,EAAE,cAAc,KAAK,OAAO,GAAG,OAAO,CAAC,OAAO,CAAC,CAAC;AACrF,MAAM,MAAM,cAAc,GAAG;IAAE,IAAI,EAAE,SAAS,CAAC;IAAC,GAAG,EAAE,cAAc,CAAA;CAAE,CAAC;AACtE,eAAO,MAAM,OAAO,QAAS,cAAc,KAAG,cAA4C,CAAC;AAyB3F,eAAO,MAAM,0BAA0B,WAAY;IACjD,iBAAiB,EAAE,MAAM,CAAC;IAC1B,OAAO,EAAE,MAAM,CAAC,MAAM,EAAE,cAAc,GAAG,cAAc,CAAC,CAAC;IACzD,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,OAAO,CAAC,EAAE,MAAM,CAAC;CAClB,cAsBiB,aAAa,kBA8J9B,CAAC;AAEF,eAAO,MAAM,mBAAmB,WAAY;IAC1C,iBAAiB,EAAE,MAAM,CAAC;IAC1B,OAAO,EAAE,MAAM,CAAC,MAAM,EAAE,cAAc,GAAG,cAAc,CAAC,CAAC;CAC1D,yEAaA,CAAC;AACF,KAAK,cAAc,GAAG;IACpB,QAAQ,EAAE,MAAM,CAAC;IACjB,KAAK,EAAE,MAAM,CAAC,MAAM,EAAE,KAAK,CAAC,CAAC;IAC7B,YAAY,EAAE,OAAO,qBAAqB,EAAE,kBAAkB,EAAE,CAAC;CAClE,CAAC"} \ No newline at end of file diff --git a/dist/index.js b/dist/index.js index e007a41..b9c237a 100644 --- a/dist/index.js +++ b/dist/index.js @@ -1,4 +1,5 @@ import http from "node:http"; +export * from "./bindings.js"; import { AsyncLocalStorage } from "node:async_hooks"; import { randomUUID } from "node:crypto"; import { create, equals } from "@bufbuild/protobuf"; @@ -131,6 +132,13 @@ export const createRuntimeContext = (camino, orch, request) => { async connect(targetObjectId) { await camino.connectEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId, targetObjectId }); }, + async disconnect(targetObjectId) { + const result = await camino.resolveEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId }); + for (const entry of result.edges) { + if (targetForEdge(entry, projectionId) === targetObjectId) + await camino.disconnectEdge({ edgeId: entry.id }); + } + }, }; ports.set(dependency.portId, edge); break; diff --git a/src/bindings.ts b/src/bindings.ts new file mode 100644 index 0000000..8fa44e9 --- /dev/null +++ b/src/bindings.ts @@ -0,0 +1,121 @@ +import { create } from "@bufbuild/protobuf"; +import { ValueSchema, ObjectValueSchema, type Value } from "./camino/api_pb.js"; +import { derived, jsToProtoValue, liveValue, protoValueToJs, + type RuntimeContext, type RuntimeHandler, type DerivedHandler } from "./index.js"; + +declare const referenceBrand: unique symbol; +export type QxObjectRef = string & { readonly [referenceBrand]: { readonly [K in Identity]: true } }; +declare const watchBrand: unique symbol; +export type QxWatchHandle = string & { readonly [watchBrand]: true }; +export type MessageBinding = { encode(value: T): Value; decode(value: Value): T }; +export type BindingValue = B extends MessageBinding ? T : never; +export type QxHandler = (context: C) => O | Promise; +export type QxDerived = { kind: "derived"; get: QxHandler }; +export const qxDerived = (get: QxHandler): QxDerived => ({ kind: "derived", get }); + +/** Versioned binding ABI. This mirrors the language-neutral value IR. */ +export type QxValueType = + | { kind: "builtin"; name: "unit" | "watch-handle" } + | { kind: "scalar"; name: string } + | { kind: "message"; descriptorId: string } + | { kind: "object-ref"; expectation: unknown } + | { kind: "optional" | "list"; value: QxValueType }; +export type QxOperationSpec = { id: string; inputType: QxValueType; outputType: QxValueType }; +export type QxPortSpec = + | { kind: "state"; id: string; valueType: QxValueType; primitives: string[] } + | { kind: "edge"; id: string; primitives: string[] } + | { kind: "interface"; id: string; operations: Record } + | { kind: "constructor"; id: string; inputType: QxValueType }; +export type QxHandlerSpec = { + inputType: QxValueType; outputType: QxValueType; eventType?: QxValueType; + ports: Record; +}; +export type QxMessages = Record>; + +// Conversion belongs at the binding boundary. It does not add orchestrator validation. +export const decodeQxValue = (type: QxValueType, value: Value | undefined, messages: QxMessages): any => { + if (type.kind === "builtin" && type.name === "unit") return null; + if (!value) throw new Error("Missing QX wire value"); + if (type.kind === "optional") return value.kind.case === "nullValue" ? null : decodeQxValue(type.value, value, messages); + if (type.kind === "list") { + if (value.kind.case !== "listValue") throw new Error("Expected QX list"); + return value.kind.value.values.map((entry) => decodeQxValue(type.value, entry, messages)); + } + if (type.kind === "message") return requireMessage(messages, type.descriptorId).decode(value); + if (type.kind === "scalar") { + if (type.name === "int64" || type.name === "uint64") { + if (value.kind.case !== "integerValue") throw new Error("Expected QX integer"); + return BigInt(value.kind.value); + } + if (type.name === "bytes") { + if (value.kind.case !== "bytesValue") throw new Error("Expected QX bytes"); + return value.kind.value; + } + } + return protoValueToJs(value); +}; + +const requireMessage = (messages: QxMessages, id: string) => { + const binding = messages[id]; + if (!binding) throw new Error(`Missing message binding ${id}`); + return binding; +}; +export const encodeQxValue = (type: QxValueType, value: any, messages: QxMessages): Value => { + if (type.kind === "builtin" && type.name === "unit") return jsToProtoValue(null); + if (type.kind === "optional") return value === null ? jsToProtoValue(null) : encodeQxValue(type.value, value, messages); + if (type.kind === "list") return jsToProtoValue(value.map((entry: unknown) => liveValue(encodeQxValue(type.value, entry, messages)))); + if (type.kind === "message") return requireMessage(messages, type.descriptorId).encode(value); + if (type.kind === "object-ref") return jsToProtoValue({ $quixosRef: value }); + return jsToProtoValue(value); +}; + +const inputValue = (context: RuntimeContext, type: QxValueType) => { + if (type.kind === "message") return create(ValueSchema, { kind: { case: "objectValue", + value: create(ObjectValueSchema, { fields: context.inputProto }) } }); + return context.inputProto.value; +}; +const inputFields = (type: QxValueType, value: unknown, messages: QxMessages): Record => { + if (type.kind === "builtin" && type.name === "unit") return {}; + const encoded = encodeQxValue(type, value, messages); + if (type.kind === "message") { + if (encoded.kind.case !== "objectValue") throw new Error("Message inputs must encode an object value"); + return Object.fromEntries(Object.entries(encoded.kind.value.fields).map(([key, entry]) => [key, liveValue(entry)])); + } + return { value: liveValue(encoded) }; +}; + +/** The sole unchecked cast connects generated contracts to the dynamic RPC runtime. */ +export const bindQxHandler = ( + spec: QxHandlerSpec, handler: QxHandler | QxDerived, messages: QxMessages, +): RuntimeHandler | DerivedHandler => { + const execute = async (raw: RuntimeContext) => { + const ports = Object.fromEntries(Object.entries(spec.ports).map(([name, port]) => { + switch (port.kind) { + case "state": { + const state = raw.state(port.id); + return [name, { + ...(port.primitives.includes("read") ? { get: async () => decodeQxValue(port.valueType, (await state.live()).$quixosValue, messages) } : {}), + ...(port.primitives.includes("write") ? { set: async (value: unknown) => state.set(liveValue(encodeQxValue(port.valueType, value, messages))) } : {}), + }]; + } + case "edge": { + const edge = raw.edge(port.id); + return [name, Object.fromEntries(port.primitives.map((primitive) => [primitive, edge[primitive as "resolve" | "connect" | "disconnect"]]))]; + } + case "interface": { + const target = raw.interface(port.id); + return [name, Object.fromEntries(Object.entries(port.operations).map(([name, operation]) => [name, + async (input: unknown) => decodeQxValue(operation.outputType, + (await target.live(operation.id, inputFields(operation.inputType, input, messages))).$quixosValue, messages), + ]))]; + } + case "constructor": return [name, { construct: (input: unknown) => raw.constructor(port.id).construct(inputFields(port.inputType, input, messages)) }]; + } + })); + const context = { objectId: raw.objectId, + input: decodeQxValue(spec.inputType, inputValue(raw, spec.inputType), messages), ports } as C; + const value = await (typeof handler === "function" ? handler(context) : handler.get(context)); + return liveValue(encodeQxValue(spec.eventType ?? spec.outputType, value, messages)); + }; + return typeof handler === "function" ? execute : derived(execute); +}; diff --git a/src/index.ts b/src/index.ts index cd7f223..8031a5a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,4 +1,5 @@ import http from "node:http"; +export * from "./bindings.js"; import { AsyncLocalStorage } from "node:async_hooks"; import { randomUUID } from "node:crypto"; import { create, equals } from "@bufbuild/protobuf"; @@ -137,6 +138,7 @@ export type EdgePort = { projectionId: string; resolve(): Promise; connect(targetObjectId: string): Promise; + disconnect(targetObjectId: string): Promise; }; export type InterfacePort = { objectId: string; @@ -210,6 +212,12 @@ export const createRuntimeContext = ( async connect(targetObjectId) { await camino.connectEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId, targetObjectId }); }, + async disconnect(targetObjectId) { + const result = await camino.resolveEdge({ objectId: dependencyObjectId, edgeTypeId, projectionId }); + for (const entry of result.edges) { + if (targetForEdge(entry, projectionId) === targetObjectId) await camino.disconnectEdge({ edgeId: entry.id }); + } + }, }; ports.set(dependency.portId, edge); break; diff --git a/test/bindings.test.mjs b/test/bindings.test.mjs new file mode 100644 index 0000000..90fd74e --- /dev/null +++ b/test/bindings.test.mjs @@ -0,0 +1,45 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { bindQxHandler, decodeQxValue, encodeQxValue, jsToProtoValue, liveValue, protoValueToJs } from "../dist/index.js"; +const scalar = (name) => ({ kind: "scalar", name }); +const unit = { kind: "builtin", name: "unit" }; +test("binding codecs round trip nested bytes, 64-bit integers, nulls, and references", () => { + const values = [[scalar("int64"), -(2n ** 63n)], [scalar("uint64"), 2n ** 64n - 1n], + [scalar("bytes"), new Uint8Array([0, 255])], + [{ kind: "list", value: { kind: "optional", value: scalar("int64") } }, [null, 2n ** 60n]], + [{ kind: "object-ref", expectation: { kind: "atom", atomId: "thing" } }, "obj:thing"]]; + for (const [type, value] of values) assert.deepEqual(decodeQxValue(type, encodeQxValue(type, value, {}), {}), value); +}); +test("typed state and interface ports preserve declared values and exact operation IDs", async () => { + const spec = { inputType: scalar("int64"), outputType: scalar("bytes"), ports: { + data: { kind: "state", id: "state-id", valueType: scalar("int64"), primitives: ["read", "write"] }, + reader: { kind: "interface", id: "interface-id", operations: { "payload.get": { id: "get-id", inputType: unit, outputType: scalar("bytes") } } }, + } }; + let written; + const handler = bindQxHandler(spec, async ({ input, ports }) => { + assert.equal(input, 2n ** 60n); + assert.equal(await ports.data.get(), 9n); + await ports.data.set(input); + return ports.reader["payload.get"](); + }, {}); + const result = await handler({ inputProto: { value: jsToProtoValue(2n ** 60n) }, objectId: "obj", + state: (id) => { assert.equal(id, "state-id"); return { + live: async () => liveValue(jsToProtoValue(9n)), set: async (value) => { written = value; }, + }; }, + interface: (id) => { assert.equal(id, "interface-id"); return { live: async (operation, input) => { + assert.equal(operation, "get-id"); assert.deepEqual(input, {}); return liveValue(jsToProtoValue(new Uint8Array([7]))); + } }; }, + }); + assert.equal(written.$quixosValue.kind.value, String(2n ** 60n)); + assert.deepEqual(result.$quixosValue.kind.value, new Uint8Array([7])); +}); +test("external message bindings and derived event types are used at the boundary", async () => { + const message = { kind: "message", descriptorId: "Payload" }; + const messages = { Payload: { encode: jsToProtoValue, decode: protoValueToJs } }; + const handler = bindQxHandler({ inputType: message, outputType: { kind: "builtin", name: "watch-handle" }, eventType: message, ports: {} }, + { kind: "derived", get: ({ input }) => ({ value: input.title }) }, messages); + assert.equal(handler.kind, "derived"); + const result = await handler.get({ objectId: "obj", inputProto: { title: jsToProtoValue("hello") } }); + assert.deepEqual(protoValueToJs(result.$quixosValue), { value: "hello" }); + assert.throws(() => encodeQxValue(message, {}, {}), /Missing message binding/); +});