import { VehicleRegistry } from "@danypops/vehicle-server"; import { createVehicleHttpApp } from "@danypops/vehicle-server/http"; import { errorResponse, healthResponse, readyResponse, requireBearerToken } from "@danypops/vehicle-server/rpc-http"; import { SERVICE_MAX_BODY_BYTES, SERVICE_MAX_RESPONSE_BYTES } from "../constants.ts"; import type { CacheEconomicsSummary } from "../observability/cache-economics.ts"; import type { ContextDelta, ContextSnapshot } from "../observability/context-delta.ts"; import { type ContextSnapshotHistory, MetricContextSnapshotHistory } from "../observability/context-snapshot-history.ts"; import type { CompactionDurationEstimate, ContextAssessment } from "../observability/context-telemetry.ts"; import type { MetricObservation, MetricQuery, StoredMetricObservation } from "../observability/metric.ts"; import type { MetricStore } from "../observability/store.ts"; import type { SubscriptionResets } from "../observability/subscription-resets.ts"; import type { TaskCostSummary } from "../observability/task-cost.ts"; import type { UsageAggregateRow } from "../observability/usage.ts"; import type { UsageImportController, UsageImportResult, UsageImportStatus } from "../observability/usage-import.ts"; import type { BenchmarkQuery, BenchmarkQueryResult, BenchmarkRefreshResult } from "../optimization/model-selection/benchmark.ts"; import type { ModelCatalogController, ModelCatalogQuery, ModelCatalogQueryResult, ModelCatalogStatus, } from "../optimization/model-selection/catalog.ts"; import type { BenchmarkController } from "../optimization/model-selection/controller.ts"; import type { ModelRanker, ModelRecommendationInput } from "../optimization/model-selection/ranker.ts"; import type { ModelRankingResult } from "../optimization/model-selection/ranking.ts"; import type { RouteOverride, RouterController, RouterStatus, TelemetryPollResult } from "../optimization/routing/controller.ts"; import type { PolicyDecision, Route } from "../optimization/routing/policy.ts"; import { InvalidSessionSecretError, type RegisterSessionIdentityResult, type SessionIdentity } from "../sessions/identity.ts"; import { routerMutationAuthorizer } from "../sessions/router-authorization.ts"; import { DisabledObservationExporter, type ObservationExporter, type ObservationExportStatus } from "../telemetry-export/exporter.ts"; import { VERSION } from "../version.ts"; import { benchmarkOperations } from "./benchmark-operations.ts"; import { cacheEconomicsOperations } from "./cache-operations.ts"; import { catalogOperations } from "./catalog-operations.ts"; import { contextOperations } from "./context-operations.ts"; import { exportOperations } from "./export-operations.ts"; import { metricsOperations } from "./metric-operations.ts"; import { modelRankingOperations } from "./model-ranking-operations.ts"; import { EXPECTED_OPERATION_NAMES, type OperationHandlerMap, type OperationName } from "./operation-types.ts"; import { registerJittorVehicleOperations } from "./registration.ts"; import { isResetOperation, type ResetInputs, type ResetOutputs } from "./reset-contracts.ts"; import { routerOperations } from "./routing-operations.ts"; import { sessionIdentityOperations } from "./session-operations.ts"; import { subscriptionResetOperations } from "./subscription-reset-operations.ts"; import { usageImportOperations } from "./usage-import-operations.ts"; export { EXPECTED_OPERATION_NAMES, type OperationName }; interface RouterScopeInput { session_id?: string; session_secret?: string; } export interface OperationInputs extends ResetInputs { "session.register": { session_id: string }; "session.release": { session_id: string; session_secret?: string }; "metrics.record": MetricObservation; "metrics.record_batch": { observations: MetricObservation[] }; "metrics.query": MetricQuery; "metrics.distinct_scopes": { source: string; since: number; until: number; limit?: number }; "metrics.usage_series": { source: string; since: number; until: number; bucketSizeMs: number; bucketCount: number; scopeLimit?: number }; "metrics.cost_by_task": { since: number; until: number }; "metrics.prune": { before: number; force?: boolean }; "benchmark.refresh": { force?: boolean }; "benchmark.status": Record; "benchmark.query": BenchmarkQuery; "catalog.refresh": { force?: boolean }; "catalog.status": Record; "catalog.query": ModelCatalogQuery; "usage.import": { dryRun?: boolean }; "usage.import_status": Record; "usage.import_cancel": Record; "export.status": Record; "export.flush": Record; "models.rank": ModelRecommendationInput & RouterScopeInput; "context.assess": { since?: number; until?: number }; "context.delta": { session_id: string }; "context.snapshot": ContextSnapshot; "compaction.estimate": Record; "service.checkpoint": Record; "telemetry.poll": Record; "router.status": RouterScopeInput; "router.decide": RouterScopeInput; "router.pause": RouterScopeInput; "router.resume": RouterScopeInput; "router.override": RouteOverride & RouterScopeInput; "router.clear_override": RouterScopeInput; "router.current_route": Route & RouterScopeInput; "router.available_routes": { routes: Route[] } & RouterScopeInput; "cache.economics": { since: number; until: number }; } export interface OperationOutputs extends ResetOutputs { "session.register": RegisterSessionIdentityResult; "session.release": { released: boolean }; "metrics.record": StoredMetricObservation; "metrics.record_batch": StoredMetricObservation[]; "metrics.query": StoredMetricObservation[]; "metrics.distinct_scopes": string[]; "metrics.usage_series": { rows: UsageAggregateRow[]; truncated: boolean }; "metrics.cost_by_task": TaskCostSummary; "metrics.prune": { deleted: number }; "benchmark.refresh": BenchmarkRefreshResult; "benchmark.status": BenchmarkRefreshResult; "benchmark.query": BenchmarkQueryResult; "catalog.refresh": ModelCatalogStatus; "catalog.status": ModelCatalogStatus; "catalog.query": ModelCatalogQueryResult; "usage.import": UsageImportResult; "usage.import_status": UsageImportStatus; "usage.import_cancel": UsageImportStatus; "export.status": ObservationExportStatus; "export.flush": ObservationExportStatus; "models.rank": ModelRankingResult; "context.assess": ContextAssessment; "context.delta": ContextDelta | null; "context.snapshot": ContextDelta; "compaction.estimate": CompactionDurationEstimate; "service.checkpoint": { ok: true }; "telemetry.poll": TelemetryPollResult; "router.status": RouterStatus; "router.decide": PolicyDecision; "router.pause": RouterStatus; "router.resume": RouterStatus; "router.override": RouterStatus; "router.clear_override": RouterStatus; "router.current_route": RouterStatus; "router.available_routes": RouterStatus; "cache.economics": CacheEconomicsSummary; } export class UnknownOperationError extends Error {} export { InvalidSessionSecretError }; class UnavailableUsageImporter implements UsageImportController { async run(): Promise { throw new Error("historical usage import is not configured"); } status(): UsageImportStatus { return { running: false, cancelRequested: false, lastResult: null }; } cancel(): UsageImportStatus { return this.status(); } } class UnavailableModelCatalog implements ModelCatalogController { async refresh(): Promise { return this.status(); } status(): ModelCatalogStatus { return { configured: false, ok: null, hasSnapshot: false, lastAttemptAt: null, lastSuccessAt: null, entries: 0, revision: null, }; } query(): ModelCatalogQueryResult { throw new Error("model catalog is not configured"); } } class UnavailableModelRanker implements ModelRanker { rank(): ModelRankingResult { throw new Error("model ranking is not configured"); } } class UnavailableBenchmarkController implements BenchmarkController { async refresh(): Promise { return this.status(); } status(): BenchmarkRefreshResult { return { observedAt: Date.now(), sources: [] }; } query(): BenchmarkQueryResult { throw new Error("benchmark evidence is not configured"); } } class UnavailableRouter implements RouterController { private readonly unavailable: RouterStatus = { ready: false, paused: false, sources: [], lastDecision: null, override: null, currentRoute: null, availableRoutes: [], }; async poll(): Promise { return { sources: [], observedAt: Date.now() }; } status(): RouterStatus { return structuredClone(this.unavailable); } decide(): PolicyDecision { return { action: "halt", pressure: Number.POSITIVE_INFINITY, reason: "router is not configured", decidedAt: Date.now(), trace: ["fail closed"], }; } pause(): RouterStatus { return this.status(); } resume(): RouterStatus { return this.status(); } setOverride(): RouterStatus { return this.status(); } clearOverride(): RouterStatus { return this.status(); } setCurrentRoute(): RouterStatus { return this.status(); } setAvailableRoutes(): RouterStatus { return this.status(); } } export class JittorService { private readonly router: RouterController; private readonly operations: OperationHandlerMap; /** Every jittor operation, also projected onto the real Vehicle protocol -- see vehicle-registration.ts. Served alongside (not replacing) the /api/v1/ops route below. */ readonly vehicleRegistry: VehicleRegistry; constructor( private readonly metrics: MetricStore, router: RouterController = new UnavailableRouter(), benchmarks: BenchmarkController = new UnavailableBenchmarkController(), modelRanker: ModelRanker = new UnavailableModelRanker(), sessionIdentity?: SessionIdentity, contextSnapshots: ContextSnapshotHistory = new MetricContextSnapshotHistory(metrics), catalog: ModelCatalogController = new UnavailableModelCatalog(), usageImporter: UsageImportController = new UnavailableUsageImporter(), exporter: ObservationExporter = new DisabledObservationExporter(), resets?: SubscriptionResets, ) { this.router = router; const authorize = routerMutationAuthorizer(sessionIdentity); // Each capability module owns a disjoint, bounded slice of EXPECTED_OPERATION_NAMES and only // the collaborators it needs -- adding a new operation domain means adding a new module here, // not another switch case in a single responsibility magnet. this.operations = { ...subscriptionResetOperations(resets), ...metricsOperations(metrics), ...benchmarkOperations(benchmarks), ...catalogOperations(catalog), ...usageImportOperations(usageImporter), ...exportOperations(exporter), ...contextOperations(metrics, contextSnapshots), ...routerOperations(router, authorize), ...modelRankingOperations(modelRanker, router, authorize), ...sessionIdentityOperations(sessionIdentity), ...cacheEconomicsOperations(metrics, catalog), }; this.vehicleRegistry = new VehicleRegistry({ name: "jittor", packageJsonUrl: new URL("../../package.json", import.meta.url), description: "Just-in-Time Token Optimization Router for Pi -- token and context observability with optimization policies.", }); // withJittorErrorParity() already converts every real handler error into a well-formed // VehicleError, so this only affects a genuine registration/binding bug. this.vehicleRegistry.setExposeHandlerFailureDetails(true); registerJittorVehicleOperations(this.vehicleRegistry, this.operations); } operationNames(): OperationName[] { return [...EXPECTED_OPERATION_NAMES]; } async execute(operation: Name, input: OperationInputs[Name]): Promise; async execute(operation: string, input: Record): Promise; async execute(operation: string, input: Record = {}): Promise { const handler = this.operations[operation]; if (!handler) throw new UnknownOperationError(`unknown operation: ${operation}`); return handler(input); } ready(): boolean { return this.router.status(undefined).ready; } close(): void { this.metrics.close(); } } export interface JittorAppOptions { service: JittorService; token: string; maxBodyBytes?: number; } /** * Bearer-check and the trivial health/ready/not-found responses now delegate to * `@danypops/vehicle-server/rpc-http` (the same handful of lines every daemon's service.ts hand-rolled). * The response-size guard below stays jittor-specific: daemon-kit's `jsonResponse` is intentionally * unbounded (it has no operation dispatch of its own to guard), while jittor's `/api/v1/ops` can * return arbitrarily large query results that must be capped (see SERVICE_MAX_RESPONSE_BYTES). */ function json(value: unknown, status = 200): Response { const body = JSON.stringify(value); if (new TextEncoder().encode(body).byteLength > SERVICE_MAX_RESPONSE_BYTES) return errorResponse("response too large", 413); return new Response(body, { status, headers: { "content-type": "application/json", "content-length": String(new TextEncoder().encode(body).byteLength) }, }); } export function createApp(options: JittorAppOptions): { fetch(request: Request): Promise } { const maxBodyBytes = options.maxBodyBytes ?? SERVICE_MAX_BODY_BYTES; // A real second transport for every jittor operation (see vehicle-registration.ts) -- // composed here rather than replacing /api/v1/ops, matching every other Vehicle-migrated // daemon in this ecosystem ("served alongside", not "instead of"). Routed before the // top-level bearer check below since createVehicleHttpApp performs its own. const permissions = [...new Set(options.service.vehicleRegistry.manifest().operations.flatMap((operation) => operation.permissions))]; const vehicleApp = createVehicleHttpApp({ registry: options.service.vehicleRegistry, token: options.token, invocationAuthority: { mode: "attested", resolve: () => ({ permissions, principal: { id: "jittor-authenticated-client" } }) }, }); return { async fetch(request: Request): Promise { const url = new URL(request.url); if (url.pathname.startsWith("/vehicle/")) return vehicleApp.fetch(request); if (!requireBearerToken(request, options.token)) return errorResponse("unauthorized", 401); if (request.method === "GET" && url.pathname === "/health") return healthResponse(VERSION); if (request.method === "GET" && url.pathname === "/ready") return readyResponse(options.service.ready()); if (request.method === "GET" && url.pathname === "/api/v1/ops") return json({ operations: options.service.operationNames() }); if (request.method !== "POST" || url.pathname !== "/api/v1/ops") return errorResponse("not found", 404); const contentLength = Number(request.headers.get("content-length") ?? 0); if (contentLength > maxBodyBytes) return errorResponse("payload too large", 413); const text = await request.text(); if (new TextEncoder().encode(text).byteLength > maxBodyBytes) return errorResponse("payload too large", 413); try { const body = JSON.parse(text) as { op?: unknown; input?: unknown }; if (typeof body.op !== "string") throw new Error("op is required"); if (isResetOperation(body.op) && (!body.input || typeof body.input !== "object" || Array.isArray(body.input))) throw new Error("Reset input must be an object"); const input = typeof body.input === "object" && body.input !== null && !Array.isArray(body.input) ? (body.input as Record) : {}; return json({ result: await options.service.execute(body.op, input) }); } catch (error) { if (error instanceof UnknownOperationError) return json({ error: error.message }, 404); if (error instanceof InvalidSessionSecretError) return json({ error: error.message }, 403); return json({ error: error instanceof Error ? error.message : String(error) }, 400); } }, }; }