chore(k8s): 下线 K8s 后端形态,移除依赖 @kubernetes/client-node
生产形态是单机 local(DEFAULT_DEPLOY_MODE=local,env 未覆盖)⇒ K8s 分支
在 local 下本来不可达;且该形态与官方 dsh 基座、插件体系均无关
=> 整体下线,并移除该形态唯一的第三方依赖。
移除(备份在 D:/github/_dsh_shenxian_K8s后端备份_20260915/,含还原命令与
「集群化方案要复用的模板清单」):
src/supervisor/{k8s-spawner,leader,reconcile}.ts
src/fs/k8s-user-fs.ts · src/tcp-bridge.ts · src/web/file-service.ts
test/{k8s-spawner,leader}.test.mjs · scripts/smoke-file-service.mjs
改写调用方:cli.ts(5 条 import / 选主块 / file-service 与 tcp-bridge 两个
子命令 / dispatch / HELP)、web/server.ts(改为 fail-loud 守卫 + 恒用
LocalSpawner)、fs/provider.ts(只留 LocalUserFs)、package.json(测试与
smoke 入口),外加 3 处指向已删类型的悬空 JSDoc。
保留(集群化方案列为未来可选):deploy/ · poc/01-04 · Dockerfile.dsh ·
docs/k8s*.md(已加「代码已下线」状态横幅)· config.ts 的 K8s 配置字段与
DeployMode 联合类型。
验证:tsc --noEmit exit 0;npm test 36 测试 / 35 通过 / 0 失败 / 1 跳过;
npm run verify exit 0;依赖与被删符号全仓 0 命中。
This commit is contained in:
1 parent
ed0c3c0ad5
commit
cf8b7b1f5c
19 files changed
+29
-2690
No files matched your search
-72
@@ -4,7 +4,6 @@
|
||||
*
|
||||
* Subcommands:
|
||||
* dshs bootstrap-admin --username <u> --password <p>
|
||||
* dshs file-service per-user file sidecar (k8s)
|
||||
* dshs [server flags]
|
||||
* @module dshs/cli
|
||||
*/
|
||||
@@ -17,11 +16,6 @@ import { createUserFs } from './fs/provider.js'
|
||||
import { homeRoot, userRoot } from './fs/workspace.js'
|
||||
import { hashPassword } from './web/auth.js'
|
||||
import { hashUid } from './isolation.js'
|
||||
import { LeaderElector } from './supervisor/leader.js'
|
||||
import { ReconcileController } from './supervisor/reconcile.js'
|
||||
import { K8sSpawner } from './supervisor/k8s-spawner.js'
|
||||
import { buildFileService, FILE_SERVICE_PORT, USER_ROOT_ENV } from './web/file-service.js'
|
||||
import { startTcpBridge } from './tcp-bridge.js'
|
||||
import { buildServer } from './web/server.js'
|
||||
|
||||
const HELP = `dshs — DSH server login orchestrator
|
||||
@@ -29,7 +23,6 @@ const HELP = `dshs — DSH server login orchestrator
|
||||
Usage:
|
||||
dshs [options] start the server
|
||||
dshs bootstrap-admin [options] create the first admin
|
||||
dshs file-service run the per-user file sidecar
|
||||
|
||||
Server options:
|
||||
--port <n> Bind port (0 = ephemeral). Default 3080.
|
||||
@@ -138,47 +131,6 @@ async function runServer(args: string[]): Promise<void> {
|
||||
app.log.info(`dshs listening on http://${config.host}:${actualPort}`)
|
||||
app.log.info(`data root: ${config.dataRoot}; db: ${config.dbPath}`)
|
||||
|
||||
// In k8s mode, only the elected leader runs the controller (reconcile + Pod
|
||||
// watch). The web layer serves on every replica.
|
||||
let controller: ReconcileController | undefined
|
||||
if (config.deployMode === 'k8s') {
|
||||
const spawner = app.supervisor as K8sSpawner
|
||||
const elector = new LeaderElector({ namespace: config.k8sNamespace, identity: config.podName })
|
||||
controller = new ReconcileController(app.db, spawner, elector)
|
||||
await controller.start()
|
||||
app.log.info(`leader election started (identity ${config.podName})`)
|
||||
}
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
app.log.info(`received ${signal}, shutting down`)
|
||||
controller?.stop()
|
||||
await app.close()
|
||||
process.exit(0)
|
||||
}
|
||||
process.on('SIGINT', () => void shutdown('SIGINT'))
|
||||
process.on('SIGTERM', () => void shutdown('SIGTERM'))
|
||||
}
|
||||
|
||||
/**
|
||||
* Run the per-user file sidecar. Serves one user's volume over HTTP on 8082 so
|
||||
* the control plane (which holds no users volume under k8s) can reach it; the
|
||||
* root comes from the Pod's env, not from argv.
|
||||
*/
|
||||
async function runFileService(): Promise<void> {
|
||||
const root = process.env[USER_ROOT_ENV]
|
||||
if (root === undefined || root === '') {
|
||||
console.error(`file-service requires ${USER_ROOT_ENV} (the user's data root inside the Pod)`)
|
||||
process.exit(2)
|
||||
}
|
||||
// The sidecar runs as the user's uid with no home dir; `homedir()` falls back
|
||||
// to `/` and `resolveConfig` would try to mkdir `/.dshs`. It only
|
||||
// needs `maxUploadBytes`/`logLevel` here, so pin dataRoot to a writable path
|
||||
// and skip the (unused) encryption-secret file.
|
||||
const config = resolveConfig({ dataRoot: '/tmp', encryptionSecret: 'file-service-unused' })
|
||||
const app = buildFileService(root, { bodyLimit: config.maxUploadBytes, logLevel: config.logLevel })
|
||||
await app.listen({ host: '0.0.0.0', port: FILE_SERVICE_PORT })
|
||||
app.log.info(`file sidecar serving ${root} on 0.0.0.0:${FILE_SERVICE_PORT}`)
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
app.log.info(`received ${signal}, shutting down`)
|
||||
await app.close()
|
||||
@@ -188,22 +140,6 @@ async function runFileService(): Promise<void> {
|
||||
process.on('SIGTERM', () => void shutdown('SIGTERM'))
|
||||
}
|
||||
|
||||
/** TCP bridge sidecar (`dshs tcp-bridge <listen> <target>`). */
|
||||
async function runTcpBridge(args: string[]): Promise<void> {
|
||||
const [listen, target] = args
|
||||
if (listen === undefined || target === undefined) {
|
||||
console.error('usage: dshs tcp-bridge <listen> <target>')
|
||||
process.exit(2)
|
||||
}
|
||||
const server = await startTcpBridge(listen, target)
|
||||
console.error(`tcp-bridge ${listen} -> ${target}`)
|
||||
const shutdown = (): void => {
|
||||
server.close(() => process.exit(0))
|
||||
}
|
||||
process.on('SIGINT', shutdown)
|
||||
process.on('SIGTERM', shutdown)
|
||||
}
|
||||
|
||||
async function uidForUserCmd(args: string[]): Promise<void> {
|
||||
const { values, positionals } = parseArgs({
|
||||
args,
|
||||
@@ -234,14 +170,6 @@ async function main(): Promise<void> {
|
||||
await uidForUserCmd(rest)
|
||||
return
|
||||
}
|
||||
if (first === 'file-service') {
|
||||
await runFileService()
|
||||
return
|
||||
}
|
||||
if (first === 'tcp-bridge') {
|
||||
await runTcpBridge(rest)
|
||||
return
|
||||
}
|
||||
await runServer(process.argv.slice(2))
|
||||
}
|
||||
|
||||
|
||||
@@ -1,167 +0,0 @@
|
||||
/**
|
||||
* HTTP {@link UserFs}: every file operation is delegated to the user's own file
|
||||
* sidecar (docs/k8s.md §4.10), because the control plane runs as uid 65532 with
|
||||
* no users volume and could not touch a `0700` user directory even if it did.
|
||||
*
|
||||
* The sidecar is addressed through its per-user Headless Service. That DNS A
|
||||
* record has a ~30s TTL, so a Pod that was just rebuilt can still resolve to
|
||||
* its old IP — connection-level failures therefore drop the keep-alive pool and
|
||||
* retry once, mirroring what the DSH proxy does (docs/k8s.md §5.4).
|
||||
* @module dshs/fs/k8s-user-fs
|
||||
*/
|
||||
|
||||
import { Agent, request as httpRequest } from 'node:http'
|
||||
import type { Endpoint } from '../supervisor/spawner.js'
|
||||
import { POSIX, PathEscapeError, resolveWithinRoot } from '../web/middleware/fs-guard.js'
|
||||
import type { PluginInfo } from './plugins.js'
|
||||
import { isUserFsErrorCode, UserFsError, type UserFs } from './user-fs.js'
|
||||
import { HOME_DIR, USERS_DIR, WORKSPACE_DIR, type FsEntry } from './workspace.js'
|
||||
|
||||
/** Data root inside every per-user Pod (mirrors `k8s-spawner`'s POD_DATA_ROOT). */
|
||||
const POD_DATA_ROOT = '/var/lib/dshs'
|
||||
|
||||
/** Errors that mean "the connection never got anywhere" — worth one retry
|
||||
* against a freshly resolved address. */
|
||||
const RETRYABLE = new Set(['ECONNREFUSED', 'ECONNRESET', 'ENOTFOUND', 'EAI_AGAIN', 'EHOSTUNREACH', 'ETIMEDOUT'])
|
||||
|
||||
/** Ensure the user's file sidecar exists and is ready; supplied by the spawner. */
|
||||
export type EnsureFileService = (userId: string) => Promise<void>
|
||||
|
||||
interface Reply {
|
||||
status: number
|
||||
body: unknown
|
||||
}
|
||||
|
||||
export class K8sUserFs implements UserFs {
|
||||
private agent = new Agent({ keepAlive: true, maxSockets: 16 })
|
||||
|
||||
/**
|
||||
* @param ensureFileService - brings the sidecar up before the first call.
|
||||
* @param endpointFor - the sidecar's address; the k8s backend supplies the
|
||||
* per-user Headless Service DNS, tests supply a loopback bind.
|
||||
*/
|
||||
constructor(
|
||||
private readonly ensureFileService: EnsureFileService,
|
||||
private readonly endpointFor: (userId: string) => Endpoint,
|
||||
) {}
|
||||
|
||||
async initUserRoot(userId: string): Promise<void> {
|
||||
// Creating the sidecar Pod *is* the initialization: its init container
|
||||
// builds `<pvc>/<userId>/{ws,home}` as the user's own uid (docs/k8s.md §4.9).
|
||||
await this.ensureFileService(userId)
|
||||
await this.call(userId, 'POST', '/fs/init')
|
||||
}
|
||||
|
||||
resolvePath(userId: string, relPath: string): string {
|
||||
const root = `${POD_DATA_ROOT}/${USERS_DIR}/${userId}/${WORKSPACE_DIR}`
|
||||
try {
|
||||
return resolveWithinRoot(root, relPath, POSIX)
|
||||
} catch (err) {
|
||||
if (err instanceof PathEscapeError) throw new UserFsError('bad_path')
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
/** The user's DSH state directory inside their Pod. */
|
||||
homePath(userId: string): string {
|
||||
return `${POD_DATA_ROOT}/${USERS_DIR}/${userId}/${HOME_DIR}`
|
||||
}
|
||||
|
||||
async listDir(userId: string, relPath: string): Promise<FsEntry[]> {
|
||||
const body = await this.call(userId, 'GET', `/fs/tree?path=${encodeURIComponent(relPath)}`)
|
||||
return (body as { entries: FsEntry[] }).entries
|
||||
}
|
||||
|
||||
async mkdir(userId: string, relPath: string): Promise<void> {
|
||||
await this.call(userId, 'POST', '/fs/mkdir', { path: relPath })
|
||||
}
|
||||
|
||||
async createEntry(userId: string, relPath: string, name: string, type: 'file' | 'dir'): Promise<string> {
|
||||
const body = await this.call(userId, 'POST', '/fs/create', { path: relPath, name, type })
|
||||
return (body as { name: string }).name
|
||||
}
|
||||
|
||||
async upload(userId: string, relPath: string, name: string, data: Buffer): Promise<string> {
|
||||
const body = await this.call(userId, 'POST', '/fs/upload', {
|
||||
path: relPath,
|
||||
name,
|
||||
data: data.toString('base64'),
|
||||
})
|
||||
return (body as { name: string }).name
|
||||
}
|
||||
|
||||
/** 档案 56:k8s 路径未验证 —— sidecar 尚无 read 端点,明确报 `unsupported` 而非静默失败。 */
|
||||
async readFile(): Promise<{ name: string; data: Buffer }> {
|
||||
throw new UserFsError('unsupported')
|
||||
}
|
||||
|
||||
async isDirectory(userId: string, relPath: string): Promise<boolean> {
|
||||
const body = await this.call(userId, 'GET', `/fs/stat?path=${encodeURIComponent(relPath)}`)
|
||||
return (body as { isDirectory: boolean }).isDirectory
|
||||
}
|
||||
|
||||
async listInstalledPlugins(userId: string): Promise<PluginInfo[]> {
|
||||
const body = await this.call(userId, 'GET', '/fs/plugins')
|
||||
return (body as { plugins: PluginInfo[] }).plugins
|
||||
}
|
||||
|
||||
async writeHandoff(userId: string, content: string): Promise<void> {
|
||||
await this.call(userId, 'POST', '/fs/handoff', { content })
|
||||
}
|
||||
|
||||
/** One sidecar call: ensure the Pod, send, retry once on a dead connection,
|
||||
* then translate a non-2xx `{error}` back into a {@link UserFsError}. */
|
||||
private async call(userId: string, method: string, path: string, payload?: unknown): Promise<unknown> {
|
||||
await this.ensureFileService(userId)
|
||||
let reply: Reply
|
||||
try {
|
||||
reply = await this.send(userId, method, path, payload)
|
||||
} catch (err) {
|
||||
const code = (err as NodeJS.ErrnoException).code
|
||||
if (code === undefined || !RETRYABLE.has(code)) throw err
|
||||
// The keep-alive pool may hold sockets to a Pod IP that no longer exists;
|
||||
// drop them so the retry re-resolves the Headless Service.
|
||||
this.agent.destroy()
|
||||
this.agent = new Agent({ keepAlive: true, maxSockets: 16 })
|
||||
reply = await this.send(userId, method, path, payload)
|
||||
}
|
||||
if (reply.status >= 200 && reply.status < 300) return reply.body
|
||||
const error = (reply.body as { error?: unknown }).error
|
||||
if (typeof error === 'string' && isUserFsErrorCode(error)) throw new UserFsError(error)
|
||||
throw new Error(`file sidecar for ${userId} returned ${reply.status}`)
|
||||
}
|
||||
|
||||
private send(userId: string, method: string, path: string, payload?: unknown): Promise<Reply> {
|
||||
const body = payload === undefined ? undefined : JSON.stringify(payload)
|
||||
const endpoint = this.endpointFor(userId)
|
||||
return new Promise<Reply>((resolve, reject) => {
|
||||
const req = httpRequest(
|
||||
{
|
||||
host: endpoint.host,
|
||||
port: endpoint.port,
|
||||
path,
|
||||
method,
|
||||
agent: this.agent,
|
||||
headers: body === undefined
|
||||
? {}
|
||||
: { 'content-type': 'application/json', 'content-length': Buffer.byteLength(body) },
|
||||
},
|
||||
(res) => {
|
||||
const chunks: Buffer[] = []
|
||||
res.on('data', (chunk: Buffer) => chunks.push(chunk))
|
||||
res.on('end', () => {
|
||||
const text = Buffer.concat(chunks).toString('utf8')
|
||||
try {
|
||||
resolve({ status: res.statusCode ?? 502, body: text === '' ? {} : JSON.parse(text) })
|
||||
} catch {
|
||||
reject(new Error(`file sidecar for ${userId} returned non-JSON (${res.statusCode})`))
|
||||
}
|
||||
})
|
||||
},
|
||||
)
|
||||
req.on('error', reject)
|
||||
if (body !== undefined) req.write(body)
|
||||
req.end()
|
||||
})
|
||||
}
|
||||
}
|
||||
+6
-24
@@ -1,37 +1,19 @@
|
||||
/**
|
||||
* {@link UserFs} factory. Picks the implementation from `deployMode`, the same
|
||||
* way {@link createDbAdapter} picks a DB backend and `buildServer` picks a
|
||||
* {@link Spawner}.
|
||||
* {@link UserFs} factory. Builds the single-machine per-user filesystem, the
|
||||
* same way {@link createDbAdapter} picks a DB backend and `buildServer` builds
|
||||
* a {@link Spawner}.
|
||||
* @module dshs/fs/provider
|
||||
*/
|
||||
|
||||
import type { ServerConfig } from '../config.js'
|
||||
import type { Endpoint } from '../supervisor/spawner.js'
|
||||
import { FILE_SERVICE_PORT } from '../web/file-service.js'
|
||||
import { K8sUserFs, type EnsureFileService } from './k8s-user-fs.js'
|
||||
import { LocalUserFs } from './local-user-fs.js'
|
||||
import type { UserFs } from './user-fs.js'
|
||||
import { userRoot } from './workspace.js'
|
||||
|
||||
/** Headless Service fronting a user's file sidecar (docs/k8s.md §4.10). */
|
||||
export function fileServiceName(userId: string): string {
|
||||
return `dsh-files-${userId}`
|
||||
}
|
||||
|
||||
/** In-cluster address of a user's file sidecar. */
|
||||
export function fileServiceEndpoint(namespace: string, userId: string): Endpoint {
|
||||
return { host: `${fileServiceName(userId)}.${namespace}.svc.cluster.local`, port: FILE_SERVICE_PORT }
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the configured per-user filesystem.
|
||||
*
|
||||
* `local` touches the users volume in-process. `k8s` delegates to each user's
|
||||
* file sidecar, which the spawner brings up on demand via `ensureFileService`.
|
||||
* Build the configured per-user filesystem. The single-machine backend
|
||||
* touches the users volume in-process.
|
||||
*/
|
||||
export function createUserFs(config: ServerConfig, ensureFileService?: EnsureFileService): UserFs {
|
||||
if (config.deployMode === 'k8s' && ensureFileService !== undefined) {
|
||||
return new K8sUserFs(ensureFileService, (userId) => fileServiceEndpoint(config.k8sNamespace, userId))
|
||||
}
|
||||
export function createUserFs(config: ServerConfig): UserFs {
|
||||
return new LocalUserFs((userId) => userRoot(config.dataRoot, userId))
|
||||
}
|
||||
+3
-6
@@ -1,13 +1,10 @@
|
||||
/**
|
||||
* The per-user filesystem seam (docs/k8s.md §4.10 / §6.0-1).
|
||||
*
|
||||
* Under `local` the control plane owns the users volume and touches it directly
|
||||
* ({@link LocalUserFs}). Under `k8s` it must not: it runs as uid 65532 while
|
||||
* each user's directory is `0700` owned by that user's uid, so every file
|
||||
* operation is delegated over HTTP to a per-user file sidecar
|
||||
* ({@link K8sUserFs}). Routes depend only on this interface.
|
||||
* The control plane owns the users volume and touches it in-process
|
||||
* ({@link LocalUserFs}); routes depend only on this interface.
|
||||
*
|
||||
* All paths crossing this interface are **workspace-relative**; each
|
||||
* All paths crossing this interface are **workspace-relative**; the
|
||||
* implementation resolves them against its own root via `resolveWithinRoot`.
|
||||
* @module dshs/fs/user-fs
|
||||
*/
|
||||
|
||||
@@ -1,687 +0,0 @@
|
||||
/**
|
||||
* 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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,254 +0,0 @@
|
||||
/**
|
||||
* 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 })
|
||||
}
|
||||
}
|
||||
@@ -1,170 +0,0 @@
|
||||
/**
|
||||
* 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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,9 @@
|
||||
/**
|
||||
* 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.
|
||||
* `LocalSpawner` spawns child processes (setuid/iptables). The route layer
|
||||
* depends only on this interface. Shared instance/status types live here so
|
||||
* no backend owns them.
|
||||
* @module dshs/supervisor/spawner
|
||||
*/
|
||||
|
||||
@@ -95,13 +93,11 @@ export interface LivePod {
|
||||
/**
|
||||
* 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).
|
||||
* `endpointFor` is spawner-specific: `127.0.0.1:<port>` for the local backend
|
||||
* (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.
|
||||
* `launch` takes the rendered cordis patch as **content**, not a path;
|
||||
* `LocalSpawner` materializes it to a file inside the user's own volume.
|
||||
*/
|
||||
export interface Spawner {
|
||||
launch(userId: string, folder: string, patch?: string, opts?: { force?: boolean }): Promise<Instance>
|
||||
|
||||
@@ -1,38 +0,0 @@
|
||||
/**
|
||||
* Tiny TCP bridge: listen on one address and forward every connection to
|
||||
* another. Replaces the `alpine/socat` sidecar in the per-user DSH Pod, so the
|
||||
* Pod no longer depends on a docker.io image (blocked on ACK, docs/k8s-deploy.md
|
||||
* §7). Runs from the control-plane image's Node runtime.
|
||||
* @module dshs/tcp-bridge
|
||||
*/
|
||||
|
||||
import { createConnection, createServer, type Server } from 'node:net'
|
||||
|
||||
/**
|
||||
* Start a TCP bridge: `listen` (e.g. `0.0.0.0:8081`) → `target`
|
||||
* (e.g. `127.0.0.1:8080`). Resolves when listening.
|
||||
*/
|
||||
export function startTcpBridge(listen: string, target: string): Promise<Server> {
|
||||
const [targetHost, targetPort] = splitHostPort(target)
|
||||
const server = createServer((socket) => {
|
||||
// A fresh outbound connection to the target — NOT `socket.connect`, which
|
||||
// would try to re-connect the already-connected inbound socket.
|
||||
const upstream = createConnection(Number(targetPort), targetHost)
|
||||
upstream.on('error', () => socket.destroy())
|
||||
socket.on('error', () => upstream.destroy())
|
||||
socket.pipe(upstream)
|
||||
upstream.pipe(socket)
|
||||
})
|
||||
const [listenHost, listenPort] = splitHostPort(listen)
|
||||
return new Promise((resolve, reject) => {
|
||||
server.once('error', reject)
|
||||
server.listen(Number(listenPort), listenHost, () => resolve(server))
|
||||
})
|
||||
}
|
||||
|
||||
/** `host:port` → `[host, port]`, defaulting the host to `0.0.0.0`. */
|
||||
function splitHostPort(addr: string): [string, string] {
|
||||
const idx = addr.lastIndexOf(':')
|
||||
if (idx === -1) return ['0.0.0.0', addr]
|
||||
return [addr.slice(0, idx), addr.slice(idx + 1)]
|
||||
}
|
||||
@@ -1,150 +0,0 @@
|
||||
/**
|
||||
* The per-user file sidecar (docs/k8s.md §4.10).
|
||||
*
|
||||
* Runs inside the user's own Pod as the user's uid, so it *can* read and write
|
||||
* their `0700` directory — the thing the control plane (uid 65532, no volume)
|
||||
* cannot do. It serves exactly one user: the root is fixed at startup, and no
|
||||
* request carries a user id.
|
||||
*
|
||||
* **No authentication.** The boundary is the NetworkPolicy that lets only the
|
||||
* control plane reach port 8082, the same argument that lets the socat sidecar
|
||||
* bridge 8081 unauthenticated (docs/k8s.md §6.1 item 4). If that policy is ever
|
||||
* widened, this service needs a token check *first*.
|
||||
* @module dshs/web/file-service
|
||||
*/
|
||||
|
||||
import Fastify, { type FastifyInstance, type FastifyReply } from 'fastify'
|
||||
import { LocalUserFs } from '../fs/local-user-fs.js'
|
||||
import { UserFsError } from '../fs/user-fs.js'
|
||||
|
||||
/** Port the sidecar binds; the per-user Service targets it directly. */
|
||||
export const FILE_SERVICE_PORT = 8082
|
||||
|
||||
/** Env var carrying the single user root this sidecar serves. */
|
||||
export const USER_ROOT_ENV = 'DSHS_USER_ROOT'
|
||||
|
||||
/** The sidecar serves one user, so the id crossing {@link UserFs} is a constant. */
|
||||
const SOLE_USER = 'self'
|
||||
|
||||
const pathSchema = {
|
||||
body: {
|
||||
type: 'object',
|
||||
required: ['path'],
|
||||
additionalProperties: false,
|
||||
properties: { path: { type: 'string', maxLength: 512 } },
|
||||
},
|
||||
} as const
|
||||
|
||||
const entrySchema = {
|
||||
body: {
|
||||
type: 'object',
|
||||
required: ['path', 'name', 'type'],
|
||||
additionalProperties: false,
|
||||
properties: {
|
||||
path: { type: 'string', maxLength: 512 },
|
||||
name: { type: 'string', maxLength: 255 },
|
||||
type: { type: 'string', enum: ['file', 'dir'] },
|
||||
},
|
||||
},
|
||||
} as const
|
||||
|
||||
const uploadSchema = {
|
||||
body: {
|
||||
type: 'object',
|
||||
required: ['path', 'name', 'data'],
|
||||
additionalProperties: false,
|
||||
properties: {
|
||||
path: { type: 'string', maxLength: 512 },
|
||||
name: { type: 'string', maxLength: 255 },
|
||||
data: { type: 'string' },
|
||||
},
|
||||
},
|
||||
} as const
|
||||
|
||||
const handoffSchema = {
|
||||
body: {
|
||||
type: 'object',
|
||||
required: ['content'],
|
||||
additionalProperties: false,
|
||||
properties: { content: { type: 'string', maxLength: 8192 } },
|
||||
},
|
||||
} as const
|
||||
|
||||
/** Reply with the seam's own wire form so the client can rebuild the error. */
|
||||
function fail(reply: FastifyReply, err: unknown): FastifyReply {
|
||||
if (err instanceof UserFsError) return reply.code(err.status).send({ error: err.code })
|
||||
throw err
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the sidecar's Fastify instance. Does not listen; the caller binds.
|
||||
* @param root - the user's data root inside the Pod.
|
||||
* @param options - `bodyLimit` must exceed the control plane's upload cap.
|
||||
*/
|
||||
export function buildFileService(root: string, options: { bodyLimit: number; logLevel: string }): FastifyInstance {
|
||||
const fs = new LocalUserFs(() => root)
|
||||
const app = Fastify({ logger: { level: options.logLevel }, bodyLimit: options.bodyLimit })
|
||||
|
||||
app.get('/healthz', async () => ({ ok: true }))
|
||||
|
||||
app.get('/fs/tree', async (request, reply) => {
|
||||
const { path = '' } = request.query as { path?: string }
|
||||
try {
|
||||
return { entries: await fs.listDir(SOLE_USER, path) }
|
||||
} catch (err) {
|
||||
return fail(reply, err)
|
||||
}
|
||||
})
|
||||
|
||||
app.get('/fs/stat', async (request, reply) => {
|
||||
const { path = '' } = request.query as { path?: string }
|
||||
try {
|
||||
return { isDirectory: await fs.isDirectory(SOLE_USER, path) }
|
||||
} catch (err) {
|
||||
return fail(reply, err)
|
||||
}
|
||||
})
|
||||
|
||||
app.get('/fs/plugins', async () => ({ plugins: await fs.listInstalledPlugins(SOLE_USER) }))
|
||||
|
||||
app.post('/fs/init', async () => {
|
||||
await fs.initUserRoot(SOLE_USER)
|
||||
return { ok: true }
|
||||
})
|
||||
|
||||
app.post('/fs/mkdir', { schema: pathSchema }, async (request, reply) => {
|
||||
const { path } = request.body as { path: string }
|
||||
try {
|
||||
await fs.mkdir(SOLE_USER, path)
|
||||
return { ok: true }
|
||||
} catch (err) {
|
||||
return fail(reply, err)
|
||||
}
|
||||
})
|
||||
|
||||
app.post('/fs/create', { schema: entrySchema }, async (request, reply) => {
|
||||
const { path, name, type } = request.body as { path: string; name: string; type: 'file' | 'dir' }
|
||||
try {
|
||||
return { name: await fs.createEntry(SOLE_USER, path, name, type) }
|
||||
} catch (err) {
|
||||
return fail(reply, err)
|
||||
}
|
||||
})
|
||||
|
||||
app.post('/fs/upload', { schema: uploadSchema }, async (request, reply) => {
|
||||
const { path, name, data } = request.body as { path: string; name: string; data: string }
|
||||
try {
|
||||
return { name: await fs.upload(SOLE_USER, path, name, Buffer.from(data, 'base64')) }
|
||||
} catch (err) {
|
||||
return fail(reply, err)
|
||||
}
|
||||
})
|
||||
|
||||
app.post('/fs/handoff', { schema: handoffSchema }, async (request) => {
|
||||
const { content } = request.body as { content: string }
|
||||
await fs.writeHandoff(SOLE_USER, content)
|
||||
return { ok: true }
|
||||
})
|
||||
|
||||
return app
|
||||
}
|
||||
+5
-6
@@ -17,7 +17,6 @@ import type { UserFs } from '../fs/user-fs.js'
|
||||
import { decrypt, deriveKey } from '../crypto.js'
|
||||
import { hashUid } from '../isolation.js'
|
||||
import { LocalSpawner } from '../supervisor/orchestrator.js'
|
||||
import { K8sSpawner } from '../supervisor/k8s-spawner.js'
|
||||
import { registerDshProxy } from '../supervisor/proxy.js'
|
||||
import type { Spawner } from '../supervisor/spawner.js'
|
||||
import {
|
||||
@@ -257,11 +256,11 @@ export async function buildServer(config: ServerConfig): Promise<FastifyInstance
|
||||
const user = await db.findUserById(userId)
|
||||
return user?.uid ?? hashUid(userId, config.baseUid)
|
||||
}
|
||||
const supervisor: Spawner =
|
||||
config.deployMode === 'k8s'
|
||||
? new K8sSpawner(config, db, resolveApiKey, resolveUid)
|
||||
: new LocalSpawner(config, resolveApiKey, resolveUid)
|
||||
const userFs = createUserFs(config, (userId) => supervisor.ensureFileService(userId))
|
||||
if (config.deployMode === 'k8s') {
|
||||
throw new Error('deployMode "k8s" is not supported by this build: only the single-machine backend ships')
|
||||
}
|
||||
const supervisor: Spawner = new LocalSpawner(config, resolveApiKey, resolveUid)
|
||||
const userFs = createUserFs(config)
|
||||
|
||||
const app = Fastify({
|
||||
logger: { level: config.logLevel },
|
||||
|
||||
Reference in new issue
Block a user