初始提交:DSH 多租户平台(dshs)
This commit is contained in:
commit
43976fea6a
167 files changed
+24456
No files matched your search
@@ -0,0 +1,112 @@
|
||||
/**
|
||||
* Crash-restart policy for a resident main DSH (档案 20 · 方案 A).
|
||||
*
|
||||
* Pure decision logic, deliberately separated from `LocalSpawner` so it can be
|
||||
* unit-tested without spawning anything: given the recent restart timestamps
|
||||
* and the current consecutive-attempt count, decide either "restart after
|
||||
* `delayMs` (exponential backoff capped at `maxDelayMs`)" or "circuit open"
|
||||
* (too many restarts inside the window — stop auto-restarting so a
|
||||
* persistently failing instance cannot spin spawns forever).
|
||||
*
|
||||
* @module dshs/supervisor/crash-policy
|
||||
*/
|
||||
|
||||
/** Tunables (see {@link ServerConfig} — all env-overridable). */
|
||||
export interface CrashPolicyConfig {
|
||||
/** Delay for the first restart (consecutive attempt 0). */
|
||||
baseDelayMs: number
|
||||
/** Upper bound for the exponential backoff. */
|
||||
maxDelayMs: number
|
||||
/** Rolling window over which restarts are counted. */
|
||||
windowMs: number
|
||||
/** Restarts allowed inside the window before the circuit opens. */
|
||||
maxRestartsInWindow: number
|
||||
/** A main that stays up this long is considered recovered: streak resets. */
|
||||
stableResetMs: number
|
||||
}
|
||||
|
||||
/** Outcome of one crash decision. */
|
||||
export type CrashDecision =
|
||||
| { action: 'restart'; delayMs: number; attempt: number; windowRestarts: number }
|
||||
| { action: 'circuit-open'; windowRestarts: number }
|
||||
|
||||
/** Drop timestamps that fell out of the rolling window. */
|
||||
export function pruneHistory(history: readonly number[], now: number, windowMs: number): number[] {
|
||||
return history.filter((at) => now - at < windowMs)
|
||||
}
|
||||
|
||||
/** Exponential backoff for a given consecutive attempt (0-based), capped. */
|
||||
export function backoffDelayMs(consecutive: number, cfg: Pick<CrashPolicyConfig, 'baseDelayMs' | 'maxDelayMs'>): number {
|
||||
const safe = Math.max(0, Math.min(consecutive, 20)) // guard against overflow
|
||||
return Math.min(cfg.baseDelayMs * 2 ** safe, cfg.maxDelayMs)
|
||||
}
|
||||
|
||||
/**
|
||||
* Decide what to do after a main crashed.
|
||||
* @param history - restart timestamps (ms epoch) already recorded for this user.
|
||||
* @param consecutive - consecutive attempts since the last stable run.
|
||||
* @param now - current time (ms epoch).
|
||||
* @param cfg - policy tunables.
|
||||
*/
|
||||
export function decideCrashAction(
|
||||
history: readonly number[],
|
||||
consecutive: number,
|
||||
now: number,
|
||||
cfg: CrashPolicyConfig,
|
||||
): CrashDecision {
|
||||
const inWindow = pruneHistory(history, now, cfg.windowMs)
|
||||
if (inWindow.length >= cfg.maxRestartsInWindow) {
|
||||
return { action: 'circuit-open', windowRestarts: inWindow.length }
|
||||
}
|
||||
return {
|
||||
action: 'restart',
|
||||
delayMs: backoffDelayMs(consecutive, cfg),
|
||||
attempt: consecutive + 1,
|
||||
windowRestarts: inWindow.length + 1,
|
||||
}
|
||||
}
|
||||
|
||||
/* ─────────────────────────────────────────────────────────────────────────────
|
||||
* 档案 78:熔断冷却(防「崩溃循环可无限重来」)
|
||||
*
|
||||
* 原始缺陷:`circuit-open` 分支里 `mains.delete` + `resetCrashState()` 会把窗口
|
||||
* 历史一并清空,于是**下一次 enter/launch 又是满额预算** —— 只要有人(用户 F5、
|
||||
* 注入脚本自愈、脚本直铺)不断重试,崩溃循环就能无限重复,且只留一行 stderr。
|
||||
*
|
||||
* 修法:每次熔断记一个**跨轮存活**的冷却窗(`opens` 递增 → 冷却指数加长),
|
||||
* 冷却期内拒绝隐式/自动启动;冷却过后只给**一次**干净预算。纯函数,便于单测。
|
||||
* ──────────────────────────────────────────────────────────────────────────── */
|
||||
|
||||
/** 熔断状态:`opens` = 该用户累计熔断次数(跨轮不清零)。 */
|
||||
export interface BreakerState {
|
||||
openedAt: number
|
||||
opens: number
|
||||
}
|
||||
|
||||
/** 冷却策略(见 ServerConfig,均可 env 覆盖)。 */
|
||||
export interface BreakerPolicy {
|
||||
baseCooldownMs: number
|
||||
maxCooldownMs: number
|
||||
}
|
||||
|
||||
/** 第 `opens` 次熔断的冷却时长:`base × 2^(opens-1)`,上限 `maxCooldownMs`。 */
|
||||
export function breakerCooldownMs(opens: number, cfg: BreakerPolicy): number {
|
||||
const safe = Math.max(1, Math.min(opens, 20))
|
||||
return Math.min(cfg.baseCooldownMs * 2 ** (safe - 1), cfg.maxCooldownMs)
|
||||
}
|
||||
|
||||
/** 冷却是否仍在生效(`now < openedAt + cooldown`)。 */
|
||||
export function breakerActive(b: BreakerState | undefined, now: number, cfg: BreakerPolicy): boolean {
|
||||
if (b === undefined) return false
|
||||
return now < b.openedAt + breakerCooldownMs(b.opens, cfg)
|
||||
}
|
||||
|
||||
/** 冷却结束时刻(无熔断时返回 0)。 */
|
||||
export function breakerUntil(b: BreakerState | undefined, cfg: BreakerPolicy): number {
|
||||
return b === undefined ? 0 : b.openedAt + breakerCooldownMs(b.opens, cfg)
|
||||
}
|
||||
|
||||
/** 再次熔断:`opens` 递增,`openedAt` 取本次时刻。 */
|
||||
export function openBreaker(prev: BreakerState | undefined, now: number): BreakerState {
|
||||
return { openedAt: now, opens: (prev?.opens ?? 0) + 1 }
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
/**
|
||||
* Loopback port guard: an iptables OUTPUT owner-match rule that lets only the
|
||||
* orchestrator (root) open a connection to a per-user DSH's loopback RPC port.
|
||||
* Because every user DSH binds the shared 127.0.0.1 loopback on a dynamic port,
|
||||
* a co-tenant local user could otherwise `curl` another user's DSH directly and
|
||||
* bypass the orchestrator's session authentication. The owner match filters on
|
||||
* the *client* uid, which is only observable on the OUTPUT chain for a
|
||||
* same-host loopback connection (an INPUT rule would see the receiving DSH's
|
||||
* uid and block nothing selectively).
|
||||
*
|
||||
* The REJECT is inserted at the *top* of OUTPUT (`-I OUTPUT 1`), not appended:
|
||||
* iptables stops at the first match, so an earlier ACCEPT (ufw or a conntrack
|
||||
* `ESTABLISHED,RELATED -j ACCEPT`) would otherwise swallow the packet before an
|
||||
* appended `-A` REJECT is ever evaluated.
|
||||
*
|
||||
* Linux + root only; opt in via `config.portGuard` and fail loud when enabled
|
||||
* on a host that cannot apply it.
|
||||
* @module dshs/supervisor/firewall
|
||||
*/
|
||||
|
||||
import { execFileSync } from 'node:child_process'
|
||||
|
||||
/** One guarded loopback port with idempotent install/remove. */
|
||||
export interface PortGuard {
|
||||
/** Add the OUTPUT owner-match REJECT rule (throws on failure). */
|
||||
install(port: number): void
|
||||
/** Remove the rule best-effort (a missing rule is not an error). */
|
||||
remove(port: number): void
|
||||
}
|
||||
|
||||
/** iptables arguments (excluding the executable) for one insert/delete.
|
||||
* Install inserts at position 1 (`-I OUTPUT 1`) to stay ahead of any earlier
|
||||
* ACCEPT (ufw / conntrack ESTABLISHED); delete matches by rule spec (`-D
|
||||
* OUTPUT`), which works regardless of position. */
|
||||
function ruleArgs(port: number, action: '-I' | '-D'): string[] {
|
||||
const insert = action === '-I' ? ['1'] : []
|
||||
return ['-t', 'filter', action, 'OUTPUT', ...insert, '-p', 'tcp', '--dport', String(port), '-m', 'owner', '!', '--uid-owner', '0', '-j', 'REJECT']
|
||||
}
|
||||
|
||||
/** Whether this process can manage the loopback OUTPUT guard. */
|
||||
function canGuard(): boolean {
|
||||
return process.platform === 'linux' && typeof process.getuid === 'function' && process.getuid() === 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Create the port guard when enabled, or undefined. Fails loud when enabled on
|
||||
* a host that cannot apply it (non-Linux or non-root): a silently absent guard
|
||||
* would leave co-tenants able to reach each other's DSH.
|
||||
* @param enabled - deployment flag (`config.portGuard`).
|
||||
*/
|
||||
export function createPortGuard(enabled: boolean): PortGuard | undefined {
|
||||
if (!enabled) return undefined
|
||||
if (!canGuard()) {
|
||||
throw new Error('portGuard requires a Linux host running as root (iptables OUTPUT owner-match)')
|
||||
}
|
||||
const guarded = new Set<number>()
|
||||
return {
|
||||
install(port) {
|
||||
if (guarded.has(port)) return
|
||||
// Insert at position 1 so an earlier ACCEPT (ufw / conntrack) can't
|
||||
// swallow the packet before the REJECT is evaluated.
|
||||
execFileSync('iptables', ruleArgs(port, '-I'))
|
||||
guarded.add(port)
|
||||
},
|
||||
remove(port) {
|
||||
if (!guarded.has(port)) return
|
||||
try {
|
||||
execFileSync('iptables', ruleArgs(port, '-D'))
|
||||
} catch {
|
||||
// Best-effort teardown: a stale rule leaves the port closed to
|
||||
// co-tenants, which is fail-safe, and a dropped rule is already gone.
|
||||
} finally {
|
||||
guarded.delete(port)
|
||||
}
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,687 @@
|
||||
/**
|
||||
* K8s backend for per-user DSH lifecycle (docs/k8s.md §5.2/§4.3/§4.4/§4.8).
|
||||
*
|
||||
* Each user gets a `dsh-<userId>` Pod (dsh + socat sidecar), a Headless Service,
|
||||
* a NetworkPolicy, and a `dsh-key-<userId>` Secret for the API key. A crash
|
||||
* pulls up a one-shot `dsh-<userId>-watchdog` Job. The control plane runs
|
||||
* in-cluster and talks to the K8s API via @kubernetes/client-node.
|
||||
* @module dshs/supervisor/k8s-spawner
|
||||
*/
|
||||
|
||||
import * as k8s from '@kubernetes/client-node'
|
||||
import type { ServerConfig } from '../config.js'
|
||||
import type { DbAdapter } from '../db/adapter.js'
|
||||
import { HANDOFF_FILE, HOME_DIR, USERS_DIR, WORKSPACE_DIR } from '../fs/workspace.js'
|
||||
import { FILE_SERVICE_PORT, USER_ROOT_ENV } from '../web/file-service.js'
|
||||
import type { Fencing } from './leader.js'
|
||||
import { AlreadyRunningError, type Endpoint, type Instance, type LivePod, type Spawner, type UserStatus } from './spawner.js'
|
||||
|
||||
/** Loopback port the dsh container binds; the socat sidecar bridges 8081 → 8080. */
|
||||
const DSH_LOOPBACK_PORT = 8080
|
||||
/** Sidecar bridge port the per-user Service targets (80 → 8081). */
|
||||
const SOCAT_PORT = 8081
|
||||
/** Task the one-shot headless watchdog runs (executes the handoff command). */
|
||||
const WATCHDOG_TASK = 'Read DSHS_HANDOFF_PATH. If it contains a JSON {"command": ...}, run that command. Then exit.'
|
||||
/** Shared RWX PVC every user's Pod mounts via subPath. */
|
||||
const USERS_PVC = 'dsh-users'
|
||||
|
||||
/** Per-user resource names (deterministic — idempotent create). */
|
||||
function names(userId: string) {
|
||||
return {
|
||||
pod: `dsh-${userId}`,
|
||||
service: `dsh-${userId}`,
|
||||
networkPolicy: `dsh-${userId}`,
|
||||
secret: `dsh-key-${userId}`,
|
||||
job: `dsh-${userId}-watchdog`,
|
||||
patch: `dsh-${userId}-patch`,
|
||||
filesPod: `dsh-files-${userId}`,
|
||||
filesService: `dsh-files-${userId}`,
|
||||
filesNetworkPolicy: `dsh-files-${userId}`,
|
||||
}
|
||||
}
|
||||
|
||||
/** The data root inside every per-user Pod, independent of the control plane's
|
||||
* own `dataRoot` (which under k8s points at a volume the Pod never sees). */
|
||||
const POD_DATA_ROOT = '/var/lib/dshs'
|
||||
|
||||
/** Home/workspace paths inside the DSH Pod (docs/k8s.md §4.3). Built with
|
||||
* POSIX separators — the control plane may be developed/tested on Windows. */
|
||||
function userPaths(userId: string): { home: string; ws: string; mount: string } {
|
||||
const mount = `${POD_DATA_ROOT}/${USERS_DIR}/${userId}`
|
||||
return { home: `${mount}/${HOME_DIR}`, ws: `${mount}/${WORKSPACE_DIR}`, mount }
|
||||
}
|
||||
|
||||
/** Pod-safe labels shared by the DSH Pod/Service/NetworkPolicy/Job. */
|
||||
function podLabels(userId: string): { app: string; user: string } {
|
||||
return { app: 'dsh', user: userId }
|
||||
}
|
||||
|
||||
/** Labels for the per-user file sidecar (a distinct `app`, so its own
|
||||
* Service/NetworkPolicy select it without touching the DSH Pod). */
|
||||
function filesLabels(userId: string): { app: string; user: string } {
|
||||
return { app: 'dsh-files', user: userId }
|
||||
}
|
||||
|
||||
/** API-key env entry, omitted when the user has no key. */
|
||||
function apiKeyEnv(userId: string, apiKey: string | null): k8s.V1EnvVar[] {
|
||||
if (apiKey === null) return []
|
||||
return [{ name: 'DEEPSEEK_API_KEY', valueFrom: { secretKeyRef: { name: names(userId).secret, key: 'key' } } }]
|
||||
}
|
||||
|
||||
/** Common per-user container security context (non-root + drop ALL + seccomp +
|
||||
* read-only rootfs — /tmp comes from an emptyDir, docs/k8s.md Phase 4). */
|
||||
function containerSecurity(uid: number): k8s.V1SecurityContext {
|
||||
return {
|
||||
runAsNonRoot: true,
|
||||
runAsUser: uid,
|
||||
allowPrivilegeEscalation: false,
|
||||
capabilities: { drop: ['ALL'] },
|
||||
seccompProfile: { type: 'RuntimeDefault' },
|
||||
readOnlyRootFilesystem: true,
|
||||
}
|
||||
}
|
||||
|
||||
/** Writable scratch volume + mount for containers whose rootfs is read-only. */
|
||||
function tmpVolume(): k8s.V1Volume[] {
|
||||
return [{ name: 'tmp', emptyDir: {} }]
|
||||
}
|
||||
function tmpMount(): k8s.V1VolumeMount {
|
||||
return { name: 'tmp', mountPath: '/tmp' }
|
||||
}
|
||||
|
||||
/** The shared RWX volume (subPath mounts per-user at use time). */
|
||||
function dataVolume(): k8s.V1Volume[] {
|
||||
return [{ name: 'data', persistentVolumeClaim: { claimName: USERS_PVC } }]
|
||||
}
|
||||
|
||||
/** The shared RWX volume mounted at its **root** (no subPath) so an init
|
||||
* container can create/chown `<pvc>/<userId>` as the user's own uid. */
|
||||
function dataRootVolume(): k8s.V1Volume[] {
|
||||
return [{ name: 'data-root', persistentVolumeClaim: { claimName: USERS_PVC } }]
|
||||
}
|
||||
|
||||
/** Every per-user Pod/Job references ACR private images, so each must carry the
|
||||
* cluster's pull secret (the control-plane Deployment has it in its manifest,
|
||||
* but generated Pods don't inherit it). */
|
||||
function pullSecrets(name: string): k8s.V1LocalObjectReference[] {
|
||||
return name === '' ? [] : [{ name }]
|
||||
}
|
||||
|
||||
/** Fencing labels stamped onto resources created by the leader, so a successor
|
||||
* can identify work a dead leader left in flight (docs/k8s.md §5.3). */
|
||||
function stampFencing(labels: Record<string, string>, fencing: Fencing | undefined): void {
|
||||
if (fencing === undefined) return
|
||||
labels['dsh.io/holder'] = fencing.holder
|
||||
labels['dsh.io/operation-id'] = String(fencing.operationId)
|
||||
}
|
||||
|
||||
/**
|
||||
* K8s backend implementing {@link Spawner}. State lives in the cluster; every
|
||||
* method is a K8s API call (or a read).
|
||||
*/
|
||||
export class K8sSpawner implements Spawner {
|
||||
private readonly core: k8s.CoreV1Api
|
||||
private readonly networking: k8s.NetworkingV1Api
|
||||
private readonly batch: k8s.BatchV1Api
|
||||
private readonly namespace: string
|
||||
private readonly kc: k8s.KubeConfig
|
||||
private fencing: Fencing | undefined
|
||||
|
||||
constructor(
|
||||
private readonly config: ServerConfig,
|
||||
private readonly db: DbAdapter,
|
||||
private readonly resolveApiKey: (userId: string) => Promise<string | null>,
|
||||
private readonly resolveUid: (userId: string) => Promise<number>,
|
||||
clients?: { core: k8s.CoreV1Api; networking: k8s.NetworkingV1Api; batch: k8s.BatchV1Api },
|
||||
) {
|
||||
// Injecting clients short-circuits cluster auth (used by tests). Otherwise
|
||||
// build them from the in-cluster config, which needs a mounted SA token.
|
||||
const kc = new k8s.KubeConfig()
|
||||
if (clients === undefined) {
|
||||
kc.loadFromCluster()
|
||||
clients = {
|
||||
core: kc.makeApiClient(k8s.CoreV1Api),
|
||||
networking: kc.makeApiClient(k8s.NetworkingV1Api),
|
||||
batch: kc.makeApiClient(k8s.BatchV1Api),
|
||||
}
|
||||
}
|
||||
this.kc = kc
|
||||
this.core = clients.core
|
||||
this.networking = clients.networking
|
||||
this.batch = clients.batch
|
||||
this.namespace = config.k8sNamespace
|
||||
}
|
||||
|
||||
async launch(userId: string, folder: string, patch?: string): Promise<Instance> {
|
||||
const n = names(userId)
|
||||
if (await this.podExists(n.pod)) throw new AlreadyRunningError(userId)
|
||||
await this.ensureFileService(userId) // the DSH Pod's subPath must already exist
|
||||
const apiKey = await this.resolveApiKey(userId)
|
||||
const uid = await this.resolveUid(userId)
|
||||
const hasPatch = this.config.enablePatch && patch !== undefined
|
||||
if (hasPatch) await this.ensurePatchConfigMap(n.patch, patch)
|
||||
await this.ensureSecret(n.secret, apiKey)
|
||||
await this.ensurePod(n.pod, userId, uid, apiKey, hasPatch ? n.patch : undefined)
|
||||
await this.ensureService(n.service, userId)
|
||||
await this.ensureNetworkPolicy(n.networkPolicy, userId)
|
||||
try {
|
||||
await this.db.upsertInstance({ id: n.pod, userId, role: 'main', status: 'starting', folder, patch })
|
||||
} catch (err) {
|
||||
this.logError(err) // reconcile needs this row; don't fail the launch over it
|
||||
}
|
||||
return { id: n.pod, userId, role: 'main', folder, status: 'starting', patch }
|
||||
}
|
||||
|
||||
async restartMain(userId: string): Promise<Instance | undefined> {
|
||||
const desired = await this.db.findUserInstance(userId, 'main')
|
||||
if (desired === undefined) return undefined
|
||||
await this.stop(userId)
|
||||
return await this.launch(userId, desired.folder ?? '', desired.patch ?? undefined)
|
||||
}
|
||||
|
||||
async restartAllMains(): Promise<void> {
|
||||
const rows = await this.db.listInstancesByRole('main')
|
||||
for (const row of rows) {
|
||||
try {
|
||||
await this.restartMain(row.userId)
|
||||
} catch (err) {
|
||||
this.logError(err) // keep broadcasting — one pod failure must not abort the rest
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async spawnWatchdog(userId: string): Promise<Instance | undefined> {
|
||||
const n = names(userId)
|
||||
const apiKey = await this.resolveApiKey(userId)
|
||||
const uid = await this.resolveUid(userId)
|
||||
const { home, ws, mount } = userPaths(userId)
|
||||
const jobLabels = podLabels(userId)
|
||||
const podTemplateLabels = podLabels(userId)
|
||||
stampFencing(jobLabels, this.fencing)
|
||||
stampFencing(podTemplateLabels, this.fencing)
|
||||
await this.batch.createNamespacedJob({ namespace: this.namespace, body: {
|
||||
apiVersion: 'batch/v1',
|
||||
kind: 'Job',
|
||||
metadata: { name: n.job, namespace: this.namespace, labels: jobLabels },
|
||||
spec: {
|
||||
ttlSecondsAfterFinished: 300,
|
||||
template: {
|
||||
metadata: { labels: podTemplateLabels },
|
||||
spec: {
|
||||
automountServiceAccountToken: false,
|
||||
imagePullSecrets: pullSecrets(this.config.imagePullSecret),
|
||||
restartPolicy: 'Never',
|
||||
securityContext: {
|
||||
runAsNonRoot: true,
|
||||
runAsUser: uid,
|
||||
fsGroup: uid,
|
||||
seccompProfile: { type: 'RuntimeDefault' },
|
||||
},
|
||||
containers: [
|
||||
{
|
||||
name: 'dsh',
|
||||
image: this.config.dshImage,
|
||||
args: ['--profile', 'headless', WATCHDOG_TASK],
|
||||
env: [
|
||||
{ name: 'HOME', value: ws },
|
||||
{ name: 'DSH_HOME', value: home },
|
||||
{ name: 'DSHS_ROLE', value: 'watchdog' },
|
||||
{ name: 'DSHS_HANDOFF_PATH', value: `${mount}/${HANDOFF_FILE}` },
|
||||
...apiKeyEnv(userId, apiKey),
|
||||
],
|
||||
volumeMounts: [{ name: 'data', mountPath: mount, subPath: userId }, tmpMount()],
|
||||
securityContext: containerSecurity(uid),
|
||||
},
|
||||
],
|
||||
volumes: [...dataVolume(), ...tmpVolume()],
|
||||
},
|
||||
},
|
||||
},
|
||||
} })
|
||||
return { id: n.job, userId, role: 'watchdog', folder: ws, status: 'starting' }
|
||||
}
|
||||
|
||||
async status(userId: string): Promise<UserStatus> {
|
||||
return { main: await this.readInstance(names(userId).pod, userId, 'main') }
|
||||
}
|
||||
|
||||
async endpointFor(userId: string): Promise<Endpoint | undefined> {
|
||||
let pod: k8s.V1Pod
|
||||
try {
|
||||
pod = await this.core.readNamespacedPod({ name: names(userId).pod, namespace: this.namespace })
|
||||
} catch (err) {
|
||||
if (this.isNotFound(err)) return undefined
|
||||
throw err // 403/500 are real failures, not "not running"
|
||||
}
|
||||
if (pod.status?.phase !== 'Running') return undefined
|
||||
// Dial the Pod IP directly (not the Headless Service DNS): re-reading the Pod
|
||||
// on every request gives the *current* IP immediately after a restart, so a
|
||||
// rebuilt Pod never leaves the proxy pointing at a stale IP for the ~30s DNS
|
||||
// TTL. Port is the sidecar's 8081 (no kube-proxy DNAT on Headless Services).
|
||||
const ip = pod.status?.podIP
|
||||
if (ip === undefined || ip === '') return undefined
|
||||
return { host: ip, port: SOCAT_PORT }
|
||||
}
|
||||
|
||||
// k8s 模式无 launch token 概念(Pod 就绪由 endpointFor 的 phase 判断),no-op。
|
||||
async waitForLaunchTokenForUser(_userId: string, _timeoutMs?: number): Promise<void> {}
|
||||
|
||||
/**
|
||||
* k8s spawner 未实现探活(本部署走本地 spawner)。返回 ok:true 保持与旧行为一致:
|
||||
* 不做回滚,交由 k8s 自身 readiness/liveness 处理。
|
||||
*/
|
||||
async restartAndProbe(_userId: string, _settleMs?: number): Promise<{ ok: boolean; reason: string }> {
|
||||
return { ok: true, reason: "k8s spawner: 探活未实现" }
|
||||
}
|
||||
|
||||
async stop(userId: string): Promise<void> {
|
||||
const n = names(userId)
|
||||
await this.ignoreNotFound(() => this.core.deleteNamespacedPod({ name: n.pod, namespace: this.namespace }))
|
||||
await this.ignoreNotFound(() => this.core.deleteNamespacedService({ name: n.service, namespace: this.namespace }))
|
||||
await this.ignoreNotFound(() => this.networking.deleteNamespacedNetworkPolicy({ name: n.networkPolicy, namespace: this.namespace }))
|
||||
await this.ignoreNotFound(() => this.batch.deleteNamespacedJob({ name: n.job, namespace: this.namespace }))
|
||||
await this.ignoreNotFound(() => this.core.deleteNamespacedConfigMap({ name: n.patch, namespace: this.namespace }))
|
||||
try {
|
||||
await this.db.deleteInstance(n.pod) // desired state is gone once stopped
|
||||
} catch (err) {
|
||||
this.logError(err)
|
||||
}
|
||||
}
|
||||
|
||||
/** No-op: per-user Pods outlive any single control-plane replica (reconcile/leader manages them). */
|
||||
async teardown(): Promise<void> {}
|
||||
|
||||
/** Activity signal is unused under k8s: the reconcile loop idles-reaps by
|
||||
* session presence (reconcile.ts Phase 4), not by proxied-traffic recency. */
|
||||
touch(_userId: string): void {}
|
||||
|
||||
/** Bring up the user's file sidecar (docs/k8s.md §4.10). The sidecar is a
|
||||
* distinct always-on Pod so the desktop is usable *before* the on-demand DSH
|
||||
* launches; its init container creates the user's subPath directory. */
|
||||
async ensureFileService(userId: string): Promise<void> {
|
||||
if (this.config.controlPlaneImage === '') {
|
||||
throw new Error('DSHS_CONTROL_PLANE_IMAGE is required in k8s mode (file sidecar image)')
|
||||
}
|
||||
const n = names(userId)
|
||||
const uid = await this.resolveUid(userId)
|
||||
await this.ensureFilesPod(n.filesPod, userId, uid)
|
||||
await this.ensureFilesService(n.filesService, userId)
|
||||
await this.ensureFilesNetworkPolicy(n.filesNetworkPolicy, userId)
|
||||
// The Headless Service only publishes an A record once the Pod is Ready;
|
||||
// the caller resolves it immediately, so block until it comes up.
|
||||
await this.waitForPodReady(n.filesPod)
|
||||
}
|
||||
|
||||
// --- reconcile / watch support (leader-only callers) ---
|
||||
|
||||
/** Stamp the current leader's fencing token onto resources created from now on. */
|
||||
setFencing(fencing: Fencing | undefined): void {
|
||||
this.fencing = fencing
|
||||
}
|
||||
|
||||
/** Main DSH Pods (`app=dsh`), mapped to the shape the controller needs. */
|
||||
async listUserPods(): Promise<LivePod[]> {
|
||||
const res = await this.core.listNamespacedPod({ namespace: this.namespace, labelSelector: 'app=dsh' })
|
||||
return (res.items ?? []).map((pod) => {
|
||||
const phase = pod.status?.phase ?? ''
|
||||
return {
|
||||
name: pod.metadata?.name ?? '',
|
||||
userId: pod.metadata?.labels?.user ?? '',
|
||||
running: phase === 'Running',
|
||||
crashed: phase === 'Failed' || phase === 'Succeeded',
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/** Recreate a lost Service/NetworkPolicy for a user whose main Pod exists. */
|
||||
async ensureUserResources(userId: string): Promise<void> {
|
||||
await this.ensureService(names(userId).service, userId)
|
||||
await this.ensureNetworkPolicy(names(userId).networkPolicy, userId)
|
||||
}
|
||||
|
||||
/** Watch main DSH Pods; `callback` fires on add/update/delete with their state. */
|
||||
watchMainPods(callback: (pod: LivePod) => void): k8s.Informer<k8s.V1Pod> & k8s.ObjectCache<k8s.V1Pod> {
|
||||
const informer = k8s.makeInformer<k8s.V1Pod>(
|
||||
this.kc,
|
||||
`/api/v1/namespaces/${this.namespace}/pods`,
|
||||
() => this.core.listNamespacedPod({ namespace: this.namespace, labelSelector: 'app=dsh' }),
|
||||
'app=dsh',
|
||||
)
|
||||
const emit = (obj: k8s.V1Pod | undefined): void => {
|
||||
const phase = obj?.status?.phase ?? ''
|
||||
callback({
|
||||
name: obj?.metadata?.name ?? '',
|
||||
userId: obj?.metadata?.labels?.user ?? '',
|
||||
running: phase === 'Running',
|
||||
crashed: phase === 'Failed' || phase === 'Succeeded',
|
||||
})
|
||||
}
|
||||
informer.on('add', (obj) => emit(obj))
|
||||
informer.on('update', (obj) => emit(obj))
|
||||
informer.on('delete', (obj) => emit(obj))
|
||||
return informer
|
||||
}
|
||||
|
||||
/** Best-effort error surface for the controller's tick loop. */
|
||||
logError(err: unknown): void {
|
||||
console.error('[dsh-reconcile]', err)
|
||||
}
|
||||
|
||||
// --- helpers ---
|
||||
|
||||
private async podExists(name: string): Promise<boolean> {
|
||||
try {
|
||||
await this.core.readNamespacedPod({ name, namespace: this.namespace })
|
||||
return true
|
||||
} catch (err) {
|
||||
if (this.isNotFound(err)) return false
|
||||
throw err // 403/500 are real failures, not "no Pod"
|
||||
}
|
||||
}
|
||||
|
||||
private async readInstance(name: string, userId: string, role: 'main' | 'watchdog'): Promise<Instance | undefined> {
|
||||
try {
|
||||
const pod = await this.core.readNamespacedPod({ name, namespace: this.namespace })
|
||||
const status = pod.status?.phase === 'Running' ? 'running' : pod.status?.phase === 'Failed' ? 'crashed' : 'starting'
|
||||
return { id: name, userId, role, folder: '', status }
|
||||
} catch (err) {
|
||||
if (this.isNotFound(err)) return undefined
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
/** Create-or-replace, so a relaunch after `stop()` (which deletes these) can
|
||||
* never 409 on a stale Secret/ConfigMap from a prior launch. */
|
||||
private async ensureSecret(name: string, apiKey: string | null): Promise<void> {
|
||||
if (apiKey === null) return
|
||||
await this.replace({
|
||||
read: () => this.core.readNamespacedSecret({ name, namespace: this.namespace }),
|
||||
create: () => this.core.createNamespacedSecret({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'Secret',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
type: 'Opaque',
|
||||
stringData: { key: apiKey },
|
||||
} }),
|
||||
del: () => this.core.deleteNamespacedSecret({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
/** Create-or-replace a generated resource. Reads first for the cheap path,
|
||||
* then treats a 409 race by deleting and recreating (we own every resource
|
||||
* this helper touches, so replacement is safe). */
|
||||
private async replace(ops: {
|
||||
read: () => Promise<unknown>
|
||||
create: () => Promise<unknown>
|
||||
del: () => Promise<unknown>
|
||||
}): Promise<void> {
|
||||
try {
|
||||
await ops.read()
|
||||
} catch (err) {
|
||||
if (!this.isNotFound(err)) throw err
|
||||
try {
|
||||
await ops.create()
|
||||
} catch (createErr) {
|
||||
if (!this.isConflict(createErr)) throw createErr
|
||||
await this.ignoreNotFound(ops.del)
|
||||
await ops.create()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async ensurePod(name: string, userId: string, uid: number, apiKey: string | null, patchConfigMapName?: string): Promise<void> {
|
||||
const { home, ws, mount } = userPaths(userId)
|
||||
const args = ['--profile', 'web', '--host', '127.0.0.1', '--port', String(DSH_LOOPBACK_PORT)]
|
||||
// The runtime plugin (dshs/runtime) is baked into the dsh image
|
||||
// but loaded via --patch; the rendered patch is mounted at /etc/dsh/patch.yml.
|
||||
if (patchConfigMapName !== undefined) args.splice(1, 0, '--patch', '/etc/dsh/patch.yml')
|
||||
const labels = podLabels(userId)
|
||||
stampFencing(labels, this.fencing)
|
||||
await this.core.createNamespacedPod({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'Pod',
|
||||
metadata: { name, namespace: this.namespace, labels },
|
||||
spec: {
|
||||
automountServiceAccountToken: false,
|
||||
imagePullSecrets: pullSecrets(this.config.imagePullSecret),
|
||||
hostNetwork: false,
|
||||
hostPID: false,
|
||||
securityContext: {
|
||||
runAsNonRoot: true,
|
||||
runAsUser: uid,
|
||||
fsGroup: uid,
|
||||
seccompProfile: { type: 'RuntimeDefault' },
|
||||
},
|
||||
containers: [
|
||||
{
|
||||
name: 'dsh',
|
||||
image: this.config.dshImage,
|
||||
args,
|
||||
// DSH binds loopback only, so a tcpSocket probe (which hits the Pod IP)
|
||||
// would never succeed; probe 127.0.0.1:8080 from inside the container.
|
||||
readinessProbe: {
|
||||
exec: { command: ['node', '-e', 'require("http").get("http://127.0.0.1:8080", r => process.exit(r.statusCode < 500 ? 0 : 1)).on("error", () => process.exit(1))'] },
|
||||
initialDelaySeconds: 5,
|
||||
periodSeconds: 3,
|
||||
},
|
||||
env: [
|
||||
{ name: 'HOME', value: ws },
|
||||
{ name: 'DSH_HOME', value: home },
|
||||
...apiKeyEnv(userId, apiKey),
|
||||
],
|
||||
volumeMounts: [
|
||||
{ name: 'data', mountPath: mount, subPath: userId },
|
||||
...(patchConfigMapName !== undefined ? [{ name: 'patch', mountPath: '/etc/dsh', readOnly: true }] : []),
|
||||
tmpMount(),
|
||||
],
|
||||
securityContext: containerSecurity(uid),
|
||||
resources: { requests: { cpu: '500m', memory: '1Gi' }, limits: { cpu: '2', memory: '4Gi' } },
|
||||
},
|
||||
{
|
||||
name: 'sidecar',
|
||||
image: this.config.controlPlaneImage,
|
||||
args: ['tcp-bridge', `0.0.0.0:${SOCAT_PORT}`, `127.0.0.1:${DSH_LOOPBACK_PORT}`],
|
||||
ports: [{ containerPort: SOCAT_PORT }],
|
||||
securityContext: containerSecurity(uid),
|
||||
resources: { requests: { cpu: '10m', memory: '32Mi' }, limits: { cpu: '100m', memory: '128Mi' } },
|
||||
},
|
||||
],
|
||||
volumes: [
|
||||
...dataVolume(),
|
||||
...(patchConfigMapName !== undefined ? [{ name: 'patch', configMap: { name: patchConfigMapName } }] : []),
|
||||
...tmpVolume(),
|
||||
],
|
||||
},
|
||||
} })
|
||||
}
|
||||
|
||||
private async ensurePatchConfigMap(name: string, patch: string): Promise<void> {
|
||||
await this.replace({
|
||||
read: () => this.core.readNamespacedConfigMap({ name, namespace: this.namespace }),
|
||||
create: () => this.core.createNamespacedConfigMap({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'ConfigMap',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
data: { 'patch.yml': patch },
|
||||
} }),
|
||||
del: () => this.core.deleteNamespacedConfigMap({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
private async ensureService(name: string, userId: string): Promise<void> {
|
||||
await this.replace({
|
||||
read: () => this.core.readNamespacedService({ name, namespace: this.namespace }),
|
||||
create: () => this.core.createNamespacedService({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'Service',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
spec: {
|
||||
clusterIP: 'None', // Headless: no ClusterIP, DNS A record → Pod IP
|
||||
selector: podLabels(userId),
|
||||
ports: [{ port: 80, targetPort: SOCAT_PORT }],
|
||||
},
|
||||
} }),
|
||||
del: () => this.core.deleteNamespacedService({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
private async ensureNetworkPolicy(name: string, userId: string): Promise<void> {
|
||||
await this.replace({
|
||||
read: () => this.networking.readNamespacedNetworkPolicy({ name, namespace: this.namespace }),
|
||||
create: () => this.networking.createNamespacedNetworkPolicy({ namespace: this.namespace, body: {
|
||||
apiVersion: 'networking.k8s.io/v1',
|
||||
kind: 'NetworkPolicy',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
spec: {
|
||||
podSelector: { matchLabels: podLabels(userId) },
|
||||
policyTypes: ['Ingress', 'Egress'],
|
||||
ingress: [{ _from: [{ podSelector: { matchLabels: { app: 'dsh-orchestrator' } } }] }],
|
||||
egress: [
|
||||
{
|
||||
// DNS to anywhere: the DNS backend varies by distribution (kube-dns /
|
||||
// coredns / node-local-dns on ACK), so don't pin a pod label.
|
||||
to: [{ ipBlock: { cidr: '0.0.0.0/0' } }],
|
||||
ports: [
|
||||
{ port: 53, protocol: 'UDP' },
|
||||
{ port: 53, protocol: 'TCP' },
|
||||
],
|
||||
},
|
||||
{
|
||||
to: this.config.egressCidrs.length > 0
|
||||
? this.config.egressCidrs.map((cidr) => ({ ipBlock: { cidr, except: ['169.254.169.254/32', '10.0.0.0/8', '172.16.0.0/12', '192.168.0.0/16'] } }))
|
||||
: [{ ipBlock: { cidr: '0.0.0.0/0', except: ['169.254.169.254/32', '10.0.0.0/8', '172.16.0.0/12', '192.168.0.0/16'] } }],
|
||||
ports: [{ port: 443 }],
|
||||
},
|
||||
],
|
||||
},
|
||||
} }),
|
||||
del: () => this.networking.deleteNamespacedNetworkPolicy({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
private async ensureFilesPod(name: string, userId: string, uid: number): Promise<void> {
|
||||
const mount = userPaths(userId).mount
|
||||
await this.replace({
|
||||
read: () => this.core.readNamespacedPod({ name, namespace: this.namespace }),
|
||||
create: () => this.core.createNamespacedPod({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'Pod',
|
||||
metadata: { name, namespace: this.namespace, labels: filesLabels(userId) },
|
||||
spec: {
|
||||
automountServiceAccountToken: false,
|
||||
imagePullSecrets: pullSecrets(this.config.imagePullSecret),
|
||||
hostNetwork: false,
|
||||
hostPID: false,
|
||||
securityContext: {
|
||||
runAsNonRoot: true,
|
||||
runAsUser: uid,
|
||||
fsGroup: uid,
|
||||
seccompProfile: { type: 'RuntimeDefault' },
|
||||
},
|
||||
// Creates `<pvc>/<userId>/{ws,home}` as the user's uid (docs/k8s.md
|
||||
// §4.9). Runs on the PVC *root* (no subPath), because the subPath dir
|
||||
// may not exist yet or be root-owned; the user's 0700 dir keeps other
|
||||
// users' files out of reach.
|
||||
initContainers: [
|
||||
{
|
||||
name: 'init-user',
|
||||
image: this.config.controlPlaneImage,
|
||||
command: ['sh', '-ec', `mkdir -p /mnt/${userId}/${WORKSPACE_DIR} /mnt/${userId}/${HOME_DIR} && chmod 0700 /mnt/${userId}`],
|
||||
volumeMounts: [{ name: 'data-root', mountPath: '/mnt' }],
|
||||
securityContext: containerSecurity(uid),
|
||||
},
|
||||
],
|
||||
containers: [
|
||||
{
|
||||
name: 'files',
|
||||
image: this.config.controlPlaneImage,
|
||||
args: ['file-service'],
|
||||
env: [{ name: USER_ROOT_ENV, value: mount }],
|
||||
ports: [{ containerPort: FILE_SERVICE_PORT }],
|
||||
readinessProbe: { tcpSocket: { port: FILE_SERVICE_PORT }, initialDelaySeconds: 2, periodSeconds: 3 },
|
||||
volumeMounts: [{ name: 'data', mountPath: mount, subPath: userId }, tmpMount()],
|
||||
securityContext: containerSecurity(uid),
|
||||
resources: { requests: { cpu: '50m', memory: '128Mi' }, limits: { cpu: '200m', memory: '256Mi' } },
|
||||
},
|
||||
],
|
||||
volumes: [...dataVolume(), ...dataRootVolume(), ...tmpVolume()],
|
||||
},
|
||||
} }),
|
||||
del: () => this.core.deleteNamespacedPod({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
private async ensureFilesService(name: string, userId: string): Promise<void> {
|
||||
await this.replace({
|
||||
read: () => this.core.readNamespacedService({ name, namespace: this.namespace }),
|
||||
create: () => this.core.createNamespacedService({ namespace: this.namespace, body: {
|
||||
apiVersion: 'v1',
|
||||
kind: 'Service',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
spec: {
|
||||
clusterIP: 'None',
|
||||
selector: filesLabels(userId),
|
||||
ports: [{ port: FILE_SERVICE_PORT, targetPort: FILE_SERVICE_PORT }],
|
||||
},
|
||||
} }),
|
||||
del: () => this.core.deleteNamespacedService({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
private async ensureFilesNetworkPolicy(name: string, userId: string): Promise<void> {
|
||||
await this.replace({
|
||||
read: () => this.networking.readNamespacedNetworkPolicy({ name, namespace: this.namespace }),
|
||||
create: () => this.networking.createNamespacedNetworkPolicy({ namespace: this.namespace, body: {
|
||||
apiVersion: 'networking.k8s.io/v1',
|
||||
kind: 'NetworkPolicy',
|
||||
metadata: { name, namespace: this.namespace },
|
||||
spec: {
|
||||
podSelector: { matchLabels: filesLabels(userId) },
|
||||
policyTypes: ['Ingress', 'Egress'],
|
||||
ingress: [{ _from: [{ podSelector: { matchLabels: { app: 'dsh-orchestrator' } } }], ports: [{ port: FILE_SERVICE_PORT }] }],
|
||||
// Files only; the sidecar never talks to the LLM API.
|
||||
egress: [
|
||||
{
|
||||
to: [{ ipBlock: { cidr: '0.0.0.0/0' } }],
|
||||
ports: [
|
||||
{ port: 53, protocol: 'UDP' },
|
||||
{ port: 53, protocol: 'TCP' },
|
||||
],
|
||||
},
|
||||
],
|
||||
},
|
||||
} }),
|
||||
del: () => this.networking.deleteNamespacedNetworkPolicy({ name, namespace: this.namespace }),
|
||||
})
|
||||
}
|
||||
|
||||
/** Poll until a Pod's Ready condition is true, or fail after `timeoutMs`. */
|
||||
private async waitForPodReady(name: string, timeoutMs = 60_000): Promise<void> {
|
||||
const deadline = Date.now() + timeoutMs
|
||||
for (;;) {
|
||||
const pod = await this.core.readNamespacedPod({ name, namespace: this.namespace })
|
||||
const ready = pod.status?.conditions?.some((c) => c.type === 'Ready' && c.status === 'True')
|
||||
if (ready) return
|
||||
if (Date.now() >= deadline) throw new Error(`Pod ${name} not ready within ${timeoutMs}ms`)
|
||||
await new Promise((resolve) => setTimeout(resolve, 1000))
|
||||
}
|
||||
}
|
||||
|
||||
private isNotFound(err: unknown): boolean {
|
||||
return (err as { code?: number }).code === 404
|
||||
}
|
||||
|
||||
private isConflict(err: unknown): boolean {
|
||||
return (err as { code?: number }).code === 409
|
||||
}
|
||||
|
||||
private async ignoreNotFound(fn: () => Promise<unknown>): Promise<void> {
|
||||
try {
|
||||
await fn()
|
||||
} catch (err) {
|
||||
// `@kubernetes/client-node` throws `ApiException` with a `code` field.
|
||||
if (this.isNotFound(err)) return
|
||||
throw err
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,254 @@
|
||||
/**
|
||||
* Hand-rolled leader election over `coordination.k8s.io/v1` Lease
|
||||
* (docs/k8s.md §5.3). `@kubernetes/client-node` ships no election helper, so
|
||||
* this holds the small state machine: try to create the lease, and when it
|
||||
* already exists either renew (we hold it) or take it over only after the
|
||||
* holder's renew time has exceeded `leaseDurationSeconds`.
|
||||
*
|
||||
* Timings follow the doc's anti-split-brain ordering
|
||||
* `LeaseDuration > RenewDeadline > RetryPeriod`.
|
||||
* @module dshs/supervisor/leader
|
||||
*/
|
||||
|
||||
import * as k8s from '@kubernetes/client-node'
|
||||
import { hostname } from 'node:os'
|
||||
|
||||
/** The fencing token a controller stamps onto resources it creates. */
|
||||
export interface Fencing {
|
||||
holder: string
|
||||
operationId: number
|
||||
}
|
||||
|
||||
export interface LeaderOptions {
|
||||
identity?: string
|
||||
leaseName?: string
|
||||
namespace: string
|
||||
leaseDurationSeconds?: number
|
||||
renewDeadlineSeconds?: number
|
||||
retryPeriodSeconds?: number
|
||||
/** Injectable clock (tests); defaults to Date.now. */
|
||||
now?: () => number
|
||||
onStartedLeading?: (fencing: Fencing) => void
|
||||
onStoppedLeading?: () => void
|
||||
/** Injectable lease client (tests); defaults to the in-cluster config. */
|
||||
coordination?: k8s.CoordinationV1Api
|
||||
}
|
||||
|
||||
/** `@kubernetes/client-node` deserializes `V1MicroTime` back to an ISO string,
|
||||
* not a Date — `renewTime.getTime()` would throw. Normalize either shape. */
|
||||
function toMillis(time: unknown): number {
|
||||
if (time === undefined || time === null) return 0
|
||||
if (typeof time === 'string') {
|
||||
const ms = Date.parse(time)
|
||||
return Number.isNaN(ms) ? 0 : ms
|
||||
}
|
||||
if (time instanceof Date) return time.getTime()
|
||||
return 0
|
||||
}
|
||||
|
||||
/** Normalize a read-back `V1MicroTime` (string) back to a Date for the replace body. */
|
||||
function toMicroTime(time: unknown): k8s.V1MicroTime | undefined {
|
||||
const ms = toMillis(time)
|
||||
return ms > 0 ? new k8s.V1MicroTime(ms) : undefined
|
||||
}
|
||||
|
||||
/** Acquire/keep a Lease, notifying callers across leadership changes. */
|
||||
export class LeaderElector {
|
||||
private readonly coordination: k8s.CoordinationV1Api
|
||||
private readonly identity: string
|
||||
private readonly leaseName: string
|
||||
private readonly namespace: string
|
||||
private readonly leaseDurationSeconds: number
|
||||
private readonly renewDeadlineSeconds: number
|
||||
private readonly retryPeriodSeconds: number
|
||||
private readonly now: () => number
|
||||
private onStartedLeading?: (fencing: Fencing) => void
|
||||
private onStoppedLeading?: () => void
|
||||
|
||||
private leading = false
|
||||
private operationId = 0
|
||||
private lastRenew = 0
|
||||
private timer: NodeJS.Timeout | undefined
|
||||
|
||||
constructor(options: LeaderOptions) {
|
||||
if (options.coordination !== undefined) {
|
||||
this.coordination = options.coordination
|
||||
} else {
|
||||
const kc = new k8s.KubeConfig()
|
||||
kc.loadFromCluster()
|
||||
this.coordination = kc.makeApiClient(k8s.CoordinationV1Api)
|
||||
}
|
||||
this.identity = options.identity ?? process.env.POD_NAME ?? hostname()
|
||||
this.leaseName = options.leaseName ?? 'dsh-orchestrator'
|
||||
this.namespace = options.namespace
|
||||
this.leaseDurationSeconds = options.leaseDurationSeconds ?? 15
|
||||
this.renewDeadlineSeconds = options.renewDeadlineSeconds ?? 10
|
||||
this.retryPeriodSeconds = options.retryPeriodSeconds ?? 2
|
||||
this.now = options.now ?? Date.now
|
||||
this.onStartedLeading = options.onStartedLeading
|
||||
this.onStoppedLeading = options.onStoppedLeading
|
||||
}
|
||||
|
||||
get isLeader(): boolean {
|
||||
return this.leading
|
||||
}
|
||||
|
||||
/** Wire (or rewire) leadership callbacks. The controller owns the reaction,
|
||||
* not the elector, so it attaches itself here. */
|
||||
setLeadershipCallbacks(onStarted?: (fencing: Fencing) => void, onStopped?: () => void): void {
|
||||
this.onStartedLeading = onStarted
|
||||
this.onStoppedLeading = onStopped
|
||||
}
|
||||
|
||||
/** The current fencing token (valid only while {@link isLeader}). */
|
||||
get fencing(): Fencing {
|
||||
return { holder: this.identity, operationId: this.operationId }
|
||||
}
|
||||
|
||||
/** Start the acquire/renew loop. Resolves once started (does not wait for leadership). */
|
||||
async start(): Promise<void> {
|
||||
if (this.timer !== undefined) return
|
||||
await this.tryAcquire()
|
||||
}
|
||||
|
||||
/** Stop the loop and yield leadership if held. */
|
||||
stop(): void {
|
||||
if (this.timer !== undefined) {
|
||||
clearTimeout(this.timer)
|
||||
this.timer = undefined
|
||||
}
|
||||
if (this.leading) {
|
||||
this.leading = false
|
||||
this.onStoppedLeading?.()
|
||||
}
|
||||
}
|
||||
|
||||
/** One acquire/renew round, rescheduling itself with retry/backoff. */
|
||||
private async tryAcquire(): Promise<void> {
|
||||
try {
|
||||
if (this.leading) {
|
||||
await this.renew()
|
||||
} else {
|
||||
await this.acquire()
|
||||
}
|
||||
} catch {
|
||||
// Transient API/network failure: retry. Leadership is only lost after the
|
||||
// renewDeadline elapses without a successful renew, checked in the loop.
|
||||
}
|
||||
this.scheduleNext()
|
||||
}
|
||||
|
||||
private scheduleNext(): void {
|
||||
if (this.timer !== undefined) return
|
||||
// When leader, renew well inside renewDeadline; otherwise back off and retry.
|
||||
const delay = this.leading ? Math.min(this.renewDeadlineSeconds / 2, this.retryPeriodSeconds) : this.retryPeriodSeconds
|
||||
this.timer = setTimeout(() => {
|
||||
this.timer = undefined
|
||||
if (this.leading && this.now() - this.lastRenew > this.renewDeadlineSeconds * 1000) {
|
||||
// Too long without a successful renew → another holder may have taken over.
|
||||
this.leading = false
|
||||
this.onStoppedLeading?.()
|
||||
}
|
||||
void this.tryAcquire()
|
||||
}, delay * 1000)
|
||||
this.timer.unref?.()
|
||||
}
|
||||
|
||||
private lease(now: number, transitions: number): k8s.V1Lease {
|
||||
return {
|
||||
apiVersion: 'coordination.k8s.io/v1',
|
||||
kind: 'Lease',
|
||||
metadata: { name: this.leaseName, namespace: this.namespace },
|
||||
spec: {
|
||||
holderIdentity: this.identity,
|
||||
leaseDurationSeconds: this.leaseDurationSeconds,
|
||||
acquireTime: new k8s.V1MicroTime(now),
|
||||
renewTime: new k8s.V1MicroTime(now),
|
||||
leaseTransitions: transitions,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
private async acquire(): Promise<void> {
|
||||
try {
|
||||
await this.coordination.createNamespacedLease({
|
||||
namespace: this.namespace,
|
||||
body: this.lease(this.now(), 0),
|
||||
})
|
||||
this.becomeLeader(0, this.now())
|
||||
} catch (err) {
|
||||
if ((err as { code?: number }).code !== 409) throw err
|
||||
// Lease exists — take over only if the holder's renew time is stale.
|
||||
const current = await this.coordination.readNamespacedLease({
|
||||
name: this.leaseName,
|
||||
namespace: this.namespace,
|
||||
})
|
||||
const holder = current.spec?.holderIdentity
|
||||
const renew = toMillis(current.spec?.renewTime)
|
||||
if (holder === this.identity) {
|
||||
// We already hold it (e.g. after a restart) — renew, then resume leading.
|
||||
await this.renew()
|
||||
this.becomeLeader(current.spec?.leaseTransitions ?? 0, this.now())
|
||||
return
|
||||
}
|
||||
const expired = this.now() - renew > this.leaseDurationSeconds * 1000
|
||||
if (!expired) return // a live leader holds it; back off
|
||||
const transitions = (current.spec?.leaseTransitions ?? 0) + 1
|
||||
const now = this.now()
|
||||
// `patchNamespacedLease` uses JSON Patch (an array); a full replace with a
|
||||
// resourceVersion gives the optimistic-concurrency takeover we want.
|
||||
await this.coordination.replaceNamespacedLease({
|
||||
name: this.leaseName,
|
||||
namespace: this.namespace,
|
||||
body: {
|
||||
apiVersion: 'coordination.k8s.io/v1',
|
||||
kind: 'Lease',
|
||||
metadata: { name: this.leaseName, namespace: this.namespace, resourceVersion: current.metadata?.resourceVersion },
|
||||
spec: {
|
||||
holderIdentity: this.identity,
|
||||
leaseDurationSeconds: this.leaseDurationSeconds,
|
||||
acquireTime: new k8s.V1MicroTime(now),
|
||||
renewTime: new k8s.V1MicroTime(now),
|
||||
leaseTransitions: transitions,
|
||||
},
|
||||
},
|
||||
})
|
||||
this.becomeLeader(transitions, now)
|
||||
}
|
||||
}
|
||||
|
||||
private async renew(): Promise<void> {
|
||||
const now = this.now()
|
||||
// Always read the latest lease so we preserve holder/acquireTime/transitions
|
||||
// and bump only renewTime.
|
||||
const current = await this.coordination.readNamespacedLease({
|
||||
name: this.leaseName,
|
||||
namespace: this.namespace,
|
||||
})
|
||||
await this.coordination.replaceNamespacedLease({
|
||||
name: this.leaseName,
|
||||
namespace: this.namespace,
|
||||
body: {
|
||||
apiVersion: 'coordination.k8s.io/v1',
|
||||
kind: 'Lease',
|
||||
metadata: { name: this.leaseName, namespace: this.namespace, resourceVersion: current.metadata?.resourceVersion },
|
||||
spec: {
|
||||
holderIdentity: current.spec?.holderIdentity ?? this.identity,
|
||||
leaseDurationSeconds: this.leaseDurationSeconds,
|
||||
acquireTime: toMicroTime(current.spec?.acquireTime),
|
||||
renewTime: new k8s.V1MicroTime(now),
|
||||
leaseTransitions: current.spec?.leaseTransitions ?? 0,
|
||||
},
|
||||
},
|
||||
})
|
||||
this.lastRenew = now
|
||||
}
|
||||
|
||||
private becomeLeader(operationId: number, now: number): void {
|
||||
const wasLeader = this.leading
|
||||
this.leading = true
|
||||
this.operationId = operationId
|
||||
this.lastRenew = now
|
||||
if (!wasLeader) this.onStartedLeading?.({ holder: this.identity, operationId })
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large.
Load diff
@@ -0,0 +1,16 @@
|
||||
/**
|
||||
* cordis patch rendering for a child DSH. Always mounts the runtime plugin
|
||||
* (`dshs/runtime`) so every child injects the watchdog contract,
|
||||
* plus one row per enabled folder plugin (id doubles as package name). The real
|
||||
* harness loads this via `--patch <file>`.
|
||||
* @module dshs/supervisor/patch
|
||||
*/
|
||||
|
||||
/** The runtime plugin patch row, mounted in every child DSH. */
|
||||
const RUNTIME_ROW = ' - id: dshs-runtime\n name: dshs/runtime'
|
||||
|
||||
/** Render a patch YAML always enabling the runtime plugin plus `enabledPlugins`. */
|
||||
export function renderPatch(enabledPlugins: readonly string[]): string {
|
||||
const rows = [RUNTIME_ROW, ...enabledPlugins.map((id) => ` - id: ${id}\n name: ${id}`)]
|
||||
return `- insert:\n${rows.join('\n')}\n`
|
||||
}
|
||||
@@ -0,0 +1,633 @@
|
||||
/**
|
||||
* Reverse proxy from the orchestrator to a running per-user DSH.
|
||||
*
|
||||
* Two entry points:
|
||||
* - subpath `/u/:slug/dsh/*` (authenticated, legacy), and
|
||||
* - per-user subdomain `<username>.<baseDomain>` (HTTP + WebSocket). The DSH's
|
||||
* absolute-path SPA requires the subdomain form: its `/assets/*` and `/api/*`
|
||||
* resolve against the host root, which only works when each DSH owns a host.
|
||||
* @module dshs/supervisor/proxy
|
||||
*/
|
||||
|
||||
import type { FastifyInstance, FastifyReply, FastifyRequest } from 'fastify'
|
||||
import { readFileSync } from 'node:fs'
|
||||
import { dirname, join } from 'node:path'
|
||||
import { fileURLToPath } from 'node:url'
|
||||
import { Agent, request as httpRequest, type IncomingHttpHeaders, type IncomingMessage } from 'node:http'
|
||||
import { connect } from 'node:net'
|
||||
import { hashSessionToken, parseCookie } from '../web/auth.js'
|
||||
import { requireAuth } from '../web/middleware/authn.js'
|
||||
import type { Endpoint } from './spawner.js'
|
||||
|
||||
// Keep-alive pool for per-user DSH upstreams. Replaced (not just destroyed) on a
|
||||
// connection error, because a Pod rebuild changes its IP and any pooled socket
|
||||
// to the old IP would keep failing (docs/k8s.md §5.4).
|
||||
let upstreamAgent = new Agent({ keepAlive: true, maxSockets: 32 })
|
||||
|
||||
// Headers the DSH's browser-trust fence must NOT see from the browser, so the
|
||||
// proxied request looks like a clean loopback client (its Host is overridden to
|
||||
// loopback; a mismatched Origin would otherwise 403).
|
||||
const STRIP_HEADERS = new Set([
|
||||
'origin',
|
||||
'referer',
|
||||
'sec-fetch-site',
|
||||
'sec-fetch-mode',
|
||||
'sec-fetch-dest',
|
||||
'sec-fetch-user',
|
||||
'x-forwarded-for',
|
||||
'x-forwarded-host',
|
||||
'x-forwarded-proto',
|
||||
])
|
||||
|
||||
|
||||
/**
|
||||
* 档案 50:注入到**实例 HTML**里的「会话过期自愈」脚本。
|
||||
*
|
||||
* 为什么需要它:实例被空闲回收后**重建**(新端口 + 新 launch token),而**已打开的页面**
|
||||
* 仍带着旧 `dsh-auth-*` cookie → 之后任何 `/api/*` XHR 都会被新实例判 **401**,
|
||||
* dsh 前端只显示 `transport failure … HTTP 401`,用户只能手动刷新。
|
||||
* 导航路径的 401 已在上面处理(302 到当前实例新 token),但 **XHR 无法靠 302 自愈**,
|
||||
* 所以给页面这段脚本:发现 `/api/*` 返回 401 就显示覆盖层并 reload ——
|
||||
* reload 会走"导航 401 → 302 新 token"这条**已经工作**的链路,从而自动恢复。
|
||||
* 仅改写响应内容,不落盘、不改官方文件(R2);README/技能已留档。
|
||||
*/
|
||||
// 注入到实例子域页面的自愈脚本(档案 50 + 档案 59 增强)。
|
||||
// ① 401 自愈:实例被回收重建后 launch token 轮换,页面内 /api/* 会拿到 401 → 亮覆盖层并整页 reload。
|
||||
// ② 慢请求提示(档案 59):服务端在实例未就绪时 hold 住请求最长 20 秒等拉起,
|
||||
// 期间浏览器端原本毫无反馈(点了重连也看不出在等什么)→ 挂起 ≥3 秒即显示「实例正在启动」。
|
||||
// ③ 回到页面自检 + 就地恢复(档案 77):判据从「进程在不在跑」改为「页面还能不能连上实例」,
|
||||
// 并把恢复过程做成**看得见**的(顶部轻提示 → 覆盖层),恢复后原地跳回而不是离开到门户域。
|
||||
// 保持多行形式便于维护;注入时整段塞进 <script>,故内容不得含反引号 / ${。
|
||||
// 档案 81 · R1-①:从**独立文件**加载(原先是 TS 模板字面量 —— 里面的 \n 会在模板求值时先被转义,
|
||||
// 曾把整段脚本写崩成 SyntaxError,而校验跳过了"求值"这一步,形同虚设)。
|
||||
// 独立文件是纯 JS:可 node --check 直接校验,不再有转义陷阱。
|
||||
/** 档案 81 · R1-①:注入脚本从 `assets/inject/` 读取(纯 JS,可 node --check;无模板转义陷阱)。
|
||||
* 找不到文件 = **启动即失败**(fail-fast):宁可平台起不来,也不要把"没有注入"的页面静默发给用户。 */
|
||||
function loadInject(file: string): string {
|
||||
// ⚠️ 本项目是 ESM(package.json "type": "module")—— **没有 __dirname**;
|
||||
// 用 import.meta.url 定位(2026-09-13 实测:写成 __dirname 会让整个平台起不来)
|
||||
const p = join(dirname(fileURLToPath(import.meta.url)), '..', '..', 'assets', 'inject', file)
|
||||
try {
|
||||
return readFileSync(p, 'utf8')
|
||||
} catch (err) {
|
||||
throw new Error(`[inject] 读取注入脚本失败: ${p}(档案 81 R1:assets/inject 必须随包部署): ${String(err)}`)
|
||||
}
|
||||
}
|
||||
const SESSION_RECOVERY_JS = loadInject('recovery.js')
|
||||
|
||||
/**
|
||||
* 档案 56:注入到实例页面的「实例助手」——右下角两个入口:
|
||||
* ① 我的文件:浏览自己的工作区并**下载**任意文件(走平台 `/api/desktop/tree` + `/api/fs/download`,
|
||||
* 不暴露宿主绝对路径;此前用户拿到的只是 `/var/lib/...` 路径,浏览器打不开);
|
||||
* ② 能力:展示平台生成的实例能力清单(与 `bundled-skills/platform-capabilities` 同源)。
|
||||
* 另外:读取 `/api/dsh/session-permission`,若**老会话仍处于 workspace-write**(沙箱后端不可用 →
|
||||
* bash 会被 fail-closed 拒绝)则顶部给一条可操作的提示(含切换办法)。
|
||||
* 纯前端注入,不改官方包、不落盘(R2);失败静默,绝不影响实例本身。
|
||||
*/
|
||||
// 档案 81 · R1-①:从**独立文件**加载(原先是 TS 模板字面量 —— 里面的 \n 会在模板求值时先被转义,
|
||||
// 曾把整段脚本写崩成 SyntaxError,而校验跳过了"求值"这一步,形同虚设)。
|
||||
// 独立文件是纯 JS:可 node --check 直接校验,不再有转义陷阱。
|
||||
const SESSION_ASSIST_JS = loadInject('assist.js')
|
||||
|
||||
/** 非 HTML / 已压缩 / 无正文的状态直接透传,避免破坏二进制或已编码内容。 */
|
||||
function injectRecovery(html: string): string {
|
||||
const tag =
|
||||
'<script>' + SESSION_RECOVERY_JS + '</script>' + '<script>' + SESSION_ASSIST_JS + '</script>'
|
||||
return html.includes('</body>') ? html.replace('</body>', tag + '</body>') : html + tag
|
||||
}
|
||||
|
||||
function buildUpstreamHeaders(headers: IncomingHttpHeaders, port: number): Record<string, string | string[]> {
|
||||
const out: Record<string, string | string[]> = {}
|
||||
for (const [key, value] of Object.entries(headers)) {
|
||||
if (value === undefined || STRIP_HEADERS.has(key.toLowerCase())) continue
|
||||
out[key] = value as string | string[]
|
||||
}
|
||||
// Keep Host loopback: DSH's /api trust fence requires the Host to be loopback
|
||||
// or a `--trusted-host` authority — a real domain would 403 every /api call.
|
||||
// DSH's absolute URLs are rewritten to the real origin in proxyHttp below.
|
||||
out.host = `127.0.0.1:${port}`
|
||||
return out
|
||||
}
|
||||
|
||||
/** 档案 51:从 freshAuthUrl 的 `?token=` 中取出当前实例的 launch token。 */
|
||||
function tokenFromAuthUrl(url: string | undefined): string | undefined {
|
||||
if (url === undefined) return undefined
|
||||
const m = /[?&]token=([^&]+)/.exec(url)
|
||||
return m === null || m[1] === undefined ? undefined : decodeURIComponent(m[1])
|
||||
}
|
||||
|
||||
/**
|
||||
* 档案 51:向**当前**实例要一份有效的浏览器凭证 cookie。
|
||||
*
|
||||
* 实例每次 (重)启动都会换 launch token **和** `dsh-auth-<随机后缀>` 的 cookie 名。
|
||||
* 已打开的页面还带着旧名字的 cookie → 新实例一律判 401。这里 GET `/?token=<当前>`
|
||||
* (实例回 303 + Set-Cookie),把 set-cookie 原样取出,供代理重放请求与回写浏览器。
|
||||
*/
|
||||
function fetchInstanceAuthCookies(endpoint: Endpoint, token: string): Promise<string[]> {
|
||||
return new Promise((resolve) => {
|
||||
let done = false
|
||||
const finish = (v: string[]): void => {
|
||||
if (!done) {
|
||||
done = true
|
||||
resolve(v)
|
||||
}
|
||||
}
|
||||
const req = httpRequest(
|
||||
{
|
||||
host: endpoint.host,
|
||||
port: endpoint.port,
|
||||
path: '/?token=' + encodeURIComponent(token),
|
||||
method: 'GET',
|
||||
headers: { host: `127.0.0.1:${endpoint.port}`, 'user-agent': 'dsh-proxy-auth-refresh' },
|
||||
},
|
||||
(res) => {
|
||||
const sc = res.headers['set-cookie']
|
||||
res.resume()
|
||||
finish(Array.isArray(sc) ? sc : sc === undefined ? [] : [sc])
|
||||
},
|
||||
)
|
||||
req.on('error', () => finish([]))
|
||||
req.setTimeout(5000, () => {
|
||||
req.destroy()
|
||||
finish([])
|
||||
})
|
||||
req.end()
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* 档案 51:把浏览器带来的 cookie 与实例的新 cookie 合并 —— 丢掉旧的 `dsh-auth-*`
|
||||
* (新实例只认自己那份;旧名留着也没用),其余(`sid` 等代理不关心)原样保留。
|
||||
*/
|
||||
function mergeCookieHeader(original: string | undefined, fresh: string[]): string {
|
||||
const keep = (original ?? '')
|
||||
.split(';')
|
||||
.map((x) => x.trim())
|
||||
.filter((x) => x !== '' && !/^dsh-auth-/.test(x))
|
||||
const add = fresh
|
||||
.map((c) => (c.split(';')[0] ?? '').trim())
|
||||
.filter((x) => x !== '')
|
||||
return [...keep, ...add].join('; ')
|
||||
}
|
||||
|
||||
/** Extract `<slug>` from `<slug>.<baseDomain>`, or null when not a match. */
|
||||
export function parseSubdomain(host: string | undefined, baseDomain: string): string | null {
|
||||
if (host === undefined || baseDomain === '') return null
|
||||
const name = host.split(':')[0] ?? ''
|
||||
if (name === baseDomain) return null
|
||||
if (name.endsWith('.' + baseDomain)) {
|
||||
const slug = name.slice(0, -(baseDomain.length + 1))
|
||||
return slug !== '' ? slug.toLowerCase() : null
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
/** The per-user subdomain for a username, or null when `baseDomain` is unset. */
|
||||
export function subdomainForUser(baseDomain: string, username: string): string | null {
|
||||
return baseDomain === '' ? null : `${username.toLowerCase()}.${baseDomain}`
|
||||
}
|
||||
|
||||
/** The browser-facing origin (`<scheme>://<host>`) this request reached us with,
|
||||
* used to rewrite DSH's loopback absolute URLs back to the real domain. */
|
||||
function realOrigin(headers: IncomingHttpHeaders): string | undefined {
|
||||
const host = clientHost(headers)
|
||||
if (host === undefined) return undefined
|
||||
const proto = headers['x-forwarded-proto']
|
||||
const scheme = typeof proto === 'string' && proto !== '' ? proto : 'https'
|
||||
return `${scheme}://${host}`
|
||||
}
|
||||
|
||||
/** A subdomain resolution: a tunnelable endpoint, an error to return, or null (not a subdomain). */
|
||||
type SubdomainAccess = { endpoint: Endpoint; userId: string } | { error: string; code: number; userId?: string } | null
|
||||
|
||||
/**
|
||||
* The client-facing host, preferring `X-Forwarded-Host`. The control plane sits
|
||||
* behind an edge proxy (Tencent nginx) that hides the real Host to sidestep the
|
||||
* cloud provider's ICP check (docs/k8s-deploy.md §7.4); the real domain arrives
|
||||
* here. The subdomain auth check still validates the cookie against the slug, so
|
||||
* a spoofed forwarded host cannot reach another user's DSH.
|
||||
*/
|
||||
function clientHost(headers: IncomingHttpHeaders): string | undefined {
|
||||
const fwd = headers['x-forwarded-host']
|
||||
if (typeof fwd === 'string' && fwd !== '') return fwd
|
||||
if (Array.isArray(fwd) && fwd[0] !== '') return fwd[0]
|
||||
return headers.host
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a subdomain Host to a DSH port, authenticating the caller: the session
|
||||
* cookie must belong to a non-disabled user whose username matches the subdomain.
|
||||
*/
|
||||
async function resolveSubdomainAccess(
|
||||
app: FastifyInstance,
|
||||
host: string | undefined,
|
||||
cookieHeader: string | undefined,
|
||||
): Promise<SubdomainAccess> {
|
||||
const slug = parseSubdomain(host, app.config.baseDomain)
|
||||
if (slug === null) return null
|
||||
const target = await app.db.findUserBySlug(slug)
|
||||
if (target === undefined) return { error: 'unknown_user', code: 404 }
|
||||
const token = parseCookie(cookieHeader, 'sid')
|
||||
const session = token === undefined ? undefined : await app.db.findSessionWithUser(hashSessionToken(token))
|
||||
if (session === undefined || session.expiresAt <= Date.now() || session.user.role === 'disabled') {
|
||||
return { error: 'unauthorized', code: 401 }
|
||||
}
|
||||
if (session.user.username.toLowerCase() !== slug) {
|
||||
return { error: 'forbidden', code: 403 }
|
||||
}
|
||||
const endpoint = await app.supervisor.endpointFor(session.user.id)
|
||||
if (endpoint === undefined) return { error: 'not_running', code: 404, userId: session.user.id }
|
||||
// userId 一并返回:实例侧 401(launch token 过期)时用它取当前实例的新 token(档案 24)。
|
||||
// Real user traffic to their own DSH counts as activity (idle-reap signal).
|
||||
app.supervisor.touch(session.user.id)
|
||||
return { endpoint, userId: session.user.id }
|
||||
}
|
||||
|
||||
function proxyHttp(
|
||||
request: FastifyRequest,
|
||||
reply: FastifyReply,
|
||||
endpoint: Endpoint,
|
||||
targetPath: string,
|
||||
rewritePrefix?: string,
|
||||
rewriteLoopbackLocation = false,
|
||||
useKeepAlive = false,
|
||||
/**
|
||||
* 实例侧 401 时的恢复入口(档案 24):浏览器导航遇到 dsh 的 "authentication required"
|
||||
* (launch token 过期 —— 实例重启/回收后刷新旧标签页)时调用,返回应 302 到的地址。
|
||||
* 返回 undefined 则回落到 '/'。
|
||||
*/
|
||||
freshAuthUrl?: () => Promise<string | undefined>,
|
||||
): void {
|
||||
reply.hijack()
|
||||
// 档案 51:为「401 透明重放」准备请求体。
|
||||
// 只在 content-length 已知且不大时才缓冲(普通 /api JSON-RPC 请求都很小);
|
||||
// 未知长度(chunked 上传)不缓冲 → 该请求退回旧行为(401 透传),不做重放。
|
||||
const MAX_REPLAY_BODY = 8 * 1024 * 1024
|
||||
const reqMethod = (request.raw.method ?? 'GET').toUpperCase()
|
||||
const mayHaveBody = reqMethod !== 'GET' && reqMethod !== 'HEAD'
|
||||
const clNum =
|
||||
typeof request.headers['content-length'] === 'string' ? Number(request.headers['content-length']) : NaN
|
||||
const bufferable = mayHaveBody && Number.isFinite(clNum) && clNum <= MAX_REPLAY_BODY
|
||||
let bodyBuf: Buffer | undefined
|
||||
// ★ 事后修正 1(2026-09-11):重放闸门**不能**用 bufferable —— 它含 mayHaveBody,
|
||||
// 会让 GET/HEAD(含 dsh 的 SSE 会话流)永远无法重放。GET 没有请求体,反而最该能重放。
|
||||
const canReplay = mayHaveBody ? false : true
|
||||
let authRetryUsed = false
|
||||
let replayCookies: string[] = []
|
||||
|
||||
const attempt = (connRetry: boolean, authCookie?: string): void => {
|
||||
// 档案 50 修正(实测定位):**HTML 导航请求必须向上游要 identity**。
|
||||
// 原因:dsh 实例对带 Accept-Encoding 的请求会 **gzip 压缩 HTML**(浏览器就带),
|
||||
// 于是响应带 content-encoding → 我下面的注入逻辑**主动跳过** → 浏览器拿到的页面
|
||||
// **没有自愈脚本** → C 方案对真实用户无效(而 curl 不带 Accept-Encoding,故此前验证
|
||||
// 是假阳性)。改法:只要**客户端接受 HTML**,就覆盖转发的 accept-encoding 为 identity,
|
||||
// 让上游回未压缩 HTML → 注入成功。体积影响可忽略(HTML ~24KB),且边缘 nginx 若开 gzip
|
||||
// 仍会对浏览器压缩,用户侧无感知。静态资源/SSE 的压缩行为不受影响。
|
||||
const upHeaders = buildUpstreamHeaders(request.headers, endpoint.port)
|
||||
if (String(request.headers.accept ?? '').includes('text/html')) {
|
||||
upHeaders['accept-encoding'] = 'identity'
|
||||
}
|
||||
// 档案 51:重放时用**当前实例的新 cookie** 覆盖浏览器带来的旧 cookie。
|
||||
if (authCookie !== undefined) upHeaders.cookie = authCookie
|
||||
if (bodyBuf !== undefined) {
|
||||
// 请求体已缓冲 → 显式带上长度并去掉可能存在的 chunked 头,保证重放一致。
|
||||
delete upHeaders['transfer-encoding']
|
||||
upHeaders['content-length'] = String(bodyBuf.length)
|
||||
}
|
||||
const upstream = httpRequest(
|
||||
{
|
||||
host: endpoint.host,
|
||||
port: endpoint.port,
|
||||
path: targetPath,
|
||||
method: request.method,
|
||||
// Keep-alive only where it is required (k8s Headless Service DNS
|
||||
// re-resolution on retry — docs/k8s.md §5.4). Local mode uses a fresh
|
||||
// connection per request: the loopback connect cost is negligible and
|
||||
// this avoids the half-open keep-alive pool that hung the proxy after
|
||||
// child-instance restarts (2026-09-09 outage: pooled sockets to a
|
||||
// since-recycled instance port never errored, so requests hung).
|
||||
agent: useKeepAlive && !connRetry ? upstreamAgent : false,
|
||||
// Host header stays loopback for the DSH trust fence; the TCP target host
|
||||
// is endpoint.host above. k8s mode overrides the header port separately.
|
||||
headers: upHeaders,
|
||||
},
|
||||
(upRes: IncomingMessage) => {
|
||||
const headers = { ...upRes.headers }
|
||||
const isNav =
|
||||
request.raw.method === 'GET' && String(request.headers.accept ?? '').includes('text/html')
|
||||
// 档案 51:**非导航**请求的实例侧 401 = 页面拿着旧实例的 dsh-auth cookie。
|
||||
// 透明重放:取当前实例新 cookie → 覆盖 cookie 头重发同一请求 → 新 cookie 回写浏览器。
|
||||
// 只重放一次(authRetryUsed),避免与实例互相刷 401 造成死循环。
|
||||
const replayReady = mayHaveBody ? bodyBuf !== undefined : canReplay
|
||||
if (upRes.statusCode === 401 && !isNav && !authRetryUsed && freshAuthUrl !== undefined && replayReady) {
|
||||
authRetryUsed = true
|
||||
upRes.resume()
|
||||
void (async () => {
|
||||
let cookieHeader: string | undefined
|
||||
try {
|
||||
const url = await freshAuthUrl()
|
||||
const tk = tokenFromAuthUrl(url)
|
||||
if (tk !== undefined) {
|
||||
const fresh = await fetchInstanceAuthCookies(endpoint, tk)
|
||||
if (fresh.length > 0) {
|
||||
replayCookies = fresh
|
||||
cookieHeader = mergeCookieHeader(request.headers.cookie, fresh)
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
/* 取不到 → 回落到原样 401 */
|
||||
}
|
||||
if (cookieHeader === undefined) {
|
||||
reply.raw.writeHead(401, headers)
|
||||
reply.raw.end('unauthorized')
|
||||
return
|
||||
}
|
||||
process.stderr.write(
|
||||
`[proxy-auth-replay] ${request.raw.method ?? 'GET'} ${targetPath}(旧 cookie → 实例新 cookie,重放)\n`,
|
||||
)
|
||||
attempt(false, cookieHeader)
|
||||
})()
|
||||
return
|
||||
}
|
||||
// 重放成功后把新 cookie 交给浏览器:之后(含 SSE 自动重连)不再需要重放。
|
||||
if (authCookie !== undefined && replayCookies.length > 0) {
|
||||
const prev = headers['set-cookie']
|
||||
headers['set-cookie'] = [
|
||||
...(Array.isArray(prev) ? prev : prev === undefined ? [] : [prev]),
|
||||
...replayCookies,
|
||||
]
|
||||
}
|
||||
const location = upRes.headers.location
|
||||
if (
|
||||
rewritePrefix !== undefined &&
|
||||
typeof location === 'string' &&
|
||||
location.startsWith('/') &&
|
||||
!location.startsWith('//') &&
|
||||
!location.startsWith(rewritePrefix)
|
||||
) {
|
||||
headers.location = rewritePrefix + location
|
||||
}
|
||||
// ── 2026-09-12(平台缓存治理):让「改了 client bundle 却看不到变化」不再发生 ──
|
||||
//
|
||||
// 背景(档案 67 v0.2.6 实测,排查成本极高):
|
||||
// dsh 的客户端模块 URL 形如 `/plugins/??<pkg>/client.js,…&rev=<内容 sha1>`,
|
||||
// `rev` 由 **dsh 官方**按内容计算(`dsh-client-modules`),**平台不得改(R2)**。
|
||||
// 平台更新 bundle 后,若 `rev` 未变 → 浏览器认为手上那份缓存仍有效 → **不重新请求**
|
||||
// ⇒ 服务端已是新版、用户却一直看到旧 UI(本次:SQL 层/磁盘/服务端三处都验过是新的)。
|
||||
//
|
||||
// 处置:对实例的**模块/静态资源路径**统一加 `Cache-Control: no-cache` ——
|
||||
// **允许缓存,但每次必须回源校验**(配合上游 ETag/Last-Modified 走 304,开销极小)。
|
||||
// ⇒ bundle 一变,浏览器下一次请求就拿到新的,**用户无需清缓存/换无痕**。
|
||||
//
|
||||
// ⚠️ 勿删:这是 UI 改动能否被验收的前置机制(`06-工作台UI规范 §6.5`)。
|
||||
if (targetPath.startsWith('/plugins/') || targetPath.startsWith('/assets/')) {
|
||||
headers['cache-control'] = 'no-cache'
|
||||
}
|
||||
// local mode only: DSH builds absolute URLs from the loopback Host we
|
||||
// forward; rewrite any 127.0.0.1 Location to the real origin so the
|
||||
// browser doesn't jump to the user's own machine. k8s mode keeps its
|
||||
// verified behavior (no rewrite).
|
||||
if (
|
||||
rewriteLoopbackLocation &&
|
||||
typeof location === 'string' &&
|
||||
realOrigin(request.headers) !== undefined
|
||||
) {
|
||||
headers.location = location.replace(/^https?:\/\/127\.0\.0\.1(:\d+)?/, realOrigin(request.headers)!)
|
||||
}
|
||||
// 2026-09-11(档案 24):实例侧 401 = launch token 过期。浏览器导航时不要停在
|
||||
// dsh 的 "dsh web authentication required; reopen the URL printed by dsh web." 死端页,
|
||||
// 而是用当前实例的新 token 302 回同一地址(拿不到则回门户)。
|
||||
if (
|
||||
upRes.statusCode === 401 &&
|
||||
request.raw.method === 'GET' &&
|
||||
String(request.headers.accept ?? '').includes('text/html') &&
|
||||
freshAuthUrl !== undefined
|
||||
) {
|
||||
upRes.resume()
|
||||
void freshAuthUrl()
|
||||
.then((url) => {
|
||||
reply.raw.writeHead(302, { location: url ?? '/' })
|
||||
reply.raw.end()
|
||||
})
|
||||
.catch(() => {
|
||||
reply.raw.writeHead(302, { location: '/' })
|
||||
reply.raw.end()
|
||||
})
|
||||
return
|
||||
}
|
||||
// 档案 50:HTML 响应缓冲后注入「会话过期自愈」脚本(其余一切原样 pipe,零影响)。
|
||||
const ctype = String(upRes.headers['content-type'] ?? '')
|
||||
const enc = upRes.headers['content-encoding']
|
||||
const status = upRes.statusCode ?? 502
|
||||
const canInject =
|
||||
ctype.includes('text/html') && enc === undefined && status !== 204 && status !== 304
|
||||
if (!canInject) {
|
||||
if (ctype.includes('text/html') && enc !== undefined) {
|
||||
process.stderr.write(`[inject-recovery] skip: content-encoding=${String(enc)}\n`)
|
||||
}
|
||||
reply.raw.writeHead(status, headers)
|
||||
upRes.pipe(reply.raw)
|
||||
return
|
||||
}
|
||||
const chunks: Buffer[] = []
|
||||
upRes.on('data', (c: Buffer) => chunks.push(c))
|
||||
upRes.on('end', () => {
|
||||
const out = injectRecovery(Buffer.concat(chunks).toString('utf8'))
|
||||
delete headers['content-length']
|
||||
delete headers['transfer-encoding']
|
||||
headers['content-length'] = String(Buffer.byteLength(out))
|
||||
reply.raw.writeHead(status, headers)
|
||||
reply.raw.end(out)
|
||||
})
|
||||
upRes.on('error', () => reply.raw.destroy())
|
||||
},
|
||||
)
|
||||
upstream.on('error', () => {
|
||||
if (connRetry) {
|
||||
reply.raw.destroy()
|
||||
return
|
||||
}
|
||||
// A stale keep-alive socket or a Pod that just restarted: drop the pool,
|
||||
// replace it with a fresh one, and retry once on a fresh connection.
|
||||
upstreamAgent.destroy()
|
||||
upstreamAgent = new Agent({ keepAlive: true, maxSockets: 32 })
|
||||
request.raw.unpipe(upstream)
|
||||
attempt(true)
|
||||
})
|
||||
if (bodyBuf !== undefined) upstream.end(bodyBuf)
|
||||
else request.raw.pipe(upstream)
|
||||
}
|
||||
if (!bufferable) {
|
||||
attempt(false)
|
||||
return
|
||||
}
|
||||
const pre: Buffer[] = []
|
||||
request.raw.on('data', (c: Buffer) => pre.push(c))
|
||||
request.raw.on('end', () => {
|
||||
bodyBuf = Buffer.concat(pre)
|
||||
attempt(false)
|
||||
})
|
||||
request.raw.on('error', () => reply.raw.destroy())
|
||||
}
|
||||
|
||||
export async function registerDshProxy(app: FastifyInstance): Promise<void> {
|
||||
// Legacy authenticated subpath proxy.
|
||||
app.all('/u/:slug/dsh/*', { preHandler: requireAuth }, async (request, reply) => {
|
||||
const slug = (request.params as { slug: string }).slug
|
||||
if (request.user === null || request.user.id !== slug) {
|
||||
reply.code(403).send({ error: 'forbidden' })
|
||||
return
|
||||
}
|
||||
const endpoint = await app.supervisor.endpointFor(slug)
|
||||
if (endpoint === undefined) {
|
||||
reply.code(404).send({ error: 'not_running' })
|
||||
return
|
||||
}
|
||||
app.supervisor.touch(request.user!.id)
|
||||
const prefix = `/u/${slug}/dsh`
|
||||
const rawUrl = request.raw.url ?? '/'
|
||||
const targetPath = rawUrl.startsWith(prefix) ? rawUrl.slice(prefix.length) || '/' : rawUrl
|
||||
const pathScheme = app.config.secureCookies ? 'https' : 'http'
|
||||
proxyHttp(
|
||||
request,
|
||||
reply,
|
||||
endpoint,
|
||||
targetPath,
|
||||
prefix,
|
||||
app.config.deployMode === 'local',
|
||||
app.config.deployMode === 'k8s',
|
||||
async () => {
|
||||
const status = await app.supervisor.status(request.user!.id)
|
||||
const token = status.main?.launchToken
|
||||
return token !== undefined && token !== ''
|
||||
? `${pathScheme}://${app.config.baseDomain}${prefix}/?token=${encodeURIComponent(token)}`
|
||||
: `${pathScheme}://${app.config.baseDomain}/`
|
||||
},
|
||||
)
|
||||
})
|
||||
|
||||
// Per-user subdomain: HTTP (intercept before normal routing).
|
||||
app.addHook('onRequest', async (request, reply) => {
|
||||
const access = await resolveSubdomainAccess(app, clientHost(request.headers), request.headers.cookie)
|
||||
if (access === null) return
|
||||
if ('error' in access) {
|
||||
// 浏览器导航(GET + 期待 HTML)的 401(会话失效/退出后刷新实例子域)→ 302 回门户
|
||||
// 登录页,而不是吐 JSON {"error":"unauthorized"} 停在原地。
|
||||
if (
|
||||
access.code === 401 &&
|
||||
request.raw.method === 'GET' &&
|
||||
(request.headers.accept ?? '').includes('text/html')
|
||||
) {
|
||||
const scheme = app.config.secureCookies ? 'https' : 'http'
|
||||
reply.redirect(`${scheme}://${app.config.baseDomain}/login.html`)
|
||||
return
|
||||
}
|
||||
// 实例未运行(后台回收/崩溃/停止/熔断后刷新会话页)→ 需要拉起。
|
||||
//
|
||||
// 2026-09-11(档案 49 · 体验修复):**浏览器导航一律"立刻" 302 到平台自有的过渡页**,
|
||||
// 由过渡页显示加载动画并自己调 /api/dsh/enter 完成「拉起 + 等 token + 跳回」。
|
||||
// 原因:此前导航请求会在**服务端阻塞等待 launch token(最长 20 秒)**,这期间浏览器
|
||||
// 只看到白屏、没有任何反馈 → 用户以为卡死,只能手动刷新(那时实例已就绪,看似"刷新才好")。
|
||||
// 现在导航分支**不阻塞**(0.2s 内返回 302),动画立刻出现,且**不会**在这里触发 launch
|
||||
//(拉起交给过渡页,单点、可重试)。
|
||||
if (access.code === 404 && access.error === 'not_running' && access.userId !== undefined) {
|
||||
const navScheme = app.config.secureCookies ? 'https' : 'http'
|
||||
const navHost = clientHost(request.headers) ?? app.config.baseDomain
|
||||
const isNavigation =
|
||||
request.raw.method === 'GET' && (request.headers.accept ?? '').includes('text/html')
|
||||
if (isNavigation) {
|
||||
const next = `${navScheme}://${navHost}${request.raw.url ?? '/'}`
|
||||
reply.redirect(
|
||||
`${navScheme}://${app.config.baseDomain}/wake.html?next=${encodeURIComponent(next)}`,
|
||||
)
|
||||
return
|
||||
}
|
||||
// 非导航(XHR / API)——**这是"页面已打开、再对话没反应"的主场景**。
|
||||
// 正解不是"回一个错误/动画",而是**等实例就绪后继续转发**:请求最终成功,
|
||||
// dsh 前端自己会保持它的 loading 态 → 用户无感,不需要刷新。
|
||||
try {
|
||||
const folderAbs = app.userFs.resolvePath(access.userId, '')
|
||||
try {
|
||||
await app.supervisor.launch(access.userId, folderAbs, undefined)
|
||||
} catch {
|
||||
// 并发进场(另一请求正在拉起 → AlreadyRunningError):等它拿到 token
|
||||
await app.supervisor.waitForLaunchTokenForUser(access.userId, 20000)
|
||||
}
|
||||
} catch {
|
||||
/* 拉起失败 → 落到下面的 503 */
|
||||
}
|
||||
const again = await resolveSubdomainAccess(app, clientHost(request.headers), request.headers.cookie)
|
||||
if (again !== null && !('error' in again)) {
|
||||
proxyHttp(
|
||||
request,
|
||||
reply,
|
||||
again.endpoint,
|
||||
request.raw.url ?? '/',
|
||||
undefined,
|
||||
app.config.deployMode === 'local',
|
||||
app.config.deployMode === 'k8s',
|
||||
async () => {
|
||||
const st = await app.supervisor.status(again.userId)
|
||||
const tk = st.main?.launchToken
|
||||
return tk !== undefined && tk !== ''
|
||||
? `${navScheme}://${navHost}/?token=${encodeURIComponent(tk)}`
|
||||
: `${navScheme}://${app.config.baseDomain}/`
|
||||
},
|
||||
)
|
||||
return
|
||||
}
|
||||
// 真的起不来:给明确的「启动中」+ Retry-After,让前端自行退避重试。
|
||||
reply.code(503).header('retry-after', '3').send({ error: 'instance_starting' })
|
||||
return
|
||||
}
|
||||
reply.code(access.code).send({ error: access.error })
|
||||
return
|
||||
}
|
||||
const navScheme = app.config.secureCookies ? 'https' : 'http'
|
||||
const navHost = clientHost(request.headers) ?? app.config.baseDomain
|
||||
proxyHttp(
|
||||
request,
|
||||
reply,
|
||||
access.endpoint,
|
||||
request.raw.url ?? '/',
|
||||
undefined,
|
||||
app.config.deployMode === 'local',
|
||||
app.config.deployMode === 'k8s',
|
||||
async () => {
|
||||
const status = await app.supervisor.status(access.userId)
|
||||
const token = status.main?.launchToken
|
||||
return token !== undefined && token !== ''
|
||||
? `${navScheme}://${navHost}/?token=${encodeURIComponent(token)}`
|
||||
: `${navScheme}://${app.config.baseDomain}/`
|
||||
},
|
||||
)
|
||||
})
|
||||
|
||||
// Per-user subdomain: WebSocket upgrade tunnel. Auth is async (DB lookup), so
|
||||
// the raw `upgrade` callback defers to an async IIFE before deciding to tunnel.
|
||||
app.server.on('upgrade', (req, socket, head) => {
|
||||
void (async () => {
|
||||
const access = await resolveSubdomainAccess(app, clientHost(req.headers), req.headers.cookie)
|
||||
if (access === null || 'error' in access) {
|
||||
socket.destroy()
|
||||
return
|
||||
}
|
||||
const upstream = connect({ host: access.endpoint.host, port: access.endpoint.port })
|
||||
upstream.on('connect', () => {
|
||||
const lines = [`${req.method} ${req.url} HTTP/${req.httpVersion}`]
|
||||
for (let i = 0; i < req.rawHeaders.length; i += 2) {
|
||||
const name = (req.rawHeaders[i] ?? '').toLowerCase()
|
||||
if (name === 'host' || STRIP_HEADERS.has(name)) continue
|
||||
lines.push(`${req.rawHeaders[i]}: ${req.rawHeaders[i + 1]}`)
|
||||
}
|
||||
lines.push(`Host: 127.0.0.1:${access.endpoint.port}`)
|
||||
upstream.write(lines.join('\r\n') + '\r\n\r\n')
|
||||
if (head !== undefined && head.length > 0) upstream.write(head)
|
||||
socket.pipe(upstream)
|
||||
upstream.pipe(socket)
|
||||
})
|
||||
upstream.on('error', () => socket.destroy())
|
||||
socket.on('error', () => upstream.destroy())
|
||||
})()
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
/**
|
||||
* Leader-only controller: reconciles the cluster against the desired state in
|
||||
* `dsh_instances` (docs/k8s.md §5.7) and watches the main DSH Pods for crashes.
|
||||
*
|
||||
* The k8s backend has no child-process `exit` event the way the local backend
|
||||
* does — a crashed Pod is observed either by the informer (here) or by the next
|
||||
* reconcile tick (a desired main whose Pod is gone). Both funnel into the same
|
||||
* repair path: mark the instance crashed, pull up the one-shot watchdog Job,
|
||||
* and let the reconcile relaunch the main if its Pod stays absent.
|
||||
* @module dshs/supervisor/reconcile
|
||||
*/
|
||||
|
||||
import { makeInformer, type Informer } from '@kubernetes/client-node'
|
||||
import type { DbAdapter } from '../db/adapter.js'
|
||||
import type { DshInstance } from '../db/types.js'
|
||||
import type { LeaderElector } from './leader.js'
|
||||
import type { K8sSpawner } from './k8s-spawner.js'
|
||||
import type { LivePod } from './spawner.js'
|
||||
|
||||
/** The result of diffing desired state against live Pods. Pure and testable. */
|
||||
export interface ReconcilePlan {
|
||||
/** Desired mains whose Pod is missing → relaunch. */
|
||||
launch: DshInstance[]
|
||||
/** Live Pods with no desired row → delete (orphans). */
|
||||
delete: LivePod[]
|
||||
}
|
||||
|
||||
/** Diff desired mains vs live Pods. */
|
||||
export function planReconcile(desired: DshInstance[], live: LivePod[]): ReconcilePlan {
|
||||
const liveByUser = new Map(live.map((pod) => [pod.userId, pod]))
|
||||
const launch: DshInstance[] = []
|
||||
for (const instance of desired) {
|
||||
if (instance.role !== 'main') continue
|
||||
const pod = liveByUser.get(instance.userId)
|
||||
if (pod === undefined || !pod.running) launch.push(instance)
|
||||
}
|
||||
const desiredUsers = new Set(desired.filter((i) => i.role === 'main').map((i) => i.userId))
|
||||
const del = live.filter((pod) => !desiredUsers.has(pod.userId))
|
||||
return { launch, delete: del }
|
||||
}
|
||||
|
||||
/** Loop that reconciles only while leader, plus a Pod informer for crashes. */
|
||||
export class ReconcileController {
|
||||
private informer: Informer<object> | undefined
|
||||
private timer: NodeJS.Timeout | undefined
|
||||
private leading = false
|
||||
private readonly watchdogFired = new Set<string>()
|
||||
private readonly intervalMs: number
|
||||
|
||||
constructor(
|
||||
private readonly db: DbAdapter,
|
||||
private readonly spawner: K8sSpawner,
|
||||
private readonly elector: LeaderElector,
|
||||
intervalMs = 10_000,
|
||||
) {
|
||||
this.intervalMs = intervalMs
|
||||
elector.setLeadershipCallbacks(
|
||||
(fencing) => {
|
||||
this.leading = true
|
||||
this.spawner.setFencing(fencing)
|
||||
this.startInformer()
|
||||
this.scheduleTick(0)
|
||||
},
|
||||
() => {
|
||||
this.leading = false
|
||||
this.stopInformer()
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
await this.elector.start()
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
this.elector.stop()
|
||||
this.stopInformer()
|
||||
if (this.timer !== undefined) {
|
||||
clearTimeout(this.timer)
|
||||
this.timer = undefined
|
||||
}
|
||||
}
|
||||
|
||||
private startInformer(): void {
|
||||
if (this.informer !== undefined) return
|
||||
this.informer = this.spawner.watchMainPods((pod) => this.onPodEvent(pod))
|
||||
void this.informer.start().catch(() => {
|
||||
// The informer is best-effort; the reconcile tick still converges.
|
||||
this.informer = undefined
|
||||
})
|
||||
}
|
||||
|
||||
private stopInformer(): void {
|
||||
void this.informer?.stop().catch(() => {})
|
||||
this.informer = undefined
|
||||
}
|
||||
|
||||
private scheduleTick(delayMs: number): void {
|
||||
if (this.timer !== undefined) return
|
||||
this.timer = setTimeout(() => {
|
||||
this.timer = undefined
|
||||
if (!this.leading) return
|
||||
void this.tick().finally(() => this.scheduleTick(this.intervalMs))
|
||||
}, delayMs)
|
||||
this.timer.unref?.()
|
||||
}
|
||||
|
||||
private async tick(): Promise<void> {
|
||||
const [desired, live] = await Promise.all([
|
||||
this.db.listInstancesByRole('main'),
|
||||
this.spawner.listUserPods(),
|
||||
])
|
||||
const plan = planReconcile(desired, live)
|
||||
|
||||
for (const orphan of plan.delete) {
|
||||
// Orphan = a main Pod with no desired row (user deleted or disabled).
|
||||
await this.safe(`orphan ${orphan.userId}`, async () => {
|
||||
await this.db.deleteInstance(orphan.name)
|
||||
await this.spawner.stop(orphan.userId)
|
||||
})
|
||||
}
|
||||
for (const instance of plan.launch) {
|
||||
await this.safe(`launch ${instance.userId}`, () =>
|
||||
this.spawner.launch(instance.userId, instance.folder ?? '', instance.patch ?? undefined))
|
||||
}
|
||||
// Recreate a lost Service/NetworkPolicy for any desired main that still has a Pod.
|
||||
for (const instance of desired) {
|
||||
if (plan.launch.includes(instance)) continue
|
||||
await this.safe(`ensure ${instance.userId}`, () => this.spawner.ensureUserResources(instance.userId))
|
||||
}
|
||||
// Idle reap (Phase 4): a desired main whose user has no active session has
|
||||
// outlived its session TTL — stop the Pod and drop the desired row so the
|
||||
// next tick does not relaunch it.
|
||||
for (const instance of desired) {
|
||||
await this.safe(`reap ${instance.userId}`, async () => {
|
||||
if (await this.db.hasActiveSession(instance.userId)) return
|
||||
await this.spawner.stop(instance.userId)
|
||||
await this.db.deleteInstance(instance.id)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/** Run one per-user step without letting its failure abort the rest of the tick. */
|
||||
private async safe(label: string, fn: () => Promise<unknown>): Promise<void> {
|
||||
try {
|
||||
await fn()
|
||||
} catch (err) {
|
||||
// One bad user (e.g. a files Pod that can't become Ready) must not block
|
||||
// launches/idle-reap for every other user.
|
||||
this.spawner.logError?.(new Error(`reconcile ${label}: ${err instanceof Error ? err.message : String(err)}`))
|
||||
}
|
||||
}
|
||||
|
||||
private async onPodEvent(pod: LivePod): Promise<void> {
|
||||
if (!pod.crashed || !pod.running) {
|
||||
// A healthy transition clears the "already fired" guard so a later crash
|
||||
// triggers the watchdog again.
|
||||
if (pod.running && !pod.crashed) this.watchdogFired.delete(pod.userId)
|
||||
return
|
||||
}
|
||||
if (this.watchdogFired.has(pod.userId)) return
|
||||
this.watchdogFired.add(pod.userId)
|
||||
try {
|
||||
await this.spawner.spawnWatchdog(pod.userId)
|
||||
} catch (err) {
|
||||
this.watchdogFired.delete(pod.userId)
|
||||
this.spawner.logError?.(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,118 @@
|
||||
/**
|
||||
* 会话权限档位读取(档案 56)。
|
||||
*
|
||||
* 背景:档位由 `DSH_PERMISSION_MODE` 决定,但 **dsh 在"会话创建时"把它播种进会话**,
|
||||
* 之后平台改默认值**不会**更新既有会话 → 老会话仍停在 `workspace-write`,而本机沙箱
|
||||
* 后端不可用(内核无 Landlock、bwrap 在平台合成根内探测失败)→ dsh **fail-closed 拒绝
|
||||
* 任何 shell**。用户只会看到「bash 不可用」,无从判断原因。
|
||||
*
|
||||
* 本模块从会话事件流里读出**最近修改的那个会话**的 `permission/preset`,交给路由层
|
||||
* 判断"是否与平台默认不一致",从而在实例页面上给出提示。
|
||||
*
|
||||
* 只读:仅解压文件头部的若干 zstd 帧,不写任何东西。
|
||||
* @module dshs/supervisor/session-preset
|
||||
*/
|
||||
import { closeSync, existsSync, openSync, readdirSync, readFileSync, readSync, statSync } from 'node:fs'
|
||||
import { join } from 'node:path'
|
||||
import { zstdDecompressSync } from 'node:zlib'
|
||||
|
||||
/** zstd 帧魔数:会话文件是**多帧拼接**,必须按它切分后逐帧解压。 */
|
||||
const MAGIC = Buffer.from([0x28, 0xb5, 0x2f, 0xfd])
|
||||
|
||||
/** 解压文件头部若干帧(够覆盖 session 元信息与 permission/preset 记录,通常在第 1 帧)。 */
|
||||
function headFrames(buf: Buffer, maxFrames = 6): string {
|
||||
const offsets: number[] = []
|
||||
for (let i = 0; i + 4 <= buf.length; i++) {
|
||||
if (buf.compare(MAGIC, 0, 4, i, i + 4) === 0) offsets.push(i)
|
||||
}
|
||||
const list = offsets.length > 0 ? offsets.slice(0, maxFrames) : [0]
|
||||
let out = ''
|
||||
for (let f = 0; f < list.length; f++) {
|
||||
const start = list[f]
|
||||
const end = f + 1 < list.length ? list[f + 1] : buf.length
|
||||
try {
|
||||
out += zstdDecompressSync(buf.subarray(start, end)).toString('utf8')
|
||||
} catch {
|
||||
/* 单帧损坏/被截断:忽略,继续下一帧 */
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
export interface SessionPreset {
|
||||
/** 会话目录名(`session-<uuid>` 或裸 uuid)。 */
|
||||
sessionId: string
|
||||
/** 该会话当前生效的权限档位(取最后一次 `permission/preset` 记录)。 */
|
||||
preset: string
|
||||
/** 会话文件最后修改时间(ms)。 */
|
||||
updatedAt: number
|
||||
}
|
||||
|
||||
/** 用户数据根下的会话目录:`<userRoot>/home/sessions`。 */
|
||||
export function sessionsDir(homeRoot: string): string {
|
||||
return join(homeRoot, 'sessions')
|
||||
}
|
||||
|
||||
/**
|
||||
* 读一个 `session.jsonl.zstd` 的权限档位。
|
||||
* 取**最后一条** `permission/preset`(用户在 UI 里切换档位会追加新记录)。
|
||||
*/
|
||||
function readPreset(file: string): string | undefined {
|
||||
let text: string
|
||||
try {
|
||||
const st = statSync(file)
|
||||
const HEAD = 256 * 1024
|
||||
const TAIL = 64 * 1024
|
||||
if (st.size <= 4 * 1024 * 1024) {
|
||||
// 常见会话 < 4MB:直接整读(切换档位的记录可能在文件尾)
|
||||
text = headFrames(readFileSync(file), 12)
|
||||
} else {
|
||||
const fd = openSync(file, 'r')
|
||||
try {
|
||||
const head = Buffer.allocUnsafe(HEAD)
|
||||
const readHead = readSync(fd, head, 0, HEAD, 0)
|
||||
const tailLen = Math.min(TAIL, st.size)
|
||||
const tail = Buffer.allocUnsafe(tailLen)
|
||||
readSync(fd, tail, 0, tailLen, st.size - tailLen)
|
||||
text = headFrames(head.subarray(0, readHead), 6) + headFrames(tail, 6)
|
||||
} finally {
|
||||
closeSync(fd)
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
return undefined
|
||||
}
|
||||
const re = /"type":"permission\/preset"[^\n]*?"preset":"([^"]+)"/g
|
||||
let last: string | undefined
|
||||
for (const m of text.matchAll(re)) last = m[1]
|
||||
return last
|
||||
}
|
||||
|
||||
/** 该用户**最近修改的**会话及其档位(跨全部工作区),无会话则返回 undefined。 */
|
||||
export function latestSessionPreset(homeRoot: string): SessionPreset | undefined {
|
||||
const dir = sessionsDir(homeRoot)
|
||||
if (!existsSync(dir)) return undefined
|
||||
let best: { file: string; id: string; mtime: number } | undefined
|
||||
for (const slug of readdirSync(dir)) {
|
||||
let sids: string[]
|
||||
try {
|
||||
sids = readdirSync(join(dir, slug))
|
||||
} catch {
|
||||
continue
|
||||
}
|
||||
for (const sid of sids) {
|
||||
const file = join(dir, slug, sid, 'session.jsonl.zstd')
|
||||
try {
|
||||
const st = statSync(file)
|
||||
if (!st.isFile()) continue
|
||||
if (best === undefined || st.mtimeMs > best.mtime) best = { file, id: sid, mtime: st.mtimeMs }
|
||||
} catch {
|
||||
/* 无会话文件(空目录)→ 跳过 */
|
||||
}
|
||||
}
|
||||
}
|
||||
if (best === undefined) return undefined
|
||||
const preset = readPreset(best.file)
|
||||
if (preset === undefined) return undefined
|
||||
return { sessionId: best.id, preset, updatedAt: best.mtime }
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
/**
|
||||
* Child-process helpers for the supervisor: env scrubbing and free-port lookup.
|
||||
*
|
||||
* Env scrubbing mirrors the harness `scrubbedParentEnv` / `SENSITIVE_ENV_PATTERN`
|
||||
* doctrine (packages/subprocess/subprocess/src/index.ts): build the child env
|
||||
* from a clean allowlist so no orchestrator secret leaks into a user DSH, then
|
||||
* inject only the resolved per-user values.
|
||||
* @module dshs/supervisor/spawn
|
||||
*/
|
||||
|
||||
import { createServer } from 'node:net'
|
||||
|
||||
const ALLOWED_ENV = new Set([
|
||||
'PATH',
|
||||
'HOME',
|
||||
'USER',
|
||||
'TMP',
|
||||
'TEMP',
|
||||
'TMPDIR',
|
||||
'SYSTEMROOT',
|
||||
'SystemRoot',
|
||||
'PATHEXT',
|
||||
'ProgramFiles',
|
||||
'ProgramFiles(x86)',
|
||||
'LANG',
|
||||
'LC_ALL',
|
||||
// Shared read-only skill dir; injected explicitly in baseEnv, allowlisted here
|
||||
// so it survives scrubEnv if ever set on the orchestrator process.
|
||||
'DSH_BUNDLED_SKILL_DIR',
|
||||
// 档案 33:实例内权限档位(dsh-base 读 DSH_PERMISSION_MODE 决定 sandbox mode + approval policy)。
|
||||
// 只做 allowlist,实际值由 orchestrator.baseEnv 注入。
|
||||
'DSH_PERMISSION_MODE',
|
||||
// 档案 44:冻结实例的基础运行时版本(Python / pip / node 一律用平台装的那份,禁止版本漂移)。
|
||||
// 这两条是**限制性** env(语义为收窄,不是扩大):
|
||||
// · PYTHONNOUSERSITE=1 —— 不把 `$HOME/.local/lib/python*/site-packages` 加进 sys.path
|
||||
// (实测:一旦该目录被 pip 创建,它就在 sys.path 里且**优先于平台 site-packages** →
|
||||
// 用户装的同名包会盖住平台包,正是"版本差异导致插件功能不可用"的来源);
|
||||
// · PYTHONUSERBASE=<只读占位位> —— 让 `pip install --user` **明确失败**而不是静默无效。
|
||||
'PYTHONNOUSERSITE',
|
||||
'PYTHONUSERBASE',
|
||||
])
|
||||
|
||||
/** Drop credential-shaped and unknown env vars; keep only a safe allowlist. */
|
||||
export function scrubEnv(env: NodeJS.ProcessEnv): Record<string, string> {
|
||||
const out: Record<string, string> = {}
|
||||
for (const [key, value] of Object.entries(env)) {
|
||||
if (ALLOWED_ENV.has(key) && value !== undefined) out[key] = value
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
/** Reserve an ephemeral loopback port, release it, and return its number. */
|
||||
export function findFreePort(): Promise<number> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const server = createServer()
|
||||
server.on('error', reject)
|
||||
server.listen(0, '127.0.0.1', () => {
|
||||
const address = server.address()
|
||||
const port = typeof address === 'object' && address !== null ? address.port : undefined
|
||||
server.close(() => {
|
||||
if (port !== undefined) resolve(port)
|
||||
else reject(new Error('could not reserve a free port'))
|
||||
})
|
||||
})
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,139 @@
|
||||
/**
|
||||
* Backend abstraction for per-user DSH lifecycle (docs/k8s.md §5.2).
|
||||
*
|
||||
* `local` spawns child processes (setuid/iptables) via `LocalSpawner`; `k8s`
|
||||
* creates/deletes per-user DSH Pods via the K8s API (`K8sSpawner`). The route
|
||||
* layer depends only on this interface, so both backends coexist behind the
|
||||
* same API. Shared instance/status types live here so neither backend owns
|
||||
* them.
|
||||
* @module dshs/supervisor/spawner
|
||||
*/
|
||||
|
||||
export type InstanceStatus = 'starting' | 'running' | 'crashed' | 'stopped' | 'failed'
|
||||
export type InstanceRole = 'main' | 'watchdog'
|
||||
|
||||
/** A tracked DSH instance (main or watchdog). `port`/`pid` are local-only. */
|
||||
export interface Instance {
|
||||
id: string
|
||||
userId: string
|
||||
role: InstanceRole
|
||||
folder: string
|
||||
port?: number
|
||||
status: InstanceStatus
|
||||
pid?: number
|
||||
exitCode?: number
|
||||
lastError?: string
|
||||
/** Rendered cordis patch **content**, not a path — see {@link Spawner.launch}. */
|
||||
patch?: string
|
||||
/** dsh web 一次性 launch token(本地模式从子进程 stdout 解析),用于拼装可直达的打开 URL。 */
|
||||
launchToken?: string
|
||||
/** systemd scope 名(bwrap 沙箱隔离时),stop 时用 systemctl stop 正确终止整个 scope。 */
|
||||
unit?: string
|
||||
/** 崩溃自动重启次数(档案 20 观测面,随实例重建归零)。 */
|
||||
restarts?: number
|
||||
/** 最近一次崩溃时间(ms epoch)。 */
|
||||
lastCrashedAt?: number
|
||||
}
|
||||
|
||||
/** Thrown when a user already has a running main DSH. */
|
||||
export class AlreadyRunningError extends Error {
|
||||
constructor(userId: string) {
|
||||
super(`user ${userId} already has a running DSH`)
|
||||
this.name = 'AlreadyRunningError'
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 档案 78:该用户的实例刚因崩溃循环被熔断,冷却期内拒绝启动。
|
||||
*
|
||||
* 为什么需要它:`circuit-open` 之后若允许立刻重来,崩溃循环可无限重复
|
||||
* (用户 F5 / 注入脚本自愈 / 脚本直铺都可能触发),且平台只留一行 stderr。
|
||||
* 冷却过后只给一次干净预算。
|
||||
*/
|
||||
export class CrashBreakerOpenError extends Error {
|
||||
readonly userId: string
|
||||
/** 冷却结束时刻(ms epoch)。 */
|
||||
readonly retryAt: number
|
||||
/** 累计熔断次数。 */
|
||||
readonly opens: number
|
||||
|
||||
constructor(userId: string, retryAt: number, opens: number) {
|
||||
const waitSec = Math.max(0, Math.ceil((retryAt - Date.now()) / 1000))
|
||||
super(`user ${userId} instance crash-breaker open (opens=${opens}, retry in ${waitSec}s)`)
|
||||
this.name = 'CrashBreakerOpenError'
|
||||
this.userId = userId
|
||||
this.retryAt = retryAt
|
||||
this.opens = opens
|
||||
}
|
||||
|
||||
/** 距可重试还剩多少毫秒。 */
|
||||
get retryAfterMs(): number {
|
||||
return Math.max(0, this.retryAt - Date.now())
|
||||
}
|
||||
}
|
||||
|
||||
/** A user's main + watchdog pair. */
|
||||
export interface UserStatus {
|
||||
main?: Instance
|
||||
watchdog?: Instance
|
||||
}
|
||||
|
||||
/** The host:port the proxy forwards a user's DSH traffic to. */
|
||||
export interface Endpoint {
|
||||
host: string
|
||||
port: number
|
||||
}
|
||||
|
||||
/** A main DSH Pod as observed by the k8s backend (for reconcile + the watch). */
|
||||
export interface LivePod {
|
||||
name: string
|
||||
userId: string
|
||||
running: boolean
|
||||
crashed: boolean
|
||||
}
|
||||
|
||||
/**
|
||||
* The lifecycle seam the route layer delegates to.
|
||||
*
|
||||
* `endpointFor` is spawner-specific: local → `127.0.0.1:<port>`, k8s → the
|
||||
* per-user Headless Service DNS (docs/k8s.md §5.4).
|
||||
*
|
||||
* `launch` takes the rendered cordis patch as **content**, not a path: under
|
||||
* k8s the control plane holds no user volume, so it can neither write the patch
|
||||
* nor read it back. `LocalSpawner` materializes it to a file (it does have the
|
||||
* volume) and `K8sSpawner` puts it straight into a ConfigMap.
|
||||
*/
|
||||
export interface Spawner {
|
||||
launch(userId: string, folder: string, patch?: string, opts?: { force?: boolean }): Promise<Instance>
|
||||
restartMain(userId: string): Promise<Instance | undefined>
|
||||
/**
|
||||
* 档案 78:熔断观测面(可选 —— 熔断是**本地模式**概念,k8s 模式没有)。
|
||||
* 返回 null = 该用户未被熔断冷却;非 null = 正在冷却(含累计熔断次数与解冻时刻)。
|
||||
*/
|
||||
breakerInfo?(userId: string): { opens: number; openedAt: number; cooldownUntil: number } | null
|
||||
|
||||
/** Restart every running main so a swapped global API key takes effect (env is a spawn-time snapshot). */
|
||||
restartAllMains(): Promise<void>
|
||||
spawnWatchdog(userId: string): Promise<Instance | undefined>
|
||||
status(userId: string): Promise<UserStatus>
|
||||
endpointFor(userId: string): Promise<Endpoint | undefined>
|
||||
stop(userId: string): Promise<void>
|
||||
teardown(): Promise<void>
|
||||
/** 等待该用户 main 实例打印 launch token(本地模式 = 启动完成的信号)。无实例 /
|
||||
* 已崩溃 / 已停 → 立即返回;k8s 模式无 token 概念 → no-op。供 enter 复用分支在返回
|
||||
* 打开 URL 前等待,避免把浏览器导向「HTTP 已监听但路由未就绪 → 404」的启动窗口。 */
|
||||
waitForLaunchTokenForUser(userId: string, timeoutMs?: number): Promise<void>
|
||||
|
||||
/**
|
||||
* 档案 34:重启该用户实例并**探活**——功能插件启用后判定实例是否还起得来。
|
||||
* `ok:false` 时调用方必须回滚/隔离该插件(否则会把实例拖进崩溃循环,档案 25 教训)。
|
||||
*/
|
||||
restartAndProbe(userId: string, settleMs?: number): Promise<{ ok: boolean; reason: string }>
|
||||
/** Record user activity (proxied traffic / entering the workspace) so idle
|
||||
* reaping keeps warm instances that are genuinely in use. Local mode tracks
|
||||
* this in memory; k8s mode relies on its own session-based reconcile. */
|
||||
touch(userId: string): void
|
||||
/** Make the user's file sidecar exist and be ready (k8s). No-op under local,
|
||||
* where the control plane touches the volume in-process. */
|
||||
ensureFileService(userId: string): Promise<void>
|
||||
}
|
||||
Reference in new issue
Block a user