Files
pikasTech-HWLAB/internal/cloud/code-agent-agentrun-adapter.ts
T
2026-06-08 11:25:54 +08:00

1987 lines
88 KiB
TypeScript

import { randomUUID } from "node:crypto";
import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts";
import { decorateCodeAgentSession } from "./code-agent-session-lifecycle.ts";
import {
parsePositiveInteger,
safeConversationId,
safeOpaqueId,
safeSessionId,
safeTraceId,
text,
truthyFlag
} from "./server-http-utils.ts";
const ADAPTER_ID = "agentrun-v01";
const DEFAULT_AGENTRUN_MGR_URL = "http://agentrun-mgr.agentrun-v01.svc.cluster.local:8080";
const DEFAULT_TENANT_ID = "hwlab";
const DEFAULT_PROJECT_ID = "pikasTech/HWLAB";
const DEFAULT_PROVIDER_ID = "G14";
const DEFAULT_REPO_URL = "http://git-mirror-http.devops-infra.svc.cluster.local/pikasTech/HWLAB.git";
const DEFAULT_RUNNER_NAMESPACE = "agentrun-v01";
const DEFAULT_TIMEOUT_MS = 1_200_000;
const AGENTRUN_BACKEND_PREFIX = ADAPTER_ID;
const AGENTRUN_RUNNER_KIND = "agentrun-v01-shared-runner";
const AGENTRUN_SESSION_MODE = "agentrun-v01-durable-session";
const AGENTRUN_IMPLEMENTATION_TYPE = "agentrun-v01-shared-execution-infra";
const AGENTRUN_CAPABILITY_LEVEL = "agentrun-v01-shared-code-agent-session";
const AGENTRUN_PROVIDER_TRACE_PROTOCOL = "agentrun-v01-jsonrpc";
const AGENTRUN_PROVIDER_TRACE_WIRE_API = "agentrun-v01-command-result";
const AGENTRUN_BACKEND_ALIASES = Object.freeze({ "codex-api": "codex", codex: "codex" });
const AGENTRUN_BACKEND_PROFILE_ID_PATTERN = /^[a-z0-9][a-z0-9-]{0,63}$/u;
const THREAD_CONTINUITY_POLICY = "hwlab-agentrun-v01-reuse-runner-thread";
const SESSION_POLICY_RUN_LOCAL = "hwlab-agentrun-v01-session-runner-reuse";
const TERMINAL_RUN_STATUSES = new Set(["completed", "failed", "blocked", "cancelled", "canceled"]);
const HWLAB_RESOURCE_GIT_BUNDLES = Object.freeze([
Object.freeze({ name: "hwlab-tools", subpath: "tools", target_path: "tools" }),
Object.freeze({ name: "hwlab-agent-skills", subpath: "skills", target_path: ".agents/skills" })
]);
const HWLAB_RESOURCE_PROMPT_REFS = Object.freeze([
Object.freeze({ name: "hwlab-v02-runtime", path: "internal/agent/prompts/hwlab-v02-runtime.md", inject: "thread-start", required: true })
]);
export function codeAgentAgentRunAdapterEnabled(env = process.env) {
const value = String(env.HWLAB_CODE_AGENT_ADAPTER ?? env.HWLAB_CODE_AGENT_PROVIDER ?? "").trim().toLowerCase();
return value === ADAPTER_ID || value === "agentrun" || value === "agentrun-v0.1";
}
export function describeAgentRunAdapterAvailability(env = process.env, options = {}) {
const params = options.params ?? {};
const backendProfile = resolveAgentRunBackendProfile(env, params);
const provider = providerForBackendProfile(backendProfile);
const backend = `${ADAPTER_ID}/${backendProfile}`;
const model = modelForBackendProfile(backendProfile, env);
const blockers = [];
let managerUrl = null;
let repoUrl = null;
let resourceBundleSourceCommit = null;
try {
managerUrl = resolveAgentRunManagerUrl(env);
} catch (error) {
blockers.push(agentRunAvailabilityBlocker(error, "manager-url"));
}
try {
repoUrl = resolveAgentRunRepoUrl(env);
} catch (error) {
blockers.push(agentRunAvailabilityBlocker(error, "repo-url"));
}
try {
resourceBundleSourceCommit = requireAgentRunSourceCommit(env);
} catch (error) {
blockers.push(agentRunAvailabilityBlocker(error, "resource-bundle-source-commit"));
}
const providerId = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID);
const runnerNamespace = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_NAMESPACE, env.AGENTRUN_RUNTIME_NAMESPACE, DEFAULT_RUNNER_NAMESPACE);
const ready = blockers.length === 0 && Boolean(providerId && runnerNamespace && backendProfile && managerUrl && repoUrl);
const managerHost = hostForUrl(managerUrl);
const repoHost = hostForUrl(repoUrl);
const status = ready ? "agentrun-ready" : "blocked";
const capabilityLevel = ready ? AGENTRUN_CAPABILITY_LEVEL : "blocked";
return {
endpoint: "POST /v1/agent/chat",
adapter: ADAPTER_ID,
provider,
model,
backend,
infrastructureBackend: backend,
mode: ADAPTER_ID,
schema: [
"conversationId",
"sessionId",
"messageId",
"status",
"createdAt",
"updatedAt",
"traceId",
"provider",
"model",
"backend",
"infrastructureBackend",
"agentRun",
"runner",
"runnerTrace",
"session",
"sessionMode",
"sessionReuse",
"threadContinuityPolicy",
"error.message",
"error.code"
],
runner: {
kind: AGENTRUN_RUNNER_KIND,
adapter: ADAPTER_ID,
backend,
provider,
model,
namespace: runnerNamespace,
managerHost,
repoHost,
mode: ADAPTER_ID,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
status: ready ? "available" : "blocked",
ready,
capabilityLevel,
longLivedSession: ready,
durableSession: ready,
durable: ready,
codexStdio: false,
delegatedToAgentRun: true,
writeCapable: ready,
readOnly: false,
runnerLimitations: ready ? [] : ["agentrun-adapter-config-blocked"],
secretMaterialRead: false,
valuesRedacted: true
},
agentRun: {
adapter: ADAPTER_ID,
backendProfile,
providerId,
runnerNamespace,
managerUrl,
managerHost,
repoUrl,
repoHost,
resourceBundleSourceCommit,
managerConfigured: Boolean(managerUrl),
repoConfigured: Boolean(repoUrl),
internalServiceDns: managerHost === "agentrun-mgr.agentrun-v01.svc.cluster.local",
ready,
valuesPrinted: false
},
status,
readinessStatus: status,
agentKind: "agentrun-v01-adapter",
capabilityStatus: status,
blocker: ready ? null : blockers[0]?.summary ?? "AgentRun adapter runtime configuration is blocked.",
reason: ready ? null : blockers[0]?.code ?? "agentrun_adapter_blocked",
summary: ready
? "HWLAB Code Agent delegates execution to AgentRun v0.1 shared infrastructure; hwlab-cloud-api keeps session/trace/API ownership and no longer manages a repo-owned runner control plane."
: "HWLAB Code Agent is configured for AgentRun v0.1, but the adapter runtime configuration is incomplete or invalid.",
missingEnv: [],
secretRefs: [],
egress: {
ready,
mode: "agentrun-managed-provider-egress",
directPublicOpenAi: false,
secretMaterialRead: false,
valuesRedacted: true
},
safety: {
secretsRead: false,
secretValuesPrinted: false,
kubeconfigRead: false,
providerCredentialsDelegatedToAgentRun: true,
repoOwnedRunnerControlPlane: false,
valuesRedacted: true
},
capabilityLevel,
workspace: repoUrl,
sandbox: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_SANDBOX, env.HWLAB_CODE_AGENT_CODEX_SANDBOX, "danger-full-access"),
session: null,
sessionStatus: null,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
runnerLimitations: ready ? [] : ["agentrun-adapter-config-blocked"],
longLivedSessionGate: {
status: ready ? "pass" : "blocked",
provider,
runnerKind: AGENTRUN_RUNNER_KIND,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
adapter: ADAPTER_ID,
codexStdio: false,
delegatedToAgentRun: true,
blockers,
secretMaterialRead: false,
valuesRedacted: true
},
blockers,
blockerCodes: blockers.map((blocker) => blocker.code).filter(Boolean),
ready,
partialReady: false
};
}
export function initialAgentRunChatResult({ params = {}, options = {}, traceId }) {
const env = options.env ?? process.env;
const timestamp = nowIso(options.now);
const backendProfile = resolveAgentRunBackendProfile(env, params);
return {
accepted: true,
status: "running",
shortConnection: true,
traceId,
conversationId: safeConversationId(params.conversationId) || null,
sessionId: safeSessionId(params.sessionId) || agentRunSessionId(traceId),
threadId: safeOpaqueId(params.threadId) || null,
messageId: `msg_${safeTraceId(traceId)?.slice(4) || randomUUID()}`,
createdAt: timestamp,
updatedAt: timestamp,
provider: providerForBackendProfile(backendProfile),
model: modelForBackendProfile(backendProfile, env),
backend: backendForBackendProfile(backendProfile),
infrastructureBackend: `agentrun-v01/${backendProfile}`,
capabilityLevel: AGENTRUN_CAPABILITY_LEVEL,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
threadContinuityPolicy: THREAD_CONTINUITY_POLICY,
agentRun: {
adapter: ADAPTER_ID,
status: "pending-dispatch",
backendProfile,
providerId: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID),
managerUrl: resolveAgentRunManagerUrl(env),
repoUrl: resolveAgentRunRepoUrl(env),
resourceBundleSourceCommit: requireAgentRunSourceCommit(env),
valuesPrinted: false
},
valuesPrinted: false
};
}
export async function submitAgentRunChatTurn({ params = {}, options = {}, traceId, traceStore = defaultCodeAgentTraceStore, results }) {
const env = options.env ?? process.env;
const managerUrl = resolveAgentRunManagerUrl(env);
const backendProfile = resolveAgentRunBackendProfile(env, params);
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
const startedAt = nowIso(options.now);
const ownerApiKey = await resolveOwnerApiKey({ params, options, now: () => nowIso(options.now) });
const toolCapabilities = await resolveToolCapabilities({ params, options });
traceStore.ensure(traceId, agentRunTraceMeta(env, params));
traceStore.append(traceId, {
type: "request",
status: "accepted",
label: "agentrun:request:accepted",
message: "HWLAB Code Agent request accepted by the AgentRun v0.1 adapter; hwlab-cloud-api will reuse an active runner when the HWLAB session has an active AgentRun reuse window, otherwise it will create run/command/runner-job over the k3s Service DNS.",
waitingFor: "agentrun-run-reuse-or-create",
adapter: ADAPTER_ID,
managerHost: new URL(managerUrl).hostname,
valuesPrinted: false
}, agentRunTraceMeta(env, params));
const reusable = await resolveReusableAgentRun({ params, options, env, managerUrl, fetchImpl, timeoutMs, backendProfile, traceId, traceStore });
if (reusable?.mapping) {
try {
const commandInput = buildAgentRunCommandInput({ params, traceId, backendProfile, sessionId: reusable.sessionId });
const command = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(reusable.mapping.runId)}/commands`, {
method: "POST",
body: commandInput,
timeoutMs
});
const commandId = requiredString(command?.id, "command.id");
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:command:created",
message: `AgentRun command ${commandId} created on reused run ${reusable.mapping.runId}; hwlab-cloud-api will ensure a runner Job is available for this turn.`,
runId: reusable.mapping.runId,
commandId,
backendProfile,
waitingFor: "agentrun-runner-job-ensure",
valuesPrinted: false
});
let runnerJob = null;
try {
const runnerJobInput = buildAgentRunRunnerJobInput({ env, traceId, commandId, ownerApiKey, toolCapabilities, backendProfile });
runnerJob = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(reusable.mapping.runId)}/runner-jobs`, {
method: "POST",
body: runnerJobInput,
timeoutMs
});
} catch (error) {
if (!isAgentRunCommandAlreadyClaimed(error)) throw error;
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:runner-job:already-active",
message: `AgentRun command ${commandId} is already claimed by an active runner; no replacement runner Job is needed for this turn.`,
errorCode: error?.code ?? "agentrun_command_already_claimed",
runId: reusable.mapping.runId,
commandId,
runnerId: reusable.mapping.runnerId ?? null,
jobName: reusable.mapping.jobName ?? null,
namespace: reusable.mapping.namespace ?? null,
waitingFor: "agentrun-result",
valuesPrinted: false
});
}
const mapping = agentRunReusedMapping({ previous: reusable.mapping, run: reusable.run, command, runnerJob, traceId, startedAt, backendProfile, managerUrl, env });
traceStore.append(traceId, {
type: "backend",
status: "running",
label: runnerJob ? "agentrun:runner-job:ensured" : "agentrun:runner-job:reused",
message: runnerJob
? `AgentRun runner Job ${mapping.jobName ?? "unknown"} ensured for reused run ${mapping.runId}; this keeps the persistent session resumable after pod replacement.`
: `AgentRun runner Job ${mapping.jobName ?? "unknown"} is already active for this HWLAB session turn.`,
runId: mapping.runId,
commandId: mapping.commandId,
attemptId: mapping.attemptId,
runnerId: mapping.runnerId,
jobName: mapping.jobName,
namespace: mapping.namespace,
waitingFor: "agentrun-result",
valuesPrinted: false
});
return decorateAgentRunRunningResult({ base: initialAgentRunChatResult({ params, options, traceId }), mapping, traceStore, traceId });
} catch (error) {
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reuse-command-failed",
message: `AgentRun reused run ${reusable.mapping.runId} rejected the new command; hwlab-cloud-api will create a fresh run/runner Job for this turn.`,
errorCode: error?.code ?? "agentrun_reuse_command_failed",
runId: reusable.mapping.runId,
commandId: reusable.mapping.commandId ?? null,
waitingFor: "agentrun-run-create",
valuesPrinted: false
});
}
}
const baseSessionId = scopedAgentRunSessionIdForParams(params, traceId, backendProfile);
let sessionId = baseSessionId;
let sessionReset = false;
let run, command, runnerJob, mapping;
for (let attempt = 0; attempt < 2; attempt += 1) {
if (sessionReset) {
sessionId = newSessionIdAfterEviction(baseSessionId, traceId);
const resetParams = { ...params, threadId: null };
traceStore.append(traceId, {
type: "backend", status: "running",
label: "agentrun:session-reset",
message: "AgentRun session storage was evicted; hwlab-cloud-api is creating a fresh sessionId " + sessionId + " with threadId=null.",
sessionId, previousSessionId: baseSessionId, backendProfile, waitingFor: "agentrun-run-create", valuesPrinted: false,
});
}
try {
await ensureAgentRunSessionPersistent({ fetchImpl, managerUrl, sessionId, env, traceId, backendProfile, traceStore });
} catch (error) {
if (attempt === 0 && await shouldResetSessionAfterEviction("session-store-evicted", error?.message)) {
sessionReset = true;
continue;
}
throw error;
}
const runInput = buildAgentRunCreateRunInput({ params, env, traceId, backendProfile, sessionId, toolCapabilities });
run = await agentRunJson(fetchImpl, managerUrl, "/api/v1/runs", { method: "POST", body: runInput, timeoutMs });
const runId = requiredString(run?.id, "run.id");
traceStore.append(traceId, {
type: "backend", status: "running",
label: "agentrun:run:created",
message: "AgentRun run " + runId + " created through internal k3s Service DNS.",
runId, backendProfile, waitingFor: "agentrun-command-create", valuesPrinted: false,
});
const commandInput = buildAgentRunCommandInput({ params, traceId, backendProfile, sessionId });
command = await agentRunJson(fetchImpl, managerUrl, "/api/v1/runs/" + encodeURIComponent(runId) + "/commands", { method: "POST", body: commandInput, timeoutMs });
const commandId = requiredString(command?.id, "command.id");
traceStore.append(traceId, {
type: "backend", status: "running",
label: "agentrun:command:created",
message: "AgentRun command " + commandId + " created; hwlab-cloud-api will start a runner Job explicitly without relying on scheduler automation.",
runId, commandId, backendProfile, waitingFor: "agentrun-runner-job-create", valuesPrinted: false,
});
const runnerJobInput = buildAgentRunRunnerJobInput({ env, traceId, commandId, ownerApiKey, toolCapabilities, backendProfile });
try {
runnerJob = await agentRunJson(fetchImpl, managerUrl, "/api/v1/runs/" + encodeURIComponent(runId) + "/runner-jobs", { method: "POST", body: runnerJobInput, timeoutMs });
} catch (error) {
if (attempt === 0 && await shouldResetSessionAfterEviction("session-store-evicted", error?.message)) {
sessionReset = true;
continue;
}
throw error;
}
mapping = agentRunMapping({ env, managerUrl, backendProfile, run, command, runnerJob, traceId, startedAt, params });
break;
}
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:runner-job:created",
message: `AgentRun runner Job ${mapping.jobName ?? "unknown"} created in namespace ${mapping.namespace ?? DEFAULT_RUNNER_NAMESPACE}.`,
runId: mapping.runId,
commandId: mapping.commandId,
attemptId: mapping.attemptId,
runnerId: mapping.runnerId,
jobName: mapping.jobName,
namespace: mapping.namespace,
waitingFor: "agentrun-result",
valuesPrinted: false
});
return decorateAgentRunRunningResult({ base: initialAgentRunChatResult({ params, options, traceId }), mapping, traceStore, traceId });
}
function isAgentRunCommandAlreadyClaimed(error) {
const statusCode = Number(error?.statusCode ?? 0);
const message = String(error?.message ?? error?.agentRunError?.message ?? "");
return statusCode === 409 && /command\s+[^\s]+\s+is not pending:/u.test(message);
}
export async function syncAgentRunChatResult({ traceId, currentResult = null, options = {}, traceStore = defaultCodeAgentTraceStore, appendResultEvent = true, refreshEvents = true }) {
const initial = currentResult ?? await loadPersistedAgentRunResult(traceId, options);
const mapped = initial ? await resolveAgentRunTraceCommandMapping({ traceId, mapped: initial, options }) : initial;
if (!mapped?.agentRun?.runId || !mapped?.agentRun?.commandId) return { result: currentResult, runnerTrace: traceStore.snapshot(traceId), found: Boolean(currentResult) };
const env = options.env ?? process.env;
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const managerUrl = resolveAgentRunManagerUrl(env, mapped.agentRun.managerUrl);
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
const eventsResponse = refreshEvents ? await fetchAgentRunEventsForTrace({ fetchImpl, managerUrl, timeoutMs, mapping: { ...mapped.agentRun, traceSummary: mapped.traceSummary } }) : null;
if (eventsResponse) appendAgentRunEventsToTrace(traceStore, traceId, eventsResponse.events, mapped.agentRun);
const result = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(mapped.agentRun.runId)}/commands/${encodeURIComponent(mapped.agentRun.commandId)}/result`, {
method: "GET",
timeoutMs
});
const nextMapping = {
...mapped.agentRun,
...agentRunResultRefs(result),
lastSeq: eventsResponse ? agentRunTraceCursorSeq(eventsResponse, mapped.agentRun.lastSeq) : mapped.agentRun.lastSeq,
status: result?.status ?? mapped.agentRun.status ?? "running",
runStatus: result?.runStatus ?? mapped.agentRun.runStatus ?? null,
commandState: result?.commandState ?? mapped.agentRun.commandState ?? null,
terminalStatus: result?.terminalStatus ?? mapped.agentRun.terminalStatus ?? null,
updatedAt: nowIso(options.now),
valuesPrinted: false
};
const base = { ...mapped, agentRun: nextMapping, updatedAt: nowIso(options.now) };
const payload = agentRunResultToCodeAgentPayload({ base, result, traceStore, traceId, appendResultEvent });
options.codeAgentChatResults?.set?.(traceId, payload);
return { result: payload, runnerTrace: payload.runnerTrace ?? traceStore.snapshot(traceId), found: true };
}
async function resolveAgentRunTraceCommandMapping({ traceId, mapped, options = {} }) {
const safeId = safeTraceId(traceId);
const agentRun = mapped?.agentRun && typeof mapped.agentRun === "object" ? mapped.agentRun : null;
if (!safeId || !agentRun?.runId) return mapped;
const agentRunTraceId = safeTraceId(agentRun.traceId);
const providerTraceId = safeTraceId(agentRun.providerTrace?.traceId);
const traceSummaryCommandId = text(mapped?.traceSummary?.agentRun?.commandId);
const commandId = text(agentRun.commandId);
const traceFieldsMatch = agentRunTraceId === safeId && providerTraceId === safeId && (!traceSummaryCommandId || traceSummaryCommandId === commandId);
if (mapped.status === "running" && traceFieldsMatch) return mapped;
const env = options.env ?? process.env;
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const managerUrl = resolveAgentRunManagerUrl(env, agentRun.managerUrl);
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
const command = await findAgentRunCommandForTrace({ fetchImpl, managerUrl, timeoutMs, runId: agentRun.runId, traceId: safeId });
if (!command?.id) {
throw Object.assign(new Error(`AgentRun command registry has no command for ${safeId} in run ${agentRun.runId}`), {
code: "agentrun_trace_command_not_found",
statusCode: 404,
traceId: safeId,
runId: agentRun.runId
});
}
if (command.id === commandId) {
return {
...mapped,
finalResponse: null,
traceSummary: null,
agentRun: {
...agentRun,
traceId: safeId,
lastSeq: 0,
providerTrace: null,
commandState: command.state ?? agentRun.commandState ?? null,
updatedAt: nowIso(options.now),
valuesPrinted: false
}
};
}
return {
...mapped,
finalResponse: null,
traceSummary: null,
agentRun: {
...agentRun,
commandId: command.id,
traceId: safeId,
lastSeq: 0,
providerTrace: null,
commandState: command.state ?? agentRun.commandState ?? null,
updatedAt: nowIso(options.now),
valuesPrinted: false
}
};
}
export async function refreshAgentRunTrace({ traceId, result = null, options = {}, traceStore = defaultCodeAgentTraceStore }) {
const initial = result ?? await loadPersistedAgentRunResult(traceId, options);
const mapped = initial ? await resolveAgentRunTraceCommandMapping({ traceId, mapped: initial, options }) : initial;
if (!mapped?.agentRun?.runId) return traceStore.snapshot(traceId);
const env = options.env ?? process.env;
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const managerUrl = resolveAgentRunManagerUrl(env, mapped.agentRun.managerUrl);
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
const eventsResponse = await fetchAgentRunEventsForTrace({ fetchImpl, managerUrl, timeoutMs, mapping: { ...mapped.agentRun, traceSummary: mapped.traceSummary } });
const events = eventsResponse.events;
appendAgentRunEventsToTrace(traceStore, traceId, events, mapped.agentRun);
const lastSeq = agentRunTraceCursorSeq(eventsResponse, mapped.agentRun.lastSeq);
if (lastSeq !== Number(mapped.agentRun.lastSeq ?? 0)) {
options.codeAgentChatResults?.set?.(traceId, { ...mapped, agentRun: { ...mapped.agentRun, lastSeq, updatedAt: nowIso(options.now) } });
}
return traceStore.snapshot(traceId);
}
export async function cancelAgentRunChatTurn({ traceId, currentResult = null, options = {}, traceStore = defaultCodeAgentTraceStore }) {
const mapped = currentResult ?? await loadPersistedAgentRunResult(traceId, options);
if (!mapped?.agentRun?.commandId) return null;
const env = options.env ?? process.env;
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const managerUrl = resolveAgentRunManagerUrl(env, mapped.agentRun.managerUrl);
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
await agentRunJson(fetchImpl, managerUrl, `/api/v1/commands/${encodeURIComponent(mapped.agentRun.commandId)}/cancel`, {
method: "POST",
body: { reason: "hwlab-user-cancel", traceId },
timeoutMs
});
traceStore.append(traceId, {
type: "cancel",
status: "canceled",
label: "agentrun:cancel:canceled",
message: "HWLAB forwarded cancel to AgentRun command cancel API.",
runId: mapped.agentRun.runId,
commandId: mapped.agentRun.commandId,
terminal: true,
valuesPrinted: false
});
const now = nowIso(options.now);
const agentRun = { ...mapped.agentRun, status: "cancelled", commandState: "cancelled", terminalStatus: "cancelled", updatedAt: now };
const payload = {
...mapped,
status: "canceled",
canceled: true,
updatedAt: now,
agentRun,
runnerTrace: traceStore.snapshot(traceId),
error: {
code: "agentrun_canceled",
layer: "agentrun",
category: "canceled",
retryable: true,
message: "AgentRun command was canceled by user request",
userMessage: "当前 AgentRun 请求已取消;traceId/runId/commandId 已保留,可重试上一条消息。",
traceId,
route: "/v1/agent/chat/cancel",
toolName: "agentrun.command.cancel"
},
userMessage: "当前 AgentRun 请求已取消;traceId/runId/commandId 已保留,可重试上一条消息。"
};
options.codeAgentChatResults?.set?.(traceId, payload);
return payload;
}
export async function steerAgentRunChatTurn({ traceId, currentResult = null, params = {}, options = {}, traceStore = defaultCodeAgentTraceStore }) {
const mapped = currentResult ?? await loadPersistedAgentRunResult(traceId, options);
if (!safeTraceId(traceId) || !mapped?.agentRun?.runId || !mapped?.agentRun?.commandId) return null;
const message = firstNonEmpty(params.message, params.prompt, params.text);
if (!message) {
throw Object.assign(new Error("steer command requires non-empty message text"), { code: "steer_message_missing", statusCode: 400 });
}
const env = options.env ?? process.env;
const fetchImpl = options.fetchImpl ?? globalThis.fetch;
const managerUrl = resolveAgentRunManagerUrl(env, mapped.agentRun.managerUrl);
const timeoutMs = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 20_000);
const steerTraceId = safeTraceId(params.steerTraceId ?? params.steerId) || `trc_steer_${randomUUID().replace(/-/gu, "")}`;
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:steer:accepted",
message: "HWLAB accepted an in-flight steer request and will create an AgentRun type=steer command on the existing run.",
runId: mapped.agentRun.runId,
commandId: mapped.agentRun.commandId,
steerTraceId,
waitingFor: "agentrun-steer-command-create",
valuesPrinted: false
});
const commandInput = buildAgentRunSteerCommandInput({ params: { ...params, message }, traceId, steerTraceId, mapped });
const command = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(mapped.agentRun.runId)}/commands`, {
method: "POST",
body: commandInput,
timeoutMs
});
const steerCommandId = requiredString(command?.id, "command.id");
const now = nowIso(options.now);
const agentRun = {
...mapped.agentRun,
lastSteerCommandId: steerCommandId,
lastSteerTraceId: steerTraceId,
updatedAt: now,
valuesPrinted: false
};
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:steer:command-created",
message: `AgentRun steer command ${steerCommandId} created on run ${mapped.agentRun.runId}; runner will apply it if the target Codex turn is still active.`,
runId: mapped.agentRun.runId,
commandId: mapped.agentRun.commandId,
steerCommandId,
targetCommandId: mapped.agentRun.commandId,
steerTraceId,
waitingFor: "agentrun-steer-apply",
valuesPrinted: false
});
const payload = {
ok: true,
accepted: true,
status: "running",
shortConnection: true,
controlSemantics: "steer-active-turn-and-poll-target-trace",
route: "/v1/agent/chat/steer",
traceId,
targetTraceId: traceId,
steerTraceId,
conversationId: safeConversationId(mapped.conversationId ?? params.conversationId) || null,
sessionId: safeSessionId(mapped.sessionId ?? params.sessionId) || null,
threadId: safeOpaqueId(mapped.threadId ?? params.threadId) || null,
resultUrl: `/v1/agent/chat/result/${encodeURIComponent(traceId)}`,
traceUrl: `/v1/agent/chat/trace/${encodeURIComponent(traceId)}`,
agentRun: {
...agentRun,
steerCommandId,
targetCommandId: mapped.agentRun.commandId,
steerCommandState: command.state ?? null,
valuesPrinted: false
},
runnerTrace: traceStore.snapshot(traceId, agentRunTraceMeta({}, {})),
updatedAt: now,
valuesPrinted: false
};
options.codeAgentChatResults?.set?.(traceId, { ...mapped, agentRun, runnerTrace: payload.runnerTrace, updatedAt: now });
return payload;
}
export async function loadPersistedAgentRunResult(traceId, options = {}) {
const safeId = safeTraceId(traceId);
if (!safeId || typeof options.accessController?.getAgentSessionByTraceId !== "function") return null;
const session = await options.accessController.getAgentSessionByTraceId(safeId);
const agentRun = agentRunSeedFromSession(session, safeId);
if (!agentRun?.runId) return null;
const backendProfile = firstNonEmpty(agentRun.backendProfile, "deepseek");
const status = firstNonEmpty(session?.status === "canceled" ? "canceled" : null, agentRun.terminalStatus, agentRun.commandState, agentRun.status, "running");
return {
accepted: true,
status,
shortConnection: true,
traceId: safeId,
conversationId: safeConversationId(session?.conversationId) || agentRun.conversationId || null,
sessionId: safeSessionId(session?.id) || agentRun.sessionId || null,
threadId: safeOpaqueId(session?.threadId) || agentRun.threadId || null,
messageId: `msg_${safeId.slice(4)}`,
createdAt: session?.startedAt ?? session?.updatedAt ?? nowIso(options.now),
updatedAt: session?.updatedAt ?? nowIso(options.now),
ownerUserId: session?.ownerUserId,
ownerRole: session?.ownerRole,
provider: providerForBackendProfile(backendProfile),
model: modelForBackendProfile(backendProfile, options.env ?? process.env),
backend: backendForBackendProfile(backendProfile),
infrastructureBackend: `agentrun-v01/${backendProfile}`,
capabilityLevel: AGENTRUN_CAPABILITY_LEVEL,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
agentRun: { ...agentRun, adapter: ADAPTER_ID, traceId: safeId, commandId: text(agentRun.commandId) || null, valuesPrinted: false },
valuesPrinted: false
};
}
function agentRunSeedFromSession(session, traceId) {
const snapshot = session?.session && typeof session.session === "object" ? session.session : null;
const traceResults = snapshot?.traceResults && typeof snapshot.traceResults === "object" ? snapshot.traceResults : null;
const traceAgentRun = traceResults?.[traceId]?.agentRun && typeof traceResults[traceId].agentRun === "object" ? traceResults[traceId].agentRun : null;
if (traceAgentRun?.runId) return traceAgentRun;
const topLevelAgentRun = snapshot?.agentRun && typeof snapshot.agentRun === "object" ? snapshot.agentRun : null;
return topLevelAgentRun?.runId ? topLevelAgentRun : null;
}
async function findAgentRunCommandForTrace({ fetchImpl, managerUrl, timeoutMs, runId, traceId }) {
if (!runId || !safeTraceId(traceId)) return null;
const commands = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(runId)}/commands?afterSeq=0&limit=100`, { method: "GET", timeoutMs });
return (Array.isArray(commands?.items) ? commands.items : []).find((item) => agentRunCommandMatchesTrace(item, traceId)) ?? null;
}
function agentRunCommandMatchesTrace(command, traceId) {
const payload = command?.payload && typeof command.payload === "object" ? command.payload : {};
return command?.idempotencyKey === traceId || payload.traceId === traceId;
}
export function agentRunSessionEvidence(payload = {}) {
if (!payload?.agentRun) return {};
return {
agentRun: {
adapter: ADAPTER_ID,
runId: payload.agentRun.runId ?? null,
commandId: payload.agentRun.commandId ?? null,
attemptId: payload.agentRun.attemptId ?? null,
runnerId: payload.agentRun.runnerId ?? null,
jobName: payload.agentRun.jobName ?? null,
namespace: payload.agentRun.namespace ?? null,
backendProfile: payload.agentRun.backendProfile ?? null,
managerUrl: payload.agentRun.managerUrl ?? DEFAULT_AGENTRUN_MGR_URL,
repoUrl: payload.agentRun.repoUrl ?? DEFAULT_REPO_URL,
status: payload.agentRun.status ?? null,
runStatus: payload.agentRun.runStatus ?? null,
commandState: payload.agentRun.commandState ?? null,
terminalStatus: payload.agentRun.terminalStatus ?? null,
traceId: payload.agentRun.traceId ?? payload.traceId ?? payload.providerTrace?.traceId ?? null,
lastSeq: payload.agentRun.lastSeq ?? 0,
providerId: payload.agentRun.providerId ?? DEFAULT_PROVIDER_ID,
reuseEligible: payload.agentRun.reuseEligible ?? false,
sessionId: payload.agentRun.sessionId ?? payload.sessionId ?? null,
conversationId: payload.conversationId ?? payload.agentRun.conversationId ?? null,
threadId: payload.threadId ?? payload.agentRun.threadId ?? null,
providerTrace: payload.providerTrace ?? payload.agentRun.providerTrace ?? null,
valuesPrinted: false
}
};
}
async function ensureAgentRunSessionPersistent({ fetchImpl, managerUrl, sessionId, env, traceId, backendProfile, traceStore }) {
const defaultPolicy = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_SESSION_STORAGE, "persistent");
if (defaultPolicy !== "persistent") return;
try {
const tenantId = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_TENANT_ID, DEFAULT_TENANT_ID);
const projectId = agentRunProjectIdForEnv(env);
const expiresInDays = parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_SESSION_TTL_DAYS, 30);
const expiresAt = new Date(Date.now() + Math.max(1, expiresInDays) * 24 * 60 * 60 * 1000).toISOString();
await agentRunJson(fetchImpl, managerUrl, "/api/v1/sessions", {
method: "POST",
timeoutMs: 15_000,
body: { sessionId, tenantId, projectId, backendProfile, expiresAt, codexRolloutSubdir: "sessions" },
});
} catch (error) {
const message = error?.message ?? String(error);
if (/evicted/i.test(message)) {
traceStore.append(traceId, {
type: "backend", status: "running",
label: "agentrun:session-storage-evicted",
message: "AgentRun session " + sessionId + " storage was previously evicted; hwlab-cloud-api will create a fresh sessionId for this turn.",
sessionId, backendProfile, waitingFor: "agentrun-session-reset", valuesPrinted: false,
});
throw error;
}
traceStore.append(traceId, {
type: "backend", status: "running",
label: "agentrun:session-storage-recover-warning",
message: "AgentRun session " + sessionId + " pre-flight could not ensure PVC; falling back to legacy metadata-only session (" + message + ").",
sessionId, backendProfile, valuesPrinted: false,
});
}
}
function agentRunProjectIdForEnv(env = process.env) {
return firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROJECT_ID, DEFAULT_PROJECT_ID);
}
async function shouldResetSessionAfterEviction(failureKind, failureMessage) {
return failureKind === "session-store-evicted" || /session store evicted/i.test(failureMessage ?? "");
}
function newSessionIdAfterEviction(baseSessionId, traceId) {
const profile = String(traceId).replace(/[^A-Za-z0-9]/gu, "").slice(0, 12) || "fresh";
return baseSessionId + "-reset-" + profile;
}
function buildAgentRunCreateRunInput({ params, env, traceId, backendProfile, sessionId, toolCapabilities = null }) {
const commitId = requireAgentRunSourceCommit(env);
const threadId = safeOpaqueId(params.threadId);
const hwlabProjectId = firstNonEmpty(params.projectId);
rejectRemovedResourceWorkspaceFiles(params.workspaceFiles ?? params.resourceWorkspaceFiles);
const resourceBundleRef = {
kind: "gitbundle",
repoUrl: resolveAgentRunRepoUrl(env),
commitId,
submodules: false,
lfs: false,
bundles: HWLAB_RESOURCE_GIT_BUNDLES,
promptRefs: HWLAB_RESOURCE_PROMPT_REFS,
};
return {
tenantId: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_TENANT_ID, DEFAULT_TENANT_ID),
projectId: agentRunProjectIdForEnv(env),
workspaceRef: {
kind: "opaque",
repo: DEFAULT_PROJECT_ID,
branch: firstNonEmpty(env.HWLAB_BOOT_REF, env.HWLAB_CODE_AGENT_AGENTRUN_BRANCH, "v0.2"),
runtimeNamespace: firstNonEmpty(env.HWLAB_RUNTIME_NAMESPACE, env.POD_NAMESPACE, env.HWLAB_NAMESPACE, "hwlab-v02"),
hwlabTraceId: traceId,
valuesPrinted: false
},
sessionRef: {
sessionId: sessionId ?? scopedAgentRunSessionIdForParams(params, traceId, backendProfile),
...(safeConversationId(params.conversationId) ? { conversationId: safeConversationId(params.conversationId) } : {}),
...(threadId ? { threadId } : {}),
metadata: {
adapter: ADAPTER_ID,
hwlabTraceId: traceId,
hwlabApi: "/v1/agent/chat",
hwlabProjectId,
hwlabSessionId: safeSessionId(params.sessionId) || null,
threadContinuityPolicy: THREAD_CONTINUITY_POLICY,
sessionPolicy: SESSION_POLICY_RUN_LOCAL,
agentRunSessionProfile: backendProfile,
agentRunSessionPolicy: "backend-profile-scoped",
valuesPrinted: false
}
},
resourceBundleRef,
providerId: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID),
backendProfile,
executionPolicy: {
sandbox: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_SANDBOX, env.HWLAB_CODE_AGENT_CODEX_SANDBOX, "danger-full-access"),
approval: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_APPROVAL, "never"),
timeoutMs: parsePositiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_TIMEOUT_MS, parsePositiveInteger(env.HWLAB_CODE_AGENT_TIMEOUT_MS, DEFAULT_TIMEOUT_MS)),
network: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_NETWORK, "enabled"),
secretScope: {
allowCredentialEcho: false,
providerCredentials: [providerCredentialSecretRef(backendProfile, env)],
toolCredentials: agentRunToolCredentials(env, toolCapabilities)
}
},
traceSink: {
kind: "hwlab-cloud-api",
traceId,
resultUrl: `/v1/agent/chat/result/${encodeURIComponent(traceId)}`,
traceUrl: `/v1/agent/chat/trace/${encodeURIComponent(traceId)}`,
valuesPrinted: false
}
};
}
function rejectRemovedResourceWorkspaceFiles(value) {
if (value == null || value === false || value === "") return [];
throw adapterError("legacy_workspace_files_removed", "workspaceFiles/resourceWorkspaceFiles are removed; use AgentRun ResourceBundleRef kind=gitbundle with bundles[]");
}
function buildAgentRunCommandInput({ params, traceId, backendProfile, sessionId }) {
const prompt = String(params.message ?? params.prompt ?? "").trim();
const threadId = safeOpaqueId(params.threadId);
return {
type: "turn",
payload: {
prompt,
message: prompt,
traceId,
projectId: firstNonEmpty(params.projectId) || null,
conversationId: safeConversationId(params.conversationId) || null,
sessionId: sessionId ?? scopedAgentRunSessionIdForParams(params, traceId, backendProfile),
hwlabSessionId: safeSessionId(params.sessionId) || null,
threadId: threadId || null,
threadContinuityPolicy: THREAD_CONTINUITY_POLICY,
sessionPolicy: SESSION_POLICY_RUN_LOCAL,
providerProfile: backendProfile,
source: "hwlab-cloud-api",
valuesPrinted: false
},
idempotencyKey: traceId
};
}
function buildAgentRunSteerCommandInput({ params, traceId, steerTraceId, mapped }) {
const prompt = String(params.message ?? params.prompt ?? params.text ?? "").trim();
return {
type: "steer",
payload: {
prompt,
message: prompt,
text: prompt,
traceId: steerTraceId,
targetTraceId: traceId,
targetCommandId: mapped.agentRun.commandId,
conversationId: safeConversationId(mapped.conversationId ?? params.conversationId) || null,
sessionId: mapped.agentRun.sessionId ?? safeSessionId(mapped.sessionId ?? params.sessionId) ?? null,
hwlabSessionId: safeSessionId(mapped.sessionId ?? params.sessionId) || null,
threadId: safeOpaqueId(mapped.threadId ?? params.threadId) || null,
source: "hwlab-cloud-api",
valuesPrinted: false
},
idempotencyKey: steerTraceId
};
}
async function resolveOwnerApiKey({ params, options, now }) {
const ownerUserId = text(params.ownerUserId);
if (!ownerUserId) return "";
const accessController = options.accessController;
if (!accessController?.store?.findActiveDefaultApiKeyForUser) return "";
const key = await accessController.store.findActiveDefaultApiKeyForUser(ownerUserId) ?? null;
if (!key) return "";
return text(key.displaySecret ?? "");
}
function buildAgentRunRunnerJobInput({ env, traceId, commandId, ownerApiKey, toolCapabilities = null, backendProfile }) {
const baseTransient = buildAgentRunTransientEnv(env, { providerProfile: backendProfile, parentTraceId: traceId });
const hwpodAllowed = toolCapabilityAllowed(toolCapabilities, "hwpod");
const transientEnv = ownerApiKey && hwpodAllowed
? baseTransient.concat([{ name: "HWLAB_API_KEY", value: ownerApiKey, sensitive: true }])
: baseTransient;
return {
commandId,
idempotencyKey: traceId,
namespace: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_NAMESPACE, env.AGENTRUN_RUNTIME_NAMESPACE, DEFAULT_RUNNER_NAMESPACE),
...(firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_IMAGE, env.AGENTRUN_RUNNER_IMAGE)
? { image: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_IMAGE, env.AGENTRUN_RUNNER_IMAGE) }
: {}),
...(firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_SERVICE_ACCOUNT, env.AGENTRUN_RUNNER_SERVICE_ACCOUNT)
? { serviceAccountName: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_RUNNER_SERVICE_ACCOUNT, env.AGENTRUN_RUNNER_SERVICE_ACCOUNT) }
: {}),
...(transientEnv.length > 0 ? { transientEnv } : {})
};
}
function buildAgentRunTransientEnv(env = process.env, options = {}) {
const providerProfile = firstNonEmpty(options.providerProfile);
const parentTraceId = firstNonEmpty(options.parentTraceId);
const entries = [];
for (const name of [
"HWLAB_RUNTIME_API_URL",
"HWLAB_RUNTIME_WEB_URL",
"HWLAB_RUNTIME_NAMESPACE",
"HWLAB_RUNTIME_LANE",
"HWLAB_RUNTIME_ENDPOINT_LOCKED",
"HWLAB_CODE_AGENT_ASSEMBLED_RUNTIME",
"UNIDESK_MAIN_SERVER_IP"
]) {
const value = firstNonEmpty(env[name]);
if (value) entries.push({ name, value, sensitive: false });
}
if (providerProfile) entries.push({ name: "HWLAB_CODE_AGENT_PROVIDER_PROFILE", value: providerProfile, sensitive: false });
if (parentTraceId) entries.push({ name: "HWLAB_CODE_AGENT_PARENT_TRACE_ID", value: parentTraceId, sensitive: false });
return entries;
}
function agentRunToolCredentials(env = process.env, toolCapabilities = null) {
const namespace = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_TOOL_SECRET_NAMESPACE, env.HWLAB_CODE_AGENT_AGENTRUN_SECRET_NAMESPACE, DEFAULT_RUNNER_NAMESPACE);
const githubSecretName = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_GITHUB_TOOL_SECRET_NAME, "agentrun-v01-tool-github-pr");
const githubSecretKey = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_GITHUB_TOOL_SECRET_KEY, "GH_TOKEN");
const unideskSecretName = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_UNIDESK_SSH_TOOL_SECRET_NAME, "agentrun-v01-tool-unidesk-ssh");
const credentials = [];
if (toolCapabilityAllowed(toolCapabilities, "github_pr")) credentials.push({
tool: "github",
purpose: "pull-request",
secretRef: { namespace, name: githubSecretName, keys: [githubSecretKey] },
projection: { kind: "env", envName: githubSecretKey, secretKey: githubSecretKey }
});
if (toolCapabilityAllowed(toolCapabilities, "unidesk_ssh")) credentials.push({
tool: "unidesk-ssh",
purpose: "ssh-passthrough",
secretRef: { namespace, name: unideskSecretName, keys: ["UNIDESK_SSH_CLIENT_TOKEN"] },
projection: { kind: "env", envName: "UNIDESK_SSH_CLIENT_TOKEN", secretKey: "UNIDESK_SSH_CLIENT_TOKEN" }
});
return credentials;
}
async function resolveToolCapabilities({ params = {}, options = {} } = {}) {
const ownerUserId = text(params.ownerUserId);
const accessController = options.accessController;
if (!ownerUserId || typeof accessController?.codeAgentToolCapabilitiesForOwner !== "function") return null;
return await accessController.codeAgentToolCapabilitiesForOwner(ownerUserId);
}
function toolCapabilityAllowed(toolCapabilities = null, toolId) {
if (!toolCapabilities || !toolCapabilities.tools) return true;
return toolCapabilities.tools?.[toolId]?.allowed === true;
}
async function resolveReusableAgentRun({ params = {}, options = {}, env = process.env, managerUrl, fetchImpl, timeoutMs, backendProfile, traceId, traceStore }) {
const hwlabSessionId = hwlabSessionIdForParams(params, traceId);
if (!safeSessionId(hwlabSessionId) || hwlabSessionId === agentRunSessionId(traceId) || typeof options.accessController?.getAgentSession !== "function") {
return null;
}
const session = await options.accessController.getAgentSession(hwlabSessionId);
if (!canReuseAgentRunSessionForOwner(session, params, options)) {
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reuse-owner-skipped",
message: "Stored AgentRun session is not visible to the current actor; a new runner Job will be created for this request.",
sessionId: hwlabSessionId,
waitingFor: "agentrun-run-create",
valuesPrinted: false
});
return null;
}
const mapping = session?.session?.agentRun;
if (!isReusableAgentRunMapping(mapping, backendProfile)) {
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reuse-skipped",
message: "No reusable AgentRun run was found for this HWLAB session; a new runner Job will be created.",
sessionId: hwlabSessionId,
waitingFor: "agentrun-run-create",
valuesPrinted: false
});
return null;
}
let run = null;
try {
run = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(mapping.runId)}`, {
method: "GET",
timeoutMs
});
} catch (error) {
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reuse-unavailable",
message: `Stored AgentRun run ${mapping.runId} could not be read; a new runner Job will be created.`,
errorCode: error?.code ?? "agentrun_run_lookup_failed",
runId: mapping.runId,
sessionId: hwlabSessionId,
waitingFor: "agentrun-run-create",
valuesPrinted: false
});
return null;
}
const reason = agentRunReuseBlocker(run, mapping);
if (reason) {
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reuse-blocked",
message: `Stored AgentRun run ${mapping.runId} cannot be reused: ${reason}; a new runner Job will be created.`,
runId: mapping.runId,
sessionId: hwlabSessionId,
waitingFor: "agentrun-run-create",
valuesPrinted: false
});
return null;
}
traceStore.append(traceId, {
type: "backend",
status: "running",
label: "agentrun:run:reused",
message: `AgentRun run ${mapping.runId} is reused for HWLAB session ${hwlabSessionId}; runner reuse window is active and no new bundle will be requested for this turn.`,
runId: mapping.runId,
commandId: mapping.commandId ?? null,
runnerId: mapping.runnerId ?? null,
jobName: mapping.jobName ?? null,
namespace: mapping.namespace ?? null,
sessionId: hwlabSessionId,
waitingFor: "agentrun-command-create",
valuesPrinted: false
});
return { sessionId: mapping.sessionId ?? scopedAgentRunSessionIdForParams(params, traceId, backendProfile), hwlabSessionId, session, mapping, run };
}
function isReusableAgentRunMapping(mapping, backendProfile) {
return Boolean(
mapping &&
typeof mapping === "object" &&
mapping.adapter === ADAPTER_ID &&
mapping.reuseEligible !== false &&
mapping.runId &&
mapping.jobName &&
(!mapping.backendProfile || mapping.backendProfile === backendProfile)
);
}
function canReuseAgentRunSessionForOwner(session, params = {}, options = {}) {
if (!session) return false;
const ownerUserId = options.actor?.id ?? params.ownerUserId ?? null;
const ownerRole = options.actor?.role ?? params.ownerRole ?? null;
if (ownerRole === "admin") return true;
if (!session.ownerUserId) return true;
return Boolean(ownerUserId && session.ownerUserId === ownerUserId);
}
function agentRunReuseBlocker(run = {}, mapping = {}) {
const status = String(run?.terminalStatus ?? run?.status ?? mapping.runStatus ?? "").toLowerCase();
if (TERMINAL_RUN_STATUSES.has(status)) return `run_status_${status}`;
if (!run?.claimedBy) return "runner_not_claimed";
if (runLeaseExpired(run.leaseExpiresAt)) return "runner_reuse_window_expired";
return null;
}
function runLeaseExpired(leaseExpiresAt) {
const expiresMs = Date.parse(String(leaseExpiresAt ?? ""));
return !Number.isFinite(expiresMs) || expiresMs <= Date.now();
}
function providerCredentialSecretRef(profile, env) {
const profileEnvKey = String(profile ?? "").toUpperCase().replace(/[^A-Z0-9]+/gu, "_");
return {
profile,
secretRef: {
namespace: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_SECRET_NAMESPACE, DEFAULT_RUNNER_NAMESPACE),
name: firstNonEmpty(env[`HWLAB_CODE_AGENT_AGENTRUN_${profileEnvKey}_SECRET_NAME`], env[`HWLAB_CODE_AGENT_AGENTRUN_${String(profile ?? "").toUpperCase()}_SECRET_NAME`], env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_SECRET_NAME, `agentrun-v01-provider-${profile}`),
keys: ["auth.json", "config.toml"],
mountPath: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_CODEX_HOME_MOUNT, "~/.codex")
}
};
}
function agentRunProviderTrace({ base = {}, result = {}, terminalStatus = null } = {}) {
return {
transport: "agentrun-v01",
protocol: AGENTRUN_PROVIDER_TRACE_PROTOCOL,
wireApi: AGENTRUN_PROVIDER_TRACE_WIRE_API,
runnerKind: AGENTRUN_RUNNER_KIND,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
command: "agentrun.v01.command.turn",
toolName: "agentrun.v01.command.turn",
source: ADAPTER_ID,
backendProfile: base.agentRun?.backendProfile ?? null,
runId: base.agentRun?.runId ?? result?.runId ?? null,
commandId: base.agentRun?.commandId ?? result?.commandId ?? null,
runnerId: base.agentRun?.runnerId ?? result?.runnerId ?? null,
jobName: base.agentRun?.jobName ?? result?.jobName ?? null,
namespace: base.agentRun?.namespace ?? result?.namespace ?? null,
traceId: base.traceId ?? base.agentRun?.traceId ?? null,
threadId: result?.sessionRef?.threadId ?? base.agentRun?.threadId ?? base.threadId ?? null,
turnId: result?.turnId ?? null,
terminalStatus: terminalStatus ?? result?.terminalStatus ?? null,
failureKind: result?.failureKind ?? null,
failureMessage: result?.failureMessage ?? result?.blocker?.message ?? null,
valuesPrinted: false
};
}
function agentRunFailureAttribution({ code, message, canceled = false } = {}) {
if (canceled) {
return {
category: "canceled",
retryable: true,
userMessage: "AgentRun 请求已取消。",
summary: message || "AgentRun command was canceled."
};
}
if (code === "provider-invalid-tool-call" || /invalid function arguments json string|invalid_prompt|tool_call_id/iu.test(String(message ?? ""))) {
return {
category: "provider_invalid_tool_call",
retryable: true,
userMessage: "AgentRun/provider 返回了无效 tool-call arguments JSON;这不是 HWPOD 或 Cloud API 端点失败,请查看 providerTrace.failureKind 和 runnerTrace 定位上游 tool_call_id。",
summary: message || "AgentRun/provider returned invalid tool-call arguments JSON."
};
}
if (code === "thread-resume-failed") {
return {
category: "thread_resume_failed",
retryable: true,
userMessage: "AgentRun 复用的 thread 已失效;当前 turn 会按失败终止,Workbench 会清理 stale continuation 指针,下一轮自动从新 session/thread 启动。",
summary: message || "AgentRun thread/resume failed for an existing thread."
};
}
return {
category: "agentrun_failed",
retryable: true,
userMessage: "AgentRun 请求失败;请查看 runnerTrace、providerTrace 和 agentRun 字段定位 runId/commandId/jobName。",
summary: message || "AgentRun command failed."
};
}
function agentRunLongLivedSessionGate(base = {}) {
const session = agentRunSessionSummary(base, "idle");
return {
status: "pass",
pass: true,
requiredCapability: AGENTRUN_CAPABILITY_LEVEL,
currentCapability: AGENTRUN_IMPLEMENTATION_TYPE,
sessionId: session.sessionId ?? null,
sessionStatus: session.status ?? null,
runnerKind: AGENTRUN_RUNNER_KIND,
provider: providerForBackendProfile(base.agentRun?.backendProfile ?? "deepseek"),
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
adapter: ADAPTER_ID,
delegatedToAgentRun: true,
codexStdio: false,
blockers: [],
valuesPrinted: false
};
}
function agentRunToolCalls(result = {}, status = "completed") {
return [{
name: "agentrun.v01.command.turn",
status,
runId: result?.runId ?? null,
commandId: result?.commandId ?? null,
outputSummary: result?.reply ? `assistantChars=${String(result.reply).length}` : null,
valuesPrinted: false
}];
}
function agentRunResultToCodeAgentPayload({ base, result, traceStore, traceId, appendResultEvent = true }) {
const terminalStatus = String(result?.terminalStatus ?? "");
const terminal = terminalStatus === "completed" || terminalStatus === "failed" || terminalStatus === "blocked" || terminalStatus === "cancelled";
if (!terminal) {
return decorateAgentRunRunningResult({ base, mapping: base.agentRun, traceStore, traceId });
}
const now = nowIso();
const runnerTrace = traceStore.snapshot(traceId, agentRunTraceMeta({}, {}));
const terminalEventCreatedAt = agentRunResultTraceCreatedAt(runnerTrace, now);
const providerTrace = agentRunProviderTrace({ base, result, terminalStatus });
if (terminalStatus === "completed" && String(result?.reply ?? "").trim()) {
const finalResponse = agentRunCompletedFinalResponse({ base, result, traceId, now });
const traceSummary = agentRunCompletedTraceSummary({ base, runnerTrace, finalResponse, traceId });
if (appendResultEvent) {
traceStore.append(traceId, {
type: "result",
status: "completed",
label: "agentrun:result:completed",
createdAt: terminalEventCreatedAt,
message: "AgentRun result is ready for HWLAB short-connection polling.",
runId: base.agentRun.runId,
commandId: base.agentRun.commandId,
terminal: true,
valuesPrinted: false
});
}
return {
...base,
status: "completed",
updatedAt: now,
workspace: base.workspace ?? "/home/agentrun/workspace",
sandbox: base.sandbox ?? "danger-full-access",
session: agentRunSessionSummary(base, "idle"),
sessionReuse: agentRunSessionReuseSummary(base, base.agentRun.reused === true),
runner: agentRunRunnerSummary(base.agentRun),
runnerTrace: traceStore.snapshot(traceId, agentRunTraceMeta({}, {})),
toolCalls: agentRunToolCalls(result, "completed"),
skills: { status: "delegated", provider: ADAPTER_ID, count: 0, items: [], valuesPrinted: false },
longLivedSessionGate: agentRunLongLivedSessionGate(base),
providerTrace,
finalResponse,
traceSummary,
reply: {
messageId: finalResponse.messageId,
role: "assistant",
content: finalResponse.text,
createdAt: finalResponse.createdAt
},
usage: null,
agentRun: { ...base.agentRun, terminalStatus, completed: true, reuseEligible: true, providerTrace, valuesPrinted: false },
valuesPrinted: false
};
}
const canceled = terminalStatus === "cancelled";
const code = canceled ? "agentrun_canceled" : result?.failureKind ?? (terminalStatus === "blocked" ? "agentrun_blocked" : "agentrun_failed");
const message = result?.failureMessage ?? result?.blocker?.message ?? (canceled ? "AgentRun command was canceled" : "AgentRun command failed");
const attribution = agentRunFailureAttribution({ code, message, canceled });
if (appendResultEvent) {
traceStore.append(traceId, {
type: "result",
status: canceled ? "canceled" : "failed",
label: `agentrun:result:${canceled ? "canceled" : terminalStatus || "failed"}`,
createdAt: terminalEventCreatedAt,
errorCode: code,
message: attribution.summary,
runId: base.agentRun.runId,
commandId: base.agentRun.commandId,
terminal: true,
valuesPrinted: false
});
}
const partialContext = partialAgentRunContext(runnerTrace);
return {
...base,
status: canceled ? "canceled" : "failed",
canceled,
updatedAt: now,
session: agentRunSessionSummary(base, canceled ? "canceled" : "failed"),
sessionReuse: agentRunSessionReuseSummary(base, base.agentRun.reused === true, { status: canceled ? undefined : "failed-requires-new-session" }),
runner: agentRunRunnerSummary(base.agentRun),
runnerTrace: partialContext ? { ...runnerTrace, partialContext } : runnerTrace,
toolCalls: agentRunToolCalls(result, canceled ? "canceled" : "failed"),
skills: { status: "delegated", provider: ADAPTER_ID, count: 0, items: [], valuesPrinted: false },
providerTrace,
error: {
code,
layer: "agentrun",
category: attribution.category,
retryable: attribution.retryable,
message,
userMessage: attribution.userMessage,
traceId,
route: "/v1/agent/chat",
toolName: "agentrun.manual-dispatch"
},
blocker: canceled ? null : {
code,
layer: "agentrun",
category: attribution.category,
retryable: attribution.retryable,
summary: attribution.summary,
traceId,
route: "/v1/agent/chat",
toolName: "agentrun.manual-dispatch"
},
agentRun: { ...base.agentRun, terminalStatus, completed: false, reuseEligible: false, providerTrace, valuesPrinted: false },
reuseEligible: false,
...(partialContext ? { partialContext } : {}),
valuesPrinted: false
};
}
function agentRunCompletedFinalResponse({ base, result, traceId, now }) {
const textValue = String(result?.reply ?? "").trim();
return {
text: textValue,
textChars: textValue.length,
role: "assistant",
status: "completed",
traceId,
messageId: base.messageId ?? base.reply?.messageId ?? `msg_${traceId.slice(4)}`,
createdAt: base.reply?.createdAt ?? base.createdAt ?? now,
updatedAt: now,
valuesPrinted: false
};
}
function agentRunCompletedTraceSummary({ base, runnerTrace, finalResponse, traceId }) {
const events = Array.isArray(runnerTrace?.events) ? runnerTrace.events : [];
const lastSeq = Math.max(0, ...events.map((event) => Number(event?.seq ?? 0)).filter(Number.isFinite));
return {
traceId,
source: "agentrun-command-result",
sourceEventCount: Number(runnerTrace?.eventCount ?? events.length ?? 0),
terminalStatus: "completed",
finalAssistantRow: {
role: finalResponse.role,
status: finalResponse.status,
textChars: finalResponse.textChars,
textPreview: finalResponse.text.slice(0, 240),
messageId: finalResponse.messageId,
valuesPrinted: false
},
agentRun: {
runId: base.agentRun?.runId ?? null,
commandId: base.agentRun?.commandId ?? null,
lastSeq,
valuesPrinted: false
},
valuesPrinted: false
};
}
function partialAgentRunContext(runnerTrace = {}) {
const events = Array.isArray(runnerTrace.events) ? runnerTrace.events : [];
const assistantMessages = events
.filter((event) => event?.type === "assistant" && String(event.text ?? event.message ?? "").trim())
.map((event) => String(event.text ?? event.message).trim())
.slice(-4);
const toolEvidence = events
.filter((event) => event?.type === "tool_call" && String(event.status ?? "") === "completed")
.map((event) => [event.toolName ?? event.label ?? "tool", event.command ?? event.outputSummary ?? event.stdoutSummary ?? event.message ?? ""].filter(Boolean).join(": "))
.filter(Boolean)
.slice(-6);
if (assistantMessages.length === 0 && toolEvidence.length === 0) return null;
return {
status: "partial-before-terminal",
summary: "AgentRun terminal 前已有可延续的 assistant/tool 上下文;后续同 conversation/session/thread 轮次必须能继续使用。",
assistantMessages,
toolEvidence,
valuesPrinted: false
};
}
function agentRunResultTraceCreatedAt(runnerTrace = {}, fallback) {
const last = String(runnerTrace?.lastEvent?.createdAt ?? runnerTrace?.updatedAt ?? "");
return Number.isFinite(Date.parse(last)) ? last : fallback;
}
function decorateAgentRunRunningResult({ base, mapping, traceStore, traceId }) {
return {
...base,
status: "running",
accepted: true,
shortConnection: true,
updatedAt: nowIso(),
session: agentRunSessionSummary({ ...base, agentRun: mapping }, "running"),
sessionReuse: agentRunSessionReuseSummary({ ...base, agentRun: mapping }, mapping.reused === true),
runner: agentRunRunnerSummary(mapping),
runnerTrace: traceStore.snapshot(traceId, agentRunTraceMeta({}, {})),
agentRun: { ...mapping, adapter: ADAPTER_ID, valuesPrinted: false },
valuesPrinted: false
};
}
function appendAgentRunEventsToTrace(traceStore, traceId, events, mapping = {}) {
for (const event of events) {
if (isForeignAgentRunCommandEvent(event, mapping)) continue;
const normalized = mapAgentRunEvent(event, mapping);
if (normalized) traceStore.append(traceId, normalized, agentRunTraceMeta({}, {}));
}
}
async function fetchAgentRunEventsForTrace({ fetchImpl, managerUrl, timeoutMs, mapping = {} }) {
const runId = requiredString(mapping.runId, "runId");
const currentCommandId = typeof mapping.commandId === "string" ? mapping.commandId : "";
const { afterSeq, endSeq } = agentRunTraceReplayWindow(mapping);
const path = `/api/v1/runs/${encodeURIComponent(runId)}/events?afterSeq=${encodeURIComponent(String(afterSeq))}&limit=500`;
const response = await agentRunJson(fetchImpl, managerUrl, path, { method: "GET", timeoutMs });
const rawEvents = Array.isArray(response?.items) ? response.items : [];
const events = rawEvents.filter((event) => agentRunEventBelongsToTrace(event, { currentCommandId, afterSeq, endSeq }));
return {
events,
afterSeq,
endSeq,
commandFiltered: Boolean(currentCommandId),
maxSeq: Math.max(afterSeq, ...rawEvents.map((event) => Number(event?.seq ?? 0))),
traceLastSeq: Math.max(afterSeq, ...events.map((event) => Number(event?.seq ?? 0)).filter(Number.isFinite))
};
}
function agentRunTraceCursorSeq(eventsResponse = {}, previousLastSeq = 0) {
const afterSeq = Number(eventsResponse.afterSeq ?? 0);
const endSeq = Number(eventsResponse.endSeq ?? 0);
const traceLastSeq = Number(eventsResponse.traceLastSeq ?? 0);
if (Number.isFinite(traceLastSeq) && traceLastSeq > afterSeq) return Math.floor(traceLastSeq);
if (Number.isFinite(endSeq) && endSeq > 0) return Math.floor(endSeq);
if (eventsResponse.commandFiltered === true) return Math.max(Number(previousLastSeq ?? 0), afterSeq);
return Math.max(Number(previousLastSeq ?? 0), Number(eventsResponse.maxSeq ?? 0));
}
function agentRunTraceReplayWindow(mapping = {}) {
const eventStartSeq = Number(mapping.eventStartSeq ?? mapping.commandStartSeq ?? mapping.startSeq ?? 0);
const summary = mapping.traceSummary && typeof mapping.traceSummary === "object" ? mapping.traceSummary : null;
const summaryAgentRun = summary?.agentRun && typeof summary.agentRun === "object" ? summary.agentRun : null;
const summaryLastSeq = Number(summaryAgentRun?.lastSeq ?? summary?.lastSeq ?? 0);
const currentLastSeq = Number(mapping.lastSeq ?? 0);
const endSeq = Number.isFinite(summaryLastSeq) && summaryLastSeq > 0 ? Math.floor(summaryLastSeq) : 0;
if (Number.isFinite(eventStartSeq) && eventStartSeq > 0) return { afterSeq: Math.max(0, Math.floor(eventStartSeq) - 1), endSeq };
if (endSeq > 0) return { afterSeq: Math.max(0, endSeq - 500), endSeq };
if (Number.isFinite(currentLastSeq) && currentLastSeq > 0) return { afterSeq: Math.floor(currentLastSeq), endSeq: 0 };
return { afterSeq: 0, endSeq: 0 };
}
function agentRunEventBelongsToTrace(event, { currentCommandId = "", afterSeq = 0, endSeq = 0 } = {}) {
const eventCommandId = agentRunEventCommandId(event);
const seq = Number(event?.seq ?? 0);
if (currentCommandId && eventCommandId && eventCommandId !== currentCommandId) return false;
if (endSeq > 0 && Number.isFinite(seq)) return seq > afterSeq && seq <= endSeq;
return true;
}
function agentRunEventCommandId(event) {
const payload = event?.payload && typeof event.payload === "object" ? event.payload : {};
return typeof payload.commandId === "string" ? payload.commandId : "";
}
function isForeignAgentRunCommandEvent(event, mapping = {}) {
const currentCommandId = typeof mapping.commandId === "string" ? mapping.commandId : "";
if (!currentCommandId) return false;
const eventCommandId = agentRunEventCommandId(event);
return Boolean(eventCommandId && eventCommandId !== currentCommandId);
}
function mapAgentRunEvent(event, mapping = {}) {
if (!event || typeof event !== "object") return null;
const payload = event.payload && typeof event.payload === "object" ? event.payload : {};
const base = {
createdAt: event.createdAt ?? null,
source: "agentrun",
sourceSeq: event.seq ?? null,
runId: event.runId ?? mapping.runId ?? null,
commandId: payload.commandId ?? mapping.commandId ?? null,
attemptId: payload.attemptId ?? mapping.attemptId ?? null,
runnerId: payload.runnerId ?? mapping.runnerId ?? null,
jobName: payload.jobName ?? mapping.jobName ?? null,
namespace: payload.namespace ?? mapping.namespace ?? null,
valuesPrinted: false
};
if (event.type === "backend_status") {
const phase = String(payload.phase ?? "status");
return {
...base,
type: "backend",
status: "running",
label: `agentrun:backend:${phase}`,
message: textPayload(payload, phase),
waitingFor: waitingForForPhase(phase),
details: agentRunBackendStatusDetails(phase, payload)
};
}
if (event.type === "assistant_message") {
const terminal = payload.replyAuthority === true || payload.final === true;
return {
...base,
type: "assistant",
status: terminal ? "completed" : "running",
label: "agentrun:assistant:message",
message: textPayload(payload, "assistant message"),
text: textPayload(payload, ""),
itemId: payload.itemId ?? null,
messageIndex: payload.messageIndex ?? null,
messageCount: payload.messageCount ?? null,
replyAuthority: payload.replyAuthority === true,
final: payload.final === true,
terminal
};
}
if (event.type === "tool_call") {
const commandExecution = agentRunCommandExecutionEvent(base, payload);
if (commandExecution) return commandExecution;
return {
...base,
type: "tool_call",
status: String(payload.status ?? (payload.method === "item/completed" ? "completed" : "running")),
label: `agentrun:tool:${String(payload.name ?? payload.method ?? "call")}`,
itemId: payload.item?.id ?? payload.itemId ?? null,
toolName: payload.item?.type ?? payload.name ?? payload.method ?? "call",
message: textPayload(payload, "tool call"),
outputBytes: typeof payload.outputBytes === "number" ? payload.outputBytes : payload.summary?.outputBytes,
outputTruncated: payload.outputTruncated === true || payload.summary?.outputTruncated === true
};
}
if (event.type === "command_output") {
const stream = payload.stream === "stderr" ? "stderr" : "stdout";
return { ...base, type: "output", status: "running", label: `agentrun:output:${stream}`, stream, message: textPayload(payload, ""), outputTruncated: payload.outputTruncated === true };
}
if (event.type === "diff") {
return { ...base, type: "diff", status: "running", label: "agentrun:diff", message: textPayload(payload, "diff") };
}
if (event.type === "error") {
return { ...base, type: "error", status: "failed", label: `agentrun:error:${String(payload.failureKind ?? "backend")}`, errorCode: payload.failureKind ?? "agentrun_error", message: textPayload(payload, "AgentRun error") };
}
if (event.type === "terminal_status") {
const terminalStatus = String(payload.terminalStatus ?? "failed");
return { ...base, type: "result", status: terminalStatus === "cancelled" ? "canceled" : terminalStatus, label: `agentrun:terminal:${terminalStatus}`, errorCode: payload.failureKind ?? null, message: textPayload(payload, `AgentRun terminal status ${terminalStatus}`), terminal: true };
}
return { ...base, type: "backend", status: "running", label: `agentrun:event:${String(event.type ?? "unknown")}`, message: textPayload(payload, "AgentRun event") };
}
function agentRunBackendStatusDetails(phase, payload = {}) {
if (phase === "initial-prompt-assembly") {
return compactObject({
initialPromptInjected: payload.initialPromptInjected === true,
reason: firstNonEmpty(payload.reason),
initialPrompt: agentRunResourceSummary(payload.initialPrompt),
valuesPrinted: false
});
}
if (phase === "resource-bundle-materialized") {
return compactObject({
kind: firstNonEmpty(payload.kind),
commitId: firstNonEmpty(payload.commitId),
bundles: agentRunResourceSummary(payload.bundles),
tools: agentRunResourceSummary(payload.tools),
promptRefs: agentRunResourceSummary(payload.promptRefs),
skillDirs: agentRunResourceSummary(payload.skillDirs),
initialPrompt: agentRunResourceSummary(payload.initialPrompt),
valuesPrinted: false
});
}
return undefined;
}
function agentRunResourceSummary(value) {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
return compactObject({
available: typeof value.available === "boolean" ? value.available : undefined,
count: Number.isInteger(value.count) ? value.count : undefined,
bytes: Number.isInteger(value.bytes) ? value.bytes : undefined,
sha256: firstNonEmpty(value.sha256),
promptRefCount: Number.isInteger(value.promptRefCount) ? value.promptRefCount : undefined,
skillCount: Number.isInteger(value.skillCount) ? value.skillCount : undefined,
toolCount: Number.isInteger(value.toolCount) ? value.toolCount : undefined,
names: Array.isArray(value.names) ? value.names.map((item) => firstNonEmpty(item)).filter(Boolean).slice(0, 20) : undefined,
targetPaths: Array.isArray(value.targetPaths) ? value.targetPaths.map((item) => firstNonEmpty(item)).filter(Boolean).slice(0, 20) : undefined,
skillNames: Array.isArray(value.skillNames) ? value.skillNames.map((item) => firstNonEmpty(item)).filter(Boolean).slice(0, 20) : undefined,
requiredMissing: Array.isArray(value.requiredMissing) ? value.requiredMissing.map((item) => firstNonEmpty(item)).filter(Boolean).slice(0, 20) : undefined,
valuesPrinted: false
});
}
function compactObject(value) {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
const result = {};
for (const [key, item] of Object.entries(value)) {
if (item === undefined || item === null) continue;
if (Array.isArray(item) && item.length === 0) continue;
if (typeof item === "object" && !Array.isArray(item) && Object.keys(item).length === 0) continue;
result[key] = item;
}
return Object.keys(result).length > 0 ? result : undefined;
}
function agentRunCommandExecutionEvent(base, payload = {}) {
const item = payload.item && typeof payload.item === "object" ? payload.item : null;
const payloadToolType = firstNonEmpty(payload.toolName, payload.type, payload.name);
if (item?.type !== "commandExecution" && payloadToolType !== "commandExecution") return null;
const method = String(payload.method ?? "");
const completed = method === "item/completed" || payload.status === "completed" || item?.status === "completed";
const output = firstNonEmpty(item?.aggregatedOutput, item?.stdout, item?.output, payload.outputSummary, payload.stdoutSummary, payload.summary?.text, payload.itemPreview);
return {
...base,
type: "tool_call",
status: completed ? "completed" : "started",
label: completed ? "item/commandExecution:completed" : "item/commandExecution:started",
toolName: "commandExecution",
itemId: item?.id ?? payload.itemId ?? null,
command: firstNonEmpty(item?.command, item?.commandLine, payload.command, payload.commandLine),
exitCode: Number.isInteger(item?.exitCode) ? item.exitCode : Number.isInteger(payload.exitCode) ? payload.exitCode : undefined,
durationMs: typeof item?.durationMs === "number" ? item.durationMs : typeof payload.durationMs === "number" ? payload.durationMs : undefined,
outputBytes: typeof payload.outputBytes === "number" ? payload.outputBytes : payload.summary?.outputBytes,
stdoutSummary: output,
outputSummary: output,
outputTruncated: payload.outputTruncated === true || payload.summary?.outputTruncated === true,
message: output || textPayload(payload, "commandExecution")
};
}
async function agentRunJson(fetchImpl, managerUrl, path, { method = "GET", body, timeoutMs = 20_000 } = {}) {
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), timeoutMs);
try {
const response = await fetchImpl(`${managerUrl}${path}`, {
method,
headers: body === undefined ? undefined : { "content-type": "application/json" },
body: body === undefined ? undefined : JSON.stringify(body),
signal: controller.signal
});
const text = await response.text();
const parsed = text ? JSON.parse(text) : {};
if (!response.ok || parsed?.ok === false) {
throw Object.assign(new Error(parsed?.message ?? `AgentRun ${method} ${path} failed with HTTP ${response.status}`), {
code: parsed?.failureKind ?? "agentrun_request_failed",
statusCode: response.status,
agentRunError: parsed
});
}
return parsed?.ok === true && Object.hasOwn(parsed, "data") ? parsed.data : parsed;
} catch (error) {
if (error?.name === "AbortError") {
throw Object.assign(new Error(`AgentRun ${method} ${path} timed out after ${timeoutMs}ms`), { code: "agentrun_timeout", statusCode: 504 });
}
throw error;
} finally {
clearTimeout(timeout);
}
}
function resolveAgentRunManagerUrl(env = process.env, override = null) {
const raw = firstNonEmpty(override, env.AGENTRUN_MGR_URL, env.HWLAB_CODE_AGENT_AGENTRUN_MGR_URL, DEFAULT_AGENTRUN_MGR_URL);
const url = new URL(raw);
const host = url.hostname.toLowerCase();
const allowedHost = "agentrun-mgr.agentrun-v01.svc.cluster.local";
const allowedByTest = truthyFlag(env.HWLAB_CODE_AGENT_AGENTRUN_ALLOW_NON_K3S_URL) && ["127.0.0.1", "localhost"].includes(host);
if (url.protocol !== "http:" || (host !== allowedHost && !allowedByTest)) {
throw Object.assign(new Error(`AGENTRUN_MGR_URL must use internal k3s Service DNS ${DEFAULT_AGENTRUN_MGR_URL}; got ${redactUrl(raw)}`), {
code: "agentrun_internal_url_required",
statusCode: 500
});
}
url.pathname = url.pathname.replace(/\/+$/u, "");
url.search = "";
url.hash = "";
return url.toString().replace(/\/+$/u, "");
}
function resolveAgentRunBackendProfile(env = process.env, params = {}) {
const requested = String(params.providerProfile ?? params.codeAgentProviderProfile ?? env.HWLAB_CODE_AGENT_AGENTRUN_BACKEND_PROFILE ?? env.HWLAB_CODE_AGENT_DEFAULT_PROVIDER_PROFILE ?? "deepseek").trim().toLowerCase();
const fallback = String(env.HWLAB_CODE_AGENT_AGENTRUN_DEFAULT_BACKEND_PROFILE ?? "deepseek").trim().toLowerCase() || "deepseek";
const resolved = requested === "runtime-default" ? fallback : requested;
if (!resolved || resolved === "runtime-default") return "deepseek";
const mapped = AGENTRUN_BACKEND_ALIASES[resolved] ?? resolved;
return AGENTRUN_BACKEND_PROFILE_ID_PATTERN.test(mapped) ? mapped : "deepseek";
}
function agentRunMapping({ env, managerUrl, backendProfile, run, command, runnerJob, traceId, startedAt, params = {} }) {
const threadReused = Boolean(safeOpaqueId(params.threadId));
return {
adapter: ADAPTER_ID,
managerUrl,
backendProfile,
providerId: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID),
runId: run.id,
commandId: command.id,
attemptId: runnerJob?.attemptId ?? runnerJob?.runner?.attemptId ?? null,
runnerId: runnerJob?.runnerId ?? runnerJob?.runner?.runnerId ?? null,
runnerJobId: runnerJob?.id ?? null,
jobName: runnerJob?.jobName ?? runnerJob?.jobIdentity?.name ?? null,
namespace: runnerJob?.namespace ?? runnerJob?.jobIdentity?.namespace ?? DEFAULT_RUNNER_NAMESPACE,
status: "runner-job-created",
runStatus: run.status ?? null,
commandState: command.state ?? null,
terminalStatus: null,
lastSeq: 0,
traceId,
sessionId: run?.sessionRef?.sessionId ?? null,
conversationId: run?.sessionRef?.conversationId ?? null,
threadId: run?.sessionRef?.threadId ?? null,
runnerReused: false,
threadReused,
persistentResume: threadReused,
reused: false,
reuseEligible: true,
createdAt: startedAt,
updatedAt: nowIso(),
valuesPrinted: false
};
}
function agentRunReusedMapping({ previous = {}, run = {}, command = {}, runnerJob = null, traceId, startedAt, backendProfile, managerUrl, env }) {
const threadReused = Boolean(safeOpaqueId(run?.sessionRef?.threadId ?? previous.threadId));
return {
...previous,
adapter: ADAPTER_ID,
managerUrl,
backendProfile,
providerId: previous.providerId ?? firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID),
runId: previous.runId ?? run?.id,
commandId: command.id,
attemptId: runnerJob?.attemptId ?? runnerJob?.runner?.attemptId ?? previous.attemptId ?? null,
runnerId: runnerJob?.runnerId ?? runnerJob?.runner?.runnerId ?? previous.runnerId ?? null,
runnerJobId: runnerJob?.id ?? previous.runnerJobId ?? null,
jobName: runnerJob?.jobName ?? runnerJob?.jobIdentity?.name ?? previous.jobName ?? null,
namespace: runnerJob?.namespace ?? runnerJob?.jobIdentity?.namespace ?? previous.namespace ?? DEFAULT_RUNNER_NAMESPACE,
status: runnerJob ? "runner-job-ensured" : "runner-job-reused",
runStatus: run?.status ?? previous.runStatus ?? null,
commandState: command.state ?? null,
terminalStatus: null,
lastSeq: previous.lastSeq ?? 0,
runnerJobCount: runnerJob ? Math.max(1, Number(previous.runnerJobCount ?? 0) + 1) : Number(previous.runnerJobCount ?? 0),
traceId,
sessionId: run?.sessionRef?.sessionId ?? previous.sessionId ?? null,
conversationId: run?.sessionRef?.conversationId ?? previous.conversationId ?? null,
threadId: run?.sessionRef?.threadId ?? previous.threadId ?? null,
runnerReused: true,
threadReused,
persistentResume: threadReused,
reused: true,
reuseEligible: true,
createdAt: previous.createdAt ?? startedAt,
updatedAt: nowIso(),
valuesPrinted: false
};
}
function agentRunResultRefs(result = {}) {
const refs = {};
for (const key of ["attemptId", "runnerId", "jobName", "namespace", "runnerJobCount"]) {
if (result?.[key] !== undefined && result?.[key] !== null) refs[key] = result[key];
}
if (result?.sessionRef && typeof result.sessionRef === "object") {
if (result.sessionRef.sessionId) refs.sessionId = result.sessionRef.sessionId;
if (result.sessionRef.conversationId) refs.conversationId = result.sessionRef.conversationId;
if (result.sessionRef.threadId) refs.threadId = result.sessionRef.threadId;
}
return refs;
}
function hwlabSessionIdForParams(params = {}, traceId) {
return safeSessionId(params.sessionId) || agentRunSessionId(traceId);
}
function scopedAgentRunSessionIdForParams(params = {}, traceId, backendProfile) {
const baseSessionId = hwlabSessionIdForParams(params, traceId);
const profile = agentRunSessionProfileToken(backendProfile);
const base = String(baseSessionId).replace(/^ses_/u, "").replace(/[^A-Za-z0-9_]+/gu, "_").replace(/^_+|_+$/gu, "") || "session";
return `ses_agentrun_${profile}_${base}`;
}
function agentRunSessionProfileToken(backendProfile) {
return String(backendProfile ?? "deepseek").trim().toLowerCase().replace(/[^a-z0-9]+/gu, "_").replace(/^_+|_+$/gu, "") || "default";
}
function agentRunSessionId(traceId) {
return `ses_agentrun_${String(safeTraceId(traceId) || `trc_${randomUUID()}`).slice(4)}`;
}
function agentRunSessionSummary(base, status) {
return decorateCodeAgentSession({
sessionId: base.sessionId ?? base.agentRun?.sessionId ?? null,
conversationId: base.conversationId ?? base.agentRun?.conversationId ?? null,
threadId: base.threadId ?? base.agentRun?.threadId ?? null,
status,
sessionMode: AGENTRUN_SESSION_MODE,
runnerKind: AGENTRUN_RUNNER_KIND,
lastTraceId: base.traceId ?? base.agentRun?.traceId ?? null,
longLivedSession: true,
codexStdio: false,
delegatedToAgentRun: true,
writeCapable: true,
durable: true,
durableSession: true,
idleTimeoutMs: parsePositiveInteger(base.agentRun?.idleTimeoutMs, 600000),
secretMaterialStored: false,
valuesRedacted: true
});
}
function agentRunThreadReused(base = {}, options = {}) {
if (options.status === "failed-requires-new-session") return false;
if (base.agentRun?.threadReused === true || base.agentRun?.persistentResume === true) return true;
return Boolean(safeOpaqueId(base.threadId));
}
function agentRunSessionReuseSummary(base, reused, options = {}) {
const runnerReused = Boolean(reused);
const threadReused = agentRunThreadReused(base, options);
const persistentResume = threadReused;
const sessionReused = runnerReused || persistentResume;
return {
sessionId: base.sessionId ?? base.agentRun?.sessionId ?? null,
conversationId: base.conversationId ?? base.agentRun?.conversationId ?? null,
threadId: base.threadId ?? base.agentRun?.threadId ?? null,
reused: sessionReused,
status: options.status ?? (runnerReused ? "reused" : persistentResume ? "thread-resumed" : "new"),
runnerReused,
threadReused,
persistentResume,
valuesRedacted: true
};
}
function agentRunRunnerSummary(mapping = {}) {
return {
kind: AGENTRUN_RUNNER_KIND,
adapter: ADAPTER_ID,
runId: mapping.runId ?? null,
commandId: mapping.commandId ?? null,
attemptId: mapping.attemptId ?? null,
runnerId: mapping.runnerId ?? null,
jobName: mapping.jobName ?? null,
namespace: mapping.namespace ?? null,
backendProfile: mapping.backendProfile ?? null,
provider: providerForBackendProfile(mapping.backendProfile ?? "deepseek"),
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE,
capabilityLevel: AGENTRUN_CAPABILITY_LEVEL,
codexStdio: false,
delegatedToAgentRun: true,
writeCapable: true,
durableSession: true,
longLivedSession: true,
sandbox: "danger-full-access",
valuesPrinted: false
};
}
function agentRunTraceMeta(env = process.env) {
return {
runnerKind: AGENTRUN_RUNNER_KIND,
workspace: resolveAgentRunRepoUrl(env),
sandbox: firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_SANDBOX, env.HWLAB_CODE_AGENT_CODEX_SANDBOX, "danger-full-access"),
protocol: AGENTRUN_PROVIDER_TRACE_PROTOCOL,
sessionMode: AGENTRUN_SESSION_MODE,
implementationType: AGENTRUN_IMPLEMENTATION_TYPE
};
}
function backendForBackendProfile(profile) {
return `${AGENTRUN_BACKEND_PREFIX}/${resolveAgentRunBackendProfile({}, { providerProfile: profile })}`;
}
function resolveAgentRunRepoUrl(env = process.env) {
const raw = firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_REPO_URL, env.HWLAB_BOOT_READ_URL, DEFAULT_REPO_URL);
const url = new URL(raw);
const host = url.hostname.toLowerCase();
const allowedHost = "git-mirror-http.devops-infra.svc.cluster.local";
const allowedByTest = truthyFlag(env.HWLAB_CODE_AGENT_AGENTRUN_ALLOW_NON_K3S_URL) && ["127.0.0.1", "localhost"].includes(host);
if (url.protocol !== "http:" || (host !== allowedHost && !allowedByTest)) {
throw Object.assign(new Error(`HWLAB_CODE_AGENT_AGENTRUN_REPO_URL must use internal k3s git mirror ${DEFAULT_REPO_URL}; got ${redactUrl(raw)}`), {
code: "agentrun_internal_repo_url_required",
statusCode: 500
});
}
url.pathname = url.pathname.replace(/\/+$/u, "");
url.search = "";
url.hash = "";
return url.toString().replace(/\/+$/u, "");
}
function requireAgentRunSourceCommit(env) {
const text = String(env.HWLAB_CODE_AGENT_AGENTRUN_SOURCE_COMMIT ?? "").trim().toLowerCase();
if (/^[0-9a-f]{40}$/u.test(text)) return text;
throw Object.assign(new Error("HWLAB_CODE_AGENT_AGENTRUN_SOURCE_COMMIT must be a full 40-character source commit before AgentRun dispatch"), {
code: "agentrun_bundle_source_commit_invalid",
statusCode: 500
});
}
function modelForBackendProfile(profile, env = process.env) {
const resolved = resolveAgentRunBackendProfile({}, { providerProfile: profile });
if (resolved === "codex") return firstNonEmpty(env.HWLAB_CODE_AGENT_CODEX_API_MODEL, env.HWLAB_CODE_AGENT_MODEL, "gpt-5.5");
if (resolved === "minimax-m3") return firstNonEmpty(env.HWLAB_CODE_AGENT_MINIMAX_M3_MODEL, "MiniMax-M3");
if (resolved === "deepseek") return firstNonEmpty(env.HWLAB_CODE_AGENT_DEEPSEEK_MODEL, "deepseek-chat");
if (resolved === "dsflash-go") return firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_DSFLASH_GO_MODEL, env.HWLAB_CODE_AGENT_DSFLASH_GO_MODEL, "deepseek-v4-flash");
const profileEnvKey = dynamicBackendProfileEnvKey(resolved);
return firstNonEmpty(env[`HWLAB_CODE_AGENT_AGENTRUN_${profileEnvKey}_MODEL`], env[`HWLAB_CODE_AGENT_${profileEnvKey}_MODEL`], resolved);
}
function providerForBackendProfile(profile) {
const resolved = resolveAgentRunBackendProfile({}, { providerProfile: profile });
if (resolved === "codex") return "codex-api";
return resolved;
}
function dynamicBackendProfileEnvKey(profile) {
return String(profile ?? "")
.trim()
.toUpperCase()
.replace(/[^A-Z0-9]+/gu, "_")
.replace(/^_+|_+$/gu, "") || "PROFILE";
}
function agentRunAvailabilityBlocker(error, scope) {
return {
code: error?.code ?? `agentrun_${String(scope).replace(/[^a-z0-9]+/giu, "_")}_blocked`,
type: "agent_blocker",
scope: `agentrun-${scope}`,
status: "open",
sourceIssue: "pikasTech/HWLAB#879",
summary: error?.message ?? `AgentRun ${scope} is not configured correctly.` ,
secretMaterialRead: false,
valuesRedacted: true
};
}
function hostForUrl(value) {
if (!value) return null;
try {
return new URL(value).hostname.toLowerCase();
} catch {
return null;
}
}
function waitingForForPhase(phase) {
if (phase.includes("runner-job")) return "agentrun-runner";
if (phase.includes("resource-bundle")) return "agentrun-resource-bundle";
if (phase.includes("backend")) return "agentrun-backend";
return "agentrun-result";
}
function textPayload(payload, fallback) {
for (const key of ["message", "text", "content", "delta", "summary", "phase"]) {
if (typeof payload?.[key] === "string" && payload[key].length > 0) return payload[key];
}
return fallback;
}
function requiredString(value, fieldName) {
if (typeof value === "string" && value.trim()) return value.trim();
throw Object.assign(new Error(`AgentRun response missing ${fieldName}`), { code: "agentrun_response_invalid", statusCode: 502 });
}
function firstNonEmpty(...values) {
for (const value of values) {
const text = String(value ?? "").trim();
if (text) return text;
}
return null;
}
function nowIso(now) {
return typeof now === "function" ? now() : new Date().toISOString();
}
function redactUrl(value) {
try {
const url = new URL(value);
if (url.username || url.password) {
url.username = "***";
url.password = "";
}
return url.toString();
} catch {
return "<invalid-url>";
}
}