Files
SnapOtter/apps/api/src/lib/feature-status.ts
T
SnapOtter b1e94bd988 fix: build script venv handling and lint fixes
- Use /opt/venv directly when --entrypoint bash bypasses entrypoint.sh
- Use sys.executable for all pip calls (not bare pip)
- Override entrypoint in CI workflow to avoid startup banner
- Fix Biome formatting (template literals, try/catch blocks)
2026-06-13 20:48:27 +08:00

741 lines
24 KiB
TypeScript

import { execFileSync } from "node:child_process";
import { randomUUID } from "node:crypto";
import {
constants,
copyFileSync,
existsSync,
mkdirSync,
openSync,
readdirSync,
readFileSync,
renameSync,
rmSync,
statSync,
unlinkSync,
writeFileSync,
} from "node:fs";
import { dirname, join, resolve } from "node:path";
import type { Readable } from "node:stream";
import { fileURLToPath } from "node:url";
import type { FeatureBundleState, FeatureStatus } from "@snapotter/shared";
import { FEATURE_BUNDLES, TOOL_BUNDLE_MAP } from "@snapotter/shared";
import * as tar from "tar";
// ── Paths ───────────────────────────────────────────────────────────────
const __dirname = dirname(fileURLToPath(import.meta.url));
const PROJECT_ROOT = resolve(__dirname, "../../../..");
const DATA_DIR = process.env.DATA_DIR || "/data";
const AI_DIR = join(DATA_DIR, "ai");
const MODELS_DIR = join(AI_DIR, "models");
const INSTALLED_PATH = join(AI_DIR, "installed.json");
const INSTALLED_TMP_PATH = `${INSTALLED_PATH}.tmp`;
const LOCK_PATH = join(AI_DIR, "install.lock");
const MANIFEST_PATH =
process.env.FEATURE_MANIFEST_PATH || join(PROJECT_ROOT, "docker/feature-manifest.json");
export function getAiDir(): string {
return AI_DIR;
}
export function getModelsDir(): string {
return MODELS_DIR;
}
export function getManifestPath(): string {
return MANIFEST_PATH;
}
export function getInstallScriptPath(): string {
return join(PROJECT_ROOT, "packages/ai/python/install_feature.py");
}
// ── Directory setup ─────────────────────────────────────────────────────
export function ensureAiDirs(): void {
if (!isDockerEnvironment()) return;
try {
mkdirSync(join(AI_DIR, "venv"), { recursive: true });
mkdirSync(MODELS_DIR, { recursive: true });
mkdirSync(join(AI_DIR, "pip-cache"), { recursive: true });
} catch (err: unknown) {
// Never refuse to boot over the AI data dir. On native checkouts the
// default DATA_DIR (/data) is often uncreatable (ENOENT/EROFS on a
// sealed macOS root, EACCES on restrictive volumes); AI tools simply
// report as not installed until DATA_DIR points somewhere writable.
const code = (err as NodeJS.ErrnoException).code;
console.error(
`WARNING: Cannot create AI directories under "${AI_DIR}" (${code}). AI features will be unavailable. Set DATA_DIR to a writable path (or check volume permissions / PUID / PGID in Docker).`,
);
}
}
// ── Docker detection ────────────────────────────────────────────────────
export function isDockerEnvironment(): boolean {
return existsSync("/.dockerenv") || existsSync(MANIFEST_PATH);
}
// ── installed.json cache ────────────────────────────────────────────────
interface InstalledBundle {
version: string;
installedAt: string;
models: string[];
}
interface InstalledData {
bundles: Record<string, InstalledBundle>;
}
let installedCache: InstalledData | null = null;
function readInstalled(): InstalledData {
if (installedCache) return installedCache;
if (!existsSync(INSTALLED_PATH)) {
installedCache = { bundles: {} };
return installedCache;
}
try {
const raw = readFileSync(INSTALLED_PATH, "utf-8");
installedCache = JSON.parse(raw) as InstalledData;
} catch {
console.warn("[feature-status] installed.json is corrupt or unreadable, treating as empty");
installedCache = { bundles: {} };
}
return installedCache;
}
function writeInstalled(data: InstalledData): void {
writeFileSync(INSTALLED_TMP_PATH, JSON.stringify(data, null, 2), "utf-8");
renameSync(INSTALLED_TMP_PATH, INSTALLED_PATH);
}
export function invalidateCache(): void {
installedCache = null;
}
// ── Install status queries ──────────────────────────────────────────────
export function isFeatureInstalled(bundleId: string): boolean {
const data = readInstalled();
return bundleId in data.bundles;
}
export function isToolInstalled(toolId: string): boolean {
const bundleId = TOOL_BUNDLE_MAP[toolId];
if (!bundleId) return true;
return isFeatureInstalled(bundleId);
}
// ── Install status mutations ────────────────────────────────────────────
export function markInstalled(bundleId: string, version: string, models: string[]): void {
const data = readInstalled();
data.bundles[bundleId] = {
version,
installedAt: new Date().toISOString(),
models,
};
writeInstalled(data);
invalidateCache();
}
export function markUninstalled(bundleId: string): void {
const data = readInstalled();
delete data.bundles[bundleId];
writeInstalled(data);
invalidateCache();
}
// ── Install lock (file-based) ───────────────────────────────────────────
interface LockData {
bundleId: string;
startedAt: string;
}
const LOCK_STALE_MS = 45 * 60 * 1000;
export function acquireInstallLock(bundleId: string): boolean {
// If a lock exists, check its age. OOM-killed processes or crashes
// may leave a stale lock file behind; treat anything older than 45
// minutes as abandoned.
if (existsSync(LOCK_PATH)) {
try {
const age = Date.now() - statSync(LOCK_PATH).mtimeMs;
if (age > LOCK_STALE_MS) {
unlinkSync(LOCK_PATH);
console.warn(
`[feature-status] Removed stale install lock (age: ${Math.round(age / 1000)}s)`,
);
}
} catch {
// stat/unlink failed -- fall through to O_EXCL which will fail too
}
}
try {
const fd = openSync(LOCK_PATH, constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL);
const lock: LockData = {
bundleId,
startedAt: new Date().toISOString(),
};
writeFileSync(fd, JSON.stringify(lock, null, 2), "utf-8");
return true;
} catch {
return false;
}
}
export function releaseInstallLock(): void {
try {
unlinkSync(LOCK_PATH);
} catch {
// Lock already gone
}
}
export function getInstallingBundle(): {
bundleId: string;
startedAt: string;
} | null {
if (!existsSync(LOCK_PATH)) return null;
try {
const raw = readFileSync(LOCK_PATH, "utf-8");
const lock = JSON.parse(raw) as LockData;
return { bundleId: lock.bundleId, startedAt: lock.startedAt };
} catch {
try {
unlinkSync(LOCK_PATH);
} catch {
// already gone
}
return null;
}
}
// ── Progress tracking (in-memory, for SSE) ──────────────────────────────
let currentProgress: {
bundleId: string;
progress: { percent: number; stage: string } | null;
error: string | null;
} | null = null;
export function setInstallProgress(
bundleId: string | null,
progress: { percent: number; stage: string } | null,
error: string | null,
): void {
if (!bundleId) {
currentProgress = null;
return;
}
currentProgress = { bundleId, progress, error };
}
// ── Manifest reading ────────────────────────────────────────────────────
interface ManifestModel {
id: string;
path?: string;
downloadFn?: string;
args?: string[];
minSize?: number;
}
interface ManifestBundle {
models: ManifestModel[];
}
interface Manifest {
bundles: Record<string, ManifestBundle>;
}
function readManifest(): Manifest | null {
if (!existsSync(MANIFEST_PATH)) return null;
try {
return JSON.parse(readFileSync(MANIFEST_PATH, "utf-8")) as Manifest;
} catch {
return null;
}
}
// ── Startup recovery ────────────────────────────────────────────────────
function deleteDownloadingFiles(dir: string): void {
if (!existsSync(dir)) return;
const entries = readdirSync(dir, { withFileTypes: true, recursive: true });
for (const entry of entries) {
if (entry.isFile() && entry.name.endsWith(".downloading")) {
const fullPath = join(entry.parentPath ?? entry.path, entry.name);
try {
unlinkSync(fullPath);
console.info(`[feature-status] Deleted partial download: ${fullPath}`);
} catch {
// best-effort
}
}
}
}
export function recoverInterruptedInstalls(): void {
// 1. Delete partial downloads
deleteDownloadingFiles(MODELS_DIR);
// 2. Delete stale tmp file
if (existsSync(INSTALLED_TMP_PATH)) {
try {
unlinkSync(INSTALLED_TMP_PATH);
console.info("[feature-status] Deleted stale installed.json.tmp");
} catch {
// best-effort
}
}
// 3. Delete bootstrapping venv
const bootstrappingDir = join(AI_DIR, "venv.bootstrapping");
if (existsSync(bootstrappingDir)) {
try {
rmSync(bootstrappingDir, { recursive: true, force: true });
console.info("[feature-status] Deleted stale venv.bootstrapping/");
} catch {
// best-effort
}
}
// 4. Delete stale lock — if the server is starting up, any previous install is dead
if (existsSync(LOCK_PATH)) {
try {
const raw = readFileSync(LOCK_PATH, "utf-8");
const lock = JSON.parse(raw) as LockData;
unlinkSync(LOCK_PATH);
console.warn(
`[feature-status] Removed stale install lock for "${lock.bundleId}" (server restarted)`,
);
} catch {
try {
unlinkSync(LOCK_PATH);
} catch {
// already gone
}
}
}
// 5. Verify installed bundles still have their model files
const manifest = readManifest();
if (manifest) {
const data = readInstalled();
for (const bundleId of Object.keys(data.bundles)) {
const manifestBundle = manifest.bundles[bundleId];
if (!manifestBundle) continue;
for (const model of manifestBundle.models) {
if (model.path) {
const modelPath = join(MODELS_DIR, model.path);
if (!existsSync(modelPath)) {
console.warn(`[feature-status] Bundle "${bundleId}" missing model file: ${model.path}`);
break;
}
if (model.minSize != null && model.minSize > 0) {
try {
const st = statSync(modelPath);
if (st.size < model.minSize) {
console.warn(
`[feature-status] Bundle "${bundleId}" model "${model.path}" is undersized (${st.size} < ${model.minSize})`,
);
break;
}
} catch {
console.warn(
`[feature-status] Bundle "${bundleId}" cannot stat model: ${model.path}`,
);
break;
}
}
} else if (model.downloadFn === "rembg_session" && model.args?.[0]) {
const filePath = join(MODELS_DIR, "rembg", `${model.args[0]}.onnx`);
if (!existsSync(filePath)) {
console.warn(
`[feature-status] Bundle "${bundleId}" missing rembg model: ${model.args[0]}`,
);
break;
}
} else if (model.downloadFn === "hf_snapshot" && model.args?.[1]) {
const dirPath = join(MODELS_DIR, model.args[1]);
if (!existsSync(dirPath)) {
console.warn(
`[feature-status] Bundle "${bundleId}" missing model directory: ${model.args[1]}`,
);
break;
}
}
}
}
}
// 6. Delete staging-{bundleId}/ directories (incomplete extraction)
try {
const aiEntries = readdirSync(AI_DIR, { withFileTypes: true });
for (const entry of aiEntries) {
if (entry.isDirectory() && entry.name.startsWith("staging-")) {
const stagingPath = join(AI_DIR, entry.name);
rmSync(stagingPath, { recursive: true, force: true });
console.info(`[feature-status] Deleted orphaned ${entry.name}/`);
}
}
} catch {
// AI_DIR may not exist yet
}
// 7. Clean up staging/ download directory (partial downloads, orphaned tars)
const downloadStaging = join(AI_DIR, "staging");
if (existsSync(downloadStaging)) {
try {
const files = readdirSync(downloadStaging);
for (const file of files) {
const filePath = join(downloadStaging, file);
if (file.endsWith(".partial") || file.endsWith(".meta")) {
unlinkSync(filePath);
console.info(`[feature-status] Deleted stale download file: ${file}`);
} else if (file.endsWith(".tar.gz")) {
unlinkSync(filePath);
console.info(`[feature-status] Deleted orphaned archive: ${file}`);
}
}
} catch {
// best-effort
}
}
invalidateCache();
}
// ── Feature states (composite view) ─────────────────────────────────────
export function verifyBundleModels(bundleId: string): string | null {
const manifest = readManifest();
if (!manifest) return null;
const manifestBundle = manifest.bundles[bundleId];
if (!manifestBundle) return null;
for (const model of manifestBundle.models) {
if (model.path) {
const modelPath = join(MODELS_DIR, model.path);
if (!existsSync(modelPath)) {
return `Missing model file: ${model.path}`;
}
if (model.minSize != null && model.minSize > 0) {
try {
const st = statSync(modelPath);
if (st.size < model.minSize) {
return `Model "${model.path}" is undersized (${st.size} < ${model.minSize})`;
}
} catch {
return `Cannot read model file: ${model.path}`;
}
}
} else if (model.downloadFn === "rembg_session" && model.args?.[0]) {
const filePath = join(MODELS_DIR, "rembg", `${model.args[0]}.onnx`);
if (!existsSync(filePath)) {
return `Missing rembg model: ${model.args[0]}`;
}
} else if (model.downloadFn === "hf_snapshot" && model.args?.[1]) {
const dirPath = join(MODELS_DIR, model.args[1]);
if (!existsSync(dirPath)) {
return `Missing model directory: ${model.args[1]}`;
}
}
}
return null;
}
export function getFeatureStates(): FeatureBundleState[] {
const installed = readInstalled();
const lock = getInstallingBundle();
return Object.values(FEATURE_BUNDLES).map((bundle) => {
const installedBundle = installed.bundles[bundle.id];
let status: FeatureStatus = "not_installed";
let error: string | null = null;
let progress: { percent: number; stage: string } | null = null;
if (lock && lock.bundleId === bundle.id) {
status = "installing";
if (currentProgress && currentProgress.bundleId === bundle.id) {
progress = currentProgress.progress;
if (currentProgress.error) {
status = "error";
error = currentProgress.error;
}
}
} else if (installedBundle) {
// Verify model files exist and are properly sized
const modelError = verifyBundleModels(bundle.id);
if (modelError) {
status = "error";
error = modelError;
} else {
status = "installed";
}
} else if (currentProgress?.bundleId === bundle.id && currentProgress.error) {
status = "error";
error = currentProgress.error;
}
return {
id: bundle.id,
name: bundle.name,
description: bundle.description,
status,
installedVersion: installedBundle?.version ?? null,
estimatedSize: bundle.estimatedSize,
enablesTools: bundle.enablesTools,
progress,
error,
};
});
}
// ── Offline bundle import ──────────────────────────────────────────────
//
// Archive format (v1):
// Gzipped tar containing:
// bundle.json - { bundleId, version, models: string[] }
// models/... - files mirroring MODELS_DIR layout
//
// v1 validates bundleId against the manifest but does NOT verify per-file
// checksums. Transport integrity is the operator's responsibility; a
// checksum manifest is a phase-3 candidate.
//
// Security: symlink/hardlink and other non-file entry types are rejected
// during extraction. Only "File" and "Directory" entries are permitted.
interface BundleDescriptor {
bundleId: string;
version: string;
models: string[];
}
const IMPORT_MAX_BYTES = 20 * 1024 * 1024 * 1024; // 20 GB cumulative
const IMPORT_MAX_ENTRIES = 10_000;
export async function importBundleArchive(
stream: Readable,
): Promise<{ bundleId: string; version: string; models: string[] }> {
const stagingId = `import-${randomUUID()}`;
const stagingDir = join(AI_DIR, stagingId);
if (!acquireInstallLock("__import__")) {
throw new ImportLockError("Another install or import is already in progress");
}
try {
mkdirSync(stagingDir, { recursive: true });
// Extract with safety guards
let cumulativeBytes = 0;
let entryCount = 0;
await new Promise<void>((res, rej) => {
const extractor = tar.extract({
cwd: stagingDir,
strip: 0,
filter: (entryPath, entry) => {
// Reject non-file entry types (symlinks, hardlinks, devices, FIFOs, etc.)
if ("type" in entry && entry.type !== "File" && entry.type !== "Directory") {
rej(new ImportValidationError(`Unsupported entry type ${entry.type}: ${entryPath}`));
return false;
}
// Reject absolute paths and path traversal
if (entryPath.startsWith("/") || entryPath.split("/").includes("..")) {
rej(new ImportValidationError(`Blocked unsafe archive entry: ${entryPath}`));
return false;
}
entryCount++;
if (entryCount > IMPORT_MAX_ENTRIES) {
rej(new ImportValidationError(`Archive exceeds ${IMPORT_MAX_ENTRIES} entry limit`));
return false;
}
// Track cumulative size from tar headers (avoids consuming
// entry data which would prevent extraction to disk)
const entrySize = "size" in entry ? (entry.size as number) : 0;
cumulativeBytes += entrySize;
if (cumulativeBytes > IMPORT_MAX_BYTES) {
rej(new ImportValidationError("Archive exceeds 20 GB cumulative size limit"));
return false;
}
return true;
},
});
stream.pipe(extractor);
extractor.on("finish", () => res());
extractor.on("error", rej);
stream.on("error", rej);
});
// Read and validate bundle.json
const bundlePath = join(stagingDir, "bundle.json");
if (!existsSync(bundlePath)) {
throw new ImportValidationError("Archive is missing bundle.json at root");
}
let descriptor: BundleDescriptor;
try {
descriptor = JSON.parse(readFileSync(bundlePath, "utf-8")) as BundleDescriptor;
} catch {
throw new ImportValidationError("bundle.json is not valid JSON");
}
if (!descriptor.bundleId || typeof descriptor.bundleId !== "string") {
throw new ImportValidationError("bundle.json: bundleId must be a non-empty string");
}
if (!descriptor.version || typeof descriptor.version !== "string") {
throw new ImportValidationError("bundle.json: version must be a non-empty string");
}
if (
!Array.isArray(descriptor.models) ||
!descriptor.models.every((m) => typeof m === "string")
) {
throw new ImportValidationError("bundle.json: models must be an array of strings");
}
// Validate individual model paths against traversal
for (const model of descriptor.models) {
if (!model || model.includes("..") || model.startsWith("/") || model.includes("\\")) {
throw new ImportValidationError(`bundle.json: invalid model path "${model}"`);
}
}
// Validate bundleId against the manifest
const manifest = readManifest();
if (!manifest) {
throw new ImportValidationError("Feature manifest not found; cannot validate bundle");
}
if (!manifest.bundles[descriptor.bundleId]) {
throw new ImportValidationError(
`Unknown bundleId "${descriptor.bundleId}"; not in feature manifest`,
);
}
// Move models/* into MODELS_DIR
const stagingModels = join(stagingDir, "models");
if (existsSync(stagingModels)) {
mkdirSync(MODELS_DIR, { recursive: true });
moveTreeRecursive(stagingModels, MODELS_DIR);
}
// Move site-packages/* into venv site-packages
const stagingSitePackages = join(stagingDir, "site-packages");
if (existsSync(stagingSitePackages)) {
const venvPath = process.env.PYTHON_VENV_PATH || join(AI_DIR, "venv");
let sitePackagesDir = "";
const libDir = join(venvPath, "lib");
if (existsSync(libDir)) {
const pyDirs = readdirSync(libDir).filter((d) => d.startsWith("python"));
if (pyDirs.length > 0) {
sitePackagesDir = join(libDir, pyDirs[0], "site-packages");
}
}
if (sitePackagesDir && existsSync(sitePackagesDir)) {
moveTreeRecursive(stagingSitePackages, sitePackagesDir);
}
}
// Apply fixups (NCCL wheel) if present
const stagingFixups = join(stagingDir, "fixups");
if (existsSync(stagingFixups)) {
const wheels = readdirSync(stagingFixups).filter((f) => f.endsWith(".whl"));
if (wheels.length > 0) {
const venvPython = `${process.env.PYTHON_VENV_PATH || join(AI_DIR, "venv")}/bin/python3`;
for (const wheel of wheels) {
try {
execFileSync(
venvPython,
[
"-m",
"pip",
"install",
"--no-index",
`--find-links=${stagingFixups}`,
wheel.split("-")[0],
],
{ stdio: "ignore", timeout: 30_000 },
);
} catch {
// Non-fatal
}
}
}
}
markInstalled(descriptor.bundleId, descriptor.version, descriptor.models);
return {
bundleId: descriptor.bundleId,
version: descriptor.version,
models: descriptor.models,
};
} finally {
// Clean up staging dir (best-effort)
try {
rmSync(stagingDir, { recursive: true, force: true });
} catch {
// staging cleanup is best-effort
}
releaseInstallLock();
}
}
/** Recursively move entries from src into dest, merging directories. */
function moveTreeRecursive(src: string, dest: string): void {
const entries = readdirSync(src, { withFileTypes: true });
for (const entry of entries) {
const srcPath = join(src, entry.name);
const destPath = join(dest, entry.name);
if (entry.isDirectory()) {
mkdirSync(destPath, { recursive: true });
moveTreeRecursive(srcPath, destPath);
} else {
mkdirSync(dirname(destPath), { recursive: true });
try {
renameSync(srcPath, destPath);
} catch (err: unknown) {
// EXDEV: cross-device link (staging on different fs than MODELS_DIR).
// Staging is under AI_DIR so same-fs is expected, but fall back to
// copy+unlink for robustness (e.g. /tmp overlay mounts).
if ((err as NodeJS.ErrnoException).code === "EXDEV") {
copyFileSync(srcPath, destPath);
unlinkSync(srcPath);
} else {
throw err;
}
}
}
}
}
// ── Import-specific error types ────────────────────────────────────────
export class ImportValidationError extends Error {
constructor(message: string) {
super(message);
this.name = "ImportValidationError";
}
}
export class ImportLockError extends Error {
constructor(message: string) {
super(message);
this.name = "ImportLockError";
}
}