mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 21:05:21 +02:00
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work > - Environment runtime drivers provide workspace, lease, and custom image behavior > - Runtime code used driver identity checks and several capability-specific members > - These checks spread capability rules across the runtime and made new drivers harder to verify > - This pull request adds one general capability classifier and one static driver support table > - The benefit is one fail-closed capability model that keeps current behavior and supports future drivers ## Linked Issues or Issue Description **What existing behavior does this improve?** Environment runtime capability checks for workspace realization, custom images, lease capabilities, and duplex authorization. **Subsystem affected** Cross-cutting (multiple of the above) **Current behavior** The runtime selects several capability paths from driver identity and separate capability members. Custom image gates also trust provider declarations without checking every matching live worker method. **Proposed behavior** The runtime uses one general capability classifier and one static support table. Custom image gates require both the provider declaration and every matching live worker method. The public capability names remain unchanged. **Reason and benefit** The change keeps capability rules in one place. It removes identity conditions from runtime consumers and makes unsupported drivers fail closed. **Breaking changes** None. The public names sandboxCapabilities, sandboxProviders, and EffectiveSandboxCapabilities remain available. ## What Changed - Add classifyEnvironmentCapabilities and static support definitions for all four driver families. - Add resolveCapabilities to every environment runtime driver. - Move driver traits into environment-driver-traits.ts and migrate runtime consumers. - Require provider declarations and matching live worker methods for all custom image gates. - Migrate duplex authorization to the general resolver and remove the dead sandbox-only member. - Delete the unused resolveEffectiveSandboxCapabilities wrapper and update its test. ## Verification - pnpm --filter @paperclipai/server typecheck - pnpm exec vitest run server/src/__tests__/environment-capability-contract.test.ts server/src/__tests__/environment-runtime.test.ts — 92 tests pass - pnpm exec vitest run server/src/__tests__/environment-driver-traits.test.ts server/src/__tests__/general-capability-classifier.test.ts — 12 tests pass - pnpm exec vitest run server/src/__tests__/environment-custom-images-service.test.ts server/src/__tests__/environment-execution-target-capabilities.test.ts server/src/__tests__/environment-execution-target-duplex.test.ts server/src/__tests__/environment-execution-target-duplex-kill-switch.test.ts — 66 tests pass ## Risks The main risk is a capability gate that denies a valid driver or permits an invalid driver. The static support matrix, live worker method checks, and regression tests reduce this risk. No database, public API, or published type name changes. ## Model Used OpenAI Codex, GPT-5, with tool use and code execution. The deployment does not provide a separate context-window value. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with Fixes: # / Closes # / Refs # OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub #NNN / github.com/paperclipai/paperclip URLs) - [x] My branch name describes the change (for example, docs/... or fix/...) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
659 lines
24 KiB
TypeScript
659 lines
24 KiB
TypeScript
/**
|
|
* Centralized environment run orchestrator.
|
|
*
|
|
* Owns the full environment lifecycle for a heartbeat run:
|
|
* 1. Resolve selected environment
|
|
* 2. Validate environment is active and allowed
|
|
* 3. Acquire or resume lease
|
|
* 4. Realize workspace in the environment
|
|
* 5. Resolve execution target for the adapter
|
|
* 6. Release / retain / fail lease according to policy
|
|
* 7. Record activity and operator-visible status
|
|
*
|
|
* Heartbeat callers delegate to this service instead of inlining
|
|
* environment resolution, lease management, workspace realization,
|
|
* and transport logic.
|
|
*/
|
|
|
|
import type { Db } from "@paperclipai/db";
|
|
import type {
|
|
Environment,
|
|
EnvironmentLease,
|
|
EnvironmentLeasePolicy,
|
|
EnvironmentLeaseStatus,
|
|
ExecutionWorkspace,
|
|
ExecutionWorkspaceConfig,
|
|
IssueExecutionWorkspaceSettings,
|
|
} from "@paperclipai/shared";
|
|
import { environmentService } from "./environments.js";
|
|
import {
|
|
environmentRuntimeService,
|
|
buildEnvironmentLeaseContext,
|
|
type EnvironmentRuntimeLeaseRecord,
|
|
type EnvironmentRuntimeService,
|
|
} from "./environment-runtime.js";
|
|
import { ENVIRONMENT_DRIVER_TRAITS } from "./environment-driver-traits.js";
|
|
import {
|
|
resolveEnvironmentExecutionTarget,
|
|
resolveEnvironmentExecutionTransport,
|
|
} from "./environment-execution-target.js";
|
|
import {
|
|
adapterExecutionTargetToRemoteSpec,
|
|
type AdapterExecutionTarget,
|
|
type AdapterRemoteExecutionSpec,
|
|
type AdapterWorkspaceRealization,
|
|
} from "@paperclipai/adapter-utils/execution-target";
|
|
import type { DuplexTelemetryRecorder } from "@paperclipai/adapter-utils/duplex-telemetry";
|
|
import type { DuplexAggregateByteLedger } from "@paperclipai/adapter-utils/duplex-aggregate-byte-ledger";
|
|
import { buildWorkspaceRealizationRequest } from "./workspace-realization.js";
|
|
import { executionWorkspaceService } from "./execution-workspaces.js";
|
|
import { logActivity } from "./activity-log.js";
|
|
import { logger } from "../middleware/logger.js";
|
|
import { parseObject } from "../adapters/utils.js";
|
|
import type { RealizedExecutionWorkspace } from "./workspace-runtime.js";
|
|
import type { PluginWorkerManager } from "./plugin-worker-manager.js";
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Error types
|
|
// ---------------------------------------------------------------------------
|
|
|
|
export type EnvironmentErrorCode =
|
|
| "environment_not_found"
|
|
| "environment_inactive"
|
|
| "unsupported_environment"
|
|
| "unsupported_adapter_environment"
|
|
| "probe_failed"
|
|
| "lease_acquire_failed"
|
|
| "workspace_realization_failed"
|
|
| "transport_resolution_failed"
|
|
| "lease_release_failed"
|
|
| "lease_cleanup_failed";
|
|
|
|
export class EnvironmentRunError extends Error {
|
|
code: EnvironmentErrorCode;
|
|
environmentId?: string;
|
|
driver?: string;
|
|
provider?: string;
|
|
cause?: unknown;
|
|
|
|
constructor(
|
|
code: EnvironmentErrorCode,
|
|
message: string,
|
|
details?: {
|
|
environmentId?: string;
|
|
driver?: string;
|
|
provider?: string;
|
|
cause?: unknown;
|
|
},
|
|
) {
|
|
super(message);
|
|
this.name = "EnvironmentRunError";
|
|
this.code = code;
|
|
this.environmentId = details?.environmentId;
|
|
this.driver = details?.driver;
|
|
this.provider = details?.provider;
|
|
this.cause = details?.cause;
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Orchestration result types
|
|
// ---------------------------------------------------------------------------
|
|
|
|
export interface EnvironmentAcquisitionResult {
|
|
environment: Environment;
|
|
lease: EnvironmentLease;
|
|
leaseContext: ReturnType<typeof buildEnvironmentLeaseContext>;
|
|
executionTransport: Record<string, unknown> | null;
|
|
}
|
|
|
|
export interface EnvironmentRealizationResult {
|
|
lease: EnvironmentLease;
|
|
workspaceRealization: Record<string, unknown>;
|
|
executionTarget: AdapterExecutionTarget | null;
|
|
remoteExecution: AdapterRemoteExecutionSpec | null;
|
|
persistedExecutionWorkspace: ExecutionWorkspace | null;
|
|
}
|
|
|
|
export interface EnvironmentReleaseResult {
|
|
released: EnvironmentRuntimeLeaseRecord[];
|
|
errors: Array<{ leaseId: string; error: unknown }>;
|
|
}
|
|
|
|
function firstNonEmptyLine(text: string | null | undefined): string | null {
|
|
if (!text) return null;
|
|
for (const rawLine of text.split(/\r?\n/)) {
|
|
const line = rawLine.trim();
|
|
if (line) return line;
|
|
}
|
|
return null;
|
|
}
|
|
|
|
function formatProvisionFailureDetail(result: {
|
|
exitCode: number | null;
|
|
signal?: string | null;
|
|
timedOut: boolean;
|
|
stdout: string;
|
|
stderr: string;
|
|
}): string {
|
|
if (result.timedOut) {
|
|
return "provision command timed out";
|
|
}
|
|
const signal = typeof result.signal === "string" && result.signal.trim().length > 0
|
|
? ` (signal ${result.signal.trim()})`
|
|
: "";
|
|
const detail = firstNonEmptyLine(result.stderr) ?? firstNonEmptyLine(result.stdout);
|
|
const status = `exit code ${result.exitCode ?? "null"}${signal}`;
|
|
return detail ? `${status}: ${detail}` : status;
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Service factory
|
|
// ---------------------------------------------------------------------------
|
|
|
|
export function environmentRunOrchestrator(
|
|
db: Db,
|
|
options: {
|
|
pluginWorkerManager?: PluginWorkerManager;
|
|
environmentRuntime?: EnvironmentRuntimeService;
|
|
/**
|
|
* The process-owned aggregate byte ledger for the sandbox duplex channel.
|
|
* The server root creates one ledger per host process and injects the same
|
|
* object here. The orchestrator stamps it onto the sandbox execution target,
|
|
* so one shared gauge bounds the aggregate retained bytes across all live
|
|
* duplex routes. Absent keeps the bridge inert for this seam.
|
|
*/
|
|
duplexAggregateByteLedger?: DuplexAggregateByteLedger | null;
|
|
} = {},
|
|
) {
|
|
const environmentsSvc = environmentService(db);
|
|
const executionWorkspacesSvc = executionWorkspaceService(db);
|
|
const environmentRuntime = options.environmentRuntime ?? environmentRuntimeService(db, {
|
|
pluginWorkerManager: options.pluginWorkerManager,
|
|
});
|
|
|
|
/**
|
|
* Resolve the selected environment for a run. The caller passes the concrete
|
|
* selected environment id plus the built-in local fallback id used to lazily
|
|
* ensure the local environment row exists.
|
|
*/
|
|
async function resolveEnvironment(input: {
|
|
companyId: string;
|
|
selectedEnvironmentId: string;
|
|
localEnvironmentId: string;
|
|
}): Promise<Environment> {
|
|
const environmentId = input.selectedEnvironmentId || input.localEnvironmentId;
|
|
|
|
const environment =
|
|
environmentId === input.localEnvironmentId
|
|
? await environmentsSvc.ensureLocalEnvironment(input.companyId)
|
|
: await environmentsSvc.getById(environmentId);
|
|
|
|
if (!environment) {
|
|
throw new EnvironmentRunError("environment_not_found", `Environment "${environmentId}" not found.`, {
|
|
environmentId,
|
|
});
|
|
}
|
|
|
|
if (environment.status !== "active") {
|
|
throw new EnvironmentRunError("environment_inactive", `Environment "${environment.name}" is not active (status: ${environment.status}).`, {
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
});
|
|
}
|
|
|
|
return environment;
|
|
}
|
|
|
|
/**
|
|
* Acquire an environment lease for a heartbeat run.
|
|
* Wraps the runtime driver's acquire call with standardized error handling.
|
|
*/
|
|
async function acquireLease(input: {
|
|
companyId: string;
|
|
environment: Environment;
|
|
issueId: string | null;
|
|
agentId: string;
|
|
heartbeatRunId: string;
|
|
persistedExecutionWorkspace: Pick<ExecutionWorkspace, "id" | "mode"> | null;
|
|
executionWorkspaceSettings: IssueExecutionWorkspaceSettings | null;
|
|
adapterType: string | null;
|
|
}): Promise<EnvironmentRuntimeLeaseRecord> {
|
|
try {
|
|
return await environmentRuntime.acquireRunLease(input);
|
|
} catch (err) {
|
|
throw new EnvironmentRunError(
|
|
"lease_acquire_failed",
|
|
`Failed to acquire lease for environment "${input.environment.name}" (${input.environment.driver}): ${err instanceof Error ? err.message : String(err)}`,
|
|
{
|
|
environmentId: input.environment.id,
|
|
driver: input.environment.driver,
|
|
cause: err,
|
|
},
|
|
);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Resolve the execution transport for an adapter based on the acquired lease.
|
|
*/
|
|
async function resolveTransport(input: {
|
|
companyId: string;
|
|
adapterType: string;
|
|
environment: Environment;
|
|
leaseMetadata: Record<string, unknown> | null;
|
|
}): Promise<Record<string, unknown> | null> {
|
|
try {
|
|
return await resolveEnvironmentExecutionTransport({
|
|
db,
|
|
companyId: input.companyId,
|
|
adapterType: input.adapterType,
|
|
environment: input.environment,
|
|
leaseMetadata: input.leaseMetadata,
|
|
});
|
|
} catch (err) {
|
|
throw new EnvironmentRunError(
|
|
"transport_resolution_failed",
|
|
`Failed to resolve execution transport for "${input.environment.name}": ${err instanceof Error ? err.message : String(err)}`,
|
|
{
|
|
environmentId: input.environment.id,
|
|
driver: input.environment.driver,
|
|
cause: err,
|
|
},
|
|
);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Full acquisition flow: resolve environment, acquire lease, resolve transport.
|
|
* This is the primary entry point for heartbeat run setup.
|
|
*/
|
|
async function acquireForRun(input: {
|
|
companyId: string;
|
|
selectedEnvironmentId: string;
|
|
localEnvironmentId: string;
|
|
adapterType: string;
|
|
issueId: string | null;
|
|
heartbeatRunId: string;
|
|
agentId: string;
|
|
persistedExecutionWorkspace: Pick<ExecutionWorkspace, "id" | "mode"> | null;
|
|
executionWorkspaceSettings: IssueExecutionWorkspaceSettings | null;
|
|
}): Promise<EnvironmentAcquisitionResult> {
|
|
// Step 1: Resolve environment
|
|
const environment = await resolveEnvironment({
|
|
companyId: input.companyId,
|
|
selectedEnvironmentId: input.selectedEnvironmentId,
|
|
localEnvironmentId: input.localEnvironmentId,
|
|
});
|
|
|
|
// Step 2: Acquire lease
|
|
const leaseRecord = await acquireLease({
|
|
companyId: input.companyId,
|
|
environment,
|
|
issueId: input.issueId,
|
|
agentId: input.agentId,
|
|
heartbeatRunId: input.heartbeatRunId,
|
|
persistedExecutionWorkspace: input.persistedExecutionWorkspace,
|
|
executionWorkspaceSettings: input.executionWorkspaceSettings,
|
|
adapterType: input.adapterType ?? null,
|
|
});
|
|
|
|
// Step 3: Log lease acquisition activity
|
|
await logActivity(db, {
|
|
companyId: input.companyId,
|
|
actorType: "agent",
|
|
actorId: input.agentId,
|
|
agentId: input.agentId,
|
|
runId: input.heartbeatRunId,
|
|
action: "environment.lease_acquired",
|
|
entityType: "environment_lease",
|
|
entityId: leaseRecord.lease.id,
|
|
issueId: input.issueId,
|
|
details: {
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
leasePolicy: leaseRecord.lease.leasePolicy,
|
|
provider: leaseRecord.lease.provider,
|
|
executionWorkspaceId: leaseRecord.leaseContext.executionWorkspaceId,
|
|
issueId: input.issueId,
|
|
networkEgress: input.executionWorkspaceSettings?.networkEgress ?? null,
|
|
},
|
|
});
|
|
|
|
// Step 4: Resolve execution transport
|
|
const executionTransport = await resolveTransport({
|
|
companyId: input.companyId,
|
|
adapterType: input.adapterType,
|
|
environment,
|
|
leaseMetadata: leaseRecord.lease.metadata,
|
|
});
|
|
|
|
return {
|
|
environment,
|
|
lease: leaseRecord.lease,
|
|
leaseContext: leaseRecord.leaseContext,
|
|
executionTransport,
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Realize workspace in the environment and resolve the execution target.
|
|
*
|
|
* After lease acquisition, this method:
|
|
* 1. Builds a workspace realization request
|
|
* 2. Calls the environment runtime driver to realize the workspace
|
|
* 3. Persists realization metadata on the lease and execution workspace
|
|
* 4. Resolves the adapter execution target (local/ssh/sandbox)
|
|
*
|
|
* Returns the updated lease, realization metadata, and the execution
|
|
* target spec that the adapter needs to run.
|
|
*/
|
|
async function realizeForRun(input: {
|
|
environment: Environment;
|
|
lease: EnvironmentLease;
|
|
adapterType: string;
|
|
companyId: string;
|
|
issueId: string | null;
|
|
heartbeatRunId: string;
|
|
executionWorkspace: RealizedExecutionWorkspace;
|
|
effectiveExecutionWorkspaceMode: string | null;
|
|
persistedExecutionWorkspace: ExecutionWorkspace | null;
|
|
/**
|
|
* The host duplex telemetry recorder for this run. The orchestrator threads
|
|
* it to `resolveEnvironmentExecutionTarget`, which stamps it on the sandbox
|
|
* target. Absent keeps the safe no-op default in the bridge.
|
|
*/
|
|
duplexTelemetryRecorder?: DuplexTelemetryRecorder | null;
|
|
}): Promise<EnvironmentRealizationResult> {
|
|
const {
|
|
environment,
|
|
adapterType,
|
|
companyId,
|
|
issueId,
|
|
heartbeatRunId,
|
|
executionWorkspace,
|
|
effectiveExecutionWorkspaceMode,
|
|
} = input;
|
|
let { lease, persistedExecutionWorkspace } = input;
|
|
|
|
// Step 1: Build workspace realization request
|
|
const workspaceRealizationRequest = buildWorkspaceRealizationRequest({
|
|
adapterType,
|
|
companyId,
|
|
environmentId: environment.id,
|
|
executionWorkspaceId: persistedExecutionWorkspace?.id ?? null,
|
|
issueId,
|
|
heartbeatRunId,
|
|
requestedMode: persistedExecutionWorkspace?.mode ?? effectiveExecutionWorkspaceMode,
|
|
workspace: executionWorkspace,
|
|
workspaceConfig: persistedExecutionWorkspace?.config ?? null,
|
|
});
|
|
|
|
// Step 2: Realize workspace in the environment via the runtime driver
|
|
let workspaceRealization: Record<string, unknown> = {};
|
|
let realizedWorkspaceCwd: string | null = null;
|
|
if (ENVIRONMENT_DRIVER_TRAITS[environment.driver].realizesWorkspace) {
|
|
try {
|
|
const remoteCwd =
|
|
typeof lease.metadata?.remoteCwd === "string" && lease.metadata.remoteCwd.trim().length > 0
|
|
? lease.metadata.remoteCwd
|
|
: undefined;
|
|
const workspaceRealizationResult = await environmentRuntime.realizeWorkspace({
|
|
environment,
|
|
lease,
|
|
workspace: {
|
|
localPath: executionWorkspace.cwd,
|
|
remotePath: remoteCwd,
|
|
mode: persistedExecutionWorkspace?.mode ?? effectiveExecutionWorkspaceMode ?? undefined,
|
|
metadata: {
|
|
workspaceRealizationRequest,
|
|
},
|
|
},
|
|
});
|
|
realizedWorkspaceCwd =
|
|
typeof workspaceRealizationResult.cwd === "string" && workspaceRealizationResult.cwd.trim().length > 0
|
|
? workspaceRealizationResult.cwd.trim()
|
|
: null;
|
|
workspaceRealization = parseObject(workspaceRealizationResult.metadata?.workspaceRealization);
|
|
} catch (err) {
|
|
throw new EnvironmentRunError(
|
|
"workspace_realization_failed",
|
|
`Failed to realize workspace for environment "${environment.name}" (${environment.driver}): ${err instanceof Error ? err.message : String(err)}`,
|
|
{
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
cause: err,
|
|
},
|
|
);
|
|
}
|
|
}
|
|
|
|
const provisionCommand = workspaceRealizationRequest.runtimeOverlay.provisionCommand?.trim() ?? "";
|
|
const realizedCwd =
|
|
realizedWorkspaceCwd ??
|
|
(typeof lease.metadata?.remoteCwd === "string" && lease.metadata.remoteCwd.trim().length > 0
|
|
? lease.metadata.remoteCwd.trim()
|
|
: executionWorkspace.cwd);
|
|
// The host `provisionCommand` runs on the host worktree during the
|
|
// `workspace_provision` step, before the run reaches the environment.
|
|
// A `sandbox`-driver environment does not receive the repo tree here. The
|
|
// sandbox driver `realizeWorkspace` step only creates the remote folder.
|
|
// The adapter uploads the provisioned tree later, in its `stage.sync` step.
|
|
// So the host command must not run inside the still-empty sandbox; it fails
|
|
// there (exit 127). Skip the step for `sandbox`, and keep the existing skip
|
|
// for `local`. Keep the step for `ssh`, which runs the command on the
|
|
// remote host that shares the workspace path.
|
|
const driverSkipsHostProvision =
|
|
environment.driver === "local" || environment.driver === "sandbox";
|
|
if (provisionCommand && environment.driver === "sandbox") {
|
|
logger.info(
|
|
{
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
issueId,
|
|
heartbeatRunId,
|
|
},
|
|
"Skip host provisionCommand for sandbox-driver environment; the adapter stage.sync step delivers the provisioned tree",
|
|
);
|
|
}
|
|
if (provisionCommand && !driverSkipsHostProvision) {
|
|
try {
|
|
const provisionResult = await environmentRuntime.execute({
|
|
environment,
|
|
lease,
|
|
command: "bash",
|
|
args: ["-lc", provisionCommand],
|
|
cwd: realizedCwd,
|
|
env: {
|
|
SHELL: "/bin/bash",
|
|
},
|
|
timeoutMs: 300_000,
|
|
// The provision command runs before the run opens its trace root, so it
|
|
// carries no run parent. A sandbox provider that opens a persistent
|
|
// session on the first command must not open the session here, or the
|
|
// session-setup span loses its parent and the span backend drops it.
|
|
// Bypass the session for this command; the session opens on the first
|
|
// in-run command instead, whose setup span parents to the run trace.
|
|
bypassSession: true,
|
|
});
|
|
if (provisionResult.exitCode !== 0 || provisionResult.timedOut) {
|
|
throw new Error(formatProvisionFailureDetail(provisionResult));
|
|
}
|
|
} catch (err) {
|
|
throw new EnvironmentRunError(
|
|
"workspace_realization_failed",
|
|
`Failed to provision workspace for environment "${environment.name}" (${environment.driver}): ${err instanceof Error ? err.message : String(err)}`,
|
|
{
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
cause: err,
|
|
},
|
|
);
|
|
}
|
|
}
|
|
|
|
// Step 3: Persist realization metadata on lease and execution workspace
|
|
if (Object.keys(workspaceRealization).length > 0) {
|
|
const nextLeaseMetadata = {
|
|
...(lease.metadata ?? {}),
|
|
workspaceRealization,
|
|
};
|
|
const updatedLease = await environmentsSvc.updateLeaseMetadata(lease.id, nextLeaseMetadata);
|
|
if (updatedLease) {
|
|
lease = updatedLease;
|
|
}
|
|
if (persistedExecutionWorkspace) {
|
|
const updatedEw = await executionWorkspacesSvc.update(persistedExecutionWorkspace.id, {
|
|
metadata: {
|
|
...(persistedExecutionWorkspace.metadata ?? {}),
|
|
workspaceRealizationRequest,
|
|
workspaceRealization,
|
|
},
|
|
});
|
|
if (updatedEw) {
|
|
persistedExecutionWorkspace = updatedEw;
|
|
}
|
|
}
|
|
}
|
|
|
|
// Step 4: Resolve execution target for the adapter
|
|
let executionTarget: AdapterExecutionTarget | null;
|
|
try {
|
|
executionTarget = await resolveEnvironmentExecutionTarget({
|
|
db,
|
|
companyId,
|
|
adapterType,
|
|
environment,
|
|
leaseId: lease.id,
|
|
leaseMetadata: (lease.metadata as Record<string, unknown> | null) ?? null,
|
|
lease,
|
|
environmentRuntime,
|
|
duplexTelemetryRecorder: input.duplexTelemetryRecorder ?? null,
|
|
duplexAggregateByteLedger: options.duplexAggregateByteLedger ?? null,
|
|
});
|
|
const realizationMode = workspaceRealization.mode === "in_place" ? "in_place" : "copy";
|
|
const authoritativeRoot =
|
|
typeof workspaceRealization.authoritativeRoot === "string" && workspaceRealization.authoritativeRoot.trim().length > 0
|
|
? workspaceRealization.authoritativeRoot.trim()
|
|
: realizedCwd;
|
|
const workspaceTargetMetadata: AdapterWorkspaceRealization = {
|
|
mode: realizationMode,
|
|
authoritativeRoot,
|
|
pathAliases: Array.isArray(workspaceRealization.pathAliases)
|
|
? workspaceRealization.pathAliases.filter(
|
|
(entry): entry is { path: string; target: string } =>
|
|
typeof entry === "object" && entry !== null &&
|
|
typeof (entry as { path?: unknown }).path === "string" &&
|
|
typeof (entry as { target?: unknown }).target === "string",
|
|
)
|
|
: [],
|
|
outboundRestorePaths: Array.isArray(workspaceRealization.outboundRestorePaths)
|
|
? workspaceRealization.outboundRestorePaths.filter((entry): entry is string => typeof entry === "string")
|
|
: [],
|
|
};
|
|
if (executionTarget) {
|
|
executionTarget = {
|
|
...executionTarget,
|
|
...(executionTarget.kind === "remote" && realizationMode === "in_place"
|
|
? { remoteCwd: authoritativeRoot }
|
|
: {}),
|
|
workspaceRealization: workspaceTargetMetadata,
|
|
} as AdapterExecutionTarget;
|
|
}
|
|
} catch (err) {
|
|
throw new EnvironmentRunError(
|
|
"transport_resolution_failed",
|
|
`Failed to resolve execution target for "${environment.name}": ${err instanceof Error ? err.message : String(err)}`,
|
|
{
|
|
environmentId: environment.id,
|
|
driver: environment.driver,
|
|
cause: err,
|
|
},
|
|
);
|
|
}
|
|
|
|
return {
|
|
lease,
|
|
workspaceRealization,
|
|
executionTarget,
|
|
remoteExecution: adapterExecutionTargetToRemoteSpec(executionTarget),
|
|
persistedExecutionWorkspace,
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Release all active leases for a heartbeat run.
|
|
* Tracks cleanup status per lease. Errors during individual lease release
|
|
* are captured but do not prevent other leases from being released.
|
|
* The original run failure (if any) is never hidden by cleanup errors.
|
|
*/
|
|
async function releaseForRun(input: {
|
|
heartbeatRunId: string;
|
|
companyId: string;
|
|
agentId: string;
|
|
status?: Extract<EnvironmentLeaseStatus, "released" | "expired" | "failed">;
|
|
failureReason?: string;
|
|
}): Promise<EnvironmentReleaseResult> {
|
|
const status = input.status ?? "released";
|
|
const result: EnvironmentReleaseResult = { released: [], errors: [] };
|
|
|
|
let releasedLeases: EnvironmentRuntimeLeaseRecord[];
|
|
try {
|
|
releasedLeases = await environmentRuntime.releaseRunLeases(
|
|
input.heartbeatRunId,
|
|
status,
|
|
(leaseId, error) => result.errors.push({ leaseId, error }),
|
|
);
|
|
} catch (err) {
|
|
result.errors.push({ leaseId: "*", error: err });
|
|
return result;
|
|
}
|
|
|
|
for (const released of releasedLeases) {
|
|
try {
|
|
await logActivity(db, {
|
|
companyId: input.companyId,
|
|
actorType: "agent",
|
|
actorId: input.agentId,
|
|
agentId: input.agentId,
|
|
runId: input.heartbeatRunId,
|
|
action: "environment.lease_released",
|
|
entityType: "environment_lease",
|
|
entityId: released.lease.id,
|
|
issueId: released.lease.issueId,
|
|
details: {
|
|
environmentId: released.lease.environmentId,
|
|
driver: released.environment.driver,
|
|
leasePolicy: released.lease.leasePolicy,
|
|
provider: released.lease.provider,
|
|
executionWorkspaceId: released.lease.executionWorkspaceId,
|
|
issueId: released.lease.issueId,
|
|
status: released.lease.status,
|
|
cleanupStatus: released.lease.cleanupStatus,
|
|
failureReason: input.failureReason ?? released.lease.failureReason,
|
|
},
|
|
});
|
|
} catch {
|
|
// Activity logging failure should not block lease release
|
|
}
|
|
result.released.push(released);
|
|
}
|
|
|
|
return result;
|
|
}
|
|
|
|
return {
|
|
resolveEnvironment,
|
|
acquireLease,
|
|
resolveTransport,
|
|
acquireForRun,
|
|
realizeForRun,
|
|
releaseForRun,
|
|
|
|
// Expose the underlying runtime for cases that need direct driver access
|
|
runtime: environmentRuntime,
|
|
};
|
|
}
|
|
|
|
export type EnvironmentRunOrchestrator = ReturnType<typeof environmentRunOrchestrator>;
|