Something went wrong. Try again.
Monorepo for Aesthetic.Computer aesthetic.computer
Something went wrong. Try again.
JavaScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156#!/usr/bin/env node// silo - data & storage dashboard for aesthetic.computer
import "dotenv/config";import express from "express";import https from "https";import http from "http";import fs from "fs";import path from "path";import { fileURLToPath } from "url";import { execFile } from "child_process";import { MongoClient, ObjectId } from "mongodb";import { createClient } from "redis";import { S3Client, ListObjectsV2Command, PutObjectCommand } from "@aws-sdk/client-s3";import { WebSocketServer } from "ws";import { IgApiClient, IgCheckpointError, IgLoginTwoFactorRequiredError, IgLoginBadPasswordError,} from "instagram-private-api";import { ingestAll as ingestBluesky, ingestFromActor as ingestBlueskyActor, getConfiguredSources as getBlueskySources } from "./bluesky-ingest.mjs";
const __dirname = path.dirname(fileURLToPath(import.meta.url));
const app = express();const PORT = process.env.PORT || 3003;const dev = process.env.NODE_ENV === "development";const SERVER_START_TIME = Date.now();
// --- Auth0 ---const AUTH0_DOMAIN = process.env.AUTH0_DOMAIN || "aesthetic.us.auth0.com";const AUTH0_CLIENT_ID = process.env.AUTH0_CLIENT_ID || "";const AUTH0_CUSTOM_DOMAIN = "hi.aesthetic.computer";const ADMIN_SUB = process.env.ADMIN_SUB || "";const PUBLISH_SECRET = process.env.PUBLISH_SECRET || "";const AUTH_CACHE_TTL = 300_000;const authCache = new Map();
// --- Activity Log ---const activityLog = [];const MAX_LOG = 200;let wss = null;
// --- Firehose (MongoDB change stream) ---let changeStream = null;const firehoseThrottle = { count: 0, resetAt: 0 };const FIREHOSE_MAX_PER_SEC = 50;
// Dedup: suppress rapid updates to the same document within a short windowconst FIREHOSE_DEDUP_MS = 2000;const firehoseDedup = new Map(); // "coll:docId" → timestamp// Clean stale dedup entries every 30s to prevent memory leaksetInterval(() => { const cutoff = Date.now() - FIREHOSE_DEDUP_MS * 2; for (const [key, ts] of firehoseDedup) { if (ts < cutoff) firehoseDedup.delete(key); }}, 30_000);
// --- Instagram (persistent session) ---let igClient = null;let igLoggedIn = false;let igSessionUsername = null;let igLastUsed = null;let igChallengeInProgress = false;
// --- TikTok (OAuth token-based) ---const TIKTOK_CLIENT_KEY = process.env.TIKTOK_CLIENT_KEY || "";const TIKTOK_CLIENT_SECRET = process.env.TIKTOK_CLIENT_SECRET || "";const TIKTOK_REDIRECT_URI = process.env.TIKTOK_REDIRECT_URI || "https://silo.aesthetic.computer/api/tiktok/callback";let tiktokToken = null; // { access_token, refresh_token, open_id, expires_at, username }let tiktokLastUsed = null;
// --- Collection Categories ---const COLLECTION_CATEGORIES = { "users": "identity", "@handles": "identity", "verifications": "identity", "paintings": "content", "pieces": "content", "kidlisp": "content", "tapes": "content", "moods": "content", "chat-system": "communication", "chat-clock": "communication", "chat-sotce": "communication", "boots": "system", "piece-runs": "system", "kidlisp-logs": "system", "oven-bakes": "system", "_firehose": "system", "insta-sessions": "system", "tiktok-sessions": "system",};const CATEGORY_META = { identity: { label: "Users & Identity", color: "#48f" }, content: { label: "Content", color: "#a6f" }, communication: { label: "Communication", color: "#4a4" }, system: { label: "System & Logs", color: "#f80" }, other: { label: "Other", color: "#888" },};
// --- Handle Cache (Auth0 sub → handle) ---const handleCache = new Map();
async function loadHandleCache() { if (!db) return; try { const docs = await db.collection("@handles").find({}).toArray(); for (const doc of docs) { if (doc._id && doc.handle) handleCache.set(String(doc._id), doc.handle); } log("info", `handle cache loaded: ${handleCache.size} entries`); } catch (err) { log("error", `handle cache load failed: ${err.message}`); }}
function resolveHandle(str) { if (!str || typeof str !== "string") return null; if (/^(auth0|google-oauth2|apple|windowslive|github)\|/.test(str)) { return handleCache.get(str) || null; } return str;}
function log(type, msg) { const entry = { time: new Date().toISOString(), type, msg }; activityLog.unshift(entry); if (activityLog.length > MAX_LOG) activityLog.pop(); if (wss?.clients) { const data = JSON.stringify({ logEntry: entry }); wss.clients.forEach((c) => c.readyState === 1 && c.send(data)); } const prefix = { info: " ", warn: "! ", error: "x " }[type] || " "; console.log(`${prefix}${msg}`);}
log("info", "silo starting...");
// Parse a MongoDB URI into a display-friendly host string (no credentials)function parseMongoHost(uri) { if (!uri) return "not set"; try { const u = new URL(uri.replace("mongodb+srv://", "https://").replace("mongodb://", "http://")); return u.hostname + (u.port ? ":" + u.port : ""); } catch { return uri.split("@").pop()?.split("/")[0] || uri; }}
// --- MongoDB (primary) ---let mongoClient, db;
let mongoRetryTimer = null;
// Retries, because a single failed connect used to be permanent: db stayed null// for the life of the process, every mongo-backed route stayed broken, and silo// went on answering 200. That is how the 2026-08-09 credential mismatch hid —// mongod was healthy the whole time and one restart would have recovered it.async function connectMongo(attempt = 0) { const uri = process.env.MONGODB_CONNECTION_STRING; if (!uri) { log("error", "MONGODB_CONNECTION_STRING not set"); return; } try { if (mongoClient) await mongoClient.close().catch(() => {}); const isLocal = uri.includes("localhost") || uri.includes("127.0.0.1"); mongoClient = new MongoClient(uri, { ...(isLocal ? {} : { tls: true }), serverSelectionTimeoutMS: 10000, connectTimeoutMS: 10000, socketTimeoutMS: 45000, maxPoolSize: 5, }); await mongoClient.connect(); db = mongoClient.db(process.env.MONGODB_NAME || "aesthetic"); log("info", `mongo connected (${parseMongoHost(uri)})`); } catch (err) { db = null; // Backs off to 5 min and keeps going. Bad credentials will never heal on // their own, but the retry line in the log is a heartbeat that says which // failure it is, instead of one buried message at boot. const delay = Math.min(300000, 2000 * 2 ** attempt); log("error", `mongo connect failed: ${err.message} — retrying in ${Math.round(delay / 1000)}s`); clearTimeout(mongoRetryTimer); mongoRetryTimer = setTimeout(() => connectMongo(attempt + 1), delay); mongoRetryTimer.unref?.(); }}
// Extract a short summary from the full document (admin-only dashboard, safe to show)function firehoseSummary(coll, op, doc) { if (!doc || op === "delete") return null; // Try to resolve user handle from common fields, resolving Auth0 subs const rawHandle = doc.handle || doc.user; const resolved = rawHandle ? resolveHandle(String(rawHandle)) : null; // Also try resolving the doc _id (often an Auth0 sub in user-related collections) const resolvedId = doc._id ? resolveHandle(String(doc._id)) : null; const who = resolved ? `@${resolved}` : (resolvedId && resolvedId !== String(doc._id) ? `@${resolvedId}` : null); switch (coll) { case "chat-system": case "chat-clock": case "chat-sotce": return (who || "anon") + (doc.text ? ` "${doc.text.slice(0, 30)}"` : ""); case "@handles": return who || (doc.handle ? `@${doc.handle}` : null); case "users": { const name = doc.name ? resolveHandle(String(doc.name)) : null; const displayName = name && name !== doc.name ? `@${name}` : doc.name; return doc.email?.split("@")[0] || displayName || who || null; } case "paintings": return (who ? who + " " : "") + (doc.slug || doc.title || ""); case "tapes": return (who ? who + " " : "") + (doc.piece || ""); case "pieces": return doc.slug || doc.name || null; case "kidlisp": return (who ? who + " " : "") + (doc.code ? "$" + doc.code : "") + (doc.name ? " " + doc.name : ""); case "moods": return (who ? who + " " : "") + (doc.mood || ""); case "verifications": return who || null; case "boots": { const m = doc.meta || {}; const rawUser = m.user?.handle || m.user?.sub || m.user; const bootUser = rawUser ? (resolveHandle(String(rawUser)) || rawUser) : null; const who = bootUser ? `@${bootUser}` : "visitor"; const status = doc.status || ""; const statusTag = status !== "started" ? ` [${status}]` : ""; // Rich info: browser, device, referrer, geo const browser = m.browser || ""; const device = m.mobile ? "mobile" : ""; const geo = doc.server?.country || ""; const referrer = m.referrer ? ` via ${m.referrer.replace(/^https?:\/\//, "").split("/")[0]}` : ""; const parts = [who, m.host + (m.path || "/"), browser, device, geo].filter(Boolean); return parts.join(" ") + referrer + statusTag; } case "kidlisp-logs": { const effect = doc.effect || ""; const type = doc.type || ""; const gpu = doc.device?.gpu?.renderer || ""; const ua = doc.device?.mobile ? "mobile" : "desktop"; const country = doc.server?.country || ""; return [type, effect, gpu, ua, country].filter(Boolean).join(" "); } case "oven-bakes": return doc.status || null; default: { return who || doc.slug || doc.name || doc.email?.split("@")[0] || null; } }}
const FIREHOSE_COLLECTION = "_firehose";const FIREHOSE_TTL_DAYS = 30;
async function ensureFirehoseCollection() { if (!db) return; try { const col = db.collection(FIREHOSE_COLLECTION); // TTL index: auto-delete events older than 30 days await col.createIndex({ time: 1 }, { expireAfterSeconds: FIREHOSE_TTL_DAYS * 86400 }); // Index for recent queries await col.createIndex({ ns: 1, time: -1 }); log("info", `_firehose collection ready (${FIREHOSE_TTL_DAYS}d TTL)`); } catch (err) { log("error", `_firehose index setup: ${err.message}`); }}
function startFirehose() { if (!db) return; const firehoseCol = db.collection(FIREHOSE_COLLECTION); try { changeStream = db.watch( [{ $match: { operationType: { $in: ["insert", "update", "replace", "delete"] }, "ns.coll": { $ne: FIREHOSE_COLLECTION }, // avoid infinite loop }, }], { fullDocument: "updateLookup" } ); log("info", "firehose change stream started");
changeStream.on("change", (change) => { const now = Date.now(); if (now > firehoseThrottle.resetAt) { firehoseThrottle.count = 0; firehoseThrottle.resetAt = now + 1000; } if (firehoseThrottle.count >= FIREHOSE_MAX_PER_SEC) return; firehoseThrottle.count++;
const coll = change.ns?.coll || "unknown"; const op = change.operationType; const docId = change.documentKey?._id?.toString() || null;
// Keep handle cache fresh when @handles collection changes if (coll === "@handles" && change.fullDocument) { const hDoc = change.fullDocument; if (hDoc._id && hDoc.handle) handleCache.set(String(hDoc._id), hDoc.handle); }
// Dedup: suppress rapid updates/replaces to the same document // (e.g. boot telemetry writes start→log→complete to same bootId) if (docId && (op === "update" || op === "replace")) { const dedupKey = `${coll}:${docId}`; const prev = firehoseDedup.get(dedupKey); if (prev && now - prev < FIREHOSE_DEDUP_MS) { firehoseDedup.set(dedupKey, now); return; // suppress duplicate } firehoseDedup.set(dedupKey, now); }
const event = { ns: coll, op, time: new Date(now), docId, summary: firehoseSummary(coll, op, change.fullDocument), };
// Persist to MongoDB (fire-and-forget) firehoseCol.insertOne(event).catch(() => {});
if (wss?.clients) { const data = JSON.stringify({ firehose: { ...event, time: now } }); wss.clients.forEach((c) => c.readyState === 1 && c.send(data)); } });
changeStream.on("error", (err) => { log("error", `firehose stream error: ${err.message}`); changeStream = null; setTimeout(startFirehose, 5000); }); } catch (err) { log("error", `firehose setup failed: ${err.message}`); }}
// --- MongoDB Atlas (for comparison) ---let atlasClient, atlasDb;
async function connectAtlas() { const uri = process.env.ATLAS_CONNECTION_STRING; if (!uri) { log("info", "ATLAS_CONNECTION_STRING not set, skipping atlas"); return; } // Skip if same as primary if (uri === process.env.MONGODB_CONNECTION_STRING) { log("info", "atlas URI same as primary, sharing connection"); atlasClient = mongoClient; atlasDb = db; return; } try { atlasClient = new MongoClient(uri, { tls: true, serverSelectionTimeoutMS: 10000, connectTimeoutMS: 10000, socketTimeoutMS: 45000, maxPoolSize: 3, }); await atlasClient.connect(); atlasDb = atlasClient.db(process.env.MONGODB_NAME || "aesthetic"); log("info", "atlas connected (comparison)"); } catch (err) { log("error", `atlas connect failed: ${err.message}`); }}
// --- Redis ---let redisClient;
async function connectRedis() { const url = process.env.REDIS_CONNECTION_STRING; if (!url) { log("info", "REDIS_CONNECTION_STRING not set, skipping redis"); return; } try { redisClient = createClient({ url }); redisClient.on("error", (err) => log("error", `redis error: ${err.message}`)); await redisClient.connect(); log("info", "redis connected"); } catch (err) { log("error", `redis connect failed: ${err.message}`); }}
// --- Redis Subscriber (for fairy:point from session server) ---let redisSub;
async function connectRedisSub() { const url = process.env.REDIS_CONNECTION_STRING; if (!url) return; try { redisSub = createClient({ url }); redisSub.on("error", (err) => log("error", `redis sub error: ${err.message}`)); await redisSub.connect();
await redisSub.subscribe("fairy:point", (message) => { if (wss?.clients) { try { const payload = JSON.stringify({ fairyPoint: JSON.parse(message) }); wss.clients.forEach((c) => c.readyState === 1 && c.send(payload)); } catch {} } });
log("info", "redis subscriber connected (fairy:point)"); } catch (err) { log("error", `redis sub connect failed: ${err.message}`); }}
// --- S3 (DigitalOcean Spaces) ---const s3 = new S3Client({ endpoint: process.env.SPACES_ENDPOINT, region: "us-east-1", credentials: { accessKeyId: process.env.SPACES_KEY || "", secretAccessKey: process.env.SPACES_SECRET || "", }, forcePathStyle: false,});const BUCKETS = (process.env.BUCKETS || "").split(",").filter(Boolean);
// --- Cached Stats ---let cachedDbStats = null, cachedStorageStats = null, cachedAtlasStats = null;let dbStatsAge = 0, storageStatsAge = 0, atlasStatsAge = 0;const DB_CACHE_TTL = 30_000;const STORAGE_CACHE_TTL = 300_000;
async function getCollectionStats(database, label) { if (!database) return null; try { const collections = await database.listCollections().toArray(); const stats = []; for (const col of collections) { try { const cs = await database.command({ collStats: col.name }); stats.push({ name: col.name, count: cs.count || 0, size: cs.size || 0, storageSize: cs.storageSize || 0, indexSize: cs.totalIndexSize || 0, category: COLLECTION_CATEGORIES[col.name] || "other", }); } catch { const count = await database.collection(col.name).estimatedDocumentCount(); stats.push({ name: col.name, count, size: 0, storageSize: 0, indexSize: 0, category: COLLECTION_CATEGORIES[col.name] || "other", }); } } stats.sort((a, b) => b.count - a.count); return stats; } catch (err) { log("error", `${label} stats error: ${err.message}`); return null; }}
async function getDbStats() { if (cachedDbStats && Date.now() - dbStatsAge < DB_CACHE_TTL) return cachedDbStats; const stats = await getCollectionStats(db, "primary"); if (stats) { cachedDbStats = stats; dbStatsAge = Date.now(); } return cachedDbStats;}
async function getAtlasStats() { if (cachedAtlasStats && Date.now() - atlasStatsAge < DB_CACHE_TTL) return cachedAtlasStats; const stats = await getCollectionStats(atlasDb, "atlas"); if (stats) { cachedAtlasStats = stats; atlasStatsAge = Date.now(); } return cachedAtlasStats;}
async function getDbSize(database) { if (!database) return null; try { const stats = await database.stats(); return { dataSize: stats.dataSize, storageSize: stats.storageSize, indexSize: stats.indexSize }; } catch { return null; }}
async function getStorageStats() { if (cachedStorageStats && Date.now() - storageStatsAge < STORAGE_CACHE_TTL) return cachedStorageStats; try { const results = []; for (const bucket of BUCKETS) { let totalSize = 0, totalObjects = 0, continuationToken; do { const resp = await s3.send(new ListObjectsV2Command({ Bucket: bucket, ContinuationToken: continuationToken, })); if (resp.Contents) { totalObjects += resp.Contents.length; totalSize += resp.Contents.reduce((sum, obj) => sum + (obj.Size || 0), 0); } continuationToken = resp.IsTruncated ? resp.NextContinuationToken : undefined; } while (continuationToken); results.push({ bucket, objects: totalObjects, bytes: totalSize, gb: (totalSize / 1e9).toFixed(2) }); } cachedStorageStats = results; storageStatsAge = Date.now(); return results; } catch (err) { log("error", `storage stats error: ${err.message}`); return cachedStorageStats || []; }}
async function getRedisStats() { if (!redisClient?.isReady) return null; try { const info = await redisClient.info(); const parse = (section, key) => { const m = info.match(new RegExp(`${key}:(.+)`)); return m ? m[1].trim() : null; }; return { connected: true, version: parse("server", "redis_version"), usedMemory: parse("memory", "used_memory_human"), peakMemory: parse("memory", "used_memory_peak_human"), totalKeys: parseInt(parse("keyspace", "keys") || "0") || await redisClient.dbSize(), connectedClients: parseInt(parse("clients", "connected_clients") || "0"), uptimeSeconds: parseInt(parse("server", "uptime_in_seconds") || "0"), hitRate: (() => { const hits = parseInt(parse("stats", "keyspace_hits") || "0"); const misses = parseInt(parse("stats", "keyspace_misses") || "0"); return hits + misses > 0 ? ((hits / (hits + misses)) * 100).toFixed(1) : "n/a"; })(), }; } catch (err) { log("error", `redis stats error: ${err.message}`); return { connected: false, error: err.message }; }}
// --- Auth ---async function validateToken(authorization) { if (!authorization) return null; const cached = authCache.get(authorization); if (cached && Date.now() - cached.timestamp < AUTH_CACHE_TTL) return cached; try { const resp = await fetch(`https://${AUTH0_DOMAIN}/userinfo`, { headers: { Authorization: authorization }, }); if (!resp.ok) return null; const user = await resp.json(); let handle = null, isAdmin = false; if (db) { const doc = await db.collection("@handles").findOne({ _id: user.sub }); handle = doc?.handle; isAdmin = !!(user.email_verified && handle === "jeffrey" && user.sub === ADMIN_SUB); } const result = { user, handle, isAdmin, timestamp: Date.now() }; authCache.set(authorization, result); if (isAdmin) log("info", `auth: @${handle} verified`); return result; } catch (err) { log("error", `token validation failed: ${err.message}`); return null; }}
async function requireAdmin(req, res, next) { // Allow publish secret for CI/CLI scripts if (PUBLISH_SECRET && req.headers["x-publish-secret"] === PUBLISH_SECRET) { req.auth = { handle: "cli", isAdmin: true }; return next(); } const auth = await validateToken(req.headers.authorization); if (!auth) return res.status(401).json({ error: "Unauthorized" }); if (!auth.isAdmin) return res.status(403).json({ error: "Admin access required" }); req.auth = auth; next();}
// --- Express Routes ---app.use(express.json());
app.get("/auth/config", (req, res) => { res.json({ domain: AUTH0_CUSTOM_DOMAIN, clientId: AUTH0_CLIENT_ID });});
app.get("/auth/me", async (req, res) => { const auth = await validateToken(req.headers.authorization); if (!auth) return res.status(401).json({ error: "Unauthorized" }); res.json({ handle: auth.handle, isAdmin: auth.isAdmin });});
// --- Instagram Session Management ---async function saveInstaSession(ig, username) { if (!db) return; try { const serialized = await ig.state.serialize(); delete serialized.constants; delete serialized.supportedCapabilities; await db.collection("insta-sessions").updateOne( { _id: username }, { $set: { _id: username, state: serialized, updatedAt: new Date() } }, { upsert: true }, ); } catch (e) { log("warn", `insta session save failed: ${e.message}`); }}
async function getInstaClient() { if (igClient && igLoggedIn) { igLastUsed = Date.now(); return igClient; }
if (!db) throw new Error("Database not connected");
// Get credentials from MongoDB secrets const secrets = await db.collection("secrets").findOne({ _id: "instagram" }); if (!secrets?.username || !secrets?.password) { throw new Error("Instagram credentials not configured in MongoDB secrets"); } const username = secrets.username;
igClient = new IgApiClient(); igClient.state.generateDevice(username); igSessionUsername = username;
// Save session after every Instagram API request igClient.request.end$.subscribe(async () => { await saveInstaSession(igClient, username); });
// Try to restore session from MongoDB try { const sessionDoc = await db .collection("insta-sessions") .findOne({ _id: username }); if (sessionDoc?.state) { await igClient.state.deserialize(sessionDoc.state); await igClient.account.currentUser(); // verify session is alive igLoggedIn = true; igLastUsed = Date.now(); log("info", `insta session restored for @${username}`); return igClient; } } catch { log("info", "insta session restore failed, need fresh login via dashboard"); }
throw new Error( "Instagram session not available. Log in via silo dashboard.", );}
function formatInstaCount(n) { if (n >= 1_000_000) return (n / 1_000_000).toFixed(1) + "m"; if (n >= 1000) return (n / 1000).toFixed(1) + "k"; return String(n);}
// --- Instagram Public Routes (CORS-enabled, no auth) ---app.use("/insta", (req, res, next) => { res.header("Access-Control-Allow-Origin", "*"); res.header("Access-Control-Allow-Methods", "GET, OPTIONS"); res.header("Access-Control-Allow-Headers", "Content-Type"); if (req.method === "OPTIONS") return res.sendStatus(200); next();});
app.get("/insta", async (req, res) => { const { action, username, max } = req.query;
if (!action || !username) { return res.status(400).json({ error: "action and username required" }); }
const cleanUsername = username.replace(/^@/, "").toLowerCase();
try { const ig = await getInstaClient();
if (action === "profile") { const userId = await ig.user.getIdByUsername(cleanUsername); const info = await ig.user.info(userId); return res.json({ username: info.username, fullName: info.full_name || "", bio: info.biography || "", profilePicUrl: info.profile_pic_url, mediaCount: info.media_count, followerCount: info.follower_count, followingCount: info.following_count, isVerified: info.is_verified, isPrivate: info.is_private, externalUrl: info.external_url || null, mediaCountFormatted: formatInstaCount(info.media_count), followerCountFormatted: formatInstaCount(info.follower_count), followingCountFormatted: formatInstaCount(info.following_count), }); }
if (action === "feed") { const maxPosts = parseInt(max || "18", 10); const userId = await ig.user.getIdByUsername(cleanUsername); const feed = ig.feed.user(userId); const items = await feed.items(); const posts = items.slice(0, maxPosts).map((item) => ({ id: item.id, shortcode: item.code, mediaType: item.media_type, caption: item.caption?.text || "", likeCount: item.like_count || 0, commentCount: item.comment_count || 0, timestamp: item.taken_at, width: item.original_width, height: item.original_height, thumbnailUrl: item.image_versions2?.candidates?.[ item.image_versions2.candidates.length - 1 ]?.url || null, likeCountFormatted: formatInstaCount(item.like_count || 0), commentCountFormatted: formatInstaCount(item.comment_count || 0), })); return res.json({ posts, hasMore: feed.isMoreAvailable() }); }
return res.status(400).json({ error: `Unknown action: ${action}` }); } catch (err) { log("error", `insta public API error: ${err.message}`);
if (err.message.includes("not configured") || err.message.includes("not available")) { return res.status(503).json({ error: err.message }); } if (err.message.includes("User not found") || err.name === "IgExactUserNotFoundError") { return res.status(404).json({ error: `User @${cleanUsername} not found` }); }
// Session may have expired igLoggedIn = false; igClient = null; return res.status(500).json({ error: "Instagram API error" }); }});
// --- TikTok Session Management ---async function saveTiktokSession(tokenData) { if (!db) return; try { await db.collection("tiktok-sessions").updateOne( { _id: "default" }, { $set: { ...tokenData, updatedAt: new Date() } }, { upsert: true }, ); } catch (e) { log("warn", `tiktok session save failed: ${e.message}`); }}
async function loadTiktokSession() { if (tiktokToken) return tiktokToken; if (!db) return null; try { const doc = await db.collection("tiktok-sessions").findOne({ _id: "default" }); if (doc?.access_token) { tiktokToken = doc; log("info", `tiktok session restored for @${doc.username || doc.open_id}`); return tiktokToken; } } catch { log("info", "tiktok session restore failed"); } return null;}
async function refreshTiktokToken() { if (!tiktokToken?.refresh_token) return false; try { const resp = await fetch("https://open.tiktokapis.com/v2/oauth/token/", { method: "POST", headers: { "Content-Type": "application/x-www-form-urlencoded" }, body: new URLSearchParams({ client_key: TIKTOK_CLIENT_KEY, client_secret: TIKTOK_CLIENT_SECRET, grant_type: "refresh_token", refresh_token: tiktokToken.refresh_token, }), }); const data = await resp.json(); if (data.access_token) { tiktokToken.access_token = data.access_token; tiktokToken.expires_at = Date.now() + data.expires_in * 1000; if (data.refresh_token) tiktokToken.refresh_token = data.refresh_token; await saveTiktokSession(tiktokToken); log("info", "tiktok token refreshed"); return true; } log("warn", `tiktok refresh failed: ${JSON.stringify(data)}`); return false; } catch (e) { log("error", `tiktok refresh error: ${e.message}`); return false; }}
async function tiktokApiFetch(endpoint, fields) { let token = await loadTiktokSession(); if (!token) throw new Error("TikTok not connected. Authorize via silo dashboard.");
// Refresh if expired or about to expire (5 min buffer) if (token.expires_at && Date.now() > token.expires_at - 300_000) { const ok = await refreshTiktokToken(); if (!ok) { tiktokToken = null; throw new Error("TikTok token expired and refresh failed. Re-authorize via dashboard."); } token = tiktokToken; }
const url = `https://open.tiktokapis.com/v2/${endpoint}${fields ? "?fields=" + fields : ""}`; const resp = await fetch(url, { headers: { Authorization: `Bearer ${token.access_token}` }, }); const data = await resp.json(); if (data.error?.code === "access_token_invalid") { // Try one refresh const ok = await refreshTiktokToken(); if (ok) { const retry = await fetch(url, { headers: { Authorization: `Bearer ${tiktokToken.access_token}` }, }); return retry.json(); } tiktokToken = null; throw new Error("TikTok token invalid. Re-authorize via dashboard."); } tiktokLastUsed = Date.now(); return data;}
function formatCount(n) { if (n >= 1_000_000) return (n / 1_000_000).toFixed(1) + "m"; if (n >= 1000) return (n / 1000).toFixed(1) + "k"; return String(n);}
// --- TikTok Public Routes (CORS-enabled, no auth) ---app.use("/tiktok", (req, res, next) => { res.header("Access-Control-Allow-Origin", "*"); res.header("Access-Control-Allow-Methods", "GET, OPTIONS"); res.header("Access-Control-Allow-Headers", "Content-Type"); if (req.method === "OPTIONS") return res.sendStatus(200); next();});
app.get("/tiktok", async (req, res) => { const { action } = req.query; if (!action) return res.status(400).json({ error: "action required (profile, videos)" });
try { if (action === "profile") { const data = await tiktokApiFetch( "user/info/", "open_id,display_name,avatar_url,bio_description,follower_count,following_count,likes_count,video_count,is_verified,username,profile_deep_link", ); const u = data.data?.user || {}; return res.json({ username: u.username || u.display_name, displayName: u.display_name, bio: u.bio_description || "", avatarUrl: u.avatar_url, followerCount: u.follower_count, followingCount: u.following_count, likesCount: u.likes_count, videoCount: u.video_count, isVerified: u.is_verified, profileUrl: u.profile_deep_link, followerCountFormatted: formatCount(u.follower_count || 0), followingCountFormatted: formatCount(u.following_count || 0), likesCountFormatted: formatCount(u.likes_count || 0), videoCountFormatted: formatCount(u.video_count || 0), }); }
if (action === "videos") { const max = parseInt(req.query.max || "20", 10); const data = await tiktokApiFetch( "video/list/", "id,title,cover_image_url,share_url,view_count,like_count,comment_count,share_count,create_time,duration", ); const videos = (data.data?.videos || []).slice(0, max).map((v) => ({ id: v.id, title: v.title || "", coverUrl: v.cover_image_url, shareUrl: v.share_url, viewCount: v.view_count, likeCount: v.like_count, commentCount: v.comment_count, shareCount: v.share_count, createTime: v.create_time, duration: v.duration, viewCountFormatted: formatCount(v.view_count || 0), likeCountFormatted: formatCount(v.like_count || 0), })); return res.json({ videos, cursor: data.data?.cursor, hasMore: data.data?.has_more }); }
return res.status(400).json({ error: `Unknown action: ${action}` }); } catch (err) { log("error", `tiktok public API error: ${err.message}`); if (err.message.includes("not connected") || err.message.includes("Re-authorize")) { return res.status(503).json({ error: err.message }); } return res.status(500).json({ error: "TikTok API error" }); }});
// --- TikTok OAuth Callback (must be before requireAdmin) ---app.get("/api/tiktok/callback", async (req, res) => { const { code, error: oauthError } = req.query; if (oauthError || !code) { return res.send(`<html><body><h2>TikTok auth failed</h2><p>${oauthError || "no code"}</p><script>setTimeout(()=>window.close(),2000)</script></body></html>`); }
try { const resp = await fetch("https://open.tiktokapis.com/v2/oauth/token/", { method: "POST", headers: { "Content-Type": "application/x-www-form-urlencoded" }, body: new URLSearchParams({ client_key: TIKTOK_CLIENT_KEY, client_secret: TIKTOK_CLIENT_SECRET, code, grant_type: "authorization_code", redirect_uri: TIKTOK_REDIRECT_URI, }), }); const data = await resp.json(); if (!data.access_token) { log("error", `tiktok oauth token exchange failed: ${JSON.stringify(data)}`); return res.send(`<html><body><h2>Token exchange failed</h2><pre>${JSON.stringify(data, null, 2)}</pre></body></html>`); }
// Fetch username let username = ""; try { const userResp = await fetch("https://open.tiktokapis.com/v2/user/info/?fields=username,display_name", { headers: { Authorization: `Bearer ${data.access_token}` }, }); const userData = await userResp.json(); username = userData.data?.user?.username || userData.data?.user?.display_name || ""; } catch {}
tiktokToken = { access_token: data.access_token, refresh_token: data.refresh_token, open_id: data.open_id, expires_at: Date.now() + data.expires_in * 1000, refresh_expires_at: data.refresh_expires_in ? Date.now() + data.refresh_expires_in * 1000 : null, scope: data.scope, username, }; tiktokLastUsed = Date.now(); await saveTiktokSession(tiktokToken); log("info", `tiktok connected for @${username || data.open_id}`);
return res.send(`<html><body style="font-family:monospace;background:#111;color:#6c6;display:flex;align-items:center;justify-content:center;height:100vh"><h2>tiktok connected ✓</h2><script>setTimeout(()=>window.close(),1500)</script></body></html>`); } catch (err) { log("error", `tiktok callback error: ${err.message}`); return res.send(`<html><body><h2>Error</h2><p>${err.message}</p></body></html>`); }});
app.use("/api", requireAdmin);
// Provenance rides along with liveness, in the fleet-wide shape from// shared/provenance.mjs: { service, sha, startedAt }. silo is deployed as loose// files with no git repo on the box, so it cannot derive its own commit — the// deploy stamps AC_GIT_SHA into .env. A null sha is reported honestly rather// than guessed, so "unknown" never passes for "current".app.get("/health", async (req, res) => { const mongoOk = db ? await db.admin().ping().then(() => true).catch(() => false) : false; // A service that cannot reach its database is not healthy, and saying so in // the status code is the whole point: `mongo:false` inside a 200 body is a // sentence nothing was reading. Body shape is unchanged either way, so a // caller after provenance still gets the sha from a degraded silo. res.status(mongoOk ? 200 : 503).json({ status: mongoOk ? "ok" : "degraded", service: "silo", sha: process.env.AC_GIT_SHA?.trim() || null, // True when the deploy shipped working-tree files that are in no commit, so // the sha above names a commit that does not contain what is running. dirty: process.env.AC_GIT_DIRTY === "1" || undefined, startedAt: new Date(SERVER_START_TIME).toISOString(), uptime: Date.now() - SERVER_START_TIME, mongo: mongoOk, });});
// Which boxes are running which commit — the question no existing check asked.//// Reachability tells you a service answers. This tells you what it answers// WITH, by joining three facts nothing previously joined: knot's tip (what// "current" means), the GitHub mirror's tip (what the pull-based deploys can// actually reach), and each service's self-reported sha.//// Reports facts rather than verdicts. silo has no checkout, so it cannot say// how far behind a box is — only whether it matches. Callers holding a repo// (`ac-fleet`, `npm run doctor`) do the ancestry.const FLEET = [ { service: "session-server", health: "https://session-server.aesthetic.computer/health" }, { service: "silo", health: "https://silo.aesthetic.computer/health" }, { service: "oven", health: "https://oven.aesthetic.computer/health" }, { service: "lith", health: "https://aesthetic.computer/api/health" },];
function remoteTip(url) { return new Promise((resolve) => { execFile("git", ["ls-remote", url, "refs/heads/main"], { timeout: 15000 }, (err, out) => resolve(err ? null : (out.trim().split(/\s+/)[0] || null))); });}
app.get("/api/fleet/drift", async (req, res) => { const [knot, mirror] = await Promise.all([ remoteTip("https://knot.aesthetic.computer/aesthetic.computer/core"), remoteTip("https://github.com/whistlegraph/aesthetic-computer.git"), ]);
const services = await Promise.all(FLEET.map(async ({ service, health }) => { try { // Read the body before judging the status: a degraded service still // reports which commit it is running, and that is the question here. const r = await fetch(health, { signal: AbortSignal.timeout(8000) }); const body = await r.json().catch(() => null); const sha = typeof body?.sha === "string" ? body.sha : null; if (!sha) return { service, sha: null, note: `HTTP ${r.status}` }; return { service, sha, startedAt: body?.startedAt ?? null, degraded: r.ok ? undefined : `HTTP ${r.status}`, // Deliberately NOT called "current". With no checkout, exact match // against knot's tip is the most this can know, and reading that as // currency marks the whole fleet stale on any unrelated commit. Callers // holding the history (ac-fleet) decide what current means. matchesKnotTip: knot ? sha === knot : null, }; } catch (e) { return { service, sha: null, note: e.name === "TimeoutError" ? "timeout" : "unreachable" }; } }));
res.json({ knot, mirror, mirrorInSync: knot && mirror ? knot === mirror : null, services, checkedAt: new Date().toISOString(), });});
app.get("/api/overview", async (req, res) => { log("info", "overview requested"); const [collections, storage, redis, dbSize] = await Promise.all([ getDbStats(), getStorageStats(), getRedisStats(), getDbSize(db), ]); const find = (name) => collections?.find((c) => c.name === name)?.count || 0; res.json({ db: { users: find("users"), handles: find("@handles"), paintings: find("paintings"), pieces: find("pieces"), kidlisp: find("kidlisp"), moods: find("moods"), chatSystem: find("chat-system"), chatClock: find("chat-clock"), verifications: find("verifications"), totalCollections: collections?.length || 0, totalDocuments: collections?.reduce((s, c) => s + c.count, 0) || 0, size: dbSize, }, storage: { buckets: storage, totalGB: storage.reduce((s, b) => s + parseFloat(b.gb), 0).toFixed(2), }, redis, uptime: Date.now() - SERVER_START_TIME, });});
// 📡 Telemetry overview — GPU/KidLisp logs + boot logs statsapp.get("/api/telemetry", async (req, res) => { if (!db) return res.status(503).json({ error: "no db" }); try { const [klStats, bootStats, klRecent] = await Promise.all([ db.collection("kidlisp-logs").aggregate([ { $facet: { total: [{ $count: "n" }], byType: [{ $group: { _id: "$type", count: { $sum: 1 } } }], byEffect: [{ $group: { _id: "$effect", count: { $sum: 1 } } }], byGpu: [{ $group: { _id: "$device.gpu.renderer", count: { $sum: 1 } } }], size: [{ $group: { _id: null, avgSize: { $avg: { $bsonSize: "$$ROOT" } } } }], }} ]).toArray(), db.collection("boots").aggregate([ { $facet: { total: [{ $count: "n" }], byStatus: [{ $group: { _id: "$status", count: { $sum: 1 } } }], }} ]).toArray(), db.collection("kidlisp-logs") .find() .sort({ createdAt: -1 }) .limit(20) .toArray(), ]);
const kl = klStats[0] || {}; const bt = bootStats[0] || {}; const klTotal = kl.total?.[0]?.n || 0; const klAvgSize = kl.size?.[0]?.avgSize || 0;
res.json({ kidlispLogs: { total: klTotal, estimatedSizeKB: Math.round((klTotal * klAvgSize) / 1024), byType: kl.byType || [], byEffect: kl.byEffect || [], byGpu: kl.byGpu || [], recent: klRecent, }, boots: { total: bt.total?.[0]?.n || 0, byStatus: bt.byStatus || [], }, }); } catch (err) { res.status(500).json({ error: err.message }); }});
// Purge kidlisp-logs from siloapp.delete("/api/telemetry/kidlisp-logs", async (req, res) => { if (!db) return res.status(503).json({ error: "no db" }); try { const result = await db.collection("kidlisp-logs").deleteMany({}); log("info", `purged ${result.deletedCount} kidlisp-logs`); res.json({ ok: true, deleted: result.deletedCount }); } catch (err) { res.status(500).json({ error: err.message }); }});
app.get("/api/db/collections", async (req, res) => { log("info", "collections requested"); const stats = await getDbStats() || []; res.json({ collections: stats, categories: CATEGORY_META });});
app.get("/api/db/backups", async (req, res) => { log("info", "backups requested"); const backups = []; // Check local filesystem for mongodump backups const backupPath = process.env.BACKUP_PATH || "/var/backups/mongodb"; try { if (fs.existsSync(backupPath)) { const entries = fs.readdirSync(backupPath, { withFileTypes: true }); for (const entry of entries) { try { const fullPath = backupPath + "/" + entry.name; const stat = fs.statSync(fullPath); let totalSize = stat.size; // For directories (mongodump output), sum up contents if (entry.isDirectory()) { totalSize = 0; const walk = (dir) => { for (const f of fs.readdirSync(dir, { withFileTypes: true })) { const fp = dir + "/" + f.name; if (f.isDirectory()) walk(fp); else totalSize += fs.statSync(fp).size; } }; walk(fullPath); } backups.push({ name: entry.name, source: "local", size: totalSize, date: stat.mtime.toISOString(), isDirectory: entry.isDirectory(), }); } catch {} } } } catch (err) { log("error", `backup scan (local): ${err.message}`); } // Check S3 for backups const backupBucket = process.env.BACKUP_BUCKET; const backupPrefix = process.env.BACKUP_PREFIX || "backups/"; if (backupBucket) { try { let continuationToken; do { const resp = await s3.send(new ListObjectsV2Command({ Bucket: backupBucket, Prefix: backupPrefix, ContinuationToken: continuationToken, })); if (resp.Contents) { for (const obj of resp.Contents) { const name = obj.Key.replace(backupPrefix, ""); if (!name) continue; backups.push({ name, source: "s3", bucket: backupBucket, size: obj.Size, date: obj.LastModified?.toISOString(), }); } } continuationToken = resp.IsTruncated ? resp.NextContinuationToken : undefined; } while (continuationToken); } catch (err) { log("error", `backup scan (s3): ${err.message}`); } } backups.sort((a, b) => new Date(b.date) - new Date(a.date)); res.json({ backups, backupPath, backupBucket: backupBucket || null });});
app.get("/api/db/compare", async (req, res) => { log("info", "db compare requested"); const [primary, atlas, primarySize, atlasSize] = await Promise.all([ getDbStats(), getAtlasStats(), getDbSize(db), getDbSize(atlasDb), ]); const primaryHost = parseMongoHost(process.env.MONGODB_CONNECTION_STRING); const atlasHost = parseMongoHost(process.env.ATLAS_CONNECTION_STRING);
// Silo is now the source of truth; Atlas is deprecated const activeDb = "silo";
const allNames = new Set([ ...(primary || []).map(c => c.name), ...(atlas || []).map(c => c.name), ]); const comparison = [...allNames].sort().map(name => { const p = primary?.find(c => c.name === name); const a = atlas?.find(c => c.name === name); return { name, primary: p?.count ?? null, atlas: a?.count ?? null, synced: p?.count === a?.count, }; }); const allSynced = comparison.every(c => c.synced);
res.json({ primaryLabel: primaryHost, atlasLabel: atlasHost, activeDb, allSynced, primarySize, atlasSize, primaryConnected: !!db, atlasConnected: !!atlasDb, sameConnection: process.env.MONGODB_CONNECTION_STRING === process.env.ATLAS_CONNECTION_STRING, collections: comparison, });});
// Sync: copy all collections from Atlas → primarylet syncInProgress = false;app.post("/api/db/sync", async (req, res) => { if (syncInProgress) return res.status(409).json({ error: "Sync already in progress" }); if (!db || !atlasDb) return res.status(500).json({ error: "Both databases must be connected" }); if (atlasDb === db) return res.status(400).json({ error: "Primary and Atlas are the same connection" });
syncInProgress = true; log("info", "sync started: atlas -> primary");
try { const collections = await atlasDb.listCollections().toArray(); const results = []; for (const col of collections) { const name = col.name; const docs = await atlasDb.collection(name).find({}).toArray(); // Drop and re-insert await db.collection(name).deleteMany({}); if (docs.length > 0) { await db.collection(name).insertMany(docs); } results.push({ name, docs: docs.length }); log("info", `synced ${name}: ${docs.length} docs`); } // Invalidate cache cachedDbStats = null; dbStatsAge = 0; log("info", `sync complete: ${results.length} collections`); res.json({ ok: true, collections: results }); } catch (err) { log("error", `sync failed: ${err.message}`); res.status(500).json({ error: err.message }); } finally { syncInProgress = false; }});
app.get("/api/db/health", async (req, res) => { if (!db) return res.json({ connected: false }); try { const status = await db.admin().serverStatus(); res.json({ connected: true, version: status.version, uptime: status.uptime, connections: status.connections, opcounters: status.opcounters, }); } catch (err) { res.json({ connected: false, error: err.message }); }});
app.get("/api/redis", async (req, res) => { log("info", "redis stats requested"); res.json(await getRedisStats() || { connected: false });});
app.get("/api/storage/buckets", async (req, res) => { log("info", "storage buckets requested"); res.json(await getStorageStats());});
app.get("/api/services/oven", async (req, res) => { const url = process.env.OVEN_URL || "https://localhost:3002"; try { const start = Date.now(); const resp = await fetch(`${url}/health`, { signal: AbortSignal.timeout(5000) }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/session", async (req, res) => { const url = process.env.SESSION_URL || "https://localhost:8889"; try { const start = Date.now(); const resp = await fetch(`${url}/health`, { signal: AbortSignal.timeout(5000) }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/feed", async (req, res) => { const url = process.env.FEED_URL || "http://localhost:8787"; try { const start = Date.now(); const resp = await fetch(`${url}/api/v1/health`, { signal: AbortSignal.timeout(5000) }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/feed/info", async (req, res) => { const url = process.env.FEED_URL || "http://localhost:8787"; try { const start = Date.now(); const resp = await fetch(`${url}/api/v1`, { signal: AbortSignal.timeout(5000) }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/feed/playlists", async (req, res) => { const url = process.env.FEED_URL || "http://localhost:8787"; const apiSecret = process.env.FEED_API_SECRET || ""; try { const start = Date.now(); const headers = apiSecret ? { Authorization: `Bearer ${apiSecret}` } : {}; const resp = await fetch(`${url}/api/v1/playlists?limit=100`, { headers, signal: AbortSignal.timeout(10000), }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/feed/channels", async (req, res) => { const url = process.env.FEED_URL || "http://localhost:8787"; const apiSecret = process.env.FEED_API_SECRET || ""; try { const start = Date.now(); const headers = apiSecret ? { Authorization: `Bearer ${apiSecret}` } : {}; const resp = await fetch(`${url}/api/v1/channels?limit=100`, { headers, signal: AbortSignal.timeout(10000), }); const data = await resp.json(); res.json({ status: "ok", responseMs: Date.now() - start, ...data }); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
// --- Lith stats proxy ---const LITH_URL = process.env.LITH_URL || "https://aesthetic.computer";
app.get("/api/services/lith/stats", async (req, res) => { try { const resp = await fetch(`${LITH_URL}/lith/stats`, { signal: AbortSignal.timeout(5000) }); res.json(await resp.json()); } catch (err) { res.json({ status: "unreachable", error: err.message }); }});
app.get("/api/services/lith/errors", async (req, res) => { try { const resp = await fetch(`${LITH_URL}/lith/errors?limit=${req.query.limit || 100}`, { signal: AbortSignal.timeout(5000) }); res.json(await resp.json()); } catch (err) { res.json({ errors: [], error: err.message }); }});
app.get("/api/services/lith/requests", async (req, res) => { try { const resp = await fetch(`${LITH_URL}/lith/requests?limit=${req.query.limit || 100}`, { signal: AbortSignal.timeout(5000) }); res.json(await resp.json()); } catch (err) { res.json({ requests: [], error: err.message }); }});
app.get("/api/services/billing", async (req, res) => { const billingUrl = process.env.BILLING_URL || "https://aesthetic.computer/api/billing"; try { const start = Date.now(); const authHeader = req.headers.authorization || ""; const resp = await fetch(billingUrl, { headers: authHeader ? { Authorization: authHeader } : {}, signal: AbortSignal.timeout(10000), });
let data = null; try { data = await resp.json(); } catch { data = null; }
if (!resp.ok) { return res.status(resp.status).json({ status: "error", responseMs: Date.now() - start, error: data?.error || `Billing API error (${resp.status})`, }); }
return res.json({ status: "ok", responseMs: Date.now() - start, ...data, }); } catch (err) { return res.json({ status: "unreachable", error: err.message }); }});
// ─────────────────── Datomic admin (kidlisp sidecar proxy) ───────────────────//// Dashboard → silo (/api/datomic/*) → sidecar (/admin/*) with admin secret.// Sidecar is bound to 127.0.0.1 on silo, so the URL is silo-internal only.
const DATOMIC_SIDECAR_URL = process.env.DATOMIC_SIDECAR_URL || "http://127.0.0.1:8891";const DATOMIC_SIDECAR_ADMIN_SECRET = process.env.DATOMIC_SIDECAR_ADMIN_SECRET;
async function datomicProxy(req, res, { method, path, body }) { if (!DATOMIC_SIDECAR_ADMIN_SECRET) { return res.status(503).json({ error: "Datomic sidecar not configured" }); } try { const resp = await fetch(`${DATOMIC_SIDECAR_URL}${path}`, { method, headers: { "content-type": "application/json", "x-sidecar-secret": DATOMIC_SIDECAR_ADMIN_SECRET, }, body: body != null ? JSON.stringify(body) : undefined, signal: AbortSignal.timeout(15000), }); const text = await resp.text(); let data; try { data = text ? JSON.parse(text) : null; } catch { data = text; } res.status(resp.status).json(data); } catch (err) { res.status(502).json({ error: err.message }); }}
app.get("/api/datomic/health", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: "/admin/health" }));
app.get("/api/datomic/schema", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: "/admin/schema" }));
app.get("/api/datomic/stats", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: "/admin/stats" }));
app.get("/api/datomic/entities/:type", requireAdmin, (req, res) => { const qs = new URLSearchParams(req.query).toString(); datomicProxy(req, res, { method: "GET", path: `/admin/entities/${encodeURIComponent(req.params.type)}${qs ? `?${qs}` : ""}`, });});
app.get("/api/datomic/entity/:eid", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: `/admin/entity/${encodeURIComponent(req.params.eid)}` }));
app.get("/api/datomic/entity/:eid/history", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: `/admin/entity/${encodeURIComponent(req.params.eid)}/history` }));
app.get("/api/datomic/tx-log", requireAdmin, (req, res) => { const qs = new URLSearchParams(req.query).toString(); datomicProxy(req, res, { method: "GET", path: `/admin/tx-log${qs ? `?${qs}` : ""}` });});
app.post("/api/datomic/query", requireAdmin, (req, res) => datomicProxy(req, res, { method: "POST", path: "/admin/query", body: req.body }));
app.get("/api/datomic/backups", requireAdmin, (req, res) => datomicProxy(req, res, { method: "GET", path: "/admin/backups" }));
// Admin-gated proxy into the sidecar's client-secret read endpoints —// used by the dashboard's kidlisp tab to render AST trees + run structural// queries. Silo holds both secrets; dashboard admin auth → silo adds the// client-secret header → sidecar serves.const DATOMIC_SIDECAR_CLIENT_SECRET = process.env.DATOMIC_SIDECAR_CLIENT_SECRET;
async function sidecarClientProxy(req, res, { method, path, body }) { if (!DATOMIC_SIDECAR_CLIENT_SECRET) { return res.status(503).json({ error: "sidecar client secret not configured on silo" }); } try { const resp = await fetch(`${DATOMIC_SIDECAR_URL}${path}`, { method, headers: { "content-type": "application/json", "x-sidecar-secret": DATOMIC_SIDECAR_CLIENT_SECRET, }, body: body != null ? JSON.stringify(body) : undefined, signal: AbortSignal.timeout(15000), }); const text = await resp.text(); let data; try { data = text ? JSON.parse(text) : null; } catch { data = text; } res.status(resp.status).json(data); } catch (err) { res.status(502).json({ error: err.message }); }}
app.get("/api/datomic/kidlisp", requireAdmin, (req, res) => { const qs = new URLSearchParams(req.query).toString(); sidecarClientProxy(req, res, { method: "GET", path: `/kidlisp${qs ? `?${qs}` : ""}` });});
app.get("/api/datomic/kidlisp/:code", requireAdmin, (req, res) => sidecarClientProxy(req, res, { method: "GET", path: `/kidlisp/${encodeURIComponent(req.params.code)}` }));
app.get("/api/datomic/kidlisp/:code/ast", requireAdmin, (req, res) => sidecarClientProxy(req, res, { method: "GET", path: `/kidlisp/${encodeURIComponent(req.params.code)}/ast` }));
app.get("/api/datomic/kidlisp/:code/lineage", requireAdmin, (req, res) => sidecarClientProxy(req, res, { method: "GET", path: `/kidlisp/${encodeURIComponent(req.params.code)}/lineage` }));
app.get("/api/datomic/structural/pieces-using", requireAdmin, (req, res) => { const qs = new URLSearchParams(req.query).toString(); sidecarClientProxy(req, res, { method: "GET", path: `/kidlisp/structural/pieces-using${qs ? `?${qs}` : ""}`, });});
// ─────────────────── Datomic sidecar — public kidlisp proxy ───────────────────//// Server-to-server surface for lith (which runs store-kidlisp-datomic.mjs).// Authenticates via the sidecar's CLIENT_SECRET header — a shared secret// held by lith and the sidecar only, never exposed to browsers. This// route transparently forwards /sidecar/<anything> to the sidecar's// /<anything> endpoint so the Node client can use a clean base URL.
app.all("/sidecar/*", async (req, res) => { // Forward the client secret header as-is; sidecar rejects if missing/wrong. const subpath = req.originalUrl.replace(/^\/sidecar/, "") || "/"; const clientSecret = req.headers["x-sidecar-secret"]; if (!clientSecret) { return res.status(401).json({ error: "missing x-sidecar-secret" }); } try { const hasBody = req.method !== "GET" && req.method !== "HEAD"; const resp = await fetch(`${DATOMIC_SIDECAR_URL}${subpath}`, { method: req.method, headers: { "content-type": "application/json", "x-sidecar-secret": clientSecret, }, body: hasBody ? JSON.stringify(req.body ?? {}) : undefined, signal: AbortSignal.timeout(15000), }); const text = await resp.text(); res.status(resp.status); res.setHeader("content-type", resp.headers.get("content-type") || "application/json"); res.send(text); } catch (err) { res.status(502).json({ error: err.message }); }});
app.get("/api/firehose/history", async (req, res) => { if (!db) return res.json([]); const limit = Math.min(parseInt(req.query.limit) || 100, 1000); const ns = req.query.ns || null; const query = ns ? { ns } : {}; try { const events = await db.collection(FIREHOSE_COLLECTION) .find(query).sort({ time: -1 }).limit(limit).toArray(); res.json(events.reverse()); } catch (err) { res.status(500).json({ error: err.message }); }});
app.get("/api/firehose/stats", async (req, res) => { if (!db) return res.json({}); try { const col = db.collection(FIREHOSE_COLLECTION); const total = await col.estimatedDocumentCount(); const pipeline = [ { $group: { _id: { ns: "$ns", op: "$op" }, count: { $sum: 1 } } }, { $sort: { count: -1 } }, ]; const breakdown = await col.aggregate(pipeline).toArray(); res.json({ total, breakdown }); } catch (err) { res.status(500).json({ error: err.message }); }});
// --- Database Browser ---app.get("/api/db/browse/:collection", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); const colName = req.params.collection; const skip = Math.max(0, parseInt(req.query.skip) || 0); const limit = Math.min(100, Math.max(1, parseInt(req.query.limit) || 25)); const sortField = req.query.sort || "_id"; const sortDir = req.query.dir === "asc" ? 1 : -1; const q = req.query.q || "";
try { const col = db.collection(colName); let filter = {}; if (q) { // Search by _id (string or ObjectId) or text match on common fields const conditions = [ { _id: q }, ]; if (ObjectId.isValid(q) && q.length === 24) { conditions.push({ _id: new ObjectId(q) }); } // Regex search on string fields const regex = { $regex: q, $options: "i" }; conditions.push( { handle: regex }, { user: regex }, { name: regex }, { slug: regex }, { text: regex }, { email: regex }, { code: regex }, { mood: regex }, { piece: regex }, ); filter = { $or: conditions }; } const [docs, total] = await Promise.all([ col.find(filter).sort({ [sortField]: sortDir }).skip(skip).limit(limit).toArray(), col.countDocuments(filter), ]); res.json({ docs, total, skip, limit }); } catch (err) { res.status(500).json({ error: err.message }); }});
app.get("/api/db/browse/:collection/:id", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); try { const col = db.collection(req.params.collection); const id = req.params.id; let doc = await col.findOne({ _id: id }); if (!doc && ObjectId.isValid(id) && id.length === 24) { doc = await col.findOne({ _id: new ObjectId(id) }); } if (!doc) return res.status(404).json({ error: "Document not found" }); res.json(doc); } catch (err) { res.status(500).json({ error: err.message }); }});
app.put("/api/db/browse/:collection/:id", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); const { _id, ...update } = req.body; try { const col = db.collection(req.params.collection); const id = req.params.id; let result = await col.updateOne({ _id: id }, { $set: update }); if (result.matchedCount === 0 && ObjectId.isValid(id) && id.length === 24) { result = await col.updateOne({ _id: new ObjectId(id) }, { $set: update }); } if (result.matchedCount === 0) return res.status(404).json({ error: "Document not found" }); log("info", `browse: updated ${req.params.collection}/${id}`); res.json({ ok: true, modified: result.modifiedCount }); } catch (err) { res.status(500).json({ error: err.message }); }});
app.delete("/api/db/browse/:collection/:id", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); try { const col = db.collection(req.params.collection); const id = req.params.id; let result = await col.deleteOne({ _id: id }); if (result.deletedCount === 0 && ObjectId.isValid(id) && id.length === 24) { result = await col.deleteOne({ _id: new ObjectId(id) }); } if (result.deletedCount === 0) return res.status(404).json({ error: "Document not found" }); log("info", `browse: deleted ${req.params.collection}/${id}`); res.json({ ok: true }); } catch (err) { res.status(500).json({ error: err.message }); }});
// --- Instagram Admin Routes ---app.get("/api/insta/status", async (req, res) => { const sessionDoc = db ? await db.collection("insta-sessions").findOne({ _id: igSessionUsername }).catch(() => null) : null; res.json({ loggedIn: igLoggedIn, username: igSessionUsername || null, lastUsed: igLastUsed ? new Date(igLastUsed).toISOString() : null, sessionAge: sessionDoc?.updatedAt ? Math.round((Date.now() - new Date(sessionDoc.updatedAt).getTime()) / 1000) : null, challengeInProgress: igChallengeInProgress, });});
app.post("/api/insta/login", async (req, res) => { if (!db) return res.status(500).json({ error: "Database not connected" });
const secrets = await db.collection("secrets").findOne({ _id: "instagram" }); if (!secrets?.username || !secrets?.password) { return res.status(400).json({ error: "Instagram credentials not in MongoDB secrets" }); }
const username = secrets.username; const password = secrets.password;
log("info", `insta login attempt for @${username}`);
igClient = new IgApiClient(); igClient.state.generateDevice(username); igSessionUsername = username;
igClient.request.end$.subscribe(async () => { await saveInstaSession(igClient, username); });
// Pre-login flow (instagram-cli pattern) try { await igClient.launcher.preLoginSync(); log("info", "insta preLoginSync OK"); } catch (e) { log("warn", `insta preLoginSync failed: ${e.message}`); }
try { await igClient.account.login(username, password); igLoggedIn = true; igLastUsed = Date.now();
// Post-login flow try { await igClient.feed.reelsTray("cold_start").request(); await igClient.feed.timeline("cold_start_fetch").request(); } catch (e) { log("warn", `insta postLoginFlow failed: ${e.message}`); }
await saveInstaSession(igClient, username); log("info", `insta login success for @${username}`); res.json({ ok: true, username }); } catch (err) { if (err instanceof IgCheckpointError) { log("info", "insta checkpoint triggered, requesting verification code"); try { await igClient.challenge.auto(true); igChallengeInProgress = true; return res.json({ ok: false, challenge: true, message: "Verification code sent (check email/SMS). Submit via /api/insta/challenge.", }); } catch (challengeErr) { return res.status(500).json({ error: `Challenge setup failed: ${challengeErr.message}` }); } }
if (err instanceof IgLoginTwoFactorRequiredError) { const twoFactorInfo = err.response.body.two_factor_info; igChallengeInProgress = true; return res.json({ ok: false, twoFactor: true, method: twoFactorInfo.totp_two_factor_on ? "authenticator" : "sms", twoFactorIdentifier: twoFactorInfo.two_factor_identifier, message: "Enter 2FA code via /api/insta/challenge.", }); }
if (err instanceof IgLoginBadPasswordError) { return res.status(401).json({ error: "Bad password. Check credentials in MongoDB secrets." }); }
log("error", `insta login failed: ${err.message}`); return res.status(500).json({ error: err.message }); }});
app.post("/api/insta/challenge", async (req, res) => { if (!igClient || !igChallengeInProgress) { return res.status(400).json({ error: "No challenge in progress" }); }
const { code, twoFactorIdentifier, method } = req.body || {}; if (!code) return res.status(400).json({ error: "code is required" });
try { if (twoFactorIdentifier) { // 2FA flow await igClient.account.twoFactorLogin({ username: igSessionUsername, verificationCode: String(code).trim(), twoFactorIdentifier, verificationMethod: method === "authenticator" ? "0" : "1", }); } else { // Checkpoint flow await igClient.challenge.sendSecurityCode(String(code).trim()); }
igLoggedIn = true; igChallengeInProgress = false; igLastUsed = Date.now();
// Post-login flow try { await igClient.feed.reelsTray("cold_start").request(); await igClient.feed.timeline("cold_start_fetch").request(); } catch (e) { log("warn", `insta postLoginFlow failed: ${e.message}`); }
await saveInstaSession(igClient, igSessionUsername); log("info", `insta challenge completed for @${igSessionUsername}`); res.json({ ok: true, username: igSessionUsername }); } catch (err) { log("error", `insta challenge failed: ${err.message}`); res.status(400).json({ error: `Challenge verification failed: ${err.message}` }); }});
app.post("/api/insta/logout", async (req, res) => { igClient = null; igLoggedIn = false; igChallengeInProgress = false; igLastUsed = null;
if (db && igSessionUsername) { await db.collection("insta-sessions").deleteOne({ _id: igSessionUsername }).catch(() => {}); }
log("info", `insta session cleared for @${igSessionUsername}`); igSessionUsername = null; res.json({ ok: true });});
// --- TikTok Admin Routes ---app.get("/api/tiktok/status", async (req, res) => { const token = await loadTiktokSession(); const sessionDoc = db ? await db.collection("tiktok-sessions").findOne({ _id: "default" }).catch(() => null) : null; res.json({ connected: !!token?.access_token, username: token?.username || null, openId: token?.open_id || null, scope: token?.scope || null, lastUsed: tiktokLastUsed ? new Date(tiktokLastUsed).toISOString() : null, expiresAt: token?.expires_at ? new Date(token.expires_at).toISOString() : null, sessionAge: sessionDoc?.updatedAt ? Math.round((Date.now() - new Date(sessionDoc.updatedAt).getTime()) / 1000) : null, });});
app.get("/api/tiktok/auth", (req, res) => { if (!TIKTOK_CLIENT_KEY) return res.status(500).json({ error: "TIKTOK_CLIENT_KEY not configured" }); const scopes = "user.info.basic,user.info.profile,user.info.stats,video.list"; const url = `https://www.tiktok.com/v2/auth/authorize/?client_key=${TIKTOK_CLIENT_KEY}&scope=${scopes}&response_type=code&redirect_uri=${encodeURIComponent(TIKTOK_REDIRECT_URI)}`; res.json({ url });});
app.post("/api/tiktok/disconnect", async (req, res) => { tiktokToken = null; tiktokLastUsed = null; if (db) { await db.collection("tiktok-sessions").deleteOne({ _id: "default" }).catch(() => {}); } log("info", "tiktok session disconnected"); res.json({ ok: true });});
app.post("/api/tiktok/refresh", async (req, res) => { const ok = await refreshTiktokToken(); res.json({ ok, token: ok ? { expiresAt: new Date(tiktokToken.expires_at).toISOString() } : null });});
// --- Desktop Release Distribution ---const DESKTOP_BUCKET = "releases-aesthetic-computer";const DESKTOP_PREFIX = "desktop/";const DESKTOP_BASE_URL = "https://releases.aesthetic.computer/desktop";
async function ensureReleasesCollection() { if (!db) return; try { const col = db.collection("releases"); await col.createIndex({ version: 1 }, { unique: true }); await col.createIndex({ current: 1 }); log("info", "releases collection ready"); } catch (err) { log("error", `releases index setup: ${err.message}`); }}
// Public routes (CORS-enabled, no auth)app.use("/desktop", (req, res, next) => { res.header("Access-Control-Allow-Origin", "*"); res.header("Access-Control-Allow-Methods", "GET, OPTIONS"); res.header("Access-Control-Allow-Headers", "Content-Type"); if (req.method === "OPTIONS") return res.sendStatus(200); next();});
app.get("/desktop/latest", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); try { const release = await db.collection("releases").findOne({ current: true }); if (!release) return res.status(404).json({ error: "No release published" }); res.json({ version: release.version, publishedAt: release.publishedAt, releaseNotes: release.releaseNotes || null, platforms: release.platforms, }); } catch (err) { res.status(500).json({ error: err.message }); }});
app.get("/desktop/download/:platform", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); const platform = req.params.platform; try { const release = await db.collection("releases").findOne({ current: true }); if (!release) return res.status(404).json({ error: "No release published" }); const plat = release.platforms?.[platform]; if (!plat?.url) return res.status(404).json({ error: `No ${platform} build available` }); res.redirect(302, plat.url); } catch (err) { res.status(500).json({ error: err.message }); }});
// Admin routesapp.post("/api/desktop/register", async (req, res) => { if (!db) return res.status(503).json({ error: "Database not connected" }); const { version, platforms, releaseNotes } = req.body; if (!version || !platforms) { return res.status(400).json({ error: "version and platforms required" }); }
// Build full URLs from filenames for (const [key, plat] of Object.entries(platforms)) { if (plat.filename && !plat.url) { plat.url = `${DESKTOP_BASE_URL}/${plat.filename}`; } }
try { // Unset current on all previous releases await db.collection("releases").updateMany({}, { $set: { current: false } }); // Upsert this release await db.collection("releases").updateOne( { version }, { $set: { version, platforms, releaseNotes: releaseNotes || null, publishedAt: new Date(), publishedBy: req.auth?.handle || "unknown", current: true, }, }, { upsert: true }, ); log("info", `desktop release registered: v${version} by @${req.auth?.handle}`); res.json({ ok: true, version }); } catch (err) { log("error", `desktop register failed: ${err.message}`); res.status(500).json({ error: err.message }); }});
app.get("/api/desktop/releases", async (req, res) => { if (!db) return res.json({ releases: [] }); try { const releases = await db.collection("releases") .find({}).sort({ publishedAt: -1 }).limit(20).toArray(); res.json({ releases }); } catch (err) { res.status(500).json({ error: err.message }); }});
// --- Dashboard (loaded from file, reloadable via SIGHUP) ---const DASHBOARD_PATH = path.join(__dirname, "dashboard.html");let dashboardHtml = "";
function loadDashboard() { try { dashboardHtml = fs.readFileSync(DASHBOARD_PATH, "utf-8"); log("info", `dashboard loaded (${(dashboardHtml.length / 1024).toFixed(0)} KB)`); } catch (err) { log("error", `dashboard load failed: ${err.message}`); }}loadDashboard();
app.get("/", (req, res) => { res.setHeader("Content-Type", "text/html"); res.send(dashboardHtml);});
// --- Bluesky ingest (external news headlines from trusted sources) ---const BSKY_INGEST_INTERVAL_MS = parseInt(process.env.NEWS_BLUESKY_INTERVAL_MS || "", 10) || 10 * 60 * 1000; // default 10 minconst BSKY_INGEST_LIMIT = parseInt(process.env.NEWS_BLUESKY_LIMIT || "", 10) || 30;let bskyIngestTimer = null;let bskyLastRun = { when: null, results: [], error: null, runCount: 0 };
async function runBlueskyIngestOnce(actor) { if (!db) { const err = "mongo not connected"; bskyLastRun = { when: new Date().toISOString(), results: [], error: err, runCount: bskyLastRun.runCount }; return bskyLastRun; } try { const results = actor ? [await ingestBlueskyActor(db, actor, { limit: BSKY_INGEST_LIMIT, log: (line) => log("info", `bsky: ${line}`) })] : await ingestBluesky(db, { limit: BSKY_INGEST_LIMIT, log: (line) => log("info", `bsky: ${line}`) }); const totalInserted = results.reduce((n, r) => n + (r.inserted || 0), 0); const totalErrors = results.reduce((n, r) => n + (r.errors?.length || 0), 0); bskyLastRun = { when: new Date().toISOString(), results, error: null, runCount: (bskyLastRun.runCount || 0) + 1, }; if (totalInserted > 0 || totalErrors > 0) { log("info", `bsky ingest: +${totalInserted} new, ${totalErrors} errors`); } return bskyLastRun; } catch (err) { bskyLastRun = { when: new Date().toISOString(), results: [], error: err.message, runCount: (bskyLastRun.runCount || 0) + 1, }; log("error", `bsky ingest failed: ${err.message}`); return bskyLastRun; }}
function startBlueskyIngestLoop() { if (bskyIngestTimer) return; const sources = getBlueskySources(); log("info", `bsky ingest: ${sources.length} source(s), every ${Math.round(BSKY_INGEST_INTERVAL_MS / 60000)}m — ${sources.join(", ")}`); // Kick off a first run ~30s after boot to avoid colliding with other startup work. setTimeout(() => runBlueskyIngestOnce().catch(() => {}), 30_000); bskyIngestTimer = setInterval(() => runBlueskyIngestOnce().catch(() => {}), BSKY_INGEST_INTERVAL_MS);}
app.get("/api/news/bluesky/status", (req, res) => { res.json({ sources: getBlueskySources(), intervalMs: BSKY_INGEST_INTERVAL_MS, limit: BSKY_INGEST_LIMIT, lastRun: bskyLastRun, });});
app.post("/api/news/bluesky/pull", async (req, res) => { const actor = typeof req.body?.actor === "string" ? req.body.actor.trim() : ""; const result = await runBlueskyIngestOnce(actor || undefined); res.json(result);});
// --- 404 ---app.use((req, res) => res.status(404).json({ error: "Not found" }));
// --- Server Start ---let server;if (dev) { try { const httpsOpts = { key: fs.readFileSync("../ssl-dev/localhost-key.pem"), cert: fs.readFileSync("../ssl-dev/localhost.pem"), }; server = https.createServer(httpsOpts, app); } catch { log("warn", "no SSL certs, falling back to HTTP"); server = http.createServer(app); }} else { server = http.createServer(app);}
// --- WebSocket ---wss = new WebSocketServer({ server, path: "/ws" });
wss.on("connection", async (ws) => { log("info", "dashboard client connected"); for (const entry of activityLog.slice(0, 30).reverse()) { ws.send(JSON.stringify({ logEntry: entry })); } ws.on("close", () => log("info", "dashboard client disconnected"));});
// --- Boot ---await connectMongo();await loadHandleCache();await ensureFirehoseCollection();await ensureReleasesCollection();startFirehose();await connectAtlas();await connectRedis();await connectRedisSub();
// Restore TikTok session on startuploadTiktokSession().catch(() => {});
// Start the Bluesky news ingest loop (pulls headlines from trusted sources).startBlueskyIngestLoop();
server.listen(PORT, () => { const proto = dev ? "https" : "http"; log("info", `silo running on ${proto}://localhost:${PORT}`);});
// --- Shutdown ---function shutdown(signal) { log("info", `received ${signal}, shutting down...`); if (bskyIngestTimer) { clearInterval(bskyIngestTimer); bskyIngestTimer = null; } if (changeStream) changeStream.close().catch(() => {}); wss.clients.forEach((ws) => ws.close()); server.close(); mongoClient?.close(); if (atlasClient && atlasClient !== mongoClient) atlasClient?.close(); redisClient?.quit().catch(() => {}); redisSub?.quit().catch(() => {}); setTimeout(() => process.exit(0), 500);}process.on("SIGTERM", () => shutdown("SIGTERM"));process.on("SIGINT", () => shutdown("SIGINT"));
// --- SIGHUP: hot-reload dashboard.html without dropping connections ---process.on("SIGHUP", () => { log("info", "SIGHUP received, reloading dashboard..."); loadDashboard();});