diff --git a/.agents/docs/gateway-extensions.md b/.agents/docs/gateway-extensions.md new file mode 100644 index 0000000..44b6ff9 --- /dev/null +++ b/.agents/docs/gateway-extensions.md @@ -0,0 +1,65 @@ +# Gateway extensions + +The `@cordisx/plugin-cli-proxy-api/extensions/v1` entrypoint is the generic +CordisX contract for extending the managed CLIProxyAPI gateway. It is not a +provider-specific integration and must not contain private tenant, domain, +header, or product policy. + +## Topology + +- CordisX runs one Host-owned CLIProxyAPI managed process. +- The gateway package supplies one native `cordisx-gateway-bridge` plugin. +- Consumer plugins register adapters and connections through the + `cliProxyGatewayExtensions` context service. +- A registration change creates one generation-fenced materialization and + managed-service restart transaction. Disposing the consumer binding revokes + only that producer generation. +- The bridge publishes models as `connectionId/modelId` and routes them to its + own static executor. It never starts a proxy process per connection. + +## Public contract + +An adapter declares bounded session extraction and request transforms. Session +sources are an HTTP header or a body JSON Pointer. Transforms may clear headers +and set headers or body values with `{{session}}`, `{{connectionId}}`, and +`{{sourceModelId}}` substitutions. + +A connection references an existing Host-managed source, a stable +`connectionId`, one adapter revision, an endpoint path, an authorization mode, +and model mappings. The Host materializer resolves the source origin and +authorization into the private process configuration. Endpoint and credential +values are neither registration fields nor renderer-visible state. + +`connectionId` is the identity shared with Host profile synchronization. Host +`providerBindings` explicitly bind an existing managed connection to a target +profile and do not duplicate credentials. This package uses that identity +semantics but does not import Host-private synchronization contracts or create +a second persistent connection database. + +## Fail-closed behavior + +Required adapters are removed from model publication when missing, disabled, +or revision-mismatched. The native bridge remains loaded even when its active +plan is empty so known protected model IDs remain owned by its router and fail +inside the executor before an upstream request. They must not fall through to +another provider with the same model ID. + +The native executor removes caller authorization, host, content-length, and +adapter-declared untrusted headers before applying connection authorization. +Missing required sessions and missing credentials fail before `host.http.do` +or `host.http.do_stream`. + +## Verification boundary + +`npm run test:native` builds the current-platform C-shared bridge and loads it +through the real CLIProxyAPI plugin host from `CLIPROXY_SOURCE` (default +`/private/tmp/cliproxy-src`). Its synthetic loopback tests cover two connections +with the same source model and different credentials, header and body session +sources, leakage checks, streaming, cancellation, required-adapter failures, +and disabled-plan fallback protection. + +The checked-in runtime artifact is Darwin arm64. `scripts/build-native.mjs` +can build Darwin, Linux, or Windows on arm64 or x64 when run on that target, but +release artifacts must declare only binaries actually included in +`runtime-manifest.json`. This test is native ABI evidence, not a real provider, +native CordisX App, credential migration, or cross-platform acceptance test. diff --git a/.gitignore b/.gitignore index edfe37a..851d3a8 100644 --- a/.gitignore +++ b/.gitignore @@ -2,3 +2,4 @@ dist/ node_modules/ .cache/sdk/ +.cache/native-gateway-bridge/ diff --git a/AGENTS.md b/AGENTS.md index 475347c..c796d50 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -26,6 +26,9 @@ for the interaction contract and older-Host capability boundary. Dependency setup: [notification migration](./.agents/docs/notifications.md). +Gateway extension architecture and verification: +[gateway extensions](./.agents/docs/gateway-extensions.md). + ## Development and release - Requires Node.js 22 or newer. Install dependencies with `npm ci`. diff --git a/THIRD_PARTY_NOTICES.md b/THIRD_PARTY_NOTICES.md index 36be809..49a036a 100644 --- a/THIRD_PARTY_NOTICES.md +++ b/THIRD_PARTY_NOTICES.md @@ -1,7 +1,11 @@ # Third-party notices -This repository does not currently vendor third-party source or assets. +The native gateway bridge statically links `gopkg.in/yaml.v3` version 3.0.1. +Files ported from libyaml are licensed under the MIT License, copyright +2006-2011 Kirill Simonov. Remaining files are licensed under the Apache License +2.0, copyright 2011-2019 Canonical Ltd. The upstream license and notice are at +. Development dependencies retain their own licenses and notices. The future -runtime package must update this file if it begins to redistribute any +runtime package must update this file if it redistributes additional third-party code or media. diff --git a/config/cli-proxy-api.yaml b/config/cli-proxy-api.yaml index 6ce48fa..da7c5db 100644 --- a/config/cli-proxy-api.yaml +++ b/config/cli-proxy-api.yaml @@ -12,3 +12,10 @@ remote-management: disable-auto-update-panel: true codex-api-key: [] openai-compatibility: [] +plugins: + enabled: true + dir: runtime/cli-proxy-plugins + configs: + cordisx-gateway-bridge: + enabled: true + active: false diff --git a/cordisx-package.json b/cordisx-package.json index 5a65647..cc3cc8d 100644 --- a/cordisx-package.json +++ b/cordisx-package.json @@ -21,6 +21,6 @@ "runtimeManifest": { "path": "./runtime-manifest.json", "schema": "https://raw.githubusercontent.com/cordisx/cordisx-protocol/main/schemas/plugin-manifest.v14.schema.json", - "digest": "sha256:a66facc4315bf2b63941e94551ee982c32c034188bfbceb92135c779f0332dbf" + "digest": "sha256:485bcb6d6d89e2d41065f563b0ecba1d556693149374505ecc69f25d8cfb39f7" } } diff --git a/native/gateway-bridge/go.mod b/native/gateway-bridge/go.mod new file mode 100644 index 0000000..bf3c802 --- /dev/null +++ b/native/gateway-bridge/go.mod @@ -0,0 +1,5 @@ +module github.com/cordisx/plugin-cli-proxy-api/native/gateway-bridge + +go 1.22 + +require gopkg.in/yaml.v3 v3.0.1 diff --git a/native/gateway-bridge/go.sum b/native/gateway-bridge/go.sum new file mode 100644 index 0000000..a62c313 --- /dev/null +++ b/native/gateway-bridge/go.sum @@ -0,0 +1,4 @@ +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/native/gateway-bridge/main.go b/native/gateway-bridge/main.go new file mode 100644 index 0000000..ee3581b --- /dev/null +++ b/native/gateway-bridge/main.go @@ -0,0 +1,751 @@ +package main + +/* +#include +#include + +typedef struct { + void* ptr; + size_t len; +} cliproxy_buffer; + +typedef int (*cliproxy_host_call_fn)(void*, const char*, const uint8_t*, size_t, cliproxy_buffer*); +typedef void (*cliproxy_host_free_fn)(void*, size_t); + +typedef struct { + uint32_t abi_version; + void* host_ctx; + cliproxy_host_call_fn call; + cliproxy_host_free_fn free_buffer; +} cliproxy_host_api; + +typedef int (*cliproxy_plugin_call_fn)(char*, uint8_t*, size_t, cliproxy_buffer*); +typedef void (*cliproxy_plugin_free_fn)(void*, size_t); +typedef void (*cliproxy_plugin_shutdown_fn)(void); + +typedef struct { + uint32_t abi_version; + cliproxy_plugin_call_fn call; + cliproxy_plugin_free_fn free_buffer; + cliproxy_plugin_shutdown_fn shutdown; +} cliproxy_plugin_api; + +extern int cliproxyPluginCall(char*, uint8_t*, size_t, cliproxy_buffer*); +extern void cliproxyPluginFree(void*, size_t); +extern void cliproxyPluginShutdown(void); + +static const cliproxy_host_api* stored_host; + +static void store_host_api(const cliproxy_host_api* host) { + stored_host = host; +} + +static int call_host_api(const char* method, const uint8_t* request, size_t request_len, cliproxy_buffer* response) { + if (stored_host == NULL || stored_host->call == NULL) { + return 1; + } + return stored_host->call(stored_host->host_ctx, method, request, request_len, response); +} + +static void free_host_buffer(void* ptr, size_t len) { + if (stored_host != NULL && stored_host->free_buffer != NULL && ptr != NULL) { + stored_host->free_buffer(ptr, len); + } +} +*/ +import "C" + +import ( + "encoding/json" + "fmt" + "net/http" + "net/url" + "slices" + "strings" + "sync" + "unsafe" + + "gopkg.in/yaml.v3" +) + +const ( + abiVersion uint32 = 1 + pluginSchema uint32 = 3 + pluginID = "cordisx-gateway-bridge" + providerID = "cordisx-gateway-bridge" + sessionTemplate = "{{session}}" + connectionTemplate = "{{connectionId}}" + modelTemplate = "{{sourceModelId}}" +) + +type envelope struct { + OK bool `json:"ok"` + Result json.RawMessage `json:"result,omitempty"` + Error *envelopeError `json:"error,omitempty"` +} + +type envelopeError struct { + Code string `json:"code"` + Message string `json:"message"` + HTTPStatus int `json:"http_status,omitempty"` +} + +type lifecycleRequest struct { + ConfigYAML []byte `json:"config_yaml"` +} + +type sessionSource struct { + Kind string `yaml:"kind"` + Name string `yaml:"name"` + Pointer string `yaml:"pointer"` +} + +type adapter struct { + AdapterID string `yaml:"adapterId"` + Revision int64 `yaml:"revision"` + Enabled bool `yaml:"enabled"` + RequiredSession bool `yaml:"requiredSession"` + SessionSources []sessionSource `yaml:"sessionSources"` + Request struct { + ClearHeaders []string `yaml:"clearHeaders"` + SetHeaders map[string]any `yaml:"setHeaders"` + SetBody map[string]any `yaml:"setBody"` + } `yaml:"request"` +} + +type connectionModel struct { + SourceModelID string `yaml:"sourceModelId"` + GatewayModelID string `yaml:"gatewayModelId"` + DisplayName string `yaml:"displayName"` + IsDefault bool `yaml:"isDefault"` +} + +type connection struct { + ConnectionID string `yaml:"connectionId"` + Revision int64 `yaml:"revision"` + Enabled bool `yaml:"enabled"` + AdapterID string `yaml:"adapterId"` + AdapterRevision int64 `yaml:"adapterRevision"` + AdapterRequired bool `yaml:"adapterRequired"` + AdapterAvailable bool `yaml:"adapterAvailable"` + WireAPI string `yaml:"wireApi"` + EndpointPath string `yaml:"endpointPath"` + Endpoint string `yaml:"endpoint"` + Authorization string `yaml:"authorization"` + Credential string `yaml:"credential"` + Models []connectionModel `yaml:"models"` +} + +type bridgeConfig struct { + Enabled bool `yaml:"enabled"` + Active bool `yaml:"active"` + Revision string `yaml:"revision"` + Adapters []adapter `yaml:"adapters"` + Connections []connection `yaml:"connections"` +} + +type runtimePlan struct { + config bridgeConfig + adapters map[string]adapter + connections map[string]connection + models map[string]connectionModel +} + +type modelRouteRequest struct { + RequestedModel string `json:"RequestedModel"` +} + +type executorRequest struct { + Model string `json:"Model"` + Format string `json:"Format"` + Stream bool `json:"Stream"` + Headers map[string][]string `json:"Headers"` + Query url.Values `json:"Query"` + OriginalRequest []byte `json:"OriginalRequest"` + SourceFormat string `json:"SourceFormat"` + Payload []byte `json:"Payload"` + StreamID string `json:"stream_id"` + HostCallbackID string `json:"host_callback_id"` +} + +type hostHTTPResponse struct { + StatusCode int `json:"StatusCode"` + Headers map[string][]string `json:"Headers"` + Body []byte `json:"Body"` +} + +type hostHTTPStreamResponse struct { + StatusCode int `json:"status_code"` + Headers map[string][]string `json:"headers"` + StreamID string `json:"stream_id"` +} + +type hostHTTPStreamReadResponse struct { + Payload []byte `json:"payload"` + Error string `json:"error"` + Done bool `json:"done"` +} + +var state = struct { + sync.RWMutex + plan runtimePlan +}{plan: emptyPlan()} + +func main() {} + +//export cliproxy_plugin_init +func cliproxy_plugin_init(host *C.cliproxy_host_api, plugin *C.cliproxy_plugin_api) C.int { + if host == nil || plugin == nil || uint32(host.abi_version) != abiVersion { + return 1 + } + C.store_host_api(host) + plugin.abi_version = C.uint32_t(abiVersion) + plugin.call = C.cliproxy_plugin_call_fn(C.cliproxyPluginCall) + plugin.free_buffer = C.cliproxy_plugin_free_fn(C.cliproxyPluginFree) + plugin.shutdown = C.cliproxy_plugin_shutdown_fn(C.cliproxyPluginShutdown) + return 0 +} + +//export cliproxyPluginCall +func cliproxyPluginCall(method *C.char, request *C.uint8_t, requestLen C.size_t, response *C.cliproxy_buffer) C.int { + if response != nil { + response.ptr = nil + response.len = 0 + } + if method == nil { + writeResponse(response, errorEnvelope("invalid_method", "method is required", 0)) + return 1 + } + rawRequest := []byte(nil) + if request != nil && requestLen > 0 { + rawRequest = C.GoBytes(unsafe.Pointer(request), C.int(requestLen)) + } + raw, errHandle := handleMethod(C.GoString(method), rawRequest) + if errHandle != nil { + writeResponse(response, errorEnvelope("gateway_bridge_error", errHandle.Error(), httpStatus(errHandle))) + return 1 + } + writeResponse(response, raw) + return 0 +} + +//export cliproxyPluginFree +func cliproxyPluginFree(ptr unsafe.Pointer, length C.size_t) { + if ptr != nil { + C.free(ptr) + } + _ = length +} + +//export cliproxyPluginShutdown +func cliproxyPluginShutdown() { + state.Lock() + state.plan = emptyPlan() + state.Unlock() +} + +func handleMethod(method string, raw []byte) ([]byte, error) { + switch method { + case "plugin.register", "plugin.reconfigure": + if err := configure(raw); err != nil { + return nil, err + } + return okEnvelope(map[string]any{ + "schema_version": pluginSchema, + "metadata": map[string]any{ + "Name": "CordisX Gateway Bridge", "Version": "1.0.0", "Author": "CordisX", + "GitHubRepository": "https://github.com/cordisx/plugin-cli-proxy-api", "ConfigFields": []any{}, + }, + "capabilities": map[string]any{ + "model_provider": true, "model_router": true, "executor": true, + "executor_model_scope": "static", + "executor_input_formats": []string{"openai", "openai-response"}, + "executor_output_formats": []string{"openai", "openai-response"}, + }, + }) + case "executor.identifier": + return okEnvelope(map[string]string{"identifier": providerID}) + case "model.static": + return staticModels() + case "model.for_auth": + return okEnvelope(map[string]any{"Provider": providerID, "Models": []any{}}) + case "model.route": + return routeModel(raw) + case "executor.execute": + return execute(raw) + case "executor.execute_stream": + return executeStream(raw) + case "executor.count_tokens", "executor.http_request": + return nil, statusError{status: http.StatusNotImplemented, message: "operation is not supported by the gateway bridge"} + default: + return nil, fmt.Errorf("unknown method: %s", method) + } +} + +func emptyPlan() runtimePlan { + return runtimePlan{adapters: map[string]adapter{}, connections: map[string]connection{}, models: map[string]connectionModel{}} +} + +func configure(raw []byte) error { + var request lifecycleRequest + if err := json.Unmarshal(raw, &request); err != nil { + return fmt.Errorf("decode lifecycle request: %w", err) + } + var config bridgeConfig + if err := yaml.Unmarshal(request.ConfigYAML, &config); err != nil { + return fmt.Errorf("decode bridge configuration: %w", err) + } + plan := emptyPlan() + plan.config = config + for _, item := range config.Adapters { + if item.AdapterID == "" || item.Revision < 1 { + return fmt.Errorf("invalid adapter configuration") + } + plan.adapters[item.AdapterID] = item + } + for _, item := range config.Connections { + if item.ConnectionID == "" || item.EndpointPath == "" || item.Revision < 1 { + return fmt.Errorf("invalid connection configuration") + } + if item.Enabled && item.Endpoint == "" { + return fmt.Errorf("invalid connection configuration") + } + if _, exists := plan.connections[item.ConnectionID]; exists { + return fmt.Errorf("duplicate connection %s", item.ConnectionID) + } + plan.connections[item.ConnectionID] = item + for _, model := range item.Models { + if _, exists := plan.models[model.GatewayModelID]; exists { + return fmt.Errorf("duplicate gateway model %s", model.GatewayModelID) + } + plan.models[model.GatewayModelID] = model + } + } + state.Lock() + state.plan = plan + state.Unlock() + return nil +} + +func snapshot() runtimePlan { + state.RLock() + defer state.RUnlock() + return state.plan +} + +func staticModels() ([]byte, error) { + plan := snapshot() + ids := make([]string, 0, len(plan.models)) + for id := range plan.models { + ids = append(ids, id) + } + slices.Sort(ids) + models := make([]map[string]any, 0, len(ids)) + for _, id := range ids { + model := plan.models[id] + connection := plan.connections[strings.SplitN(id, "/", 2)[0]] + if !connectionReady(plan, connection) { + continue + } + models = append(models, map[string]any{ + "ID": id, "Object": "model", "OwnedBy": providerID, + "DisplayName": model.DisplayName, "Name": model.SourceModelID, + }) + } + return okEnvelope(map[string]any{"Provider": providerID, "Models": models}) +} + +func routeModel(raw []byte) ([]byte, error) { + var request modelRouteRequest + if err := json.Unmarshal(raw, &request); err != nil { + return nil, fmt.Errorf("decode route request: %w", err) + } + plan := snapshot() + _, exists := plan.models[request.RequestedModel] + if !exists { + return okEnvelope(map[string]any{"Handled": false}) + } + return okEnvelope(map[string]any{ + "Handled": true, "TargetKind": "self", "Reason": "cordisx protected connection", + }) +} + +func connectionReady(plan runtimePlan, connection connection) bool { + if !plan.config.Active || !connection.Enabled || connection.ConnectionID == "" { + return false + } + adapter, exists := plan.adapters[connection.AdapterID] + return !connection.AdapterRequired || (exists && adapter.Enabled && adapter.Revision == connection.AdapterRevision && connection.AdapterAvailable) +} + +func execute(raw []byte) ([]byte, error) { + request, connection, adapter, model, err := prepare(raw) + if err != nil { + return nil, err + } + upstream, err := outboundRequest(request, connection, adapter, model) + if err != nil { + return nil, err + } + var response hostHTTPResponse + if err := callHost("host.http.do", upstream, &response); err != nil { + return nil, err + } + if response.StatusCode < 200 || response.StatusCode >= 300 { + return nil, statusError{status: response.StatusCode, message: "gateway upstream rejected request"} + } + return okEnvelope(map[string]any{"Payload": response.Body, "Headers": response.Headers}) +} + +func executeStream(raw []byte) ([]byte, error) { + request, connection, adapter, model, err := prepare(raw) + if err != nil { + return nil, err + } + if request.StreamID == "" { + return nil, fmt.Errorf("stream id is required") + } + upstream, err := outboundRequest(request, connection, adapter, model) + if err != nil { + return nil, err + } + var response hostHTTPStreamResponse + if err := callHost("host.http.do_stream", upstream, &response); err != nil { + return nil, err + } + if response.StatusCode < 200 || response.StatusCode >= 300 { + _ = callHost("host.http.stream_close", map[string]string{"stream_id": response.StreamID}, nil) + return nil, statusError{status: response.StatusCode, message: "gateway upstream rejected stream"} + } + go forwardStream(request.StreamID, response.StreamID) + return okEnvelope(map[string]any{"headers": response.Headers}) +} + +func prepare(raw []byte) (executorRequest, connection, adapter, connectionModel, error) { + var request executorRequest + if err := json.Unmarshal(raw, &request); err != nil { + return request, connection{}, adapter{}, connectionModel{}, fmt.Errorf("decode executor request: %w", err) + } + plan := snapshot() + model, exists := plan.models[request.Model] + if !exists { + return request, connection{}, adapter{}, model, statusError{status: http.StatusNotFound, message: "gateway model is not registered"} + } + connectionID := strings.SplitN(model.GatewayModelID, "/", 2)[0] + connection := plan.connections[connectionID] + if !connectionReady(plan, connection) { + return request, connection, adapter{}, model, statusError{status: http.StatusServiceUnavailable, message: "required gateway adapter is unavailable"} + } + selectedAdapter, exists := plan.adapters[connection.AdapterID] + if connection.AdapterRequired && (!exists || !selectedAdapter.Enabled || selectedAdapter.Revision != connection.AdapterRevision) { + return request, connection, adapter{}, model, statusError{status: http.StatusServiceUnavailable, message: "gateway adapter is unavailable"} + } + if !exists || !selectedAdapter.Enabled || selectedAdapter.Revision != connection.AdapterRevision { + selectedAdapter = adapter{} + } + return request, connection, selectedAdapter, model, nil +} + +func outboundRequest(request executorRequest, connection connection, adapter adapter, model connectionModel) (map[string]any, error) { + body := request.Payload + if len(body) == 0 { + body = request.OriginalRequest + } + var document any + if err := json.Unmarshal(body, &document); err != nil { + return nil, statusError{status: http.StatusBadRequest, message: "gateway request body must be JSON"} + } + session := sessionValue(adapter.SessionSources, request.Headers, document) + if adapter.RequiredSession && session == "" { + return nil, statusError{status: http.StatusBadRequest, message: "required gateway session is missing"} + } + if err := setPointer(&document, "/model", model.SourceModelID); err != nil { + return nil, err + } + vars := map[string]string{sessionTemplate: session, connectionTemplate: connection.ConnectionID, modelTemplate: model.SourceModelID} + for pointer, template := range adapter.Request.SetBody { + if err := setPointer(&document, pointer, expandTemplate(template, vars)); err != nil { + return nil, statusError{status: http.StatusBadRequest, message: "invalid gateway body transform"} + } + } + encoded, err := json.Marshal(document) + if err != nil { + return nil, fmt.Errorf("encode gateway request: %w", err) + } + headers := cloneHeaders(request.Headers) + deleteHeader(headers, "Authorization") + deleteHeader(headers, "Host") + deleteHeader(headers, "Content-Length") + for _, name := range adapter.Request.ClearHeaders { + deleteHeader(headers, name) + } + for name, template := range adapter.Request.SetHeaders { + value := expandTemplate(template, vars) + text, ok := value.(string) + if !ok { + raw, err := json.Marshal(value) + if err != nil { + return nil, fmt.Errorf("encode gateway header transform: %w", err) + } + text = string(raw) + } + headers[http.CanonicalHeaderKey(name)] = []string{text} + } + if connection.Authorization == "bearer" { + if connection.Credential == "" { + return nil, statusError{status: http.StatusServiceUnavailable, message: "gateway connection credential is unavailable"} + } + headers["Authorization"] = []string{"Bearer " + connection.Credential} + } + headers["Content-Type"] = []string{"application/json"} + target, err := joinEndpoint(connection.Endpoint, connection.EndpointPath) + if err != nil { + return nil, statusError{status: http.StatusServiceUnavailable, message: "gateway connection endpoint is invalid"} + } + return map[string]any{ + "host_callback_id": request.HostCallbackID, + "method": http.MethodPost, + "url": target, + "headers": headers, + "body": encoded, + }, nil +} + +func forwardStream(pluginStreamID, upstreamStreamID string) { + defer func() { + _ = callHost("host.http.stream_close", map[string]string{"stream_id": upstreamStreamID}, nil) + _ = callHost("host.stream.close", map[string]string{"stream_id": pluginStreamID}, nil) + }() + for { + var chunk hostHTTPStreamReadResponse + if err := callHost("host.http.stream_read", map[string]string{"stream_id": upstreamStreamID}, &chunk); err != nil { + _ = callHost("host.stream.close", map[string]string{"stream_id": pluginStreamID, "error": err.Error()}, nil) + return + } + if chunk.Error != "" { + _ = callHost("host.stream.close", map[string]string{"stream_id": pluginStreamID, "error": chunk.Error}, nil) + return + } + if len(chunk.Payload) > 0 { + if err := callHost("host.stream.emit", map[string]any{"stream_id": pluginStreamID, "payload": chunk.Payload}, nil); err != nil { + return + } + } + if chunk.Done { + return + } + } +} + +func sessionValue(sources []sessionSource, headers map[string][]string, body any) string { + for _, source := range sources { + if source.Kind == "header" { + for name, values := range headers { + if strings.EqualFold(name, source.Name) && len(values) > 0 && strings.TrimSpace(values[0]) != "" { + return strings.TrimSpace(values[0]) + } + } + continue + } + if source.Kind == "body-json-pointer" { + if value, ok := getPointer(body, source.Pointer).(string); ok && strings.TrimSpace(value) != "" { + return strings.TrimSpace(value) + } + } + } + return "" +} + +func expandTemplate(value any, vars map[string]string) any { + switch item := value.(type) { + case string: + if replacement, ok := vars[item]; ok { + return replacement + } + for token, replacement := range vars { + item = strings.ReplaceAll(item, token, replacement) + } + return item + case []any: + out := make([]any, len(item)) + for index, child := range item { + out[index] = expandTemplate(child, vars) + } + return out + case map[string]any: + out := make(map[string]any, len(item)) + for key, child := range item { + out[key] = expandTemplate(child, vars) + } + return out + default: + return value + } +} + +func pointerTokens(pointer string) ([]string, error) { + if pointer == "" || pointer[0] != '/' { + return nil, fmt.Errorf("invalid JSON pointer") + } + parts := strings.Split(pointer[1:], "/") + for index := range parts { + parts[index] = strings.ReplaceAll(strings.ReplaceAll(parts[index], "~1", "/"), "~0", "~") + } + return parts, nil +} + +func getPointer(root any, pointer string) any { + parts, err := pointerTokens(pointer) + if err != nil { + return nil + } + current := root + for _, part := range parts { + object, ok := current.(map[string]any) + if !ok { + return nil + } + current, ok = object[part] + if !ok { + return nil + } + } + return current +} + +func setPointer(root *any, pointer string, value any) error { + parts, err := pointerTokens(pointer) + if err != nil { + return err + } + current, ok := (*root).(map[string]any) + if !ok { + return fmt.Errorf("request body is not an object") + } + for _, part := range parts[:len(parts)-1] { + next, exists := current[part] + if !exists { + child := map[string]any{} + current[part] = child + current = child + continue + } + child, ok := next.(map[string]any) + if !ok { + return fmt.Errorf("JSON pointer crosses a non-object value") + } + current = child + } + current[parts[len(parts)-1]] = value + return nil +} + +func joinEndpoint(origin, endpointPath string) (string, error) { + base, err := url.Parse(origin) + if err != nil || base.Scheme == "" || base.Host == "" { + return "", fmt.Errorf("invalid origin") + } + reference, err := url.Parse(strings.TrimPrefix(endpointPath, "/")) + if err != nil { + return "", err + } + if !strings.HasSuffix(base.Path, "/") { + base.Path += "/" + } + return base.ResolveReference(reference).String(), nil +} + +func cloneHeaders(source map[string][]string) map[string][]string { + out := make(map[string][]string, len(source)) + for name, values := range source { + out[http.CanonicalHeaderKey(name)] = append([]string(nil), values...) + } + return out +} + +func deleteHeader(headers map[string][]string, target string) { + for name := range headers { + if strings.EqualFold(name, target) { + delete(headers, name) + } + } +} + +func callHost(method string, request any, response any) error { + raw, err := json.Marshal(request) + if err != nil { + return err + } + cMethod := C.CString(method) + defer C.free(unsafe.Pointer(cMethod)) + var cRequest unsafe.Pointer + if len(raw) > 0 { + cRequest = C.CBytes(raw) + defer C.free(cRequest) + } + var output C.cliproxy_buffer + rc := C.call_host_api(cMethod, (*C.uint8_t)(cRequest), C.size_t(len(raw)), &output) + var result []byte + if output.ptr != nil && output.len > 0 { + result = C.GoBytes(output.ptr, C.int(output.len)) + } + if output.ptr != nil { + C.free_host_buffer(output.ptr, output.len) + } + if rc != 0 { + return fmt.Errorf("host callback %s failed", method) + } + var wrapper envelope + if err := json.Unmarshal(result, &wrapper); err != nil { + return fmt.Errorf("decode host callback %s: %w", method, err) + } + if !wrapper.OK { + if wrapper.Error != nil { + return statusError{status: wrapper.Error.HTTPStatus, message: wrapper.Error.Message} + } + return fmt.Errorf("host callback %s failed", method) + } + if response == nil || len(wrapper.Result) == 0 { + return nil + } + return json.Unmarshal(wrapper.Result, response) +} + +type statusError struct { + status int + message string +} + +func (error_ statusError) Error() string { return error_.message } + +func httpStatus(err error) int { + if value, ok := err.(statusError); ok { + return value.status + } + return 0 +} + +func okEnvelope(result any) ([]byte, error) { + raw, err := json.Marshal(result) + if err != nil { + return nil, err + } + return json.Marshal(envelope{OK: true, Result: raw}) +} + +func errorEnvelope(code, message string, status int) []byte { + raw, _ := json.Marshal(envelope{OK: false, Error: &envelopeError{Code: code, Message: message, HTTPStatus: status}}) + return raw +} + +func writeResponse(response *C.cliproxy_buffer, raw []byte) { + if response == nil || len(raw) == 0 { + return + } + ptr := C.CBytes(raw) + if ptr == nil { + return + } + response.ptr = ptr + response.len = C.size_t(len(raw)) +} diff --git a/package.json b/package.json index abfd20c..4dcb8a7 100644 --- a/package.json +++ b/package.json @@ -9,6 +9,7 @@ "dist", "assets", "config", + "runtime", "schemas", "cordisx-package.json", "runtime-manifest.json", @@ -34,10 +35,15 @@ "types": "./dist/types/service/gateway.d.ts", "import": "./dist/gateway.mjs", "default": "./dist/gateway.mjs" + }, + "./extensions/v1": { + "types": "./dist/types/service/extension-registry.d.ts", + "import": "./dist/extensions.mjs", + "default": "./dist/extensions.mjs" } }, "scripts": { - "build": "rm -rf dist && node scripts/sync-package-manifest.mjs && vite build && tsc -p tsconfig.json --emitDeclarationOnly && node scripts/build-service.mjs", + "build": "rm -rf dist && node scripts/build-native.mjs && node scripts/sync-runtime-resources.mjs && node scripts/sync-package-manifest.mjs && vite build && tsc -p tsconfig.json --emitDeclarationOnly && node scripts/build-service.mjs", "check": "npm run typecheck && npm run build && npm run typecheck:consumer && npm test && npm run lint && npm run stylelint && npm run format:check", "dev": "cordisx dev --config cordisx.config.json", "dev:dry-run": "cordisx dev --config cordisx.config.json --dry-run", @@ -46,6 +52,7 @@ "lint": "eslint .", "stylelint": "stylelint \"src/**/*.css\" --allow-empty-input", "test": "node --test test/*.mjs", + "test:native": "node scripts/build-native.mjs && node scripts/test-native-integration.mjs", "typecheck": "tsc -p tsconfig.json --noEmit", "typecheck:consumer": "tsc -p test/upstream-consumer/tsconfig.json --noEmit", "prepare": "npm run build" diff --git a/runtime-manifest.json b/runtime-manifest.json index 4f23395..d92b6f5 100644 --- a/runtime-manifest.json +++ b/runtime-manifest.json @@ -69,8 +69,8 @@ { "path": "./config/cli-proxy-api.yaml", "mode": "data", - "digest": "sha256:4e2806a2c9a4a490f8685ffc96cf1dc40010248a8d25437b718249f3275b198e", - "byteLength": 268 + "digest": "sha256:1dee049ec7bc4b9ff9a6bf77e8d8a88136045d3fd9d25ac15acb2925502e424c", + "byteLength": 405 }, { "path": "./schemas/cli-proxy-gateway-model-catalog.v1.schema.json", @@ -90,6 +90,24 @@ "digest": "sha256:8f226eeaf778394f452c0c224bd163c1c68806f56c63982d0394c03d6be4483b", "byteLength": 1925 }, + { + "path": "./schemas/cli-proxy-gateway-adapter.v1.schema.json", + "mode": "data", + "digest": "sha256:178a7ad887b0400fdea8e143eb748e51bf491fff92381a9de186c9b41cf6a5a5", + "byteLength": 3011 + }, + { + "path": "./schemas/cli-proxy-gateway-connection.v1.schema.json", + "mode": "data", + "digest": "sha256:bbe81ec1a835586e67f7c9e7565a888a8120e55cbb1cc0a8d299b36831173abf", + "byteLength": 2371 + }, + { + "path": "./schemas/cli-proxy-gateway-extension-plan.v1.schema.json", + "mode": "data", + "digest": "sha256:398b8f8b05e16a39808f20829edc4f641f33409030faed18ab32bc2e420e05d1", + "byteLength": 2761 + }, { "path": "./schemas/cli-proxy-management-auth-files.v1.schema.json", "mode": "data", @@ -101,6 +119,18 @@ "mode": "data", "digest": "sha256:c1a9e12cf48ee53825948b053ea364ce70d09122321502560defebeb3d8d0c12", "byteLength": 218 + }, + { + "path": "./runtime/cli-proxy-plugins/darwin/arm64/cordisx-gateway-bridge.dylib", + "mode": "executable", + "digest": "sha256:da5786f5347f39f7ec1b19b60655a899885e884238a95a20b01d3b74fa1c2882", + "byteLength": 4553202, + "platforms": [ + "darwin" + ], + "architectures": [ + "arm64" + ] } ], "consumerGrants": [ diff --git a/runtime/cli-proxy-plugins/darwin/arm64/cordisx-gateway-bridge.dylib b/runtime/cli-proxy-plugins/darwin/arm64/cordisx-gateway-bridge.dylib new file mode 100644 index 0000000..54353e1 Binary files /dev/null and b/runtime/cli-proxy-plugins/darwin/arm64/cordisx-gateway-bridge.dylib differ diff --git a/schemas/cli-proxy-gateway-adapter.v1.schema.json b/schemas/cli-proxy-gateway-adapter.v1.schema.json new file mode 100644 index 0000000..64c5b78 --- /dev/null +++ b/schemas/cli-proxy-gateway-adapter.v1.schema.json @@ -0,0 +1,94 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-adapter.v1.schema.json", + "title": "CordisX CLIProxy Gateway Adapter v1", + "type": "object", + "additionalProperties": false, + "required": [ + "$schema", + "contract", + "schemaVersion", + "adapterId", + "revision", + "enabled", + "requiredSession", + "sessionSources", + "request" + ], + "properties": { + "$schema": { + "const": "https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-adapter.v1.schema.json" + }, + "contract": { "const": "cordisx.cli-proxy-gateway-adapter/v1" }, + "schemaVersion": { "const": 1 }, + "adapterId": { "$ref": "#/$defs/id" }, + "revision": { "type": "integer", "minimum": 1 }, + "enabled": { "type": "boolean" }, + "requiredSession": { "type": "boolean" }, + "sessionSources": { + "type": "array", + "maxItems": 16, + "items": { + "oneOf": [ + { + "type": "object", + "additionalProperties": false, + "required": ["kind", "name"], + "properties": { + "kind": { "const": "header" }, + "name": { "$ref": "#/$defs/header" } + } + }, + { + "type": "object", + "additionalProperties": false, + "required": ["kind", "pointer"], + "properties": { + "kind": { "const": "body-json-pointer" }, + "pointer": { "$ref": "#/$defs/pointer" } + } + } + ] + } + }, + "request": { + "type": "object", + "additionalProperties": false, + "properties": { + "clearHeaders": { + "type": "array", + "maxItems": 64, + "uniqueItems": true, + "items": { "$ref": "#/$defs/header" } + }, + "setHeaders": { + "type": "object", + "maxProperties": 64, + "propertyNames": { "$ref": "#/$defs/header" }, + "additionalProperties": { "$ref": "#/$defs/template" } + }, + "setBody": { + "type": "object", + "maxProperties": 64, + "propertyNames": { "$ref": "#/$defs/pointer" }, + "additionalProperties": { "$ref": "#/$defs/template" } + } + } + } + }, + "$defs": { + "id": { "type": "string", "pattern": "^[a-z0-9][a-z0-9._-]{0,95}$" }, + "header": { "type": "string", "pattern": "^[!#$%&'*+.^_`|~0-9A-Za-z-]{1,128}$" }, + "pointer": { "type": "string", "pattern": "^(?:/(?:[^~/]|~0|~1)*)+$", "maxLength": 512 }, + "template": { + "oneOf": [ + { "type": "null" }, + { "type": "boolean" }, + { "type": "number" }, + { "type": "string", "maxLength": 8192 }, + { "type": "array", "maxItems": 128, "items": { "$ref": "#/$defs/template" } }, + { "type": "object", "maxProperties": 128, "additionalProperties": { "$ref": "#/$defs/template" } } + ] + } + } +} diff --git a/schemas/cli-proxy-gateway-connection.v1.schema.json b/schemas/cli-proxy-gateway-connection.v1.schema.json new file mode 100644 index 0000000..0fa51ff --- /dev/null +++ b/schemas/cli-proxy-gateway-connection.v1.schema.json @@ -0,0 +1,66 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-connection.v1.schema.json", + "title": "CordisX CLIProxy Gateway Connection v1", + "type": "object", + "additionalProperties": false, + "required": [ + "$schema", + "contract", + "schemaVersion", + "connectionId", + "revision", + "enabled", + "source", + "wireApi", + "endpointPath", + "authorization", + "adapter", + "models" + ], + "properties": { + "$schema": { + "const": "https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-connection.v1.schema.json" + }, + "contract": { "const": "cordisx.cli-proxy-gateway-connection/v1" }, + "schemaVersion": { "const": 1 }, + "connectionId": { "$ref": "#/$defs/id" }, + "revision": { "type": "integer", "minimum": 1 }, + "enabled": { "type": "boolean" }, + "source": {}, + "wireApi": { "enum": ["responses", "chat-completions"] }, + "endpointPath": { "type": "string", "pattern": "^/(?!/)[^?#\\u0000-\\u001f\\u007f]{0,511}$" }, + "authorization": { "enum": ["none", "bearer"] }, + "adapter": { + "type": "object", + "additionalProperties": false, + "required": ["adapterId", "revision", "required"], + "properties": { + "adapterId": { "$ref": "#/$defs/id" }, + "revision": { "type": "integer", "minimum": 1 }, + "required": { "type": "boolean" } + } + }, + "models": { + "type": "array", + "minItems": 1, + "maxItems": 256, + "items": { + "type": "object", + "additionalProperties": false, + "required": ["sourceModelId", "modelId", "enabled", "isDefault"], + "properties": { + "sourceModelId": { "$ref": "#/$defs/model" }, + "modelId": { "type": "string", "minLength": 1, "maxLength": 256, "pattern": "^[^/\\u0000-\\u001f\\u007f]+$" }, + "displayName": { "type": "string", "minLength": 1, "maxLength": 200, "pattern": "^\\S(?:.*\\S)?$" }, + "enabled": { "type": "boolean" }, + "isDefault": { "type": "boolean" } + } + } + } + }, + "$defs": { + "id": { "type": "string", "pattern": "^[a-z0-9][a-z0-9._-]{0,95}$" }, + "model": { "type": "string", "minLength": 1, "maxLength": 256, "pattern": "^[^\\u0000-\\u001f\\u007f]+$" } + } +} diff --git a/schemas/cli-proxy-gateway-extension-plan.v1.schema.json b/schemas/cli-proxy-gateway-extension-plan.v1.schema.json new file mode 100644 index 0000000..d0fd94a --- /dev/null +++ b/schemas/cli-proxy-gateway-extension-plan.v1.schema.json @@ -0,0 +1,71 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-extension-plan.v1.schema.json", + "title": "CLIProxyAPI gateway bridge configuration v1", + "type": "object", + "additionalProperties": false, + "required": ["enabled", "active", "revision", "adapters", "connections"], + "properties": { + "enabled": { "const": true }, + "active": { "type": "boolean" }, + "revision": { "type": "string", "pattern": "^sha256:[a-f0-9]{64}$" }, + "adapters": { + "type": "array", + "maxItems": 64, + "items": { "$ref": "cli-proxy-gateway-adapter.v1.schema.json" } + }, + "connections": { + "type": "array", + "maxItems": 128, + "items": { + "type": "object", + "additionalProperties": false, + "required": [ + "connectionId", + "revision", + "enabled", + "adapterId", + "adapterRevision", + "adapterRequired", + "adapterAvailable", + "wireApi", + "endpointPath", + "endpoint", + "authorization", + "credential", + "models" + ], + "properties": { + "connectionId": { "type": "string", "pattern": "^[a-z0-9][a-z0-9._-]{0,95}$" }, + "revision": { "type": "integer", "minimum": 1 }, + "enabled": { "type": "boolean" }, + "adapterId": { "type": "string", "pattern": "^[a-z0-9][a-z0-9._-]{0,95}$" }, + "adapterRevision": { "type": "integer", "minimum": 1 }, + "adapterRequired": { "type": "boolean" }, + "adapterAvailable": { "type": "boolean" }, + "wireApi": { "enum": ["responses", "chat-completions"] }, + "endpointPath": { "type": "string", "pattern": "^/(?!/)[^?#\\u0000-\\u001f\\u007f]{0,511}$" }, + "endpoint": { "type": ["string", "null"], "format": "uri" }, + "authorization": { "enum": ["none", "bearer"] }, + "credential": { "type": ["string", "null"], "maxLength": 16384 }, + "models": { + "type": "array", + "minItems": 1, + "maxItems": 256, + "items": { + "type": "object", + "additionalProperties": false, + "required": ["sourceModelId", "gatewayModelId", "displayName", "isDefault"], + "properties": { + "sourceModelId": { "type": "string", "minLength": 1, "maxLength": 256 }, + "gatewayModelId": { "type": "string", "minLength": 3, "maxLength": 353 }, + "displayName": { "type": "string", "minLength": 1, "maxLength": 200 }, + "isDefault": { "type": "boolean" } + } + } + } + } + } + } + } +} diff --git a/scripts/build-native.mjs b/scripts/build-native.mjs new file mode 100644 index 0000000..4ce0112 --- /dev/null +++ b/scripts/build-native.mjs @@ -0,0 +1,23 @@ +import { execFile } from 'node:child_process' +import { mkdir, rm } from 'node:fs/promises' +import path from 'node:path' +import { promisify } from 'node:util' +import { fileURLToPath } from 'node:url' + +const execFileAsync = promisify(execFile) +const root = fileURLToPath(new URL('../', import.meta.url)) +const source = path.join(root, 'native/gateway-bridge') +const extension = process.platform === 'darwin' ? '.dylib' : process.platform === 'win32' ? '.dll' : '.so' +const directory = path.join(root, 'runtime/cli-proxy-plugins', process.platform, process.arch) +const output = path.join(directory, `cordisx-gateway-bridge${extension}`) + +if (!['darwin', 'linux', 'win32'].includes(process.platform) || !['arm64', 'x64'].includes(process.arch)) { + throw new Error(`Unsupported native gateway bridge target ${process.platform}/${process.arch}`) +} + +await mkdir(directory, { recursive: true }) +await execFileAsync('go', ['build', '-buildmode=c-shared', '-trimpath', '-o', output, '.'], { + cwd: source, + env: { ...process.env, CGO_ENABLED: '1' }, +}) +await rm(output.replace(new RegExp(`${extension.replace('.', '\\.')}$$`), '.h'), { force: true }) diff --git a/scripts/build-service.mjs b/scripts/build-service.mjs index 800f13f..094edf2 100644 --- a/scripts/build-service.mjs +++ b/scripts/build-service.mjs @@ -6,6 +6,7 @@ await build({ entryPoints: { service: new URL('../src/service/index.ts', import.meta.url).pathname, gateway: new URL('../src/service/gateway.ts', import.meta.url).pathname, + extensions: new URL('../src/service/extension-registry.ts', import.meta.url).pathname, }, outdir: new URL('../dist/', import.meta.url).pathname, outExtension: { '.js': '.mjs' }, diff --git a/scripts/sync-runtime-resources.mjs b/scripts/sync-runtime-resources.mjs new file mode 100644 index 0000000..c209276 --- /dev/null +++ b/scripts/sync-runtime-resources.mjs @@ -0,0 +1,51 @@ +import { createHash } from 'node:crypto' +import { readdir, readFile, writeFile } from 'node:fs/promises' +import path from 'node:path' +import { fileURLToPath } from 'node:url' + +const root = fileURLToPath(new URL('../', import.meta.url)) +const runtimeManifestPath = path.join(root, 'runtime-manifest.json') +const runtimeManifest = JSON.parse(await readFile(runtimeManifestPath, 'utf8')) +const gateway = runtimeManifest.services.find(service => service.id === 'gateway-runtime') +if (gateway === undefined) throw new Error('gateway-runtime service is missing') + +const dataPaths = [ + './config/cli-proxy-api.yaml', + './schemas/cli-proxy-gateway-model-catalog.v1.schema.json', + './schemas/cli-proxy-gateway-codex-upstream.v1.schema.json', + './schemas/cli-proxy-gateway-openai-upstream.v1.schema.json', + './schemas/cli-proxy-gateway-adapter.v1.schema.json', + './schemas/cli-proxy-gateway-connection.v1.schema.json', + './schemas/cli-proxy-gateway-extension-plan.v1.schema.json', + './schemas/cli-proxy-management-auth-files.v1.schema.json', + './schemas/cli-proxy-management-operation.v1.schema.json', +] + +const resources = [] +for (const relative of dataPaths) resources.push(await declaration(relative, 'data')) +const nativeRoot = path.join(root, 'runtime/cli-proxy-plugins') +for (const platform of await readdir(nativeRoot).catch(() => [])) { + for (const architecture of await readdir(path.join(nativeRoot, platform)).catch(() => [])) { + const directory = path.join(nativeRoot, platform, architecture) + for (const file of await readdir(directory)) { + if (!/^cordisx-gateway-bridge(?:\.dylib|\.so|\.dll)$/.test(file)) continue + resources.push({ + ...await declaration(`./runtime/cli-proxy-plugins/${platform}/${architecture}/${file}`, 'executable'), + platforms: [platform], + architectures: [architecture], + }) + } + } +} +gateway.runtimeResources = resources +await writeFile(runtimeManifestPath, `${JSON.stringify(runtimeManifest, null, 2)}\n`) + +async function declaration(relative, mode) { + const bytes = await readFile(path.join(root, relative.slice(2))) + return { + path: relative, + mode, + digest: `sha256:${createHash('sha256').update(bytes).digest('hex')}`, + byteLength: bytes.byteLength, + } +} diff --git a/scripts/test-native-integration.mjs b/scripts/test-native-integration.mjs new file mode 100644 index 0000000..7f641a9 --- /dev/null +++ b/scripts/test-native-integration.mjs @@ -0,0 +1,31 @@ +import { execFile } from 'node:child_process' +import { copyFile, mkdir, readFile, rm, writeFile } from 'node:fs/promises' +import path from 'node:path' +import { promisify } from 'node:util' +import { fileURLToPath } from 'node:url' + +const execFileAsync = promisify(execFile) +const root = fileURLToPath(new URL('../', import.meta.url)) +const source = process.env.CLIPROXY_SOURCE ?? '/private/tmp/cliproxy-src' +const fixture = path.join(root, 'test/native-gateway-bridge') +const work = path.join(root, '.cache/native-gateway-bridge') + +await rm(work, { recursive: true, force: true }) +await mkdir(work, { recursive: true }) +await Promise.all([ + copyFile(path.join(fixture, 'gateway_bridge_test.go'), path.join(work, 'gateway_bridge_test.go')), + copyFile(path.join(fixture, 'go.mod'), path.join(work, 'go.mod')), +]) +const moduleText = await readFile(path.join(work, 'go.mod'), 'utf8') +await writeFile( + path.join(work, 'go.mod'), + `${moduleText.trimEnd()}\n\nreplace github.com/router-for-me/CLIProxyAPI/v7 => ${source}\n`, +) +await execFileAsync('go', ['test', '-mod=mod', '-count=1', '-timeout', '90s', '-v', '.'], { + cwd: work, + env: { + ...process.env, + CORDISX_GATEWAY_REPO: root, + }, + maxBuffer: 16 * 1024 * 1024, +}) diff --git a/src/service/extension-context.ts b/src/service/extension-context.ts new file mode 100644 index 0000000..e774b83 --- /dev/null +++ b/src/service/extension-context.ts @@ -0,0 +1,71 @@ +import type { + ManagedServiceContextAuthorityV1, + ManagedServiceContextBindingV1, + ManagedServiceContextProviderRootV1, + ManagedServiceContextProviderV1, +} from '@cordisx/protocol/managed-service-context/v1' +import type { ManagedServiceIdentityV1 } from '@cordisx/protocol/managed-service-runtime/v1' +import { + CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1, + CliProxyGatewayExtensionRegistryV1, + createCliProxyGatewayExtensionRegistrarV1, +} from './extension-registry.js' + +function sourceKey(source: ManagedServiceIdentityV1): string { + return JSON.stringify([source.source, source.pluginId, source.serviceId]) +} + +export function createCliProxyGatewayExtensionContextProviderV1(): ManagedServiceContextProviderV1 { + return Object.freeze({ + service: CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1, + create: ({ target, signal }: { + readonly target: ManagedServiceContextAuthorityV1 + readonly signal: AbortSignal + }): ManagedServiceContextProviderRootV1 => { + const registry = new CliProxyGatewayExtensionRegistryV1( + target.pluginGeneration, + sourceKey, + ) + const bindings = new Set() + let disposed = false + const dispose = (): void => { + if (disposed) return + disposed = true + for (const binding of [...bindings]) binding.dispose() + bindings.clear() + registry.dispose() + } + signal.addEventListener('abort', dispose, { once: true }) + return Object.freeze({ + providerValue: registry, + bind: ({ producer, target: bindingTarget, signal: bindingSignal }: { + readonly producer: { readonly pluginId: string; readonly pluginGeneration: string } + readonly target: ManagedServiceContextAuthorityV1 + readonly signal: AbortSignal + }) => { + if ( + disposed || bindingSignal.aborted || bindingTarget.pluginId !== target.pluginId + || bindingTarget.pluginGeneration !== target.pluginGeneration + || bindingTarget.serviceId !== target.serviceId + || bindingTarget.serviceGeneration !== target.serviceGeneration + ) throw new Error('CLIProxyAPI gateway extension context generation is stale') + const registrar = createCliProxyGatewayExtensionRegistrarV1(registry, producer) + let bindingDisposed = false + const binding: ManagedServiceContextBindingV1 = Object.freeze({ + value: registrar, + dispose: () => { + if (bindingDisposed) return + bindingDisposed = true + bindings.delete(binding) + registrar.dispose() + }, + }) + bindings.add(binding) + bindingSignal.addEventListener('abort', () => binding.dispose(), { once: true }) + return binding + }, + dispose, + }) + }, + }) +} diff --git a/src/service/extension-registry.ts b/src/service/extension-registry.ts new file mode 100644 index 0000000..f729deb --- /dev/null +++ b/src/service/extension-registry.ts @@ -0,0 +1,677 @@ +import { createHash } from 'node:crypto' + +export const CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1 = + 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-adapter.v1.schema.json' +export const CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1 = + 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-connection.v1.schema.json' +export const CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1 = 'cordisx.cli-proxy-gateway-adapter/v1' +export const CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1 = 'cordisx.cli-proxy-gateway-connection/v1' +export const CLI_PROXY_GATEWAY_EXTENSION_PLAN_CONTRACT_V1 = 'cordisx.cli-proxy-gateway-extension-plan/v1' +export const CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1 = 'cliProxyGatewayExtensions' + +const ID_PATTERN = /^[a-z0-9][a-z0-9._-]{0,95}$/ +const HEADER_PATTERN = /^[!#$%&'*+.^_`|~0-9A-Za-z-]{1,128}$/ +const POINTER_PATTERN = /^(?:\/(?:[^~/]|~0|~1)*)+$/ +const MAX_ADAPTERS = 64 +const MAX_CONNECTIONS = 128 +const MAX_MODELS = 256 + +export type CliProxyGatewayWireApiV1 = 'responses' | 'chat-completions' +export type CliProxyGatewayTemplateV1 = + | null + | boolean + | number + | string + | readonly CliProxyGatewayTemplateV1[] + | { readonly [key: string]: CliProxyGatewayTemplateV1 } + +export type CliProxyGatewaySessionSourceV1 = + | { readonly kind: 'header'; readonly name: string } + | { readonly kind: 'body-json-pointer'; readonly pointer: `/${string}` } + +export interface CliProxyGatewayAdapterV1 { + readonly $schema: typeof CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1 + readonly contract: typeof CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1 + readonly schemaVersion: 1 + readonly adapterId: string + readonly revision: number + readonly enabled: boolean + readonly requiredSession: boolean + readonly sessionSources: readonly CliProxyGatewaySessionSourceV1[] + readonly request: { + readonly clearHeaders?: readonly string[] + readonly setHeaders?: Readonly> + readonly setBody?: Readonly> + } +} + +export interface CliProxyGatewayConnectionModelV1 { + readonly sourceModelId: string + readonly modelId: string + readonly displayName?: string + readonly enabled: boolean + readonly isDefault: boolean +} + +export interface CliProxyGatewayConnectionV1 { + readonly $schema: typeof CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1 + readonly contract: typeof CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1 + readonly schemaVersion: 1 + readonly connectionId: string + readonly revision: number + readonly enabled: boolean + readonly source: Source + readonly wireApi: CliProxyGatewayWireApiV1 + readonly endpointPath: `/${string}` + readonly authorization: 'none' | 'bearer' + readonly adapter: { + readonly adapterId: string + readonly revision: number + readonly required: boolean + } + readonly models: readonly CliProxyGatewayConnectionModelV1[] +} + +export interface CliProxyGatewayExtensionProducerV1 { + readonly pluginId: string + readonly pluginGeneration: string +} + +export type CliProxyGatewayExtensionDiagnosticCodeV1 = + | 'invalid-registration' + | 'capacity-exceeded' + | 'conflicting-producer' + | 'conflicting-source' + | 'conflicting-model' + | 'stale-revision' + | 'stale-generation' + | 'not-found' + +export interface CliProxyGatewayExtensionDiagnosticV1 { + readonly code: CliProxyGatewayExtensionDiagnosticCodeV1 + readonly field?: string + readonly message: string +} + +export type CliProxyGatewayExtensionMutationResultV1 = + | { readonly status: 'registered' | 'updated' | 'revoked'; readonly registryRevision: number } + | { readonly status: 'unchanged'; readonly registryRevision: number } + | { readonly status: 'rejected' | 'stale-generation'; readonly diagnostic: CliProxyGatewayExtensionDiagnosticV1 } + +export interface CliProxyGatewayExtensionRegistrarV1 { + readonly producer: CliProxyGatewayExtensionProducerV1 + registerAdapter(adapter: CliProxyGatewayAdapterV1): CliProxyGatewayExtensionMutationResultV1 + revokeAdapter( + input: { readonly adapterId: string; readonly expectedRevision: number }, + ): CliProxyGatewayExtensionMutationResultV1 + registerConnection(connection: CliProxyGatewayConnectionV1): CliProxyGatewayExtensionMutationResultV1 + revokeConnection(input: { + readonly connectionId: string + readonly expectedRevision: number + }): CliProxyGatewayExtensionMutationResultV1 + dispose(): CliProxyGatewayExtensionMutationResultV1 +} + +export interface CliProxyGatewayExtensionConsumerContextV1 { + readonly cliProxyGatewayExtensions: CliProxyGatewayExtensionRegistrarV1 +} + +interface AdapterEntry { + readonly producer: CliProxyGatewayExtensionProducerV1 + readonly adapter: CliProxyGatewayAdapterV1 +} + +interface ConnectionEntry { + readonly producer: CliProxyGatewayExtensionProducerV1 + readonly sourceKey: string + readonly connection: CliProxyGatewayConnectionV1 +} + +export interface CliProxyGatewayExtensionPlanV1 { + readonly contract: typeof CLI_PROXY_GATEWAY_EXTENSION_PLAN_CONTRACT_V1 + readonly schemaVersion: 1 + readonly pluginGeneration: string + readonly registryRevision: number + readonly planDigest: `sha256:${string}` + readonly configuration: { + readonly adapters: readonly CliProxyGatewayAdapterV1[] + readonly connections: readonly { + readonly connectionId: string + readonly revision: number + readonly enabled: boolean + readonly adapterId: string + readonly adapterRevision: number + readonly adapterRequired: boolean + readonly adapterAvailable: boolean + readonly wireApi: CliProxyGatewayWireApiV1 + readonly endpointPath: `/${string}` + readonly endpoint: null + readonly authorization: 'none' | 'bearer' + readonly credential: null + readonly models: readonly { + readonly sourceModelId: string + readonly gatewayModelId: string + readonly displayName: string + readonly isDefault: boolean + }[] + }[] + } + readonly composition: readonly { + readonly connectionId: string + readonly sourceAlias: string + readonly source: Source + readonly originPointer: `/${string}` + readonly authorizationPointer?: `/${string}` + }[] +} + +interface ValidationFailure { + readonly code: CliProxyGatewayExtensionDiagnosticCodeV1 + readonly field?: string + readonly message: string +} + +function isText(value: string, maximum: number): boolean { + return value.length > 0 && value.length <= maximum && value === value.trim() && !/[\u0000-\u001f\u007f]/.test(value) +} + +function stableValue(value: unknown): unknown { + if (Array.isArray(value)) return value.map(stableValue) + if (value === null || typeof value !== 'object') return value + return Object.fromEntries( + Object.entries(value as Record).sort(([left], [right]) => left.localeCompare(right)).map( + ([key, item]) => [key, stableValue(item)], + ), + ) +} + +function stableStringify(value: unknown): string { + return JSON.stringify(stableValue(value)) +} + +function digest(value: unknown): `sha256:${string}` { + return `sha256:${createHash('sha256').update(stableStringify(value)).digest('hex')}` +} + +function rejection(failure: ValidationFailure): CliProxyGatewayExtensionMutationResultV1 { + return { status: 'rejected', diagnostic: failure } +} + +function validRevision(value: number): boolean { + return Number.isSafeInteger(value) && value > 0 +} + +function validTemplate(value: CliProxyGatewayTemplateV1, depth = 0): boolean { + if (depth > 12) return false + if (value === null || typeof value === 'boolean') return true + if (typeof value === 'number') return Number.isFinite(value) + if (typeof value === 'string') return value.length <= 8_192 && !value.includes('\u0000') + if (Array.isArray(value)) return value.length <= 128 && value.every(item => validTemplate(item, depth + 1)) + const entries = Object.entries(value) + return entries.length <= 128 && entries.every(([key, item]) => isText(key, 256) && validTemplate(item, depth + 1)) +} + +function validateAdapter(adapter: CliProxyGatewayAdapterV1): ValidationFailure | undefined { + if ( + adapter.$schema !== CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1 + || adapter.contract !== CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1 + || adapter.schemaVersion !== 1 + ) return { code: 'invalid-registration', field: 'contract', message: 'Unsupported gateway adapter contract' } + if (!ID_PATTERN.test(adapter.adapterId)) { + return { code: 'invalid-registration', field: 'adapterId', message: 'Invalid adapter id' } + } + if (!validRevision(adapter.revision)) { + return { code: 'invalid-registration', field: 'revision', message: 'Revision must be a positive safe integer' } + } + if (adapter.sessionSources.length > 16) { + return { code: 'invalid-registration', field: 'sessionSources', message: 'Too many session sources' } + } + for (const source of adapter.sessionSources) { + if ( + source.kind === 'header' + ? !HEADER_PATTERN.test(source.name) + : !POINTER_PATTERN.test(source.pointer) || source.pointer.length > 512 + ) return { code: 'invalid-registration', field: 'sessionSources', message: 'Invalid session source' } + } + for (const name of adapter.request.clearHeaders ?? []) { + if (!HEADER_PATTERN.test(name)) { + return { code: 'invalid-registration', field: 'request.clearHeaders', message: 'Invalid header name' } + } + } + for (const [name, value] of Object.entries(adapter.request.setHeaders ?? {})) { + if (!HEADER_PATTERN.test(name) || !validTemplate(value)) { + return { code: 'invalid-registration', field: 'request.setHeaders', message: 'Invalid header transform' } + } + } + for (const [pointer, value] of Object.entries(adapter.request.setBody ?? {})) { + if (!POINTER_PATTERN.test(pointer) || pointer.length > 512 || !validTemplate(value)) { + return { code: 'invalid-registration', field: 'request.setBody', message: 'Invalid body transform' } + } + } + if (stableStringify(adapter).length > 64 * 1_024) { + return { code: 'invalid-registration', message: 'Adapter exceeds the supported size' } + } + return undefined +} + +function validateConnection(connection: CliProxyGatewayConnectionV1): ValidationFailure | undefined { + if ( + connection.$schema !== CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1 + || connection.contract !== CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1 + || connection.schemaVersion !== 1 + ) return { code: 'invalid-registration', field: 'contract', message: 'Unsupported gateway connection contract' } + if (!ID_PATTERN.test(connection.connectionId)) { + return { code: 'invalid-registration', field: 'connectionId', message: 'Invalid connection id' } + } + if (!validRevision(connection.revision)) { + return { code: 'invalid-registration', field: 'revision', message: 'Revision must be a positive safe integer' } + } + if (connection.wireApi !== 'responses' && connection.wireApi !== 'chat-completions') { + return { code: 'invalid-registration', field: 'wireApi', message: 'Unsupported wire API' } + } + if ( + !connection.endpointPath.startsWith('/') || connection.endpointPath.startsWith('//') + || connection.endpointPath.length > 512 || /[?#\u0000-\u001f\u007f]/.test(connection.endpointPath) + ) return { code: 'invalid-registration', field: 'endpointPath', message: 'Invalid endpoint path' } + if (!ID_PATTERN.test(connection.adapter.adapterId) || !validRevision(connection.adapter.revision)) { + return { code: 'invalid-registration', field: 'adapter', message: 'Invalid adapter reference' } + } + if (connection.models.length === 0 || connection.models.length > MAX_MODELS) { + return { code: 'invalid-registration', field: 'models', message: 'Model routes are outside the supported range' } + } + const gatewayModels = new Set() + const sourceModels = new Set() + let defaults = 0 + for (const model of connection.models) { + if (!isText(model.sourceModelId, 256) || !isText(model.modelId, 256) || model.modelId.includes('/')) { + return { code: 'invalid-registration', field: 'models', message: 'Invalid model identity' } + } + if (model.displayName !== undefined && !isText(model.displayName, 200)) { + return { code: 'invalid-registration', field: 'models.displayName', message: 'Invalid model display name' } + } + if (gatewayModels.has(model.modelId) || sourceModels.has(model.sourceModelId)) { + return { code: 'conflicting-model', field: 'models', message: 'Duplicate connection model route' } + } + gatewayModels.add(model.modelId) + sourceModels.add(model.sourceModelId) + if (model.enabled && model.isDefault) defaults += 1 + } + if (!connection.models.some(model => model.enabled)) { + return { code: 'invalid-registration', field: 'models', message: 'At least one enabled model route is required' } + } + if (defaults > 1) { + return { code: 'conflicting-model', field: 'models.isDefault', message: 'Only one default model is allowed' } + } + return undefined +} + +export class CliProxyGatewayExtensionRegistryV1 { + private readonly adapters = new Map() + private readonly connections = new Map>() + private readonly producerGenerations = new Map() + private readonly listeners = new Set<(revision: number) => void>() + private revision = 0 + private disposed = false + + constructor( + readonly pluginGeneration: string, + private readonly sourceKey: (source: Source) => string, + ) { + if (!isText(pluginGeneration, 256)) throw new Error('Invalid plugin generation') + } + + activateProducer(producer: CliProxyGatewayExtensionProducerV1, registryGeneration: string): void { + if (this.disposed || registryGeneration !== this.pluginGeneration) { + throw new Error('Gateway extension registry generation is stale') + } + if (!ID_PATTERN.test(producer.pluginId) || !isText(producer.pluginGeneration, 256)) { + throw new Error('Invalid gateway extension producer authority') + } + if (this.producerGenerations.get(producer.pluginId) === producer.pluginGeneration) return + this.producerGenerations.set(producer.pluginId, producer.pluginGeneration) + if (this.removeProducerEntries(producer.pluginId)) this.changed() + } + + registerAdapter( + adapter: CliProxyGatewayAdapterV1, + registryGeneration: string, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 { + const generationFailure = this.generationFailure(registryGeneration, producer) + if (generationFailure !== undefined) return generationFailure + const validation = validateAdapter(adapter) + if (validation !== undefined) return rejection(validation) + const current = this.adapters.get(adapter.adapterId) + const ownership = this.ownershipFailure(current?.producer, producer) + if (ownership !== undefined) return ownership + if (current !== undefined) { + const revision = this.revisionResult(current.adapter.revision, adapter.revision, current.adapter, adapter) + if (revision !== undefined) return revision + } else if (this.adapters.size >= MAX_ADAPTERS) { + return rejection({ code: 'capacity-exceeded', message: 'Gateway adapter capacity was reached' }) + } + this.adapters.set(adapter.adapterId, { producer, adapter: structuredClone(adapter) }) + this.changed() + return { status: current === undefined ? 'registered' : 'updated', registryRevision: this.revision } + } + + registerConnection( + connection: CliProxyGatewayConnectionV1, + registryGeneration: string, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 { + const generationFailure = this.generationFailure(registryGeneration, producer) + if (generationFailure !== undefined) return generationFailure + const validation = validateConnection(connection) + if (validation !== undefined) return rejection(validation) + const current = this.connections.get(connection.connectionId) + const ownership = this.ownershipFailure(current?.producer, producer) + if (ownership !== undefined) return ownership + const sourceKey = this.sourceKey(connection.source) + if (!isText(sourceKey, 512)) { + return rejection({ + code: 'invalid-registration', + field: 'source', + message: 'Invalid connection source reference', + }) + } + if (current !== undefined) { + const revision = this.revisionResult( + current.connection.revision, + connection.revision, + { ...current.connection, source: current.sourceKey }, + { ...connection, source: sourceKey }, + ) + if (revision !== undefined) return revision + } else if (this.connections.size >= MAX_CONNECTIONS) { + return rejection({ code: 'capacity-exceeded', message: 'Gateway connection capacity was reached' }) + } + for (const [connectionId, entry] of this.connections) { + if (connectionId === connection.connectionId) continue + if (entry.sourceKey === sourceKey) { + return rejection({ + code: 'conflicting-source', + field: 'source', + message: 'Connection source is already registered', + }) + } + const existingModels = new Set( + entry.connection.models.filter(model => entry.connection.enabled && model.enabled).map(model => + `${entry.connection.connectionId}/${model.modelId}` + ), + ) + if ( + connection.models.some(model => + connection.enabled && model.enabled && existingModels.has( + `${connection.connectionId}/${model.modelId}`, + ) + ) + ) return rejection({ code: 'conflicting-model', field: 'models', message: 'Gateway model is already registered' }) + } + this.connections.set(connection.connectionId, { + producer, + sourceKey, + connection: structuredClone(connection), + }) + this.changed() + return { status: current === undefined ? 'registered' : 'updated', registryRevision: this.revision } + } + + revokeAdapter( + input: { readonly adapterId: string; readonly expectedRevision: number }, + registryGeneration: string, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 { + const generationFailure = this.generationFailure(registryGeneration, producer) + if (generationFailure !== undefined) return generationFailure + const current = this.adapters.get(input.adapterId) + if (current === undefined) { + return rejection({ code: 'not-found', field: 'adapterId', message: 'Adapter is not registered' }) + } + const ownership = this.ownershipFailure(current.producer, producer) + if (ownership !== undefined) return ownership + if (current.adapter.revision !== input.expectedRevision) { + return rejection({ code: 'stale-revision', field: 'expectedRevision', message: 'Adapter revision is stale' }) + } + this.adapters.delete(input.adapterId) + this.changed() + return { status: 'revoked', registryRevision: this.revision } + } + + revokeConnection( + input: { readonly connectionId: string; readonly expectedRevision: number }, + registryGeneration: string, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 { + const generationFailure = this.generationFailure(registryGeneration, producer) + if (generationFailure !== undefined) return generationFailure + const current = this.connections.get(input.connectionId) + if (current === undefined) { + return rejection({ code: 'not-found', field: 'connectionId', message: 'Connection is not registered' }) + } + const ownership = this.ownershipFailure(current.producer, producer) + if (ownership !== undefined) return ownership + if (current.connection.revision !== input.expectedRevision) { + return rejection({ code: 'stale-revision', field: 'expectedRevision', message: 'Connection revision is stale' }) + } + this.connections.delete(input.connectionId) + this.changed() + return { status: 'revoked', registryRevision: this.revision } + } + + disposeProducer( + producer: CliProxyGatewayExtensionProducerV1, + registryGeneration: string, + ): CliProxyGatewayExtensionMutationResultV1 { + const generationFailure = this.generationFailure(registryGeneration, producer) + if (generationFailure !== undefined) return generationFailure + this.producerGenerations.delete(producer.pluginId) + if (!this.removeProducerEntries(producer.pluginId, producer.pluginGeneration)) { + return { status: 'unchanged', registryRevision: this.revision } + } + this.changed() + return { status: 'revoked', registryRevision: this.revision } + } + + subscribe(listener: (revision: number) => void): () => void { + if (this.disposed) throw new Error('Gateway extension registry generation is stale') + this.listeners.add(listener) + return () => this.listeners.delete(listener) + } + + plan(pluginGeneration: string): CliProxyGatewayExtensionPlanV1 | CliProxyGatewayExtensionDiagnosticV1 { + if (this.disposed || pluginGeneration !== this.pluginGeneration) { + return { code: 'stale-generation', message: 'Gateway extension plan generation is stale' } + } + const adapters = [...this.adapters.values()].map(entry => structuredClone(entry.adapter)).sort((left, right) => + left.adapterId.localeCompare(right.adapterId) + ) + const connections: CliProxyGatewayExtensionPlanV1['configuration']['connections'][number][] = [] + const composition: CliProxyGatewayExtensionPlanV1['composition'][number][] = [] + for ( + const entry of [...this.connections.values()].sort((left, right) => + left.connection.connectionId.localeCompare(right.connection.connectionId) + ) + ) { + const connection = entry.connection + const adapter = this.adapters.get(connection.adapter.adapterId)?.adapter + const adapterAvailable = adapter?.enabled === true && adapter.revision === connection.adapter.revision + const enabled = connection.enabled && (!connection.adapter.required || adapterAvailable) + const index = connections.length + connections.push({ + connectionId: connection.connectionId, + revision: connection.revision, + enabled, + adapterId: connection.adapter.adapterId, + adapterRevision: connection.adapter.revision, + adapterRequired: connection.adapter.required, + adapterAvailable, + wireApi: connection.wireApi, + endpointPath: connection.endpointPath, + endpoint: null, + authorization: connection.authorization, + credential: null, + models: connection.models.filter(model => model.enabled).sort((left, right) => + left.modelId.localeCompare(right.modelId) + ).map(model => ({ + sourceModelId: model.sourceModelId, + gatewayModelId: `${connection.connectionId}/${model.modelId}`, + displayName: model.displayName ?? model.modelId, + isDefault: model.isDefault, + })), + }) + if (!enabled) continue + const sourceAlias = `extension-${index}` + composition.push({ + connectionId: connection.connectionId, + sourceAlias, + source: connection.source, + originPointer: `/connections/${index}/endpoint`, + ...(connection.authorization === 'bearer' + ? { authorizationPointer: `/connections/${index}/credential` as const } + : {}), + }) + } + const base = { + contract: CLI_PROXY_GATEWAY_EXTENSION_PLAN_CONTRACT_V1, + schemaVersion: 1 as const, + pluginGeneration: this.pluginGeneration, + registryRevision: this.revision, + configuration: { adapters, connections }, + composition, + } satisfies Omit, 'planDigest'> + return { ...base, planDigest: digest({ ...base, composition: composition.map(({ source: _, ...item }) => item) }) } + } + + dispose(): void { + if (this.disposed) return + this.disposed = true + this.adapters.clear() + this.connections.clear() + this.producerGenerations.clear() + this.listeners.clear() + this.revision += 1 + } + + private revisionResult( + currentRevision: number, + nextRevision: number, + current: unknown, + next: unknown, + ): CliProxyGatewayExtensionMutationResultV1 | undefined { + if (nextRevision < currentRevision) { + return rejection({ code: 'stale-revision', field: 'revision', message: 'Extension revision is stale' }) + } + if (nextRevision !== currentRevision) return undefined + return stableStringify(current) === stableStringify(next) + ? { status: 'unchanged', registryRevision: this.revision } + : rejection({ code: 'stale-revision', field: 'revision', message: 'Revision was reused with different content' }) + } + + private ownershipFailure( + current: CliProxyGatewayExtensionProducerV1 | undefined, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 | undefined { + if (current === undefined) return undefined + if (current.pluginId !== producer.pluginId) { + return rejection({ code: 'conflicting-producer', message: 'Extension is owned by another producer' }) + } + if (current.pluginGeneration !== producer.pluginGeneration) { + return { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Extension producer generation is stale' }, + } + } + return undefined + } + + private generationFailure( + registryGeneration: string, + producer: CliProxyGatewayExtensionProducerV1, + ): CliProxyGatewayExtensionMutationResultV1 | undefined { + if ( + this.disposed || registryGeneration !== this.pluginGeneration + || this.producerGenerations.get(producer.pluginId) !== producer.pluginGeneration + ) { + return { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Gateway extension registry generation is stale' }, + } + } + return undefined + } + + private removeProducerEntries(pluginId: string, pluginGeneration?: string): boolean { + let removed = false + for (const [adapterId, entry] of this.adapters) { + if ( + entry.producer.pluginId !== pluginId + || (pluginGeneration !== undefined && entry.producer.pluginGeneration !== pluginGeneration) + ) continue + this.adapters.delete(adapterId) + removed = true + } + for (const [connectionId, entry] of this.connections) { + if ( + entry.producer.pluginId !== pluginId + || (pluginGeneration !== undefined && entry.producer.pluginGeneration !== pluginGeneration) + ) continue + this.connections.delete(connectionId) + removed = true + } + return removed + } + + private changed(): void { + this.revision += 1 + for (const listener of this.listeners) listener(this.revision) + } +} + +export function createCliProxyGatewayExtensionRegistrarV1( + registry: CliProxyGatewayExtensionRegistryV1, + producer: CliProxyGatewayExtensionProducerV1, +): CliProxyGatewayExtensionRegistrarV1 { + registry.activateProducer(producer, registry.pluginGeneration) + let disposed = false + const currentRevision = (): number => { + const plan = registry.plan(registry.pluginGeneration) + return 'code' in plan ? 0 : plan.registryRevision + } + return Object.freeze({ + producer, + registerAdapter: (adapter: CliProxyGatewayAdapterV1) => + disposed + ? { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Extension registrar is disposed' }, + } + : registry.registerAdapter(adapter, registry.pluginGeneration, producer), + revokeAdapter: (input: { readonly adapterId: string; readonly expectedRevision: number }) => + disposed + ? { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Extension registrar is disposed' }, + } + : registry.revokeAdapter(input, registry.pluginGeneration, producer), + registerConnection: (connection: CliProxyGatewayConnectionV1) => + disposed + ? { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Extension registrar is disposed' }, + } + : registry.registerConnection(connection, registry.pluginGeneration, producer), + revokeConnection: (input: { readonly connectionId: string; readonly expectedRevision: number }) => + disposed + ? { + status: 'stale-generation', + diagnostic: { code: 'stale-generation', message: 'Extension registrar is disposed' }, + } + : registry.revokeConnection(input, registry.pluginGeneration, producer), + dispose: () => { + if (disposed) return { status: 'unchanged', registryRevision: currentRevision() } + disposed = true + return registry.disposeProducer(producer, registry.pluginGeneration) + }, + }) as CliProxyGatewayExtensionRegistrarV1 +} diff --git a/src/service/gateway-definition.ts b/src/service/gateway-definition.ts index a329501..c6ec8c0 100644 --- a/src/service/gateway-definition.ts +++ b/src/service/gateway-definition.ts @@ -7,6 +7,8 @@ export const CLI_PROXY_GATEWAY_CODEX_UPSTREAMS_SCHEMA_V1 = 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-codex-upstream.v1.schema.json' as const export const CLI_PROXY_GATEWAY_OPENAI_UPSTREAMS_SCHEMA_V1 = 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-openai-upstream.v1.schema.json' as const +export const CLI_PROXY_GATEWAY_EXTENSION_PLAN_SCHEMA_V1 = + 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-gateway-extension-plan.v1.schema.json' as const const CLI_PROXY_MANAGEMENT_AUTH_FILES_SCHEMA_V1 = 'https://raw.githubusercontent.com/cordisx/plugin-cli-proxy-api/main/schemas/cli-proxy-management-auth-files.v1.schema.json' as const const CLI_PROXY_MANAGEMENT_OPERATION_SCHEMA_V1 = @@ -69,6 +71,13 @@ export const cliProxyGatewayDefinition = { pointer: '/openai-compatibility', valueSchema: CLI_PROXY_GATEWAY_OPENAI_UPSTREAMS_SCHEMA_V1, }, + { + slot: 'gateway-extensions', + source: 'composition', + target: 'configuration', + pointer: '/plugins/configs/cordisx-gateway-bridge', + valueSchema: CLI_PROXY_GATEWAY_EXTENSION_PLAN_SCHEMA_V1, + }, ], authentication: { mode: 'none' }, httpAuthentication: { mode: 'authorization-header', scheme: 'Bearer', slot: 'gateway-key' }, diff --git a/src/service/gateway.ts b/src/service/gateway.ts index 96c5d77..e23ecc2 100644 --- a/src/service/gateway.ts +++ b/src/service/gateway.ts @@ -14,6 +14,12 @@ import type { } from '@cordisx/protocol/managed-service-runtime/v1' import { cliProxyGatewayDefinition } from './gateway-definition.js' import { bindGatewayUIState, managedServiceUI } from './gateway-ui.js' +import { createCliProxyGatewayExtensionContextProviderV1 } from './extension-context.js' +import { + CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1, + type CliProxyGatewayExtensionPlanV1, + CliProxyGatewayExtensionRegistryV1, +} from './extension-registry.js' import { createCliProxyUpstreamContextProviderV1 } from './upstream-context.js' import { CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1, @@ -24,16 +30,19 @@ import { } from './upstream-registry.js' export * from './gateway-definition.js' +export * from './extension-registry.js' export * from './upstream-registry.js' export const contextServices: readonly ManagedServiceContextProviderV1[] = Object.freeze([ createCliProxyUpstreamContextProviderV1(), + createCliProxyGatewayExtensionContextProviderV1(), ]) export { managedServiceUI } export const CLI_PROXY_NATIVE_PROVIDER_ID = 'cli-proxy-api' export type CliProxyGatewayContextV1 = ManagedServiceServiceContextV1 & Pick & { readonly cliProxyUpstreams: CliProxyUpstreamRegistryV1 + readonly cliProxyGatewayExtensions: CliProxyGatewayExtensionRegistryV1 } function activationError( @@ -49,6 +58,20 @@ function assertPlan( if ('code' in value) throw new Error(`CLIProxyAPI gateway plan failed: ${value.code}`) } +function assertExtensionPlan( + value: CliProxyGatewayExtensionPlanV1 | { readonly code: string }, +): asserts value is CliProxyGatewayExtensionPlanV1 { + if ('code' in value) throw new Error(`CLIProxyAPI gateway extension plan failed: ${value.code}`) +} + +function combinedRevision(upstreamDigest: string, extensionDigest: string): `sha256:${string}` { + return `sha256:${createHash('sha256').update(JSON.stringify([upstreamDigest, extensionDigest])).digest('hex')}` +} + +function safeValue(value: unknown): ManagedServiceSafeValueV1 { + return JSON.parse(JSON.stringify(value)) as ManagedServiceSafeValueV1 +} + function modelIds(value: ManagedServiceSafeValueV1): ReadonlySet { if (value === null || typeof value !== 'object' || Array.isArray(value)) return new Set() const data = (value as { readonly [key: string]: ManagedServiceSafeValueV1 }).data @@ -99,9 +122,11 @@ export async function activateCliProxyGateway( input: ManagedServiceApplyInputV1, ): Promise { const registry = context.cliProxyUpstreams - let active: { readonly revision: number; readonly dispose: () => Promise } | undefined + const extensionRegistry = context.cliProxyGatewayExtensions + let active: { readonly revision: string; readonly dispose: () => Promise } | undefined let transition = Promise.resolve() let unsubscribe = (): void => undefined + let unsubscribeExtensions = (): void => undefined let disposePromise: Promise | undefined let gatewayHandle: ManagedServiceRegistrationHandleV1 | undefined let publication: ManagedNativeProviderPublicationHandleV1 | undefined @@ -126,9 +151,11 @@ export async function activateCliProxyGateway( const activate = async ( catalog: CliProxyModelCatalogV1, plan: CliProxyGatewayPlanV1, + extensionPlan: CliProxyGatewayExtensionPlanV1, + revision: `sha256:${string}`, ): Promise<() => Promise> => { gatewayHandle ??= await context.managedServices.register(cliProxyGatewayDefinition, { - revision: plan.catalogDigest, + revision, }) const handle = gatewayHandle const leases: ManagedServiceLeaseV1[] = [] @@ -176,6 +203,14 @@ export async function activateCliProxyGateway( ) } } + for (const composition of extensionPlan.composition) { + const acquired = await input.client.acquire(composition.source, { signal: input.signal }) + if (acquired.status !== 'ready') { + throw activationError(`${composition.connectionId} extension acquisition`, acquired) + } + leases.push(acquired.lease) + sources.push({ source: composition.sourceAlias, lease: acquired.lease }) + } const bindings: Parameters[0]['bindings'][number][] = [ { @@ -186,6 +221,19 @@ export async function activateCliProxyGateway( targetSlot: 'openai-upstreams', source: { kind: 'safe-literal', value: plan.configuration['openai-compatibility'] }, }, + { + targetSlot: 'gateway-extensions', + source: { + kind: 'safe-literal', + value: safeValue({ + enabled: true, + active: extensionPlan.composition.length > 0, + revision: extensionPlan.planDigest, + adapters: extensionPlan.configuration.adapters, + connections: extensionPlan.configuration.connections, + }), + }, + }, ] for (const composition of plan.composition) { const collection = composition.origin.targetPointer.startsWith('/codex-api-key/') @@ -204,10 +252,24 @@ export async function activateCliProxyGateway( }) } } + for (const composition of extensionPlan.composition) { + bindings.push({ + targetSlot: 'gateway-extensions', + targetPointer: composition.originPointer, + source: { kind: 'source-origin', source: composition.sourceAlias, origin: 'api' }, + }) + if (composition.authorizationPointer !== undefined) { + bindings.push({ + targetSlot: 'gateway-extensions', + targetPointer: composition.authorizationPointer, + source: { kind: 'source-authorization', source: composition.sourceAlias }, + }) + } + } const materialized = await input.client.materialize({ target: handle.binding, - revision: plan.catalogDigest, + revision, sources, bindings, }, { signal: input.signal }) @@ -225,6 +287,10 @@ export async function activateCliProxyGateway( const actualGatewayModels = await readGatewayModels() const expectedGatewayModels = catalog.groups.flatMap(group => group.providers.flatMap(provider => provider.models.map(model => model.gatewayModelId)) + ).concat( + extensionPlan.configuration.connections.flatMap(connection => + connection.enabled ? connection.models.map(model => model.gatewayModelId) : [] + ), ) const missingGatewayModels = expectedGatewayModels.filter(model => !actualGatewayModels.has(model)) if (missingGatewayModels.length > 0) { @@ -296,7 +362,10 @@ export async function activateCliProxyGateway( if (input.signal.aborted) return const catalog = registry.catalog(input.owner.pluginGeneration) if ('code' in catalog) throw new Error(`CLIProxyAPI catalog failed: ${catalog.code}`) - if (active?.revision === catalog.registryRevision) return + const extensionPlan = extensionRegistry.plan(input.owner.pluginGeneration) + assertExtensionPlan(extensionPlan) + const revision = combinedRevision(catalog.catalogDigest, extensionPlan.planDigest) + if (active?.revision === revision) return const previous = active active = undefined activeRevision = undefined @@ -306,18 +375,18 @@ export async function activateCliProxyGateway( if (gatewayHandle === undefined) { gatewayHandle = await context.managedServices.register(cliProxyGatewayDefinition, { - revision: plan.catalogDigest, + revision, }) } if (input.signal.aborted) return - const cleanup = await activate(catalog, plan) + const cleanup = await activate(catalog, plan, extensionPlan, revision) if (input.signal.aborted) { await cleanup() return } - active = { revision: catalog.registryRevision, dispose: cleanup } + active = { revision, dispose: cleanup } } const schedule = (): Promise => { @@ -339,6 +408,7 @@ export async function activateCliProxyGateway( if (disposePromise !== undefined) return disposePromise disposePromise = (async () => { unsubscribe() + unsubscribeExtensions() unbindGatewayUIState() await transition const failures: unknown[] = [] @@ -370,6 +440,9 @@ export async function activateCliProxyGateway( unsubscribe = registry.subscribe(() => { void schedule().catch(() => undefined) }) + unsubscribeExtensions = extensionRegistry.subscribe(() => { + void schedule().catch(() => undefined) + }) const onAbort = (): void => { void dispose().catch(() => undefined) } @@ -402,5 +475,9 @@ const gatewayApply: ManagedServiceApplyV1 = async (context, input) => { } export const apply = Object.assign(gatewayApply, { - inject: ['managedServices', CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1] as const, + inject: [ + 'managedServices', + CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1, + CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1, + ] as const, }) diff --git a/src/service/index.ts b/src/service/index.ts index 3629131..3de4cd0 100644 --- a/src/service/index.ts +++ b/src/service/index.ts @@ -6,6 +6,7 @@ import { CliProxyPlatformProviderAdapter } from './adapter.js' import { BROKER_BINDINGS } from './bindings.js' export * from './upstream-registry.js' +export * from './extension-registry.js' const OPERATIONS = [ 'models.list', diff --git a/test/extension-registry.mjs b/test/extension-registry.mjs new file mode 100644 index 0000000..35087cd --- /dev/null +++ b/test/extension-registry.mjs @@ -0,0 +1,154 @@ +import assert from 'node:assert/strict' +import { readFile } from 'node:fs/promises' +import test from 'node:test' +import Ajv2020 from 'ajv/dist/2020.js' + +const extension = await import('../dist/extensions.mjs') +const { + CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1, + CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1, + CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1, + CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1, + CliProxyGatewayExtensionRegistryV1, + createCliProxyGatewayExtensionRegistrarV1, +} = extension + +const generation = 'gateway-generation' +const producer = (pluginId = 'synthetic-extension', pluginGeneration = 'producer-generation-1') => ({ + pluginId, + pluginGeneration, +}) +const source = id => ({ identityHandle: `source_${id}` }) +const adapter = (overrides = {}) => ({ + $schema: CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1, + contract: CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1, + schemaVersion: 1, + adapterId: 'session-header', + revision: 1, + enabled: true, + requiredSession: true, + sessionSources: [{ kind: 'header', name: 'Session-Id' }], + request: { + clearHeaders: ['X-Untrusted'], + setHeaders: { 'X-Session': '{{session}}', 'X-Connection': '{{connectionId}}' }, + setBody: { '/metadata/session': '{{session}}' }, + }, + ...overrides, +}) +const connection = (id = 'connection-a', overrides = {}) => ({ + $schema: CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1, + contract: CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1, + schemaVersion: 1, + connectionId: id, + revision: 1, + enabled: true, + source: source(id), + wireApi: 'responses', + endpointPath: '/responses', + authorization: 'bearer', + adapter: { adapterId: 'session-header', revision: 1, required: true }, + models: [{ sourceModelId: 'shared-model', modelId: 'shared', enabled: true, isDefault: true }], + ...overrides, +}) + +function registry() { + return new CliProxyGatewayExtensionRegistryV1(generation, value => value.identityHandle) +} + +function registrar(extensions, authority = producer()) { + return createCliProxyGatewayExtensionRegistrarV1(extensions, authority) +} + +test('adapter and connection schemas reject private transport fields', async () => { + const ajv = new Ajv2020({ allErrors: true, strict: true }) + const adapterSchema = JSON.parse( + await readFile(new URL('../schemas/cli-proxy-gateway-adapter.v1.schema.json', import.meta.url), 'utf8'), + ) + const connectionSchema = JSON.parse( + await readFile(new URL('../schemas/cli-proxy-gateway-connection.v1.schema.json', import.meta.url), 'utf8'), + ) + const validateAdapter = ajv.compile(adapterSchema) + const validateConnection = ajv.compile(connectionSchema) + assert.equal(validateAdapter(adapter()), true, JSON.stringify(validateAdapter.errors)) + assert.equal(validateAdapter({ ...adapter(), tenantId: 'private' }), false) + assert.equal(validateConnection(connection()), true, JSON.stringify(validateConnection.errors)) + assert.equal(validateConnection({ ...connection(), endpoint: 'https://private.invalid' }), false) + assert.equal(validateConnection({ ...connection(), token: 'secret' }), false) +}) + +test('registers, updates, revokes, disposes, and fences producer generations', () => { + const extensions = registry() + const first = registrar(extensions) + assert.equal(first.registerAdapter(adapter()).status, 'registered') + assert.equal(first.registerConnection(connection()).status, 'registered') + assert.equal(first.registerAdapter(adapter()).status, 'unchanged') + assert.equal(first.registerConnection(connection()).status, 'unchanged') + assert.equal(first.registerAdapter(adapter({ revision: 2, requiredSession: false })).status, 'updated') + assert.equal( + first.registerConnection(connection('connection-a', { + revision: 2, + adapter: { adapterId: 'session-header', revision: 2, required: true }, + })).status, + 'updated', + ) + + const replacement = registrar(extensions, producer('synthetic-extension', 'producer-generation-2')) + assert.equal(first.revokeConnection({ connectionId: 'connection-a', expectedRevision: 2 }).status, 'stale-generation') + assert.equal(first.dispose().status, 'stale-generation') + assert.equal(replacement.registerAdapter(adapter({ revision: 1 })).status, 'registered') + assert.equal(replacement.registerConnection(connection()).status, 'registered') + assert.equal(replacement.revokeConnection({ connectionId: 'connection-a', expectedRevision: 1 }).status, 'revoked') + assert.equal(replacement.revokeAdapter({ adapterId: 'session-header', expectedRevision: 1 }).status, 'revoked') + assert.equal(replacement.dispose().status, 'unchanged') +}) + +test('plans two exact connections sharing one source model without credential or endpoint values', () => { + const extensions = registry() + const registration = registrar(extensions) + registration.registerAdapter(adapter()) + registration.registerAdapter(adapter({ + adapterId: 'session-body', + sessionSources: [{ kind: 'body-json-pointer', pointer: '/metadata/thread' }], + request: { setHeaders: { 'X-Session': '{{session}}' } }, + })) + registration.registerConnection(connection('connection-b', { + source: source('b'), + adapter: { adapterId: 'session-body', revision: 1, required: true }, + })) + registration.registerConnection(connection('connection-a', { source: source('a') })) + + const plan = extensions.plan(generation) + assert.equal(plan.contract, 'cordisx.cli-proxy-gateway-extension-plan/v1') + assert.deepEqual(plan.configuration.connections.map(item => item.connectionId), ['connection-a', 'connection-b']) + assert.deepEqual( + plan.configuration.connections.map(item => item.models[0].gatewayModelId), + ['connection-a/shared', 'connection-b/shared'], + ) + assert.equal(plan.configuration.connections.every(item => item.endpoint === null), true) + assert.equal(plan.configuration.connections.every(item => item.credential === null), true) + assert.deepEqual(plan.composition.map(item => item.sourceAlias), ['extension-0', 'extension-1']) +}) + +test('required adapter mismatch disables the protected route before materialization', () => { + const extensions = registry() + const registration = registrar(extensions) + registration.registerAdapter(adapter({ enabled: false })) + registration.registerConnection(connection()) + const plan = extensions.plan(generation) + assert.equal(plan.configuration.connections[0].adapterAvailable, false) + assert.equal(plan.configuration.connections[0].enabled, false) + assert.deepEqual(plan.composition, []) +}) + +test('rejects cross-producer ownership and duplicate source identities', () => { + const extensions = registry() + const owner = registrar(extensions, producer('owner')) + const intruder = registrar(extensions, producer('intruder')) + owner.registerAdapter(adapter()) + owner.registerConnection(connection()) + assert.equal(intruder.registerAdapter(adapter({ revision: 2 })).diagnostic.code, 'conflicting-producer') + assert.equal( + intruder.registerConnection(connection('connection-b', { source: source('connection-a') })).diagnostic.code, + 'conflicting-source', + ) +}) diff --git a/test/gateway.mjs b/test/gateway.mjs index 6961467..3e1d305 100644 --- a/test/gateway.mjs +++ b/test/gateway.mjs @@ -1,11 +1,26 @@ import assert from 'node:assert/strict' -import { readFile } from 'node:fs/promises' import test from 'node:test' -import Ajv2020 from 'ajv/dist/2020.js' import { binding, expectedNativeCatalog, identity, lease, projection, uiBinding } from './support/gateway-fixture.mjs' +import { + activationInput, + compositionLeaves, + descriptor, + expectedRootCollections, + extensionAdapter, + extensionConnection, + providerRoot, + rootCollections, + syntheticExtensions, + syntheticSources, + validateMaterializationRequest, +} from './support/gateway-extension-fixture.mjs' const gatewayModule = await import('../dist/gateway.mjs') const { + CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1, + CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1, + CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1, + CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1, CLI_PROXY_UPSTREAM_DESCRIPTOR_CONTRACT_V1, CLI_PROXY_UPSTREAM_DESCRIPTOR_SCHEMA_V1, CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1, @@ -39,117 +54,6 @@ test('declares a separate Host-private CLIProxyAPI management credential', () => ) }) -const compositionSchemaFiles = new Map([ - ['codex-upstreams', '../schemas/cli-proxy-gateway-codex-upstream.v1.schema.json'], - ['openai-upstreams', '../schemas/cli-proxy-gateway-openai-upstream.v1.schema.json'], -]) -const compositionValidators = new Map( - await Promise.all([...compositionSchemaFiles].map(async ([slot, file]) => { - const schema = JSON.parse(await readFile(new URL(file, import.meta.url), 'utf8')) - return [slot, new Ajv2020({ allErrors: true, strict: true }).compile(schema)] - })), -) - -const syntheticSources = Object.freeze({ - responses: Object.freeze({ - pluginId: 'synthetic-responses', - pluginGeneration: 'synthetic-responses-one', - serviceId: 'synthetic-responses-service', - serviceGeneration: 'synthetic-responses-service-one', - upstreamId: 'responses-source', - displayName: 'Synthetic Responses', - prefix: 'responses', - wireApi: 'responses', - sourceModelId: 'responses-model', - modelId: 'family/fast', - authenticated: false, - }), - chat: Object.freeze({ - pluginId: 'synthetic-chat', - pluginGeneration: 'synthetic-chat-one', - serviceId: 'synthetic-chat-service', - serviceGeneration: 'synthetic-chat-service-one', - upstreamId: 'chat-source', - displayName: 'Synthetic Chat', - prefix: 'chat', - wireApi: 'chat-completions', - sourceModelId: 'chat-model', - modelId: 'family/balanced', - authenticated: true, - }), -}) - -const descriptor = (source, overrides = {}) => ({ - $schema: CLI_PROXY_UPSTREAM_DESCRIPTOR_SCHEMA_V1, - contract: CLI_PROXY_UPSTREAM_DESCRIPTOR_CONTRACT_V1, - schemaVersion: 1, - revision: 1, - upstreamId: source.upstreamId, - displayName: source.displayName, - enabled: true, - wireApi: source.wireApi, - prefix: source.prefix, - order: source.wireApi === 'responses' ? 10 : 20, - requestTimeoutMs: 30_000, - group: { groupId: 'synthetic', displayName: 'Synthetic', order: 10 }, - models: [{ - sourceModelId: source.sourceModelId, - modelId: source.modelId, - enabled: true, - isDefault: source.wireApi === 'responses', - ...(source.wireApi === 'chat-completions' ? { inputModalities: ['text'] } : {}), - }], - ...overrides, -}) - -function providerRoot() { - const controller = new AbortController() - const target = { - pluginId: 'cli-proxy-api', - pluginGeneration: 'gateway-generation', - serviceId: 'gateway-runtime', - serviceGeneration: 'gateway-service-generation', - } - return { - controller, - target, - root: contextServices[0].create({ target, signal: controller.signal }), - } -} - -function validateMaterializationRequest(request) { - const compositionSlots = cliProxyGatewayDefinition.protectedBindings.filter(item => item.source === 'composition') - const declaredSources = new Set(request.sources.map(source => source.source)) - for (const slot of compositionSlots) { - const slotBindings = request.bindings.filter(binding => binding.targetSlot === slot.slot) - assert.ok(slotBindings.length > 0, `missing composition slot ${slot.slot}`) - const roots = slotBindings.filter(binding => binding.targetPointer === undefined) - assert.equal(roots.length, 1, `composition slot ${slot.slot} must have exactly one root binding`) - assert.equal(roots[0].source.kind, 'safe-literal') - const validate = compositionValidators.get(slot.slot) - assert.equal(validate(roots[0].source.value), true, JSON.stringify(validate.errors)) - } - for (const item of request.bindings) { - if (item.source.kind === 'source-origin' || item.source.kind === 'source-authorization') { - assert.equal(declaredSources.has(item.source.source), true, `undeclared composition source ${item.source.source}`) - } - } -} - -function activationInput(client, signal) { - return { - owner: { - ownerHandle: 'mso_gateway', - pluginId: 'cli-proxy-api', - sourceDigest: `sha256:${'a'.repeat(64)}`, - hostGeneration: 'host-one', - pluginGeneration: 'gateway-generation', - }, - client, - signal, - } -} - function gatewayFixture({ kinds = ['responses', 'chat'], gatewayModels = [], @@ -164,7 +68,7 @@ function gatewayFixture({ publicationDisposeError, publicationDisposeDiagnostic, } = {}) { - const { root, target } = providerRoot() + const { root, extensionRoot, target } = providerRoot() const sources = [] let ownedGatewayModels = [...gatewayModels] let oauthPollIndex = 0 @@ -185,6 +89,35 @@ function gatewayFixture({ }) assert.equal(result.status, 'registered') } + const registerExtension = ({ adapter, connection, source: extensionSource }) => { + const sourceLease = lease( + extensionSource.pluginId, + extensionSource.serviceId, + extensionSource.serviceGeneration, + ) + sourceLeases.set(extensionSource.serviceId, sourceLease) + sourceByService.set(extensionSource.serviceId, extensionSource) + const consumer = extensionRoot.bind({ + producer: { pluginId: extensionSource.pluginId, pluginGeneration: extensionSource.pluginGeneration }, + target, + signal: new AbortController().signal, + }) + assert.match(consumer.value.registerAdapter(adapter).status, /^(?:registered|updated|unchanged)$/) + assert.match( + consumer.value.registerConnection({ + ...connection, + source: identity(extensionSource.pluginId, extensionSource.serviceId), + }).status, + /^(?:registered|updated|unchanged)$/, + ) + if (connection.enabled) { + ownedGatewayModels = [ + ...ownedGatewayModels, + ...connection.models.filter(model => model.enabled).map(model => `${connection.connectionId}/${model.modelId}`), + ] + } + return consumer + } for (const kind of kinds) registerSource(syntheticSources[kind]) const events = [] @@ -401,6 +334,7 @@ function gatewayFixture({ }, }, cliProxyUpstreams: root.providerValue, + cliProxyGatewayExtensions: extensionRoot.providerValue, } const hostRevokePublication = providerId => { activePublications.delete(providerId) @@ -423,6 +357,7 @@ function gatewayFixture({ publicationDisposeResults, publicationHandles, registerSource, + registerExtension, registration, released, sources, @@ -431,56 +366,14 @@ function gatewayFixture({ } } -function rootCollections(request) { - return Object.fromEntries( - request.bindings.filter(item => item.targetPointer === undefined).map(item => [ - item.targetSlot, - item.source.value, - ]), - ) -} - -function compositionLeaves(request) { - return request.bindings.filter(item => item.targetPointer !== undefined) -} - -function expectedRootCollections(sources) { - return { - 'codex-upstreams': sources.filter(source => source.wireApi === 'responses').map(source => ({ - 'api-key': null, - prefix: source.prefix, - 'base-url': null, - 'request-retry': 0, - models: [{ - name: source.sourceModelId, - alias: source.modelId, - 'display-name': source.modelId, - 'force-mapping': true, - }], - })), - 'openai-upstreams': sources.filter(source => source.wireApi === 'chat-completions').map(source => ({ - name: source.upstreamId, - disabled: false, - prefix: source.prefix, - 'base-url': null, - 'api-key-entries': [{ 'api-key': null }], - 'request-retry': 0, - models: [{ - name: source.sourceModelId, - alias: source.modelId, - 'display-name': source.modelId, - 'force-mapping': true, - 'input-modalities': ['text'], - }], - })), - } -} - -test('exports one versioned context service with owner-local registry access', () => { - assert.equal(contextServices.length, 1) - assert.equal(contextServices[0].service, CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1) - const { root } = providerRoot() +test('exports two versioned context services with owner-local registry access', () => { + assert.deepEqual(contextServices.map(service => service.service), [ + CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1, + 'cliProxyGatewayExtensions', + ]) + const { root, extensionRoot } = providerRoot() assert.equal(root.providerValue.pluginGeneration, 'gateway-generation') + assert.equal(extensionRoot.providerValue.pluginGeneration, 'gateway-generation') }) test('projects the active upstream registry through the Node UI source without exposing secrets', async () => { @@ -727,7 +620,8 @@ for ( assert.equal(request.revision, fixture.events.find(event => event[0] === 'register')[2].revision) assert.deepEqual(request.sources.map(source => source.source), fixture.sources.map(source => source.upstreamId)) const roots = rootCollections(request) - assert.deepEqual(roots, expectedRootCollections(fixture.sources)) + assert.match(roots['gateway-extensions'].revision, /^sha256:[0-9a-f]{64}$/) + assert.deepEqual(roots, expectedRootCollections(fixture.sources, roots['gateway-extensions'].revision)) assert.deepEqual( compositionLeaves(request), fixture.sources.flatMap(source => { @@ -777,9 +671,18 @@ test('starts without upstreams and publishes account-backed gateway models immed assert.equal(fixture.events.filter(event => event[0] === 'register').length, 1) assert.equal(fixture.events.filter(event => event[0] === 'materialize').length, 1) - assert.deepEqual(rootCollections(fixture.events.find(event => event[0] === 'materialize')[1]), { + const roots = rootCollections(fixture.events.find(event => event[0] === 'materialize')[1]) + assert.match(roots['gateway-extensions'].revision, /^sha256:[0-9a-f]{64}$/) + assert.deepEqual(roots, { 'codex-upstreams': [], 'openai-upstreams': [], + 'gateway-extensions': { + enabled: true, + active: false, + revision: roots['gateway-extensions'].revision, + adapters: [], + connections: [], + }, }) assert.deepEqual(fixture.events.filter(event => event[0] === 'publish-native').map(event => event[1]), [{ providerId: CLI_PROXY_NATIVE_PROVIDER_ID, @@ -838,6 +741,83 @@ test('registry changes rematerialize and republish only the CLIProxyAPI provider await fixture.lifecycle.cleanup() }) +test('extension registry materializes two exact protected connections and serializes each restart', async () => { + const fixture = gatewayFixture({ kinds: [] }) + await activateCliProxyGateway( + fixture.context, + activationInput(fixture.client, new AbortController().signal), + ) + + fixture.registerExtension({ + adapter: extensionAdapter('header-session', [{ kind: 'header', name: 'Session-Id' }]), + connection: extensionConnection('connection-a', 'header-session'), + source: syntheticExtensions.header, + }) + await new Promise(resolve => setTimeout(resolve, 0)) + fixture.registerExtension({ + adapter: extensionAdapter('body-session', [{ kind: 'body-json-pointer', pointer: '/metadata/thread' }]), + connection: extensionConnection('connection-b', 'body-session'), + source: syntheticExtensions.body, + }) + await new Promise(resolve => setTimeout(resolve, 0)) + + const materializations = fixture.events.filter(event => event[0] === 'materialize') + assert.equal(materializations.length, 3) + const request = materializations.at(-1)[1] + assert.deepEqual(request.sources.map(source => source.source), ['extension-0', 'extension-1']) + assert.deepEqual(compositionLeaves(request), [ + { + targetSlot: 'gateway-extensions', + targetPointer: '/connections/0/endpoint', + source: { kind: 'source-origin', source: 'extension-0', origin: 'api' }, + }, + { + targetSlot: 'gateway-extensions', + targetPointer: '/connections/0/credential', + source: { kind: 'source-authorization', source: 'extension-0' }, + }, + { + targetSlot: 'gateway-extensions', + targetPointer: '/connections/1/endpoint', + source: { kind: 'source-origin', source: 'extension-1', origin: 'api' }, + }, + { + targetSlot: 'gateway-extensions', + targetPointer: '/connections/1/credential', + source: { kind: 'source-authorization', source: 'extension-1' }, + }, + ]) + const roots = rootCollections(request) + assert.equal(roots['gateway-extensions'].enabled, true) + assert.equal(roots['gateway-extensions'].active, true) + assert.deepEqual(roots['gateway-extensions'].adapters.map(adapter => adapter.adapterId), [ + 'body-session', + 'header-session', + ]) + assert.deepEqual( + roots['gateway-extensions'].connections.map(connection => ({ + connectionId: connection.connectionId, + gatewayModelId: connection.models[0].gatewayModelId, + endpoint: connection.endpoint, + credential: connection.credential, + })), + [ + { connectionId: 'connection-a', gatewayModelId: 'connection-a/shared', endpoint: null, credential: null }, + { connectionId: 'connection-b', gatewayModelId: 'connection-b/shared', endpoint: null, credential: null }, + ], + ) + assert.equal(fixture.events.filter(event => event[0] === 'ensure-ready').length, 3) + assert.deepEqual( + fixture.events.filter(event => event[0] === 'publish-native').map(event => + event[1].catalog.routes.map(route => route.alias) + ), + [['connection-a/shared'], ['connection-a/shared', 'connection-b/shared']], + ) + assert.equal(fixture.events.filter(event => event[0] === 'publication-dispose').length, 1) + + await fixture.lifecycle.cleanup() +}) + test('shares async disposal across abort and effect cleanup while completing every cleanup', async () => { const gatewayDisposeError = new Error('synthetic gateway dispose failed') const clientDisposeError = new Error('synthetic client dispose failed') diff --git a/test/native-gateway-bridge/gateway_bridge_test.go b/test/native-gateway-bridge/gateway_bridge_test.go new file mode 100644 index 0000000..ae7c364 --- /dev/null +++ b/test/native-gateway-bridge/gateway_bridge_test.go @@ -0,0 +1,412 @@ +package cordisxgatewayintegration + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "runtime" + "strings" + "sync" + "testing" + "time" + + "github.com/gin-gonic/gin" + "github.com/router-for-me/CLIProxyAPI/v7/internal/config" + "github.com/router-for-me/CLIProxyAPI/v7/internal/pluginhost" + "github.com/router-for-me/CLIProxyAPI/v7/internal/registry" + "github.com/router-for-me/CLIProxyAPI/v7/sdk/api/handlers" + "gopkg.in/yaml.v3" +) + +type receipt struct { + Connection string + Path string + Auth string + Session string + Model string + Body map[string]any + Headers http.Header +} + +type bridgeFixture struct { + t *testing.T + host *pluginhost.Host + handler *handlers.BaseAPIHandler + receipts []receipt + mu sync.Mutex + servers []*httptest.Server +} + +func TestGatewayBridgeUsesRealCLIProxyAPIPluginHost(t *testing.T) { + gin.SetMode(gin.TestMode) + fixture := newBridgeFixture(t, true, true) + + t.Run("two-connections-same-model", func(t *testing.T) { + bodyA, okA := fixture.execute("connection-a/shared", http.Header{ + "Session-Id": {"header-session"}, + "Authorization": {"Bearer attacker"}, + "X-Untrusted-Session": {"attacker-session"}, + }, map[string]any{"metadata": map[string]any{"thread": "body-attacker"}}) + if !okA || string(bodyA) != `{"connection":"connection-a"}` { + t.Fatalf("connection-a response = %s, ok=%v", bodyA, okA) + } + + bodyB, okB := fixture.execute("connection-b/shared", http.Header{ + "Authorization": {"Bearer attacker"}, + "Session-Id": {"header-attacker"}, + }, map[string]any{"metadata": map[string]any{"thread": "body-session"}}) + if !okB || string(bodyB) != `{"connection":"connection-b"}` { + t.Fatalf("connection-b response = %s, ok=%v", bodyB, okB) + } + + got := fixture.snapshot() + if len(got) != 2 { + t.Fatalf("receipts = %d, want 2", len(got)) + } + assertReceipt(t, got[0], "connection-a", "fixture-key-a", "header-session") + assertReceipt(t, got[1], "connection-b", "fixture-key-b", "body-session") + if got[0].Headers.Get("Session-Id") != "" || got[0].Headers.Get("X-Untrusted-Session") != "" || got[1].Headers.Get("Session-Id") != "" { + t.Fatalf("untrusted session carrier reached upstream: %#v", got) + } + if got[0].Headers.Get("Authorization") == "Bearer attacker" || got[1].Headers.Get("Authorization") == "Bearer attacker" { + t.Fatalf("untrusted authorization reached upstream: %#v", got) + } + }) + + t.Run("streaming", func(t *testing.T) { + start := fixture.count() + raw := []byte(`{"model":"connection-a/shared","stream":true,"metadata":{"thread":"ignored"}}`) + ctx := requestContext(context.Background(), raw, http.Header{"Session-Id": {"stream-session"}}) + data, _, errors := fixture.handler.ExecuteStreamWithAuthManager(ctx, "openai", "connection-a/shared", raw, "") + var output strings.Builder + for data != nil || errors != nil { + select { + case chunk, ok := <-data: + if !ok { + data = nil + } else { + output.Write(chunk) + } + case err, ok := <-errors: + if !ok { + errors = nil + } else if err != nil { + t.Fatalf("stream error: %+v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("stream timed out") + } + } + if fixture.count() != start+1 || !strings.Contains(output.String(), "stream-one") || !strings.Contains(output.String(), "[DONE]") { + t.Fatalf("unexpected stream output %q", output.String()) + } + }) + + t.Run("cancellation-closes-host-http-stream", func(t *testing.T) { + started := make(chan struct{}) + canceled := make(chan struct{}) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + close(started) + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(http.StatusOK) + if flusher, ok := w.(http.Flusher); ok { + flusher.Flush() + } + <-r.Context().Done() + close(canceled) + })) + defer server.Close() + fixture.reconfigure(true, true, server.URL) + + raw := []byte(`{"model":"connection-a/shared","stream":true}`) + base, cancel := context.WithCancel(context.Background()) + ctx := requestContext(base, raw, http.Header{"Session-Id": {"cancel-session"}}) + data, _, errors := fixture.handler.ExecuteStreamWithAuthManager(ctx, "openai", "connection-a/shared", raw, "") + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("upstream stream did not start") + } + cancel() + for data != nil || errors != nil { + select { + case _, ok := <-data: + if !ok { + data = nil + } + case _, ok := <-errors: + if !ok { + errors = nil + } + case <-time.After(5 * time.Second): + t.Fatal("canceled downstream stream did not close") + } + } + select { + case <-canceled: + case <-time.After(5 * time.Second): + t.Fatal("host HTTP stream remained open after cancellation") + } + fixture.reconfigure(true, true, "") + }) + + t.Run("required-adapter-failures-never-reach-upstream", func(t *testing.T) { + start := fixture.count() + _, ok := fixture.execute("connection-a/shared", nil, map[string]any{}) + if ok || fixture.count() != start { + t.Fatal("missing required session reached upstream") + } + + fixture.reconfigure(false, true, "") + _, ok = fixture.execute("connection-a/shared", http.Header{"Session-Id": {"blocked"}}, map[string]any{}) + if ok || fixture.count() != start { + t.Fatal("disabled required adapter reached upstream") + } + + fixture.reconfigure(true, false, "") + _, ok = fixture.execute("connection-a/shared", http.Header{"Session-Id": {"blocked"}}, map[string]any{}) + if ok || fixture.count() != start { + t.Fatal("mismatched required adapter reached upstream") + } + + fixture.reconfigure(true, true, "") + body, ok := fixture.execute("connection-a/shared", http.Header{"Session-Id": {"restored"}}, map[string]any{}) + if !ok || string(body) != `{"connection":"connection-a"}` || fixture.count() != start+1 { + t.Fatalf("restored adapter did not resume exact route: %s, ok=%v", body, ok) + } + }) + + t.Run("disabled-bridge-owns-protected-model-without-fallthrough", func(t *testing.T) { + start := fixture.count() + fixture.reconfigure(true, true, "", false) + registry.GetGlobalRegistry().RegisterClient("fallback-client", "openai", []*registry.ModelInfo{{ + ID: "connection-a/shared", Object: "model", + }}) + defer registry.GetGlobalRegistry().UnregisterClient("fallback-client") + _, ok := fixture.execute("connection-a/shared", http.Header{"Session-Id": {"blocked"}}, map[string]any{}) + if ok || fixture.count() != start { + t.Fatal("disabled bridge fell through to a same-id provider") + } + }) +} + +func newBridgeFixture(t *testing.T, adapterEnabled, matchingRevision bool) *bridgeFixture { + t.Helper() + root := os.Getenv("CORDISX_GATEWAY_REPO") + if root == "" { + var err error + root, err = filepath.Abs(filepath.Join("..", "..")) + if err != nil { + t.Fatal(err) + } + } + pluginRoot := filepath.Join(root, "runtime", "cli-proxy-plugins") + extension := pluginhost.PluginExtension(runtime.GOOS) + pluginPath := filepath.Join(pluginRoot, runtime.GOOS, runtime.GOARCH, "cordisx-gateway-bridge"+extension) + if _, err := os.Stat(pluginPath); err != nil { + t.Fatalf("native bridge is missing at %s: %v", pluginPath, err) + } + fixture := &bridgeFixture{t: t} + makeServer := func(connection string) *httptest.Server { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + raw, _ := io.ReadAll(r.Body) + body := map[string]any{} + _ = json.Unmarshal(raw, &body) + metadata, _ := body["metadata"].(map[string]any) + fixture.mu.Lock() + fixture.receipts = append(fixture.receipts, receipt{ + Connection: connection, + Path: r.URL.Path, + Auth: r.Header.Get("Authorization"), + Session: r.Header.Get("X-Session"), + Model: fmt.Sprint(body["model"]), + Body: body, + Headers: r.Header.Clone(), + }) + fixture.mu.Unlock() + w.Header().Set("Content-Type", "text/event-stream") + if fmt.Sprint(body["stream"]) == "true" { + _, _ = io.WriteString(w, "data: {\"choices\":[{\"delta\":{\"content\":\"stream-one\"}}]}\n\n") + _, _ = io.WriteString(w, "data: [DONE]\n\n") + return + } + _ = metadata + _, _ = io.WriteString(w, fmt.Sprintf(`{"connection":%q}`, connection)) + })) + fixture.servers = append(fixture.servers, server) + return server + } + a, b := makeServer("connection-a"), makeServer("connection-b") + fixture.host = pluginhost.New() + initialConfig := config.Config{} + fixture.handler = handlers.NewBaseAPIHandlers(&initialConfig.SDKConfig, nil) + fixture.handler.SetPluginHost(fixture.host) + fixture.reconfigureWithEndpoints(pluginRoot, a.URL, b.URL, adapterEnabled, matchingRevision, true) + t.Cleanup(func() { + fixture.host.ShutdownAll() + for _, server := range fixture.servers { + server.Close() + } + }) + return fixture +} + +func (fixture *bridgeFixture) reconfigure(adapterEnabled, matchingRevision bool, overrideA string, enabled ...bool) { + fixture.t.Helper() + pluginRoot := filepath.Dir(filepath.Dir(filepath.Dir(fixture.pluginPath()))) + a := fixture.servers[0].URL + if overrideA != "" { + a = overrideA + } + bridgeEnabled := true + if len(enabled) > 0 { + bridgeEnabled = enabled[0] + } + fixture.reconfigureWithEndpoints(pluginRoot, a, fixture.servers[1].URL, adapterEnabled, matchingRevision, bridgeEnabled) +} + +func (fixture *bridgeFixture) reconfigureWithEndpoints(pluginRoot, endpointA, endpointB string, adapterEnabled, matchingRevision, bridgeEnabled bool) { + fixture.t.Helper() + adapterRevision := 1 + if !matchingRevision { + adapterRevision = 2 + } + connectionAEnabled := adapterEnabled && matchingRevision + if !connectionAEnabled { + endpointA = "" + } + configYAML := fmt.Sprintf(`plugins: + enabled: true + dir: %q + configs: + cordisx-gateway-bridge: + enabled: true + active: %t + priority: 100 + revision: fixture + adapters: + - adapterId: header-session + revision: 1 + enabled: %t + requiredSession: true + sessionSources: + - kind: header + name: Session-Id + request: + clearHeaders: [Session-Id, X-Untrusted-Session] + setHeaders: + X-Session: "{{session}}" + X-Connection: "{{connectionId}}" + setBody: + /metadata/session: "{{session}}" + - adapterId: body-session + revision: 1 + enabled: true + requiredSession: true + sessionSources: + - kind: body-json-pointer + pointer: /metadata/thread + request: + clearHeaders: [Session-Id] + setHeaders: + X-Session: "{{session}}" + setBody: + /metadata/session: "{{session}}" + connections: + - connectionId: connection-a + revision: 1 + enabled: %t + adapterId: header-session + adapterRevision: %d + adapterRequired: true + adapterAvailable: %t + wireApi: responses + endpointPath: /responses + endpoint: %q + authorization: bearer + credential: fixture-key-a + models: + - sourceModelId: shared-model + gatewayModelId: connection-a/shared + displayName: Shared A + isDefault: true + - connectionId: connection-b + revision: 1 + enabled: true + adapterId: body-session + adapterRevision: 1 + adapterRequired: true + adapterAvailable: true + wireApi: responses + endpointPath: /responses + endpoint: %q + authorization: bearer + credential: fixture-key-b + models: + - sourceModelId: shared-model + gatewayModelId: connection-b/shared + displayName: Shared B + isDefault: true +`, pluginRoot, bridgeEnabled, adapterEnabled, connectionAEnabled, adapterRevision, matchingRevision, endpointA, endpointB) + var cfg config.Config + if err := yaml.Unmarshal([]byte(configYAML), &cfg); err != nil { + fixture.t.Fatal(err) + } + fixture.host.ApplyConfig(context.Background(), &cfg) + fixture.handler.UpdateClients(&cfg.SDKConfig) +} + +func (fixture *bridgeFixture) pluginPath() string { + root := os.Getenv("CORDISX_GATEWAY_REPO") + if root == "" { + root, _ = filepath.Abs(filepath.Join("..", "..")) + } + return filepath.Join(root, "runtime", "cli-proxy-plugins", runtime.GOOS, runtime.GOARCH, "cordisx-gateway-bridge"+pluginhost.PluginExtension(runtime.GOOS)) +} + +func (fixture *bridgeFixture) execute(model string, headers http.Header, body map[string]any) ([]byte, bool) { + fixture.t.Helper() + body["model"] = model + raw, err := json.Marshal(body) + if err != nil { + fixture.t.Fatal(err) + } + ctx := requestContext(context.Background(), raw, headers) + response, _, errMessage := fixture.handler.ExecuteWithAuthManager(ctx, "openai-response", model, raw, "") + return response, errMessage == nil +} + +func requestContext(parent context.Context, body []byte, headers http.Header) context.Context { + ginContext, _ := gin.CreateTestContext(httptest.NewRecorder()) + ginContext.Request = httptest.NewRequest(http.MethodPost, "http://127.0.0.1/v1/responses", strings.NewReader(string(body))).WithContext(parent) + ginContext.Request.Header = headers.Clone() + return context.WithValue(parent, "gin", ginContext) +} + +func (fixture *bridgeFixture) count() int { + fixture.mu.Lock() + defer fixture.mu.Unlock() + return len(fixture.receipts) +} + +func (fixture *bridgeFixture) snapshot() []receipt { + fixture.mu.Lock() + defer fixture.mu.Unlock() + return append([]receipt(nil), fixture.receipts...) +} + +func assertReceipt(t *testing.T, got receipt, connection, credential, session string) { + t.Helper() + if got.Connection != connection || got.Path != "/responses" || got.Auth != "Bearer "+credential || got.Session != session || got.Model != "shared-model" { + t.Fatalf("unexpected receipt: %#v", got) + } + metadata, _ := got.Body["metadata"].(map[string]any) + if fmt.Sprint(metadata["session"]) != session { + t.Fatalf("body session = %#v, want %q", metadata["session"], session) + } +} diff --git a/test/native-gateway-bridge/go.mod b/test/native-gateway-bridge/go.mod new file mode 100644 index 0000000..18e538c --- /dev/null +++ b/test/native-gateway-bridge/go.mod @@ -0,0 +1,9 @@ +module github.com/router-for-me/CLIProxyAPI/v7/cordisxgatewayintegration + +go 1.26.0 + +require ( + github.com/gin-gonic/gin v1.10.1 + github.com/router-for-me/CLIProxyAPI/v7 v7.2.137 + gopkg.in/yaml.v3 v3.0.1 +) diff --git a/test/package-artifact.mjs b/test/package-artifact.mjs index 5c9e412..79e842b 100644 --- a/test/package-artifact.mjs +++ b/test/package-artifact.mjs @@ -13,44 +13,19 @@ const execFileAsync = promisify(execFile) const schemas = process.env.CORDISX_PROTOCOL_ROOT === undefined ? new URL('../node_modules/@cordisx/protocol/schemas/', import.meta.url) : new URL('schemas/', pathToFileURL(`${resolve(process.env.CORDISX_PROTOCOL_ROOT)}/`)) -const gatewayRuntimeResources = [ - { - path: './config/cli-proxy-api.yaml', - mode: 'data', - digest: 'sha256:4e2806a2c9a4a490f8685ffc96cf1dc40010248a8d25437b718249f3275b198e', - byteLength: 268, - }, - { - path: './schemas/cli-proxy-gateway-model-catalog.v1.schema.json', - mode: 'data', - digest: 'sha256:b20bf7d388ae1417eaa1e86186eab241c754fc54af3c3b4d5439308202194b30', - byteLength: 639, - }, - { - path: './schemas/cli-proxy-gateway-codex-upstream.v1.schema.json', - mode: 'data', - digest: 'sha256:a9908de439af944ac24744721c117fd0c0a18c18e61355d343f8ad57311fafba', - byteLength: 1253, - }, - { - path: './schemas/cli-proxy-gateway-openai-upstream.v1.schema.json', - mode: 'data', - digest: 'sha256:8f226eeaf778394f452c0c224bd163c1c68806f56c63982d0394c03d6be4483b', - byteLength: 1925, - }, - { - path: './schemas/cli-proxy-management-auth-files.v1.schema.json', - mode: 'data', - digest: 'sha256:9785b09906567848143054d522b141f5f06234f7b16a8dc2a02ef31e6d58e76c', - byteLength: 370, - }, - { - path: './schemas/cli-proxy-management-operation.v1.schema.json', - mode: 'data', - digest: 'sha256:c1a9e12cf48ee53825948b053ea364ce70d09122321502560defebeb3d8d0c12', - byteLength: 218, - }, +const gatewayResourcePaths = [ + './config/cli-proxy-api.yaml', + './schemas/cli-proxy-gateway-model-catalog.v1.schema.json', + './schemas/cli-proxy-gateway-codex-upstream.v1.schema.json', + './schemas/cli-proxy-gateway-openai-upstream.v1.schema.json', + './schemas/cli-proxy-gateway-adapter.v1.schema.json', + './schemas/cli-proxy-gateway-connection.v1.schema.json', + './schemas/cli-proxy-gateway-extension-plan.v1.schema.json', + './schemas/cli-proxy-management-auth-files.v1.schema.json', + './schemas/cli-proxy-management-operation.v1.schema.json', ] +const nativeResourceRoot = new URL('../runtime/cli-proxy-plugins/', import.meta.url) +const nativeExtensions = { darwin: 'dylib', linux: 'so', win32: 'dll' } const gatewayConsumerOperations = [ 'gateway.models.list', @@ -106,19 +81,48 @@ test('publishes schema-valid v14 package and runtime manifests with exact reques for (const service of runtimeManifest.services) { await readFile(new URL(`..${service.entry.slice(1)}`, import.meta.url)) } - assert.deepEqual(runtimeManifest.services[1], { + const gateway = runtimeManifest.services[1] + assert.deepEqual({ ...gateway, runtimeResources: undefined }, { id: 'gateway-runtime', kind: 'managed-backend', owner: 'host', entry: './dist/gateway.mjs', definitionSchema: 'https://raw.githubusercontent.com/cordisx/cordisx-protocol/main/schemas/managed-service-definition.v1.schema.json', - runtimeResources: gatewayRuntimeResources, + runtimeResources: undefined, consumerGrants: [{ pluginId: 'cli-proxy-api', operations: gatewayConsumerOperations }], }) - for (const resource of runtimeManifest.services[1].runtimeResources) { + assert.deepEqual( + gateway.runtimeResources.filter(resource => resource.mode === 'data').map(resource => resource.path), + gatewayResourcePaths, + ) + const executableResources = gateway.runtimeResources.filter(resource => resource.mode === 'executable') + const actualExecutablePaths = [] + for (const platform of await readdir(nativeResourceRoot)) { + for (const architecture of await readdir(new URL(`${platform}/`, nativeResourceRoot))) { + const directory = new URL(`${platform}/${architecture}/`, nativeResourceRoot) + for (const file of await readdir(directory)) { + if (/^cordisx-gateway-bridge\.(?:dylib|so|dll)$/.test(file)) { + actualExecutablePaths.push(`./runtime/cli-proxy-plugins/${platform}/${architecture}/${file}`) + } + } + } + } + assert.deepEqual(executableResources.map(resource => resource.path).sort(), actualExecutablePaths.sort()) + for (const resource of executableResources) { + const match = + /^\.\/runtime\/cli-proxy-plugins\/(darwin|linux|win32)\/(arm64|x64)\/cordisx-gateway-bridge\.(dylib|so|dll)$/ + .exec( + resource.path, + ) + assert.notEqual(match, null, `invalid native resource path ${resource.path}`) + const [, platform, architecture, extension] = match + assert.equal(extension, nativeExtensions[platform], `native resource extension does not match ${platform}`) + assert.deepEqual(resource.platforms, [platform]) + assert.deepEqual(resource.architectures, [architecture]) + } + for (const resource of gateway.runtimeResources) { const bytes = await readFile(new URL(`../${resource.path.slice(2)}`, import.meta.url)) - assert.equal(resource.mode, 'data') assert.equal(resource.byteLength, bytes.byteLength) assert.equal(resource.digest, `sha256:${createHash('sha256').update(bytes).digest('hex')}`) } @@ -132,6 +136,9 @@ test('packs gateway runtime data resources into the tarball', async () => { const [{ files }] = JSON.parse(stdout) const archiveFiles = new Set(files.map(file => file.path)) assert.equal(archiveFiles.has('assets/icon.png'), true, 'tarball is missing the plugin brand PNG') + const runtimeManifest = JSON.parse(await readFile(new URL('../runtime-manifest.json', import.meta.url), 'utf8')) + const gatewayRuntimeResources = + runtimeManifest.services.find(service => service.id === 'gateway-runtime').runtimeResources for (const resource of gatewayRuntimeResources) { const path = resource.path.slice(2) assert.equal(archiveFiles.has(path), true, `tarball is missing ${path}`) @@ -146,7 +153,7 @@ test('packs gateway runtime data resources into the tarball', async () => { } }) -test('publishes the versioned upstream registration entrypoint', async () => { +test('publishes the versioned upstream and gateway extension registration entrypoints', async () => { const packageManifest = JSON.parse(await readFile(new URL('../package.json', import.meta.url), 'utf8')) assert.deepEqual(packageManifest.exports['./upstream/v1'], { types: './dist/types/service/upstream-registry.d.ts', @@ -156,6 +163,14 @@ test('publishes the versioned upstream registration entrypoint', async () => { const upstream = await import('@cordisx/plugin-cli-proxy-api/upstream/v1') assert.equal(upstream.CLI_PROXY_UPSTREAM_REGISTRY_SERVICE_V1, 'cliProxyUpstreams') assert.equal(typeof upstream.CliProxyUpstreamRegistryV1, 'function') + assert.deepEqual(packageManifest.exports['./extensions/v1'], { + types: './dist/types/service/extension-registry.d.ts', + import: './dist/extensions.mjs', + default: './dist/extensions.mjs', + }) + const extensions = await import('@cordisx/plugin-cli-proxy-api/extensions/v1') + assert.equal(extensions.CLI_PROXY_GATEWAY_EXTENSION_REGISTRY_SERVICE_V1, 'cliProxyGatewayExtensions') + assert.equal(typeof extensions.CliProxyGatewayExtensionRegistryV1, 'function') }) test('publishes the managed gateway Node module and public Protocol context ABI', async () => { @@ -166,8 +181,11 @@ test('publishes the managed gateway Node module and public Protocol context ABI' default: './dist/gateway.mjs', }) const gateway = await import('@cordisx/plugin-cli-proxy-api/gateway/v1') - assert.equal(gateway.contextServices[0].service, 'cliProxyUpstreams') - assert.deepEqual(gateway.apply.inject, ['managedServices', 'cliProxyUpstreams']) + assert.deepEqual(gateway.contextServices.map(service => service.service), [ + 'cliProxyUpstreams', + 'cliProxyGatewayExtensions', + ]) + assert.deepEqual(gateway.apply.inject, ['managedServices', 'cliProxyUpstreams', 'cliProxyGatewayExtensions']) }) test('publishes a schema-valid managed gateway definition', async () => { diff --git a/test/support/gateway-extension-fixture.mjs b/test/support/gateway-extension-fixture.mjs new file mode 100644 index 0000000..8fab24e --- /dev/null +++ b/test/support/gateway-extension-fixture.mjs @@ -0,0 +1,247 @@ +import assert from 'node:assert/strict' +import { readFile } from 'node:fs/promises' +import Ajv2020 from 'ajv/dist/2020.js' +import addFormats from 'ajv-formats' +import { identity } from './gateway-fixture.mjs' + +const gatewayModule = await import('../../dist/gateway.mjs') +const { + CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1, + CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1, + CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1, + CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1, + CLI_PROXY_UPSTREAM_DESCRIPTOR_CONTRACT_V1, + CLI_PROXY_UPSTREAM_DESCRIPTOR_SCHEMA_V1, + cliProxyGatewayDefinition, + contextServices, +} = gatewayModule + +const compositionSchemaFiles = new Map([ + ['codex-upstreams', '../../schemas/cli-proxy-gateway-codex-upstream.v1.schema.json'], + ['openai-upstreams', '../../schemas/cli-proxy-gateway-openai-upstream.v1.schema.json'], + ['gateway-extensions', '../../schemas/cli-proxy-gateway-extension-plan.v1.schema.json'], +]) +const compositionValidators = new Map( + await Promise.all([...compositionSchemaFiles].map(async ([slot, file]) => { + const ajv = new Ajv2020({ allErrors: true, strict: true }) + addFormats(ajv) + for ( + const dependency of [ + '../../schemas/cli-proxy-gateway-adapter.v1.schema.json', + '../../schemas/cli-proxy-gateway-connection.v1.schema.json', + ] + ) ajv.addSchema(JSON.parse(await readFile(new URL(dependency, import.meta.url), 'utf8'))) + const schema = JSON.parse(await readFile(new URL(file, import.meta.url), 'utf8')) + return [slot, ajv.compile(schema)] + })), +) + +export const syntheticSources = Object.freeze({ + responses: Object.freeze({ + pluginId: 'synthetic-responses', + pluginGeneration: 'synthetic-responses-one', + serviceId: 'synthetic-responses-service', + serviceGeneration: 'synthetic-responses-service-one', + upstreamId: 'responses-source', + displayName: 'Synthetic Responses', + prefix: 'responses', + wireApi: 'responses', + sourceModelId: 'responses-model', + modelId: 'family/fast', + authenticated: false, + }), + chat: Object.freeze({ + pluginId: 'synthetic-chat', + pluginGeneration: 'synthetic-chat-one', + serviceId: 'synthetic-chat-service', + serviceGeneration: 'synthetic-chat-service-one', + upstreamId: 'chat-source', + displayName: 'Synthetic Chat', + prefix: 'chat', + wireApi: 'chat-completions', + sourceModelId: 'chat-model', + modelId: 'family/balanced', + authenticated: true, + }), +}) + +export const syntheticExtensions = Object.freeze({ + header: Object.freeze({ + pluginId: 'synthetic-header-extension', + pluginGeneration: 'synthetic-header-extension-one', + serviceId: 'synthetic-header-connection', + serviceGeneration: 'synthetic-header-connection-one', + authenticated: true, + }), + body: Object.freeze({ + pluginId: 'synthetic-body-extension', + pluginGeneration: 'synthetic-body-extension-one', + serviceId: 'synthetic-body-connection', + serviceGeneration: 'synthetic-body-connection-one', + authenticated: true, + }), +}) + +export const extensionAdapter = (adapterId, sessionSources) => ({ + $schema: CLI_PROXY_GATEWAY_ADAPTER_SCHEMA_V1, + contract: CLI_PROXY_GATEWAY_ADAPTER_CONTRACT_V1, + schemaVersion: 1, + adapterId, + revision: 1, + enabled: true, + requiredSession: true, + sessionSources, + request: { + clearHeaders: ['X-Untrusted-Session'], + setHeaders: { 'X-Session': '{{session}}', 'X-Connection': '{{connectionId}}' }, + setBody: { '/metadata/session': '{{session}}' }, + }, +}) + +export const extensionConnection = (connectionId, adapterId) => ({ + $schema: CLI_PROXY_GATEWAY_CONNECTION_SCHEMA_V1, + contract: CLI_PROXY_GATEWAY_CONNECTION_CONTRACT_V1, + schemaVersion: 1, + connectionId, + revision: 1, + enabled: true, + wireApi: 'responses', + endpointPath: '/v1/responses', + authorization: 'bearer', + adapter: { adapterId, revision: 1, required: true }, + models: [{ + sourceModelId: 'shared-model', + modelId: 'shared', + displayName: `Shared ${connectionId}`, + enabled: true, + isDefault: true, + }], +}) + +export const descriptor = (source, overrides = {}) => ({ + $schema: CLI_PROXY_UPSTREAM_DESCRIPTOR_SCHEMA_V1, + contract: CLI_PROXY_UPSTREAM_DESCRIPTOR_CONTRACT_V1, + schemaVersion: 1, + revision: 1, + upstreamId: source.upstreamId, + displayName: source.displayName, + enabled: true, + wireApi: source.wireApi, + prefix: source.prefix, + order: source.wireApi === 'responses' ? 10 : 20, + requestTimeoutMs: 30_000, + group: { groupId: 'synthetic', displayName: 'Synthetic', order: 10 }, + models: [{ + sourceModelId: source.sourceModelId, + modelId: source.modelId, + enabled: true, + isDefault: source.wireApi === 'responses', + ...(source.wireApi === 'chat-completions' ? { inputModalities: ['text'] } : {}), + }], + ...overrides, +}) + +export function providerRoot() { + const controller = new AbortController() + const target = { + pluginId: 'cli-proxy-api', + pluginGeneration: 'gateway-generation', + serviceId: 'gateway-runtime', + serviceGeneration: 'gateway-service-generation', + } + return { + controller, + target, + root: contextServices[0].create({ target, signal: controller.signal }), + extensionRoot: contextServices[1].create({ target, signal: controller.signal }), + } +} + +export function validateMaterializationRequest(request) { + const compositionSlots = cliProxyGatewayDefinition.protectedBindings.filter(item => item.source === 'composition') + const declaredSources = new Set(request.sources.map(source => source.source)) + for (const slot of compositionSlots) { + const slotBindings = request.bindings.filter(binding => binding.targetSlot === slot.slot) + assert.ok(slotBindings.length > 0, `missing composition slot ${slot.slot}`) + const roots = slotBindings.filter(binding => binding.targetPointer === undefined) + assert.equal(roots.length, 1, `composition slot ${slot.slot} must have exactly one root binding`) + assert.equal(roots[0].source.kind, 'safe-literal') + const validate = compositionValidators.get(slot.slot) + assert.equal(validate(roots[0].source.value), true, JSON.stringify(validate.errors)) + } + for (const item of request.bindings) { + if (item.source.kind === 'source-origin' || item.source.kind === 'source-authorization') { + assert.equal(declaredSources.has(item.source.source), true, `undeclared composition source ${item.source.source}`) + } + } +} + +export function activationInput(client, signal) { + return { + owner: { + ownerHandle: 'mso_gateway', + pluginId: 'cli-proxy-api', + sourceDigest: `sha256:${'a'.repeat(64)}`, + hostGeneration: 'host-one', + pluginGeneration: 'gateway-generation', + }, + client, + signal, + } +} + +export function rootCollections(request) { + return Object.fromEntries( + request.bindings.filter(item => item.targetPointer === undefined).map(item => [ + item.targetSlot, + item.source.value, + ]), + ) +} + +export function compositionLeaves(request) { + return request.bindings.filter(item => item.targetPointer !== undefined) +} + +export function expectedRootCollections(sources, extensionRevision) { + return { + 'codex-upstreams': sources.filter(source => source.wireApi === 'responses').map(source => ({ + 'api-key': null, + prefix: source.prefix, + 'base-url': null, + 'request-retry': 0, + models: [{ + name: source.sourceModelId, + alias: source.modelId, + 'display-name': source.modelId, + 'force-mapping': true, + }], + })), + 'openai-upstreams': sources.filter(source => source.wireApi === 'chat-completions').map(source => ({ + name: source.upstreamId, + disabled: false, + prefix: source.prefix, + 'base-url': null, + 'api-key-entries': [{ 'api-key': null }], + 'request-retry': 0, + models: [{ + name: source.sourceModelId, + alias: source.modelId, + 'display-name': source.modelId, + 'force-mapping': true, + 'input-modalities': ['text'], + }], + })), + 'gateway-extensions': { + enabled: true, + active: false, + revision: extensionRevision, + adapters: [], + connections: [], + }, + } +} + +export function extensionSourceIdentity(source) { + return identity(source.pluginId, source.serviceId) +}