mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
feat: show live cloud worker pool in settings (#157200)
* feat: show ready cloud worker pool in settings * test: complete cloud worker cleanup fixtures * test: wait for sidebar before plugin registration check * test: run composed gateway scenarios in source lanes * fix(gateway): opt in to live-authorized pool details * test: document native TypeScript shutdown noise workaround
This commit is contained in:
parent
2d42a1da81
commit
809808b76a
33 changed files with 1474 additions and 174 deletions
|
|
@ -6965,34 +6965,47 @@ public struct EnvironmentsListParams: Codable, Sendable {
|
|||
public let runtimeid: String?
|
||||
public let projection: String?
|
||||
public let includedesktopsetup: Bool?
|
||||
public let includeprepareddetails: Bool?
|
||||
|
||||
public init(
|
||||
runtimeid: String? = nil,
|
||||
projection: String? = nil,
|
||||
includedesktopsetup: Bool? = nil)
|
||||
includedesktopsetup: Bool? = nil,
|
||||
includeprepareddetails: Bool? = nil)
|
||||
{
|
||||
self.runtimeid = runtimeid
|
||||
self.projection = projection
|
||||
self.includedesktopsetup = includedesktopsetup
|
||||
self.includeprepareddetails = includeprepareddetails
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case runtimeid = "runtimeId"
|
||||
case projection
|
||||
case includedesktopsetup = "includeDesktopSetup"
|
||||
case includeprepareddetails = "includePreparedDetails"
|
||||
}
|
||||
}
|
||||
|
||||
public struct EnvironmentsListResult: Codable, Sendable {
|
||||
public let environments: [EnvironmentSummary]
|
||||
public let profiles: [[String: AnyCodable]]?
|
||||
public let preparedpool: [String: AnyCodable]?
|
||||
|
||||
public init(
|
||||
environments: [EnvironmentSummary],
|
||||
profiles: [[String: AnyCodable]]? = nil)
|
||||
profiles: [[String: AnyCodable]]? = nil,
|
||||
preparedpool: [String: AnyCodable]? = nil)
|
||||
{
|
||||
self.environments = environments
|
||||
self.profiles = profiles
|
||||
self.preparedpool = preparedpool
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case environments
|
||||
case profiles
|
||||
case preparedpool = "preparedPool"
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -7162,15 +7175,19 @@ public struct EnvironmentsSessionStatusParams: Codable, Sendable {
|
|||
|
||||
public struct EnvironmentsStatusParams: Codable, Sendable {
|
||||
public let environmentid: String
|
||||
public let includeprepareddetails: Bool?
|
||||
|
||||
public init(
|
||||
environmentid: String)
|
||||
environmentid: String,
|
||||
includeprepareddetails: Bool? = nil)
|
||||
{
|
||||
self.environmentid = environmentid
|
||||
self.includeprepareddetails = includeprepareddetails
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case environmentid = "environmentId"
|
||||
case includeprepareddetails = "includePreparedDetails"
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -25316,6 +25333,7 @@ public struct WorkerEnvironmentMetadata: Codable, Sendable {
|
|||
public let state: WorkerEnvironmentState
|
||||
public let agems: Int
|
||||
public let idlems: Int?
|
||||
public let destroyrequestedatms: Int?
|
||||
public let attachedsessionids: [String]
|
||||
public let tunnelstatus: WorkerTunnelStatus
|
||||
public let error: String?
|
||||
|
|
@ -25329,6 +25347,7 @@ public struct WorkerEnvironmentMetadata: Codable, Sendable {
|
|||
state: WorkerEnvironmentState,
|
||||
agems: Int,
|
||||
idlems: Int? = nil,
|
||||
destroyrequestedatms: Int? = nil,
|
||||
attachedsessionids: [String],
|
||||
tunnelstatus: WorkerTunnelStatus,
|
||||
error: String? = nil,
|
||||
|
|
@ -25341,6 +25360,7 @@ public struct WorkerEnvironmentMetadata: Codable, Sendable {
|
|||
self.state = state
|
||||
self.agems = agems
|
||||
self.idlems = idlems
|
||||
self.destroyrequestedatms = destroyrequestedatms
|
||||
self.attachedsessionids = attachedsessionids
|
||||
self.tunnelstatus = tunnelstatus
|
||||
self.error = error
|
||||
|
|
@ -25355,6 +25375,7 @@ public struct WorkerEnvironmentMetadata: Codable, Sendable {
|
|||
case state
|
||||
case agems = "ageMs"
|
||||
case idlems = "idleMs"
|
||||
case destroyrequestedatms = "destroyRequestedAtMs"
|
||||
case attachedsessionids = "attachedSessionIds"
|
||||
case tunnelstatus = "tunnelStatus"
|
||||
case error
|
||||
|
|
|
|||
|
|
@ -97,6 +97,19 @@ retirement unless pinned or still held by an allocation.
|
|||
|
||||
### Ready workers
|
||||
|
||||
Open **Settings → Connections → Cloud workers → Pool** to inspect running
|
||||
prepared workers. The view groups workers by profile and shows ready, preparing,
|
||||
releasing, and attention counts, the project and prepared commit, expiry, and
|
||||
recorded failures. Capacity includes preparation and unconfirmed cleanup;
|
||||
consumed workers leave this view unless pending cleanup still reserves capacity. The inventory
|
||||
refreshes every 10 seconds while the view is visible. A failed refresh keeps the
|
||||
last result visible with a warning. The **Profiles** tab controls reserve targets
|
||||
and the shared pool limit; **Snapshots** manages the reusable disk images.
|
||||
|
||||
Pool details require current administrator access. API clients request them with
|
||||
`includePreparedDetails: true` on `environments.list` or `environments.status`;
|
||||
default responses retain the existing inventory shape for older clients.
|
||||
|
||||
For an eligible local Git project or repository-only session, a successful session activation can prepare a
|
||||
dedicated worker for the next session in the background. The default target is
|
||||
one unassigned worker per project and profile, with a Gateway-wide cap of four.
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ import {
|
|||
validateEnvironmentsListParams,
|
||||
validateEnvironmentsPrepareParams,
|
||||
validateEnvironmentsPrepareResult,
|
||||
validateEnvironmentsStatusParams,
|
||||
validateWorkerDesktopLaunchParams,
|
||||
validateWorkerDesktopLaunchResult,
|
||||
WorkerEnvironmentStateSchema,
|
||||
|
|
@ -51,6 +52,20 @@ function workerSummary(
|
|||
}
|
||||
|
||||
describe("worker environment protocol schemas", () => {
|
||||
it("accepts only boolean opt-in for prepared details in list and status requests", () => {
|
||||
for (const includePreparedDetails of [undefined, false, true]) {
|
||||
const option = includePreparedDetails === undefined ? {} : { includePreparedDetails };
|
||||
expect(validateEnvironmentsListParams(option)).toBe(true);
|
||||
expect(validateEnvironmentsStatusParams({ environmentId: "worker-1", ...option })).toBe(true);
|
||||
}
|
||||
for (const includePreparedDetails of [null, "true", 1]) {
|
||||
expect(validateEnvironmentsListParams({ includePreparedDetails })).toBe(false);
|
||||
expect(
|
||||
validateEnvironmentsStatusParams({ environmentId: "worker-1", includePreparedDetails }),
|
||||
).toBe(false);
|
||||
}
|
||||
});
|
||||
|
||||
it("accepts opt-in desktop setup discovery with a closed credential-free result", () => {
|
||||
expect(validateEnvironmentsListParams({ includeDesktopSetup: true })).toBe(true);
|
||||
expect(validateEnvironmentsListParams({ includeDesktopSetup: false })).toBe(true);
|
||||
|
|
@ -79,7 +94,7 @@ describe("worker environment protocol schemas", () => {
|
|||
}
|
||||
});
|
||||
|
||||
it("allows only bounded readonly profile display IDs, never settings", () => {
|
||||
it("allows only bounded readonly profile metadata, never settings", () => {
|
||||
const check = (profile: Record<string, unknown>) =>
|
||||
Value.Check(EnvironmentsListResultSchema, {
|
||||
environments: [],
|
||||
|
|
@ -88,12 +103,40 @@ describe("worker environment protocol schemas", () => {
|
|||
expect(check({})).toBe(true);
|
||||
expect(check({ providerDisplayId: "aws" })).toBe(true);
|
||||
expect(check({ providerDisplayId: "google-cloud" })).toBe(true);
|
||||
expect(check({ readyWorkers: 0 })).toBe(true);
|
||||
expect(check({ readyWorkers: 2 })).toBe(true);
|
||||
for (const providerDisplayId of ["", "AWS", "aws\n", "a".repeat(65), "aws/token", 42, {}]) {
|
||||
expect(check({ providerDisplayId })).toBe(false);
|
||||
}
|
||||
for (const readyWorkers of [-1, 0.5, "2", null]) {
|
||||
expect(check({ readyWorkers })).toBe(false);
|
||||
}
|
||||
expect(check({ providerDisplayId: "aws", settings: { provider: "aws" } })).toBe(false);
|
||||
});
|
||||
|
||||
it("reports unique prepared reservation identities even above a reduced cap", () => {
|
||||
const check = (preparedPool: unknown) =>
|
||||
Value.Check(EnvironmentsListResultSchema, { environments: [], preparedPool });
|
||||
expect(Value.Check(EnvironmentsListResultSchema, { environments: [] })).toBe(true);
|
||||
expect(check({ maxTotal: 0, reservedEnvironmentIds: [] })).toBe(true);
|
||||
expect(check({ maxTotal: 4, reservedEnvironmentIds: ["worker:one", "worker:two"] })).toBe(true);
|
||||
expect(check({ maxTotal: 0, reservedEnvironmentIds: ["worker:one", "worker:two"] })).toBe(true);
|
||||
for (const preparedPool of [
|
||||
{},
|
||||
{ maxTotal: 4 },
|
||||
{ reservedEnvironmentIds: [] },
|
||||
{ maxTotal: -1, reservedEnvironmentIds: [] },
|
||||
{ maxTotal: 0.5, reservedEnvironmentIds: [] },
|
||||
{ maxTotal: 4, reservedEnvironmentIds: [""] },
|
||||
{ maxTotal: 4, reservedEnvironmentIds: [42] },
|
||||
{ maxTotal: 4, reservedEnvironmentIds: "worker:one" },
|
||||
{ maxTotal: 4, reservedEnvironmentIds: ["worker:one", "worker:one"] },
|
||||
{ maxTotal: 4, reservedEnvironmentIds: [], reserved: 0 },
|
||||
]) {
|
||||
expect(check(preparedPool)).toBe(false);
|
||||
}
|
||||
});
|
||||
|
||||
it("accepts bounded desktop availability in environment lists and status responses", () => {
|
||||
const base = { id: "node:mac-1", type: "node", status: "available" };
|
||||
for (const state of ["locked", "unlocked", "unknown"]) {
|
||||
|
|
@ -124,22 +167,70 @@ describe("worker environment protocol schemas", () => {
|
|||
expect(validateEnvironmentsPrepareResult({ ...result, preparationKey: "" })).toBe(false);
|
||||
});
|
||||
|
||||
it("exposes only preparation purpose and key in list and status summaries", () => {
|
||||
it("preserves basic preparation summaries with optional closed lifecycle details", () => {
|
||||
const preparation = { purpose: "build", key: "project-key" };
|
||||
const details = {
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 2_000,
|
||||
consumedAtMs: null,
|
||||
};
|
||||
const baseCommit = "a".repeat(40);
|
||||
for (const purpose of ["build", "reserve"] as const) {
|
||||
const summary = {
|
||||
...workerSummary("requested"),
|
||||
preparation: { purpose, key: "project-key" },
|
||||
preparation: { ...preparation, purpose },
|
||||
};
|
||||
expect(Value.Check(EnvironmentsListResultSchema, { environments: [summary] })).toBe(true);
|
||||
expect(Value.Check(EnvironmentsStatusResultSchema, summary)).toBe(true);
|
||||
}
|
||||
for (const preparation of [
|
||||
{ purpose: "unknown", key: "project-key" },
|
||||
{ purpose: "build", key: "" },
|
||||
{ purpose: "build", key: "project-key", projectPath: "/projects/app" },
|
||||
for (const detail of [
|
||||
details,
|
||||
{ ...details, project: { baseCommit } },
|
||||
{ ...details, consumedAtMs: 1_500, project: { label: "openclaw", baseCommit } },
|
||||
]) {
|
||||
expect(
|
||||
Value.Check(EnvironmentSummarySchema, { ...workerSummary("requested"), preparation }),
|
||||
Value.Check(EnvironmentSummarySchema, {
|
||||
...workerSummary("attached"),
|
||||
preparation: { ...preparation, details: detail },
|
||||
}),
|
||||
).toBe(true);
|
||||
}
|
||||
const { demandAtMs: _demandAtMs, ...withoutDemand } = details;
|
||||
const { expiresAtMs: _expiresAtMs, ...withoutExpiry } = details;
|
||||
const { consumedAtMs: _consumedAtMs, ...withoutConsumption } = details;
|
||||
for (const invalid of [
|
||||
{},
|
||||
withoutDemand,
|
||||
withoutExpiry,
|
||||
withoutConsumption,
|
||||
{ ...details, demandAtMs: -1 },
|
||||
{ ...details, expiresAtMs: 0.5 },
|
||||
{ ...details, consumedAtMs: -1 },
|
||||
{ ...details, consumedAtMs: "1500" },
|
||||
{ ...details, projectPath: "/projects/app" },
|
||||
{ ...details, project: {} },
|
||||
{ ...details, project: { baseCommit: "" } },
|
||||
{ ...details, project: { baseCommit, label: "" } },
|
||||
{ ...details, project: { baseCommit, root: "/projects/app" } },
|
||||
]) {
|
||||
expect(
|
||||
Value.Check(EnvironmentSummarySchema, {
|
||||
...workerSummary("requested"),
|
||||
preparation: { ...preparation, details: invalid },
|
||||
}),
|
||||
).toBe(false);
|
||||
}
|
||||
for (const invalid of [
|
||||
{ ...preparation, purpose: "unknown" },
|
||||
{ ...preparation, key: "" },
|
||||
{ ...preparation, projectPath: "/projects/app" },
|
||||
{ ...preparation, demandAtMs: 1_000 },
|
||||
]) {
|
||||
expect(
|
||||
Value.Check(EnvironmentSummarySchema, {
|
||||
...workerSummary("requested"),
|
||||
preparation: invalid,
|
||||
}),
|
||||
).toBe(false);
|
||||
}
|
||||
});
|
||||
|
|
@ -193,6 +284,7 @@ describe("worker environment protocol schemas", () => {
|
|||
...destroyedBase.worker,
|
||||
leaseId: "lease-1",
|
||||
idleMs: 50,
|
||||
destroyRequestedAtMs: 2_000,
|
||||
error: "provider teardown failed",
|
||||
},
|
||||
};
|
||||
|
|
@ -540,6 +632,15 @@ describe("worker environment protocol schemas", () => {
|
|||
worker: { ...workerSummary("ready", "available").worker, ageMs: -1 },
|
||||
}),
|
||||
).toBe(false);
|
||||
for (const destroyRequestedAtMs of [-1, 0.5, "2000", null]) {
|
||||
const summary = workerSummary("destroying");
|
||||
expect(
|
||||
Value.Check(EnvironmentSummarySchema, {
|
||||
...summary,
|
||||
worker: { ...summary.worker, destroyRequestedAtMs },
|
||||
}),
|
||||
).toBe(false);
|
||||
}
|
||||
expect(
|
||||
Value.Check(EnvironmentSummarySchema, {
|
||||
...workerSummary("attached", "available"),
|
||||
|
|
|
|||
|
|
@ -111,6 +111,7 @@ export const WorkerEnvironmentMetadataSchema = closedObject({
|
|||
state: WorkerEnvironmentStateSchema,
|
||||
ageMs: Type.Integer({ minimum: 0 }),
|
||||
idleMs: Type.Optional(Type.Integer({ minimum: 0 })),
|
||||
destroyRequestedAtMs: Type.Optional(Type.Integer({ minimum: 0 })),
|
||||
attachedSessionIds: Type.Array(NonEmptyString),
|
||||
tunnelStatus: WorkerTunnelStatusSchema,
|
||||
error: Type.Optional(NonEmptyString),
|
||||
|
|
@ -150,6 +151,16 @@ function createEnvironmentSummaryProperties() {
|
|||
closedObject({
|
||||
purpose: Type.Union([Type.Literal("reserve"), Type.Literal("build")]),
|
||||
key: NonEmptyString,
|
||||
details: Type.Optional(
|
||||
closedObject({
|
||||
demandAtMs: Type.Integer({ minimum: 0 }),
|
||||
expiresAtMs: Type.Integer({ minimum: 0 }),
|
||||
consumedAtMs: Type.Union([Type.Integer({ minimum: 0 }), Type.Null()]),
|
||||
project: Type.Optional(
|
||||
closedObject({ label: Type.Optional(NonEmptyString), baseCommit: NonEmptyString }),
|
||||
),
|
||||
}),
|
||||
),
|
||||
}),
|
||||
),
|
||||
};
|
||||
|
|
@ -181,6 +192,7 @@ export const EnvironmentsListParamsSchema = closedObject({
|
|||
runtimeId: Type.Optional(Type.String({ minLength: 1, maxLength: 128 })),
|
||||
projection: Type.Optional(Type.Literal("profiles")),
|
||||
includeDesktopSetup: Type.Optional(Type.Boolean()),
|
||||
includePreparedDetails: Type.Optional(Type.Boolean()),
|
||||
});
|
||||
|
||||
/** Provider-authored machine choice for one configured worker profile. */
|
||||
|
|
@ -216,6 +228,7 @@ export const WorkerExecutionModeSchema = Type.Union([
|
|||
const WorkerEnvironmentProfileSummarySchema = closedObject({
|
||||
id: NonEmptyString,
|
||||
providerId: NonEmptyString,
|
||||
readyWorkers: Type.Optional(Type.Integer({ minimum: 0 })),
|
||||
providerDisplayId: Type.Optional(
|
||||
Type.String({ pattern: "^[a-z][a-z0-9-]{0,63}(?![\\s\\S])", maxLength: 64 }),
|
||||
),
|
||||
|
|
@ -237,10 +250,19 @@ const WorkerEnvironmentProfileSummarySchema = closedObject({
|
|||
export const EnvironmentsListResultSchema = closedObject({
|
||||
environments: Type.Array(EnvironmentSummarySchema),
|
||||
profiles: Type.Optional(Type.Array(WorkerEnvironmentProfileSummarySchema)),
|
||||
preparedPool: Type.Optional(
|
||||
closedObject({
|
||||
maxTotal: Type.Integer({ minimum: 0 }),
|
||||
reservedEnvironmentIds: Type.Array(NonEmptyString, { uniqueItems: true }),
|
||||
}),
|
||||
),
|
||||
});
|
||||
|
||||
/** Status lookup request for one environment id. */
|
||||
export const EnvironmentsStatusParamsSchema = closedObject({ environmentId: NonEmptyString });
|
||||
export const EnvironmentsStatusParamsSchema = closedObject({
|
||||
environmentId: NonEmptyString,
|
||||
includePreparedDetails: Type.Optional(Type.Boolean()),
|
||||
});
|
||||
|
||||
/** Status lookup result for one environment id. */
|
||||
export const EnvironmentsStatusResultSchema = createEnvironmentSummarySchema();
|
||||
|
|
|
|||
|
|
@ -27,11 +27,11 @@ scenario:
|
|||
- docs/help/testing.md
|
||||
codeRefs:
|
||||
- packages/sdk/src/client.ts
|
||||
- packages/sdk/src/app-sdk-composed-resources.e2e.test.ts
|
||||
- src/gateway/server.app-sdk-composed-resources.test.ts
|
||||
- src/gateway/agent-turn/agent-job.ts
|
||||
- src/gateway/server-methods/artifacts.ts
|
||||
- src/gateway/server-methods/environments.ts
|
||||
execution:
|
||||
kind: vitest
|
||||
path: packages/sdk/src/app-sdk-composed-resources.e2e.test.ts
|
||||
path: src/gateway/server.app-sdk-composed-resources.test.ts
|
||||
summary: Run handler-backed deterministic assertions and real-Gateway schema and behavior proof through the source-internal preview SDK.
|
||||
|
|
|
|||
|
|
@ -22,8 +22,8 @@ scenario:
|
|||
- src/gateway/server-methods/environments.ts
|
||||
- src/gateway/server-methods/artifacts.ts
|
||||
- src/gateway/managed-image-attachments.ts
|
||||
- test/e2e/qa-lab/runtime/gateway-agent-artifact-apis.e2e.test.ts
|
||||
- src/gateway/server.agent-artifact-apis.test.ts
|
||||
execution:
|
||||
kind: vitest
|
||||
path: test/e2e/qa-lab/runtime/gateway-agent-artifact-apis.e2e.test.ts
|
||||
path: src/gateway/server.agent-artifact-apis.test.ts
|
||||
summary: Run a real Gateway and operator client across config-backed agents, editable files, injected environment lifecycle, SQLite task transcripts, artifact RPCs, and ticketed HTTP bytes.
|
||||
|
|
|
|||
|
|
@ -48,6 +48,7 @@ function workerRecord(
|
|||
ownerEpoch: 1,
|
||||
createdAtMs: 1_000,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
|
|||
291
src/gateway/server-methods/environments.pool.test.ts
Normal file
291
src/gateway/server-methods/environments.pool.test.ts
Normal file
|
|
@ -0,0 +1,291 @@
|
|||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { GATEWAY_OWNER_PROFILE_ID } from "../../../packages/gateway-protocol/src/schema/users.js";
|
||||
import { listDevicePairing } from "../../infra/device-pairing.js";
|
||||
import { createDeferredCore } from "../../shared/deferred.js";
|
||||
import { collectNodeCatalogRuntimeState } from "../node-registry-private.js";
|
||||
import { handleGatewayRequest } from "../server-methods.js";
|
||||
import {
|
||||
createContext,
|
||||
createOperatorClient,
|
||||
} from "../server-plugin-in-process-dispatch.test-support.js";
|
||||
import { environmentsHandlers } from "./environments.js";
|
||||
import {
|
||||
callEnvironmentMethod,
|
||||
mockContext,
|
||||
pairedNodeDevice,
|
||||
workerRecord,
|
||||
workerService,
|
||||
} from "./environments.test-support.js";
|
||||
|
||||
vi.mock("../../infra/device-pairing.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../../infra/device-pairing.js")>()),
|
||||
listDevicePairing: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("../node-registry-private.js", () => ({
|
||||
collectNodeCatalogRuntimeState: vi.fn(),
|
||||
}));
|
||||
|
||||
beforeEach(() => {
|
||||
vi.spyOn(Date, "now").mockReturnValue(10_000);
|
||||
vi.mocked(collectNodeCatalogRuntimeState).mockReturnValue({
|
||||
sessionHostNodeIds: new Set(),
|
||||
issuesByNodeId: new Map(),
|
||||
workerSlotsByNodeId: new Map(),
|
||||
workerBundleByNodeId: new Map(),
|
||||
});
|
||||
vi.mocked(listDevicePairing).mockResolvedValue({
|
||||
pending: [],
|
||||
paired: [
|
||||
pairedNodeDevice("node-live", { commands: ["system.run"] }),
|
||||
pairedNodeDevice("node-offline", {
|
||||
displayName: "Offline Node",
|
||||
caps: ["screen"],
|
||||
commands: ["camera.snap"],
|
||||
}),
|
||||
],
|
||||
});
|
||||
});
|
||||
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
|
||||
const preparationRequests = [
|
||||
{ scope: "operator.read", requested: true, includeDetails: false },
|
||||
{ scope: "operator.write", requested: true, includeDetails: false },
|
||||
{ scope: "operator.admin", requested: undefined, includeDetails: false },
|
||||
{ scope: "operator.admin", requested: false, includeDetails: false },
|
||||
{ scope: "operator.admin", requested: true, includeDetails: true },
|
||||
];
|
||||
|
||||
describe("prepared worker pool projection", () => {
|
||||
it.each(preparationRequests)(
|
||||
"projects worker metadata for $scope with details requested=$requested",
|
||||
async ({ scope, requested, includeDetails }) => {
|
||||
const project = {
|
||||
label: "example/prepared",
|
||||
baseCommit: "a".repeat(40),
|
||||
root: "/private/workspace",
|
||||
};
|
||||
const preparation = {
|
||||
purpose: "reserve" as const,
|
||||
key: "prepared-project-key",
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 60_000,
|
||||
consumedAtMs: 5_000,
|
||||
project,
|
||||
};
|
||||
const service = workerService({
|
||||
readPreparedPoolSummary: vi.fn(() => ({
|
||||
maxTotal: 4,
|
||||
reservedEnvironmentIds: ["worker-1"],
|
||||
})),
|
||||
readReadyWorkerTarget: vi.fn((profileId) => (profileId === "aws" ? 2 : 0)),
|
||||
list: vi.fn(() => [
|
||||
workerRecord({
|
||||
state: "idle",
|
||||
attachedSessionIds: ["session-z", "session-a", "session-z", " "],
|
||||
idleSinceAtMs: 6_000,
|
||||
destroyRequestedAtMs: 9_000,
|
||||
preparation,
|
||||
}),
|
||||
]),
|
||||
});
|
||||
const [ok, payload] = await callEnvironmentMethod(
|
||||
"environments.list",
|
||||
requested === undefined ? {} : { includePreparedDetails: requested },
|
||||
{
|
||||
service,
|
||||
scopes: [scope],
|
||||
},
|
||||
);
|
||||
|
||||
expect(ok).toBe(true);
|
||||
expect(payload).toMatchObject({
|
||||
profiles: [
|
||||
{ id: "aws", providerId: "crabbox" },
|
||||
{ id: "zeta", providerId: "static-ssh" },
|
||||
],
|
||||
environments: [
|
||||
{ id: "gateway", type: "local" },
|
||||
{ id: "node:node-live", type: "node" },
|
||||
{ id: "node:node-offline", type: "node" },
|
||||
{
|
||||
id: "worker-1",
|
||||
type: "worker",
|
||||
status: "available",
|
||||
trust: "disposable",
|
||||
worker: {
|
||||
providerId: "static-ssh",
|
||||
leaseId: "lease-1",
|
||||
state: "idle",
|
||||
ageMs: 9_000,
|
||||
idleMs: 4_000,
|
||||
attachedSessionIds: ["session-a", "session-z"],
|
||||
tunnelStatus: "stopped",
|
||||
},
|
||||
},
|
||||
],
|
||||
});
|
||||
const worker = (payload as { environments: Array<Record<string, unknown>> }).environments.at(
|
||||
-1,
|
||||
);
|
||||
expect(worker).not.toHaveProperty("sshEndpoint");
|
||||
expect(worker?.worker).not.toHaveProperty("sshEndpoint");
|
||||
expect(worker?.worker).not.toHaveProperty("keyRef");
|
||||
expect(worker?.preparation).toEqual({
|
||||
purpose: "reserve",
|
||||
key: preparation.key,
|
||||
...(includeDetails
|
||||
? {
|
||||
details: {
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 60_000,
|
||||
consumedAtMs: 5_000,
|
||||
project: { label: project.label, baseCommit: project.baseCommit },
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
});
|
||||
if (includeDetails) {
|
||||
expect(worker?.worker).toHaveProperty("destroyRequestedAtMs", 9_000);
|
||||
expect(payload.preparedPool).toEqual({ maxTotal: 4, reservedEnvironmentIds: ["worker-1"] });
|
||||
expect(payload.profiles).toEqual([
|
||||
{ id: "aws", providerId: "crabbox", readyWorkers: 2 },
|
||||
{ id: "zeta", providerId: "static-ssh", readyWorkers: 0 },
|
||||
]);
|
||||
expect(service.readPreparedPoolSummary).toHaveBeenCalledOnce();
|
||||
} else {
|
||||
expect(worker?.worker).not.toHaveProperty("destroyRequestedAtMs");
|
||||
expect(payload).not.toHaveProperty("preparedPool");
|
||||
expect(payload.profiles).toEqual([
|
||||
{ id: "aws", providerId: "crabbox" },
|
||||
{ id: "zeta", providerId: "static-ssh" },
|
||||
]);
|
||||
expect(service.readPreparedPoolSummary).not.toHaveBeenCalled();
|
||||
expect(service.readReadyWorkerTarget).not.toHaveBeenCalled();
|
||||
}
|
||||
expect(service.list).toHaveBeenCalledOnce();
|
||||
for (const profile of (payload as { profiles: Array<Record<string, unknown>> }).profiles) {
|
||||
expect(profile).not.toHaveProperty("executionMode");
|
||||
expect(profile).not.toHaveProperty("executionModes");
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it.each(preparationRequests)(
|
||||
"returns worker status for $scope with details requested=$requested",
|
||||
async ({ scope, requested, includeDetails }) => {
|
||||
const get = vi.fn(() =>
|
||||
workerRecord({
|
||||
state: "attached",
|
||||
preparation: {
|
||||
purpose: "build",
|
||||
key: "build-key",
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 60_000,
|
||||
consumedAtMs: null,
|
||||
},
|
||||
}),
|
||||
);
|
||||
const service = workerService({ get });
|
||||
const [ok, payload] = await callEnvironmentMethod(
|
||||
"environments.status",
|
||||
{
|
||||
environmentId: "worker-1",
|
||||
...(requested === undefined ? {} : { includePreparedDetails: requested }),
|
||||
},
|
||||
{ service, scopes: [scope] },
|
||||
);
|
||||
|
||||
expect(ok).toBe(true);
|
||||
expect(payload).toMatchObject({
|
||||
id: "worker-1",
|
||||
status: "available",
|
||||
trust: "disposable",
|
||||
worker: { state: "attached", ageMs: 9_000 },
|
||||
});
|
||||
expect(get).toHaveBeenCalledWith("worker-1");
|
||||
expect(payload.preparation).toEqual({
|
||||
purpose: "build",
|
||||
key: "build-key",
|
||||
...(includeDetails
|
||||
? {
|
||||
details: { demandAtMs: 1_000, expiresAtMs: 60_000, consumedAtMs: null },
|
||||
}
|
||||
: {}),
|
||||
});
|
||||
expect(payload.worker).not.toHaveProperty("destroyRequestedAtMs");
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["environments.list", "environments.status"] as const)(
|
||||
"withholds %s after administrator scopes change during discovery",
|
||||
async (method) => {
|
||||
const entered = createDeferredCore();
|
||||
const release = createDeferredCore();
|
||||
const worker = workerRecord({
|
||||
preparation: {
|
||||
purpose: "reserve",
|
||||
key: "prepared-project-key",
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 60_000,
|
||||
consumedAtMs: null,
|
||||
project: { label: "private-project", baseCommit: "a".repeat(40) },
|
||||
},
|
||||
});
|
||||
const service = workerService({ list: vi.fn(() => [worker]), get: vi.fn(() => worker) });
|
||||
const waitForDiscovery = async () => {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
};
|
||||
if (method === "environments.list") {
|
||||
vi.mocked(service.listMachineOptions).mockImplementation(async () => {
|
||||
await waitForDiscovery();
|
||||
return undefined;
|
||||
});
|
||||
} else {
|
||||
vi.mocked(listDevicePairing).mockImplementation(async () => {
|
||||
await waitForDiscovery();
|
||||
return { pending: [], paired: [] };
|
||||
});
|
||||
}
|
||||
const client = createOperatorClient({
|
||||
profileId: GATEWAY_OWNER_PROFILE_ID,
|
||||
scopes: ["operator.admin"],
|
||||
});
|
||||
const respond = vi.fn();
|
||||
const pending = handleGatewayRequest({
|
||||
req: {
|
||||
type: "req",
|
||||
id: `prepared-details-${method}`,
|
||||
method,
|
||||
params: {
|
||||
includePreparedDetails: true,
|
||||
...(method === "environments.status" ? { environmentId: worker.environmentId } : {}),
|
||||
},
|
||||
},
|
||||
client,
|
||||
context: Object.assign(createContext(), mockContext(service)),
|
||||
respond,
|
||||
isWebchatConnect: () => false,
|
||||
extraHandlers: environmentsHandlers,
|
||||
});
|
||||
try {
|
||||
await Promise.race([entered.promise, pending]);
|
||||
expect(respond).not.toHaveBeenCalled();
|
||||
client.connect.scopes = ["operator.read"];
|
||||
} finally {
|
||||
release.resolve();
|
||||
}
|
||||
await pending;
|
||||
expect(respond).toHaveBeenCalledExactlyOnceWith(
|
||||
false,
|
||||
undefined,
|
||||
expect.objectContaining({
|
||||
code: "FORBIDDEN",
|
||||
message: "Gateway requester authority changed",
|
||||
}),
|
||||
);
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
@ -87,7 +87,7 @@ describe("environments.prepare", () => {
|
|||
]);
|
||||
});
|
||||
|
||||
it("projects preparation identity without the durable demand or expiry fields", () => {
|
||||
it("omits administrator preparation details by default", () => {
|
||||
const preparation = {
|
||||
purpose: "build" as const,
|
||||
key: "project-key",
|
||||
|
|
|
|||
|
|
@ -104,6 +104,8 @@ export function workerRecord(overrides: Partial<TestWorkerRecord> = {}): TestWor
|
|||
updatedAtMs: 1_000,
|
||||
stateChangedAtMs: 1_000,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
preparation: null,
|
||||
lastError: null,
|
||||
tunnelStatus: "stopped",
|
||||
desktopAvailable: false,
|
||||
|
|
@ -134,6 +136,8 @@ export const workerService = (overrides: Partial<TestWorkerService> = {}) => ({
|
|||
throw new Error("No attached portal fixture");
|
||||
}),
|
||||
list: vi.fn(() => []),
|
||||
readPreparedPoolSummary: vi.fn(() => ({ maxTotal: 4, reservedEnvironmentIds: [] })),
|
||||
readReadyWorkerTarget: vi.fn(() => 1),
|
||||
get: vi.fn(() => undefined),
|
||||
inventoryVersion: vi.fn(() => 0),
|
||||
readMachineShape: () => undefined,
|
||||
|
|
@ -177,12 +181,14 @@ export async function callEnvironmentMethod(
|
|||
onCleanupError?: (error: unknown) => void,
|
||||
) => Promise<TestWorkerRecord>;
|
||||
connectedNodes?: unknown[];
|
||||
scopes?: string[];
|
||||
} = {},
|
||||
) {
|
||||
const respond = vi.fn();
|
||||
await environmentsHandlers[method]?.({
|
||||
params: params as Record<string, unknown>,
|
||||
respond,
|
||||
...(options.scopes ? { client: { connect: { scopes: options.scopes } } } : {}),
|
||||
context: mockContext(
|
||||
options.service,
|
||||
options.reconcileActive,
|
||||
|
|
|
|||
|
|
@ -2,7 +2,10 @@
|
|||
* Tests for environment gateway methods and configured environment discovery.
|
||||
*/
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { ErrorCodes } from "../../../packages/gateway-protocol/src/index.js";
|
||||
import {
|
||||
ErrorCodes,
|
||||
type EnvironmentsListResult,
|
||||
} from "../../../packages/gateway-protocol/src/index.js";
|
||||
import { listDevicePairing } from "../../infra/device-pairing.js";
|
||||
import { NODE_RUNNER_UPDATE_REQUIRED_ISSUE } from "../../infra/node-runner-inventory.js";
|
||||
import { NODE_DESKTOP_STREAM_COMMAND } from "../../shared/node-desktop-stream.js";
|
||||
|
|
@ -364,58 +367,6 @@ describe("environment gateway methods", () => {
|
|||
expect(environments.find((entry) => entry.id === "node:node-offline")?.desktop).toBeUndefined();
|
||||
});
|
||||
|
||||
it("appends worker metadata with stable sessions and elapsed times", async () => {
|
||||
const service = workerService({
|
||||
list: vi.fn(() => [
|
||||
workerRecord({
|
||||
state: "idle",
|
||||
attachedSessionIds: ["session-z", "session-a", "session-z", " "],
|
||||
idleSinceAtMs: 6_000,
|
||||
}),
|
||||
]),
|
||||
});
|
||||
const [ok, payload] = await callEnvironmentMethod("environments.list", {}, { service });
|
||||
|
||||
expect(ok).toBe(true);
|
||||
expect(payload).toMatchObject({
|
||||
profiles: [
|
||||
{ id: "aws", providerId: "crabbox" },
|
||||
{ id: "zeta", providerId: "static-ssh" },
|
||||
],
|
||||
environments: [
|
||||
{ id: "gateway", type: "local" },
|
||||
{ id: "node:node-live", type: "node" },
|
||||
{ id: "node:node-offline", type: "node" },
|
||||
{
|
||||
id: "worker-1",
|
||||
type: "worker",
|
||||
status: "available",
|
||||
trust: "disposable",
|
||||
worker: {
|
||||
providerId: "static-ssh",
|
||||
leaseId: "lease-1",
|
||||
state: "idle",
|
||||
ageMs: 9_000,
|
||||
idleMs: 4_000,
|
||||
attachedSessionIds: ["session-a", "session-z"],
|
||||
tunnelStatus: "stopped",
|
||||
},
|
||||
},
|
||||
],
|
||||
});
|
||||
const worker = (payload as { environments: Array<Record<string, unknown>> }).environments.at(
|
||||
-1,
|
||||
);
|
||||
expect(worker).not.toHaveProperty("sshEndpoint");
|
||||
expect(worker?.worker).not.toHaveProperty("sshEndpoint");
|
||||
expect(worker?.worker).not.toHaveProperty("keyRef");
|
||||
expect(service.list).toHaveBeenCalledOnce();
|
||||
for (const profile of (payload as { profiles: Array<Record<string, unknown>> }).profiles) {
|
||||
expect(profile).not.toHaveProperty("executionMode");
|
||||
expect(profile).not.toHaveProperty("executionModes");
|
||||
}
|
||||
});
|
||||
|
||||
it("projects display identity without exposing settings or changing routing", async () => {
|
||||
const service = workerService({
|
||||
readProviderDisplayId: vi.fn((id) => (id === "aws" ? "azure" : undefined)),
|
||||
|
|
@ -469,7 +420,8 @@ describe("environment gateway methods", () => {
|
|||
const [ok, payload] = await callEnvironmentMethod("environments.list", {}, { service });
|
||||
|
||||
expect(ok).toBe(true);
|
||||
const profiles = (payload as { profiles: unknown[] }).profiles;
|
||||
const profiles = (payload as EnvironmentsListResult).profiles ?? [];
|
||||
expect(payload).not.toHaveProperty("preparedPool");
|
||||
expect(profiles).toMatchObject([
|
||||
{
|
||||
id: "aws",
|
||||
|
|
@ -495,16 +447,24 @@ describe("environment gateway methods", () => {
|
|||
vi.mocked(listDevicePairing).mockClear();
|
||||
const projected = await callEnvironmentMethod(
|
||||
"environments.list",
|
||||
{ projection: "profiles" },
|
||||
{ service },
|
||||
{ projection: "profiles", includePreparedDetails: true },
|
||||
{ service, scopes: ["operator.admin"] },
|
||||
);
|
||||
expect(projected).toEqual([true, { environments: [], profiles }, undefined]);
|
||||
expect(projected).toEqual([
|
||||
true,
|
||||
{
|
||||
environments: [],
|
||||
profiles: profiles.map((profile) => Object.assign({}, profile, { readyWorkers: 1 })),
|
||||
},
|
||||
undefined,
|
||||
]);
|
||||
expect(service.list).not.toHaveBeenCalled();
|
||||
expect(service.readPreparedPoolSummary).not.toHaveBeenCalled();
|
||||
expect(listDevicePairing).not.toHaveBeenCalled();
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["worker", "pairing"])(
|
||||
it.each(["worker", "pairing", "prepared pool"])(
|
||||
"isolates profile discovery from %s inventory failures without hiding default errors",
|
||||
async (unavailable) => {
|
||||
const list = vi.fn(() => {
|
||||
|
|
@ -513,43 +473,60 @@ describe("environment gateway methods", () => {
|
|||
}
|
||||
return [];
|
||||
});
|
||||
const service = workerService({ list });
|
||||
const service = workerService({
|
||||
list,
|
||||
readPreparedPoolSummary: vi.fn(() => {
|
||||
if (unavailable === "prepared pool") {
|
||||
throw new Error("private prepared inventory failure");
|
||||
}
|
||||
return { maxTotal: 4, reservedEnvironmentIds: [] };
|
||||
}),
|
||||
});
|
||||
if (unavailable === "pairing") {
|
||||
vi.mocked(listDevicePairing).mockRejectedValue(new Error("pairing inventory unavailable"));
|
||||
}
|
||||
const full = await callEnvironmentMethod("environments.list", {}, { service });
|
||||
const full = await callEnvironmentMethod(
|
||||
"environments.list",
|
||||
{ includePreparedDetails: true },
|
||||
{
|
||||
service,
|
||||
scopes: ["operator.admin"],
|
||||
},
|
||||
);
|
||||
expect(full).toEqual([
|
||||
false,
|
||||
undefined,
|
||||
{
|
||||
code: ErrorCodes.UNAVAILABLE,
|
||||
message:
|
||||
unavailable === "worker"
|
||||
? "Error: environment inventory unavailable"
|
||||
: "Error: pairing inventory unavailable",
|
||||
unavailable === "pairing"
|
||||
? "Error: pairing inventory unavailable"
|
||||
: "Error: environment inventory unavailable",
|
||||
},
|
||||
]);
|
||||
expect(service.listMachineOptions).not.toHaveBeenCalled();
|
||||
list.mockClear();
|
||||
vi.mocked(service.readPreparedPoolSummary).mockClear();
|
||||
vi.mocked(listDevicePairing).mockClear();
|
||||
|
||||
const projected = await callEnvironmentMethod(
|
||||
"environments.list",
|
||||
{ projection: "profiles" },
|
||||
{ service },
|
||||
{ projection: "profiles", includePreparedDetails: true },
|
||||
{ service, scopes: ["operator.admin"] },
|
||||
);
|
||||
expect(projected).toEqual([
|
||||
true,
|
||||
{
|
||||
environments: [],
|
||||
profiles: [
|
||||
{ id: "aws", providerId: "crabbox" },
|
||||
{ id: "zeta", providerId: "static-ssh" },
|
||||
{ id: "aws", providerId: "crabbox", readyWorkers: 1 },
|
||||
{ id: "zeta", providerId: "static-ssh", readyWorkers: 1 },
|
||||
],
|
||||
},
|
||||
undefined,
|
||||
]);
|
||||
expect(service.list).not.toHaveBeenCalled();
|
||||
expect(service.readPreparedPoolSummary).not.toHaveBeenCalled();
|
||||
expect(listDevicePairing).not.toHaveBeenCalled();
|
||||
},
|
||||
);
|
||||
|
|
@ -654,25 +631,6 @@ describe("environment gateway methods", () => {
|
|||
});
|
||||
});
|
||||
|
||||
it("returns status for one worker", async () => {
|
||||
const get = vi.fn(() => workerRecord({ state: "attached" }));
|
||||
const service = workerService({ get });
|
||||
const [ok, payload] = await callEnvironmentMethod(
|
||||
"environments.status",
|
||||
{ environmentId: "worker-1" },
|
||||
{ service },
|
||||
);
|
||||
|
||||
expect(ok).toBe(true);
|
||||
expect(payload).toMatchObject({
|
||||
id: "worker-1",
|
||||
status: "available",
|
||||
trust: "disposable",
|
||||
worker: { state: "attached", ageMs: 9_000 },
|
||||
});
|
||||
expect(get).toHaveBeenCalledWith("worker-1");
|
||||
});
|
||||
|
||||
it("rejects unknown environment ids", async () => {
|
||||
const [ok, , error] = await callEnvironmentMethod("environments.status", {
|
||||
environmentId: "missing",
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
import { normalizeSortedUniqueTrimmedStringList } from "@openclaw/normalization-core/string-normalization";
|
||||
import {
|
||||
type EnvironmentSummary,
|
||||
type EnvironmentsListResult,
|
||||
ErrorCodes,
|
||||
errorShape,
|
||||
validateDesktopLaunchParams,
|
||||
|
|
@ -20,7 +21,11 @@ import { listDevicePairing } from "../../infra/device-pairing.js";
|
|||
import { NODE_DESKTOP_STREAM_COMMAND } from "../../shared/node-desktop-stream.js";
|
||||
import type { NodeListNode } from "../../shared/node-list-types.js";
|
||||
import { resolveDesktopObserveRequester } from "../desktop/observe-requester.js";
|
||||
import { WRITE_SCOPE, authorizeOperatorScopesForRequiredScope } from "../method-scopes.js";
|
||||
import {
|
||||
ADMIN_SCOPE,
|
||||
WRITE_SCOPE,
|
||||
authorizeOperatorScopesForRequiredScope,
|
||||
} from "../method-scopes.js";
|
||||
import { createKnownNodeCatalog, listKnownNodes } from "../node-catalog.js";
|
||||
import {
|
||||
isNodeCommandAllowed,
|
||||
|
|
@ -121,7 +126,7 @@ function summarizeNodeEnvironment(
|
|||
}
|
||||
export async function listGatewayEnvironments(
|
||||
context: GatewayRequestContext,
|
||||
workers = listWorkerEnvironments(context),
|
||||
workers = readWorkerInventory(context, false).workers,
|
||||
runtimeId?: string,
|
||||
includeDesktopSetup = false,
|
||||
): Promise<EnvironmentSummary[]> {
|
||||
|
|
@ -177,9 +182,14 @@ export async function listGatewayEnvironments(
|
|||
),
|
||||
];
|
||||
}
|
||||
function listWorkerEnvironments(context: GatewayRequestContext): WorkerEnvironmentServiceRecord[] {
|
||||
function readWorkerInventory(context: GatewayRequestContext, includePreparedDetails: boolean) {
|
||||
try {
|
||||
return context.workerEnvironmentService?.list() ?? [];
|
||||
return {
|
||||
workers: context.workerEnvironmentService?.list() ?? [],
|
||||
preparedPool: includePreparedDetails
|
||||
? context.workerEnvironmentService?.readPreparedPoolSummary()
|
||||
: undefined,
|
||||
};
|
||||
} catch {
|
||||
throw new Error("environment inventory unavailable");
|
||||
}
|
||||
|
|
@ -254,12 +264,17 @@ async function respondWorkerMutation(
|
|||
export const environmentsHandlers: GatewayRequestHandlers = {
|
||||
...environmentsSessionHandlers,
|
||||
...environmentsSessionExecHandlers,
|
||||
"environments.list": async ({ params, respond, client, context }) => {
|
||||
"environments.list": async (options) => {
|
||||
const { params, respond, client, context } = options;
|
||||
if (!assertValidParams(params, validateEnvironmentsListParams, "environments.list", respond)) {
|
||||
return;
|
||||
}
|
||||
const scopes = Array.isArray(client?.connect.scopes) ? client.connect.scopes : [];
|
||||
const includePreparedDetails =
|
||||
params.includePreparedDetails === true &&
|
||||
authorizeOperatorScopesForRequiredScope(ADMIN_SCOPE, scopes).allowed;
|
||||
const authority = readGatewayRequestMutationAuthority(options);
|
||||
if (params.runtimeId) {
|
||||
const scopes = Array.isArray(client?.connect.scopes) ? client.connect.scopes : [];
|
||||
const access = authorizeOperatorScopesForRequiredScope(WRITE_SCOPE, scopes);
|
||||
if (!access.allowed) {
|
||||
respond(
|
||||
|
|
@ -272,33 +287,70 @@ export const environmentsHandlers: GatewayRequestHandlers = {
|
|||
}
|
||||
await respondUnavailableOnThrow(respond, async () => {
|
||||
let environments: EnvironmentSummary[] = [];
|
||||
let workers: WorkerEnvironmentServiceRecord[] = [];
|
||||
let preparedPool: EnvironmentsListResult["preparedPool"];
|
||||
if (params.projection !== "profiles") {
|
||||
const workers = listWorkerEnvironments(context);
|
||||
const inventory = readWorkerInventory(context, includePreparedDetails);
|
||||
workers = inventory.workers;
|
||||
preparedPool = inventory.preparedPool;
|
||||
environments = await listGatewayEnvironments(
|
||||
context,
|
||||
workers,
|
||||
params.runtimeId,
|
||||
params.includeDesktopSetup,
|
||||
);
|
||||
const summarizedAtMs = Date.now();
|
||||
environments.push(
|
||||
...workers.map((record) => summarizeWorkerEnvironment(record, summarizedAtMs)),
|
||||
);
|
||||
}
|
||||
const profiles = await listWorkerProfilesWithMachines(context);
|
||||
respond(true, { environments, ...(profiles.length > 0 ? { profiles } : {}) }, undefined);
|
||||
authority.assertCurrent();
|
||||
const includeCurrentPreparedDetails =
|
||||
includePreparedDetails &&
|
||||
authorizeOperatorScopesForRequiredScope(
|
||||
ADMIN_SCOPE,
|
||||
Array.isArray(client?.connect.scopes) ? client.connect.scopes : [],
|
||||
).allowed;
|
||||
const summarizedAtMs = Date.now();
|
||||
environments.push(
|
||||
...workers.map((record) =>
|
||||
summarizeWorkerEnvironment(record, summarizedAtMs, {
|
||||
includePreparedDetails: includeCurrentPreparedDetails,
|
||||
}),
|
||||
),
|
||||
);
|
||||
respond(
|
||||
true,
|
||||
{
|
||||
environments,
|
||||
...(profiles.length > 0
|
||||
? {
|
||||
profiles: includeCurrentPreparedDetails
|
||||
? profiles.map((profile) => ({
|
||||
...profile,
|
||||
readyWorkers: context.workerEnvironmentService?.readReadyWorkerTarget(
|
||||
profile.id,
|
||||
),
|
||||
}))
|
||||
: profiles,
|
||||
}
|
||||
: {}),
|
||||
...(includeCurrentPreparedDetails && preparedPool ? { preparedPool } : {}),
|
||||
},
|
||||
undefined,
|
||||
);
|
||||
});
|
||||
},
|
||||
"environments.status": async ({ params, respond, context }) => {
|
||||
"environments.status": async (options) => {
|
||||
const { params, respond, client, context } = options;
|
||||
if (
|
||||
!assertValidParams(params, validateEnvironmentsStatusParams, "environments.status", respond)
|
||||
) {
|
||||
return;
|
||||
}
|
||||
const authority = readGatewayRequestMutationAuthority(options);
|
||||
await respondUnavailableOnThrow(respond, async () => {
|
||||
const environment = (await listGatewayEnvironments(context)).find(
|
||||
(entry) => entry.id === params.environmentId,
|
||||
);
|
||||
authority.assertCurrent();
|
||||
if (environment) {
|
||||
respond(true, environment, undefined);
|
||||
return;
|
||||
|
|
@ -316,7 +368,16 @@ export const environmentsHandlers: GatewayRequestHandlers = {
|
|||
}
|
||||
respond(
|
||||
Boolean(worker),
|
||||
worker ? summarizeWorkerEnvironment(worker) : undefined,
|
||||
worker
|
||||
? summarizeWorkerEnvironment(worker, Date.now(), {
|
||||
includePreparedDetails:
|
||||
params.includePreparedDetails === true &&
|
||||
authorizeOperatorScopesForRequiredScope(
|
||||
ADMIN_SCOPE,
|
||||
Array.isArray(client?.connect.scopes) ? client.connect.scopes : [],
|
||||
).allowed,
|
||||
})
|
||||
: undefined,
|
||||
worker ? undefined : errorShape(ErrorCodes.INVALID_REQUEST, "unknown environmentId"),
|
||||
);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -700,6 +700,7 @@ describe("sessions.dispatch", () => {
|
|||
ownerEpoch,
|
||||
createdAtMs: 1,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [sessionId],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
|
|||
|
|
@ -3,36 +3,36 @@ import { createHash } from "node:crypto";
|
|||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { clearConfigCache, clearRuntimeConfigSnapshot } from "../../../../src/config/config.js";
|
||||
import { resolveSessionStorePathCore } from "../../../../src/config/sessions/paths.js";
|
||||
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
|
||||
import { clearConfigCache, clearRuntimeConfigSnapshot } from "../config/config.js";
|
||||
import { resolveSessionStorePathCore } from "../config/sessions/paths.js";
|
||||
import {
|
||||
appendTranscriptMessage,
|
||||
upsertSessionEntryCore,
|
||||
} from "../../../../src/config/sessions/session-accessor.js";
|
||||
import { clearSessionStoreCacheForTest } from "../../../../src/config/sessions/store-writer-state.js";
|
||||
} from "../config/sessions/session-accessor.js";
|
||||
import { clearSessionStoreCacheForTest } from "../config/sessions/store-writer-state.js";
|
||||
import { closeOpenClawAgentDatabasesForTest } from "../state/openclaw-agent-db.js";
|
||||
import {
|
||||
closeOpenClawStateDatabaseAsync,
|
||||
closeOpenClawStateDatabaseForTest,
|
||||
} from "../state/openclaw-state-db.js";
|
||||
import { createTaskRecord, deleteTaskRecordById } from "../tasks/task-registry.js";
|
||||
import { captureEnv, setTestEnvValue } from "../test-utils/env.js";
|
||||
import {
|
||||
attachManagedOutgoingMediaToMessage,
|
||||
cleanupManagedOutgoingMediaRecords,
|
||||
createManagedOutgoingMediaBlocks,
|
||||
} from "../../../../src/gateway/managed-image-attachments.js";
|
||||
import { listManagedImageRecordEntries } from "../../../../src/gateway/managed-image-record-store.js";
|
||||
import { ADMIN_SCOPE, READ_SCOPE } from "../../../../src/gateway/method-scopes.js";
|
||||
import { startGatewayServer } from "../../../../src/gateway/server.js";
|
||||
} from "./managed-image-attachments.js";
|
||||
import { listManagedImageRecordEntries } from "./managed-image-record-store.js";
|
||||
import { ADMIN_SCOPE, READ_SCOPE } from "./method-scopes.js";
|
||||
import { startGatewayServer } from "./server.js";
|
||||
import {
|
||||
connectGatewayClient,
|
||||
disconnectGatewayClient,
|
||||
getGatewayE2ePortBlock,
|
||||
} from "../../../../src/gateway/test-helpers.e2e.js";
|
||||
import { GATEWAY_STARTUP_MUTATED_ENV_KEYS } from "../../../../src/gateway/test-helpers.env.js";
|
||||
import type { WorkerEnvironmentServiceRecord } from "../../../../src/gateway/worker-environments/service-contract.js";
|
||||
import { closeOpenClawAgentDatabasesForTest } from "../../../../src/state/openclaw-agent-db.js";
|
||||
import {
|
||||
closeOpenClawStateDatabaseAsync,
|
||||
closeOpenClawStateDatabaseForTest,
|
||||
} from "../../../../src/state/openclaw-state-db.js";
|
||||
import { createTaskRecord, deleteTaskRecordById } from "../../../../src/tasks/task-registry.js";
|
||||
import { captureEnv, setTestEnvValue } from "../../../../src/test-utils/env.js";
|
||||
import { useAutoCleanupTempDirTracker } from "../../../helpers/temp-dir.js";
|
||||
} from "./test-helpers.e2e.js";
|
||||
import { GATEWAY_STARTUP_MUTATED_ENV_KEYS } from "./test-helpers.env.js";
|
||||
import type { WorkerEnvironmentServiceRecord } from "./worker-environments/service-contract.js";
|
||||
|
||||
const injectedWorkerService = vi.hoisted(() => {
|
||||
const records = new Map<string, WorkerEnvironmentServiceRecord>();
|
||||
|
|
@ -42,6 +42,8 @@ const injectedWorkerService = vi.hoisted(() => {
|
|||
const service = {
|
||||
list: () => [...records.values()],
|
||||
get: (environmentId: string) => records.get(environmentId),
|
||||
readPreparedPoolSummary: () => ({ maxTotal: 4, reservedEnvironmentIds: [] }),
|
||||
readReadyWorkerTarget: () => 1,
|
||||
create: async (profileId: string, idempotencyKey: string) => {
|
||||
const existingId = idempotency.get(idempotencyKey);
|
||||
if (existingId) {
|
||||
|
|
@ -59,6 +61,7 @@ const injectedWorkerService = vi.hoisted(() => {
|
|||
ownerEpoch: 1,
|
||||
createdAtMs: 1_800_000_000_000,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
@ -92,10 +95,10 @@ const injectedWorkerService = vi.hoisted(() => {
|
|||
};
|
||||
});
|
||||
|
||||
vi.mock("../../../../src/gateway/server-request-context.js", async () => {
|
||||
const actual = await vi.importActual<
|
||||
typeof import("../../../../src/gateway/server-request-context.js")
|
||||
>("../../../../src/gateway/server-request-context.js");
|
||||
vi.mock("./server-request-context.js", async () => {
|
||||
const actual = await vi.importActual<typeof import("./server-request-context.js")>(
|
||||
"./server-request-context.js",
|
||||
);
|
||||
return {
|
||||
...actual,
|
||||
createGatewayRequestContext: (
|
||||
|
|
@ -21,35 +21,36 @@ import {
|
|||
EnvironmentsListResultSchema,
|
||||
EnvironmentsStatusParamsSchema,
|
||||
EnvironmentsStatusResultSchema,
|
||||
} from "../../../packages/gateway-protocol/src/index.js";
|
||||
import { AgentWaitParamsSchema } from "../../../packages/gateway-protocol/src/schema/agent.js";
|
||||
} from "../../packages/gateway-protocol/src/index.js";
|
||||
import { AgentWaitParamsSchema } from "../../packages/gateway-protocol/src/schema/agent.js";
|
||||
import {
|
||||
ArtifactsDownloadResultSchema,
|
||||
ArtifactsGetResultSchema,
|
||||
ArtifactsListResultSchema,
|
||||
} from "../../../packages/gateway-protocol/src/schema/artifacts.js";
|
||||
import { environmentsHandlers } from "../../../src/gateway/server-methods/environments.js";
|
||||
import type {
|
||||
GatewayRequestHandlerOptions,
|
||||
RespondFn,
|
||||
} from "../../../src/gateway/server-methods/types.js";
|
||||
} from "../../packages/gateway-protocol/src/schema/artifacts.js";
|
||||
import {
|
||||
GatewayClientTransport,
|
||||
OpenClaw,
|
||||
type OpenClawEvent,
|
||||
} from "../../packages/sdk/src/index.js";
|
||||
import { emitAgentEvent } from "../infra/agent-events.js";
|
||||
import { registerAgentRunContext } from "../infra/agent-run-registry.js";
|
||||
import { withTimeout } from "../utils/with-timeout.js";
|
||||
import { environmentsHandlers } from "./server-methods/environments.js";
|
||||
import type { GatewayRequestHandlerOptions, RespondFn } from "./server-methods/types.js";
|
||||
import {
|
||||
installGatewayTestHooks,
|
||||
startServer,
|
||||
testState,
|
||||
writeSessionStore,
|
||||
} from "../../../src/gateway/test-helpers.js";
|
||||
} from "./test-helpers.js";
|
||||
import type {
|
||||
WorkerEnvironmentServiceContract,
|
||||
WorkerEnvironmentServiceRecord,
|
||||
} from "../../../src/gateway/worker-environments/service-contract.js";
|
||||
import { emitAgentEvent } from "../../../src/infra/agent-events.js";
|
||||
import { registerAgentRunContext } from "../../../src/infra/agent-run-registry.js";
|
||||
import { withTimeout } from "../../../src/utils/with-timeout.js";
|
||||
import { GatewayClientTransport, OpenClaw, type OpenClawEvent } from "./index.js";
|
||||
} from "./worker-environments/service-contract.js";
|
||||
|
||||
vi.mock("../../../src/infra/device-pairing.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../../../src/infra/device-pairing.js")>();
|
||||
vi.mock("../infra/device-pairing.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../infra/device-pairing.js")>();
|
||||
return {
|
||||
...actual,
|
||||
listDevicePairing: vi.fn(async () => ({ pending: [], paired: [] })),
|
||||
|
|
@ -57,8 +58,8 @@ vi.mock("../../../src/infra/device-pairing.js", async (importOriginal) => {
|
|||
};
|
||||
});
|
||||
|
||||
vi.mock("../../../src/infra/device-pairing-node.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../../../src/infra/device-pairing-node.js")>();
|
||||
vi.mock("../infra/device-pairing-node.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../infra/device-pairing-node.js")>();
|
||||
return {
|
||||
...actual,
|
||||
listNodePairing: vi.fn(async () => ({ pending: [], paired: [] })),
|
||||
|
|
@ -116,6 +117,7 @@ function workerRecord(state: "requested" | "ready" | "destroyed"): WorkerEnviron
|
|||
ownerEpoch: 1,
|
||||
createdAtMs: 1_000,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: ["session-sdk-e2e"],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
@ -160,6 +162,8 @@ async function createFakeGateway(): Promise<FakeGateway> {
|
|||
supportsExecutionMode: (profileId, mode) =>
|
||||
profileId === "development" && mode === "worker-turn",
|
||||
readProviderDisplayId: () => undefined,
|
||||
readPreparedPoolSummary: () => ({ maxTotal: 4, reservedEnvironmentIds: [] }),
|
||||
readReadyWorkerTarget: () => 1,
|
||||
listMachineOptions: async () => undefined,
|
||||
listOperatingSystems: async () => undefined,
|
||||
prepare: async () => {
|
||||
|
|
@ -457,6 +457,7 @@ function decorate(row: GatewaySessionRow, fixture: RowFixture, cfg: OpenClawConf
|
|||
sharedHost: null,
|
||||
createdAtMs: START,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import { normalizeCloudRepo } from "../../config/cloud-worker-project-profiles.js";
|
||||
import type { OpenClawConfig } from "../../config/types.js";
|
||||
import { withTimeout } from "../../infra/fs-safe.js";
|
||||
import type { WorkerProvider } from "../../plugins/types.js";
|
||||
|
|
@ -151,8 +152,27 @@ export function createWorkerEnvironmentAccess(options: WorkerEnvironmentAccessOp
|
|||
inState(record, "ready", "idle", "attached") &&
|
||||
record.desktop !== null;
|
||||
const nodeTunnelStatus = nodeTunnels?.status(record.environmentId);
|
||||
const preparedProject = record.preparation
|
||||
? readWorkerProjectSnapshot(record.profileSnapshot.project)
|
||||
: undefined;
|
||||
const projectLabel = preparedProject
|
||||
? "source" in preparedProject
|
||||
? normalizeCloudRepo(preparedProject.source.url)
|
||||
: preparedProject.label
|
||||
: undefined;
|
||||
return {
|
||||
...record,
|
||||
...(record.preparation && preparedProject
|
||||
? {
|
||||
preparation: {
|
||||
...record.preparation,
|
||||
project: {
|
||||
...(projectLabel ? { label: projectLabel } : {}),
|
||||
baseCommit: preparedProject.baseCommit,
|
||||
},
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
...((record.state === "failed" || record.state === "orphaned") && record.lastError
|
||||
? { error: boundedError(record.lastError) }
|
||||
: {}),
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ const WORKER_STATUS: Record<WorkerEnvironmentState, EnvironmentSummary["status"]
|
|||
export function summarizeWorkerEnvironment(
|
||||
record: WorkerEnvironmentServiceRecord,
|
||||
now = Date.now(),
|
||||
options: { includePreparedDetails?: boolean } = {},
|
||||
): EnvironmentSummary {
|
||||
return {
|
||||
id: record.environmentId,
|
||||
|
|
@ -31,13 +32,40 @@ export function summarizeWorkerEnvironment(
|
|||
: { trust: record.sharedHost ? "persistent" : "disposable" }),
|
||||
...(record.desktopAvailable ? { desktop: true } : {}),
|
||||
...(record.preparation
|
||||
? { preparation: { purpose: record.preparation.purpose, key: record.preparation.key } }
|
||||
? {
|
||||
preparation: {
|
||||
purpose: record.preparation.purpose,
|
||||
key: record.preparation.key,
|
||||
...(options.includePreparedDetails
|
||||
? {
|
||||
details: {
|
||||
demandAtMs: record.preparation.demandAtMs,
|
||||
expiresAtMs: record.preparation.expiresAtMs,
|
||||
consumedAtMs: record.preparation.consumedAtMs,
|
||||
...(record.preparation.project
|
||||
? {
|
||||
project: {
|
||||
...(record.preparation.project.label
|
||||
? { label: record.preparation.project.label }
|
||||
: {}),
|
||||
baseCommit: record.preparation.project.baseCommit,
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
worker: {
|
||||
profileId: record.profileId,
|
||||
providerId: record.providerId,
|
||||
...(record.leaseId ? { leaseId: record.leaseId } : {}),
|
||||
state: record.state,
|
||||
...(options.includePreparedDetails && record.destroyRequestedAtMs !== null
|
||||
? { destroyRequestedAtMs: record.destroyRequestedAtMs }
|
||||
: {}),
|
||||
ageMs: Math.max(0, Math.trunc(now - record.createdAtMs)),
|
||||
...(record.state === "idle" && record.idleSinceAtMs !== null
|
||||
? { idleMs: Math.max(0, Math.trunc(now - record.idleSinceAtMs)) }
|
||||
|
|
|
|||
|
|
@ -84,6 +84,7 @@ describe("worker placement projection", () => {
|
|||
sharedHost: null,
|
||||
createdAtMs: 1,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
@ -184,6 +185,7 @@ describe("worker placement projection", () => {
|
|||
ownerEpoch: active.activeOwnerEpoch,
|
||||
createdAtMs: 1,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [active.sessionId],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
@ -230,6 +232,7 @@ describe("worker placement projection", () => {
|
|||
ownerEpoch: active.activeOwnerEpoch,
|
||||
createdAtMs: 1,
|
||||
idleSinceAtMs: null,
|
||||
destroyRequestedAtMs: null,
|
||||
attachedSessionIds: [active.sessionId],
|
||||
desktopAvailable: false,
|
||||
desktopApps: [],
|
||||
|
|
|
|||
|
|
@ -9,6 +9,75 @@ import {
|
|||
describe("prepared pool retention and source admission", () => {
|
||||
const fixture = usePreparedPoolFixture();
|
||||
|
||||
it("publishes current reservation identities through consumption, cleanup, and policy changes", async () => {
|
||||
const owner = fixture.pool();
|
||||
expect(owner.summary()).toEqual({ maxTotal: 4, reservedEnvironmentIds: [] });
|
||||
expect(owner.target("development")).toBe(1);
|
||||
expect(owner.target("missing")).toBe(0);
|
||||
|
||||
const preparing = await fixture.seed("preparing", { purpose: "build" });
|
||||
const reserve = await fixture.ready(await fixture.seed("reserve", { reserve: true }));
|
||||
expect(owner.summary()).toEqual({
|
||||
maxTotal: 4,
|
||||
reservedEnvironmentIds: ["preparing", "reserve"],
|
||||
});
|
||||
|
||||
const consumed = await fixture.attach(reserve);
|
||||
expect(owner.summary()).toEqual({ maxTotal: 4, reservedEnvironmentIds: ["preparing"] });
|
||||
await fixture.store.requestDestroy({
|
||||
environmentId: consumed.environmentId,
|
||||
state: consumed.state,
|
||||
});
|
||||
expect(owner.summary()).toEqual({
|
||||
maxTotal: 4,
|
||||
reservedEnvironmentIds: ["reserve", "preparing"],
|
||||
});
|
||||
await fixture.teardown(consumed);
|
||||
expect(owner.summary()).toEqual({ maxTotal: 4, reservedEnvironmentIds: ["preparing"] });
|
||||
|
||||
const orphaned = await fixture.attach(
|
||||
await fixture.ready(await fixture.seed("orphaned", { reserve: true })),
|
||||
);
|
||||
await fixture.store.transition({
|
||||
environmentId: orphaned.environmentId,
|
||||
from: "attached",
|
||||
to: "orphaned",
|
||||
});
|
||||
const failed = await fixture.seed("failed", { reserve: true });
|
||||
await fixture.store.transition({
|
||||
environmentId: failed.environmentId,
|
||||
from: "requested",
|
||||
to: "failed",
|
||||
});
|
||||
expect(owner.summary()).toEqual({
|
||||
maxTotal: 4,
|
||||
reservedEnvironmentIds: ["orphaned", "preparing"],
|
||||
});
|
||||
|
||||
fixture.nowMs = preparing.preparation!.expiresAtMs;
|
||||
await fixture.store.requestPreparedDestroy({
|
||||
environmentId: preparing.environmentId,
|
||||
ownerEpoch: preparing.ownerEpoch,
|
||||
preparationKey: preparing.preparation!.key,
|
||||
reason: "expired",
|
||||
assertCurrent: () => {},
|
||||
});
|
||||
expect(owner.summary()).toEqual({
|
||||
maxTotal: 4,
|
||||
reservedEnvironmentIds: ["orphaned", "preparing"],
|
||||
});
|
||||
|
||||
fixture.developmentProfile.readyWorkers = 3;
|
||||
expect(owner.target("development")).toBe(3);
|
||||
fixture.developmentProfile.readyWorkers = 0;
|
||||
fixture.config.cloudWorkers!.preparedPool = { maxTotal: 0 };
|
||||
expect(owner.target("development")).toBe(0);
|
||||
expect(owner.summary()).toEqual({
|
||||
maxTotal: 0,
|
||||
reservedEnvironmentIds: ["orphaned", "preparing"],
|
||||
});
|
||||
});
|
||||
|
||||
it("does not read source admission while ready capacity is full after restart", async () => {
|
||||
await fixture.attach(await fixture.ready(await fixture.seed("source")));
|
||||
const reserve = await fixture.ready(await fixture.seed("reserve", { reserve: true }));
|
||||
|
|
|
|||
|
|
@ -58,17 +58,24 @@ export function createPreparedWorkerPool(options: PoolOptions) {
|
|||
let requested = false;
|
||||
const preparations = new Map<string, AbortController>();
|
||||
const current = () => signal.throwIfAborted();
|
||||
const policy = (record: Pick<WorkerEnvironmentRecord, "profileId" | "providerId">) => {
|
||||
const configuredPolicy = (profileId: string) => {
|
||||
const config = options.getConfig().cloudWorkers;
|
||||
const profile = config?.profiles?.[record.profileId];
|
||||
const configured =
|
||||
profile && normalizeCapabilityProviderId(profile.provider) === record.providerId;
|
||||
const profile = config?.profiles?.[profileId];
|
||||
return {
|
||||
configured: Boolean(configured),
|
||||
target: configured ? (profile.readyWorkers ?? DEFAULT_READY_WORKERS) : 0,
|
||||
providerId: profile ? normalizeCapabilityProviderId(profile.provider) : undefined,
|
||||
target: profile ? (profile.readyWorkers ?? DEFAULT_READY_WORKERS) : 0,
|
||||
maxTotal: config?.preparedPool?.maxTotal ?? DEFAULT_MAX_TOTAL,
|
||||
};
|
||||
};
|
||||
const policy = (record: Pick<WorkerEnvironmentRecord, "profileId" | "providerId">) => {
|
||||
const config = configuredPolicy(record.profileId);
|
||||
const configured = config.providerId === record.providerId;
|
||||
return {
|
||||
configured,
|
||||
target: configured ? config.target : 0,
|
||||
maxTotal: config.maxTotal,
|
||||
};
|
||||
};
|
||||
const groupKey = (record: WorkerEnvironmentRecord) => {
|
||||
const project = readWorkerProjectSnapshot(record.profileSnapshot.project);
|
||||
return project ? JSON.stringify([record.providerId, record.profileId, project.key]) : undefined;
|
||||
|
|
@ -575,5 +582,17 @@ export function createPreparedWorkerPool(options: PoolOptions) {
|
|||
controller?.abort();
|
||||
}
|
||||
};
|
||||
return { schedule, noteDemand, candidates, maintain, canPruneDemand, cancelPreparation };
|
||||
return {
|
||||
schedule,
|
||||
noteDemand,
|
||||
candidates,
|
||||
maintain,
|
||||
canPruneDemand,
|
||||
cancelPreparation,
|
||||
summary: () => ({
|
||||
maxTotal: options.getConfig().cloudWorkers?.preparedPool?.maxTotal ?? DEFAULT_MAX_TOTAL,
|
||||
reservedEnvironmentIds: store.preparedReservationEnvironmentIds(),
|
||||
}),
|
||||
target: (profileId: string) => configuredPolicy(profileId).target,
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ import type {
|
|||
WorkerProfile,
|
||||
} from "../../plugins/capability-provider.types.js";
|
||||
import type { DesktopObserveRequester } from "../desktop/observe-requester.js";
|
||||
import type { WorkerEnvironmentPreparation } from "./environment-record.js";
|
||||
import type {
|
||||
WorkerPlacementMoveSource,
|
||||
WorkerPlacementMoveTarget,
|
||||
|
|
@ -57,11 +58,14 @@ export type WorkerEnvironmentServiceRecord = {
|
|||
ownerEpoch: number;
|
||||
createdAtMs: number;
|
||||
idleSinceAtMs: number | null;
|
||||
destroyRequestedAtMs: number | null;
|
||||
attachedSessionIds: readonly string[];
|
||||
desktopAvailable: boolean;
|
||||
desktopApps: readonly WorkerDesktopApp["id"][];
|
||||
tunnelStatus: WorkerTunnelStatus;
|
||||
preparation?: { purpose: "reserve" | "build"; key: string } | null;
|
||||
preparation?:
|
||||
| (WorkerEnvironmentPreparation & { project?: { label?: string; baseCommit: string } })
|
||||
| null;
|
||||
error?: string;
|
||||
};
|
||||
|
||||
|
|
@ -134,6 +138,8 @@ export type WorkerEnvironmentServiceContract = {
|
|||
close: () => Promise<void>;
|
||||
}>;
|
||||
list(): WorkerEnvironmentServiceRecord[];
|
||||
readPreparedPoolSummary(): { maxTotal: number; reservedEnvironmentIds: string[] };
|
||||
readReadyWorkerTarget(profileId: string): number;
|
||||
get(environmentId: string): WorkerEnvironmentServiceRecord | undefined;
|
||||
inventoryVersion(): number;
|
||||
readMachineShape(
|
||||
|
|
|
|||
|
|
@ -77,6 +77,7 @@ describe("on-demand prepared worker admission", () => {
|
|||
const f = await fixture();
|
||||
support.getDevelopmentProfile().readyWorkers = 0;
|
||||
const result = await f.service.prepare(f.request);
|
||||
const baseCommit = await requireGit(f.projectPath, ["rev-parse", "HEAD"]);
|
||||
const record = support.testState.store.get(result.environmentId)!;
|
||||
expect(result).toEqual({
|
||||
environmentId: record.environmentId,
|
||||
|
|
@ -90,7 +91,7 @@ describe("on-demand prepared worker admission", () => {
|
|||
executionMode: "worker-turn",
|
||||
project: {
|
||||
root: f.projectPath,
|
||||
baseCommit: await requireGit(f.projectPath, ["rev-parse", "HEAD"]),
|
||||
baseCommit,
|
||||
},
|
||||
},
|
||||
preparation: { purpose: "build", demandAtMs: 1_000, expiresAtMs: 11_000, consumedAtMs: null },
|
||||
|
|
@ -100,9 +101,13 @@ describe("on-demand prepared worker admission", () => {
|
|||
);
|
||||
await support.waitForFast(() => expect(f.provision).toHaveBeenCalledOnce());
|
||||
expect(support.testState.store.get(record.environmentId)?.destroyRequestedAtMs).toBeNull();
|
||||
expect(f.service.list()[0]?.preparation).toMatchObject({
|
||||
expect(f.service.list()[0]?.preparation).toEqual({
|
||||
purpose: "build",
|
||||
key: result.preparationKey,
|
||||
demandAtMs: 1_000,
|
||||
expiresAtMs: 11_000,
|
||||
consumedAtMs: null,
|
||||
project: { label: "project", baseCommit },
|
||||
});
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -650,6 +650,8 @@ export function createWorkerEnvironmentService(options: WorkerEnvironmentService
|
|||
isStopping: () => stopping,
|
||||
recordError: saveError,
|
||||
list: environmentAccess.list,
|
||||
readPreparedPoolSummary: preparedPool.summary,
|
||||
readReadyWorkerTarget: preparedPool.target,
|
||||
supportsProviderExecutionMode: providerSupportsExecutionMode,
|
||||
supportsExecutionMode: (profileId: string, mode: WorkerExecutionMode) => {
|
||||
const profile = options.getConfig().cloudWorkers?.profiles?.[profileId];
|
||||
|
|
|
|||
|
|
@ -346,6 +346,8 @@ export async function createWorkerEnvironmentStore(
|
|||
read(() => owner.hasPendingNodeEnrollmentSetup(setup, device)),
|
||||
preparedCapacity: (input: Parameters<typeof preparedCapacityFromReservations>[1]) =>
|
||||
read(() => preparedCapacityFromReservations(prepared(), input)),
|
||||
preparedReservationEnvironmentIds: () =>
|
||||
read(() => prepared().map((record) => record.environmentId)),
|
||||
isPreparedIntentWithinCapacity: (
|
||||
input: Parameters<typeof isPreparedReservationWithinCapacity>[1],
|
||||
) =>
|
||||
|
|
|
|||
|
|
@ -144,7 +144,9 @@ describe("check-max-lines-ratchet", () => {
|
|||
},
|
||||
);
|
||||
expect(result.status, result.stderr).toBe(status);
|
||||
expect(result.stderr).toBe(stderr);
|
||||
// TypeScript 7.0.2 can emit this standalone line while closing its native parser.
|
||||
// Remove after upgrading past https://github.com/microsoft/TypeScript/pull/64276.
|
||||
expect(result.stderr.replace(/^context canceled\n/m, "")).toBe(stderr);
|
||||
expect(result.stdout).toBe(
|
||||
mode === "max-lines failure first"
|
||||
? ""
|
||||
|
|
|
|||
|
|
@ -384,6 +384,7 @@ export const gatewayMethodsTestExclude = [
|
|||
|
||||
// Gateway server tests that need a private module graph and the plain Vitest runner.
|
||||
export const gatewayServerIsolatedTestFiles = [
|
||||
"src/gateway/server.agent-artifact-apis.test.ts",
|
||||
"src/gateway/server-worker-environment-startup.state.test.ts",
|
||||
// A failed native close permanently fences this process's metadata owner.
|
||||
"src/gateway/server-close.agent-databases.test.ts",
|
||||
|
|
|
|||
|
|
@ -450,6 +450,8 @@ suite.define(() => {
|
|||
});
|
||||
await gateway.waitForRequest("plugins.controlUi.list");
|
||||
await expectLoading();
|
||||
// The sidebar owns manager registration and loads independently of the plugin page.
|
||||
await page.getByRole("link", { name: "Plugins", exact: true }).waitFor();
|
||||
expect(
|
||||
await page.evaluate(() => ({
|
||||
contributions: Boolean(customElements.get("openclaw-plugin-contributions")),
|
||||
|
|
|
|||
|
|
@ -351,6 +351,38 @@ const enSettings = {
|
|||
},
|
||||
},
|
||||
cloudWorkersPage: {
|
||||
pool: {
|
||||
tab: "Pool",
|
||||
title: "Ready pool",
|
||||
description:
|
||||
"Workers preparing for upcoming sessions and spares already running. Refreshes every 10 seconds while this view is visible.",
|
||||
ready: "Ready",
|
||||
preparing: "Preparing",
|
||||
releasing: "Releasing",
|
||||
attention: "Needs attention",
|
||||
expired: "Expired",
|
||||
unavailable: "Unavailable",
|
||||
unknownProject: "Unknown project",
|
||||
reserve: "Automatic reserve",
|
||||
build: "On-demand build",
|
||||
age: "Age: {age}",
|
||||
expiresAt: "Expires {time}",
|
||||
expiredAt: "Expired {time}",
|
||||
offline: "Connect to the Gateway to view the ready pool.",
|
||||
adminRequired: "Administrator access is required to view the ready pool.",
|
||||
refreshFailed: "Could not refresh the pool: {error}.",
|
||||
lastUpdated: "Showing the last update from {time}.",
|
||||
capacity: "{used} of {limit} reserve slots in use",
|
||||
disabledCapacity: "Unused workers awaiting release: {count}",
|
||||
capacityHelp:
|
||||
"Preparing workers and pending cleanup count toward the limit. Unused running workers incur machine charges.",
|
||||
disabled:
|
||||
"The pool is disabled. Unused workers are being released; active sessions continue.",
|
||||
inventoryUnavailable: "Pool inventory unavailable",
|
||||
target: "Ready workers per eligible project: {count}",
|
||||
empty:
|
||||
"No unassigned prepared workers. Eligible sessions prepare a reserve after activation; Profiles controls the target and pool limit.",
|
||||
},
|
||||
snapshots: {
|
||||
title: "Snapshots",
|
||||
viewLabel: "Cloud worker view",
|
||||
|
|
|
|||
352
ui/src/pages/cloud-workers/cloud-worker-pool.test.ts
Normal file
352
ui/src/pages/cloud-workers/cloud-worker-pool.test.ts
Normal file
|
|
@ -0,0 +1,352 @@
|
|||
/* @vitest-environment jsdom */
|
||||
import { expectDefined } from "@openclaw/normalization-core/expect";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type {
|
||||
EnvironmentSummary,
|
||||
EnvironmentsListResult,
|
||||
} from "../../../../packages/gateway-protocol/src/index.js";
|
||||
import { createDeferred } from "../../../../test/helpers/promise.js";
|
||||
import { gatewayHelloForMethods } from "../../test-helpers/gateway-methods.ts";
|
||||
import {
|
||||
button,
|
||||
mountPage,
|
||||
setupSnapshotsDomSuite,
|
||||
} from "./cloud-worker-snapshots-dom.test-support.ts";
|
||||
|
||||
vi.mock("../../components/confirm-dialog.ts", () => ({ showConfirmDialog: vi.fn() }));
|
||||
|
||||
const NOW = Date.parse("2026-09-24T12:00:00.000Z");
|
||||
type Preparation = NonNullable<EnvironmentSummary["preparation"]>;
|
||||
|
||||
function preparedWorker(
|
||||
id: string,
|
||||
options: {
|
||||
status?: EnvironmentSummary["status"];
|
||||
worker?: Partial<NonNullable<EnvironmentSummary["worker"]>>;
|
||||
purpose?: Preparation["purpose"];
|
||||
details?: Partial<NonNullable<Preparation["details"]>>;
|
||||
} = {},
|
||||
): EnvironmentSummary {
|
||||
return {
|
||||
id,
|
||||
type: "worker",
|
||||
status: options.status ?? "available",
|
||||
worker: {
|
||||
providerId: "crabbox",
|
||||
profileId: "linux",
|
||||
state: "ready",
|
||||
ageMs: 60_000,
|
||||
attachedSessionIds: [],
|
||||
tunnelStatus: "stopped",
|
||||
...options.worker,
|
||||
},
|
||||
preparation: {
|
||||
purpose: options.purpose ?? "reserve",
|
||||
key: `preparation-${id}`,
|
||||
details: {
|
||||
demandAtMs: NOW - 60_000,
|
||||
expiresAtMs: NOW + 3_600_000,
|
||||
consumedAtMs: null,
|
||||
project: { label: id, baseCommit: "0123456789abcdef0123456789abcdef01234567" },
|
||||
...options.details,
|
||||
},
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function inventory(
|
||||
environments = [preparedWorker("App")],
|
||||
reservedEnvironmentIds = environments.map((environment) => environment.id),
|
||||
): EnvironmentsListResult {
|
||||
return {
|
||||
environments,
|
||||
profiles: [
|
||||
{ id: "linux", providerId: "crabbox", readyWorkers: 1 },
|
||||
{ id: "linux-build", providerId: "crabbox", readyWorkers: 2 },
|
||||
],
|
||||
preparedPool: { maxTotal: 4, reservedEnvironmentIds },
|
||||
};
|
||||
}
|
||||
|
||||
function mountPool(
|
||||
readInventory: () => EnvironmentsListResult | Promise<EnvironmentsListResult>,
|
||||
scopes?: string[],
|
||||
) {
|
||||
return mountPage(["environments.list"], {
|
||||
scopes,
|
||||
response: (method, params) =>
|
||||
method === "environments.list" && params?.includePreparedDetails === true
|
||||
? readInventory()
|
||||
: undefined,
|
||||
});
|
||||
}
|
||||
|
||||
async function enterPool(fixture: ReturnType<typeof mountPage>) {
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
button(fixture.page, "Pool").click();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
return expectDefined(fixture.page.querySelector("openclaw-cloud-worker-pool"), "Pool view");
|
||||
}
|
||||
|
||||
setupSnapshotsDomSuite();
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(NOW);
|
||||
vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible");
|
||||
});
|
||||
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
|
||||
describe("Cloud worker pool", () => {
|
||||
it("groups unused preparations and keeps reserve capacity separate from visible states", async () => {
|
||||
const fixture = mountPool(() =>
|
||||
inventory(
|
||||
[
|
||||
preparedWorker("Ready app"),
|
||||
preparedWorker("Building app", {
|
||||
status: "starting",
|
||||
purpose: "build",
|
||||
worker: { profileId: "linux-build", state: "provisioning" },
|
||||
}),
|
||||
preparedWorker("Releasing app", {
|
||||
worker: { destroyRequestedAtMs: NOW - 1_000 },
|
||||
details: { consumedAtMs: NOW - 2_000 },
|
||||
}),
|
||||
preparedWorker("Failed app", {
|
||||
status: "error",
|
||||
worker: { state: "failed", error: "Project preparation failed" },
|
||||
}),
|
||||
preparedWorker("Expired app", { details: { expiresAtMs: NOW - 1_000 } }),
|
||||
preparedWorker("Consumed app", { details: { consumedAtMs: NOW - 1_000 } }),
|
||||
preparedWorker("Attached app", { worker: { attachedSessionIds: ["session-app"] } }),
|
||||
preparedWorker("Destroyed app", { worker: { state: "destroyed" } }),
|
||||
{ id: "Local", type: "local", status: "available" },
|
||||
],
|
||||
["Ready app", "Building app", "Releasing app", "Expired app"],
|
||||
),
|
||||
);
|
||||
try {
|
||||
const pool = await enterPool(fixture);
|
||||
expect(pool.textContent).toContain("4 of 4 reserve slots in use");
|
||||
expect(
|
||||
[...pool.querySelectorAll(".settings-summary dt")].map((entry) => entry.textContent),
|
||||
).toEqual(["Ready", "Preparing", "Releasing", "Needs attention"]);
|
||||
expect(
|
||||
[...pool.querySelectorAll(".settings-summary dd")].map((entry) => entry.textContent),
|
||||
).toEqual(["1", "1", "1", "2"]);
|
||||
const profileSections = [...pool.querySelectorAll(".settings-section")].filter(
|
||||
(section) => section.querySelector("h2")?.textContent?.trim() !== "Ready pool",
|
||||
);
|
||||
expect(
|
||||
profileSections.map((section) =>
|
||||
section.querySelector("h2")?.textContent?.replace(/\s+/g, " ").trim(),
|
||||
),
|
||||
).toEqual(["linux 4", "linux-build 1"]);
|
||||
const linux = expectDefined(profileSections[0], "Linux profile");
|
||||
expect(linux.textContent).toContain("Ready workers per eligible project: 1");
|
||||
expect(
|
||||
[...linux.querySelectorAll(".settings-row__title")].map((row) => row.textContent?.trim()),
|
||||
).toEqual(["Ready app", "Releasing app", "Failed app", "Expired app"]);
|
||||
expect(linux.textContent).toContain("01234567");
|
||||
expect(linux.textContent).toContain("Automatic reserve");
|
||||
expect(linux.textContent).toContain("Age: 1m");
|
||||
expect(linux.textContent).toContain("Project preparation failed");
|
||||
expect(linux.textContent).toContain("Expired");
|
||||
expect([...linux.querySelectorAll("time")].map((entry) => entry.dateTime)).toEqual([
|
||||
"2026-09-24T13:00:00.000Z",
|
||||
"2026-09-24T13:00:00.000Z",
|
||||
"2026-09-24T13:00:00.000Z",
|
||||
"2026-09-24T11:59:59.000Z",
|
||||
]);
|
||||
const build = expectDefined(profileSections[1], "Build profile");
|
||||
expect(build.textContent).toContain("Ready workers per eligible project: 2");
|
||||
expect(build.textContent).toContain("Building app");
|
||||
expect(build.textContent).toContain("On-demand build");
|
||||
for (const excluded of ["Consumed app", "Attached app", "Destroyed app", "Local"]) {
|
||||
expect(pool.textContent).not.toContain(excluded);
|
||||
}
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["success", "error"] as const)(
|
||||
"ignores an old request's %s after a same-client reconnect",
|
||||
async (outcome) => {
|
||||
const oldRequest = createDeferred<EnvironmentsListResult>();
|
||||
const currentRequest = createDeferred<EnvironmentsListResult>();
|
||||
const readInventory = vi
|
||||
.fn(() => currentRequest.promise)
|
||||
.mockReturnValueOnce(oldRequest.promise);
|
||||
const fixture = mountPool(readInventory);
|
||||
try {
|
||||
const pool = await enterPool(fixture);
|
||||
expect(readInventory).toHaveBeenCalledTimes(1);
|
||||
fixture.harness.publish(
|
||||
false,
|
||||
fixture.client,
|
||||
gatewayHelloForMethods(["environments.list"]),
|
||||
);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(pool.textContent).toContain("Connect to the Gateway to view the ready pool.");
|
||||
fixture.harness.publish(
|
||||
true,
|
||||
fixture.client,
|
||||
gatewayHelloForMethods(["environments.list"]),
|
||||
);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(readInventory).toHaveBeenCalledTimes(2);
|
||||
if (outcome === "success") {
|
||||
oldRequest.resolve(inventory([preparedWorker("Old connection app")]));
|
||||
} else {
|
||||
oldRequest.reject(new Error("Old connection failure"));
|
||||
}
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(pool.textContent).not.toContain("Old connection");
|
||||
expect(pool.querySelector('[role="alert"]')).toBeNull();
|
||||
expect(button(pool, "Refresh").disabled).toBe(true);
|
||||
currentRequest.resolve(inventory([preparedWorker("Current connection app")]));
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(pool.textContent).toContain("Current connection app");
|
||||
expect(button(pool, "Refresh").disabled).toBe(false);
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("keeps the last inventory when a refresh fails and clears the warning on recovery", async () => {
|
||||
const readInventory = vi
|
||||
.fn<() => Promise<EnvironmentsListResult>>()
|
||||
.mockResolvedValueOnce(inventory())
|
||||
.mockRejectedValueOnce(new Error("Worker inventory is unavailable"))
|
||||
.mockResolvedValueOnce(inventory([preparedWorker("Recovered app")]));
|
||||
const fixture = mountPool(readInventory);
|
||||
try {
|
||||
const pool = await enterPool(fixture);
|
||||
button(pool, "Refresh").click();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(pool.textContent).toContain("App");
|
||||
expect(pool.textContent).toContain("1 of 4 reserve slots in use");
|
||||
const warning = pool.querySelector('[role="alert"]');
|
||||
expect(warning?.textContent).toContain(
|
||||
"Could not refresh the pool: Worker inventory is unavailable.",
|
||||
);
|
||||
expect(warning?.textContent).toContain("Showing the last update from");
|
||||
expect(button(pool, "Refresh").disabled).toBe(false);
|
||||
button(pool, "Refresh").click();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(pool.textContent).toContain("Recovered app");
|
||||
expect(pool.querySelector('[role="alert"]')).toBeNull();
|
||||
expect(readInventory).toHaveBeenCalledTimes(3);
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
});
|
||||
|
||||
it("refreshes every ten seconds only while the Pool tab is visible and mounted", async () => {
|
||||
const readInventory = vi.fn(() => inventory());
|
||||
const fixture = mountPool(readInventory);
|
||||
try {
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(readInventory).not.toHaveBeenCalled();
|
||||
await enterPool(fixture);
|
||||
expect(readInventory).toHaveBeenCalledTimes(1);
|
||||
await vi.advanceTimersByTimeAsync(9_999);
|
||||
expect(readInventory).toHaveBeenCalledTimes(1);
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
expect(readInventory).toHaveBeenCalledTimes(2);
|
||||
vi.spyOn(document, "visibilityState", "get").mockReturnValue("hidden");
|
||||
document.dispatchEvent(new Event("visibilitychange"));
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
expect(readInventory).toHaveBeenCalledTimes(2);
|
||||
vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible");
|
||||
document.dispatchEvent(new Event("visibilitychange"));
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(readInventory).toHaveBeenCalledTimes(3);
|
||||
button(fixture.page, "Profiles").click();
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
expect(readInventory).toHaveBeenCalledTimes(3);
|
||||
await enterPool(fixture);
|
||||
expect(readInventory).toHaveBeenCalledTimes(4);
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
document.dispatchEvent(new Event("visibilitychange"));
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(readInventory).toHaveBeenCalledTimes(4);
|
||||
});
|
||||
|
||||
it("requires administrator scope before requesting pool inventory", async () => {
|
||||
const readInventory = vi.fn(() => inventory());
|
||||
const fixture = mountPool(readInventory, ["operator.read"]);
|
||||
try {
|
||||
const pool = await enterPool(fixture);
|
||||
expect(pool.textContent).toContain(
|
||||
"Administrator access is required to view the ready pool.",
|
||||
);
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
expect(readInventory).not.toHaveBeenCalled();
|
||||
expect(pool.querySelector("button")).toBeNull();
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
});
|
||||
|
||||
it("shows pending cleanup after the pool has been disabled", async () => {
|
||||
const result = inventory([
|
||||
preparedWorker("Cleanup app", {
|
||||
worker: { destroyRequestedAtMs: NOW - 1_000 },
|
||||
details: { consumedAtMs: NOW - 2_000 },
|
||||
}),
|
||||
]);
|
||||
result.preparedPool = { maxTotal: 0, reservedEnvironmentIds: ["Cleanup app"] };
|
||||
const fixture = mountPool(() => result);
|
||||
try {
|
||||
const pool = await enterPool(fixture);
|
||||
expect(pool.textContent).toContain("Unused workers awaiting release: 1");
|
||||
expect(pool.textContent).toContain("The pool is disabled.");
|
||||
expect(pool.textContent).toContain("Cleanup app");
|
||||
expect(pool.textContent).not.toContain("1 of 0");
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
});
|
||||
|
||||
it("defers hidden entry and retires an unfinished inventory request on tab exit", async () => {
|
||||
vi.spyOn(document, "visibilityState", "get").mockReturnValue("hidden");
|
||||
const pending = createDeferred<EnvironmentsListResult>();
|
||||
const readInventory = vi.fn(() => Promise.resolve(inventory([preparedWorker("Fresh app")])));
|
||||
readInventory.mockReturnValueOnce(pending.promise);
|
||||
const fixture = mountPool(readInventory);
|
||||
try {
|
||||
await enterPool(fixture);
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
expect(readInventory).not.toHaveBeenCalled();
|
||||
vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible");
|
||||
document.dispatchEvent(new Event("visibilitychange"));
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(readInventory).toHaveBeenCalledTimes(1);
|
||||
button(fixture.page, "Profiles").click();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
expect(fixture.request).toHaveBeenCalledWith(
|
||||
"environments.list",
|
||||
{ includePreparedDetails: true },
|
||||
{
|
||||
signal: expect.objectContaining({ aborted: true }),
|
||||
},
|
||||
);
|
||||
pending.resolve(inventory([preparedWorker("Retired app")]));
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
expect(readInventory).toHaveBeenCalledTimes(1);
|
||||
const pool = await enterPool(fixture);
|
||||
expect(readInventory).toHaveBeenCalledTimes(2);
|
||||
expect(pool.textContent).toContain("Fresh app");
|
||||
expect(pool.textContent).not.toContain("Retired app");
|
||||
} finally {
|
||||
fixture.dispose();
|
||||
}
|
||||
});
|
||||
});
|
||||
266
ui/src/pages/cloud-workers/cloud-worker-pool.ts
Normal file
266
ui/src/pages/cloud-workers/cloud-worker-pool.ts
Normal file
|
|
@ -0,0 +1,266 @@
|
|||
import { consume } from "@lit/context";
|
||||
import { html, nothing } from "lit";
|
||||
import { state } from "lit/decorators.js";
|
||||
import type {
|
||||
EnvironmentSummary,
|
||||
EnvironmentsListResult,
|
||||
} from "../../../../packages/gateway-protocol/src/index.js";
|
||||
import { applicationContext, type ApplicationContext } from "../../app/context.ts";
|
||||
import {
|
||||
renderSettingsEmpty,
|
||||
renderSettingsPage,
|
||||
renderSettingsRow,
|
||||
renderSettingsSection,
|
||||
renderSettingsStatus,
|
||||
renderSettingsSummary,
|
||||
} from "../../components/settings-ui.ts";
|
||||
import { t } from "../../i18n/index.ts";
|
||||
import { registerSettingsEnglish } from "../../i18n/locales/en-settings.ts";
|
||||
import { formatDurationHuman } from "../../lib/format-duration.ts";
|
||||
import { formatUiError } from "../../lib/format-error.ts";
|
||||
import { formatRelativeTimestamp } from "../../lib/format.ts";
|
||||
import { canCallGatewayMethod } from "../../lib/gateway-methods.ts";
|
||||
import { GatewayPageController } from "../../lit/gateway-page-controller.ts";
|
||||
import { OpenClawLightDomElement } from "../../lit/openclaw-element.ts";
|
||||
import { PollController } from "../../lit/poll-controller.ts";
|
||||
|
||||
registerSettingsEnglish();
|
||||
|
||||
type PoolState = "ready" | "preparing" | "releasing" | "attention";
|
||||
type Preparation = NonNullable<EnvironmentSummary["preparation"]>;
|
||||
type PreparedEnvironment = EnvironmentSummary & {
|
||||
preparation: Preparation & { details: NonNullable<Preparation["details"]> };
|
||||
worker: NonNullable<EnvironmentSummary["worker"]>;
|
||||
};
|
||||
|
||||
function poolState(environment: PreparedEnvironment, now: number): PoolState {
|
||||
const worker = environment.worker;
|
||||
if (environment.status === "error") {
|
||||
return "attention";
|
||||
}
|
||||
if (worker.destroyRequestedAtMs !== undefined || environment.status === "stopping") {
|
||||
return "releasing";
|
||||
}
|
||||
if (environment.preparation.details.expiresAtMs <= now || environment.status === "unavailable") {
|
||||
return "attention";
|
||||
}
|
||||
if (worker.state === "ready") {
|
||||
return "ready";
|
||||
}
|
||||
return environment.status === "starting" ? "preparing" : "attention";
|
||||
}
|
||||
|
||||
class CloudWorkerPool extends OpenClawLightDomElement {
|
||||
@consume({ context: applicationContext, subscribe: true })
|
||||
private context!: ApplicationContext;
|
||||
|
||||
@state() private result: EnvironmentsListResult | null = null;
|
||||
@state() private error: string | null = null;
|
||||
@state() private loading = false;
|
||||
@state() private updatedAt: number | null = null;
|
||||
private request: AbortController | undefined;
|
||||
private readonly polling = new PollController(
|
||||
this,
|
||||
10_000,
|
||||
() => void this.load(),
|
||||
false,
|
||||
"visible",
|
||||
);
|
||||
|
||||
private readonly gateway = new GatewayPageController(this, {
|
||||
getGateway: () => this.context?.gateway,
|
||||
invalidateRequests: () => {
|
||||
this.request?.abort();
|
||||
this.request = undefined;
|
||||
this.polling.stop();
|
||||
this.result = null;
|
||||
this.error = null;
|
||||
this.loading = false;
|
||||
this.updatedAt = null;
|
||||
},
|
||||
ensureInitialData: () => {
|
||||
if (this.canRead()) {
|
||||
this.polling.start();
|
||||
void this.load();
|
||||
}
|
||||
},
|
||||
onPageActivation: () => void this.load(),
|
||||
});
|
||||
|
||||
private canRead() {
|
||||
return canCallGatewayMethod(this.gateway.snapshot, "environments.list", "operator.admin");
|
||||
}
|
||||
|
||||
private async load() {
|
||||
const scope = this.gateway.capture();
|
||||
if (!scope || !this.canRead() || this.loading || document.visibilityState === "hidden") {
|
||||
return;
|
||||
}
|
||||
const request = new AbortController();
|
||||
this.request = request;
|
||||
this.loading = true;
|
||||
try {
|
||||
const result = await scope.client.request<EnvironmentsListResult>(
|
||||
"environments.list",
|
||||
{ includePreparedDetails: true },
|
||||
{ signal: request.signal },
|
||||
);
|
||||
if (this.gateway.isCurrent(scope)) {
|
||||
this.result = result;
|
||||
this.error = null;
|
||||
this.updatedAt = Date.now();
|
||||
}
|
||||
} catch (error) {
|
||||
if (this.gateway.isCurrent(scope)) {
|
||||
this.error = formatUiError(error);
|
||||
}
|
||||
} finally {
|
||||
if (this.gateway.isCurrent(scope)) {
|
||||
this.loading = false;
|
||||
this.request = undefined;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private renderWorker(environment: PreparedEnvironment, now: number) {
|
||||
const { worker, preparation } = environment;
|
||||
const { details } = preparation;
|
||||
const category = poolState(environment, now);
|
||||
const expired = details.expiresAtMs <= now;
|
||||
const label =
|
||||
category === "releasing"
|
||||
? t("cloudWorkersPage.pool.releasing")
|
||||
: worker.state === "failed" || worker.state === "orphaned"
|
||||
? t(`cloudWorkersPage.snapshots.buildStates.${worker.state}`)
|
||||
: expired
|
||||
? t("cloudWorkersPage.pool.expired")
|
||||
: category === "attention"
|
||||
? t("cloudWorkersPage.pool.unavailable")
|
||||
: t(`cloudWorkersPage.pool.${category}`);
|
||||
return renderSettingsRow({
|
||||
title: details.project?.label ?? t("cloudWorkersPage.pool.unknownProject"),
|
||||
description: html`
|
||||
${
|
||||
details.project?.baseCommit
|
||||
? html`<code>${details.project.baseCommit.slice(0, 8)}</code> · `
|
||||
: nothing
|
||||
}
|
||||
${t(`cloudWorkersPage.pool.${preparation.purpose}`)} ·
|
||||
${t("cloudWorkersPage.pool.age", { age: formatDurationHuman(worker.ageMs) })}
|
||||
<br />
|
||||
<time datetime=${new Date(details.expiresAtMs).toISOString()}>
|
||||
${t(expired ? "cloudWorkersPage.pool.expiredAt" : "cloudWorkersPage.pool.expiresAt", {
|
||||
time: formatRelativeTimestamp(details.expiresAtMs),
|
||||
})}
|
||||
</time>
|
||||
${worker.error ? html`<br /><span>${worker.error}</span>` : nothing}
|
||||
`,
|
||||
stackedOnNarrow: true,
|
||||
control: renderSettingsStatus({
|
||||
kind: category === "ready" ? "ok" : category === "attention" ? "warn" : "accent",
|
||||
label,
|
||||
}),
|
||||
});
|
||||
}
|
||||
|
||||
override render() {
|
||||
if (!this.gateway.connected) {
|
||||
return renderSettingsPage(renderSettingsEmpty(t("cloudWorkersPage.pool.offline")));
|
||||
}
|
||||
if (!this.canRead()) {
|
||||
return renderSettingsPage(renderSettingsEmpty(t("cloudWorkersPage.pool.adminRequired")));
|
||||
}
|
||||
const now = Date.now();
|
||||
const pool = this.result?.preparedPool;
|
||||
const reserved = new Set(pool?.reservedEnvironmentIds);
|
||||
const rows = (this.result?.environments ?? []).filter(
|
||||
(environment): environment is PreparedEnvironment =>
|
||||
environment.preparation?.details !== undefined &&
|
||||
environment.worker !== undefined &&
|
||||
environment.worker.state !== "destroyed" &&
|
||||
(reserved.has(environment.id) ||
|
||||
(environment.preparation.details.consumedAtMs === null &&
|
||||
environment.worker.attachedSessionIds.length === 0)),
|
||||
);
|
||||
const groups = new Map<string, PreparedEnvironment[]>();
|
||||
for (const row of rows) {
|
||||
const profileId = row.worker.profileId ?? "";
|
||||
const group = groups.get(profileId) ?? [];
|
||||
group.push(row);
|
||||
groups.set(profileId, group);
|
||||
}
|
||||
return renderSettingsPage(html`
|
||||
${renderSettingsSection(
|
||||
{
|
||||
title: t("cloudWorkersPage.pool.title"),
|
||||
description: t("cloudWorkersPage.pool.description"),
|
||||
actions: html`<button
|
||||
class="btn btn--sm"
|
||||
type="button"
|
||||
?disabled=${this.loading}
|
||||
@click=${() => void this.load()}
|
||||
>
|
||||
${t("common.refresh")}
|
||||
</button>`,
|
||||
notice: this.error
|
||||
? html`<div class="callout warning" role="alert">
|
||||
${t("cloudWorkersPage.pool.refreshFailed", { error: this.error })}
|
||||
${
|
||||
this.updatedAt !== null
|
||||
? t("cloudWorkersPage.pool.lastUpdated", {
|
||||
time: formatRelativeTimestamp(this.updatedAt),
|
||||
})
|
||||
: nothing
|
||||
}
|
||||
</div>`
|
||||
: nothing,
|
||||
},
|
||||
renderSettingsRow({
|
||||
title: pool
|
||||
? pool.maxTotal === 0
|
||||
? t("cloudWorkersPage.pool.disabledCapacity", { count: String(reserved.size) })
|
||||
: t("cloudWorkersPage.pool.capacity", {
|
||||
used: String(reserved.size),
|
||||
limit: String(pool.maxTotal),
|
||||
})
|
||||
: t(this.loading ? "common.loading" : "cloudWorkersPage.pool.inventoryUnavailable"),
|
||||
description:
|
||||
pool?.maxTotal === 0
|
||||
? t("cloudWorkersPage.pool.disabled")
|
||||
: t("cloudWorkersPage.pool.capacityHelp"),
|
||||
}),
|
||||
)}
|
||||
${
|
||||
pool
|
||||
? renderSettingsSummary(
|
||||
(["ready", "preparing", "releasing", "attention"] as const).map((category) => ({
|
||||
label: t(`cloudWorkersPage.pool.${category}`),
|
||||
value: rows.filter((row) => poolState(row, now) === category).length,
|
||||
})),
|
||||
)
|
||||
: nothing
|
||||
}
|
||||
${[...groups]
|
||||
.toSorted(([a], [b]) => a.localeCompare(b))
|
||||
.map(([profileId, workers]) => {
|
||||
const profile = this.result?.profiles?.find((entry) => entry.id === profileId);
|
||||
return renderSettingsSection(
|
||||
{
|
||||
title: profileId || t("cloudWorkersPage.snapshots.unlabeledProfile"),
|
||||
description:
|
||||
profile?.readyWorkers !== undefined
|
||||
? t("cloudWorkersPage.pool.target", { count: String(profile.readyWorkers) })
|
||||
: undefined,
|
||||
count: workers.length,
|
||||
},
|
||||
workers.map((worker) => this.renderWorker(worker, now)),
|
||||
);
|
||||
})}
|
||||
${pool && rows.length === 0 ? renderSettingsEmpty(t("cloudWorkersPage.pool.empty")) : nothing}
|
||||
`);
|
||||
}
|
||||
}
|
||||
|
||||
if (!customElements.get("openclaw-cloud-worker-pool")) {
|
||||
customElements.define("openclaw-cloud-worker-pool", CloudWorkerPool);
|
||||
}
|
||||
|
|
@ -41,14 +41,14 @@ export function mountPage(
|
|||
result?: ReturnType<typeof snapshotListFixture>;
|
||||
config?: Record<string, unknown>;
|
||||
failMutation?: boolean;
|
||||
response?: (method: string) => unknown;
|
||||
response?: (method: string, params?: Record<string, unknown>) => unknown;
|
||||
scopes?: string[];
|
||||
} = {},
|
||||
) {
|
||||
let result = options.result ?? snapshotListFixture();
|
||||
let config = options.config ?? {};
|
||||
const request = vi.fn(async (method: string, params?: Record<string, unknown>) => {
|
||||
const response = options.response?.(method);
|
||||
const response = options.response?.(method, params);
|
||||
if (response !== undefined) {
|
||||
return response;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -39,6 +39,7 @@ import {
|
|||
type ConfiguredCloudWorkerProfile,
|
||||
} from "./cloud-worker-config.ts";
|
||||
import { renderCloudWorkerRepositories } from "./cloud-worker-repositories.ts";
|
||||
import "./cloud-worker-pool.ts";
|
||||
import "./cloud-worker-snapshots.ts";
|
||||
|
||||
registerSettingsEnglish();
|
||||
|
|
@ -60,7 +61,7 @@ class CloudWorkersPage extends OpenClawLightDomElement {
|
|||
@consume({ context: applicationContext, subscribe: true })
|
||||
private context!: ApplicationContext;
|
||||
|
||||
@state() private view: "profiles" | "snapshots" = "profiles";
|
||||
@state() private view: "profiles" | "pool" | "snapshots" = "profiles";
|
||||
@state() private editor: EditorState = null;
|
||||
@state() private draft: CloudWorkerProfileDraft = createCloudWorkerDraft();
|
||||
|
||||
|
|
@ -582,6 +583,7 @@ class CloudWorkersPage extends OpenClawLightDomElement {
|
|||
ariaLabel: t("cloudWorkersPage.snapshots.viewLabel"),
|
||||
options: [
|
||||
{ value: "profiles", label: t("cloudWorkersPage.sectionTitle") },
|
||||
{ value: "pool", label: t("cloudWorkersPage.pool.tab") },
|
||||
{ value: "snapshots", label: t("cloudWorkersPage.snapshots.title") },
|
||||
],
|
||||
onChange: (value) => {
|
||||
|
|
@ -589,7 +591,13 @@ class CloudWorkersPage extends OpenClawLightDomElement {
|
|||
},
|
||||
}),
|
||||
)}
|
||||
${this.view === "profiles" ? body : html`<openclaw-cloud-worker-snapshots></openclaw-cloud-worker-snapshots>`}
|
||||
${
|
||||
this.view === "profiles"
|
||||
? body
|
||||
: this.view === "pool"
|
||||
? html`<openclaw-cloud-worker-pool></openclaw-cloud-worker-pool>`
|
||||
: html`<openclaw-cloud-worker-snapshots></openclaw-cloud-worker-snapshots>`
|
||||
}
|
||||
`)}
|
||||
`;
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue