mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-08 21:03:51 +02:00
Add typed task operations to environment plugins
Provide a negotiated task RPC for environment providers that own Runner admission. Dispatch from persisted company-scoped leases and validate typed requests and receipts without coupling the host to a provider API. Cover the worker contract, host-derived identity, capability negotiation, and failed or mismatched receipts. Document retry and credential handling. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
13 files changed
+348
No files matched your search
@@ -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.
|
||||
@@ -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.
|
||||
@@ -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<PluginEnvironmentValidationResult>;
|
||||
|
||||
/** Admit, inspect, connect to, or finish a typed task on a persisted lease. */
|
||||
onEnvironmentTask?(params: PluginEnvironmentTaskParams): Promise<PluginEnvironmentTaskResult>;
|
||||
|
||||
/** Called to test reachability or readiness of a plugin-hosted environment. */
|
||||
onEnvironmentProbe?(
|
||||
params: PluginEnvironmentProbeParams,
|
||||
|
||||
@@ -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<typeof environmentTaskOperationSchema>;
|
||||
export type PluginEnvironmentTaskResult = z.infer<typeof environmentTaskResultSchema>;
|
||||
|
||||
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;
|
||||
}
|
||||
@@ -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";
|
||||
@@ -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",
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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;
|
||||
/**
|
||||
|
||||
@@ -177,6 +177,7 @@ export type SandboxProviderCapabilitiesInput = z.infer<typeof sandboxProviderCap
|
||||
// validation error, not a silently dropped field. A misspelled capability name
|
||||
// (for example `supportsLoginPTY`) fails validation instead of dropping.
|
||||
export const pluginEnvironmentDriverDeclarationSchema = z.object({
|
||||
supportsTasks: z.boolean().optional(),
|
||||
driverKey: z.string().min(1).regex(
|
||||
/^[a-z0-9][a-z0-9._-]*$/,
|
||||
"Environment driver key must start with a lowercase alphanumeric and contain only lowercase letters, digits, dots, hyphens, or underscores",
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { PgDialect } from "drizzle-orm/pg-core";
|
||||
import type { Db } from "@paperclipai/db";
|
||||
import { environmentTaskOperationSchema, parseEnvironmentTaskResult } from "@paperclipai/plugin-sdk";
|
||||
import { executeEnvironmentTask } from "../services/environment-task-runtime.js";
|
||||
|
||||
const state = vi.hoisted(() => ({ 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<typeof row> | 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");
|
||||
});
|
||||
});
|
||||
@@ -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() {},
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
const driverConfig = config.pluginKey === plugin.pluginKey && config.driverKey === metadata.driverKey
|
||||
? (config.driverConfig as Record<string, unknown> | 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");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user