import { BadRequestException, ForbiddenException, HttpException, HttpStatus, Injectable, InternalServerErrorException, NotFoundException } from '@nestjs/common'; import { chmodSync, closeSync, existsSync, mkdirSync, openSync, readFileSync, readdirSync, realpathSync, renameSync, rmSync, statSync, unlinkSync, writeFileSync } from 'node:fs'; import { execFileSync, spawn } from 'node:child_process'; import { randomUUID } from 'node:crypto'; import { join } from 'node:path'; import { APP_ROOT } from '../common/app-root.js'; import type { SettingsActor } from '../settings/settings.service.js'; export interface GoalStartInput { goal: string; projectId: string; } export interface GoalProject { project_id: string; domain: string; domain_root: string; context_roots: string[]; } export interface GoalProjectCreateInput { projectId: string; domain: string; } export interface GoalStage { id: string; status: string; detail: string; provider: string; model: string; updated_at?: string; } export interface GoalJob { id: string; trace_id: string; goal: string; status: 'queued' | 'running' | 'completed' | 'degraded' | 'failed' | 'requires_approval'; actor: string; tenant: string; project: string; created_at: string; updated_at: string; started_at?: string; finished_at?: string; local_provider: string; local_model: string; cloud_provider: string; cloud_model: string; stages: GoalStage[]; local_draft?: string; result?: string; error?: string; audit_hash?: string; local_usage?: Record; cloud_usage?: Record; workspace?: GoalProject; context_manifest?: { files: number; characters: number; truncated: boolean; path?: string }; approval?: { id: string; status: string; action: string }; reviewer_attempts?: GoalReviewerAttempt[]; patch_artifact?: { path: string; sha256: string; files: string[]; bytes: number; status: string; preview?: string; approval_id?: string; applied_at?: string; applied_by?: string }; verification?: Array<{ command: string; exit_code: number; output: string }>; } export interface GoalReviewerAttempt { attempt: number; provider: string; model: string; status: 'pass' | 'failed'; reason: string; retryable: boolean; started_at: string; finished_at: string; latency_ms: number; } interface ModelConnection { id: string; kind: 'local' | 'cloud' | 'gateway'; connected: boolean; models: string[]; defaultModel: string; } interface ConnectionList { success: boolean; connections: ModelConnection[]; } interface AccountProviderStatus { id: 'codex' | 'claude'; available: boolean; loggedIn: boolean; } const HARNESS_BIN = join(APP_ROOT, 'packages', 'casan-harness', 'scripts', 'bash'); const CONNECTIONS_CLI = join(HARNESS_BIN, 'model-connections.py'); const ORCHESTRATOR_CLI = join(HARNESS_BIN, 'goal-orchestrator.py'); const PATCH_EXECUTOR_CLI = join(HARNESS_BIN, 'goal-patch-executor.py'); const RBAC_CLI = join(HARNESS_BIN, 'rbac-check.py'); const PROJECT_REGISTRY = join(APP_ROOT, 'packages', 'casan-harness', 'level5', 'project-registry.json'); function parseJson(value: string): T | null { try { return JSON.parse(value) as T; } catch { return null; } } function safeTenant(value: string): string { const safe = value.replace(/[^a-zA-Z0-9._-]/g, '_').slice(0, 80); return safe || 'default'; } @Injectable() export class GoalsService { private readonly startWindows = new Map(); async start(input: GoalStartInput, actor: SettingsActor): Promise { const workspace = this.resolveProject(String(input.projectId ?? '')); this.requireRead(actor, workspace.project_id); const goal = String(input.goal ?? '').trim(); if (goal.length < 10 || goal.length > 8000) { throw new BadRequestException('GOAL_LENGTH_INVALID'); } this.enforceStartLimit(actor); const connections = this.connections(actor); const local = connections.find((connection) => connection.connected && connection.kind === 'local'); const cloudConnections = connections .filter((connection) => connection.connected && connection.kind === 'cloud' && (connection.defaultModel || connection.models[0])) .sort((left, right) => ({ openai: 0, anthropic: 1 }[left.id] ?? 9) - ({ openai: 0, anthropic: 1 }[right.id] ?? 9)); const cloud = cloudConnections[0]; const gateway = connections.find((connection) => connection.connected && connection.kind === 'gateway'); const account = await this.accountReviewer(); const localModel = local?.defaultModel || local?.models[0] || 'ornith:9b'; const localRuntime = local ? this.localRuntime(this.runtime(local.id, localModel, actor)) : { CASAN_CHAT_SELECTED_MODEL: `ollama:${localModel}`, CASAN_OLLAMA_HOST: process.env.CASAN_OLLAMA_HOST || 'host.docker.internal:11434', OLLAMA_HOST: process.env.OLLAMA_HOST || 'host.docker.internal:11434', }; const selectedReviewer = cloud ?? gateway; // Runtime credentials for a saved connection must use that connection's // discovered model, even when an environment-provided direct key has won // priority for this run. const cloudRuntime = cloud ? this.runtime(cloud.id, cloud.defaultModel || cloud.models[0] || '', actor) : {}; const gatewayRuntime = gateway ? this.runtime(gateway.id, gateway.defaultModel || gateway.models[0] || '', actor) : {}; const gatewayCredentials = Object.fromEntries( Object.entries(gatewayRuntime).filter(([key]) => key.startsWith('CASAN_OPENAI_COMPATIBLE_')), ); const savedCloudCandidates = cloudConnections.map((connection) => { const model = connection.defaultModel || connection.models[0]; const runtime = this.runtime(connection.id, model, actor); return { model: String(runtime.CASAN_CHAT_SELECTED_MODEL || model), runtime }; }); const envCloudCandidates = [ process.env.OPENAI_API_KEY ? { model: `openai:${process.env.CASAN_GOAL_OPENAI_MODEL || 'gpt-4.1'}`, runtime: {} } : null, process.env.ANTHROPIC_API_KEY ? { model: `anthropic:${process.env.CASAN_GOAL_ANTHROPIC_MODEL || 'claude-3-5-sonnet-latest'}`, runtime: {} } : null, ].filter((candidate): candidate is { model: string; runtime: Record } => candidate !== null); const cloudCandidates = [...envCloudCandidates, ...savedCloudCandidates].filter((candidate, index, rows) => rows.findIndex((row) => row.model === candidate.model) === index); const cloudCredentials = Object.assign({}, ...savedCloudCandidates.map(({ runtime }) => runtime)); const gatewayModels = gateway ? [gateway.defaultModel, ...gateway.models].filter((model, index, rows) => Boolean(model) && rows.indexOf(model) === index).slice(0, 5).map((model) => `openai-compatible:${model}`) : []; const preferredCloudModel = cloudCandidates[0]?.model || ''; const preferredCloudProvider = preferredCloudModel.startsWith('openai:') ? 'openai' : preferredCloudModel.startsWith('anthropic:') ? 'anthropic' : (cloud?.id || 'unavailable'); const cloudModel = preferredCloudModel || gateway?.defaultModel || gateway?.models[0] || ''; const id = randomUUID(); const timestamp = new Date().toISOString(); const job: GoalJob = { id, trace_id: id, goal, status: 'queued', actor: actor.actor, tenant: actor.tenant, project: workspace.project_id, workspace, created_at: timestamp, updated_at: timestamp, local_provider: local?.id || 'local-policy', local_model: String(localRuntime.CASAN_CHAT_SELECTED_MODEL || `ollama:${localModel}`), cloud_provider: account ? `${account}-account` : (preferredCloudProvider !== 'unavailable' ? preferredCloudProvider : (selectedReviewer?.id || 'unavailable')), cloud_model: account ? `${account}-account-default` : preferredCloudModel || String(cloudRuntime.CASAN_CHAT_SELECTED_MODEL || ''), stages: [ { id: 'local-worker', status: 'queued', detail: 'Waiting for local worker', provider: local?.id || 'local-policy', model: localModel }, { id: 'cloud-reviewer', status: 'queued', detail: account || preferredCloudModel || selectedReviewer ? 'Waiting for independent reviewer' : 'Cloud unavailable; local reviewer will be used', provider: account ? `${account}-account` : (preferredCloudProvider !== 'unavailable' ? preferredCloudProvider : (selectedReviewer?.id || 'local-policy')), model: account ? `${account}-account-default` : cloudModel || localModel }, ], }; const jobFile = this.jobPath(actor.tenant, id); mkdirSync(join(APP_ROOT, '.specify', 'state', 'goals', safeTenant(actor.tenant)), { recursive: true, mode: 0o700 }); writeFileSync(jobFile, `${JSON.stringify(job, null, 2)}\n`, { encoding: 'utf8', mode: 0o600 }); chmodSync(jobFile, 0o600); const child = spawn('python3', [ORCHESTRATOR_CLI, '--job-file', jobFile], { cwd: APP_ROOT, env: { ...process.env, ...localRuntime, ...cloudCredentials, ...cloudRuntime, ...gatewayCredentials, CASAN_TENANT_ID: actor.tenant || 'default', CASAN_GOAL_LOCAL_MODEL: job.local_model, CASAN_GOAL_CLOUD_MODEL: job.cloud_model, CASAN_GOAL_LOCAL_PROVIDER: job.local_provider, CASAN_GOAL_CLOUD_PROVIDER: job.cloud_provider, CASAN_GOAL_ACCOUNT_PROVIDER: account || '', CASAN_GOAL_CLOUD_FALLBACK_MODEL: String(cloudRuntime.CASAN_CHAT_SELECTED_MODEL || ''), CASAN_GOAL_CLOUD_MODELS: cloudCandidates.map(({ model }) => model).join(','), CASAN_GOAL_OMNIROUTE_MODELS: gatewayModels.join(','), // H2 recovery tries direct cloud credentials before gateway and local. CASAN_GOAL_PATCH_REPAIR_MODELS: [...cloudCandidates.map(({ model }) => model), ...gatewayModels, String(localRuntime.CASAN_CHAT_SELECTED_MODEL || `ollama:${localModel}`)].filter((model, index, rows) => rows.indexOf(model) === index).join(','), CASAN_GOAL_LOCAL_REVIEWER_MODEL: String(localRuntime.CASAN_CHAT_SELECTED_MODEL || `ollama:${localModel}`), CASAN_GOAL_REVIEWER_MAX_ATTEMPTS: process.env.CASAN_GOAL_REVIEWER_MAX_ATTEMPTS || '8', CASAN_GOAL_REVIEWER_DEADLINE_SEC: process.env.CASAN_GOAL_REVIEWER_DEADLINE_SEC || '600', }, stdio: 'ignore', }); child.on('error', () => { const failed = { ...job, status: 'failed' as const, error: 'GOAL_ORCHESTRATOR_START_FAILED', updated_at: new Date().toISOString() }; writeFileSync(jobFile, `${JSON.stringify(failed, null, 2)}\n`, { encoding: 'utf8', mode: 0o600 }); }); child.unref(); return job; } projects(actor: SettingsActor): { count: number; projects: GoalProject[] } { const projects = this.registeredProjects().filter((project) => { try { this.requireRead(actor, project.project_id); return true; } catch { return false; } }); return { count: projects.length, projects }; } createProject(input: GoalProjectCreateInput, actor: SettingsActor): GoalProject { if (actor.role !== 'org-admin') throw new ForbiddenException('GOAL_PROJECT_CREATE_DENIED'); const projectId = String(input.projectId ?? '').trim(); const domain = String(input.domain ?? '').trim(); if (!/^[A-Za-z][A-Za-z0-9._-]{2,63}$/.test(projectId)) { throw new BadRequestException('GOAL_PROJECT_ID_INVALID'); } if (domain.length < 3 || domain.length > 100) { throw new BadRequestException('GOAL_PROJECT_DOMAIN_INVALID'); } const lockPath = `${PROJECT_REGISTRY}.lock`; let lockFd: number | undefined; for (let attempt = 0; attempt < 100 && lockFd === undefined; attempt += 1) { try { lockFd = openSync(lockPath, 'wx', 0o600); } catch { Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 25); } } if (lockFd === undefined) throw new HttpException('GOAL_PROJECT_REGISTRY_BUSY', HttpStatus.CONFLICT); let absoluteRoot = ''; let createdRoot = false; try { const registry = parseJson<{ projects?: Array> }>(readFileSync(PROJECT_REGISTRY, 'utf8')); if (!registry || !Array.isArray(registry.projects)) throw new InternalServerErrorException('GOAL_PROJECT_REGISTRY_INVALID'); if (registry.projects.some((entry) => String(entry.project_id) === projectId)) throw new BadRequestException('GOAL_PROJECT_ALREADY_EXISTS'); const canonicalRoot = realpathSync(APP_ROOT); const appsRoot = realpathSync(join(APP_ROOT, 'apps')); if (!appsRoot.startsWith(`${canonicalRoot}/`)) throw new ForbiddenException('GOAL_PROJECT_PARENT_DENIED'); const projectsRoot = join(appsRoot, 'projects'); if (!existsSync(projectsRoot)) mkdirSync(projectsRoot, { mode: 0o750 }); const canonicalProjects = realpathSync(projectsRoot); if (!canonicalProjects.startsWith(`${canonicalRoot}/`)) throw new ForbiddenException('GOAL_PROJECT_PARENT_DENIED'); absoluteRoot = join(canonicalProjects, projectId); mkdirSync(absoluteRoot, { mode: 0o750 }); createdRoot = true; const relativeRoot = join('apps', 'projects', projectId); const entry = { project_id: projectId, domain, domain_root: relativeRoot, context_roots: [relativeRoot], harness_package: 'fpt-casan-sdd-harness', harness_version: '1.0.0', status: 'active' }; const updated = { ...registry, projects: [...registry.projects, entry] }; const temporary = `${PROJECT_REGISTRY}.${process.pid}.${randomUUID()}.tmp`; writeFileSync(temporary, `${JSON.stringify(updated, null, 2)}\n`, { encoding: 'utf8', mode: 0o600, flag: 'wx' }); renameSync(temporary, PROJECT_REGISTRY); return { project_id: projectId, domain, domain_root: relativeRoot, context_roots: [relativeRoot] }; } catch (error) { if (createdRoot && absoluteRoot) rmSync(absoluteRoot, { recursive: true, force: true }); if (error instanceof HttpException) throw error; throw new InternalServerErrorException('GOAL_PROJECT_CREATE_FAILED'); } finally { closeSync(lockFd); unlinkSync(lockPath); } } get(id: string, actor: SettingsActor): GoalJob { if (!/^[a-f0-9-]{36}$/.test(id)) throw new NotFoundException('GOAL_NOT_FOUND'); const path = this.jobPath(actor.tenant, id); if (!existsSync(path)) throw new NotFoundException('GOAL_NOT_FOUND'); const job = parseJson(readFileSync(path, 'utf8')); if (!job || job.tenant !== actor.tenant) throw new NotFoundException('GOAL_NOT_FOUND'); this.requireRead(actor, job.project); return job; } apply(id: string, actor: SettingsActor): GoalJob { const job = this.get(id, actor); if (!job.patch_artifact || job.status !== 'requires_approval') { throw new BadRequestException('GOAL_PATCH_NOT_READY'); } try { const output = execFileSync('python3', [PATCH_EXECUTOR_CLI, '--job-file', this.jobPath(actor.tenant, id), '--actor', actor.actor], { cwd: APP_ROOT, env: { ...process.env, CASAN_TENANT_ID: actor.tenant || 'default' }, encoding: 'utf8', stdio: ['ignore', 'pipe', 'pipe'], timeout: 10 * 60_000, }); const applied = parseJson(output); if (!applied) throw new InternalServerErrorException('GOAL_APPLY_RESPONSE_INVALID'); return applied; } catch (error: unknown) { if (error instanceof HttpException) throw error; const detail = error as { status?: number; stderr?: string | Buffer; stdout?: string | Buffer }; const message = String(detail.stderr || detail.stdout || 'GOAL_APPLY_FAILED').trim(); if (Number(detail.status) === 3) throw new ForbiddenException(message); throw new InternalServerErrorException(message); } } list(actor: SettingsActor, limit = 20): { count: number; goals: GoalJob[] } { this.requireRead(actor, actor.project); const directory = join(APP_ROOT, '.specify', 'state', 'goals', safeTenant(actor.tenant)); if (!existsSync(directory)) return { count: 0, goals: [] }; const goals = readdirSync(directory) .filter((name) => /^[a-f0-9-]{36}\.json$/.test(name)) .map((name) => parseJson(readFileSync(join(directory, name), 'utf8'))) .filter((job): job is GoalJob => Boolean(job && job.tenant === actor.tenant)) .filter((job) => { try { this.requireRead(actor, job.project); return true; } catch { return false; } }) .sort((left, right) => right.created_at.localeCompare(left.created_at)); return { count: goals.length, goals: goals.slice(0, Math.max(1, Math.min(limit, 100))) }; } private jobPath(tenant: string, id: string): string { return join(APP_ROOT, '.specify', 'state', 'goals', safeTenant(tenant), `${id}.json`); } private registeredProjects(): GoalProject[] { const parsed = parseJson<{ projects?: Array> }>(readFileSync(PROJECT_REGISTRY, 'utf8')); const root = realpathSync(APP_ROOT); return (parsed?.projects ?? []).filter((entry) => entry.status === 'active').map((entry) => { const projectId = String(entry.project_id ?? ''); const domainRoot = String(entry.domain_root ?? ''); const rawRoots = Array.isArray(entry.context_roots) ? entry.context_roots.map(String) : [domainRoot]; if (!/^[A-Za-z0-9._-]+$/.test(projectId) || !domainRoot || rawRoots.length === 0) { throw new InternalServerErrorException('GOAL_PROJECT_REGISTRY_INVALID'); } const contextRoots = rawRoots.map((relative) => { const absolute = realpathSync(join(APP_ROOT, relative)); if (!(absolute === root || absolute.startsWith(`${root}/`)) || !statSync(absolute).isDirectory() && !statSync(absolute).isFile()) { throw new InternalServerErrorException('GOAL_PROJECT_CONTEXT_ROOT_DENIED'); } return relative; }); return { project_id: projectId, domain: String(entry.domain ?? projectId), domain_root: domainRoot, context_roots: contextRoots }; }); } private resolveProject(projectId: string): GoalProject { if (!projectId) throw new BadRequestException('GOAL_PROJECT_REQUIRED'); const project = this.registeredProjects().find((entry) => entry.project_id === projectId); if (!project) throw new BadRequestException('GOAL_PROJECT_NOT_ALLOWED'); return project; } private connections(actor: SettingsActor): ModelConnection[] { const payload = this.runPython(CONNECTIONS_CLI, ['list'], { CASAN_TENANT_ID: actor.tenant || 'default' }); const parsed = parseJson(payload); if (!parsed?.success || !Array.isArray(parsed.connections)) { throw new InternalServerErrorException('GOAL_MODEL_CONNECTIONS_UNAVAILABLE'); } return parsed.connections; } private runtime(provider: string, model: string, actor: SettingsActor): Record { const payload = this.runPython( CONNECTIONS_CLI, ['runtime-env', '--provider', provider, '--model', model], { CASAN_TENANT_ID: actor.tenant || 'default' }, ); const parsed = parseJson<{ success?: boolean; env?: Record }>(payload); if (!parsed?.success || !parsed.env) throw new BadRequestException('GOAL_MODEL_RUNTIME_UNAVAILABLE'); return parsed.env; } private localRuntime(runtime: Record): Record { if (existsSync('/.dockerenv')) return runtime; const normalize = (value: string): string => value .replace('http://host.docker.internal:', 'http://127.0.0.1:') .replace('https://host.docker.internal:', 'https://127.0.0.1:') .replace(/^host\.docker\.internal:/, '127.0.0.1:'); return Object.fromEntries(Object.entries(runtime).map(([key, value]) => [key, normalize(value)])); } private async accountReviewer(): Promise<'claude' | 'codex' | ''> { if (process.env.CASAN_PROVIDER_ACCOUNT_AUTH_ENABLED !== '1') return ''; const bridgeUrl = (process.env.CASAN_AUTH_BRIDGE_URL || '').replace(/\/$/, ''); const bridgeToken = process.env.CASAN_AUTH_BRIDGE_TOKEN || ''; if (!bridgeUrl || !bridgeToken) return ''; try { const response = await fetch(`${bridgeUrl}/v1/auth/providers`, { headers: { 'X-CASAN-Bridge-Token': bridgeToken }, signal: AbortSignal.timeout(10_000), }); if (!response.ok) return ''; const payload = await response.json() as { providers?: AccountProviderStatus[] }; const available = payload.providers?.filter((provider) => provider.available && provider.loggedIn) ?? []; if (available.some((provider) => provider.id === 'claude')) return 'claude'; if (available.some((provider) => provider.id === 'codex')) return 'codex'; } catch { return ''; } return ''; } private enforceStartLimit(actor: SettingsActor): void { const key = `${safeTenant(actor.tenant)}:${actor.actor}`; const timestamp = Date.now(); const recent = (this.startWindows.get(key) ?? []).filter((value) => timestamp - value < 10 * 60_000); if (recent.length >= 5) { throw new HttpException('GOAL_RATE_LIMITED', HttpStatus.TOO_MANY_REQUESTS); } const directory = join(APP_ROOT, '.specify', 'state', 'goals', safeTenant(actor.tenant)); if (existsSync(directory)) { const active = readdirSync(directory) .filter((name) => /^[a-f0-9-]{36}\.json$/.test(name)) .map((name) => parseJson(readFileSync(join(directory, name), 'utf8'))) .filter((job) => job?.actor === actor.actor && (job.status === 'queued' || job.status === 'running')); if (active.length >= 2) { throw new HttpException('GOAL_CONCURRENCY_LIMITED', HttpStatus.TOO_MANY_REQUESTS); } } recent.push(timestamp); this.startWindows.set(key, recent); } private runPython(script: string, args: string[], environment: NodeJS.ProcessEnv): string { try { return execFileSync('python3', [script, ...args], { cwd: APP_ROOT, env: { ...process.env, ...environment }, encoding: 'utf8', stdio: ['ignore', 'pipe', 'pipe'], timeout: 20_000, }).trim(); } catch (error: unknown) { const detail = error as { stderr?: string | Buffer; stdout?: string | Buffer }; throw new InternalServerErrorException(String(detail.stderr || detail.stdout || 'GOAL_RUNTIME_FAILED').trim()); } } private requireRead(actor: SettingsActor, targetProject: string): void { try { execFileSync('python3', [RBAC_CLI, 'check', '--role', actor.role, '--resource', 'monitoring', '--action', 'read', '--role-project', actor.project, '--target-project', targetProject, '--role-tenant', actor.tenant, '--target-tenant', actor.tenant], { cwd: APP_ROOT, env: process.env, stdio: ['ignore', 'pipe', 'pipe'], }); } catch { throw new ForbiddenException('GOAL_RBAC_DENIED'); } } }