diff --git a/README.md b/README.md index 1ba7799..9cb2248 100644 --- a/README.md +++ b/README.md @@ -81,8 +81,10 @@ repo. (`ate.dev/v1alpha1`), point the gVisor `SandboxConfig` at a runsc build that survives a heavy multi-process Node.js actor, and apply the control-plane changes below. Actors use the **atespace** model - (`kubectl ate create atespace ; kubectl ate create actor -a --template-ref `) - and are reached at `..actors.resources.substrate.ate.dev`. + (`kubectl ate create atespace ; kubectl ate create actor -a --template-ref `, + `--template` on Substrate main) and are reached through the atenet router, addressed + by Host `..actors.resources.substrate.ate.dev` on release-0.1 and by + the `ate-target-actor: /` header on main. The gateway sends both. ### Substrate configuration this needs diff --git a/demo/dashboard/dashboard.js b/demo/dashboard/dashboard.js index 8ba7c23..90ea580 100644 --- a/demo/dashboard/dashboard.js +++ b/demo/dashboard/dashboard.js @@ -14,6 +14,7 @@ import { Hono } from "hono"; import { serve } from "@hono/node-server"; import { exec } from "node:child_process"; +import http from "node:http"; const app = new Hono(); @@ -22,6 +23,7 @@ const ATE_NS = "ate-system"; const ATE_ENDPOINT = process.env.ATE_ENDPOINT || "api.ate-system.svc.cluster.local:443"; const GATEWAY_URL = process.env.GATEWAY_URL || `http://openclaw-gateway.${NS}.svc.cluster.local:18789`; const KUBECTL_ATE = process.env.KUBECTL_ATE || "kubectl-ate"; +const ROUTER_URL = process.env.ROUTER_URL || "http://atenet-router.ate-system.svc.cluster.local:80"; // Atespace(s) the demo actors live in (current OSS actor model). Comma-separated. const ATESPACES = (process.env.ATESPACES || "openclaw-demo").split(",").map(s => s.trim()).filter(Boolean); @@ -444,7 +446,56 @@ app.post("/api/reset-view", (c) => { // and costs both of those, so the cap is the fleet rather than a round number. const FLEET_SIZE = 18; -const churn = { running: false, stop: false, until: 0, cycles: 0, refused: 0, reportedCycles: 0 }; +const churn = { running: false, stop: false, until: 0, cycles: 0, refused: 0, failed: 0, reportedCycles: 0 }; + +// GET /healthz on an actor through atenet, resolving to the HTTP status once +// the response headers arrive, or 0 if the request never got one. It carries +// both addresses atenet has used: release-0.1 routes by Host and main only by +// ate-target-actor, and each ignores the other. node:http rather than fetch, +// because fetch quietly replaces a Host it is given with the URL's. +function pingActor(name, atespace, timeoutMs) { + return new Promise((resolve) => { + const req = http.get( + new URL("/healthz", ROUTER_URL), + { + headers: { + Host: `${name}.${atespace}.actors.resources.substrate.ate.dev`, + "ate-target-actor": `${atespace}/${name}`, + }, + timeout: timeoutMs, + }, + (res) => { + res.resume(); + resolve(res.statusCode || 0); + } + ); + req.on("timeout", () => req.destroy(new Error("timeout"))); + req.on("error", () => resolve(0)); + }); +} + +// release-0.1's kubectl-ate spells the flag --template-ref and main's spells it +// --template (substrate#1536); this image pins the first, a main checkout builds +// the second. Try both, count "already exists" as success, and put any other +// failure on the timeline: an actor that was silently never created is what +// used to make burst report success while the pod map stayed empty. +function createFleetActor(name, atespace) { + const attempt = (flag) => + new Promise((resolve) => + exec( + `${KUBECTL_ATE} create actor ${name} -a ${atespace} ${flag} openclaw-agent`, + { timeout: 15000 }, + (error, _stdout, stderr) => resolve(error ? (stderr || error.message).trim() : "") + ) + ); + return (async () => { + let err = await attempt("--template-ref"); + if (/unknown flag/i.test(err)) err = await attempt("--template"); + if (!err || /already exists/i.test(err)) return true; + addEvent("substrate", `${name}: create failed: ${err.split("\n")[0]}`); + return false; + })(); +} function sleep(ms) { return new Promise((r) => setTimeout(r, ms)); @@ -477,21 +528,37 @@ function reportRefusals() { ); churn.refused = 0; } + if (churn.failed > 0) { + addEvent( + "substrate", + `${churn.failed} resume${churn.failed === 1 ? "" : "s"} failed with something other than a refusal, and retried` + ); + churn.failed = 0; + } } async function churnActor(name, atespace, hold) { - const url = `http://${name}.${atespace}.actors.resources.substrate.ate.dev/healthz`; // Start somewhere random inside the cycle so twenty loops do not fire on the // same tick. Without this they synchronise into a slow pulse, which looks // staged and hides the refusals. await sleep(Math.random() * hold * 2); while (!churn.stop && Date.now() < churn.until) { - let served = false; - try { - const r = await fetch(url, { signal: AbortSignal.timeout(30000) }); - if (r.ok) served = true; - else if (r.status === 503) churn.refused++; - } catch {} + const status = await pingActor(name, atespace, 30000); + const served = status >= 200 && status < 300; + if (status === 503) churn.refused++; + else if (status && !served) { + // Anything else is not the pool saying no. A 404 in particular is atenet + // not recognising the actor address, and every loop would otherwise + // retry it silently forever with an empty panel. Say so once per run. + if (!churn.reportedFailure) { + churn.reportedFailure = true; + // Only a 404 points at routing. A 504 is a restore that outran the + // router's timeout, which a cold node can do. + const hint = status === 404 ? "; check actor routing" : ""; + addEvent("substrate", `${name}: atenet answered HTTP ${status}, not a refusal${hint}`); + } + churn.failed++; + } if (!served) { // Refused or timed out. Back off, with jitter so the waiting actors do @@ -526,12 +593,9 @@ app.post("/api/churn", async (c) => { const names = []; for (let i = 1; i <= count; i++) { const name = `oc-agent-${i}`; - await runCmd( - `${KUBECTL_ATE} create actor ${name} -a ${atespace} --template-ref openclaw-agent 2>/dev/null || true`, - 15000 - ); - names.push(name); + if (await createFleetActor(name, atespace)) names.push(name); } + if (names.length === 0) return c.json({ ok: false, error: "no actors could be created" }, 500); Object.assign(churn, { running: true, @@ -539,6 +603,8 @@ app.post("/api/churn", async (c) => { until: Date.now() + seconds * 1000, cycles: 0, refused: 0, + failed: 0, + reportedFailure: false, reportedCycles: 0, }); addEvent( @@ -583,6 +649,8 @@ app.post("/api/churn", async (c) => { return c.json({ ok: true, count, seconds, hold, actors: names }); }); +// Only raises the flag. The loops see it within one cycle, and the sweep that +// runs when they finish parks the fleet, so stopping needs no sweep of its own. app.post("/api/churn/stop", (c) => { churn.stop = true; return c.json({ ok: true, wasRunning: churn.running, cycles: churn.cycles }); @@ -605,31 +673,20 @@ app.post("/api/burst", async (c) => { // wakes actors that were already sitting there suspended. It also keeps the // fleet panel readable on camera, where "oc-burst-3" looks like scaffolding. const name = `oc-agent-${i}`; - // Idempotent: create the actor from the golden template if it doesn't exist. - // - // The flag is --template-ref, and it resolves the name inside --atespace, so - // it takes a bare name. This used to pass `--template openclaw/openclaw-agent` - // -- a flag the CLI doesn't have, and a namespace-qualified reference it would - // reject anyway -- with the error swallowed by `|| true`. Burst then fired - // HTTP requests at actors that had never been created, and the pod map stayed - // empty while the button reported success. - await runCmd( - `${KUBECTL_ATE} create actor ${name} -a ${atespace} --template-ref openclaw-agent 2>/dev/null || true`, - 15000 - ); - names.push(name); + // Idempotent: create the actor from the golden template if it doesn't + // exist. Only actors that exist get a request fired at them. + if (await createFleetActor(name, atespace)) names.push(name); } // Fire resume-on-demand at each actor (async, so it doesn't block the HTTP response). - // atenet routes ..actors.resources.substrate.ate.dev to a worker. for (const name of names) { - const url = `http://${name}.${atespace}.actors.resources.substrate.ate.dev/healthz`; // This request is what causes the wake, and the response is the actor // serving, so the round trip is the resume-on-demand latency with nothing // inferred. It is the only place in the dashboard that can honestly time a // resume: everywhere else is reading a 2s poll. const t0 = Date.now(); - fetch(url, { signal: AbortSignal.timeout(120000) }) - .then((r) => { + pingActor(name, atespace, 120000) + .then((status) => { + const r = { status, ok: status >= 200 && status < 300 }; // Only a served response is a sample. A 503 measures how fast the pool // said no, and a 504 is the timeout, not the restore. if (r.ok) { @@ -658,14 +715,15 @@ app.post("/api/burst", async (c) => { // screen. Calling that "no worker free" contradicts the panel next to // it. addEvent("substrate", `${name}: restore outran the request timeout (HTTP 504); actor is still coming up`); + } else if (r.status === 0) { + addEvent("substrate", `${name}: resume request got no response from atenet`); } else if (!r.ok) { addEvent("substrate", `${name}: resume request failed HTTP ${r.status}`); } - }) - .catch(() => {}); + }); } - addEvent("substrate", `Burst: fired ${count} tasks, actors now multiplexing onto the worker pool`); - return c.json({ ok: true, count, actors: names }); + addEvent("substrate", `Burst: fired ${names.length} tasks, actors now multiplexing onto the worker pool`); + return c.json({ ok: names.length > 0, count: names.length, actors: names }); }); app.get("/", (c) => diff --git a/demo/deploy-demo.sh b/demo/deploy-demo.sh index b953959..838c024 100755 --- a/demo/deploy-demo.sh +++ b/demo/deploy-demo.sh @@ -75,6 +75,12 @@ echo " OK: kubectl, gcloud, Substrate, kubectl-ate, bucket." if [ -z "$GEMINI_API_KEY" ]; then read -rsp "Enter your Gemini API key: " GEMINI_API_KEY; echo "" fi +# Reuse the token from a previous run: the gateway pod only reads it at start, so +# rotating it here would leave the running gateway unable to call new actors. +if [ -z "${OPENCLAW_GATEWAY_TOKEN:-}" ]; then + OPENCLAW_GATEWAY_TOKEN=$(kubectl -n "$NAMESPACE" get secret openclaw-secrets \ + -o jsonpath='{.data.gateway-token}' 2>/dev/null | base64 -d 2>/dev/null || true) +fi GATEWAY_TOKEN="${OPENCLAW_GATEWAY_TOKEN:-$(openssl rand -hex 32)}" # --- [1/8] Build images (optional) --- @@ -160,6 +166,18 @@ render "$PARENT_DIR/manifests/workerpool.yaml" | kubectl apply -f - # (there is no update verb), so a re-run that changes the image, bucket or key # has to delete and recreate. kubectl ate create atespace "$ATESPACE" 2>/dev/null || echo " atespace $ATESPACE exists, continuing..." +# The demo actor goes first, because the template is always (re)created below +# and an actor that outlives its template points at a golden that is gone. +# Park it before deleting. A delete refuses an actor that is not suspended, and +# forcing it with --any-state has to tear down the live sandbox, which on gVisor +# can fail in `runsc delete` and leave the actor stuck DELETING. Suspend waits +# until the actor is parked and is a no-op on one that already is. --any-state +# stays for a CRASHED actor, which cannot be suspended. +if kubectl ate get actor "$ACTOR_NAME" -a "$ATESPACE" &>/dev/null; then + echo " Deleting existing actor $ACTOR_NAME so it is recreated from the new template..." + kubectl ate suspend actor "$ACTOR_NAME" -a "$ATESPACE" &>/dev/null || true + kubectl ate delete actor "$ACTOR_NAME" -a "$ATESPACE" --any-state +fi if kubectl ate get actor-template "$TEMPLATE" -a "$ATESPACE" &>/dev/null; then echo " Replacing existing actor template (templates are immutable)..." kubectl ate delete actor-template "$TEMPLATE" -a "$ATESPACE" @@ -205,10 +223,11 @@ golden_ready() { local deadline=$((SECONDS + 420)) json snapshot err while ((SECONDS < deadline)); do if json=$(kubectl ate get actor-template "$TEMPLATE" -a "$ATESPACE" -o json 2>/dev/null); then - # ExternalSnapshot identifies itself by URI; there is no name field. - snapshot=$(jq -r '.actorTemplates[0].status.goldenSnapshotStatus.goldenSnapshot.snapshotUri // empty' <<<"$json") + # main prints a single get bare and names the golden by tag; release-0.1 + # wraps it in .actorTemplates[0] and names it by snapshot URI. + snapshot=$(jq -r '(.status.goldenSnapshotStatus.goldenTag.name // .status.goldenSnapshotStatus.goldenSnapshot.snapshotUri // .actorTemplates[0].status.goldenSnapshotStatus.goldenSnapshot.snapshotUri // empty)' <<<"$json") [ -n "$snapshot" ] && { echo " golden snapshot ready: $snapshot"; return 0; } - err=$(jq -r '.actorTemplates[0].status.goldenSnapshotStatus.errorMessage // empty' <<<"$json") + err=$(jq -r '(.status.goldenSnapshotStatus.errorMessage // .actorTemplates[0].status.goldenSnapshotStatus.errorMessage // empty)' <<<"$json") [ -n "$err" ] && { echo " golden snapshot FAILED: $err" >&2; return 1; } fi sleep 5 @@ -218,8 +237,19 @@ golden_ready() { } golden_ready || echo " (continuing anyway, the actor will resume once the golden is ready)" # The template is resolved in the actor's atespace, so both live in $ATESPACE. -kubectl ate create actor "$ACTOR_NAME" --template-ref "$TEMPLATE" --atespace "$ATESPACE" 2>/dev/null \ - || echo " actor exists, continuing..." +# release-0.1's kubectl-ate spells the flag --template-ref and main's spells it +# --template (substrate#1536), so try both. Only "already exists" is fine; any +# other failure stops here rather than leaving a demo with no actor in it. +create_actor() { + local out + out=$(kubectl ate create actor "$ACTOR_NAME" --template-ref "$TEMPLATE" --atespace "$ATESPACE" 2>&1) && return 0 + if grep -qi "unknown flag" <<<"$out"; then + out=$(kubectl ate create actor "$ACTOR_NAME" --template "$TEMPLATE" --atespace "$ATESPACE" 2>&1) && return 0 + fi + grep -qi "already exists" <<<"$out" && { echo " actor exists, continuing..."; return 0; } + die "creating actor $ACTOR_NAME failed: $out" +} +create_actor # --- [7b] Optional cron status pings (run on the always-on gateway agent) --- WHATSAPP_PEER="${WHATSAPP_PEER:-}" diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index a077d4b..1328709 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -236,9 +236,12 @@ Targets **current OSS Substrate** (`agent-substrate/substrate`, CRD group `ate.d - `WorkerPool` (`ate.dev/v1alpha1`, gVisor `SandboxClass`) + `ActorTemplate` (golden snapshot to GCS). runsc comes from the cluster gVisor `SandboxConfig`. - Actors use the **atespace** model: `kubectl ate create atespace ` then - `kubectl ate create actor -a --template-ref openclaw-agent`, and are - addressed at `..actors.resources.substrate.ate.dev` (atenet routes - restore-on-demand to worker port 80). + `kubectl ate create actor -a --template-ref openclaw-agent` (`--template` on + Substrate main), and are addressed through the atenet router, which routes + restore-on-demand to worker port 80. release-0.1 picks the actor from Host + `..actors.resources.substrate.ate.dev`; main dropped that and picks + it from the `ate-target-actor: /` header. The gateway and dashboard + send both, so either works. - Real-time dashboard: lists live actors (via `kubectl-ate` per atespace), the golden-snapshot status, the worker-pod map (which actor is restored where), and live gateway/WhatsApp status. diff --git a/extensions/substrate/README.md b/extensions/substrate/README.md index f093c0e..a0947d8 100644 --- a/extensions/substrate/README.md +++ b/extensions/substrate/README.md @@ -64,13 +64,16 @@ serving `/v1/chat/completions`; it does not load this plugin. | `atespace` | *(none)* | atespace the per-conversation actors are created in | | `template` | *(none)* | golden ActorTemplate new actors derive from, by bare name | | `templateForAgent` | *(none)* | optional per-persona template override, keyed by agentId | -| `actorDomain` | `actors.resources.substrate.ate.dev` | DNS domain actors are addressed under | +| `actorDomain` | `actors.resources.substrate.ate.dev` | domain in the `Host` sent to atenet (release-0.1 routes by it; main routes by the `ate-target-actor` header, which is always sent too) | | `actorToken` | *(none)* | bearer token for the actor's HTTP API | | `provisioner` | `ateapi` | `ateapi` (in-band gRPC, needs a podcert) or `kubectl-ate` (shell out) | | `kubectlAtePath` | `kubectl-ate` | path to the binary when `provisioner: "kubectl-ate"` | | `idleTimeoutSeconds` | 120 | idle time before the gateway suspends the actor | | `ateapiAddress` | `api.ate-system.svc.cluster.local:443` | Substrate control plane gRPC | +Turns go to the atenet router, `http://atenet-router.ate-system.svc.cluster.local:80` +unless the gateway's `ROUTER_URL` environment variable overrides it. + ## Files | File | Purpose | diff --git a/extensions/substrate/acp-runtime.ts b/extensions/substrate/acp-runtime.ts index d358623..b6c1d94 100644 --- a/extensions/substrate/acp-runtime.ts +++ b/extensions/substrate/acp-runtime.ts @@ -27,7 +27,7 @@ import type { AcpRuntimeTurnInput, AcpRuntimeTurnResult, } from "openclaw/plugin-sdk/acp-runtime-backend"; -import { actorNameForConversation, actorUrlFor, DEFAULT_ACTOR_DOMAIN } from "./actor-router.js"; +import { actorNameForConversation, DEFAULT_ACTOR_DOMAIN, fetchActor, type ActorRef } from "./actor-router.js"; import type { Provisioner } from "./actor-provisioner.js"; export type SubstrateAcpRuntimeConfig = { @@ -75,9 +75,12 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac const templateFor = (agent?: string): string => (agent && config.templateForAgent?.[agent]) || config.template; // Placement is a pure function of the conversation key, so both ensureSession - // and startTurn derive the same actor URL with no side map. - const urlForSession = (sessionKey: string): string => - actorUrlFor(actorNameForConversation(sessionKey), config.atespace, domain); + // and startTurn derive the same actor with no side map. + const refForSession = (sessionKey: string): ActorRef => ({ + name: actorNameForConversation(sessionKey), + atespace: config.atespace, + domain, + }); // The gateway's ACP session key uses reserved internal namespaces // (e.g. "agent:main:acp:binding:..."), which the actor's // /v1/chat/completions rejects via X-OpenClaw-Session-Key ("reserved @@ -85,6 +88,11 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac // key the in-actor session by the (stable, non-reserved) actor name instead. const actorSessionKey = (sessionKey: string): string => actorNameForConversation(sessionKey); + const killSession = (sessionKey: string): Promise => + fetchActor(refForSession(sessionKey), `/sessions/${encodeURIComponent(actorSessionKey(sessionKey))}/kill`, { + method: "POST", + headers: authHeaders(), + }).catch(() => {}); const runtime: AcpRuntime = { async ensureSession(input: AcpRuntimeEnsureInput): Promise { @@ -99,16 +107,12 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac backend: "substrate", runtimeSessionName: input.sessionKey, cwd: input.cwd, - // Stash the resolved URL on the handle; startTurn also recomputes it. - ...({ actorUrl: urlForSession(input.sessionKey) } as object), }; }, startTurn(input: AcpRuntimeTurnInput): AcpRuntimeTurn { - // Deterministic from the conversation key; handle carries it as a fast path. - const baseUrl = - (input.handle as { actorUrl?: string }).actorUrl ?? urlForSession(input.handle.sessionKey); - const turnActor = actorNameForConversation(input.handle.sessionKey); + const actor = refForSession(input.handle.sessionKey); + const turnActor = actor.name; config.onTurnStart?.(turnActor); const abort = new AbortController(); input.signal?.addEventListener("abort", () => abort.abort(input.signal?.reason)); @@ -116,7 +120,7 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac const result = new Promise((r) => (resolveResult = r)); result.finally(() => config.onTurnEnd?.(turnActor)).catch(() => {}); const events = streamTurn( - baseUrl, + actor, input, actorSessionKey(input.handle.sessionKey), authHeaders(), @@ -130,10 +134,7 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac result, async cancel() { abort.abort("cancelled"); - await fetch( - `${baseUrl}/sessions/${encodeURIComponent(actorSessionKey(input.handle.sessionKey))}/kill`, - { method: "POST", headers: authHeaders() }, - ).catch(() => {}); + await killSession(input.handle.sessionKey); }, async closeStream() { abort.abort("stream closed"); @@ -150,12 +151,7 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac }, async cancel(input) { - const baseUrl = - (input.handle as { actorUrl?: string }).actorUrl ?? urlForSession(input.handle.sessionKey); - await fetch( - `${baseUrl}/sessions/${encodeURIComponent(actorSessionKey(input.handle.sessionKey))}/kill`, - { method: "POST", headers: authHeaders() }, - ).catch(() => {}); + await killSession(input.handle.sessionKey); }, async close() { @@ -166,7 +162,7 @@ export function createSubstrateAcpRuntime(config: SubstrateAcpRuntimeConfig): Ac } async function* streamTurn( - baseUrl: string, + actor: ActorRef, input: AcpRuntimeTurnInput, sessionKey: string, headers: Record, @@ -176,7 +172,7 @@ async function* streamTurn( ): AsyncIterable { let res: Response; try { - res = await fetch(`${baseUrl}/v1/chat/completions`, { + res = await fetchActor(actor, "/v1/chat/completions", { method: "POST", headers: { "Content-Type": "application/json", diff --git a/extensions/substrate/actor-router.ts b/extensions/substrate/actor-router.ts index b93059c..b101f6c 100644 --- a/extensions/substrate/actor-router.ts +++ b/extensions/substrate/actor-router.ts @@ -23,8 +23,17 @@ * its state persists across suspends in that actor's snapshot. */ import { createHash } from "node:crypto"; +import http from "node:http"; +import https from "node:https"; +import { Readable } from "node:stream"; export const DEFAULT_ACTOR_DOMAIN = "actors.resources.substrate.ate.dev"; +export const DEFAULT_ROUTER_URL = "http://atenet-router.ate-system.svc.cluster.local:80"; + +/** The atenet router, overridable with ROUTER_URL. */ +export function routerUrl(): string { + return process.env.ROUTER_URL || DEFAULT_ROUTER_URL; +} /** conv-: deterministic, DNS-safe, stable. */ export function actorNameForConversation(canonicalSessionKey: string): string { @@ -32,7 +41,66 @@ export function actorNameForConversation(canonicalSessionKey: string): string { return `conv-${h}`; } -/** atenet Host-routed URL for an actor. atenet routes by Host to actor port 80. */ -export function actorUrlFor(name: string, atespace: string, domain: string = DEFAULT_ACTOR_DOMAIN): string { - return `http://${name}.${atespace}.${domain}`; +/** The hostname release-0.1 atenet routes an actor by. */ +export function actorHostFor(name: string, atespace: string, domain: string = DEFAULT_ACTOR_DOMAIN): string { + return `${name}.${atespace}.${domain}`; +} + +export type ActorRef = { name: string; atespace: string; domain?: string }; + +/** + * Sends a request to an actor through the atenet router, addressed both ways + * atenet has understood. release-0.1 routes by Host + * (`..`, which its CoreDNS stub resolves to the router) + * and ignores the header. main dropped that DNS and routes only by + * `ate-target-actor: /`, ignoring Host. Sending both to the + * router service works on either. + * + * node:http rather than fetch, because fetch silently replaces a caller's Host + * with the URL's, and release-0.1 would then route nothing. Connection failures + * reject with a TypeError, as fetch's do, so callers keep treating them as + * retryable. + */ +export function fetchActor( + ref: ActorRef, + path: string, + init: { method?: string; headers?: Record; body?: string; signal?: AbortSignal } = {}, +): Promise { + const url = new URL(path, routerUrl()); + const client = url.protocol === "https:" ? https : http; + return new Promise((resolve, reject) => { + const req = client.request( + url, + { + method: init.method ?? "GET", + headers: { + ...init.headers, + Host: actorHostFor(ref.name, ref.atespace, ref.domain), + "ate-target-actor": `${ref.atespace}/${ref.name}`, + }, + signal: init.signal, + }, + (res) => { + const headers = new Headers(); + for (const [k, v] of Object.entries(res.headers)) { + if (v !== undefined) headers.set(k, Array.isArray(v) ? v.join(", ") : v); + } + const status = res.statusCode ?? 502; + const nullBody = [101, 204, 205, 304].includes(status); + if (nullBody) res.resume(); + resolve( + new Response(nullBody ? null : (Readable.toWeb(res) as ReadableStream), { + status, + statusText: res.statusMessage, + headers, + }), + ); + }, + ); + req.on("error", (err) => { + if (init.signal?.aborted) reject(err); + else reject(new TypeError(`actor request failed: ${err.message}`, { cause: err })); + }); + req.end(init.body); + }); } \ No newline at end of file diff --git a/extensions/substrate/kubectl-ate-client.ts b/extensions/substrate/kubectl-ate-client.ts index b504508..936dd68 100644 --- a/extensions/substrate/kubectl-ate-client.ts +++ b/extensions/substrate/kubectl-ate-client.ts @@ -44,13 +44,26 @@ export function createKubectlAteClient(cfg: { binPath?: string; endpoint?: strin await run([ "create", "actor", ref.name, "-a", ref.atespace, - // `--template /` is gone. ActorTemplate stopped being a - // Kubernetes CRD, so there is no namespace to qualify it with: - // --template-ref resolves the name inside the actor's own atespace. + // release-0.1 spells it --template-ref; main renamed it to --template + // (substrate#1536). Both resolve a bare name inside the actor's atespace. "--template-ref", templateName, ]); } catch (err) { const msg = String((err as { stderr?: string })?.stderr ?? (err as Error)?.message ?? err); + if (/unknown flag: --template-ref/i.test(msg)) { + try { + await run([ + "create", "actor", ref.name, + "-a", ref.atespace, + "--template", templateName, + ]); + return; + } catch (err2) { + const msg2 = String((err2 as { stderr?: string })?.stderr ?? (err2 as Error)?.message ?? err2); + if (/already exists/i.test(msg2)) return; + throw err2; + } + } if (/already exists/i.test(msg)) return; // idempotent throw err; } diff --git a/manifests/actortemplate.yaml b/manifests/actortemplate.yaml index b47166b..6e8f13d 100644 --- a/manifests/actortemplate.yaml +++ b/manifests/actortemplate.yaml @@ -31,8 +31,9 @@ metadata: # The atespace replaces the old metadata.namespace, and it must already exist. # It has to be the SAME atespace the actors live in: `kubectl ate create actor - # --template-ref` resolves the template in the actor's atespace, so a template - # parked in the k8s namespace instead is simply not found. + # --template-ref` (`--template` on main) resolves the template in the actor's + # atespace, so a template parked in the k8s namespace instead is simply not + # found. atespace: openclaw-demo name: openclaw-agent workerSelector: