diff --git a/doc/plugins/ENVIRONMENT_TASKS.md b/doc/plugins/ENVIRONMENT_TASKS.md new file mode 100644 index 0000000000..033f78521b --- /dev/null +++ b/doc/plugins/ENVIRONMENT_TASKS.md @@ -0,0 +1,61 @@ +# Typed environment tasks + +An environment driver can own task admission without exposing a shell or creating +one disposable resource per run. Declare `supportsTasks: true` on the driver and +implement `onEnvironmentTask`. The worker advertises `environmentTask` during +initialization. Both declarations must be present. The existing +`environment.drivers.register` capability applies. + +The server-only `environmentRuntime.task({companyId, leaseId, operation})` method +loads the persisted lease and its run. It resolves the exact plugin recorded at +acquisition. It supplies the agent, issue, and project context. It does not accept +these identities or the plugin identity from the operation. An environment edit +cannot redirect an existing task to a replacement provider. The caller must +already authorize access to the company. + +This contract supplies execution capabilities. It does not select a provider for +heartbeat runs, stage runtime assets, or replace native Runner startup. Consumers +must integrate those steps before enabling task execution. Readiness checks remain +separate. No browser endpoint or automatic provider selection is added. + +## Lifetime and operations + +Acquire the environment lease before submitting a task. The provider lease ID is +the task's durable attempt ID. Save it before a remote call. Do not use a persistent +machine ID as this task ID. Multiple attempts can use the same underlying resource. + +- `submit`: supplies typed Runner identity, source revision, harness, optional + outbound WSS URL, and a transient bootstrap ticket. The Runner run and lease IDs + must match the persisted host records. Only a running run with an active, + unexpired lease can submit. +- `status`: returns the phase and optional exit code. Optional `executionStopped` + is live provider evidence that the complete task process tree has stopped. +- `connection`: returns a private authenticated WebSocket endpoint for a running + task. Transport credentials remain on the server. PRP authentication still binds + the Runner to its authorized run. +- `complete`: releases task grants according to the provider's policy. It does not + imply that the process or its descendants stopped. +- `stop`: requests cancellation of the task. An accepted request is not a process + termination receipt. Reconcile status before treating execution as stopped. + +Submission, completion and stop return `accepted`. Acceptance does not imply +readiness or success. Status and connection return distinct typed results. Every +result echoes the task ID and is checked by both the worker SDK and the host. + +## Retry and credentials + +A timeout may follow successful admission. Retain the same task ID and launch +configuration, inspect status, and retry the exact request when needed. Never +allocate a second attempt merely because the response was lost. Providers must +reject conflicting reuse and deduplicate identical submissions, including after a +terminal outcome. A new execution attempt requires a new lease and task ID. + +Do not store bootstrap tickets, connection headers, or provider credentials in +lease metadata, plugin state, run profiles, logs, or browser responses. The host +sanitizes worker errors. Treat a connection result as credential-bearing material. +An endpoint is transport access, not a replacement for the Runner protocol's +identity and artifact verification. + +Providers own resource-specific credential lookup, mounts, task status, and cleanup. +Task lease cleanup must never destroy a longer-lived resource as an implicit +fallback. Unsupported operations and unavailable providers fail closed. diff --git a/doc/plugins/PLUGIN_AUTHORING_GUIDE.md b/doc/plugins/PLUGIN_AUTHORING_GUIDE.md index a3728720e1..26b5e0926f 100644 --- a/doc/plugins/PLUGIN_AUTHORING_GUIDE.md +++ b/doc/plugins/PLUGIN_AUTHORING_GUIDE.md @@ -679,3 +679,9 @@ The host resets plugin state on account/company changes and keeps its built-in menu when no unique contribution exists, discovery fails, the module is missing, or rendering throws. The slot is a React-only contract; do not use a custom element export. This replaces only the menu, not company policy or authorization. + +## Typed task admission + +Environment drivers can implement [typed task operations](ENVIRONMENT_TASKS.md) for +provider-owned Runner launch, status, connections, completion, and cancellation. +This contract is separate from shell execution and readiness checks. diff --git a/packages/plugins/sdk/src/define-plugin.ts b/packages/plugins/sdk/src/define-plugin.ts index 947a67cbb5..61fb69db79 100644 --- a/packages/plugins/sdk/src/define-plugin.ts +++ b/packages/plugins/sdk/src/define-plugin.ts @@ -1,3 +1,4 @@ +import type { PluginEnvironmentTaskParams, PluginEnvironmentTaskResult } from "./environment-tasks.js"; import type { AiConnectionRouterRequest, AiConnectionRouterResult } from "@paperclipai/shared"; /** * `definePlugin` — the top-level helper for authoring a Paperclip plugin. @@ -371,6 +372,9 @@ export interface PluginDefinition { params: PluginEnvironmentValidateConfigParams, ): Promise; + /** Admit, inspect, connect to, or finish a typed task on a persisted lease. */ + onEnvironmentTask?(params: PluginEnvironmentTaskParams): Promise; + /** Called to test reachability or readiness of a plugin-hosted environment. */ onEnvironmentProbe?( params: PluginEnvironmentProbeParams, diff --git a/packages/plugins/sdk/src/environment-tasks.ts b/packages/plugins/sdk/src/environment-tasks.ts new file mode 100644 index 0000000000..dedd75765c --- /dev/null +++ b/packages/plugins/sdk/src/environment-tasks.ts @@ -0,0 +1,70 @@ +import { z } from "zod"; +import type { PluginEnvironmentDriverBaseParams, PluginEnvironmentLease } from "./protocol.js"; + +const identifier = z.string().regex(/^[a-zA-Z0-9][a-zA-Z0-9_-]{0,95}$/); +const secureUrl = z.string().url().refine(value => { + const url = new URL(value); + return url.protocol === "wss:" && !url.username && !url.password && !url.search && !url.hash; +}, "Expected a credential-free WSS URL"); + +/** One durable execution attempt. Secrets are transient RPC input, never lease metadata. */ +export const environmentTaskOperationSchema = z.discriminatedUnion("kind", [ + z.object({ + kind: z.literal("submit"), + runner: z.object({ + revision: z.string().regex(/^[a-f0-9]{40}$/), + harness: identifier, + runnerId: identifier, leaseId: identifier, runId: identifier, + sessionId: identifier, turnId: identifier, itemId: identifier, + connectUrl: secureUrl.optional(), + }).strict(), + bootstrapTicket: z.string().min(1).max(65_536), + }).strict(), + z.object({ kind: z.literal("status") }).strict(), + z.object({ kind: z.literal("connection") }).strict(), + z.object({ kind: z.literal("complete") }).strict(), + z.object({ kind: z.literal("stop") }).strict(), +]); + +export const environmentTaskResultSchema = z.discriminatedUnion("kind", [ + z.object({ kind: z.literal("accepted"), taskId: identifier }).strict(), + z.object({ + kind: z.literal("status"), taskId: identifier, + phase: z.enum(["preparing", "running", "completed", "failed", "cancelled", "interrupted"]), + exitCode: z.number().int().optional(), + /** Provider observed the task and all descendants stopped; terminal phase alone is insufficient. */ + executionStopped: z.boolean().optional(), + }).strict(), + z.object({ + kind: z.literal("connection"), taskId: identifier, + endpoint: z.object({ + kind: z.literal("authenticated_websocket"), websocketUrl: secureUrl, + generation: identifier, + secretHeaders: z.array(z.object({ + name: z.string().regex(/^[!#$%&'*+.^_`|~0-9A-Za-z-]+$/), + value: z.string().min(1).max(65_536).regex(/^[^\r\n]+$/), + }).strict()).max(16), + }).strict(), + }).strict(), +]); + +export type PluginEnvironmentTaskOperation = z.infer; +export type PluginEnvironmentTaskResult = z.infer; + +export interface PluginEnvironmentTaskParams extends PluginEnvironmentDriverBaseParams { + lease: PluginEnvironmentLease; + /** Provider-issued task identifier persisted in the lease, stable across ambiguous submission retries. */ + taskId: string; + runId: string; + agentId: string; + projectId: string | null; + operation: PluginEnvironmentTaskOperation; +} + +/** An accepted stop is not proof of process termination. */ +export function parseEnvironmentTaskResult(operation: PluginEnvironmentTaskOperation, taskId: string, value: unknown): PluginEnvironmentTaskResult { + const result = environmentTaskResultSchema.parse(value); + const expected = operation.kind === "status" || operation.kind === "connection" ? operation.kind : "accepted"; + if (result.taskId !== taskId || result.kind !== expected) throw new Error("Invalid environment task response"); + return result; +} diff --git a/packages/plugins/sdk/src/index.ts b/packages/plugins/sdk/src/index.ts index b834be1a60..4911149d91 100644 --- a/packages/plugins/sdk/src/index.ts +++ b/packages/plugins/sdk/src/index.ts @@ -454,3 +454,6 @@ export { PluginEnvironmentCreationCleanupError, environmentCreationCleanupErrorD export type { PluginEnvironmentCreationCleanup } from "./environment-creation-cleanup.js"; export type { AiConnectionPool, AiConnectionPoolConfig, AiConnectionPoolMember, AiConnectionRouterRequest, AiConnectionRouterResult, AiConnectionRouterSelection } from "@paperclipai/shared"; + +export { environmentTaskOperationSchema, environmentTaskResultSchema, parseEnvironmentTaskResult } from "./environment-tasks.js"; +export type { PluginEnvironmentTaskOperation, PluginEnvironmentTaskParams, PluginEnvironmentTaskResult } from "./environment-tasks.js"; diff --git a/packages/plugins/sdk/src/protocol.ts b/packages/plugins/sdk/src/protocol.ts index 76f0cd4228..c742e2e859 100644 --- a/packages/plugins/sdk/src/protocol.ts +++ b/packages/plugins/sdk/src/protocol.ts @@ -1,3 +1,4 @@ +import type { PluginEnvironmentTaskParams, PluginEnvironmentTaskResult } from "./environment-tasks.js"; import type { AiConnectionRouterRequest, AiConnectionRouterResult } from "@paperclipai/shared"; /** * JSON-RPC 2.0 message types and protocol helpers for the host ↔ worker IPC @@ -1365,6 +1366,7 @@ export interface HostToWorkerMethods { params: PluginEnvironmentValidateConfigParams, result: PluginEnvironmentValidationResult, ]; + environmentTask: [params: PluginEnvironmentTaskParams, result: PluginEnvironmentTaskResult]; environmentProbe: [ params: PluginEnvironmentProbeParams, result: PluginEnvironmentProbeResult, @@ -1485,6 +1487,7 @@ export const HOST_TO_WORKER_OPTIONAL_METHODS: readonly HostToWorkerMethodName[] "resolveExternalObject", "refreshExternalObjects", "environmentValidateConfig", + "environmentTask", "environmentProbe", "environmentAcquireLease", "environmentResumeLease", diff --git a/packages/plugins/sdk/src/worker-rpc-host.ts b/packages/plugins/sdk/src/worker-rpc-host.ts index 6684f8fdd1..8becacd273 100644 --- a/packages/plugins/sdk/src/worker-rpc-host.ts +++ b/packages/plugins/sdk/src/worker-rpc-host.ts @@ -1,3 +1,4 @@ +import { environmentTaskOperationSchema, parseEnvironmentTaskResult, type PluginEnvironmentTaskParams } from "./environment-tasks.js"; import type { AiConnectionRouterRequest, AiConnectionRouterResult } from "@paperclipai/shared"; import { environmentCreationCleanupErrorData } from "./environment-creation-cleanup.js"; /** @@ -1638,6 +1639,13 @@ export function startWorkerRpcHost(options: WorkerRpcHostOptions): WorkerRpcHost case "environmentValidateConfig": return handleEnvironmentValidateConfig(params as PluginEnvironmentValidateConfigParams); + case "environmentTask": { + if (!plugin.definition.onEnvironmentTask) throw methodNotImplemented("environmentTask"); + const input = params as PluginEnvironmentTaskParams; + const operation = environmentTaskOperationSchema.parse(input.operation); + return parseEnvironmentTaskResult(operation, input.taskId, await plugin.definition.onEnvironmentTask({ ...input, operation })); + } + case "environmentProbe": return handleEnvironmentProbe(params as PluginEnvironmentProbeParams); @@ -1750,6 +1758,7 @@ export function startWorkerRpcHost(options: WorkerRpcHostOptions): WorkerRpcHost if (plugin.definition.onResolveExternalObject) supportedMethods.push("resolveExternalObject"); if (plugin.definition.onRefreshExternalObjects) supportedMethods.push("refreshExternalObjects"); if (plugin.definition.onEnvironmentValidateConfig) supportedMethods.push("environmentValidateConfig"); + if (plugin.definition.onEnvironmentTask) supportedMethods.push("environmentTask"); if (plugin.definition.onEnvironmentProbe) supportedMethods.push("environmentProbe"); if (plugin.definition.onEnvironmentAcquireLease) supportedMethods.push("environmentAcquireLease"); if (plugin.definition.onEnvironmentResumeLease) supportedMethods.push("environmentResumeLease"); diff --git a/packages/shared/src/types/plugin.ts b/packages/shared/src/types/plugin.ts index 42643d3b00..f3425667f1 100644 --- a/packages/shared/src/types/plugin.ts +++ b/packages/shared/src/types/plugin.ts @@ -191,6 +191,8 @@ export interface SandboxProviderCapabilities { } export interface PluginEnvironmentDriverDeclaration { + /** Implements the typed environmentTask worker RPC; no shell execution is implied. */ + supportsTasks?: boolean; /** Stable driver key, unique within the plugin. Namespaced by plugin ID at runtime. */ driverKey: string; /** diff --git a/packages/shared/src/validators/plugin.ts b/packages/shared/src/validators/plugin.ts index 0024977164..b936139bba 100644 --- a/packages/shared/src/validators/plugin.ts +++ b/packages/shared/src/validators/plugin.ts @@ -177,6 +177,7 @@ export type SandboxProviderCapabilitiesInput = z.infer ({ plugin: {} as any })); +vi.mock("../services/plugin-registry.js", () => ({ pluginRegistryService: () => ({ getById: vi.fn(async () => state.plugin) }) })); +const leaseId = "10000000-0000-4000-8000-000000000001"; +const row = () => ({ + lease: { id: leaseId, companyId: "company", environmentId: "environment", providerLeaseId: "attempt-1", heartbeatRunId: "run", issueId: "issue", status: "active", expiresAt: null, + metadata: { driver: "plugin", pluginId: "original", driverKey: "tasks" } }, + environment: { id: "environment", config: { pluginKey: "test.provider", driverKey: "tasks", driverConfig: {} } }, + run: { id: "run", agentId: "agent", status: "running" }, +}); +function database(value: ReturnType | null = row()) { + const results = [value ? [value] : [], [{ projectId: "project" }]]; + const query: any = { from: () => query, innerJoin: () => query, where: vi.fn(() => query), limit: () => Promise.resolve(results.shift()) }; + return { db: { select: () => query } as unknown as Db, query }; +} +function worker(result: unknown = { kind: "accepted", taskId: "attempt-1" }) { + return { getWorker: vi.fn(() => ({ supportedMethods: ["environmentTask"] })), call: vi.fn(async () => result) }; +} +const submit = { kind: "submit" as const, runner: { revision: "a".repeat(40), harness: "codex", runnerId: "runner", leaseId, runId: "run", sessionId: "session", turnId: "turn", itemId: "item" }, bootstrapTicket: "transient-test-ticket" }; + +beforeEach(() => { + state.plugin = { id: "original", pluginKey: "test.provider", status: "ready", manifestJson: { capabilities: ["environment.drivers.register"], environmentDrivers: [{ driverKey: "tasks", supportsTasks: true }] } }; +}); +describe("typed environment task admission", () => { + it("dispatches with host-derived scope and a persisted attempt identity", async () => { + const { db, query } = database(); const workers = worker(); + await expect(executeEnvironmentTask(db, workers as never, { companyId: "company", leaseId, operation: submit })).resolves.toEqual({ kind: "accepted", taskId: "attempt-1" }); + const scope = new PgDialect().sqlToQuery(query.where.mock.calls[0][0]); + expect(scope.sql).toContain('"environment_leases"."company_id" ='); + expect(scope.sql).toContain('"environment_leases"."id" ='); + expect(scope.params).toEqual([leaseId, "company"]); + expect(workers.call).toHaveBeenCalledWith("original", "environmentTask", expect.objectContaining({ + taskId: "attempt-1", companyId: "company", agentId: "agent", projectId: "project", runId: "run", operation: submit, + }), 15_000); + }); + it("rejects an absent or cross-company lease before calling a worker", async () => { + const workers = worker(); + await expect(executeEnvironmentTask(database(null).db, workers as never, { companyId: "other", leaseId, operation: submit })).rejects.toThrow("lease unavailable"); + expect(workers.call).not.toHaveBeenCalled(); + }); + it("rejects mismatched Runner binding and inactive submission", async () => { + const workers = worker(); + await expect(executeEnvironmentTask(database().db, workers as never, { companyId: "company", leaseId, operation: { ...submit, runner: { ...submit.runner, runId: "other" } } })).rejects.toThrow("identity mismatch"); + const value = row(); value.lease.status = "released"; + await expect(executeEnvironmentTask(database(value).db, workers as never, { companyId: "company", leaseId, operation: submit })).rejects.toThrow("not active"); + expect(workers.call).not.toHaveBeenCalled(); + }); + it("uses the pinned provider for cleanup after environment edits", async () => { + const value = row(); value.lease.status = "released"; value.environment.config.pluginKey = "replacement"; + const workers = worker(); + await executeEnvironmentTask(database(value).db, workers as never, { companyId: "company", leaseId, operation: { kind: "stop" } }); + expect(workers.call).toHaveBeenCalledWith("original", "environmentTask", expect.objectContaining({ config: {}, operation: { kind: "stop" } }), 15_000); + }); + it.each(["manifest", "worker"])("requires live %s support", async kind => { + const workers = worker(); + if (kind === "manifest") state.plugin.manifestJson.environmentDrivers[0].supportsTasks = false; + else workers.getWorker.mockReturnValue({ supportedMethods: [] }); + await expect(executeEnvironmentTask(database().db, workers as never, { companyId: "company", leaseId, operation: submit })).rejects.toThrow("provider unavailable"); + expect(workers.call).not.toHaveBeenCalled(); + }); + it("rejects wrong task receipts and sanitizes worker errors", async () => { + await expect(executeEnvironmentTask(database().db, worker({ kind: "accepted", taskId: "other" }) as never, { companyId: "company", leaseId, operation: submit })).rejects.toThrow("reconcile the same task"); + const workers = worker(); workers.call.mockRejectedValue(new Error("private-credential")); + await expect(executeEnvironmentTask(database().db, workers as never, { companyId: "company", leaseId, operation: submit })).rejects.toThrow(/^Environment task operation unavailable; reconcile the same task before retrying$/); + }); + it("validates operation and result shape without making acceptance mean readiness", () => { + expect(environmentTaskOperationSchema.safeParse({ ...submit, surprise: true }).success).toBe(false); + expect(environmentTaskOperationSchema.safeParse({ ...submit, runner: { ...submit.runner, connectUrl: "wss://user:secret@example.test/path" } }).success).toBe(false); + expect(() => parseEnvironmentTaskResult({ kind: "status" }, "attempt-1", { kind: "accepted", taskId: "attempt-1" })).toThrow(); + expect(parseEnvironmentTaskResult(submit, "attempt-1", { kind: "accepted", taskId: "attempt-1" }).kind).toBe("accepted"); + }); +}); diff --git a/server/src/__tests__/plugin-environment-driver-seam.test.ts b/server/src/__tests__/plugin-environment-driver-seam.test.ts index 56c13d091e..1025509ccc 100644 --- a/server/src/__tests__/plugin-environment-driver-seam.test.ts +++ b/server/src/__tests__/plugin-environment-driver-seam.test.ts @@ -94,6 +94,48 @@ describe("plugin environment driver seam", () => { }); }); + it("negotiates and validates typed task RPCs", async () => { + const calls: unknown[] = []; + const plugin = definePlugin({ + async setup() {}, + async onEnvironmentTask(input) { + calls.push(input); + return { kind: "accepted", taskId: input.taskId }; + }, + }); + const stdin = new PassThrough(); + const stdout = new PassThrough(); + const host = startWorkerRpcHost({ plugin, stdin, stdout }); + const responses: unknown[] = []; + stdout.on("data", chunk => { + for (const line of String(chunk).split("\n").filter(Boolean)) responses.push(parseMessage(line)); + }); + try { + const manifest = { ...baseManifest, environmentDrivers: [{ ...baseManifest.environmentDrivers![0], supportsTasks: true }] }; + expect(pluginManifestV1Schema.parse(manifest).environmentDrivers![0].supportsTasks).toBe(true); + stdin.write(serializeMessage(createRequest("initialize", { + manifest, config: {}, instanceInfo: { instanceId: "instance-1", hostVersion: "1.0.0" }, apiVersion: 1, + }, 1))); + await waitForResponses(responses, 1); + const initialized = responses[0]; + expect(isJsonRpcSuccessResponse(initialized)).toBe(true); + if (!isJsonRpcSuccessResponse(initialized)) return; + expect(initialized.result.supportedMethods).toContain("environmentTask"); + const params = { driverKey: "fake-plugin", companyId: "company", environmentId: "environment", config: {}, + taskId: "attempt", runId: "run", agentId: "agent", projectId: null, lease: { providerLeaseId: "attempt" }, operation: { kind: "stop" } }; + stdin.write(serializeMessage(createRequest("environmentTask", params, 2))); + await waitForResponses(responses, 2); + expect(responses[1]).toMatchObject({ result: { kind: "accepted", taskId: "attempt" } }); + stdin.write(serializeMessage(createRequest("environmentTask", { ...params, operation: { kind: "status" } }, 3))); + await waitForResponses(responses, 3); + expect(isJsonRpcErrorResponse(responses[2])).toBe(true); // Wrong response kind. + stdin.write(serializeMessage(createRequest("environmentTask", { ...params, operation: { kind: "shell", command: "echo test" } }, 4))); + await waitForResponses(responses, 4); + expect(isJsonRpcErrorResponse(responses[3])).toBe(true); + expect(calls).toHaveLength(2); // Invalid operation did not reach the plugin. + } finally { host.stop(); } + }); + it("dispatches environment driver worker hooks and reports support", async () => { const plugin = definePlugin({ async setup() {}, diff --git a/server/src/services/environment-runtime.ts b/server/src/services/environment-runtime.ts index 4105ec6642..5efe010e60 100644 --- a/server/src/services/environment-runtime.ts +++ b/server/src/services/environment-runtime.ts @@ -1,3 +1,5 @@ +import { executeEnvironmentTask } from "./environment-task-runtime.js"; +import type { PluginEnvironmentTaskOperation } from "@paperclipai/plugin-sdk"; import { hasStopOnlyCleanup, prepareSandboxStopAndRetain, readStopOnlyCleanup, settleStopOnlyCleanup, stopOnlyCleanupKey } from "./sandbox-stop-and-retain.js"; import { readEnvironmentCreationCleanupError } from "@paperclipai/plugin-sdk"; import { remoteTerminationReceipt } from "./remote-execution-termination.js"; @@ -3870,6 +3872,12 @@ export function environmentRuntimeService( return resolveSandboxDuplexBridgeInput(experimental); }, + /** Typed task providers own process admission; this does not execute a shell command. */ + async task(input: { companyId: string; leaseId: string; operation: PluginEnvironmentTaskOperation }) { + if (!options.pluginWorkerManager) throw new Error("Environment task worker manager unavailable"); + return executeEnvironmentTask(db, options.pluginWorkerManager, input); + }, + async acquireRunLease(input: { companyId: string; environment: Environment; diff --git a/server/src/services/environment-task-runtime.ts b/server/src/services/environment-task-runtime.ts new file mode 100644 index 0000000000..b18bdd2706 --- /dev/null +++ b/server/src/services/environment-task-runtime.ts @@ -0,0 +1,62 @@ +import { and, eq } from "drizzle-orm"; +import { environmentLeases, environments, heartbeatRuns, issues, type Db } from "@paperclipai/db"; +import { environmentTaskOperationSchema, parseEnvironmentTaskResult, type PluginEnvironmentTaskOperation } from "@paperclipai/plugin-sdk"; +import { pluginRegistryService } from "./plugin-registry.js"; +import type { PluginWorkerManager } from "./plugin-worker-manager.js"; + +/** Server-only dispatch. The caller supplies an authorized company and a persisted lease, + * never a plugin identity, resource endpoint, or agent/project authority. + * A timeout is ambiguous: retry the same operation against the same lease. + */ +export async function executeEnvironmentTask(db: Db, workers: PluginWorkerManager, input: { + companyId: string; + leaseId: string; + operation: PluginEnvironmentTaskOperation; +}) { + const operation = environmentTaskOperationSchema.parse(input.operation); + const [row] = await db.select({ lease: environmentLeases, environment: environments, run: heartbeatRuns }) + .from(environmentLeases) + .innerJoin(environments, eq(environments.id, environmentLeases.environmentId)) + .innerJoin(heartbeatRuns, and(eq(heartbeatRuns.id, environmentLeases.heartbeatRunId), eq(heartbeatRuns.companyId, environmentLeases.companyId))) + .where(and(eq(environmentLeases.id, input.leaseId), eq(environmentLeases.companyId, input.companyId))).limit(1); + if (!row || !row.lease.providerLeaseId) throw new Error("Environment task lease unavailable"); + const { lease, environment, run } = row; + const taskId = row.lease.providerLeaseId; + if ((operation.kind === "submit" || operation.kind === "connection") && + (lease.status !== "active" || (lease.expiresAt && lease.expiresAt.getTime() <= Date.now()) || + run.status !== "running")) throw new Error("Environment task lease is not active"); + if (operation.kind === "submit" && (operation.runner.runId !== run.id || operation.runner.leaseId !== lease.id)) { + throw new Error("Environment task runner identity mismatch"); + } + const metadata = lease.metadata ?? {}; + if (metadata.driver !== "plugin" || typeof metadata.pluginId !== "string" || typeof metadata.driverKey !== "string") { + throw new Error("Environment task lease has no pinned plugin driver"); + } + const plugin = await pluginRegistryService(db).getById(metadata.pluginId); + const driver = plugin?.manifestJson.environmentDrivers?.find(value => value.driverKey === metadata.driverKey); + if (!plugin || plugin.status !== "ready" || !driver?.supportsTasks || + !plugin.manifestJson.capabilities.includes("environment.drivers.register") || + !workers.getWorker(plugin.id)?.supportedMethods.includes("environmentTask")) { + throw new Error("Environment task provider unavailable"); + } + const project = lease.issueId ? await db.select({ projectId: issues.projectId }).from(issues) + .where(and(eq(issues.id, lease.issueId), eq(issues.companyId, input.companyId))).limit(1).then(rows => rows[0]) : null; + if (lease.issueId && !project) throw new Error("Environment task issue unavailable"); + // Provider identity comes from the lease. Editing the environment must never + // redirect status or cleanup to a replacement plugin. + const config = environment.config as Record; + const driverConfig = config.pluginKey === plugin.pluginKey && config.driverKey === metadata.driverKey + ? (config.driverConfig as Record | undefined) ?? {} : {}; + try { + const result = await workers.call(plugin.id, "environmentTask", { + driverKey: metadata.driverKey, companyId: input.companyId, environmentId: environment.id, + issueId: lease.issueId, config: driverConfig, + taskId, runId: run.id, agentId: run.agentId, projectId: project?.projectId ?? null, + lease: { providerLeaseId: lease.providerLeaseId, metadata: lease.metadata ?? undefined }, operation, + }, 15_000); + return parseEnvironmentTaskResult(operation, taskId, result); + } catch { + // Worker errors and schema diagnostics can contain request credentials. + throw new Error("Environment task operation unavailable; reconcile the same task before retrying"); + } +}