1325 lines
76 KiB
TypeScript
1325 lines
76 KiB
TypeScript
import { createHash, randomBytes, randomUUID } from "node:crypto";
|
|
|
|
import { CLOUD_API_SERVICE_ID } from "../audit/index.mjs";
|
|
import { getHeader, readBody, sendJson, truthyFlag } from "./server-http-utils.ts";
|
|
|
|
const SESSION_COOKIE = "hwlab_session";
|
|
const SESSION_MAX_AGE_SECONDS = 60 * 60 * 24 * 7;
|
|
const DEVICE_JOB_CONTRACT_VERSION = "device-pod-job-v1";
|
|
const DEVICE_JOB_INTENTS = new Set([
|
|
"workspace.ls",
|
|
"workspace.cat",
|
|
"workspace.rg",
|
|
"workspace.apply-patch",
|
|
"workspace.put",
|
|
"workspace.rm",
|
|
"workspace.rmdir",
|
|
"workspace.keil",
|
|
"workspace.build",
|
|
"debug.status",
|
|
"debug.chip-id",
|
|
"debug.download",
|
|
"debug.reset",
|
|
"io.ports",
|
|
"io.uart.read",
|
|
"io.uart.read-after-launch-flash",
|
|
"io.uart.write",
|
|
"io.uart.jsonrpc"
|
|
]);
|
|
const ACCESS_SCHEMA_STATEMENTS = Object.freeze([
|
|
`CREATE TABLE IF NOT EXISTS users (
|
|
id TEXT PRIMARY KEY,
|
|
username TEXT NOT NULL UNIQUE,
|
|
display_name TEXT NOT NULL DEFAULT '',
|
|
role TEXT NOT NULL CHECK (role IN ('admin', 'user')),
|
|
status TEXT NOT NULL DEFAULT 'active' CHECK (status IN ('active', 'disabled')),
|
|
password_hash TEXT,
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS user_sessions (
|
|
id TEXT PRIMARY KEY,
|
|
user_id TEXT NOT NULL REFERENCES users(id),
|
|
session_token_hash TEXT NOT NULL UNIQUE,
|
|
created_at TEXT NOT NULL,
|
|
last_seen_at TEXT NOT NULL,
|
|
expires_at TEXT NOT NULL,
|
|
revoked_at TEXT
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS agent_sessions (
|
|
id TEXT PRIMARY KEY,
|
|
project_id TEXT,
|
|
agent_id TEXT NOT NULL,
|
|
status TEXT NOT NULL DEFAULT 'pending',
|
|
started_at TEXT,
|
|
ended_at TEXT,
|
|
owner_user_id TEXT REFERENCES users(id),
|
|
conversation_id TEXT,
|
|
thread_id TEXT,
|
|
last_trace_id TEXT,
|
|
session_json TEXT NOT NULL DEFAULT '{}',
|
|
updated_at TEXT
|
|
)`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS owner_user_id TEXT REFERENCES users(id)`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS conversation_id TEXT`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS thread_id TEXT`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS last_trace_id TEXT`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS session_json TEXT NOT NULL DEFAULT '{}'`,
|
|
`ALTER TABLE agent_sessions ADD COLUMN IF NOT EXISTS updated_at TEXT`,
|
|
`CREATE INDEX IF NOT EXISTS idx_agent_sessions_owner ON agent_sessions(owner_user_id)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_agent_sessions_conversation ON agent_sessions(conversation_id)`,
|
|
`CREATE TABLE IF NOT EXISTS device_pods (
|
|
id TEXT PRIMARY KEY,
|
|
name TEXT NOT NULL DEFAULT '',
|
|
status TEXT NOT NULL DEFAULT 'active' CHECK (status IN ('active', 'disabled')),
|
|
profile_json TEXT NOT NULL DEFAULT '{}',
|
|
profile_hash TEXT NOT NULL DEFAULT '',
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS device_pod_grants (
|
|
device_pod_id TEXT NOT NULL REFERENCES device_pods(id) ON DELETE CASCADE,
|
|
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
|
created_by_admin_id TEXT NOT NULL REFERENCES users(id),
|
|
created_at TEXT NOT NULL,
|
|
PRIMARY KEY (device_pod_id, user_id)
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS device_leases (
|
|
device_pod_id TEXT PRIMARY KEY REFERENCES device_pods(id) ON DELETE CASCADE,
|
|
holder_session_id TEXT NOT NULL REFERENCES agent_sessions(id) ON DELETE CASCADE,
|
|
holder_user_id TEXT NOT NULL REFERENCES users(id),
|
|
lease_token_hash TEXT NOT NULL UNIQUE,
|
|
created_at TEXT NOT NULL,
|
|
expires_at TEXT NOT NULL,
|
|
released_at TEXT
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS device_pod_jobs (
|
|
id TEXT PRIMARY KEY,
|
|
device_pod_id TEXT NOT NULL REFERENCES device_pods(id) ON DELETE CASCADE,
|
|
owner_user_id TEXT NOT NULL REFERENCES users(id),
|
|
status TEXT NOT NULL,
|
|
intent TEXT NOT NULL,
|
|
args_json TEXT NOT NULL DEFAULT '{}',
|
|
reason TEXT NOT NULL DEFAULT '',
|
|
trace_id TEXT NOT NULL,
|
|
operation_id TEXT NOT NULL,
|
|
output_json TEXT NOT NULL DEFAULT '{}',
|
|
blocker_json TEXT NOT NULL DEFAULT '{}',
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL,
|
|
completed_at TEXT
|
|
)`
|
|
]);
|
|
|
|
const MUTATING_INTENTS = new Set([
|
|
"workspace.apply-patch",
|
|
"workspace.build",
|
|
"workspace.put",
|
|
"workspace.rm",
|
|
"workspace.rmdir",
|
|
"workspace.keil",
|
|
"debug.download",
|
|
"debug.reset",
|
|
"io.uart.read-after-launch-flash",
|
|
"io.uart.write",
|
|
"io.uart.jsonrpc"
|
|
]);
|
|
const DEFAULT_DEVICE_LEASE_TTL_SECONDS = 15 * 60;
|
|
const MAX_DEVICE_LEASE_TTL_SECONDS = 60 * 60;
|
|
const DEVICE_JOB_OUTPUT_MAX_BYTES = 12000;
|
|
const DEVICE_POD_PROFILE_SECRET_KEY_PATTERN = /(?:token|secret|password|passphrase|credential|api[_-]?key|access[_-]?key|private[_-]?key|git[_-]?key|cloud[_-]?token|kubeconfig|database[_-]?url|db[_-]?url|connection[_-]?string)/iu;
|
|
const DEVICE_POD_PROFILE_SECRET_VALUE_PATTERN = /(?:-----BEGIN [A-Z ]*PRIVATE KEY-----|postgres(?:ql)?:\/\/|mysql:\/\/|mongodb(?:\+srv)?:\/\/|redis:\/\/|gh[pousr]_[A-Za-z0-9_]{20,}|sk-[A-Za-z0-9_-]{20,})/u;
|
|
|
|
export function createAccessController(options = {}) {
|
|
return new AccessController({
|
|
...options,
|
|
store: options.store ?? accessStoreForRuntime(options.runtimeStore, options)
|
|
});
|
|
}
|
|
|
|
function accessStoreForRuntime(runtimeStore, options = {}) {
|
|
if (runtimeStore && typeof runtimeStore.query === "function") {
|
|
return new PostgresAccessStore({ query: runtimeStore.query.bind(runtimeStore), now: options.now });
|
|
}
|
|
return new MemoryAccessStore({ now: options.now });
|
|
}
|
|
|
|
class AccessController {
|
|
constructor({ store, env = process.env, gatewayRegistry = null, fetchImpl = fetch, devicePodExecutorUrl = env.HWLAB_DEVICE_POD_URL, devicePodExecutorTimeoutMs = 1200, now = () => new Date().toISOString(), required = truthyFlag(env.HWLAB_ACCESS_CONTROL_REQUIRED) } = {}) {
|
|
this.store = store;
|
|
this.env = env;
|
|
this.gatewayRegistry = gatewayRegistry;
|
|
this.fetchImpl = fetchImpl;
|
|
this.devicePodExecutorUrl = normalizeBaseUrl(devicePodExecutorUrl);
|
|
this.devicePodInternalToken = textOr(env.HWLAB_DEVICE_POD_INTERNAL_TOKEN, "");
|
|
this.devicePodExecutorTimeoutMs = Number.parseInt(String(devicePodExecutorTimeoutMs), 10) || 1200;
|
|
this.now = now;
|
|
this.required = required;
|
|
this.bootstrapAttempted = false;
|
|
}
|
|
|
|
async handleAuthRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
if (request.method === "GET" && url.pathname === "/auth/session") {
|
|
sendJson(response, 200, await this.sessionPayload(request));
|
|
return;
|
|
}
|
|
|
|
if (request.method === "POST" && url.pathname === "/auth/login") {
|
|
await this.handleLogin(request, response);
|
|
return;
|
|
}
|
|
|
|
if (request.method === "POST" && url.pathname === "/auth/logout") {
|
|
await this.handleLogout(request, response);
|
|
return;
|
|
}
|
|
|
|
sendJson(response, 404, errorPayload("not_found", "Auth route is not implemented", 404));
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async handleUserRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
if (request.method === "GET" && url.pathname === "/v1/users/me") {
|
|
const auth = await this.authenticate(request, { required: true });
|
|
if (!auth.ok) return sendJson(response, auth.status, auth);
|
|
return sendJson(response, 200, {
|
|
authenticated: true,
|
|
actor: publicActor(auth.actor),
|
|
session: auth.session,
|
|
contractVersion: "user-access-v1"
|
|
});
|
|
}
|
|
|
|
sendJson(response, 404, errorPayload("not_found", "User route is not implemented", 404));
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async handleAccessStatusRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
if (request.method === "GET" && url.pathname === "/v1/access/status") {
|
|
return sendJson(response, 200, await this.accessStatusPayload(request));
|
|
}
|
|
if (request.method === "GET" && url.pathname === "/v1/setup/status") {
|
|
return sendJson(response, 200, await this.setupStatusPayload());
|
|
}
|
|
|
|
sendJson(response, 404, errorPayload("not_found", "Access status route is not implemented", 404));
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async handleSetupRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
if (request.method !== "POST" || url.pathname !== "/v1/setup/first-admin") {
|
|
return sendJson(response, 404, errorPayload("not_found", "Setup route is not implemented", 404));
|
|
}
|
|
|
|
if (!(await this.setupRequired())) {
|
|
return sendJson(response, 409, errorPayload("setup_already_completed", "First admin setup is closed because at least one user already exists", 409));
|
|
}
|
|
|
|
const body = await jsonBody(request);
|
|
const password = requiredText(body.password, "password");
|
|
const devicePodSeeds = firstAdminDevicePodSeeds(body);
|
|
const now = this.now();
|
|
const user = await this.store.createFirstAdmin?.({
|
|
id: this.env.HWLAB_BOOTSTRAP_ADMIN_ID || "usr_bootstrap_admin",
|
|
username: textOr(body.username, this.env.HWLAB_BOOTSTRAP_ADMIN_USERNAME || "admin"),
|
|
displayName: textOr(body.displayName, this.env.HWLAB_BOOTSTRAP_ADMIN_DISPLAY_NAME || "HWLAB Admin"),
|
|
passwordHash: hashPassword(password),
|
|
now
|
|
}) ?? null;
|
|
if (!user) {
|
|
return sendJson(response, 409, errorPayload("setup_already_completed", "First admin setup is closed because at least one user already exists", 409));
|
|
}
|
|
const devicePodsInitialized = [];
|
|
for (const seed of devicePodSeeds) {
|
|
const pod = await this.store.upsertDevicePod({ ...seed, now });
|
|
const grant = await this.store.grantDevicePod({
|
|
devicePodId: pod.id,
|
|
userId: user.id,
|
|
createdByAdminId: user.id,
|
|
now
|
|
});
|
|
devicePodsInitialized.push({
|
|
devicePod: publicDevicePod(pod, { includeAdmin: true }),
|
|
grant
|
|
});
|
|
}
|
|
const token = randomBytes(32).toString("base64url");
|
|
const session = await this.store.createSession({
|
|
userId: user.id,
|
|
tokenHash: sha256(token),
|
|
now,
|
|
expiresAt: new Date(Date.parse(now) + SESSION_MAX_AGE_SECONDS * 1000).toISOString()
|
|
});
|
|
setSessionCookie(response, token, SESSION_MAX_AGE_SECONDS);
|
|
return sendJson(response, 201, {
|
|
created: true,
|
|
authenticated: true,
|
|
actor: publicActor(user),
|
|
session: publicSession({ ...session, user }),
|
|
setupRequired: false,
|
|
devicePodBootstrap: {
|
|
requested: devicePodSeeds.length,
|
|
initialized: devicePodsInitialized.length
|
|
},
|
|
devicePodsInitialized,
|
|
contractVersion: "user-access-v1"
|
|
});
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async sessionPayload(request) {
|
|
const auth = await this.authenticate(request, { required: false });
|
|
return {
|
|
authenticated: Boolean(auth.actor),
|
|
actor: auth.actor ? publicActor(auth.actor) : null,
|
|
session: auth.session ?? null,
|
|
setupRequired: await this.setupRequired(),
|
|
contractVersion: "user-access-v1"
|
|
};
|
|
}
|
|
|
|
async accessStatusPayload(request) {
|
|
const session = await this.sessionPayload(request);
|
|
return {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
...session,
|
|
accessControlRequired: this.required,
|
|
roles: ["admin", "user"],
|
|
routes: accessRoutes()
|
|
};
|
|
}
|
|
|
|
async setupStatusPayload() {
|
|
return {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "user-access-v1",
|
|
accessControlRequired: this.required,
|
|
setupRequired: await this.setupRequired(),
|
|
bootstrapAdminConfigured: Boolean(textOr(this.env.HWLAB_BOOTSTRAP_ADMIN_PASSWORD_HASH, "") || textOr(this.env.HWLAB_BOOTSTRAP_ADMIN_PASSWORD, "")),
|
|
routes: accessRoutes()
|
|
};
|
|
}
|
|
|
|
async handleAdminRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
const auth = await this.authenticate(request, { required: true });
|
|
if (!auth.ok) return sendJson(response, auth.status, auth);
|
|
if (auth.actor.role !== "admin") {
|
|
return sendJson(response, 403, errorPayload("admin_required", "Only admin users can call v0.2 admin APIs", 403));
|
|
}
|
|
|
|
if (request.method === "POST" && url.pathname === "/v1/admin/users") {
|
|
const body = await jsonBody(request);
|
|
const user = await this.store.createUser({
|
|
username: requiredText(body.username, "username"),
|
|
displayName: textOr(body.displayName, body.username),
|
|
role: body.role === "admin" ? "admin" : "user",
|
|
status: body.status === "disabled" ? "disabled" : "active",
|
|
passwordHash: hashPassword(requiredText(body.password, "password")),
|
|
now: this.now()
|
|
});
|
|
return sendJson(response, 201, { created: true, user: redactedUser(user) });
|
|
}
|
|
|
|
if ((request.method === "POST" && url.pathname === "/v1/admin/device-pods") || (request.method === "PUT" && url.pathname.startsWith("/v1/admin/device-pods/"))) {
|
|
const body = await jsonBody(request);
|
|
const devicePodId = request.method === "PUT"
|
|
? decodeURIComponent(url.pathname.split("/").filter(Boolean)[3] ?? "")
|
|
: textOr(body.devicePodId ?? body.id, "");
|
|
const profile = assertDevicePodProfileSafe(normalizeObject(body.profile ?? body.profileJson), "profile");
|
|
const pod = await this.store.upsertDevicePod({
|
|
id: requiredText(devicePodId, "devicePodId"),
|
|
name: textOr(body.name, devicePodId),
|
|
status: body.status === "disabled" ? "disabled" : "active",
|
|
profile,
|
|
now: this.now()
|
|
});
|
|
return sendJson(response, request.method === "POST" ? 201 : 200, {
|
|
upserted: true,
|
|
devicePod: publicDevicePod(pod, { includeAdmin: true })
|
|
});
|
|
}
|
|
|
|
if (request.method === "POST" && url.pathname === "/v1/admin/device-pod-grants") {
|
|
const body = await jsonBody(request);
|
|
const devicePodId = requiredText(body.devicePodId, "devicePodId");
|
|
const userId = requiredText(body.userId, "userId");
|
|
await this.requireGrantTargets({ devicePodId, userId });
|
|
const grant = await this.store.grantDevicePod({
|
|
devicePodId,
|
|
userId,
|
|
createdByAdminId: auth.actor.id,
|
|
now: this.now()
|
|
});
|
|
return sendJson(response, 201, { granted: true, grant });
|
|
}
|
|
|
|
if (request.method === "DELETE" && url.pathname.startsWith("/v1/admin/device-pod-grants/")) {
|
|
const [, , , devicePodId, userId] = url.pathname.split("/").filter(Boolean).map((part) => decodeURIComponent(part));
|
|
await this.store.revokeDevicePodGrant({ devicePodId: requiredText(devicePodId, "devicePodId"), userId: requiredText(userId, "userId") });
|
|
const releasedLease = await this.store.releaseDeviceLease?.({ devicePodId, holderUserId: userId, now: this.now(), allowAdmin: false, revokeGrant: true });
|
|
return sendJson(response, 200, { revoked: true, devicePodId, userId, releasedLease: publicDeviceLease(releasedLease) });
|
|
}
|
|
|
|
sendJson(response, 404, errorPayload("not_found", "Admin route is not implemented", 404));
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async handleDevicePodRoute(request, response, url) {
|
|
try {
|
|
await this.ensureBootstrap();
|
|
const auth = await this.authenticate(request, { required: true });
|
|
if (!auth.ok) return sendJson(response, auth.status, auth);
|
|
|
|
if (request.method === "GET" && url.pathname === "/v1/device-pods") {
|
|
const pods = await this.store.listVisibleDevicePods(auth.actor);
|
|
return sendJson(response, 200, {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-authority-v1",
|
|
status: "ok",
|
|
source: authoritySource(),
|
|
actor: publicActor(auth.actor),
|
|
devicePods: pods.map((pod) => publicDevicePod(pod)),
|
|
selectedDevicePodId: pods[0]?.id ?? null,
|
|
observedAt: this.now()
|
|
});
|
|
}
|
|
|
|
const parsed = parseDevicePodPath(url.pathname);
|
|
if (!parsed) return sendJson(response, 404, errorPayload("not_found", "Device Pod route is not implemented", 404));
|
|
const pod = await this.store.getVisibleDevicePod(auth.actor, parsed.devicePodId);
|
|
if (!pod) {
|
|
const existing = await this.store.getDevicePod(parsed.devicePodId);
|
|
if (existing?.status === "active") {
|
|
return sendJson(response, 403, errorPayload("device_pod_forbidden", `Device Pod ${parsed.devicePodId} is not authorized for the current actor`, 403));
|
|
}
|
|
return sendJson(response, 404, errorPayload("device_pod_not_found", `Device Pod ${parsed.devicePodId} was not found`, 404));
|
|
}
|
|
|
|
if (request.method === "GET" && parsed.route === "status") return sendJson(response, 200, this.devicePodStatus(pod, auth.actor));
|
|
if (request.method === "GET" && parsed.route === "events") return sendJson(response, 200, await this.devicePodEvents(pod, url.searchParams));
|
|
if (request.method === "GET" && parsed.route === "debug-probe/chip-id") return this.createDevicePodProbeJob(response, pod, auth.actor, { interfaceName: "debug-probe", intent: "debug.chip-id" });
|
|
if (request.method === "GET" && parsed.route === "io-probe/uart/1") return this.createDevicePodProbeJob(response, pod, auth.actor, { interfaceName: "io-probe", intent: "io.ports", args: { uartId: "uart/1" } });
|
|
if (request.method === "GET" && parsed.route === "io-probe/uart/1/tail") return this.createDevicePodProbeJob(response, pod, auth.actor, { interfaceName: "io-probe", intent: "io.uart.read", args: uartTailArgs(url.searchParams, "uart/1"), outputMaxBytes: boundedOutputMaxBytes(url.searchParams) });
|
|
if (request.method === "POST" && parsed.route === "leases") return this.acquireDevicePodLease(request, response, pod, auth.actor);
|
|
if (request.method === "GET" && parsed.route === "leases/current") return this.getDevicePodLease(response, pod, auth.actor);
|
|
if (request.method === "DELETE" && parsed.route === "leases/current") return this.releaseDevicePodLease(request, response, pod, auth.actor);
|
|
if (request.method === "POST" && parsed.route === "jobs") return this.createDevicePodJob(request, response, pod, auth.actor);
|
|
if (request.method === "GET" && parsed.route.startsWith("jobs/")) return this.getDevicePodJob(response, pod, parsed.route, auth.actor);
|
|
if (request.method === "POST" && parsed.route.startsWith("jobs/") && parsed.route.endsWith("/cancel")) return this.cancelDevicePodJob(response, pod, parsed.route, auth.actor);
|
|
|
|
sendJson(response, 404, errorPayload("not_found", "Device Pod route is not implemented", 404));
|
|
} catch (error) {
|
|
sendAccessError(response, error);
|
|
}
|
|
}
|
|
|
|
async authenticate(request, { required = this.required } = {}) {
|
|
await this.ensureBootstrap();
|
|
const token = sessionTokenFromRequest(request);
|
|
if (!token) {
|
|
return required ? errorPayload("auth_required", "Authentication is required", 401) : { ok: true, actor: null };
|
|
}
|
|
const session = await this.store.findSessionByTokenHash(sha256(token), this.now());
|
|
if (!session?.user || session.user.status !== "active") {
|
|
return errorPayload("auth_session_invalid", "Session is missing, expired, revoked, or disabled", 401);
|
|
}
|
|
await this.store.touchSession(session.id, this.now());
|
|
return { ok: true, actor: session.user, session: publicSession(session) };
|
|
}
|
|
|
|
async handleLogin(request, response) {
|
|
const body = await jsonBody(request);
|
|
const user = await this.store.findUserByUsername(requiredText(body.username, "username"));
|
|
if (!user || user.status !== "active" || !verifyPassword(user.passwordHash, requiredText(body.password, "password"))) {
|
|
return sendJson(response, 401, errorPayload("invalid_credentials", "Username or password is invalid", 401));
|
|
}
|
|
const token = randomBytes(32).toString("base64url");
|
|
const now = this.now();
|
|
const session = await this.store.createSession({
|
|
userId: user.id,
|
|
tokenHash: sha256(token),
|
|
now,
|
|
expiresAt: new Date(Date.parse(now) + SESSION_MAX_AGE_SECONDS * 1000).toISOString()
|
|
});
|
|
setSessionCookie(response, token, SESSION_MAX_AGE_SECONDS);
|
|
return sendJson(response, 200, {
|
|
authenticated: true,
|
|
actor: publicActor(user),
|
|
session: publicSession({ ...session, user })
|
|
});
|
|
}
|
|
|
|
async handleLogout(request, response) {
|
|
const token = sessionTokenFromRequest(request);
|
|
if (token) await this.store.revokeSessionByTokenHash(sha256(token), this.now());
|
|
clearSessionCookie(response);
|
|
return sendJson(response, 200, { authenticated: false, loggedOut: true });
|
|
}
|
|
|
|
async ensureBootstrap() {
|
|
if (this.bootstrapAttempted) return;
|
|
this.bootstrapAttempted = true;
|
|
await this.store.ensureSchema?.();
|
|
const count = await this.store.countUsers();
|
|
const passwordHash = textOr(this.env.HWLAB_BOOTSTRAP_ADMIN_PASSWORD_HASH, "") || (this.env.HWLAB_BOOTSTRAP_ADMIN_PASSWORD ? hashPassword(this.env.HWLAB_BOOTSTRAP_ADMIN_PASSWORD) : "");
|
|
if (count > 0) {
|
|
if (passwordHash) await this.store.syncBootstrapAdminPassword?.({
|
|
id: this.env.HWLAB_BOOTSTRAP_ADMIN_ID || "usr_bootstrap_admin",
|
|
username: this.env.HWLAB_BOOTSTRAP_ADMIN_USERNAME || "admin",
|
|
passwordHash,
|
|
now: this.now()
|
|
});
|
|
return;
|
|
}
|
|
if (!passwordHash) return;
|
|
await this.store.createUser({
|
|
id: this.env.HWLAB_BOOTSTRAP_ADMIN_ID || "usr_bootstrap_admin",
|
|
username: this.env.HWLAB_BOOTSTRAP_ADMIN_USERNAME || "admin",
|
|
displayName: this.env.HWLAB_BOOTSTRAP_ADMIN_DISPLAY_NAME || "HWLAB Admin",
|
|
role: "admin",
|
|
status: "active",
|
|
passwordHash,
|
|
now: this.now()
|
|
});
|
|
}
|
|
|
|
async setupRequired() {
|
|
return (await this.store.countUsers()) === 0;
|
|
}
|
|
|
|
async recordAgentSessionOwner(input = {}) {
|
|
await this.ensureBootstrap();
|
|
const ownerUserId = textOr(input.ownerUserId, "");
|
|
const sessionId = agentSessionIdForRecord(input.sessionId, input.traceId);
|
|
if (!ownerUserId || !sessionId) return null;
|
|
return this.store.recordAgentSessionOwner?.({
|
|
...input,
|
|
sessionId,
|
|
ownerUserId,
|
|
projectId: textOr(input.projectId, "prj_v02_code_agent"),
|
|
agentId: textOr(input.agentId, "hwlab-code-agent"),
|
|
status: textOr(input.status, "active"),
|
|
now: this.now()
|
|
}) ?? null;
|
|
}
|
|
|
|
async requireGrantTargets({ devicePodId, userId }) {
|
|
const pod = await this.store.getDevicePod(devicePodId);
|
|
if (!pod || pod.status !== "active") {
|
|
throw Object.assign(new Error(`Device Pod ${devicePodId} does not exist or is disabled`), { statusCode: 404, code: "device_pod_not_found" });
|
|
}
|
|
const user = await this.store.getUserById(userId);
|
|
if (!user || user.status !== "active") {
|
|
throw Object.assign(new Error(`User ${userId} does not exist or is disabled`), { statusCode: 404, code: "user_not_found" });
|
|
}
|
|
}
|
|
|
|
async acquireDevicePodLease(request, response, pod, actor) {
|
|
const body = await jsonBody(request);
|
|
const now = this.now();
|
|
const ttlSeconds = boundedLeaseTtlSeconds(body.ttlSeconds);
|
|
const expiresAt = new Date(Date.parse(now) + ttlSeconds * 1000).toISOString();
|
|
const token = randomBytes(32).toString("base64url");
|
|
const holderSessionId = safeAgentSessionId(body.agentSessionId) || `ses_devicepod_${randomUUID()}`;
|
|
const current = await this.store.getActiveDeviceLease?.(pod.id, now);
|
|
if (current && current.holderUserId !== actor.id) {
|
|
return sendJson(response, 409, {
|
|
ok: false,
|
|
status: 409,
|
|
error: {
|
|
code: "device_lease_conflict",
|
|
message: "Device Pod already has an active lease held by another actor"
|
|
},
|
|
lease: publicDeviceLease(current),
|
|
devicePodId: pod.id,
|
|
profileHash: pod.profileHash
|
|
});
|
|
}
|
|
await this.store.recordAgentSessionOwner?.({
|
|
sessionId: holderSessionId,
|
|
ownerUserId: actor.id,
|
|
projectId: textOr(body.projectId, "prj_v02_device_pod"),
|
|
agentId: "hwlab-device-pod",
|
|
status: "active",
|
|
session: { source: "device-pod-lease", reason: textOr(body.reason, "") },
|
|
now
|
|
});
|
|
const lease = await this.store.acquireDeviceLease?.({
|
|
devicePodId: pod.id,
|
|
holderSessionId,
|
|
holderUserId: actor.id,
|
|
leaseTokenHash: sha256(token),
|
|
now,
|
|
expiresAt
|
|
});
|
|
if (!lease) {
|
|
const active = await this.store.getActiveDeviceLease?.(pod.id, now);
|
|
return sendJson(response, 409, {
|
|
ok: false,
|
|
status: 409,
|
|
error: {
|
|
code: "device_lease_conflict",
|
|
message: "Device Pod already has an active lease held by another actor"
|
|
},
|
|
lease: publicDeviceLease(active),
|
|
devicePodId: pod.id,
|
|
profileHash: pod.profileHash
|
|
});
|
|
}
|
|
return sendJson(response, 201, {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-lease-v1",
|
|
acquired: true,
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
profileHash: pod.profileHash,
|
|
lease: publicDeviceLease(lease),
|
|
leaseToken: token,
|
|
tokenUsage: "Pass as body.leaseToken or x-hwlab-device-lease-token when creating mutating jobs.",
|
|
source: authoritySource()
|
|
});
|
|
}
|
|
|
|
async getDevicePodLease(response, pod, actor) {
|
|
const lease = await this.store.getActiveDeviceLease?.(pod.id, this.now());
|
|
if (!lease || (actor.role !== "admin" && lease.holderUserId !== actor.id)) {
|
|
return sendJson(response, 200, {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-lease-v1",
|
|
devicePodId: pod.id,
|
|
profileHash: pod.profileHash,
|
|
lease: null,
|
|
source: authoritySource()
|
|
});
|
|
}
|
|
return sendJson(response, 200, {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-lease-v1",
|
|
devicePodId: pod.id,
|
|
profileHash: pod.profileHash,
|
|
lease: publicDeviceLease(lease),
|
|
source: authoritySource()
|
|
});
|
|
}
|
|
|
|
async releaseDevicePodLease(request, response, pod, actor) {
|
|
const token = leaseTokenFromRequest(request, null);
|
|
const released = await this.store.releaseDeviceLease?.({
|
|
devicePodId: pod.id,
|
|
holderUserId: actor.id,
|
|
tokenHash: token ? sha256(token) : null,
|
|
now: this.now(),
|
|
allowAdmin: actor.role === "admin"
|
|
});
|
|
return sendJson(response, released ? 200 : 404, {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-lease-v1",
|
|
released: Boolean(released),
|
|
devicePodId: pod.id,
|
|
profileHash: pod.profileHash,
|
|
lease: publicDeviceLease(released),
|
|
source: authoritySource()
|
|
});
|
|
}
|
|
|
|
devicePodStatus(pod, actor) {
|
|
const route = pod.profile.route ?? {};
|
|
const gatewaySessionId = textOr(route.gatewaySessionId, "");
|
|
const gatewayOnline = gatewaySessionId && this.gatewayRegistry?.isOnline?.(gatewaySessionId) === true;
|
|
const blocker = gatewayOnline ? null : gatewayDispatchBlocker(gatewaySessionId ? "gateway session is not connected" : "profile route.gatewaySessionId is missing");
|
|
const observedAt = this.now();
|
|
const status = blocker ? "blocked" : "ok";
|
|
const output = boundedOutput({ text: blocker?.summary ?? "device-pod status ok", summary: blocker?.summary ?? "device-pod status ok" });
|
|
const freshnessPayload = freshness(observedAt, blocker);
|
|
return {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-authority-v1",
|
|
status,
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
profileHash: pod.profileHash,
|
|
traceId: `trc_devicepod_${randomUUID()}`,
|
|
operationId: `op_devicepod_${randomUUID()}`,
|
|
observedAt,
|
|
source: authoritySource(),
|
|
actor: publicActor(actor),
|
|
devicePod: publicDevicePod(pod),
|
|
blocker,
|
|
freshness: freshnessPayload,
|
|
output: output.output,
|
|
text: output.text,
|
|
bytes: output.bytes,
|
|
truncation: output.truncation,
|
|
summary: {
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
status,
|
|
freshness: freshnessPayload,
|
|
profileHash: pod.profileHash,
|
|
blocker
|
|
}
|
|
};
|
|
}
|
|
|
|
async devicePodEvents(pod, searchParams) {
|
|
const limit = Math.min(Math.max(Number.parseInt(searchParams.get("limit") ?? "80", 10) || 80, 1), 1000);
|
|
const jobs = await this.store.listDevicePodJobs(pod.id, limit);
|
|
const events = jobs.map((job) => ({
|
|
eventId: `evt_${job.id}`,
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
ts: job.updatedAt,
|
|
level: job.status === "blocked" || job.status === "failed" ? "warn" : "info",
|
|
scope: "job",
|
|
intent: job.intent,
|
|
status: job.status,
|
|
summary: normalizeBlocker(job.blocker)?.summary ?? `job ${job.id} ${job.status}`,
|
|
blocker: normalizeBlocker(job.blocker),
|
|
refs: jobRefs(job)
|
|
}));
|
|
return {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: "device-pod-authority-v1",
|
|
status: "ok",
|
|
observedAt: this.now(),
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
source: authoritySource(),
|
|
events,
|
|
lines: events.map(formatEventLine),
|
|
truncation: { limit, returned: events.length, truncated: jobs.length >= limit }
|
|
};
|
|
}
|
|
|
|
async createDevicePodProbeJob(response, pod, actor, { interfaceName, intent, args = {}, outputMaxBytes = DEVICE_JOB_OUTPUT_MAX_BYTES }) {
|
|
const result = await this.createReadOnlyDevicePodJob({ pod, actor, intent, args });
|
|
const payload = this.jobOutputPayload(result.job, pod, { maxBytes: outputMaxBytes });
|
|
return sendJson(response, result.httpStatus, {
|
|
...payload,
|
|
interface: interfaceName,
|
|
intent,
|
|
source: authoritySource()
|
|
});
|
|
}
|
|
|
|
async createReadOnlyDevicePodJob({ pod, actor, intent, args = {} }) {
|
|
if (!DEVICE_JOB_INTENTS.has(intent) || MUTATING_INTENTS.has(intent)) {
|
|
throw Object.assign(new Error(`Device probe intent ${intent} is not supported as a read-only probe`), { statusCode: 400, code: "unsupported_device_probe_intent" });
|
|
}
|
|
const traceId = `trc_devicepod_${randomUUID()}`;
|
|
const operationId = `op_devicepod_${randomUUID()}`;
|
|
const now = this.now();
|
|
const job = await this.store.createDevicePodJob({
|
|
id: `job_devicepod_${randomUUID()}`,
|
|
devicePodId: pod.id,
|
|
ownerUserId: actor.id,
|
|
status: "running",
|
|
intent,
|
|
args: normalizeObject(args),
|
|
reason: "",
|
|
traceId,
|
|
operationId,
|
|
output: { text: "" },
|
|
blocker: null,
|
|
now,
|
|
completedAt: null
|
|
});
|
|
if (this.devicePodExecutorUrl) return this.dispatchDevicePodExecutorJob({ job, pod, actor });
|
|
return this.blockDevicePodJobWithoutExecutor({ job, pod });
|
|
}
|
|
|
|
async createDevicePodJob(request, response, pod, actor) {
|
|
const body = await jsonBody(request);
|
|
const intent = requiredText(body.intent, "intent");
|
|
if (!DEVICE_JOB_INTENTS.has(intent)) {
|
|
return sendJson(response, 400, errorPayload("unsupported_device_job_intent", `Device job intent ${intent} is not supported`, 400));
|
|
}
|
|
const reason = textOr(body.reason, "");
|
|
if (MUTATING_INTENTS.has(intent) && !reason) {
|
|
return this.rejectDevicePodJob(response, pod, actor, {
|
|
httpStatus: 400,
|
|
intent,
|
|
args: normalizeObject(body.args),
|
|
blocker: deviceJobBlocked("device_job_reason_required", `Device job intent ${intent} requires reason`, false)
|
|
});
|
|
}
|
|
let lease = null;
|
|
if (MUTATING_INTENTS.has(intent)) {
|
|
const leaseToken = leaseTokenFromRequest(request, body);
|
|
if (!leaseToken) {
|
|
return this.rejectDevicePodJob(response, pod, actor, {
|
|
httpStatus: 409,
|
|
intent,
|
|
reason,
|
|
args: normalizeObject(body.args),
|
|
blocker: deviceJobBlocked("device_lease_required", `Device job intent ${intent} requires an active device lease`, true)
|
|
});
|
|
}
|
|
lease = await this.store.getDeviceLeaseByToken?.({ devicePodId: pod.id, tokenHash: sha256(leaseToken), now: this.now() });
|
|
if (!lease || (lease.holderUserId !== actor.id && actor.role !== "admin")) {
|
|
return this.rejectDevicePodJob(response, pod, actor, {
|
|
httpStatus: 409,
|
|
intent,
|
|
reason,
|
|
args: normalizeObject(body.args),
|
|
blocker: deviceJobBlocked("device_lease_invalid", "Device lease is missing, expired, released, or held by another actor", true)
|
|
});
|
|
}
|
|
}
|
|
const traceId = `trc_devicepod_${randomUUID()}`;
|
|
const operationId = `op_devicepod_${randomUUID()}`;
|
|
const now = this.now();
|
|
const args = normalizeObject(body.args);
|
|
const job = await this.store.createDevicePodJob({
|
|
id: `job_devicepod_${randomUUID()}`,
|
|
devicePodId: pod.id,
|
|
ownerUserId: actor.id,
|
|
status: "running",
|
|
intent,
|
|
args,
|
|
reason,
|
|
traceId,
|
|
operationId,
|
|
output: {},
|
|
blocker: null,
|
|
now,
|
|
completedAt: null
|
|
});
|
|
const dispatched = this.devicePodExecutorUrl
|
|
? await this.dispatchDevicePodExecutorJob({ job, pod, actor })
|
|
: await this.blockDevicePodJobWithoutExecutor({ job, pod });
|
|
sendJson(response, dispatched.httpStatus, this.jobPayload(dispatched.job, pod, { lease }));
|
|
}
|
|
|
|
async rejectDevicePodJob(response, pod, actor, { httpStatus, intent, reason = "", args = {}, blocker }) {
|
|
const now = this.now();
|
|
const job = await this.store.createDevicePodJob({
|
|
id: `job_devicepod_${randomUUID()}`,
|
|
devicePodId: pod.id,
|
|
ownerUserId: actor.id,
|
|
status: "blocked",
|
|
intent,
|
|
args,
|
|
reason,
|
|
traceId: `trc_devicepod_${randomUUID()}`,
|
|
operationId: `op_devicepod_${randomUUID()}`,
|
|
output: { error: blocker.summary },
|
|
blocker,
|
|
now,
|
|
completedAt: now
|
|
});
|
|
return sendJson(response, httpStatus, {
|
|
...this.jobPayload(job, pod),
|
|
error: { code: blocker.code, message: blocker.summary }
|
|
});
|
|
}
|
|
|
|
async dispatchDevicePodExecutorJob({ job, pod, actor }) {
|
|
const target = `${this.devicePodExecutorUrl}/v1/device-pods/${encodeURIComponent(pod.id)}/jobs`;
|
|
try {
|
|
const response = await fetchJsonWithTimeout(this.fetchImpl, target, {
|
|
method: "POST",
|
|
headers: {
|
|
accept: "application/json",
|
|
"content-type": "application/json",
|
|
...devicePodInternalHeaders(CLOUD_API_SERVICE_ID, this.devicePodInternalToken)
|
|
},
|
|
body: stableJson({
|
|
jobId: job.id,
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
profileHash: pod.profileHash,
|
|
profile: pod.profile,
|
|
ownerUserId: actor.id,
|
|
intent: job.intent,
|
|
args: job.args,
|
|
reason: job.reason,
|
|
traceId: job.traceId,
|
|
operationId: job.operationId
|
|
})
|
|
}, this.devicePodExecutorTimeoutMs);
|
|
const updated = await this.updateDevicePodJobFromExecutorResponse({ job, pod, response });
|
|
const status = updated.status;
|
|
return { job: updated ?? job, httpStatus: response.status >= 400 ? response.status : status === "running" ? 202 : 200 };
|
|
} catch (error) {
|
|
const blocker = devicePodExecutorBlocker(error?.message ?? "device-pod executor request failed");
|
|
const updated = await this.store.updateDevicePodJob(pod.id, job.id, {
|
|
status: "blocked",
|
|
output: { error: blocker.summary },
|
|
blocker,
|
|
completedAt: this.now(),
|
|
updatedAt: this.now()
|
|
});
|
|
return { job: updated ?? job, httpStatus: 409 };
|
|
}
|
|
}
|
|
|
|
async getDevicePodJob(response, pod, route, actor) {
|
|
const parts = route.split("/");
|
|
const jobId = decodeURIComponent(parts[1] ?? "");
|
|
let job = await this.store.getDevicePodJob(pod.id, jobId);
|
|
if (!job) return sendJson(response, 404, errorPayload("device_job_not_found", `Device job ${jobId} was not found`, 404));
|
|
if (job.ownerUserId !== actor.id && actor.role !== "admin") return sendJson(response, 403, errorPayload("device_job_owner_required", "Only the owner or admin can inspect the job", 403));
|
|
if (this.devicePodExecutorUrl && !terminalJobStatus(job.status)) {
|
|
job = await this.refreshDevicePodExecutorJob({ job, pod, output: parts[2] === "output" });
|
|
}
|
|
if (parts[2] === "output") return sendJson(response, 200, this.jobOutputPayload(job, pod));
|
|
return sendJson(response, 200, this.jobPayload(job, pod));
|
|
}
|
|
|
|
async cancelDevicePodJob(response, pod, route, actor) {
|
|
const jobId = decodeURIComponent(route.split("/")[1] ?? "");
|
|
const job = await this.store.getDevicePodJob(pod.id, jobId);
|
|
if (!job) return sendJson(response, 404, errorPayload("device_job_not_found", `Device job ${jobId} was not found`, 404));
|
|
if (job.ownerUserId !== actor.id && actor.role !== "admin") return sendJson(response, 403, errorPayload("device_job_owner_required", "Only the owner or admin can cancel the job", 403));
|
|
if (this.devicePodExecutorUrl && !terminalJobStatus(job.status)) {
|
|
const canceledByExecutor = await this.cancelDevicePodExecutorJob({ job, pod });
|
|
return sendJson(response, 200, this.jobPayload(canceledByExecutor, pod));
|
|
}
|
|
const canceled = await this.store.updateDevicePodJob(pod.id, jobId, {
|
|
status: terminalJobStatus(job.status) ? job.status : "canceled",
|
|
blocker: terminalJobStatus(job.status) ? job.blocker : { code: "device_job_canceled", summary: "Device job was canceled by user" },
|
|
completedAt: terminalJobStatus(job.status) ? job.completedAt : this.now(),
|
|
updatedAt: this.now()
|
|
});
|
|
return sendJson(response, 200, this.jobPayload(canceled, pod));
|
|
}
|
|
|
|
async refreshDevicePodExecutorJob({ job, pod, output = false }) {
|
|
const suffix = output ? "/output" : "";
|
|
const target = `${this.devicePodExecutorUrl}/v1/device-pods/${encodeURIComponent(pod.id)}/jobs/${encodeURIComponent(job.id)}${suffix}`;
|
|
try {
|
|
const response = await fetchJsonWithTimeout(this.fetchImpl, target, {
|
|
method: "GET",
|
|
headers: { accept: "application/json", ...devicePodInternalHeaders(CLOUD_API_SERVICE_ID, this.devicePodInternalToken) }
|
|
}, this.devicePodExecutorTimeoutMs);
|
|
return await this.updateDevicePodJobFromExecutorResponse({ job, pod, response });
|
|
} catch (error) {
|
|
return await this.blockDevicePodJobOnExecutorError({ job, pod, error });
|
|
}
|
|
}
|
|
|
|
async cancelDevicePodExecutorJob({ job, pod }) {
|
|
const target = `${this.devicePodExecutorUrl}/v1/device-pods/${encodeURIComponent(pod.id)}/jobs/${encodeURIComponent(job.id)}/cancel`;
|
|
try {
|
|
const response = await fetchJsonWithTimeout(this.fetchImpl, target, {
|
|
method: "POST",
|
|
headers: { accept: "application/json", ...devicePodInternalHeaders(CLOUD_API_SERVICE_ID, this.devicePodInternalToken) }
|
|
}, this.devicePodExecutorTimeoutMs);
|
|
return await this.updateDevicePodJobFromExecutorResponse({ job, pod, response });
|
|
} catch (error) {
|
|
return await this.blockDevicePodJobOnExecutorError({ job, pod, error });
|
|
}
|
|
}
|
|
|
|
async updateDevicePodJobFromExecutorResponse({ job, pod, response }) {
|
|
const status = normalizeDeviceJobStatus(response.body?.status ?? response.body?.job?.status, response.status);
|
|
const blocker = response.body?.blocker ?? (status === "completed" || status === "running" || status === "queued" ? null : devicePodExecutorBlocker(`device-pod executor returned HTTP ${response.status}`));
|
|
const updated = await this.store.updateDevicePodJob(pod.id, job.id, {
|
|
status,
|
|
output: executorOutputPayload(response.body, response.status),
|
|
blocker,
|
|
completedAt: terminalJobStatus(status) ? job.completedAt ?? this.now() : null,
|
|
updatedAt: this.now()
|
|
});
|
|
return updated ?? job;
|
|
}
|
|
|
|
async blockDevicePodJobOnExecutorError({ job, pod, error }) {
|
|
const blocker = devicePodExecutorBlocker(error?.message ?? "device-pod executor request failed");
|
|
const updated = await this.store.updateDevicePodJob(pod.id, job.id, {
|
|
status: "blocked",
|
|
output: { error: blocker.summary },
|
|
blocker,
|
|
completedAt: this.now(),
|
|
updatedAt: this.now()
|
|
});
|
|
return updated ?? job;
|
|
}
|
|
|
|
async blockDevicePodJobWithoutExecutor({ job, pod }) {
|
|
const blocker = devicePodExecutorBlocker("HWLAB_DEVICE_POD_URL is not configured; cloud-api will not bypass hwlab-device-pod executor");
|
|
const updated = await this.store.updateDevicePodJob(pod.id, job.id, {
|
|
status: "blocked",
|
|
output: { error: blocker.summary },
|
|
blocker,
|
|
completedAt: this.now(),
|
|
updatedAt: this.now()
|
|
});
|
|
return { job: updated ?? job, httpStatus: 409 };
|
|
}
|
|
|
|
jobPayload(job, pod, { lease = null } = {}) {
|
|
const blocker = normalizeBlocker(job.blocker);
|
|
return {
|
|
serviceId: CLOUD_API_SERVICE_ID,
|
|
contractVersion: DEVICE_JOB_CONTRACT_VERSION,
|
|
accepted: !["blocked", "failed"].includes(job.status),
|
|
status: job.status,
|
|
devicePodId: pod.id,
|
|
targetId: targetIdFromProfile(pod.profile),
|
|
profileHash: pod.profileHash,
|
|
traceId: job.traceId,
|
|
operationId: job.operationId,
|
|
job: publicJob(job),
|
|
lease: lease ? publicDeviceLease(lease) : null,
|
|
blocker,
|
|
freshness: freshness(job.updatedAt, blocker),
|
|
outputUrl: `/v1/device-pods/${encodeURIComponent(pod.id)}/jobs/${encodeURIComponent(job.id)}/output`,
|
|
cancelUrl: `/v1/device-pods/${encodeURIComponent(pod.id)}/jobs/${encodeURIComponent(job.id)}/cancel`
|
|
};
|
|
}
|
|
|
|
jobOutputPayload(job, pod, { maxBytes = DEVICE_JOB_OUTPUT_MAX_BYTES } = {}) {
|
|
const bounded = boundedOutput(job.output, maxBytes);
|
|
return {
|
|
...this.jobPayload(job, pod),
|
|
output: bounded.output,
|
|
text: bounded.text,
|
|
bytes: bounded.bytes,
|
|
truncation: bounded.truncation
|
|
};
|
|
}
|
|
}
|
|
|
|
class MemoryAccessStore {
|
|
constructor({ now = () => new Date().toISOString() } = {}) {
|
|
this.now = now;
|
|
this.users = new Map();
|
|
this.sessions = new Map();
|
|
this.devicePods = new Map();
|
|
this.grants = new Set();
|
|
this.leases = new Map();
|
|
this.jobs = new Map();
|
|
this.agentSessions = new Map();
|
|
}
|
|
|
|
async countUsers() { return this.users.size; }
|
|
async getUserById(id) { return this.users.get(id) ?? null; }
|
|
async findUserByUsername(username) { return [...this.users.values()].find((user) => user.username === username) ?? null; }
|
|
async syncBootstrapAdminPassword(input) {
|
|
const existing = await this.findUserByUsername(input.username) ?? await this.getUserById(input.id);
|
|
if (!existing || existing.role !== "admin") return null;
|
|
const user = { ...existing, passwordHash: input.passwordHash, updatedAt: input.now ?? this.now() };
|
|
this.users.set(user.id, user);
|
|
return user;
|
|
}
|
|
async createUser(input) {
|
|
const now = input.now ?? this.now();
|
|
const existing = await this.findUserByUsername(input.username);
|
|
const user = { id: input.id ?? existing?.id ?? `usr_${randomUUID()}`, username: input.username, displayName: input.displayName ?? existing?.displayName ?? "", role: input.role, status: input.status ?? "active", passwordHash: input.passwordHash ?? existing?.passwordHash ?? null, createdAt: existing?.createdAt ?? now, updatedAt: now };
|
|
this.users.set(user.id, user);
|
|
return user;
|
|
}
|
|
async createFirstAdmin(input) {
|
|
if (this.users.size > 0) return null;
|
|
return this.createUser({ ...input, role: "admin", status: "active" });
|
|
}
|
|
async createSession(input) {
|
|
const session = { id: `uss_${randomUUID()}`, userId: input.userId, tokenHash: input.tokenHash, createdAt: input.now, lastSeenAt: input.now, expiresAt: input.expiresAt, revokedAt: null };
|
|
this.sessions.set(session.id, session);
|
|
return session;
|
|
}
|
|
async findSessionByTokenHash(tokenHash, now) {
|
|
const session = [...this.sessions.values()].find((item) => item.tokenHash === tokenHash && !item.revokedAt && item.expiresAt > now) ?? null;
|
|
return session ? { ...session, user: this.users.get(session.userId) ?? null } : null;
|
|
}
|
|
async touchSession(id, now) { const session = this.sessions.get(id); if (session) session.lastSeenAt = now; }
|
|
async revokeSessionByTokenHash(tokenHash, now) { for (const session of this.sessions.values()) if (session.tokenHash === tokenHash) session.revokedAt = now; }
|
|
async upsertDevicePod(input) {
|
|
const existing = this.devicePods.get(input.id);
|
|
const profile = devicePodProfileForAuthority(input.id, normalizeObject(input.profile));
|
|
const now = input.now ?? this.now();
|
|
const pod = { id: input.id, name: input.name ?? existing?.name ?? input.id, status: input.status ?? existing?.status ?? "active", profile, profileHash: profileHash(profile), createdAt: existing?.createdAt ?? now, updatedAt: now };
|
|
this.devicePods.set(pod.id, pod);
|
|
return pod;
|
|
}
|
|
async grantDevicePod(input) { const grant = { devicePodId: input.devicePodId, userId: input.userId, createdByAdminId: input.createdByAdminId, createdAt: input.now ?? this.now() }; this.grants.add(grantKey(grant.devicePodId, grant.userId)); return grant; }
|
|
async revokeDevicePodGrant(input) { this.grants.delete(grantKey(input.devicePodId, input.userId)); }
|
|
async getDevicePod(id) { return this.devicePods.get(id) ?? null; }
|
|
async listVisibleDevicePods(actor) { return [...this.devicePods.values()].filter((pod) => pod.status === "active" && (actor.role === "admin" || this.grants.has(grantKey(pod.id, actor.id)))); }
|
|
async getVisibleDevicePod(actor, id) { return (await this.listVisibleDevicePods(actor)).find((pod) => pod.id === id) ?? null; }
|
|
async acquireDeviceLease(input) {
|
|
const now = input.now ?? this.now();
|
|
const active = await this.getActiveDeviceLease(input.devicePodId, now);
|
|
if (active && active.holderUserId !== input.holderUserId) return null;
|
|
const lease = { devicePodId: input.devicePodId, holderSessionId: input.holderSessionId, holderUserId: input.holderUserId, leaseTokenHash: input.leaseTokenHash, createdAt: active?.createdAt ?? now, expiresAt: input.expiresAt, releasedAt: null };
|
|
this.leases.set(input.devicePodId, lease);
|
|
return lease;
|
|
}
|
|
async getActiveDeviceLease(devicePodId, now = this.now()) {
|
|
const lease = this.leases.get(devicePodId) ?? null;
|
|
return lease && !lease.releasedAt && lease.expiresAt > now ? lease : null;
|
|
}
|
|
async getDeviceLeaseByToken(input) {
|
|
const lease = await this.getActiveDeviceLease(input.devicePodId, input.now ?? this.now());
|
|
return lease?.leaseTokenHash === input.tokenHash ? lease : null;
|
|
}
|
|
async releaseDeviceLease(input) {
|
|
const lease = await this.getActiveDeviceLease(input.devicePodId, input.now ?? this.now());
|
|
if (!lease) return null;
|
|
const tokenMatches = input.tokenHash && lease.leaseTokenHash === input.tokenHash;
|
|
if (!input.allowAdmin && !input.revokeGrant && (lease.holderUserId !== input.holderUserId || !tokenMatches)) return null;
|
|
if (!input.allowAdmin && input.revokeGrant && lease.holderUserId !== input.holderUserId) return null;
|
|
const released = { ...lease, releasedAt: input.now ?? this.now() };
|
|
this.leases.set(input.devicePodId, released);
|
|
return released;
|
|
}
|
|
async createDevicePodJob(input) { const job = normalizeJob(input); this.jobs.set(job.id, job); return job; }
|
|
async updateDevicePodJob(devicePodId, jobId, patch) { const job = this.jobs.get(jobId); if (!job || job.devicePodId !== devicePodId) return null; const next = { ...job, ...patch, updatedAt: patch.updatedAt ?? this.now() }; this.jobs.set(jobId, next); return next; }
|
|
async getDevicePodJob(devicePodId, jobId) { const job = this.jobs.get(jobId); return job?.devicePodId === devicePodId ? job : null; }
|
|
async listDevicePodJobs(devicePodId, limit) { return [...this.jobs.values()].filter((job) => job.devicePodId === devicePodId).sort((a, b) => b.updatedAt.localeCompare(a.updatedAt)).slice(0, limit); }
|
|
async recordAgentSessionOwner(input) {
|
|
const now = input.now ?? this.now();
|
|
const existing = this.agentSessions.get(input.sessionId) ?? null;
|
|
const session = normalizeAgentSessionOwnerRecord(input, existing, now);
|
|
this.agentSessions.set(session.id, session);
|
|
return session;
|
|
}
|
|
async getAgentSession(id) { return this.agentSessions.get(id) ?? null; }
|
|
}
|
|
|
|
class PostgresAccessStore extends MemoryAccessStore {
|
|
constructor({ query, now } = {}) {
|
|
super({ now });
|
|
this.query = query;
|
|
this.schemaReady = false;
|
|
}
|
|
|
|
async ensureSchema() {
|
|
if (this.schemaReady) return;
|
|
for (const sql of ACCESS_SCHEMA_STATEMENTS) await this.query(sql, []);
|
|
this.schemaReady = true;
|
|
}
|
|
|
|
async countUsers() { await this.ensureSchema(); const result = await this.query("SELECT COUNT(*)::int AS count FROM users", []); return Number(result.rows?.[0]?.count ?? 0); }
|
|
async getUserById(id) { await this.ensureSchema(); const result = await this.query("SELECT id, username, display_name, role, status, password_hash, created_at, updated_at FROM users WHERE id = $1 LIMIT 1", [id]); return pgUser(result.rows?.[0]); }
|
|
async findUserByUsername(username) { await this.ensureSchema(); const result = await this.query("SELECT id, username, display_name, role, status, password_hash, created_at, updated_at FROM users WHERE username = $1 LIMIT 1", [username]); return pgUser(result.rows?.[0]); }
|
|
async syncBootstrapAdminPassword(input) {
|
|
await this.ensureSchema();
|
|
const result = await this.query("UPDATE users SET password_hash = $3, updated_at = $4 WHERE role = 'admin' AND (username = $1 OR id = $2) RETURNING id, username, display_name, role, status, password_hash, created_at, updated_at", [input.username, input.id, input.passwordHash, input.now ?? this.now()]);
|
|
return pgUser(result.rows?.[0]);
|
|
}
|
|
async createUser(input) {
|
|
await this.ensureSchema();
|
|
const now = input.now ?? this.now();
|
|
const user = { id: input.id ?? `usr_${randomUUID()}`, username: input.username, displayName: input.displayName ?? "", role: input.role, status: input.status ?? "active", passwordHash: input.passwordHash ?? null, createdAt: now, updatedAt: now };
|
|
const result = await this.query("INSERT INTO users (id, username, display_name, role, status, password_hash, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8) ON CONFLICT (username) DO UPDATE SET display_name = EXCLUDED.display_name, role = EXCLUDED.role, status = EXCLUDED.status, password_hash = EXCLUDED.password_hash, updated_at = EXCLUDED.updated_at RETURNING id, username, display_name, role, status, password_hash, created_at, updated_at", [user.id, user.username, user.displayName, user.role, user.status, user.passwordHash, user.createdAt, user.updatedAt]);
|
|
return pgUser(result.rows?.[0]);
|
|
}
|
|
async createFirstAdmin(input) {
|
|
await this.ensureSchema();
|
|
const now = input.now ?? this.now();
|
|
const user = { id: input.id, username: input.username, displayName: input.displayName ?? "", role: "admin", status: "active", passwordHash: input.passwordHash ?? null, createdAt: now, updatedAt: now };
|
|
const result = await this.query("INSERT INTO users (id, username, display_name, role, status, password_hash, created_at, updated_at) SELECT $1,$2,$3,$4,$5,$6,$7,$8 WHERE NOT EXISTS (SELECT 1 FROM users) ON CONFLICT DO NOTHING RETURNING id, username, display_name, role, status, password_hash, created_at, updated_at", [user.id, user.username, user.displayName, user.role, user.status, user.passwordHash, user.createdAt, user.updatedAt]);
|
|
return pgUser(result.rows?.[0]);
|
|
}
|
|
async createSession(input) { await this.ensureSchema(); const session = await super.createSession(input); await this.query("INSERT INTO user_sessions (id, user_id, session_token_hash, created_at, last_seen_at, expires_at, revoked_at) VALUES ($1,$2,$3,$4,$5,$6,$7)", [session.id, session.userId, session.tokenHash, session.createdAt, session.lastSeenAt, session.expiresAt, session.revokedAt]); return session; }
|
|
async findSessionByTokenHash(tokenHash, now) { await this.ensureSchema(); const result = await this.query("SELECT s.id, s.user_id, s.session_token_hash, s.created_at, s.last_seen_at, s.expires_at, s.revoked_at, u.username, u.display_name, u.role, u.status, u.password_hash, u.created_at AS user_created_at, u.updated_at AS user_updated_at FROM user_sessions s JOIN users u ON u.id = s.user_id WHERE s.session_token_hash = $1 AND s.revoked_at IS NULL AND s.expires_at > $2 LIMIT 1", [tokenHash, now]); const row = result.rows?.[0]; return row ? pgSession(row) : null; }
|
|
async touchSession(id, now) { await this.ensureSchema(); await this.query("UPDATE user_sessions SET last_seen_at = $2 WHERE id = $1", [id, now]); }
|
|
async revokeSessionByTokenHash(tokenHash, now) { await this.ensureSchema(); await this.query("UPDATE user_sessions SET revoked_at = $2 WHERE session_token_hash = $1", [tokenHash, now]); }
|
|
async upsertDevicePod(input) { await this.ensureSchema(); const pod = await super.upsertDevicePod(input); const result = await this.query("INSERT INTO device_pods (id, name, status, profile_json, profile_hash, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (id) DO UPDATE SET name = EXCLUDED.name, status = EXCLUDED.status, profile_json = EXCLUDED.profile_json, profile_hash = EXCLUDED.profile_hash, updated_at = EXCLUDED.updated_at RETURNING *", [pod.id, pod.name, pod.status, stableJson(pod.profile), pod.profileHash, pod.createdAt, pod.updatedAt]); return pgDevicePod(result.rows?.[0]); }
|
|
async grantDevicePod(input) { await this.ensureSchema(); const grant = await super.grantDevicePod(input); await this.query("INSERT INTO device_pod_grants (device_pod_id, user_id, created_by_admin_id, created_at) VALUES ($1,$2,$3,$4) ON CONFLICT (device_pod_id, user_id) DO UPDATE SET created_by_admin_id = EXCLUDED.created_by_admin_id, created_at = EXCLUDED.created_at", [grant.devicePodId, grant.userId, grant.createdByAdminId, grant.createdAt]); return grant; }
|
|
async revokeDevicePodGrant(input) { await this.ensureSchema(); await this.query("DELETE FROM device_pod_grants WHERE device_pod_id = $1 AND user_id = $2", [input.devicePodId, input.userId]); }
|
|
async getDevicePod(id) { await this.ensureSchema(); const result = await this.query("SELECT * FROM device_pods WHERE id = $1 LIMIT 1", [id]); return pgDevicePod(result.rows?.[0]); }
|
|
async listVisibleDevicePods(actor) { await this.ensureSchema(); const sql = actor.role === "admin" ? "SELECT * FROM device_pods WHERE status = 'active' ORDER BY id" : "SELECT p.* FROM device_pods p JOIN device_pod_grants g ON g.device_pod_id = p.id WHERE p.status = 'active' AND g.user_id = $1 ORDER BY p.id"; const result = await this.query(sql, actor.role === "admin" ? [] : [actor.id]); return result.rows.map(pgDevicePod); }
|
|
async acquireDeviceLease(input) {
|
|
await this.ensureSchema();
|
|
const active = await this.getActiveDeviceLease(input.devicePodId, input.now);
|
|
if (active && active.holderUserId !== input.holderUserId) return null;
|
|
const result = await this.query("INSERT INTO device_leases (device_pod_id, holder_session_id, holder_user_id, lease_token_hash, created_at, expires_at, released_at) VALUES ($1,$2,$3,$4,$5,$6,NULL) ON CONFLICT (device_pod_id) DO UPDATE SET holder_session_id = EXCLUDED.holder_session_id, holder_user_id = EXCLUDED.holder_user_id, lease_token_hash = EXCLUDED.lease_token_hash, expires_at = EXCLUDED.expires_at, released_at = NULL WHERE device_leases.released_at IS NOT NULL OR device_leases.expires_at <= $5 OR device_leases.holder_user_id = $3 RETURNING *", [input.devicePodId, input.holderSessionId, input.holderUserId, input.leaseTokenHash, input.now, input.expiresAt]);
|
|
return pgLease(result.rows?.[0]);
|
|
}
|
|
async getActiveDeviceLease(devicePodId, now = this.now()) { await this.ensureSchema(); const result = await this.query("SELECT * FROM device_leases WHERE device_pod_id = $1 AND released_at IS NULL AND expires_at > $2 LIMIT 1", [devicePodId, now]); return pgLease(result.rows?.[0]); }
|
|
async getDeviceLeaseByToken(input) { await this.ensureSchema(); const result = await this.query("SELECT * FROM device_leases WHERE device_pod_id = $1 AND lease_token_hash = $2 AND released_at IS NULL AND expires_at > $3 LIMIT 1", [input.devicePodId, input.tokenHash, input.now]); return pgLease(result.rows?.[0]); }
|
|
async releaseDeviceLease(input) {
|
|
await this.ensureSchema();
|
|
const params = input.allowAdmin
|
|
? [input.devicePodId, input.now]
|
|
: input.revokeGrant
|
|
? [input.devicePodId, input.now, input.holderUserId]
|
|
: [input.devicePodId, input.now, input.holderUserId, input.tokenHash];
|
|
const sql = input.allowAdmin
|
|
? "UPDATE device_leases SET released_at = $2 WHERE device_pod_id = $1 AND released_at IS NULL RETURNING *"
|
|
: input.revokeGrant
|
|
? "UPDATE device_leases SET released_at = $2 WHERE device_pod_id = $1 AND holder_user_id = $3 AND released_at IS NULL RETURNING *"
|
|
: "UPDATE device_leases SET released_at = $2 WHERE device_pod_id = $1 AND holder_user_id = $3 AND lease_token_hash = $4 AND released_at IS NULL RETURNING *";
|
|
const result = await this.query(sql, params);
|
|
return pgLease(result.rows?.[0]);
|
|
}
|
|
async createDevicePodJob(input) { await this.ensureSchema(); const job = normalizeJob(input); await this.query("INSERT INTO device_pod_jobs (id, device_pod_id, owner_user_id, status, intent, args_json, reason, trace_id, operation_id, output_json, blocker_json, created_at, updated_at, completed_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14)", jobParams(job)); return job; }
|
|
async updateDevicePodJob(devicePodId, jobId, patch) { await this.ensureSchema(); const current = await this.getDevicePodJob(devicePodId, jobId); if (!current) return null; const next = { ...current, ...patch, updatedAt: patch.updatedAt ?? this.now() }; await this.query("UPDATE device_pod_jobs SET status=$3, output_json=$4, blocker_json=$5, updated_at=$6, completed_at=$7 WHERE device_pod_id=$1 AND id=$2", [devicePodId, jobId, next.status, stableJson(next.output ?? {}), stableJson(next.blocker ?? {}), next.updatedAt, next.completedAt]); return next; }
|
|
async getDevicePodJob(devicePodId, jobId) { await this.ensureSchema(); const result = await this.query("SELECT * FROM device_pod_jobs WHERE device_pod_id = $1 AND id = $2 LIMIT 1", [devicePodId, jobId]); return pgJob(result.rows?.[0]); }
|
|
async listDevicePodJobs(devicePodId, limit) { await this.ensureSchema(); const result = await this.query("SELECT * FROM device_pod_jobs WHERE device_pod_id = $1 ORDER BY updated_at DESC LIMIT $2", [devicePodId, limit]); return result.rows.map(pgJob); }
|
|
async recordAgentSessionOwner(input) {
|
|
await this.ensureSchema();
|
|
const now = input.now ?? this.now();
|
|
const session = normalizeAgentSessionOwnerRecord(input, null, now);
|
|
const result = await this.query("INSERT INTO agent_sessions (id, project_id, agent_id, status, started_at, ended_at, owner_user_id, conversation_id, thread_id, last_trace_id, session_json, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12) ON CONFLICT (id) DO UPDATE SET project_id = COALESCE(EXCLUDED.project_id, agent_sessions.project_id), agent_id = EXCLUDED.agent_id, status = EXCLUDED.status, ended_at = EXCLUDED.ended_at, owner_user_id = EXCLUDED.owner_user_id, conversation_id = COALESCE(EXCLUDED.conversation_id, agent_sessions.conversation_id), thread_id = COALESCE(EXCLUDED.thread_id, agent_sessions.thread_id), last_trace_id = COALESCE(EXCLUDED.last_trace_id, agent_sessions.last_trace_id), session_json = EXCLUDED.session_json, updated_at = EXCLUDED.updated_at RETURNING id, project_id, agent_id, status, started_at, ended_at, owner_user_id, conversation_id, thread_id, last_trace_id, session_json, updated_at", [session.id, session.projectId, session.agentId, session.status, session.startedAt, session.endedAt, session.ownerUserId, session.conversationId, session.threadId, session.lastTraceId, stableJson(session.session ?? {}), session.updatedAt]);
|
|
return pgAgentSession(result.rows?.[0]);
|
|
}
|
|
}
|
|
|
|
function parseDevicePodPath(pathname) {
|
|
const prefix = "/v1/device-pods/";
|
|
if (!pathname.startsWith(prefix)) return null;
|
|
const [devicePodId, ...rest] = pathname.slice(prefix.length).split("/").map((part) => decodeURIComponent(part)).filter(Boolean);
|
|
return devicePodId ? { devicePodId, route: rest.join("/") } : null;
|
|
}
|
|
|
|
async function jsonBody(request) {
|
|
const text = await readBody(request);
|
|
try { return text ? JSON.parse(text) : {}; } catch (error) { throw Object.assign(new Error("Invalid JSON body"), { statusCode: 400, code: "parse_error", reason: error.message }); }
|
|
}
|
|
|
|
function requiredText(value, field) { const text = textOr(value, ""); if (!text) throw Object.assign(new Error(`${field} is required`), { statusCode: 400, code: "invalid_params" }); return text; }
|
|
function textOr(value, fallback) { const text = String(value ?? "").trim(); return text || fallback; }
|
|
function normalizeObject(value) { return value && typeof value === "object" && !Array.isArray(value) ? { ...value } : {}; }
|
|
function assertDevicePodProfileSafe(profile, field = "profile") {
|
|
const findings = [];
|
|
collectDevicePodProfileSecretFindings(profile, field, findings);
|
|
if (findings.length > 0) {
|
|
const paths = findings.slice(0, 3).join(", ");
|
|
throw Object.assign(new Error(`${field} contains forbidden secret material at ${paths}`), { statusCode: 400, code: "device_pod_profile_secret_forbidden" });
|
|
}
|
|
return profile;
|
|
}
|
|
function devicePodProfileForAuthority(devicePodId, profile) { return { ...assertDevicePodProfileSafe(profile, "profile"), devicePodId }; }
|
|
function collectDevicePodProfileSecretFindings(value, path, findings) {
|
|
if (findings.length >= 5) return;
|
|
if (Array.isArray(value)) {
|
|
value.forEach((item, index) => collectDevicePodProfileSecretFindings(item, `${path}[${index}]`, findings));
|
|
return;
|
|
}
|
|
if (value && typeof value === "object") {
|
|
for (const [key, child] of Object.entries(value)) {
|
|
const childPath = `${path}.${key}`;
|
|
if (DEVICE_POD_PROFILE_SECRET_KEY_PATTERN.test(key)) findings.push(childPath);
|
|
collectDevicePodProfileSecretFindings(child, childPath, findings);
|
|
if (findings.length >= 5) return;
|
|
}
|
|
return;
|
|
}
|
|
if (typeof value === "string" && DEVICE_POD_PROFILE_SECRET_VALUE_PATTERN.test(value)) findings.push(path);
|
|
}
|
|
function firstAdminDevicePodSeeds(body = {}) {
|
|
const hasSingle = Object.hasOwn(body, "devicePod");
|
|
const hasMany = Object.hasOwn(body, "devicePods");
|
|
if (hasSingle && hasMany) throw Object.assign(new Error("Use devicePod or devicePods, not both"), { statusCode: 400, code: "invalid_params" });
|
|
if (!hasSingle && !hasMany) return [];
|
|
const values = hasMany ? body.devicePods : [body.devicePod];
|
|
if (!Array.isArray(values)) throw Object.assign(new Error("devicePods must be an array"), { statusCode: 400, code: "invalid_params" });
|
|
return values.map((item, index) => firstAdminDevicePodSeed(item, hasSingle ? "devicePod" : `devicePods[${index}]`));
|
|
}
|
|
function firstAdminDevicePodSeed(value, field) {
|
|
if (!value || typeof value !== "object" || Array.isArray(value)) throw Object.assign(new Error(`${field} must be an object`), { statusCode: 400, code: "invalid_params" });
|
|
const id = requiredText(value.devicePodId ?? value.id, `${field}.devicePodId`);
|
|
const profile = normalizeObject(value.profile ?? value.profileJson);
|
|
if (Object.keys(profile).length === 0) throw Object.assign(new Error(`${field}.profile is required`), { statusCode: 400, code: "invalid_params" });
|
|
assertDevicePodProfileSafe(profile, `${field}.profile`);
|
|
return {
|
|
id,
|
|
name: textOr(value.name, id),
|
|
status: value.status === "disabled" ? "disabled" : "active",
|
|
profile
|
|
};
|
|
}
|
|
function normalizeBaseUrl(value) { const text = textOr(value, "").replace(/\/+$/u, ""); return /^https?:\/\//u.test(text) ? text : ""; }
|
|
function boundedLeaseTtlSeconds(value) { const parsed = Number.parseInt(String(value ?? ""), 10); return Math.min(Math.max(Number.isInteger(parsed) && parsed > 0 ? parsed : DEFAULT_DEVICE_LEASE_TTL_SECONDS, 1), MAX_DEVICE_LEASE_TTL_SECONDS); }
|
|
function boundedOutputMaxBytes(searchParams) { return Math.min(Math.max(Number.parseInt(searchParams.get("maxBytes") ?? "12000", 10) || 12000, 1), DEVICE_JOB_OUTPUT_MAX_BYTES); }
|
|
function boundedDurationMs(value, fallback = 1000) { const parsed = Number.parseInt(String(value ?? ""), 10); return Math.min(Math.max(Number.isInteger(parsed) && parsed > 0 ? parsed : fallback, 1), 60000); }
|
|
function uartTailArgs(searchParams, uartId) { return { uartId, durationMs: boundedDurationMs(searchParams.get("durationMs") ?? searchParams.get("duration-ms")), maxBytes: boundedOutputMaxBytes(searchParams) }; }
|
|
function safeAgentSessionId(value) { const text = textOr(value, ""); return /^ses_[A-Za-z0-9_.:-]+$/u.test(text) ? text : ""; }
|
|
function leaseTokenFromRequest(request, body) { return textOr(body?.leaseToken, "") || textOr(getHeader(request, "x-hwlab-device-lease-token"), ""); }
|
|
function devicePodInternalHeaders(serviceId, token) { return token ? { "x-hwlab-internal-service": serviceId, "x-hwlab-internal-token": token } : { "x-hwlab-internal-service": serviceId }; }
|
|
function hashPassword(password) { const salt = randomBytes(16).toString("hex"); return `sha256:${salt}:${sha256(`${salt}:${password}`)}`; }
|
|
function verifyPassword(stored, password) { const [, salt, digest] = String(stored ?? "").split(":"); return Boolean(salt && digest && sha256(`${salt}:${password}`) === digest); }
|
|
function sha256(value) { return createHash("sha256").update(String(value)).digest("hex"); }
|
|
function stableJson(value) { return JSON.stringify(sortJson(value)); }
|
|
function sortJson(value) { if (Array.isArray(value)) return value.map(sortJson); if (value && typeof value === "object") return Object.fromEntries(Object.keys(value).sort().map((key) => [key, sortJson(value[key])])); return value; }
|
|
function profileHash(profile) { return `sha256:${sha256(stableJson(profile))}`; }
|
|
function grantKey(devicePodId, userId) { return `${devicePodId}\u0000${userId}`; }
|
|
function authoritySource() { return { kind: "CLOUD_API_PROFILE_AUTHORITY", serviceId: CLOUD_API_SERVICE_ID, fake: false, devLiveEvidence: false, note: "Profile, grant, lease, and job authority are owned by hwlab-cloud-api; this is not fake data." }; }
|
|
function accessRoutes() { return { session: "/auth/session", login: "/auth/login", logout: "/auth/logout", v1Session: "/v1/auth/session", currentUser: "/v1/users/me", accessStatus: "/v1/access/status", setupStatus: "/v1/setup/status", firstAdminSetup: "/v1/setup/first-admin", devicePods: "/v1/device-pods" }; }
|
|
function deviceJobBlocked(code, summary, retryable) { return { code, layer: "device-pod", retryable, summary, userMessage: summary }; }
|
|
function devicePodExecutorBlocker(summary) { return { code: "device_pod_executor_unavailable", layer: "device-pod", retryable: true, summary, userMessage: "Device Pod 已授权,但内部执行服务当前不可用。" }; }
|
|
function gatewayDispatchBlocker(summary) { return { code: "gateway_dispatch_unavailable", layer: "device-pod", retryable: true, summary, userMessage: "Device Pod 已授权,但当前没有可用 gateway/device-host-cli 执行通道。" }; }
|
|
function normalizeDeviceJobStatus(status, httpStatus) { const text = textOr(status, ""); return ["queued", "running", "completed", "failed", "blocked", "canceled"].includes(text) ? text : httpStatus >= 400 ? "blocked" : "running"; }
|
|
function executorOutputPayload(body, httpStatus) { return { executor: body ?? {}, output: normalizeObject(body?.output), text: typeof body?.text === "string" ? body.text : typeof body?.output?.text === "string" ? body.output.text : "", httpStatus }; }
|
|
function boundedOutput(output, maxBytes = DEVICE_JOB_OUTPUT_MAX_BYTES) {
|
|
const source = normalizeObject(output);
|
|
const text = typeof source.text === "string" ? source.text : stableJson(source);
|
|
const buffer = Buffer.from(text, "utf8");
|
|
const clipped = buffer.length > maxBytes;
|
|
const boundedText = clipped ? buffer.subarray(0, maxBytes).toString("utf8") : text;
|
|
const bounded = { ...source, text: boundedText };
|
|
if (clipped) {
|
|
delete bounded.executor;
|
|
delete bounded.output;
|
|
bounded.omitted = { reason: "device_job_output_truncated", originalBytes: buffer.length };
|
|
}
|
|
return { output: bounded, text: boundedText, bytes: Math.min(buffer.length, maxBytes), truncation: { maxBytes, truncated: clipped, originalBytes: buffer.length } };
|
|
}
|
|
function freshness(observedAt, blocker) { return { observedAt, ageMs: 0, stale: Boolean(blocker), source: blocker ? "blocked" : "cloud-api" }; }
|
|
function publicActor(user) { return user ? { id: user.id, username: user.username, displayName: user.displayName, role: user.role, status: user.status } : null; }
|
|
function redactedUser(user) { return publicActor(user); }
|
|
function publicSession(session) { return { id: session.id, userId: session.userId, createdAt: session.createdAt, lastSeenAt: session.lastSeenAt, expiresAt: session.expiresAt, revoked: Boolean(session.revokedAt) }; }
|
|
function publicDeviceLease(lease) { return lease ? { devicePodId: lease.devicePodId, holderSessionId: lease.holderSessionId, holderUserId: lease.holderUserId, createdAt: lease.createdAt, expiresAt: lease.expiresAt, releasedAt: lease.releasedAt ?? null, active: !lease.releasedAt } : null; }
|
|
function publicDevicePod(pod, { includeAdmin = false } = {}) { return { devicePodId: pod.id, name: pod.name, status: pod.status, targetId: targetIdFromProfile(pod.profile), profileHash: pod.profileHash, profile: redactedProfile(pod.profile), createdAt: pod.createdAt, updatedAt: pod.updatedAt, blocker: null, ...(includeAdmin ? { admin: { profileStored: true, routeStored: Boolean(pod.profile.route) } } : {}) }; }
|
|
function redactedProfile(profile = {}) { return { schemaVersion: profile.schemaVersion ?? null, target: { id: profile.target?.id ?? null }, projectWorkspace: { projectPath: profile.projectWorkspace?.projectPath ?? null, targetName: profile.projectWorkspace?.targetName ?? null, hexPath: profile.projectWorkspace?.hexPath ?? null }, debugInterface: { type: profile.debugInterface?.type ?? null }, ioInterface: { uartCount: Array.isArray(profile.ioInterface?.uart) ? profile.ioInterface.uart.length : 0 }, route: { configured: Boolean(profile.route?.gatewaySessionId), gatewaySessionId: "redacted", resourceId: profile.route?.resourceId ? "redacted" : null, capabilityId: profile.route?.capabilityId ? "redacted" : null } }; }
|
|
function targetIdFromProfile(profile = {}) { return profile.target?.id ?? profile.targetId ?? null; }
|
|
function terminalJobStatus(status) { return ["completed", "failed", "blocked", "canceled"].includes(status); }
|
|
function publicJob(job) { return { id: job.id, devicePodId: job.devicePodId, ownerUserId: job.ownerUserId, status: job.status, intent: job.intent, reason: job.reason, traceId: job.traceId, operationId: job.operationId, createdAt: job.createdAt, updatedAt: job.updatedAt, completedAt: job.completedAt }; }
|
|
function jobRefs(job) { return { jobId: job.id, traceId: job.traceId, operationId: job.operationId }; }
|
|
function formatEventLine(event) { return [event.ts?.slice(11, 19) ?? "00:00:00", event.scope?.toUpperCase() ?? "JOB", event.status, event.intent, event.summary, event.refs?.traceId ? `trace=${event.refs.traceId}` : null, event.blocker?.code ? `blocker=${event.blocker.code}` : null].filter(Boolean).join(" "); }
|
|
function normalizeJob(input) { return { id: input.id, devicePodId: input.devicePodId, ownerUserId: input.ownerUserId, status: input.status, intent: input.intent, args: normalizeObject(input.args), reason: input.reason ?? "", traceId: input.traceId, operationId: input.operationId, output: normalizeObject(input.output), blocker: input.blocker ?? null, createdAt: input.now, updatedAt: input.now, completedAt: input.completedAt ?? null }; }
|
|
function normalizeAgentSessionOwnerRecord(input, existing, now) { return { id: input.sessionId, projectId: input.projectId ?? existing?.projectId ?? "prj_v02_code_agent", agentId: input.agentId ?? existing?.agentId ?? "hwlab-code-agent", status: input.status ?? existing?.status ?? "active", startedAt: existing?.startedAt ?? input.startedAt ?? now, endedAt: input.endedAt ?? existing?.endedAt ?? null, ownerUserId: input.ownerUserId, conversationId: input.conversationId ?? existing?.conversationId ?? null, threadId: input.threadId ?? existing?.threadId ?? null, lastTraceId: input.traceId ?? input.lastTraceId ?? existing?.lastTraceId ?? null, session: normalizeObject(input.session ?? existing?.session), updatedAt: now }; }
|
|
function agentSessionIdForRecord(sessionId, traceId) { const sessionText = textOr(sessionId, ""); if (/^ses_[A-Za-z0-9_.:-]+$/u.test(sessionText)) return sessionText; const traceText = textOr(traceId, ""); return /^trc_[A-Za-z0-9_.:-]+$/u.test(traceText) ? `ses_pending_${traceText.slice(4)}` : null; }
|
|
function sessionTokenFromRequest(request) { const auth = getHeader(request, "authorization"); if (/^Bearer\s+/iu.test(String(auth ?? ""))) return String(auth).replace(/^Bearer\s+/iu, "").trim(); const header = textOr(getHeader(request, "x-hwlab-session-token"), ""); if (header) return header; const cookie = parseCookie(getHeader(request, "cookie")); return cookie[SESSION_COOKIE] ?? ""; }
|
|
function parseCookie(value) { return Object.fromEntries(String(value ?? "").split(";").map((part) => part.trim()).filter(Boolean).map((part) => { const index = part.indexOf("="); return index > 0 ? [part.slice(0, index), decodeURIComponent(part.slice(index + 1))] : [part, ""]; })); }
|
|
function setSessionCookie(response, token, maxAge) { response.setHeader("set-cookie", `${SESSION_COOKIE}=${encodeURIComponent(token)}; Path=/; HttpOnly; SameSite=Lax; Max-Age=${maxAge}`); }
|
|
function clearSessionCookie(response) { response.setHeader("set-cookie", `${SESSION_COOKIE}=; Path=/; HttpOnly; SameSite=Lax; Max-Age=0`); }
|
|
function errorPayload(code, message, status) { return { ok: false, status, error: { code, message } }; }
|
|
function sendAccessError(response, error) { const status = error?.statusCode ?? 500; sendJson(response, status, errorPayload(error?.code ?? "access_control_error", error?.message ?? "Access control request failed", status)); }
|
|
function pgUser(row) { return row ? { id: row.id, username: row.username, displayName: row.display_name, role: row.role, status: row.status, passwordHash: row.password_hash, createdAt: row.created_at, updatedAt: row.updated_at } : null; }
|
|
function pgSession(row) { return { id: row.id, userId: row.user_id, tokenHash: row.session_token_hash, createdAt: row.created_at, lastSeenAt: row.last_seen_at, expiresAt: row.expires_at, revokedAt: row.revoked_at, user: { id: row.user_id, username: row.username, displayName: row.display_name, role: row.role, status: row.status, passwordHash: row.password_hash, createdAt: row.user_created_at, updatedAt: row.user_updated_at } }; }
|
|
function pgDevicePod(row) { return { id: row.id, name: row.name, status: row.status, profile: parseJson(row.profile_json, {}), profileHash: row.profile_hash, createdAt: row.created_at, updatedAt: row.updated_at }; }
|
|
function pgLease(row) { return row ? { devicePodId: row.device_pod_id, holderSessionId: row.holder_session_id, holderUserId: row.holder_user_id, leaseTokenHash: row.lease_token_hash, createdAt: row.created_at, expiresAt: row.expires_at, releasedAt: row.released_at } : null; }
|
|
function pgJob(row) { return row ? { id: row.id, devicePodId: row.device_pod_id, ownerUserId: row.owner_user_id, status: row.status, intent: row.intent, args: parseJson(row.args_json, {}), reason: row.reason, traceId: row.trace_id, operationId: row.operation_id, output: parseJson(row.output_json, {}), blocker: normalizeBlocker(parseJson(row.blocker_json, null)), createdAt: row.created_at, updatedAt: row.updated_at, completedAt: row.completed_at } : null; }
|
|
function pgAgentSession(row) { return row ? { id: row.id, projectId: row.project_id, agentId: row.agent_id, status: row.status, startedAt: row.started_at, endedAt: row.ended_at, ownerUserId: row.owner_user_id, conversationId: row.conversation_id, threadId: row.thread_id, lastTraceId: row.last_trace_id, session: parseJson(row.session_json, {}), updatedAt: row.updated_at } : null; }
|
|
function jobParams(job) { return [job.id, job.devicePodId, job.ownerUserId, job.status, job.intent, stableJson(job.args), job.reason, job.traceId, job.operationId, stableJson(job.output), stableJson(job.blocker ?? {}), job.createdAt, job.updatedAt, job.completedAt]; }
|
|
function parseJson(value, fallback) { if (!value) return fallback; if (typeof value === "object") return value; try { return JSON.parse(String(value)); } catch { return fallback; } }
|
|
function normalizeBlocker(value) { return value && typeof value === "object" && !Array.isArray(value) && Object.keys(value).length > 0 ? value : null; }
|
|
|
|
async function fetchJsonWithTimeout(fetchImpl, url, options, timeoutMs) {
|
|
const controller = new AbortController();
|
|
const timeout = setTimeout(() => controller.abort(), timeoutMs);
|
|
try {
|
|
const response = await fetchImpl(url, { ...options, signal: controller.signal });
|
|
const text = await response.text();
|
|
return { status: response.status, body: parseJson(text, {}) };
|
|
} finally {
|
|
clearTimeout(timeout);
|
|
}
|
|
}
|