diff --git a/.env.example b/.env.example index d51001a..c211dbe 100644 --- a/.env.example +++ b/.env.example @@ -291,6 +291,24 @@ ORCAROUTER_ENDPOINT=https://api.orcarouter.ai/v1/chat/completions # DESCRIPTION: Per-request timeout in ms. # CODEX_TIMEOUT=120000 +# ------------------------------------------------------------------------------ +# Cursor (uses your Cursor Pro subscription via local cursor-agent CLI) +# ------------------------------------------------------------------------------ +# DESCRIPTION: Enable the Cursor local provider (requires `cursor-agent` CLI installed + `agent login`). +# Values: true | false +# CURSOR_ENABLED=true +# DESCRIPTION: Cursor model id (run `cursor-agent models` to list). +# CURSOR_MODEL=composer-2.5 +# DESCRIPTION: Path to the `cursor-agent` binary; auto-detected if unset. +# CURSOR_BINARY_PATH=cursor-agent +# DESCRIPTION: Per-request timeout in ms. +# CURSOR_TIMEOUT=120000 +# DESCRIPTION: Auto-approve ALL cursor-agent permission requests, shell execution +# included (passes --force). The sandbox cwd is not a security boundary while this +# is on. Set false to approve MCPs only and reject edit/execute/delete requests. +# Values: true | false +# CURSOR_AUTO_APPROVE=true + # ------------------------------------------------------------------------------ # Embeddings provider override # ------------------------------------------------------------------------------ diff --git a/config/difficulty-anchors.json b/config/difficulty-anchors.json index 12bcad0..eed92ab 100644 --- a/config/difficulty-anchors.json +++ b/config/difficulty-anchors.json @@ -28,5 +28,17 @@ "design a horizontally scalable architecture for this service with failure-mode analysis", "Compare how databricks.js and openai-format.js each handle non-2xx upstream responses, then propose a unified error-normalization helper both could share", "Code review the PR #84 routing hardening changes" + ], + "frontier": [ + "Prove the correctness of this lock-free queue implementation and identify any ABA hazards", + "Design a novel conflict-free replicated data type for collaborative rich-text editing and argue its convergence from first principles", + "Derive the optimal cache eviction policy for this access distribution and prove its competitive ratio", + "Diagnose this heisenbug: a race that only reproduces under load across three services, given these interleaved logs", + "Design a Byzantine-fault-tolerant consensus protocol variant that tolerates f faults with 2f+1 replicas by weakening liveness, and analyze exactly which guarantees are lost", + "Formally verify that this state machine can never deadlock, or produce a counterexample trace", + "Given these profiler traces, determine whether the tail latency is queueing-theoretic or GC-driven, and prove which by constructing a discriminating experiment", + "Reason through the security of this key-rotation scheme under an adversary who can observe but not modify traffic, and find the weakest assumption it relies on", + "Work out the exact memory ordering constraints this concurrent hashmap needs on ARM, and justify each barrier from the C++ memory model", + "Given this failing distributed transaction trace, determine whether the anomaly is write skew or lost update, and design the minimal isolation-level change that eliminates it" ] } diff --git a/package.json b/package.json index 929fe46..9f5fdcf 100644 --- a/package.json +++ b/package.json @@ -38,7 +38,7 @@ "start:supervised": "while true; do node index.js 2>&1 | npx pino-pretty --sync; code=$?; echo \"[supervisor] lynkr exited ($code) — restarting in 3s\"; sleep 3; done", "lint": "eslint src index.js", "test": "npm run test:unit && npm run test:performance", - "test:unit": "LYNKR_KNN_DIR=/tmp/lynkr-test-knn DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com LOG_FILE_ENABLED=false node --test test/routing.test.js test/hybrid-routing-integration.test.js test/retry-logic.test.js test/sse-transformer.test.js test/passthrough-stream.test.js test/passthrough-mode.test.js test/openrouter-error-resilience.test.js test/format-conversion.test.js test/azure-openai-config.test.js test/azure-openai-format-conversion.test.js test/azure-openai-routing.test.js test/azure-openai-streaming.test.js test/azure-openai-error-resilience.test.js test/azure-openai-integration.test.js test/openai-integration.test.js test/atlas-integration.test.js test/toon-compression.test.js test/gcf-compression.test.js test/llamacpp-integration.test.js test/resilience.test.js test/telemetry-routing.test.js test/memory/store.test.js test/memory/surprise.test.js test/memory/extractor.test.js test/memory/search.test.js test/memory/retriever.test.js test/memory/distiller.test.js test/memory/distiller-freeze.test.js test/memory/wiki.test.js test/memory/skills-cache.test.js test/memory/tencentdb-launcher.test.js test/distill.test.js test/large-payload.test.js test/prompt-cache-injection.test.js test/risk-analyzer.test.js test/interaction-block.test.js test/preflight.test.js test/token-reduction.test.js test/session-affinity.test.js test/cache-state.test.js test/cache-switch-cost.test.js test/lens-recommendations.test.js test/model-registry-cost.test.js test/output-format-guard.test.js test/tier-fallback.test.js test/wrap.test.js test/init.test.js test/tool-call-response-metadata.test.js test/degradation.test.js test/routing-telemetry-columns.test.js test/sticky-routing.test.js test/knn-ambiguous-escalate.test.js test/deescalator.test.js test/client-profiles.test.js test/strip-internal-fields.test.js test/complexity-tool-subtraction.test.js test/bandit.test.js test/routing-propensity.test.js test/reward-pipeline.test.js test/knn-cold-start.test.js test/calibration.test.js test/feedback-loop.test.js test/session-fingerprint.test.js test/side-request-guards.test.js test/verifier.test.js test/intent-score.test.js test/difficulty-classifier.test.js test/classifier-setup.test.js test/usage-stats.test.js test/loop-guard.test.js test/moonshot-model-mapping.test.js test/baidu-model-mapping.test.js test/tenant-policy-ingress-parity.test.js test/decide.test.js test/embeddings-degradation.test.js test/health-probe.test.js test/stuck-detector.test.js test/onnx-embedder.test.js test/ope.test.js test/hierarchical-budget.test.js test/token-rate-limit.test.js test/otel-export.test.js test/mcp-broker.test.js test/compression-budget.test.js test/gpt-utils.test.js test/dedup-observe-only.test.js test/context-window-header.test.js test/token-budget-auto.test.js test/opencode-setup.test.js test/auth-mode-first-party.test.js test/task-ledger.test.js test/jev-router.test.js test/jev-routing.test.js test/force-patterns.test.js test/upstream-fidelity.test.js test/tool-schema-compression.test.js test/passthrough-route.test.js test/quota-ledger.test.js", + "test:unit": "LYNKR_KNN_DIR=/tmp/lynkr-test-knn DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com LOG_FILE_ENABLED=false node --test test/routing.test.js test/hybrid-routing-integration.test.js test/retry-logic.test.js test/sse-transformer.test.js test/passthrough-stream.test.js test/passthrough-mode.test.js test/openrouter-error-resilience.test.js test/format-conversion.test.js test/azure-openai-config.test.js test/azure-openai-format-conversion.test.js test/azure-openai-routing.test.js test/azure-openai-streaming.test.js test/azure-openai-error-resilience.test.js test/azure-openai-integration.test.js test/openai-integration.test.js test/atlas-integration.test.js test/toon-compression.test.js test/gcf-compression.test.js test/llamacpp-integration.test.js test/resilience.test.js test/telemetry-routing.test.js test/memory/store.test.js test/memory/surprise.test.js test/memory/extractor.test.js test/memory/search.test.js test/memory/retriever.test.js test/memory/distiller.test.js test/memory/distiller-freeze.test.js test/memory/wiki.test.js test/memory/skills-cache.test.js test/memory/tencentdb-launcher.test.js test/distill.test.js test/large-payload.test.js test/prompt-cache-injection.test.js test/risk-analyzer.test.js test/interaction-block.test.js test/preflight.test.js test/token-reduction.test.js test/session-affinity.test.js test/cache-state.test.js test/cache-switch-cost.test.js test/lens-recommendations.test.js test/model-registry-cost.test.js test/output-format-guard.test.js test/tier-fallback.test.js test/wrap.test.js test/init.test.js test/tool-call-response-metadata.test.js test/degradation.test.js test/routing-telemetry-columns.test.js test/sticky-routing.test.js test/knn-ambiguous-escalate.test.js test/deescalator.test.js test/client-profiles.test.js test/strip-internal-fields.test.js test/complexity-tool-subtraction.test.js test/bandit.test.js test/routing-propensity.test.js test/reward-pipeline.test.js test/knn-cold-start.test.js test/calibration.test.js test/feedback-loop.test.js test/session-fingerprint.test.js test/side-request-guards.test.js test/verifier.test.js test/intent-score.test.js test/difficulty-classifier.test.js test/classifier-setup.test.js test/usage-stats.test.js test/loop-guard.test.js test/moonshot-model-mapping.test.js test/baidu-model-mapping.test.js test/tenant-policy-ingress-parity.test.js test/decide.test.js test/embeddings-degradation.test.js test/health-probe.test.js test/stuck-detector.test.js test/onnx-embedder.test.js test/ope.test.js test/hierarchical-budget.test.js test/token-rate-limit.test.js test/otel-export.test.js test/mcp-broker.test.js test/compression-budget.test.js test/gpt-utils.test.js test/dedup-observe-only.test.js test/context-window-header.test.js test/token-budget-auto.test.js test/opencode-setup.test.js test/auth-mode-first-party.test.js test/harness-envelope.test.js test/task-ledger.test.js test/jev-router.test.js test/jev-routing.test.js test/force-patterns.test.js test/upstream-fidelity.test.js test/tool-schema-compression.test.js test/passthrough-route.test.js test/quota-ledger.test.js", "test:memory": "LYNKR_KNN_DIR=/tmp/lynkr-test-knn DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com node --test test/memory/store.test.js test/memory/surprise.test.js test/memory/extractor.test.js test/memory/search.test.js test/memory/retriever.test.js test/memory/distiller.test.js test/memory/distiller-freeze.test.js test/memory/wiki.test.js test/memory/skills-cache.test.js test/memory/tencentdb-launcher.test.js", "test:new-features": "LYNKR_KNN_DIR=/tmp/lynkr-test-knn DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com node --test test/passthrough-mode.test.js test/openrouter-error-resilience.test.js test/format-conversion.test.js", "test:performance": "LYNKR_KNN_DIR=/tmp/lynkr-test-knn DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com node test/hybrid-routing-performance.test.js && DATABRICKS_API_KEY=test-key DATABRICKS_API_BASE=http://test.com node test/performance-tests.js", diff --git a/src/api/router.js b/src/api/router.js index 65cfc92..f78b73e 100644 --- a/src/api/router.js +++ b/src/api/router.js @@ -10,6 +10,7 @@ const providersRouter = require("./providers-handler"); const claudeDesktopGatewayRouter = require("./claude-desktop-gateway"); const { getRoutingHeaders, getRoutingStats, analyzeComplexity, getModelTierSelector, analyzeRisk, checkSessionPin, writeSessionPin, checkPinScoreDrift } = require("../routing"); const { resolveTierForModelId } = require("../routing/model-slots"); +const { stripHarnessEnvelope } = require("../routing/harness-envelope"); // Upstream streams can die without a clean end (reader.read() never // resolves on a dropped socket), hanging the client forever. Every @@ -1409,7 +1410,11 @@ router.post("/v1/messages", rateLimiter, async (req, res, next) => { } return ''; })(); + // Suggestion-mode detection reads the RAW text (the marker is itself + // wrapper text); the force/risk probes below get envelope-stripped text + // so Cursor's / blocks can't fire triggers on "Hi". const isSuggestionMode = _lastUserText.includes('[SUGGESTION MODE:'); + const _lastUserAskClean = stripHarnessEnvelope(_lastUserText); // Tool-lessness alone is NOT harness evidence: generic API clients // (curl, benchmarks, SDKs) legitimately send bare messages and must get // full routing — live regression 2026-07-08: a benchmark's security- @@ -1535,7 +1540,7 @@ router.post("/v1/messages", rateLimiter, async (req, res, next) => { let _pinForceBypass = null; if (!sideTier && pinCheck.serve && pinCheck.reason === 'guards_passed' && !isSideRequest) { try { - const _probe = { messages: [{ role: 'user', content: _lastUserText || '' }] }; + const _probe = { messages: [{ role: 'user', content: _lastUserAskClean || '' }] }; const ca = require("../routing/complexity-analyzer"); if (_lastUserText && ca.shouldForceReasoning(_probe)) _pinForceBypass = 'force_reasoning'; else if (_lastUserText && ca.shouldForceCloud(_probe)) _pinForceBypass = 'force_cloud'; @@ -1883,6 +1888,15 @@ router.post("/v1/messages", rateLimiter, async (req, res, next) => { tierMethod: tier?.method || null, tierPinned: tier?.pinned ?? null, cacheState: _pinCacheState, + // Mid tool-exchange frames must keep serving the in-flight model + // (same invariant as the orchestrator's tool_history pin serves). + hasToolHistory: (() => { + try { return require("../routing/session-affinity").payloadHasToolHistory(req.body); } + catch { return false; } + })(), + // This branch IS the flat-fee subscription: the gate's dollar + // break-even leg has no premium to amortize here. + flatRate: true, // Quota pressure shortens the hold horizon inside the gate (high // sustained burn → descents clear sooner). Zero when the ledger has // no history — the gate behaves exactly as before. @@ -1919,6 +1933,36 @@ router.post("/v1/messages", rateLimiter, async (req, res, next) => { logger.debug({ err: err.message }, '[Routing] Pin-hold restore failed (non-fatal)'); } } + // Upgrades must persist the SERVED model into the pin, mirroring the + // pin_hold restore above — the write at the fresh-decision site stored + // the pre-rewrite tier model, and a pin that forgets the upgrade lets + // later frames descend below what actually served. Unlike pin_hold, + // an upgrade usually happens with no prior pin, so the payload comes + // from the fresh tier object. Floored (+continuation_inherit) and + // risk-lifted tiers stay unpinned (writeSessionPin re-checks, but the + // method is overridden to 'session_pin' here, so guard at the source). + if ( + _passthroughRoute.action === 'upgrade' + && _passthroughRoute.model + && pinCheck?.sessionId + && tier?.provider + && !String(tier?.method || '').includes('+continuation_inherit') + && tier?.escalation_source !== 'risk' + ) { + try { + const { writeSessionPin } = require("../routing/index"); + writeSessionPin(pinCheck.sessionId, { + provider: tier.provider, + model: _passthroughRoute.model, + tier: tier.tier, + score: tier.score ?? null, + method: 'session_pin', + reason: 'passthrough_upgrade', + }, req.body); + } catch (err) { + logger.debug({ err: err.message }, '[Routing] Upgrade pin persist failed (non-fatal)'); + } + } if (_passthroughRoute.action !== 'verbatim' && _passthroughRoute.model) { logger.info({ from: _passthroughClientModel, diff --git a/src/cache/embeddings.js b/src/cache/embeddings.js index a11afe1..bf2b731 100644 --- a/src/cache/embeddings.js +++ b/src/cache/embeddings.js @@ -145,6 +145,9 @@ const STRICT = process.env.LYNKR_EMBEDDINGS_STRICT === 'true'; function _noteFallback(providerName, err) { fallbackCount += 1; lastProviderError = err?.message || String(err); + // Cooldown counts from the LAST failed dial, not the first attempt's start — + // the in-place retry delay must not eat into the cooldown window. + lastProviderAttempt = Date.now(); if (embeddingProviderAvailable !== false) { embeddingProviderAvailable = false; degradedSince = Date.now(); @@ -168,6 +171,14 @@ function _noteRecovery(providerName) { degradedSince = null; } +// One in-place retry before tripping degradation: a single transient failure +// (Ollama cold-loading the model at boot, a blip mid-restart) was flipping +// the WHOLE provider to hash embeddings for a full cooldown after every +// server start (recurring live incident, last 2026-09-26). A short backoff +// absorbs the blip; a provider that is actually down fails twice and +// degrades exactly as before. +const TRANSIENT_RETRY_DELAY_MS = 1500; + function _wrapProvider(providerName, providerFn) { return async (text) => { // While degraded, only re-attempt the provider after the cooldown; serve @@ -184,9 +195,25 @@ function _wrapProvider(providerName, providerFn) { const result = await providerFn(text); _noteRecovery(providerName); return result; - } catch (err) { - _noteFallback(providerName, err); - if (STRICT) throw err; + } catch (firstErr) { + // Already-degraded providers get no second chance (this attempt WAS the + // post-cooldown probe); healthy ones earn one retry before the flip. + if (embeddingProviderAvailable !== false) { + await new Promise((r) => setTimeout(r, TRANSIENT_RETRY_DELAY_MS)); + try { + const result = await providerFn(text); + _noteRecovery(providerName); + logger.debug({ provider: providerName, error: firstErr?.message }, + '[Embeddings] Transient provider failure absorbed by retry'); + return result; + } catch (secondErr) { + _noteFallback(providerName, secondErr); + if (STRICT) throw secondErr; + return generateHashEmbedding(text); + } + } + _noteFallback(providerName, firstErr); + if (STRICT) throw firstErr; return generateHashEmbedding(text); } }; diff --git a/src/clients/cursor-utils.js b/src/clients/cursor-utils.js new file mode 100644 index 0000000..6871a58 --- /dev/null +++ b/src/clients/cursor-utils.js @@ -0,0 +1,791 @@ +/** + * Cursor CLI Format Conversion + Invocation Utilities. + * + * Mirrors `src/clients/codex-utils.js` / `invokeCodex` in + * `src/clients/databricks.js`, but for the official `cursor-agent` CLI. + * + * Invocation model (reworked 2026-09-26 after the ETIMEDOUT incident): + * - async `execFile` — a slow CLI call must never block the proxy's event + * loop (the old `execFileSync` froze EVERY session for up to 120s per + * call, and a 3-candidate tier fallback froze it for 6 minutes). + * - prompt travels over STDIN, not argv — argv has an OS size ceiling and + * agentic payloads exceed it; stdin has none. + * - session resume: the CLI returns a `session_id`; subsequent turns of + * the same Lynkr session send `--resume ` with ONLY the newest + * user turn instead of re-flattening the whole conversation. This is + * what makes heavy agentic traffic viable (bounded payload per turn, + * CLI-side prompt cache reuse). A failed resume falls back to one + * fresh full-prompt attempt and re-seeds the session. + * - timeout scales with payload size (base + per-KB), hard-capped. + * - at most MAX_CONCURRENT_PROCS CLI processes run at once; extra calls + * queue (each spawn is a full Node worker). + * - a one-time background warmup absorbs the ~20s worker cold-start so + * the first real request doesn't pay it. + * + * Auth is inherited from the user's own login (`agent login` session or + * `CURSOR_API_KEY` env) — Lynkr never reads Cursor's token store, it just + * spawns the official binary with an inherited environment. Auth failures + * are detected from CLI stderr and reported as such; a timeout is reported + * as a timeout (the old code blamed every failure on the subscription). + * + * @module clients/cursor-utils + */ + +const { execFile, execSync } = require("node:child_process"); +const crypto = require("node:crypto"); +const fs = require("node:fs"); +const os = require("node:os"); +const path = require("node:path"); +const logger = require("../logger"); + +const DEFAULT_MODEL = "composer-2.5"; +const DEFAULT_BINARY = "cursor-agent"; +// Timeout: base covers worker spin-up + small prompts; large payloads earn +// more, capped hard. Constants, not env knobs (house convention). +const BASE_TIMEOUT_MS = 120_000; +const PER_KB_TIMEOUT_MS = 250; +const MAX_TIMEOUT_MS = 600_000; +const DEFAULT_TIMEOUT_MS = BASE_TIMEOUT_MS; // kept for existing callers/tests +const MAX_BUFFER_BYTES = 32 * 1024 * 1024; +const MAX_CONCURRENT_PROCS = 2; +const SESSION_CACHE_MAX = 200; +const WARMUP_TIMEOUT_MS = 45_000; +// Operator decision 2026-09-27 ("No blocking shell"): auto-approve ALL agent +// permission requests — shell execution included — on both the ACP and +// one-shot paths. NOTE: with shell allowed, the sandbox cwd is a default +// directory, not a security boundary. CURSOR_AUTO_APPROVE=false restores the +// deny-mutations policy (MCPs-only approval, no --force). +const CURSOR_AUTO_APPROVE = process.env.CURSOR_AUTO_APPROVE?.trim().toLowerCase() !== "false"; + +/** + * Resolve the binary to spawn. Test-overridable via env only — + * no config import here (keeps this module require-cycle free; + * databricks.js passes model/timeout in). + */ +function getBinaryPath(configCursor) { + return configCursor?.binaryPath?.trim() || process.env.CURSOR_BINARY_PATH?.trim() || DEFAULT_BINARY; +} + +/** + * True when the `cursor-agent` binary exists on PATH. On the real path + * (no injected which), a successful check also kicks the one-time + * background warmup so the worker cold-start is paid before traffic. + * @param {Function} [whichFn] - injectable for tests. + */ +function isAvailable(whichFn) { + const run = whichFn || ((bin, opts) => require("node:child_process").execFileSync("which", [bin], opts)); + try { + const binary = process.env.CURSOR_BINARY_PATH?.trim() || DEFAULT_BINARY; + run(binary, { stdio: "ignore" }); + if (!whichFn) setImmediate(() => warmupCursorAgent().catch(() => {})); + return true; + } catch { + return false; + } +} + +function extractText(message) { + if (!message) return ""; + const content = message.content; + if (typeof content === "string") return content; + if (!Array.isArray(content)) return ""; + return content + .map((block) => { + if (block.type === "text") return block.text || ""; + if (block.type === "tool_result") { + const result = typeof block.content === "string" ? block.content : JSON.stringify(block.content); + return `[Tool Result: ${result}]`; + } + if (block.type === "tool_use") { + return `[Tool Call: ${block.name}(${JSON.stringify(block.input)})]`; + } + return ""; + }) + .filter(Boolean) + .join("\n"); +} + +/** + * Flatten Anthropic system + history into a single CLI prompt. + * Same strategy as convertAnthropicToCodexPrompt (last user message is + * the prompt, prior turns become context). Used for FRESH sessions; + * resumed sessions send latestUserTurnPrompt instead. + */ +function convertAnthropicToCursorPrompt(body) { + const systemContext = body.system || null; + const messages = body.messages || []; + if (messages.length === 0) return { prompt: "", systemContext }; + if (messages.length === 1 && messages[0].role === "user") { + return { prompt: extractText(messages[0]), systemContext }; + } + let lastUserIndex = -1; + for (let i = messages.length - 1; i >= 0; i--) { + if (messages[i]?.role === "user") { + lastUserIndex = i; + break; + } + } + if (lastUserIndex === -1) { + return { prompt: extractText(messages[messages.length - 1]), systemContext }; + } + const lastUserMessage = extractText(messages[lastUserIndex]); + const priorMessages = messages.slice(0, lastUserIndex); + if (priorMessages.length === 0) return { prompt: lastUserMessage, systemContext }; + const contextParts = priorMessages + .map((m) => { + const text = extractText(m); + if (!text) return null; + return `${m.role === "user" ? "User" : "Assistant"}: ${text}`; + }) + .filter(Boolean); + const conversationContext = contextParts.join("\n\n"); + return { + prompt: conversationContext ? `Previous conversation:\n${conversationContext}\n\nUser: ${lastUserMessage}` : lastUserMessage, + systemContext, + }; +} + +/** + * The newest user turn only — what a RESUMED CLI session receives (its own + * transcript already holds the earlier turns). + */ +function latestUserTurnPrompt(body) { + const messages = body?.messages || []; + for (let i = messages.length - 1; i >= 0; i--) { + if (messages[i]?.role === "user") return extractText(messages[i]); + } + return ""; +} + +/** + * Stable per-conversation key for the resume cache: the caller's session id + * when present, else a hash of the first message (stable across turns — + * flattened prompts are NOT, their prefix changes shape after turn one). + */ +function deriveSessionKey(body) { + if (body?._sessionId && typeof body._sessionId === "string") return `sid:${body._sessionId}`; + const first = body?.messages?.[0]; + const text = first ? extractText(first) : ""; + if (!text) return null; + return `msg1:${crypto.createHash("sha1").update(text).digest("hex").slice(0, 32)}`; +} + +/** + * Pull plain text out of `cursor-agent -p --output-format json` stdout. + * Handles: {text|result|output|content|string}, stream-json event lines, + * and falls back to the raw trimmed stdout. + */ +function extractCursorText(stdout) { + const trimmed = String(stdout || "").trim(); + if (!trimmed) return ""; + try { + const parsed = JSON.parse(trimmed); + if (typeof parsed === "string") return parsed; + for (const key of ["text", "result", "output", "content", "response", "message"]) { + if (typeof parsed?.[key] === "string" && parsed[key].trim()) return parsed[key]; + } + if (Array.isArray(parsed?.content)) { + const t = parsed.content + .map((b) => (typeof b === "string" ? b : b?.text || "")) + .filter(Boolean) + .join("\n"); + if (t.trim()) return t; + } + return trimmed; + } catch { + // Possibly stream-json (one JSON object per line) — collect assistant deltas. + const lines = trimmed.split("\n"); + if (lines.length > 1) { + const parts = []; + for (const line of lines) { + const l = line.trim(); + if (!l) continue; + try { + const ev = JSON.parse(l); + const t = + ev?.message?.content?.[0]?.text || + ev?.content?.[0]?.text || + (typeof ev?.text === "string" ? ev.text : "") || + (typeof ev?.delta === "string" ? ev.delta : ""); + if (t) parts.push(t); + } catch { + parts.push(l); + } + } + if (parts.length) return parts.join(""); + } + return trimmed; + } +} + +/** + * Structured view of one CLI run: text, the session id (resume handle), + * and token usage when the CLI reports it. + */ +function parseCursorResult(stdout) { + const text = extractCursorText(stdout); + let obj = null; + const trimmed = String(stdout || "").trim(); + try { + obj = JSON.parse(trimmed); + } catch { + const lines = trimmed.split("\n"); + for (let i = lines.length - 1; i >= 0; i--) { + try { + const candidate = JSON.parse(lines[i].trim()); + if (candidate && typeof candidate === "object") { + obj = candidate; + break; + } + } catch { /* keep walking back */ } + } + } + const usage = obj?.usage && typeof obj.usage === "object" + ? { + inputTokens: Number(obj.usage.inputTokens) || 0, + outputTokens: Number(obj.usage.outputTokens) || 0, + cacheReadTokens: Number(obj.usage.cacheReadTokens) || 0, + cacheWriteTokens: Number(obj.usage.cacheWriteTokens) || 0, + } + : null; + return { text, sessionId: typeof obj?.session_id === "string" ? obj.session_id : null, usage }; +} + +function convertCursorResponseToAnthropic(text, model, usage = null, thinking = "") { + const estimatedOutputTokens = Math.ceil(String(text || "").length / 4); + const content = []; + if (thinking) content.push({ type: "thinking", thinking }); + content.push({ type: "text", text: text || "" }); + return { + id: `msg_cursor_${Date.now()}`, + type: "message", + role: "assistant", + model: model || "cursor", + content, + stop_reason: "end_turn", + stop_sequence: null, + usage: { + input_tokens: usage?.inputTokens || 0, + output_tokens: usage?.outputTokens ?? estimatedOutputTokens, + cache_creation_input_tokens: usage?.cacheWriteTokens || 0, + cache_read_input_tokens: usage?.cacheReadTokens || 0, + }, + }; +} + +/** Payload-scaled timeout: base + per-KB allowance, hard-capped. */ +function scaleTimeoutMs(payloadBytes, baseMs) { + const base = Number(baseMs) > 0 ? Number(baseMs) : BASE_TIMEOUT_MS; + const scaled = base + Math.ceil((Number(payloadBytes) || 0) / 1024) * PER_KB_TIMEOUT_MS; + return Math.min(MAX_TIMEOUT_MS, Math.max(base, scaled)); +} + +// --- concurrency gate (each spawn is a full Node worker) --------------------- +let _inFlight = 0; +const _waiters = []; +async function _acquire() { + if (_inFlight < MAX_CONCURRENT_PROCS) { + _inFlight++; + return; + } + await new Promise((resolve) => _waiters.push(resolve)); + _inFlight++; +} +function _release() { + _inFlight--; + const next = _waiters.shift(); + if (next) next(); +} + +// --- resume cache: Lynkr session key → CLI session_id ------------------------ +const _sessionCache = new Map(); +function _sessionCacheSet(key, chatId) { + if (!key || !chatId) return; + if (_sessionCache.has(key)) _sessionCache.delete(key); + _sessionCache.set(key, chatId); + while (_sessionCache.size > SESSION_CACHE_MAX) { + _sessionCache.delete(_sessionCache.keys().next().value); + } +} + +/** + * Honest error classification: timeouts are timeouts, auth is auth, missing + * binary is missing binary. The old code stamped every failure with the + * subscription hint, which misdiagnosed a plain timeout in production. + */ +function classifyCursorError(err, { payloadBytes = 0, timeoutMs = 0 } = {}) { + const raw = String(err?.message || err || ""); + const stderr = String(err?.stderr || ""); + let hint; + if (/ENOENT|not found|no such file/i.test(raw)) { + hint = "is `cursor-agent` installed? run `cursor-agent --version`; set CURSOR_BINARY_PATH if needed"; + } else if (/not (currently )?logged in|sign ?in|unauthoriz|401|no active subscription|login required/i.test(`${raw}\n${stderr}`)) { + hint = "cursor-agent is not authenticated — run `cursor-agent login` and verify with `cursor-agent status`"; + } else if (err?.killed || err?.signal === "SIGKILL" || /ETIMEDOUT|timed? ?out/i.test(raw)) { + hint = `timed out after ${Math.round(timeoutMs / 1000)}s with a ${Math.ceil(payloadBytes / 1024)}KB payload (worker cold-start or oversized turn)`; + } + const e = new Error(hint ? `${raw} (${hint})` : raw); + e.cursorHint = hint || null; + e.stderr = stderr; + return e; +} + +/** + * The CLI is an AGENT with file tools, and a spawned child inherits the + * server's cwd — which is the Lynkr repo itself. Without an explicit + * workspace, every "review this project" routed here would read (and could + * write, on approval) the server's own checkout, .env included (live + * incident 2026-09-26). No workspace ⇒ empty sandbox dir, always. + */ +let _sandboxDir = null; +function getSandboxDir() { + if (_sandboxDir) return _sandboxDir; + const dir = path.join(os.tmpdir(), "lynkr-cursor-sandbox"); + try { fs.mkdirSync(dir, { recursive: true }); } catch { /* tmpdir exists */ } + _sandboxDir = dir; + return dir; +} + +/** Default async exec: stdin transport, hard timeout, stderr captured. */ +function _defaultExec({ binaryPath, args, stdin, timeoutMs, maxBuffer, cwd }) { + return new Promise((resolve, reject) => { + const child = execFile( + binaryPath, + args, + { timeout: timeoutMs, maxBuffer, encoding: "utf8", env: { ...process.env }, cwd: cwd || getSandboxDir(), killSignal: "SIGKILL" }, + (err, stdout, stderr) => { + if (err) { + err.stderr = String(stderr || ""); + reject(err); + return; + } + resolve(String(stdout || "")); + } + ); + if (child.stdin) { + child.stdin.on("error", () => {}); + child.stdin.end(stdin ?? ""); + } + }); +} + +// --- ACP persistent channel (spawn-tax elimination, 2026-09-27) -------------- +// `cursor-agent acp` is an official JSON-RPC-over-stdio server mode +// (cursor.com/docs/cli/acp): one long-lived process, one ndjson line per +// request, streaming updates, session load/resume, per-session set_model. +// Spike-measured: warm turns ~1.6s vs 6-9s per one-shot spawn. The one-shot +// `-p` path below is retained verbatim as the automatic fallback (and is +// still what injected execFn tests exercise). +const ACP_ENABLED = true; // kill switch: flip + redeploy (no env knobs) +const ACP_INIT_TIMEOUT_MS = 30_000; +const ACP_REQUEST_ID_BASE = 1; + +class AcpClient { + constructor({ binaryPath, cwd }) { + this.binaryPath = binaryPath; + this.cwd = cwd; + this.child = null; + this.buf = ""; + this.nextId = ACP_REQUEST_ID_BASE; + this.pending = new Map(); + this.collectors = new Map(); // sessionId → {chunks:[]} + this.knownSessions = new Set(); + this.toolCalls = new Map(); // toolCallId → {kind,title,command,path,start,announced,done} + this.models = []; + this.currentModel = new Map(); // sessionId → acpModelId + this.dead = false; + this.initPromise = null; + this.promptQueue = Promise.resolve(); // serialize prompts per process + } + + _spawn() { + const { spawn } = require("node:child_process"); + this.child = spawn(this.binaryPath, ["acp"], { cwd: this.cwd, env: { ...process.env } }); + this.child.on("exit", (code) => { + this.dead = true; + for (const [, p] of this.pending) p.rej(new Error(`ACP process exited (${code})`)); + this.pending.clear(); + logger.warn({ code, cwd: this.cwd }, "[Cursor/ACP] server process exited"); + }); + this.child.stdin.on("error", () => {}); + this.child.stderr.on("data", () => {}); + this.child.stdout.on("data", (d) => this._onData(String(d))); + } + + _onData(s) { + this.buf += s; + let nl; + while ((nl = this.buf.indexOf("\n")) >= 0) { + const line = this.buf.slice(0, nl); this.buf = this.buf.slice(nl + 1); + if (!line.trim()) continue; + let msg; try { msg = JSON.parse(line); } catch { continue; } + if (msg.id !== undefined && this.pending.has(msg.id)) { + const p = this.pending.get(msg.id); this.pending.delete(msg.id); + if (msg.error) p.rej(new Error(`${p.method}: ${JSON.stringify(msg.error).slice(0, 200)}`)); + else p.res(msg.result); + } else if (msg.method === "session/update") { + const sid = msg.params?.sessionId; + const u = msg.params?.update || {}; + const c = this.collectors.get(sid); + if (!c) continue; + try { + if (u.sessionUpdate === "agent_message_chunk" && u.content?.text) { + c.chunks.push(u.content.text); + c.onChunk?.(u.content.text, "text"); + } else if (u.sessionUpdate === "agent_thought_chunk" && u.content?.text) { + c.thoughts.push(u.content.text); + c.onChunk?.(u.content.text, "thought"); + } else if (u.sessionUpdate === "tool_call" || u.sessionUpdate === "tool_call_update") { + // Correlated narration (2026-09-27): the initial tool_call often + // carries only a generic title; the real target (locations/ + // rawInput) arrives in enrichment updates, and status updates are + // title-less. Merge everything per toolCallId and emit exactly + // two meaningful lines per tool — "▶ verb: target" once running, + // "✓/✗ verb: target (Ns)" on terminal. All else is swallowed. + const id = u.toolCallId || "?"; + let entry = this.toolCalls.get(id); + if (!entry) { + entry = { kind: null, title: null, command: null, path: null, start: Date.now(), announced: false, done: false }; + this.toolCalls.set(id, entry); + while (this.toolCalls.size > 200) this.toolCalls.delete(this.toolCalls.keys().next().value); + } + if (u.kind) entry.kind = u.kind; + if (u.title) entry.title = u.title; + if (u.rawInput?.command) entry.command = u.rawInput.command; + const loc = u.locations?.[0]?.path || u.rawInput?.path; + if (loc) entry.path = loc; + + const emit = (line) => { + c.thoughts.push(`\n[${line}]`); + c.onChunk?.(line, "tool"); + }; + const label = () => { + const verb = { execute: "run", read: "read", edit: "edit", delete: "delete", move: "move", search: "search", fetch: "fetch", think: "think" }[entry.kind] || "tool"; + let target = entry.command || entry.path || String(entry.title || "").replace(/`/g, "").trim(); + if (target.startsWith(this.cwd)) target = target.slice(this.cwd.length + 1) || target; + else if (target.startsWith("/") && target.split("/").length > 3) target = "…/" + target.split("/").slice(-2).join("/"); + const capped = target.slice(0, 80); + return `${verb}: ${capped}` + (target.length > 80 ? "…" : ""); + }; + + if (u.status === "in_progress" && !entry.announced && !entry.done) { + entry.announced = true; + emit(`▶ ${label()}`); + } else if ((u.status === "completed" || u.status === "failed" || u.status === "cancelled") && !entry.done) { + entry.done = true; + const dur = ((Date.now() - entry.start) / 1000).toFixed(1); + if (u.status === "completed") { + emit(`✓ ${label()} (${dur}s)`); + } else { + const err = String(u.rawOutput?.stderr || "").trim().split("\n")[0].slice(0, 80); + const code = u.rawOutput?.exitCode != null ? ` exit ${u.rawOutput.exitCode}` : ""; + emit(`✗ ${label()} (${dur}s)${code}${err ? " — " + err : ""}`); + } + this.toolCalls.delete(id); + } + } + } catch { /* consumer errors must not kill the reader */ } + } else if (msg.id !== undefined && msg.method === "session/request_permission") { + const kind = msg.params?.toolCall?.kind || ""; + const options = msg.params?.options || []; + let pick; + if (CURSOR_AUTO_APPROVE) { + // Allow everything (operator-approved yolo): prefer allow_always to + // cut repeat prompts, fall back to any allow option. + pick = options.find((o) => /allow_always/i.test(o.kind || o.optionId || "")) + || options.find((o) => /allow/i.test(o.kind || o.optionId || "")); + } else { + const mutating = /edit|execute|delete|move|write/i.test(kind); + pick = mutating + ? options.find((o) => /reject/i.test(o.kind || o.optionId || "")) + : options.find((o) => /allow/i.test(o.kind || o.optionId || "")); + } + const chosen = pick || options[0]; + this._write({ jsonrpc: "2.0", id: msg.id, result: { outcome: { outcome: "selected", optionId: chosen?.optionId } } }); + } else if (msg.id !== undefined && msg.method) { + // Unsupported agent→client request (fs/* etc — we advertise no fs). + this._write({ jsonrpc: "2.0", id: msg.id, error: { code: -32601, message: "unsupported by lynkr acp client" } }); + } + } + } + + _write(obj) { + try { this.child.stdin.write(JSON.stringify(obj) + "\n"); } catch { /* exit handler rejects pendings */ } + } + + _request(method, params, timeoutMs) { + if (this.dead) return Promise.reject(new Error("ACP process dead")); + return new Promise((res, rej) => { + const id = this.nextId++; + this.pending.set(id, { res, rej, method }); + this._write({ jsonrpc: "2.0", id, method, params }); + const cap = timeoutMs || ACP_INIT_TIMEOUT_MS; + setTimeout(() => { + if (this.pending.has(id)) { this.pending.delete(id); rej(new Error(`${method} timed out after ${cap}ms`)); } + }, cap).unref?.(); + }); + } + + async init() { + if (this.initPromise) return this.initPromise; + this.initPromise = (async () => { + this._spawn(); + await this._request("initialize", { + protocolVersion: 1, + clientCapabilities: { fs: { readTextFile: false, writeTextFile: false } }, + }); + })(); + return this.initPromise; + } + + /** Map a Lynkr/CLI-style id (composer-2.5-fast, grok-4.7-medium-fast) to + * the ACP catalog's bracketed id via longest base-name prefix match. */ + resolveModelId(tierModel) { + if (!tierModel || tierModel === "auto") return null; + let best = null; + for (const m of this.models) { + const base = String(m.modelId).split("[")[0]; + if (tierModel === base || tierModel.startsWith(base)) { + if (!best || base.length > best.base.length) best = { base, id: m.modelId }; + } + } + return best ? best.id : null; + } + + async newSession() { + const r = await this._request("session/new", { cwd: this.cwd, mcpServers: [] }); + if (Array.isArray(r?.models?.availableModels)) this.models = r.models.availableModels; + this.knownSessions.add(r.sessionId); + return r.sessionId; + } + + async ensureSession(sessionId) { + if (this.knownSessions.has(sessionId)) return true; + await this._request("session/load", { sessionId, cwd: this.cwd, mcpServers: [] }, ACP_INIT_TIMEOUT_MS); + this.knownSessions.add(sessionId); + return true; + } + + async setModel(sessionId, tierModel) { + const acpId = this.resolveModelId(tierModel); + if (!acpId || this.currentModel.get(sessionId) === acpId) return; + try { + await this._request("session/set_model", { sessionId, modelId: acpId }); + this.currentModel.set(sessionId, acpId); + } catch (err) { + logger.debug({ tierModel, acpId, err: err.message }, "[Cursor/ACP] set_model failed — session default serves"); + } + } + + prompt(sessionId, text, timeoutMs, onChunk) { + // One prompt at a time per process — concurrency across sessions in a + // single ACP server is unverified; the queue keeps ordering sane and + // warm turns are ~1.6s so the wait is small. + const run = this.promptQueue.then(async () => { + const collector = { chunks: [], thoughts: [], onChunk }; + this.collectors.set(sessionId, collector); + try { + const r = await this._request("session/prompt", { + sessionId, + prompt: [{ type: "text", text }], + }, timeoutMs); + return { text: collector.chunks.join(""), thinking: collector.thoughts.join(""), stopReason: r?.stopReason || "end_turn" }; + } finally { + this.collectors.delete(sessionId); + } + }); + this.promptQueue = run.catch(() => {}); + return run; + } +} + +// One warm server per cwd (model is per-session, so one process serves every +// tier). Sandbox cwd is the default — same invariant as the one-shot path. +const _acpClients = new Map(); +async function _getAcpClient(binaryPath, cwd) { + const key = `${binaryPath}\u0000${cwd}`; + let client = _acpClients.get(key); + if (client && !client.dead) return client; + client = new AcpClient({ binaryPath, cwd }); + _acpClients.set(key, client); + await client.init(); + return client; +} + +/** + * ACP-path serve. Same return contract as the one-shot path; throws on any + * ACP-layer failure so runCursorAgent can fall back to the spawn path. + */ +async function _runViaAcp({ prompt, resumePrompt, sessionKey, model, baseTimeoutMs, binaryPath, workspace, onDelta }) { + const cwd = workspace || getSandboxDir(); + const client = await _getAcpClient(binaryPath || DEFAULT_BINARY, cwd); + const cached = sessionKey ? _sessionCache.get(sessionKey) : null; + const cachedAcp = cached && typeof cached === "object" && cached.acp ? cached.acp : null; + + let sessionId = null; + let resumed = false; + let stdinText = prompt || ""; + if (cachedAcp && resumePrompt) { + try { + await client.ensureSession(cachedAcp); + sessionId = cachedAcp; + stdinText = resumePrompt; + resumed = true; + } catch { /* stale/unloadable — fall through to a fresh session */ } + } + if (!sessionId) { + if (sessionKey) _sessionCache.delete(sessionKey); + sessionId = await client.newSession(); + } + await client.setModel(sessionId, model); + const timeoutMs = scaleTimeoutMs(Buffer.byteLength(stdinText, "utf8"), baseTimeoutMs || BASE_TIMEOUT_MS); + const t0 = Date.now(); + const { text, thinking, stopReason } = await client.prompt(sessionId, stdinText, timeoutMs, onDelta); + if (sessionKey) _sessionCacheSet(sessionKey, { acp: sessionId }); + logger.debug({ resumed, model, ms: Date.now() - t0, stopReason }, "[Cursor/ACP] prompt served"); + return { + stdout: JSON.stringify({ type: "result", subtype: "acp", result: text, session_id: sessionId }), + text, + thinking: thinking || "", + sessionId, + usage: null, // ACP updates carry no token usage — converter estimates + resumed, + }; +} + +let _warmupDone = false; +/** + * One-time background warmup: absorbs the CLI worker's ~20s cold start so + * the first real request doesn't spend its timeout budget on it. + */ +async function warmupCursorAgent(binaryPath) { + if (_warmupDone) return; + _warmupDone = true; + const bin = binaryPath || getBinaryPath(null); + try { + if (ACP_ENABLED) { + // Boot the persistent ACP server (pays the worker cold-start once). + await _getAcpClient(bin, getSandboxDir()); + logger.debug("[Cursor] ACP server warm"); + return; + } + await _defaultExec({ + binaryPath: bin, + args: ["-p", "--output-format", "json", "--model", DEFAULT_MODEL, "--trust", "--approve-mcps"], + stdin: "Reply with OK.", + timeoutMs: WARMUP_TIMEOUT_MS, + cwd: getSandboxDir(), + maxBuffer: MAX_BUFFER_BYTES, + }); + logger.debug("[Cursor] warmup complete — worker hot"); + } catch (err) { + logger.warn({ err: err.message }, "[Cursor] warmup failed (non-fatal) — first request will pay the cold start"); + } +} + +/** + * Spawn `cursor-agent -p` (async) and return the parsed run. + * + * @param {Object} args + * @param {string} args.prompt - full flattened prompt (fresh sessions) + * @param {string} [args.resumePrompt] - newest user turn (resumed sessions) + * @param {string|null} [args.sessionKey] - stable conversation key (deriveSessionKey) + * @param {string} args.model - Cursor model id + * @param {number} [args.baseTimeoutMs] - base before payload scaling (legacy alias: timeoutMs) + * @param {string} args.binaryPath + * @param {string|null} args.workspace - passed as --workspace when set + * @param {Function} [args.execFn] - injectable for tests: async ({binaryPath,args,stdin,timeoutMs,maxBuffer}) => stdout + * @returns {Promise<{stdout:string, text:string, sessionId:string|null, usage:Object|null, resumed:boolean}>} + */ +async function runCursorAgent({ prompt, resumePrompt, sessionKey, model, baseTimeoutMs, timeoutMs, binaryPath, workspace, execFn, onDelta }) { + const exec = execFn || _defaultExec; + const base = baseTimeoutMs || timeoutMs || BASE_TIMEOUT_MS; + // Persistent ACP channel first (real path only — injected execFn keeps the + // one-shot contract for tests). Any ACP failure falls through to the + // one-shot spawn below. + if (ACP_ENABLED && !execFn) { + try { + return await _runViaAcp({ prompt, resumePrompt, sessionKey, model, baseTimeoutMs: base, binaryPath, workspace, onDelta }); + } catch (err) { + logger.warn({ err: err.message }, "[Cursor/ACP] persistent channel failed — falling back to one-shot spawn"); + } + } + const cachedRaw = sessionKey ? _sessionCache.get(sessionKey) || null : null; + // Legacy --resume takes a chat id STRING; ACP-era cache entries are + // objects and must not leak into the spawn path. + const cachedChat = typeof cachedRaw === "string" ? cachedRaw : null; + const useResume = Boolean(cachedChat && resumePrompt); + + const buildArgs = (resumeChatId) => { + const a = ["-p", "--output-format", "json"]; + if (model) a.push("--model", model); + a.push("--trust"); + // Required for headless operation (no one can click approve): MCP servers + // must be pre-approved via --approve-mcps or the run stalls on an + // approval prompt until timeout. Shell/write approvals follow + // CURSOR_AUTO_APPROVE — the operator made that call explicitly + // (2026-09-27); it is not a silent default. + a.push("--approve-mcps"); + if (CURSOR_AUTO_APPROVE) a.push("--force"); + if (workspace) a.push("--workspace", workspace); + if (resumeChatId) a.push("--resume", resumeChatId); + return a; + }; + + await _acquire(); + try { + let stdout; + let resumed = useResume; + const firstStdin = useResume ? resumePrompt : prompt || ""; + const firstTimeout = scaleTimeoutMs(Buffer.byteLength(firstStdin, "utf8"), base); + logger.debug( + { binary: binaryPath, model, payloadBytes: Buffer.byteLength(firstStdin, "utf8"), resumed: useResume, timeoutMs: firstTimeout }, + "[Cursor] Spawning cursor-agent" + ); + try { + stdout = await exec({ binaryPath, args: buildArgs(useResume ? cachedChat : null), stdin: firstStdin, timeoutMs: firstTimeout, maxBuffer: MAX_BUFFER_BYTES, cwd: workspace || getSandboxDir() }); + } catch (err) { + if (!useResume) throw classifyCursorError(err, { payloadBytes: Buffer.byteLength(firstStdin, "utf8"), timeoutMs: firstTimeout }); + // Stale/failed resume → drop the handle, retry once fresh with the + // full flattened prompt so the conversation re-seeds. + logger.warn({ sessionKey, err: err.message }, "[Cursor] resume failed — retrying fresh"); + _sessionCache.delete(sessionKey); + resumed = false; + const freshTimeout = scaleTimeoutMs(Buffer.byteLength(prompt || "", "utf8"), base); + try { + stdout = await exec({ binaryPath, args: buildArgs(null), stdin: prompt || "", timeoutMs: freshTimeout, maxBuffer: MAX_BUFFER_BYTES, cwd: workspace || getSandboxDir() }); + } catch (err2) { + throw classifyCursorError(err2, { payloadBytes: Buffer.byteLength(prompt || "", "utf8"), timeoutMs: freshTimeout }); + } + } + const parsed = parseCursorResult(stdout); + if (sessionKey && parsed.sessionId) _sessionCacheSet(sessionKey, parsed.sessionId); + return { stdout, text: parsed.text, sessionId: parsed.sessionId, usage: parsed.usage, resumed }; + } finally { + _release(); + } +} + +module.exports = { + DEFAULT_MODEL, + DEFAULT_TIMEOUT_MS, + BASE_TIMEOUT_MS, + MAX_TIMEOUT_MS, + MAX_CONCURRENT_PROCS, + DEFAULT_BINARY, + getBinaryPath, + isAvailable, + extractText, + convertAnthropicToCursorPrompt, + latestUserTurnPrompt, + deriveSessionKey, + extractCursorText, + parseCursorResult, + convertCursorResponseToAnthropic, + scaleTimeoutMs, + classifyCursorError, + getSandboxDir, + CURSOR_AUTO_APPROVE, + warmupCursorAgent, + runCursorAgent, +}; diff --git a/src/clients/databricks.js b/src/clients/databricks.js index 5ff8263..96dcde5 100644 --- a/src/clients/databricks.js +++ b/src/clients/databricks.js @@ -3382,6 +3382,121 @@ async function invokeCodex(body, _incomingHeaders = {}) { }; } +async function invokeCursor(body, _incomingHeaders = {}) { + const cursorUtils = require("./cursor-utils"); + + if (config.cursor?.enabled === false) { + throw new Error("Cursor provider is disabled (set CURSOR_ENABLED=true to use cursor: tiers)"); + } + + const model = body._tierModel || config.cursor?.model || cursorUtils.DEFAULT_MODEL; + const { prompt, systemContext } = cursorUtils.convertAnthropicToCursorPrompt(body); + + if (!prompt) { + throw new Error("Cursor: no prompt content to send"); + } + + const fullPrompt = systemContext ? `System context:\n${systemContext}\n\nUser request:\n${prompt}` : prompt; + const binaryPath = cursorUtils.getBinaryPath(config.cursor); + const baseTimeoutMs = config.cursor?.timeout || cursorUtils.DEFAULT_TIMEOUT_MS; + const workspace = body._workspace || null; + + // Streaming (Phase 3, 2026-09-27): synthesize an OpenAI SSE stream from + // ACP session/update deltas. "cursor" is registered in + // DEFAULT_OPENAI_SSE_PROVIDERS, so the orchestrator either transforms this + // to Anthropic SSE (badge injection included) or passes it through to + // OpenAI-surface clients. The legacy spawn fallback emits no deltas — its + // full text arrives as one chunk at the end, which is still a valid + // stream. Errors after SSE has started cannot re-enter tier fallback (the + // 200 is already committed) — same semantics as every streaming provider. + if (body.stream === true) { + const { PassThrough } = require("node:stream"); + const stream = new PassThrough(); + const chunkId = `chatcmpl-cursor-${Date.now()}`; + let deltasSent = false; + const writeChunk = (delta, finishReason = null) => { + stream.write(`data: ${JSON.stringify({ + id: chunkId, + object: "chat.completion.chunk", + created: Math.floor(Date.now() / 1000), + model, + choices: [{ index: 0, delta, finish_reason: finishReason }], + })}\n\n`); + }; + writeChunk({ role: "assistant" }); + // SSE comment keep-alive: harmless to parsers, keeps idle-timeout-prone + // clients on passthrough surfaces alive during silent gaps. + const keepAlive = setInterval(() => { try { stream.write(": ka\n\n"); } catch { /* stream gone */ } }, 15000); + keepAlive.unref?.(); + cursorUtils.runCursorAgent({ + prompt: fullPrompt, + resumePrompt: cursorUtils.latestUserTurnPrompt(body), + sessionKey: cursorUtils.deriveSessionKey(body), + model, + baseTimeoutMs, + binaryPath, + workspace, + onDelta: (text, kind) => { + deltasSent = true; + if (kind === "thought") { + writeChunk({ reasoning_content: text }); + } else if (kind === "tool") { + // Visible in EVERY client (Cursor chat ignores reasoning_content), + // badge-styled so stripLynkrBadges removes it from resubmitted + // history — narration renders once and never re-enters context. + writeChunk({ content: `\n\n*[Lynkr] ${String(text).replace(/[*\n]+/g, " ").trim()}*\n\n` }); + } else { + writeChunk({ content: text }); + } + }, + }).then((run) => { + clearInterval(keepAlive); + if (!deltasSent && run.text) writeChunk({ content: run.text }); + writeChunk({}, "stop"); + stream.write("data: [DONE]\n\n"); + stream.end(); + }).catch((err) => { + clearInterval(keepAlive); + logger.warn({ err: err.message }, "[Cursor] streaming serve failed mid-flight"); + writeChunk({ content: `\n[cursor provider error: ${String(err.message || err).slice(0, 200)}]` }, "stop"); + stream.write("data: [DONE]\n\n"); + stream.end(); + }); + return { ok: true, status: 200, stream, contentType: "text/event-stream", headers: {} }; + } + + let run; + try { + run = await cursorUtils.runCursorAgent({ + prompt: fullPrompt, + resumePrompt: cursorUtils.latestUserTurnPrompt(body), + sessionKey: cursorUtils.deriveSessionKey(body), + model, + baseTimeoutMs, + binaryPath, + workspace, + }); + } catch (err) { + // classifyCursorError already attached the honest cause (timeout vs + // auth vs missing binary) — no blanket subscription hint here. + throw new Error(`Cursor CLI failed: ${err.message || err}`); + } + + if (!run.text) { + throw new Error("Cursor: empty response from cursor-agent"); + } + + const anthropicJson = cursorUtils.convertCursorResponseToAnthropic(run.text, model, run.usage, run.thinking); + + return { + ok: true, + status: 200, + json: anthropicJson, + text: JSON.stringify(anthropicJson), + contentType: "application/json", + }; +} + /** * Compute request cost in USD from model pricing × token usage. * Registry returns per-1M-token prices ({ input, output }); returns null when @@ -3457,6 +3572,13 @@ function captureResponseText(resultJson) { // non-greedy lazy match is unnecessary — match up to (and including) the // closing `*` plus trailing whitespace. const LYNKR_BADGE_PREFIX_RE = /^\*\[Lynkr\][^*\n]*\*\s*/; +// Mid-message badge LINES (tool narration, 2026-09-27) — strip anywhere in +// assistant content, swallowing surrounding blank lines so paragraphs reflow. +const LYNKR_BADGE_LINE_RE = /\n{0,2}\*\[Lynkr\][^*\n]*\*[ \t]*(?=\n|$)/g; + +function _stripBadgeText(text) { + return text.replace(LYNKR_BADGE_PREFIX_RE, "").replace(LYNKR_BADGE_LINE_RE, ""); +} function stripLynkrBadges(messages) { if (!Array.isArray(messages)) return messages; @@ -3468,8 +3590,8 @@ function stripLynkrBadges(messages) { // what the orchestrator's OpenAI-format response branch produces, and // it's where badges actually leak in the Ollama agent loop. if (typeof msg.content === 'string') { - if (!LYNKR_BADGE_PREFIX_RE.test(msg.content)) return msg; - const stripped = msg.content.replace(LYNKR_BADGE_PREFIX_RE, ''); + const stripped = _stripBadgeText(msg.content); + if (stripped === msg.content) return msg; mutated = true; // Badge-only content must not become an empty string — Anthropic // rejects empty assistant content (this is the interrupted-response @@ -3487,9 +3609,9 @@ function stripLynkrBadges(messages) { let changed = false; const rebuilt = []; for (const b of msg.content) { - if (b?.type === 'text' && typeof b.text === 'string' && LYNKR_BADGE_PREFIX_RE.test(b.text)) { + if (b?.type === 'text' && typeof b.text === 'string' && _stripBadgeText(b.text) !== b.text) { changed = true; - const strippedText = b.text.replace(LYNKR_BADGE_PREFIX_RE, ''); + const strippedText = _stripBadgeText(b.text); if (strippedText.trim()) rebuilt.push({ ...b, text: strippedText }); // badge-only block → drop } else { @@ -3533,6 +3655,7 @@ const PROVIDER_INVOKERS = { vertex: invokeVertex, moonshot: invokeMoonshot, codex: invokeCodex, + cursor: invokeCursor, baidu: invokeBaidu, fireworks: invokeFireworks, orcarouter: invokeOrcaRouter, @@ -4423,6 +4546,7 @@ module.exports = { invokeFireworks, invokeAtlas, invokeOrcaRouter, + invokeCursor, PROVIDER_INVOKERS, stripLynkrBadges, destroyHttpAgents, diff --git a/src/clients/health-probe.js b/src/clients/health-probe.js index b168c71..17deaaf 100644 --- a/src/clients/health-probe.js +++ b/src/clients/health-probe.js @@ -60,6 +60,16 @@ function _registerBuiltins() { : {}; probes.set('lmstudio', () => _cheapGet(`${config.lmstudio.endpoint}/v1/models`, headers)); } + if (config.cursor?.enabled === true) { + // Cheap local check: the CLI binary responds to --version. Auth status + // (subscription) is NOT probed here — a logged-out CLI still prints a + // version, and per-request failures surface via normal tier fallback. + probes.set('cursor', async () => { + const { execFileSync } = require('node:child_process'); + const binary = config.cursor?.binaryPath?.trim() || process.env.CURSOR_BINARY_PATH?.trim() || 'cursor-agent'; + execFileSync(binary, ['--version'], { timeout: PROBE_TIMEOUT_MS, stdio: 'ignore' }); + }); + } } /** diff --git a/src/config/index.js b/src/config/index.js index 34a73d6..fc9b636 100644 --- a/src/config/index.js +++ b/src/config/index.js @@ -727,6 +727,12 @@ var config = { model: process.env.CODEX_MODEL?.trim() || "gpt-5.3-codex", timeout: Number.parseInt(process.env.CODEX_TIMEOUT || "120000", 10) || 120000, }, + cursor: { + enabled: process.env.CURSOR_ENABLED === "true", + binaryPath: process.env.CURSOR_BINARY_PATH?.trim() || "cursor-agent", + model: process.env.CURSOR_MODEL?.trim() || "composer-2.5", + timeout: Number.parseInt(process.env.CURSOR_TIMEOUT || "120000", 10) || 120000, + }, hotReload: { enabled: hotReloadEnabled, debounceMs: Number.isNaN(hotReloadDebounceMs) ? 1000 : hotReloadDebounceMs, @@ -1197,6 +1203,10 @@ function reloadConfig() { config.orcarouter.authBaseUrl = orcaAuth; config.orcarouter.apiBaseUrl = orcaApi; config.orcarouter.endpoint = process.env.ORCAROUTER_ENDPOINT?.trim() || `${orcaApi}/v1/chat/completions`; + config.cursor.enabled = process.env.CURSOR_ENABLED === "true"; + config.cursor.binaryPath = process.env.CURSOR_BINARY_PATH?.trim() || "cursor-agent"; + config.cursor.model = process.env.CURSOR_MODEL?.trim() || "composer-2.5"; + config.cursor.timeout = Number.parseInt(process.env.CURSOR_TIMEOUT || "120000", 10) || 120000; // Model provider settings const newProvider = (process.env.MODEL_PROVIDER ?? "databricks").toLowerCase(); diff --git a/src/orchestrator/index.js b/src/orchestrator/index.js index 5c71b2b..ff67ef5 100644 --- a/src/orchestrator/index.js +++ b/src/orchestrator/index.js @@ -65,6 +65,8 @@ function getDestinationUrl(providerType) { return config.orcarouter?.endpoint ?? 'https://api.orcarouter.ai/v1/chat/completions'; case 'codex': return 'codex://app-server (local process)'; + case 'cursor': + return 'cursor://cursor-agent (local CLI)'; default: return 'unknown'; } @@ -2783,6 +2785,12 @@ IMPORTANT TOOL USAGE RULES: if (Array.isArray(anthropicPayload?.content)) { anthropicPayload.content = policy.sanitiseContent(anthropicPayload.content); } + } else if (actualProvider === "cursor") { + // Cursor CLI responses are already in Anthropic format from invokeCursor + anthropicPayload = databricksResponse.json; + if (Array.isArray(anthropicPayload?.content)) { + anthropicPayload.content = policy.sanitiseContent(anthropicPayload.content); + } } else if (databricksResponse.json?.type === "message" && Array.isArray(databricksResponse.json?.content)) { // Shape-detected: already Anthropic (some clients convert upstream). // Re-converting via toAnthropicResponse reads the absent choices[] diff --git a/src/orchestrator/sse-transformer.js b/src/orchestrator/sse-transformer.js index 3acefb2..8c4c7fb 100644 --- a/src/orchestrator/sse-transformer.js +++ b/src/orchestrator/sse-transformer.js @@ -25,9 +25,9 @@ const logger = require("../logger"); // invoke fn converts buffered responses to Anthropic, and its Anthropic-format // endpoint has the native passthrough path instead). moonshot joined the list // once invokeMoonshot learned to return the raw stream (its old forced -// stream:false predated this transformer). Caveat: reasoning_content deltas -// (kimi thinking) are not reshaped — thinking text is dropped from streamed -// responses; the buffered path still lifts it into thinking blocks. +// stream:false predated this transformer). reasoning_content deltas (kimi/ +// glm/gpt-oss/cursor thoughts) are reshaped into Anthropic thinking blocks +// (2026-09-27) so the thinking phase streams live instead of dropping. // baidu (Qianfan's /v2/chat/completions) already returns the raw stream the // same way moonshot does (invokeBaidu's `if (response?.stream) return // response;`) and its endpoint is documented as OpenAI-compatible — but @@ -38,6 +38,7 @@ const logger = require("../logger"); // shared transformer for one provider's quirk. const DEFAULT_OPENAI_SSE_PROVIDERS = [ "openai", + "cursor", "atlas", "azure-openai", "openrouter", @@ -163,6 +164,7 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { let model = fallbackModel; let nextIndex = 0; let textIndex = null; // open text block index, null when closed + let thinkingIndex = null; // open thinking block index, null when closed let finishReason = null; const toolAcc = new Map(); // openai tool index -> { id, name, args } @@ -239,6 +241,12 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { return out; }; + const closeThinkingBlock = function* () { + if (thinkingIndex === null) return; + yield _sse("content_block_stop", { type: "content_block_stop", index: thinkingIndex }); + thinkingIndex = null; + }; + const closeTextBlock = function* () { // Flush a held-back fragment that turned out not to be a tag. if (_tagTail && !_inThink && textIndex !== null) { @@ -306,11 +314,34 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { yield* startMessage(); + // Reasoning deltas (kimi/glm/gpt-oss/cursor-ACP thoughts) become a + // proper Anthropic thinking block so clients show live activity during + // the thinking phase instead of dead air (2026-09-26 complaint: badge, + // then minutes of silence). Thinking and text are sequential blocks — + // close one before opening the other. + if (typeof delta.reasoning_content === "string" && delta.reasoning_content.length > 0) { + yield* closeTextBlock(); + if (thinkingIndex === null) { + thinkingIndex = nextIndex++; + yield _sse("content_block_start", { + type: "content_block_start", + index: thinkingIndex, + content_block: { type: "thinking", thinking: "" }, + }); + } + yield _sse("content_block_delta", { + type: "content_block_delta", + index: thinkingIndex, + delta: { type: "thinking_delta", thinking: delta.reasoning_content }, + }); + } + // Text deltas are straightforward: open a block on first text, then // emit a text_delta per chunk. if (typeof delta.content === "string" && delta.content.length > 0) { const visible = filterThink(delta.content); if (visible.length > 0) { + yield* closeThinkingBlock(); if (textIndex === null) { textIndex = nextIndex++; yield _sse("content_block_start", { @@ -331,6 +362,7 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { // function.arguments are NOT parseable individually — only the full // concatenation is valid JSON, so blocks are emitted at stream end. if (Array.isArray(delta.tool_calls)) { + yield* closeThinkingBlock(); yield* closeTextBlock(); for (const tc of delta.tool_calls) { const idx = tc.index ?? 0; @@ -345,6 +377,7 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { } catch (err) { logger.warn({ err: err.message }, "[SSETransform] Upstream stream failed mid-flight"); yield* startMessage(); + yield* closeThinkingBlock(); yield* closeTextBlock(); yield _sse("error", { type: "error", @@ -358,6 +391,7 @@ async function* _openaiToAnthropicEvents(upstream, opts = {}) { // Normal end of stream ([DONE] or upstream EOF). yield* startMessage(); + yield* closeThinkingBlock(); yield* closeTextBlock(); // An EOF with no finish_reason is a dropped upstream, not a completed diff --git a/src/routing/capabilities.js b/src/routing/capabilities.js index eaf9029..2bd968b 100644 --- a/src/routing/capabilities.js +++ b/src/routing/capabilities.js @@ -71,6 +71,12 @@ function buildRequirementVector({ dimensions = {}, agenticResult = null } = {}) // Agentic work always needs tool orchestration — floor only, never lower. if (agenticResult?.isAgentic) { toolUse = Math.max(toolUse, 0.6); + } else if (agenticResult && agenticResult.isAgentic === false) { + // The detector already scored this turn NON-agentic using EFFECTIVE + // tools (harness baseline subtracted, WS3.2). Attached-tool dims must + // not out-vote that verdict: a bare "Hi" from a harness with 9 default + // tools is not a tool-orchestration workload. Cap, never raise. + toolUse = Math.min(toolUse, 0.25); } const round3 = (v) => Math.round(_clamp01(v) * 1000) / 1000; diff --git a/src/routing/complexity-analyzer.js b/src/routing/complexity-analyzer.js index 2a99cc1..09a4ce9 100644 --- a/src/routing/complexity-analyzer.js +++ b/src/routing/complexity-analyzer.js @@ -422,18 +422,22 @@ function extractContent(payload) { return ''; } - // Get last user message + // Get last user message. Harness envelopes (Cursor's // + // attached context inside the user message) are stripped so force patterns, + // risk keywords (risk-analyzer imports this) and text dims score the ASK, + // not the harness boilerplate. + const { stripHarnessEnvelope } = require('./harness-envelope'); for (let i = payload.messages.length - 1; i >= 0; i--) { const msg = payload.messages[i]; if (msg?.role === 'user') { if (typeof msg.content === 'string') { - return msg.content; + return stripHarnessEnvelope(msg.content); } if (Array.isArray(msg.content)) { - return msg.content + return stripHarnessEnvelope(msg.content .filter(block => block?.type === 'text') .map(block => block.text || '') - .join(' '); + .join(' ')); } } } diff --git a/src/routing/harness-envelope.js b/src/routing/harness-envelope.js new file mode 100644 index 0000000..1ff1352 --- /dev/null +++ b/src/routing/harness-envelope.js @@ -0,0 +1,71 @@ +/** + * Harness-envelope hygiene for trigger/risk/dimension text. + * + * GUI harnesses (Cursor at minimum) deliver the user's ask embedded in a + * context envelope INSIDE the user message: , , + * (workspace rules — routinely contain phrases like "never expose + * API keys" or "do not deploy to production"), , attached + * file contents. Trigger-style scanners (force patterns, risk keywords) and + * text dimensions must evaluate the ASK, not the harness's boilerplate — a + * workspace rule about keys must not make "Hi" high-risk (live incident + * 2026-09-26: routing_method=risk REASONING serve on a bare greeting). + * + * Claude Code payloads never carry these tags (its wrapper text rides in + * system-reminders, stripped elsewhere) — this is a no-op for them. + * + * Pure functions, no I/O. Never throws. + */ + +const ENVELOPE_TAGS = [ + 'user_info', + 'git_status', + 'agent_transcripts', + 'rules', + 'always_applied_workspace_rules', + 'uuid', + 'project_layout', + 'attached_files', + 'file_contents', + 'additional_data', + 'custom_instructions', + 'linter_errors', + 'recently_viewed_files', + 'open_files', + 'user_rules', + 'memories', + 'workspace_rules', +]; + +const PAIRED_RES = ENVELOPE_TAGS.map( + (t) => new RegExp(`<${t}(?:\\s[^>]*)?>[\\s\\S]*?<\\/${t}>`, 'gi') +); +// Unclosed blocks that OPEN at line start swallow to end-of-string (telemetry +// and upstream truncation cut envelopes mid-block). Mid-line opens are +// preserved — a user QUOTING a tag is content, not envelope (same rationale +// as the jev-router cleaner's line-start rule). +const UNCLOSED_RES = ENVELOPE_TAGS.map( + (t) => new RegExp(`(?:^|\\n)<${t}(?:\\s[^>]*)?>[\\s\\S]*$`, 'i') +); +const USER_QUERY_RE = /]*)?>([\s\S]*?)<\/user_query>/gi; + +/** + * @param {string} text - one user message's text content + * @returns {string} the user's ask with harness envelope blocks removed; + * when the harness marks the ask explicitly (), that wins. + */ +function stripHarnessEnvelope(text) { + if (typeof text !== 'string' || text.length === 0) return typeof text === 'string' ? text : ''; + try { + if (!text.includes('<')) return text; + const queries = [...text.matchAll(USER_QUERY_RE)].map((m) => m[1].trim()).filter(Boolean); + if (queries.length > 0) return queries.join(' '); + let out = text; + for (const re of PAIRED_RES) out = out.replace(re, ' '); + for (const re of UNCLOSED_RES) out = out.replace(re, ' '); + return out.replace(/[ \t]{2,}/g, ' ').replace(/\n{3,}/g, '\n\n').trim(); + } catch { + return text; + } +} + +module.exports = { stripHarnessEnvelope, ENVELOPE_TAGS }; diff --git a/src/routing/index.js b/src/routing/index.js index 75d9a8a..ef21a0b 100644 --- a/src/routing/index.js +++ b/src/routing/index.js @@ -98,6 +98,7 @@ function _enabledProviders() { if (config.ollama?.endpoint) out.push('ollama'); if (config.llamacpp?.endpoint) out.push('llamacpp'); if (config.lmstudio?.endpoint) out.push('lmstudio'); + if (config.cursor?.enabled === true) out.push('cursor'); return out; } @@ -1431,11 +1432,29 @@ async function _determineProviderSmartInner(payload, options = {}) { }); const result = sf.selectByShortfall(req, candidates); if (result) { - const agreed = result.selected.provider === provider && result.selected.model === selectedModel; + // Serve-path selection is bounded to ONE band above the legacy tier + // (same philosophy as the Jev one-band-up cap): capability estimates + // are soft evidence and must not buy multi-band jumps in one turn. + // The uncapped result above still shadow-logs what shortfall wanted. + const _tierLadder = ['SIMPLE', 'MEDIUM', 'COMPLEX', 'REASONING']; + const _legacyIdx = Math.max(0, _tierLadder.indexOf(tier)); + const _capTier = _tierLadder[Math.min(_legacyIdx + 1, _tierLadder.length - 1)]; + const serveResult = sf.selectByShortfall(req, candidates, { maxTier: _capTier }) || result; + // Humility gate: if the LEGACY pick's capabilities resolved to pure + // tier ignorance ('tier'/'tier-fallback' — we know nothing about the + // configured model), an "incapable" verdict is unfalsifiable; never + // serve an escalation from it (the composer-2.5 → grok incident). + const _legacyRow = (result.shortfalls || []).find( + (r) => r.provider === provider && r.model === selectedModel + ); + const _legacyBlind = /^tier/.test(String(_legacyRow?.source || '')); + const _serveEscalates = (_tierLadder.indexOf(serveResult.selected.tier) > _legacyIdx); + const agreed = serveResult.selected.provider === provider && serveResult.selected.model === selectedModel; shortfallInfo = { req, tau: result.tau, - selected: result.selected, + selected: serveResult.selected, + wanted: result.selected, agreed, legacy: { provider, model: selectedModel, tier }, }; @@ -1444,14 +1463,16 @@ async function _determineProviderSmartInner(payload, options = {}) { tau: result.tau, legacy: `${tier}:${provider}:${selectedModel}`, shortfall: `${result.selected.tier}:${result.selected.provider}:${result.selected.model}`, + served: `${serveResult.selected.tier}:${serveResult.selected.provider}:${serveResult.selected.model}`, + legacyCapSource: _legacyRow?.source || null, agreed, }, '[Routing] Shortfall shadow compare'); - if (sf.isEnabled() && !agreed) { + if (sf.isEnabled() && !agreed && !(_legacyBlind && _serveEscalates)) { const fromTier = tier; const fromModel = selectedModel; - provider = result.selected.provider; - selectedModel = result.selected.model; - tier = result.selected.tier; + provider = serveResult.selected.provider; + selectedModel = serveResult.selected.model; + tier = serveResult.selected.tier; analysis.tier = tier; method = method + '+shortfall'; if ((TIER_DEFINITIONS[tier]?.priority || 0) > (TIER_DEFINITIONS[fromTier]?.priority || 0)) { diff --git a/src/routing/intent-score.js b/src/routing/intent-score.js index 3f41da5..83090b6 100644 --- a/src/routing/intent-score.js +++ b/src/routing/intent-score.js @@ -134,11 +134,18 @@ function cosine(a, b) { return denom > 0 ? dot / denom : 0; } +// Classes scoring cannot run without. Optional classes (frontier) may be +// absent from an anchors file — classify()/blendScore() already handle a +// missing frontier centroid (sim -1, excluded below FRONTIER_MIN_SIM), so +// scoring degrades to the 3-class baseline instead of dying outright. +const REQUIRED_ANCHOR_CLASSES = ['trivial', 'substantive', 'heavyweight']; + /** * Embed every anchor text and mean them per class. * @param {Object} anchorsByClass * @param {(text:string)=>Promise} embedFn - * @returns {Promise|null>} null if any class has no vectors + * @returns {Promise|null>} null when a required class + * is missing from the file or produced no vectors (embedder down) */ async function buildCentroids(anchorsByClass, embedFn) { const centroids = {}; @@ -151,13 +158,25 @@ async function buildCentroids(anchorsByClass, embedFn) { if (Array.isArray(v) && v.length > 0) vectors.push(v); } catch { /* embed never throws by contract, belt-and-braces */ } } - if (vectors.length === 0) return null; // a class with no anchors is unusable + if (vectors.length === 0) { + logger.warn({ cls, anchorCount: texts.length }, '[IntentScore] Anchor class produced no embeddings (embedder down or degraded?)'); + return null; + } const dim = vectors[0].length; const mean = new Array(dim).fill(0); for (const v of vectors) for (let i = 0; i < dim; i++) mean[i] += v[i] / vectors.length; centroids[cls] = mean; } - return Object.keys(centroids).length === Object.keys(CLASS_VALUES).length ? centroids : null; + const missingRequired = REQUIRED_ANCHOR_CLASSES.filter((c) => !centroids[c]); + if (missingRequired.length > 0) { + logger.warn({ missing: missingRequired }, '[IntentScore] Anchors file missing required class(es) — anchor mode unusable'); + return null; + } + const missingOptional = Object.keys(CLASS_VALUES).filter((c) => !centroids[c]); + if (missingOptional.length > 0) { + logger.warn({ missing: missingOptional }, '[IntentScore] Optional anchor class(es) missing — reduced-class blend (frontier absent: REASONING reachable via force triggers only)'); + } + return centroids; } /** @@ -242,7 +261,7 @@ async function _loadDefaultCentroids() { const router = getKnnRouter(); const centroids = await buildCentroids(anchors, (t) => router.embed(t)); if (!centroids) { - logger.warn('[IntentScore] Anchor embedding failed (Ollama down?) — lexical fallback until next attempt'); + logger.warn('[IntentScore] Anchor centroids unavailable (cause logged above) — lexical fallback until next attempt'); return null; } try { diff --git a/src/routing/jev-router.js b/src/routing/jev-router.js index 31a2fcc..7748ff6 100644 --- a/src/routing/jev-router.js +++ b/src/routing/jev-router.js @@ -146,6 +146,9 @@ function _cleanUserText(msg) { } for (const re of _PAIRED_TAG_RES) text = text.replace(re, ''); for (const re of _UNCLOSED_TAG_RES) text = text.replace(re, ''); + // GUI-harness envelopes (Cursor's //attached context) + // ride INSIDE user messages — strip them so the ledger anchors on asks. + text = require('./harness-envelope').stripHarnessEnvelope(text); text = text.replace(/^\s*\[Lynkr\][^\n]*$/gm, '').trim(); if (!text) return null; if (/^\s*(\[SYSTEM NOTIFICATION|]|]|\[Request interrupted|This session is being continued from a previous conversation)/i.test(text)) return null; diff --git a/src/routing/model-registry.js b/src/routing/model-registry.js index 06ca60d..3ba8a7a 100644 --- a/src/routing/model-registry.js +++ b/src/routing/model-registry.js @@ -277,11 +277,14 @@ class ModelRegistry { output: info.cost?.output || 0, cacheRead: info.cost?.cache_read, cacheWrite: info.cost?.cache_write, - context: info.context || 128000, - maxOutput: info.output || 4096, + // models.dev nests limits under `limit` and modalities under + // `modalities` — reading them flat silently gave every entry the + // 128000/4096 defaults and vision:false. + context: info.limit?.context || 128000, + maxOutput: info.limit?.output || 4096, toolCall: info.tool_call ?? false, reasoning: info.reasoning ?? false, - vision: Array.isArray(info.input) && info.input.includes('image'), + vision: Array.isArray(info.modalities?.input) && info.modalities.input.includes('image'), source: 'models.dev', }; diff --git a/src/routing/model-tiers.js b/src/routing/model-tiers.js index 8f2200e..ef8bb67 100644 --- a/src/routing/model-tiers.js +++ b/src/routing/model-tiers.js @@ -370,6 +370,8 @@ class ModelTierSelector { return config.orcarouter?.model || null; case 'codex': return config.codex?.model || null; + case 'cursor': + return config.cursor?.model || null; case 'vertex': return config.vertex?.model || null; case 'databricks': diff --git a/src/routing/passthrough-route.js b/src/routing/passthrough-route.js index e7a76d0..9caa5f4 100644 --- a/src/routing/passthrough-route.js +++ b/src/routing/passthrough-route.js @@ -159,6 +159,13 @@ function shouldHoldForCache(pinModel, cacheState, opts = {}) { if (warm < holdMinPrefixTokens()) { return { hold: false, reason: 'downgrade_prefix_small', warmPrefixTokens: warm }; } + // Flat-rate subscription: the per-token premium the break-even evaluator + // amortizes does not exist (marginal cost $0), so the dollar leg can never + // legitimately clear a descent. With the cold/stale/small outs already + // taken above, a warm prefix on flat rate always holds. + if (opts.flatRate) { + return { hold: true, reason: 'hold_flat_rate_warm', warmPrefixTokens: warm }; + } try { const evaluate = opts.evaluateSwitch || require('./cache-switch-cost').evaluateSwitch; @@ -207,7 +214,7 @@ function shouldHoldForCache(pinModel, cacheState, opts = {}) { * @param {function|null} [args.evaluateSwitch] - injectable evaluator (tests). * @returns {{model:string|null, action:'verbatim'|'upgrade'|'pin_hold', reason:string, warmPrefixTokens:number|null}} */ -function decidePassthroughModel({ tierModel = null, clientModel = null, pinModel = null, tierMethod = null, tierPinned = null, cacheState = null, remainingTurns = null, sessionBurnPressure = null, evaluateSwitch = null } = {}) { +function decidePassthroughModel({ tierModel = null, clientModel = null, pinModel = null, tierMethod = null, tierPinned = null, cacheState = null, remainingTurns = null, sessionBurnPressure = null, evaluateSwitch = null, hasToolHistory = false, flatRate = false } = {}) { try { if (!isRoutingEnabled()) { return { model: clientModel, action: 'verbatim', reason: 'routing_disabled', warmPrefixTokens: null }; @@ -221,6 +228,16 @@ function decidePassthroughModel({ tierModel = null, clientModel = null, pinModel } const tierRank = familyRank(tierModel); const pinRank = familyRank(pinModel); + // Tool-loop invariant (same contract the orchestrator enforces): a frame + // carrying tool history must keep serving whatever model is mid-exchange. + // The downgrade gate has no business re-litigating the pin between a + // tool_use and its tool_result — that is where mid-loop descents were + // demoting COMPLEX tasks to the client model frame-by-frame. + if (hasToolHistory && pinRank !== null) { + return pinModel === clientModel + ? { model: clientModel, action: 'verbatim', reason: 'tool_loop_hold', warmPrefixTokens: null } + : { model: pinModel, action: 'pin_hold', reason: 'tool_loop_hold', warmPrefixTokens: null }; + } if (tierRank !== null && tierRank > clientRank) { // Rule 4a — the upgrade branch can also be a step DOWN from the pin // (pin Opus, fresh verdict Sonnet, client Haiku). Consult the gate for @@ -230,6 +247,7 @@ function decidePassthroughModel({ tierModel = null, clientModel = null, pinModel downgradeModel: tierModel, remainingTurns, sessionBurnPressure, + flatRate, ...(evaluateSwitch ? { evaluateSwitch } : {}), }); if (gate && gate.hold) { @@ -243,6 +261,7 @@ function decidePassthroughModel({ tierModel = null, clientModel = null, pinModel downgradeModel: clientModel, remainingTurns, sessionBurnPressure, + flatRate, ...(evaluateSwitch ? { evaluateSwitch } : {}), }); if (gate && gate.hold) { diff --git a/src/routing/risk-analyzer.js b/src/routing/risk-analyzer.js index 704097d..57b57ab 100644 --- a/src/routing/risk-analyzer.js +++ b/src/routing/risk-analyzer.js @@ -9,6 +9,7 @@ */ const { extractContent } = require('./complexity-analyzer'); +const { stripHarnessEnvelope } = require('./harness-envelope'); // Substring keywords found in file paths or instruction text. // Matched case-insensitively as raw substrings, so "auth" hits @@ -143,7 +144,13 @@ function findHits(keywords, haystack) { const HARNESS_BOILERPLATE_LINE = /^\s*(authentication successful[.!]?(\s+connected to .+)?|connected to \S+ MCP\.?|(\d+\s+)?MCP servers? need authentication\b.*|run \/mcp\b.*)\s*$/i; function stripSystemReminders(text) { if (typeof text !== 'string' || !text) return ''; - return text + // Cursor harness wrapper (live 2026-09-26): every turn's user message is + // prefixed with , and (persistent + // memories + ). Workspace rules routinely + // mention migration/permission/security and paths like auth/schema — a + // bare "Hi" scored high_risk_forced_tier → REASONING 100 on all of them. + // Shared with the other scorers so new harness tags land in one place. + return stripHarnessEnvelope(text) .replace(/[\s\S]*?<\/system-reminder>/g, ' ') // Codex harness blocks arrive as user-role messages and get merged into // the typed text by the orchestrator's consecutive-role coalescing. diff --git a/src/routing/shortfall.js b/src/routing/shortfall.js index c72e41c..15540ed 100644 --- a/src/routing/shortfall.js +++ b/src/routing/shortfall.js @@ -225,7 +225,8 @@ function _costValue(c) { * @param {object} req — requirement vector {reasoning, codegen, debugging, tool_use} in [0,1] * @param {Array<{provider, model, tier, cost?}>} candidates — catalog-constrained set * (callers pass getAllConfiguredModels() + model-registry costs) - * @param {object} [opts] — { tau, weights } + * @param {object} [opts] — { tau, weights, maxTier } — maxTier caps the pool at + * that tier's priority (serve-path escalation bound; shadow calls omit it) * @returns {null | { selected, shortfalls: Array<{provider, model, tier, shortfall, cost, source}>, tau }} */ function selectByShortfall(req, candidates, opts = {}) { @@ -234,9 +235,11 @@ function selectByShortfall(req, candidates, opts = {}) { if (!Array.isArray(candidates) || candidates.length === 0) return null; const tau = opts.tau ?? getTau(); const weights = opts.weights ?? getWeights(); + const maxPri = opts.maxTier ? (TIER_PRIORITY[opts.maxTier] || 4) : null; const rows = candidates .filter((c) => c && c.provider && c.model) + .filter((c) => maxPri === null || (TIER_PRIORITY[c.tier] || 0) <= maxPri) .map((c) => { const { caps, source } = resolveCapabilitiesWithSource(c); return { @@ -254,11 +257,14 @@ function selectByShortfall(req, candidates, opts = {}) { const pool = covering.length > 0 ? covering : rows; pool.sort((a, b) => { if (covering.length > 0) { - // Cheapest covering wins; cost tie (incl. all-unknown) breaks toward - // the LOWER tier — a covering lower tier is sufficient by definition, - // so prefer it over excess headroom (avoids over-provisioning). - if (a.cost !== b.cost) return a.cost - b.cost; - return (TIER_PRIORITY[a.tier] || 0) - (TIER_PRIORITY[b.tier] || 0); + // A covering lower tier is sufficient by definition, so tier wins + // BEFORE cost: unknown prices resolve to Infinity, and letting cost + // dominate let a seed-priced top-tier model beat the operator's own + // unpriced mid tiers (the "Hi → grok-4.7-high" incident). Cost only + // breaks ties within a tier. + const tp = (TIER_PRIORITY[a.tier] || 0) - (TIER_PRIORITY[b.tier] || 0); + if (tp !== 0) return tp; + return a.cost - b.cost; } // Nothing covers: minimal shortfall wins, tiebreak higher tier then cheaper. if (a.shortfall !== b.shortfall) return a.shortfall - b.shortfall; diff --git a/src/routing/vision.js b/src/routing/vision.js index aa14469..35d633b 100644 --- a/src/routing/vision.js +++ b/src/routing/vision.js @@ -14,7 +14,7 @@ const VISION_IMAGE_TOKEN_ESTIMATE = 1500; // Providers that can never transport image bytes (no vision path, even after // converter fixes). Everything else either forwards natively or converts. -const TRANSPORTLESS_PROVIDERS = new Set(['llamacpp', 'lmstudio', 'codex']); +const TRANSPORTLESS_PROVIDERS = new Set(['llamacpp', 'lmstudio', 'codex', 'cursor']); /** * Does a single content block carry image data? diff --git a/test/cursor-provider.test.js b/test/cursor-provider.test.js new file mode 100644 index 0000000..33845c2 --- /dev/null +++ b/test/cursor-provider.test.js @@ -0,0 +1,228 @@ +/** + * Cursor CLI provider — unit tests. + * + * Covers prompt flattening, CLI output parsing, Anthropic conversion, + * one-shot spawn arg shape (via injected execFn — never spawns a real + * `cursor-agent`), availability detection, and dispatch registration. + */ + +const test = require("node:test"); +const assert = require("node:assert/strict"); + +process.env.DATABRICKS_API_KEY = process.env.DATABRICKS_API_KEY || "test-key"; +process.env.DATABRICKS_API_BASE = process.env.DATABRICKS_API_BASE || "http://test.com"; +process.env.LOG_FILE_ENABLED = "false"; + +const cursorUtils = require("../src/clients/cursor-utils"); + +test("convertAnthropicToCursorPrompt passes a single user message through", () => { + const { prompt } = cursorUtils.convertAnthropicToCursorPrompt({ + messages: [{ role: "user", content: "Fix this bug" }], + }); + assert.equal(prompt, "Fix this bug"); +}); + +test("convertAnthropicToCursorPrompt flattens history with last user message first", () => { + const { prompt } = cursorUtils.convertAnthropicToCursorPrompt({ + messages: [ + { role: "user", content: "Here is my file" }, + { role: "assistant", content: [{ type: "text", text: "Got it" }] }, + { role: "user", content: "Now refactor it" }, + ], + }); + assert.match(prompt, /Previous conversation:/); + assert.match(prompt, /Now refactor it/); +}); + +test("extractCursorText handles json / text / stream-json shapes", () => { + assert.equal(cursorUtils.extractCursorText(JSON.stringify({ text: "hello" })), "hello"); + assert.equal(cursorUtils.extractCursorText(JSON.stringify({ result: "done" })), "done"); + assert.equal(cursorUtils.extractCursorText("plain output"), "plain output"); + assert.equal(cursorUtils.extractCursorText(""), ""); + const stream = [ + JSON.stringify({ type: "assistant", message: { content: [{ text: "he" }] } }), + JSON.stringify({ type: "assistant", message: { content: [{ text: "llo" }] } }), + ].join("\n"); + assert.equal(cursorUtils.extractCursorText(stream), "hello"); +}); + +test("convertCursorResponseToAnthropic returns a message block", () => { + const msg = cursorUtils.convertCursorResponseToAnthropic("hi there", "composer-2.5"); + assert.equal(msg.type, "message"); + assert.equal(msg.model, "composer-2.5"); + assert.deepEqual(msg.content, [{ type: "text", text: "hi there" }]); + assert.ok(msg.usage.output_tokens > 0); +}); + +test("runCursorAgent builds the expected CLI argv, prompt over STDIN (no --force)", async () => { + let seen = null; + const execFn = async ({ binaryPath, args, stdin, timeoutMs }) => { + seen = { binaryPath, args, stdin, timeoutMs }; + return JSON.stringify({ type: "result", result: "agent answer", session_id: "chat-1" }); + }; + const out = await cursorUtils.runCursorAgent({ + prompt: "do the thing", + model: "composer-2.5", + baseTimeoutMs: 5000, + binaryPath: "cursor-agent", + workspace: "/repo", + execFn, + }); + assert.equal(out.text, "agent answer"); + assert.equal(out.sessionId, "chat-1"); + assert.equal(out.resumed, false); + assert.equal(seen.binaryPath, "cursor-agent"); + assert.ok(seen.args.includes("-p"), "must run in print mode"); + assert.ok(seen.args.includes("--output-format"), "must request structured output"); + assert.ok(seen.args.includes("composer-2.5"), "must pin the tier model"); + assert.ok(seen.args.includes("--trust"), "headless runs need --trust"); + assert.ok(seen.args.includes("--approve-mcps"), "headless runs need pre-approved MCPs or they stall"); + assert.equal(seen.args.includes("--force"), cursorUtils.CURSOR_AUTO_APPROVE, "--force must track the CURSOR_AUTO_APPROVE operator policy"); + assert.ok(seen.args.includes("/repo"), "workspace must be threaded through"); + assert.equal(seen.stdin, "do the thing", "prompt must travel over stdin, not argv"); + assert.ok(!seen.args.includes("do the thing"), "argv must not carry the prompt (OS size ceiling)"); + assert.ok(seen.timeoutMs >= 5000, "timeout scales up from base, never below it"); +}); + +test("session resume: second turn sends --resume with only the new turn; stale resume falls back fresh", async () => { + const calls = []; + let failNextResume = false; + const execFn = async ({ args, stdin }) => { + calls.push({ args: [...args], stdin }); + if (args.includes("--resume") && failNextResume) { + const err = new Error("session not found"); + throw err; + } + return JSON.stringify({ result: "ok", session_id: "chat-abc" }); + }; + const common = { model: "composer-2.5", binaryPath: "cursor-agent", execFn, sessionKey: "sid:test-resume" }; + + const first = await cursorUtils.runCursorAgent({ ...common, prompt: "full conversation flatten", resumePrompt: "turn 1" }); + assert.equal(first.resumed, false, "no cached chat yet → fresh"); + assert.ok(!calls[0].args.includes("--resume")); + + const second = await cursorUtils.runCursorAgent({ ...common, prompt: "full conversation flatten v2", resumePrompt: "just the new turn" }); + assert.equal(second.resumed, true); + assert.ok(calls[1].args.includes("--resume"), "second turn must resume"); + assert.equal(calls[1].args[calls[1].args.indexOf("--resume") + 1], "chat-abc"); + assert.equal(calls[1].stdin, "just the new turn", "resumed sessions send only the newest turn"); + + failNextResume = true; + const third = await cursorUtils.runCursorAgent({ ...common, prompt: "full flatten v3", resumePrompt: "newest" }); + assert.equal(third.resumed, false, "failed resume must fall back to a fresh full-prompt attempt"); + const last = calls[calls.length - 1]; + assert.ok(!last.args.includes("--resume")); + assert.equal(last.stdin, "full flatten v3"); +}); + +test("timeout scales with payload and is hard-capped", () => { + const base = cursorUtils.scaleTimeoutMs(0, 120_000); + assert.equal(base, 120_000); + const big = cursorUtils.scaleTimeoutMs(1024 * 1024, 120_000); // 1MB + assert.ok(big > 120_000, "large payloads earn more time"); + const huge = cursorUtils.scaleTimeoutMs(100 * 1024 * 1024, 120_000); + assert.equal(huge, cursorUtils.MAX_TIMEOUT_MS, "hard cap holds"); +}); + +test("at most MAX_CONCURRENT_PROCS CLI spawns run at once; extras queue", async () => { + let inFlight = 0; + let peak = 0; + const execFn = async () => { + inFlight++; + peak = Math.max(peak, inFlight); + await new Promise((r) => setTimeout(r, 30)); + inFlight--; + return JSON.stringify({ result: "ok" }); + }; + await Promise.all( + Array.from({ length: 5 }, (_, i) => + cursorUtils.runCursorAgent({ prompt: `p${i}`, model: "m", binaryPath: "b", execFn }) + ) + ); + assert.ok(peak <= cursorUtils.MAX_CONCURRENT_PROCS, `peak ${peak} must respect the gate`); +}); + +test("classifyCursorError reports timeouts as timeouts and auth as auth", () => { + const t = new Error("spawnSync cursor-agent ETIMEDOUT"); + t.killed = true; + const classified = cursorUtils.classifyCursorError(t, { payloadBytes: 250 * 1024, timeoutMs: 120_000 }); + assert.match(classified.message, /timed out after 120s/); + assert.match(classified.message, /250KB/); + assert.ok(!/subscription/.test(classified.message), "a timeout must not be blamed on the subscription"); + + const a = new Error("exit 1"); + a.stderr = "Error: not logged in. Please sign in."; + const auth = cursorUtils.classifyCursorError(a, {}); + assert.match(auth.message, /cursor-agent login/); +}); + +test("parseCursorResult surfaces session id and real token usage", () => { + const out = JSON.stringify({ + type: "result", result: "OK", session_id: "s-1", + usage: { inputTokens: 8408, outputTokens: 32, cacheReadTokens: 3904, cacheWriteTokens: 0 }, + }); + const parsed = cursorUtils.parseCursorResult(out); + assert.equal(parsed.sessionId, "s-1"); + assert.equal(parsed.usage.inputTokens, 8408); + const msg = cursorUtils.convertCursorResponseToAnthropic(parsed.text, "composer-2.5", parsed.usage); + assert.equal(msg.usage.input_tokens, 8408); + assert.equal(msg.usage.cache_read_input_tokens, 3904); +}); + +test("deriveSessionKey: _sessionId wins, first-message hash is the fallback, empty body → null", () => { + assert.equal(cursorUtils.deriveSessionKey({ _sessionId: "abc" }), "sid:abc"); + const k1 = cursorUtils.deriveSessionKey({ messages: [{ role: "user", content: "hello world" }] }); + const k2 = cursorUtils.deriveSessionKey({ messages: [{ role: "user", content: "hello world" }, { role: "assistant", content: "hi" }] }); + assert.equal(k1, k2, "key must be stable as the conversation grows"); + assert.equal(cursorUtils.deriveSessionKey({ messages: [] }), null); +}); + +test("isAvailable honors the injected which implementation", () => { + assert.equal(cursorUtils.isAvailable(() => {}), true); + assert.equal( + cursorUtils.isAvailable(() => { + throw new Error("not found"); + }), + false + ); +}); + +test("getBinaryPath prefers explicit config, then env, then default", () => { + assert.equal(cursorUtils.getBinaryPath({ binaryPath: "/opt/cursor-agent" }), "/opt/cursor-agent"); + const prev = process.env.CURSOR_BINARY_PATH; + process.env.CURSOR_BINARY_PATH = "my-agent"; + assert.equal(cursorUtils.getBinaryPath({}), "my-agent"); + if (prev === undefined) delete process.env.CURSOR_BINARY_PATH; + else process.env.CURSOR_BINARY_PATH = prev; + assert.equal(cursorUtils.getBinaryPath({}), "cursor-agent"); +}); + +test("cursor invoker is registered and refuses empty prompts", async () => { + const { PROVIDER_INVOKERS, invokeCursor } = require("../src/clients/databricks"); + assert.equal(typeof PROVIDER_INVOKERS.cursor, "function"); + assert.equal(PROVIDER_INVOKERS.cursor, invokeCursor); + const config = require("../src/config"); + const prev = config.cursor?.enabled; + if (config.cursor) config.cursor.enabled = true; + try { + await assert.rejects(() => invokeCursor({ messages: [] }), /no prompt content/); + } finally { + if (config.cursor) config.cursor.enabled = prev; + } +}); + +test("cursor is a transportless (text-only) provider like codex", () => { + const fs = require("node:fs"); + const src = fs.readFileSync(`${__dirname}/../src/routing/vision.js`, "utf8"); + assert.match(src, /'cursor'/); +}); + +test("agent spawns in an empty sandbox cwd unless a workspace is explicit (server-repo leak fix)", async () => { + const seen = []; + const execFn = async ({ cwd }) => { seen.push(cwd); return JSON.stringify({ result: "ok" }); }; + await cursorUtils.runCursorAgent({ prompt: "review this project", model: "m", binaryPath: "b", execFn }); + await cursorUtils.runCursorAgent({ prompt: "p", model: "m", binaryPath: "b", workspace: "/client/repo", execFn }); + assert.equal(seen[0], cursorUtils.getSandboxDir(), "no workspace → sandbox dir, never the server cwd"); + assert.ok(!seen[0].includes("claude-code"), "sandbox must not be the Lynkr repo"); + assert.equal(seen[1], "/client/repo", "explicit workspace is honored"); +}); diff --git a/test/embeddings-degradation.test.js b/test/embeddings-degradation.test.js index f9bffee..c79ad30 100644 --- a/test/embeddings-degradation.test.js +++ b/test/embeddings-degradation.test.js @@ -99,3 +99,12 @@ test('degradation is not a permanent latch: provider recovery is automatic after assert.equal(status.providerAvailable, true); assert.equal(status.degradedSince, null); }); + +test('a transient blip is absorbed by the in-place retry — no degradation flip', async () => { + failMode = true; + setTimeout(() => { failMode = false; }, 100); // provider recovers before the 1.5s retry fires + const vec = await embeddings.generateEmbedding('boot-race probe'); + assert.equal(vec.length, 768, 'retry must return the real embedding, not the hash fallback'); + const status = embeddings.getEmbeddingStatus(); + assert.equal(status.providerAvailable, true, 'a single blip must not trip 60s of degraded routing'); +}); diff --git a/test/harness-envelope.test.js b/test/harness-envelope.test.js new file mode 100644 index 0000000..b8f2ec2 --- /dev/null +++ b/test/harness-envelope.test.js @@ -0,0 +1,81 @@ +const assert = require('assert'); +const { describe, it } = require('node:test'); +const { stripHarnessEnvelope } = require('../src/routing/harness-envelope'); +const { analyzeRisk } = require('../src/routing/risk-analyzer'); +const ca = require('../src/routing/complexity-analyzer'); + +// Modeled on the live 2026-09-26 incident: Cursor's envelope rides INSIDE the +// user message; the workspace block contains risk/force vocabulary the +// user never typed, and telemetry-style truncation leaves the last block +// unclosed. The ask ("Hi") trails the envelope. +const ENVELOPE = ` +OS Version: darwin 25.5.0 +Shell: zsh +Workspace Path: /Users/someone/opencode-wrap + + + +On branch main. Modified: server.js, package.json + + + +Never expose API keys or credentials. Do not deploy to production without +approval. Always think step by step and verify your work before migrating +the system or restructuring the codebase per the migration plan. + + +Hi`; + +describe('stripHarnessEnvelope', () => { + it('removes paired envelope blocks, keeps the trailing ask', () => { + const out = stripHarnessEnvelope(ENVELOPE); + assert.strictEqual(out, 'Hi'); + }); + + it('unclosed block at line start (truncated envelope) swallows to end', () => { + const out = stripHarnessEnvelope('Hi again\n\nnever expose api keys and then it truncat'); + assert.strictEqual(out, 'Hi again'); + }); + + it('mid-line tag mention is user content, preserved', () => { + const t = 'why does in cursor payloads break my parser?'; + assert.strictEqual(stripHarnessEnvelope(t), t); + }); + + it(' wins outright when present', () => { + const out = stripHarnessEnvelope('xfix the bugdeploy stuff'); + assert.strictEqual(out, 'fix the bug'); + }); + + it('plain text and non-strings pass through safely', () => { + assert.strictEqual(stripHarnessEnvelope('just a normal ask'), 'just a normal ask'); + assert.strictEqual(stripHarnessEnvelope(null), ''); + }); + + it('claude-code wrappers are untouched (different mechanism, different strip)', () => { + const t = '[SUGGESTION MODE: xyz] some text'; + assert.strictEqual(stripHarnessEnvelope(t), t); + }); +}); + +describe('risk/force triggers scoped to the ask, not the envelope', () => { + const payload = (content) => ({ messages: [{ role: 'user', content }] }); + + it('incident regression: envelope rules vocab + "Hi" → risk low', () => { + const r = analyzeRisk(payload(ENVELOPE)); + assert.strictEqual(r.level, 'low', JSON.stringify(r)); + }); + + it('a genuinely risky ASK after the envelope still fires', () => { + const r = analyzeRisk(payload(ENVELOPE.replace(/Hi$/, 'rotate the production credentials in .env now'))); + assert.notStrictEqual(r.level, 'low'); + }); + + it('force-reasoning vocab inside the envelope does not force', () => { + assert.strictEqual(ca.shouldForceReasoning(payload(ENVELOPE)), false); + }); + + it('force-reasoning in the actual ask still forces', () => { + assert.strictEqual(ca.shouldForceReasoning(payload(ENVELOPE.replace(/Hi$/, 'ultrathink: prove this queue is correct'))), true); + }); +}); diff --git a/test/intent-score.test.js b/test/intent-score.test.js index 0cdb577..335bbc7 100644 --- a/test/intent-score.test.js +++ b/test/intent-score.test.js @@ -298,3 +298,45 @@ describe("WS7.3 — rung containment (REASONING reachable via frontier class + f assert.strictEqual(r, null); }); }); + +// --- reduced-class anchors (global installs ship without frontier) ---------- +// A stock npm install loads config/difficulty-anchors.json; if that file +// lacks an optional class, buildCentroids must degrade to the reduced blend, +// not reject every centroid (the 9.14.15 bug: 3-class config + 4-class +// CLASS_VALUES → null → permanent lexical fallback → everything SIMPLE). +describe("buildCentroids reduced-class tolerance", () => { + const dim3 = { + trivial: ["hi"], + substantive: ["review this helper"], + heavyweight: ["architecture review of the orchestrator"], + }; + const axes = { hi: [1, 0, 0], "review this helper": [0, 1, 0], "architecture review of the orchestrator": [0, 0, 1] }; + const embed = async (t) => axes[t] || [0.3, 0.3, 0.3]; + + it("builds centroids when only the three required classes exist", async () => { + const c = await buildCentroids(dim3, embed); + assert.ok(c, "3-class anchors must not be rejected"); + assert.deepStrictEqual(Object.keys(c).sort(), ["heavyweight", "substantive", "trivial"]); + }); + + it("missing a REQUIRED class → null (anchor mode unusable)", async () => { + const c = await buildCentroids({ trivial: ["hi"], substantive: ["review this helper"] }, embed); + assert.strictEqual(c, null); + }); + + it("a class whose embeds all fail → null (embedder down)", async () => { + const c = await buildCentroids(dim3, async () => null); + assert.strictEqual(c, null); + }); + + it("scoring with 3-class centroids: substantive text lands MEDIUM, REASONING unreachable", async () => { + const c = await buildCentroids(dim3, embed); + const { cls, sims } = classify([0, 1, 0], c); + assert.strictEqual(cls, "substantive"); + const score = blendScore(sims); + assert.ok(score >= 26 && score <= 50, `expected MEDIUM band, got ${score}`); + // frontier centroid absent → sim -1 → below FRONTIER_MIN_SIM → excluded: + const heavy = blendScore(classify([0, 0, 1], c).sims); + assert.ok(heavy <= 75, `3-class blend must stay out of REASONING band, got ${heavy}`); + }); +}); diff --git a/test/model-registry-cost.test.js b/test/model-registry-cost.test.js index 8c00b87..c6af17f 100644 --- a/test/model-registry-cost.test.js +++ b/test/model-registry-cost.test.js @@ -47,11 +47,10 @@ const PRICING_FIXTURE = { models: { "5.2": { cost: { input: 1.75, output: 14, cache_read: 0.175 }, - context: 128000, - output: 4096, + limit: { context: 128000, output: 4096 }, tool_call: true, reasoning: true, - input: ["text"], + modalities: { input: ["text"], output: ["text"] }, }, }, }, diff --git a/test/passthrough-route.test.js b/test/passthrough-route.test.js index 3805808..54c3687 100644 --- a/test/passthrough-route.test.js +++ b/test/passthrough-route.test.js @@ -300,3 +300,75 @@ describe('resolveTierModel (label wins ties)', () => { ); }); }); + +// --- tool-loop invariant + flat-rate gate (2026-09-24 mid-loop descent bug) -- +// Live failure: COMPLEX ask upgraded haiku→sonnet, then every tool_result +// frame descended back to haiku (downgrade_break_even_cleared once, then +// downgrade_cache_stale forever) because the gate ran mid tool-loop and its +// break-even leg priced a flat-fee subscription at metered rates. +describe('decidePassthroughModel tool-loop and flat-rate holds', () => { + const HAIKU2 = 'claude-haiku-4-5-20251001'; + const SONNET2 = 'claude-sonnet-4-5'; + const warmSonnet = (tokens) => ({ + warmPrefixTokens: tokens, provider: 'azure-anthropic', model: SONNET2, + lastRequestAt: 1, ttlMs: 300000, cold: false, + }); + + it('tool-history frame holds the pin unconditionally, even over a stale cache', () => { + const r = route.decidePassthroughModel({ + tierModel: HAIKU2, clientModel: HAIKU2, pinModel: SONNET2, + cacheState: { ...warmSonnet(50), model: HAIKU2 }, // stale — would descend + hasToolHistory: true, flatRate: true, + }); + assert.strictEqual(r.action, 'pin_hold'); + assert.strictEqual(r.model, SONNET2); + assert.strictEqual(r.reason, 'tool_loop_hold'); + }); + + it('tool-history frame with pin === client stays verbatim (no rewrite churn)', () => { + const r = route.decidePassthroughModel({ + tierModel: SONNET2, clientModel: HAIKU2, pinModel: HAIKU2, hasToolHistory: true, + }); + assert.strictEqual(r.action, 'verbatim'); + assert.strictEqual(r.reason, 'tool_loop_hold'); + assert.strictEqual(r.model, HAIKU2); + }); + + it('tool-history frame with NO pin falls through to normal rules', () => { + const r = route.decidePassthroughModel({ + tierModel: SONNET2, clientModel: HAIKU2, pinModel: null, hasToolHistory: true, + }); + assert.strictEqual(r.action, 'upgrade'); + assert.strictEqual(r.model, SONNET2); + }); + + it('flat rate + warm same-model prefix holds without consulting break-even', () => { + const r = route.decidePassthroughModel({ + tierModel: HAIKU2, clientModel: HAIKU2, pinModel: SONNET2, + cacheState: warmSonnet(12500), flatRate: true, + evaluateSwitch: () => { throw new Error('break-even must not run on flat rate'); }, + }); + assert.strictEqual(r.action, 'pin_hold'); + assert.strictEqual(r.reason, 'hold_flat_rate_warm'); + }); + + it('flat rate still descends for honest cache reasons (cold)', () => { + const r = route.decidePassthroughModel({ + tierModel: HAIKU2, clientModel: HAIKU2, pinModel: SONNET2, + cacheState: { ...warmSonnet(12500), cold: true }, flatRate: true, + }); + assert.strictEqual(r.action, 'verbatim'); + assert.strictEqual(r.reason, 'downgrade_cache_cold'); + assert.strictEqual(r.model, HAIKU2); + }); + + it('metered path (flatRate absent) keeps legacy break-even behavior', () => { + const r = route.decidePassthroughModel({ + tierModel: HAIKU2, clientModel: HAIKU2, pinModel: SONNET2, + cacheState: warmSonnet(12500), + evaluateSwitch: () => ({ switchAllowed: true, breakEvenTurns: 1, expectedRemainingTurns: 9 }), + }); + assert.strictEqual(r.action, 'verbatim'); + assert.strictEqual(r.reason, 'downgrade_break_even_cleared'); + }); +}); diff --git a/test/risk-analyzer.test.js b/test/risk-analyzer.test.js index 9528d05..1d63d9c 100644 --- a/test/risk-analyzer.test.js +++ b/test/risk-analyzer.test.js @@ -215,4 +215,33 @@ describe('analyzeRisk', () => { assert.strictEqual(r.level, 'high'); }); }); + + // Live incident (2026-09-26): Cursor prefixes every turn's user message + // with (OS/shell/store paths), and + // (persistent memories + ). + // Workspace rules mentioning migration/permission/security and paths + // like auth/schema/subscription forced a bare "Hi" to + // high_risk_forced_tier → REASONING 100 on every Cursor turn. + describe('cursor harness wrapper stripping', () => { + const WRAPPER = + '\nOS Version: darwin\nShell: zsh\n\n\n' + + '\nAgent transcripts live in /Users/x/.cursor/projects/y.\n\n\n' + + '\nHandle the security migration carefully.\n' + + '\n' + + 'Check src/auth/schema.ts and the subscription permission flow.\n' + + '\n'; + + it('trivial message inside cursor wrapper stays low', () => { + const r = analyzeRisk(userPayload(`${WRAPPER}\nHi`)); + assert.strictEqual(r.level, 'low', JSON.stringify(r)); + assert.deepStrictEqual(r.instructionHits, []); + assert.deepStrictEqual(r.pathHits, []); + }); + + it('genuinely risky typed text still fires inside a wrapper', () => { + const r = analyzeRisk(userPayload(`${WRAPPER}\ndisable the authentication check`)); + assert.strictEqual(r.level, 'high'); + assert.ok(r.instructionHits.includes('authentication')); + }); + }); }); diff --git a/test/shortfall.test.js b/test/shortfall.test.js index cef1889..094eab3 100644 --- a/test/shortfall.test.js +++ b/test/shortfall.test.js @@ -121,3 +121,84 @@ describe('shortfall matching', () => { assert.strictEqual(shortfall.selectByShortfall({ reasoning: 0.5 }, []), null); }); }); + +// --- 2026-09-25 "Hi → grok-4.7-high" incident fixes -------------------------- +// A Cursor "Hi" with 9 harness tools escalated SIMPLE→REASONING because: +// (a) attached-tool dims out-voted the detector's non-agentic verdict, +// (b) covering selection sorted cost before tier (unknown prices = Infinity, +// so the only seed-priced model — the top tier — always won), +// (c) nothing bounded serve-path escalation or required the legacy pick's +// capabilities to be actually KNOWN before declaring it incapable. +describe('non-agentic tool_use cap', () => { + it('caps tool_use when the detector says NOT agentic, despite heavy tool dims', () => { + const v = buildRequirementVector({ + dimensions: dims({ toolCount: 80, toolComplexity: 70, priorToolUsage: 60 }), + agenticResult: { isAgentic: false }, + }); + assert.ok(v.tool_use <= 0.25, `expected <=0.25, got ${v.tool_use}`); + }); + + it('agentic floor still wins at 0.6', () => { + const v = buildRequirementVector({ + dimensions: dims({ toolCount: 0, toolComplexity: 0 }), + agenticResult: { isAgentic: true }, + }); + assert.ok(v.tool_use >= 0.6); + }); + + it('no agenticResult → dims stand unmodified', () => { + const v = buildRequirementVector({ dimensions: dims({ toolCount: 80, toolComplexity: 70 }) }); + assert.ok(v.tool_use > 0.25); + }); +}); + +describe('selectByShortfall tier-before-cost + maxTier', () => { + beforeEach(() => { + // Empty profiles: every candidate resolves via tier-priority fallback + // (SIMPLE .25 / MEDIUM .5 / COMPLEX .75 / REASONING 1.0), deterministic. + shortfall._setProfilesForTests({ tau: 0.02, enabled: true, tierProfiles: {}, modelOverrides: {} }); + }); + afterEach(() => shortfall._resetProfilesCache()); + + const CANDIDATES = [ + { provider: 'cursor', model: 'composer-2.5', tier: 'SIMPLE', cost: Infinity }, + { provider: 'cursor', model: 'grok-4.6-medium', tier: 'MEDIUM', cost: Infinity }, + { provider: 'cursor', model: 'claude-opus-5-high', tier: 'COMPLEX', cost: Infinity }, + { provider: 'cursor', model: 'grok-4.7-high', tier: 'REASONING', cost: 5 }, + ]; + const REQ = { reasoning: 0.1, codegen: 0.1, debugging: 0.1, tool_use: 0.4 }; + + it('covering pool prefers the LOWEST covering tier over a cheaper higher tier', () => { + // SIMPLE (.25 caps) misses tool_use 0.4; MEDIUM+ cover. Old order picked + // REASONING (only finite price). Tier-first picks MEDIUM. + const r = shortfall.selectByShortfall(REQ, CANDIDATES, { tau: 0.02 }); + assert.strictEqual(r.selected.tier, 'MEDIUM'); + assert.strictEqual(r.selected.model, 'grok-4.6-medium'); + }); + + it('cost still breaks ties WITHIN a tier', () => { + const r = shortfall.selectByShortfall(REQ, [ + { provider: 'a', model: 'm-pricey', tier: 'MEDIUM', cost: 9 }, + { provider: 'b', model: 'm-cheap', tier: 'MEDIUM', cost: 2 }, + ], { tau: 0.02 }); + assert.strictEqual(r.selected.model, 'm-cheap'); + }); + + it('maxTier bounds the pool (serve-path one-band cap)', () => { + const r = shortfall.selectByShortfall( + { reasoning: 0.9, codegen: 0.9, debugging: 0.9, tool_use: 0.9 }, + CANDIDATES, + { tau: 0.02, maxTier: 'MEDIUM' } + ); + assert.ok(['SIMPLE', 'MEDIUM'].includes(r.selected.tier), `got ${r.selected.tier}`); + }); + + it('incident end-to-end: non-agentic Hi requirement is covered by the SIMPLE slot', () => { + const req = buildRequirementVector({ + dimensions: dims({ toolCount: 60, toolComplexity: 50, promptComplexity: 5, multiStepReasoning: 5, analysisDepth: 5, codeGeneration: 0, technicalDepth: 5 }), + agenticResult: { isAgentic: false }, + }); + const r = shortfall.selectByShortfall(req, CANDIDATES, { tau: 0.24 }); + assert.strictEqual(r.selected.tier, 'SIMPLE', `expected SIMPLE, got ${r.selected.tier} (req=${JSON.stringify(req)})`); + }); +}); diff --git a/test/token-budget-auto.test.js b/test/token-budget-auto.test.js index 2a82ed1..7938f96 100644 --- a/test/token-budget-auto.test.js +++ b/test/token-budget-auto.test.js @@ -69,7 +69,9 @@ test('virtual name without a pin falls to the conservative floor, never above', assert.ok(['min-tier', 'default'].includes(b.source), `expected floor source, got ${b.source}`); assert.ok(b.modelContextWindow > 0); assert.equal(b.effectiveMax, Math.floor(b.modelContextWindow * 0.85)); - assert.ok(b.effectiveMax <= 850000, 'floor must never exceed a big-model budget'); + // The configured ladder is now all 1M-class models (glm-5.2 / muse-spark / + // gpt-5.6-sol), so the min-tier floor can legitimately reach 0.85 × 1048576. + assert.ok(b.effectiveMax <= Math.floor(1048576 * 0.85), 'floor must never exceed a 1M-class budget'); }); test('invalid explicit cap (0 / garbage) is ignored, auto budget applies', () => {