Something went wrong. Try again.
Monorepo for Aesthetic.Computer aesthetic.computer
Something went wrong. Try again.
141 kB · 4166 lines
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766276727682769277027712772277327742775277627772778277927802781278227832784278527862787278827892790279127922793279427952796279727982799280028012802280328042805280628072808280928102811281228132814281528162817281828192820282128222823282428252826282728282829283028312832283328342835283628372838283928402841284228432844284528462847284828492850285128522853285428552856285728582859286028612862286328642865286628672868286928702871287228732874287528762877287828792880288128822883288428852886288728882889289028912892289328942895289628972898289929002901290229032904290529062907290829092910291129122913291429152916291729182919292029212922292329242925292629272928292929302931293229332934293529362937293829392940294129422943294429452946294729482949295029512952295329542955295629572958295929602961296229632964296529662967296829692970297129722973297429752976297729782979298029812982298329842985298629872988298929902991299229932994299529962997299829993000300130023003300430053006300730083009301030113012301330143015301630173018301930203021302230233024302530263027302830293030303130323033303430353036303730383039304030413042304330443045304630473048304930503051305230533054305530563057305830593060306130623063306430653066306730683069307030713072307330743075307630773078307930803081308230833084308530863087308830893090309130923093309430953096309730983099310031013102310331043105310631073108310931103111311231133114311531163117311831193120312131223123312431253126312731283129313031313132313331343135313631373138313931403141314231433144314531463147314831493150315131523153315431553156315731583159316031613162316331643165316631673168316931703171317231733174317531763177317831793180318131823183318431853186318731883189319031913192319331943195319631973198319932003201320232033204320532063207320832093210321132123213321432153216321732183219322032213222322332243225322632273228322932303231323232333234323532363237323832393240324132423243324432453246324732483249325032513252325332543255325632573258325932603261326232633264326532663267326832693270327132723273327432753276327732783279328032813282328332843285328632873288328932903291329232933294329532963297329832993300330133023303330433053306330733083309331033113312331333143315331633173318331933203321332233233324332533263327332833293330333133323333333433353336333733383339334033413342334333443345334633473348334933503351335233533354335533563357335833593360336133623363336433653366336733683369337033713372337333743375337633773378337933803381338233833384338533863387338833893390339133923393339433953396339733983399340034013402340334043405340634073408340934103411341234133414341534163417341834193420342134223423342434253426342734283429343034313432343334343435343634373438343934403441344234433444344534463447344834493450345134523453345434553456345734583459346034613462346334643465346634673468346934703471347234733474347534763477347834793480348134823483348434853486348734883489349034913492349334943495349634973498349935003501350235033504350535063507350835093510351135123513351435153516351735183519352035213522352335243525352635273528352935303531353235333534353535363537353835393540354135423543354435453546354735483549355035513552355335543555355635573558355935603561356235633564356535663567356835693570357135723573357435753576357735783579358035813582358335843585358635873588358935903591359235933594359535963597359835993600360136023603360436053606360736083609361036113612361336143615361636173618361936203621362236233624362536263627362836293630363136323633363436353636363736383639364036413642364336443645364636473648364936503651365236533654365536563657365836593660366136623663366436653666366736683669367036713672367336743675367636773678367936803681368236833684368536863687368836893690369136923693369436953696369736983699370037013702370337043705370637073708370937103711371237133714371537163717371837193720372137223723372437253726372737283729373037313732373337343735373637373738373937403741374237433744374537463747374837493750375137523753375437553756375737583759376037613762376337643765376637673768376937703771377237733774377537763777377837793780378137823783378437853786378737883789379037913792379337943795379637973798379938003801380238033804380538063807380838093810381138123813381438153816381738183819382038213822382338243825382638273828382938303831383238333834383538363837383838393840384138423843384438453846384738483849385038513852385338543855385638573858385938603861386238633864386538663867386838693870387138723873387438753876387738783879388038813882388338843885388638873888388938903891389238933894389538963897389838993900390139023903390439053906390739083909391039113912391339143915391639173918391939203921392239233924392539263927392839293930393139323933393439353936393739383939394039413942394339443945394639473948394939503951395239533954395539563957395839593960396139623963396439653966396739683969397039713972397339743975397639773978397939803981398239833984398539863987398839893990399139923993399439953996399739983999400040014002400340044005400640074008400940104011401240134014401540164017401840194020402140224023402440254026402740284029403040314032403340344035403640374038403940404041404240434044404540464047404840494050405140524053405440554056405740584059406040614062406340644065406640674068406940704071407240734074407540764077407840794080408140824083408440854086408740884089409040914092409340944095409640974098409941004101410241034104410541064107410841094110411141124113411441154116411741184119412041214122412341244125412641274128412941304131413241334134413541364137413841394140414141424143414441454146414741484149415041514152415341544155415641574158415941604161416241634164416541664167// Session Server, 23.12.04.14.57// Represents a "room" or user or "client" backend// which at the moment is run once for every "piece"// that requests it.
/* #region todo + Now - [-] Fix live reloading of in-production udp. + Done - [c] `code.channel` should return a promise, and wait for a `code-channel:subbed`. event here? This way users get better confirmation if the socket doesn't go through or if there is a server issue. 23.07.04.18.01 (Might not actually be that necessary.) - [x] Add `obscenity` filter. - [x] Conditional redis sub to dev updates. (Will save bandwidth if extension gets lots of use, also would be more secure.) - [x] Secure the "code" path to require a special string. - [x] Secure the "reload" path (must be in dev mode, sorta okay) - [c] Speed up developer reload by using redis pub/sub. - [x] Send a signal to everyone once a user leaves. - [x] Get "developer" live reloading working again. - [x] Add sockets back. - [x] Make a "local" option. - [x] Read through: https://redis.io/docs/data-types#endregion */
// Add redis pub/sub here...
import Fastify from "fastify";import geckos from "@geckos.io/server";import geoip from "geoip-lite";import { WebSocket, WebSocketServer } from "ws";import ip from "ip";import chokidar from "chokidar";import fs from "fs";import path from "path";import crypto from "crypto";import dotenv from "dotenv";import dgram from "dgram";dotenv.config();
// Module streaming - path to public directoryconst PUBLIC_DIR = path.resolve(process.cwd(), "../system/public/aesthetic.computer");
// Module hash cache (invalidated on file change)const moduleHashes = new Map(); // path -> { hash, content, mtime }
// Compute hash for a module filefunction getModuleHash(modulePath) { const fullPath = path.join(PUBLIC_DIR, modulePath); try { const stats = fs.statSync(fullPath); const cached = moduleHashes.get(modulePath); // Return cached if mtime matches if (cached && cached.mtime === stats.mtimeMs) { return cached; } // Read and hash const content = fs.readFileSync(fullPath, "utf8"); const hash = crypto.createHash("sha256").update(content).digest("hex").slice(0, 16); const entry = { hash, content, mtime: stats.mtimeMs }; moduleHashes.set(modulePath, entry); return entry; } catch (err) { return null; }}
// Fairy:point throttle (for silo firehose visualization)const fairyThrottle = new Map(); // channelId -> last publish timestampconst FAIRY_THROTTLE_MS = 100; // 10Hz max per connection
// Raw UDP fairy relay (for native bare-metal clients)const udpRelay = dgram.createSocket("udp4");const udpClients = new Map(); // key "ip:port" → { address, port, handle, lastSeen }const UDP_MIDI_SOURCE_TTL_MS = 20000;const notepatMidiSources = new Map(); // key "@handle:machine" -> source metadataconst notepatMidiSubscribers = new Map(); // connection id -> { ws, all, handle, machineId }// UDP-side subscribers for notepat:midi fan-out over geckos.io. The WS map// above handles reliable subscription handshakes; this map mirrors the same// filter model against geckos channels so we can emit events twice (once// reliably over WS, once low-latency over UDP) to consumers that opened both.const notepatMidiUdpSubscribers = new Map(); // channel id -> { channel, all, handle, machineId }
// Error logging ring buffer (for dashboard display)const errorLog = [];const MAX_ERRORS = 50;const ERROR_RETENTION_MS = 60 * 60 * 1000; // 1 hour
function logError(level, message) { const entry = { level, message: typeof message === 'string' ? message : JSON.stringify(message), timestamp: new Date().toISOString() }; errorLog.push(entry); if (errorLog.length > MAX_ERRORS) errorLog.shift();}
// Capture uncaught errorsprocess.on('uncaughtException', (err) => { logError('error', `Uncaught: ${err.message}`); console.error('Uncaught Exception:', err);});
process.on('unhandledRejection', (reason, promise) => { logError('error', `Unhandled Rejection: ${reason}`); console.error('Unhandled Rejection:', reason);});
import { exec } from "child_process";
// 🔔 Push notifications (standard Web Push + direct APNs — no Firebase).import { broadcastToTopic } from "../shared/push.mjs";
// Initialize ChatManager for multi-instance chat supportconst chatManager = new ChatManager({ dev: process.env.NODE_ENV === "development" });await chatManager.init();
// Graceful shutdown — persist in-memory chat messages before exitlet shuttingDown = false;async function gracefulShutdown(signal) { if (shuttingDown) return; shuttingDown = true; console.log(`\n${signal} received, persisting chat messages...`); try { await chatManager.shutdown(); } catch (err) { console.error("Shutdown error:", err); } process.exit(0);}process.on("SIGTERM", () => gracefulShutdown("SIGTERM"));process.on("SIGINT", () => gracefulShutdown("SIGINT"));
// Helper function to get handles of users currently on a specific piece// Used by chatManager to determine who's actually viewing the chat piecefunction getHandlesOnPiece(pieceName) { const handles = []; for (const [id, client] of Object.entries(clients)) { if (client.location === pieceName && client.handle) { handles.push(client.handle); } } return [...new Set(handles)]; // Remove duplicates}
// Expose the function to chatManagerchatManager.setPresenceResolver(getHandlesOnPiece);
// 🎯 Duel Manager — server-authoritative game for dumduel piececonst duelManager = new DuelManager();const fightManager = new FightManager();// 🌐 World Managers — Q3-style server-authoritative multiplayer, one// instance per 3D world. Wire protocol is shared; only the event prefix// differs (`arena:*`, `land:*`). Routing below is generic over this map —// adding a world is one line here plus a cfg in lib/<name>-world.mjs.const worldManagers = { arena: new WorldManager({ prefix: "arena", cfg: ARENA_CFG, spawns: ARENA_SPAWNS, icon: "🏟️" }), land: new WorldManager({ prefix: "land", cfg: LAND_CFG, spawns: LAND_SPAWNS, icon: "🌾" }),};
import { filter } from "./filter.mjs"; // Profanity filtering.import { ChatManager } from "./chat-manager.mjs"; // Multi-instance chat support.import { DuelManager } from "./duel-manager.mjs"; // Server-authoritative duel game.import { FightManager } from "./fight-manager.mjs"; // Fight presence + seat queue.import { WorldManager, ARENA_CFG, ARENA_SPAWNS } from "./world-manager.mjs"; // Server-authoritative 3D worlds.import { LAND_CFG, LAND_SPAWNS } from "../system/public/aesthetic.computer/lib/land-world.mjs";
// *** AC Machines — remote device monitoring ***// Devices connect via /machines?role=device&machineId=X&token=Y// Viewers connect via /machines?role=viewer&token=Yimport { MongoClient } from "mongodb";
const machinesDevices = new Map(); // machineId → { ws, user, handle, machineId, info, lastHeartbeat }const machinesViewers = new Map(); // userSub → Set<ws>let machinesDb = null;
async function getMachinesDb() { if (machinesDb) return machinesDb; const connStr = process.env.MONGODB_CONNECTION_STRING; if (!connStr) return null; try { const client = new MongoClient(connStr); await client.connect(); machinesDb = client.db(process.env.MONGODB_NAME || "aesthetic"); return machinesDb; } catch (e) { error("[machines] MongoDB connect error:", e.message); return null; }}
let machineTokenSecret = null;let machineTokenSecretAt = 0;const MACHINE_SECRET_TTL = 5 * 60 * 1000; // refresh from DB every 5 min
async function loadMachineTokenSecret() { const now = Date.now(); if (machineTokenSecret && now - machineTokenSecretAt < MACHINE_SECRET_TTL) { return machineTokenSecret; } try { const db = await getMachinesDb(); if (!db) return machineTokenSecret; const doc = await db.collection("secrets").findOne({ _id: "machine-token" }); if (doc?.secret) { machineTokenSecret = doc.secret; machineTokenSecretAt = now; } } catch (e) { error("[machines] Failed to load machine-token secret:", e.message); } return machineTokenSecret;}
async function verifyMachineToken(token) { if (!token) return null; const secret = await loadMachineTokenSecret(); if (!secret) return null; try { const [payloadB64, sigB64] = token.split("."); if (!payloadB64 || !sigB64) return null; const expectedSig = crypto .createHmac("sha256", secret) .update(payloadB64) .digest("base64url"); if (sigB64 !== expectedSig) return null; return JSON.parse(Buffer.from(payloadB64, "base64url").toString()); } catch { return null; }}
// Verify an AC auth token (Bearer token from authorize()) by calling Auth0 userinfoasync function verifyACToken(token) { if (!token) return null; try { const res = await fetch("https://aesthetic.us.auth0.com/userinfo", { headers: { Authorization: `Bearer ${token}` }, }); if (!res.ok) return null; return await res.json(); // { sub, nickname, name, ... } } catch { return null; }}
function broadcastToMachineViewers(userSub, msg) { const viewers = machinesViewers.get(userSub); if (!viewers) return; const data = JSON.stringify(msg); for (const v of viewers) { if (v.readyState === WebSocket.OPEN) v.send(data); }}
async function upsertMachine(userSub, machineId, info) { const db = await getMachinesDb(); if (!db) return; const col = db.collection("ac-machines"); const now = new Date(); await col.updateOne( { user: userSub, machineId }, { $set: { user: userSub, machineId, ...info, status: "online", linked: true, lastSeen: now, updatedAt: now, }, $setOnInsert: { createdAt: now, bootCount: 0 }, $inc: { bootCount: 1 }, }, { upsert: true }, );}
async function updateMachineHeartbeat(userSub, machineId, uptime, currentPiece) { const db = await getMachinesDb(); if (!db) return; await db.collection("ac-machines").updateOne( { user: userSub, machineId }, { $set: { lastSeen: new Date(), uptime, currentPiece, status: "online" } }, );}
async function insertMachineLog(userSub, machineId, msg) { const db = await getMachinesDb(); if (!db) return; await db.collection("ac-machine-logs").insertOne({ machineId, user: userSub, type: msg.logType || "log", level: msg.level || "info", message: msg.message, data: msg.data || null, crashInfo: msg.crashInfo || null, when: msg.when ? new Date(msg.when) : new Date(), receivedAt: new Date(), });}
async function setMachineOffline(userSub, machineId) { const db = await getMachinesDb(); if (!db) return; await db.collection("ac-machines").updateOne( { user: userSub, machineId }, { $set: { status: "offline", updatedAt: new Date() } }, );}
// *** SockLogs - Remote console log forwarding from devices ***// Devices with ?socklogs param send logs via WebSocket// Viewers (CLI or web) can subscribe to see device logs in real-timeconst socklogsDevices = new Map(); // deviceId -> { ws, lastLog, logCount }const socklogsViewers = new Set(); // Set of viewer WebSockets
function socklogsBroadcast(deviceId, logEntry) { const message = JSON.stringify({ type: 'log', deviceId, ...logEntry, serverTime: Date.now() }); for (const viewer of socklogsViewers) { if (viewer.readyState === WebSocket.OPEN) { viewer.send(message); } }}
function socklogsStatus() { return { devices: Array.from(socklogsDevices.entries()).map(([id, info]) => ({ deviceId: id, logCount: info.logCount, lastLog: info.lastLog, connectedAt: info.connectedAt })), viewerCount: socklogsViewers.size };}
import { createClient } from "redis";const redisConnectionString = process.env.REDIS_CONNECTION_STRING;const dev = process.env.NODE_ENV === "development";
// Dev log file for remote debuggingconst DEV_LOG_FILE = path.join(process.cwd(), "../system/public/aesthetic.computer/dev-logs.txt");
const { keys } = Object;let fastify; //, termkit, term;
if (dev) { // Load local ssl certs in development mode. fastify = Fastify({ https: { // allowHTTP1: true, key: fs.readFileSync("../ssl-dev/localhost-key.pem"), cert: fs.readFileSync("../ssl-dev/localhost.pem"), }, logger: true, });
// Import the `terminal-kit` library if dev is true. // try { // termkit = (await import("terminal-kit")).default; // } catch (err) { // error("Failed to load terminal-kit", error); // }} else { fastify = Fastify({ logger: true }); // Still log in production. No reason not to?}
// Insert `cors` headers as needed. 23.12.19.16.31// TODO: Is this even necessary?fastify.options("*", async (req, reply) => { const allowedOrigins = [ "https://aesthetic.local:8888", "https://aesthetic.computer", "https://notepat.com", "https://laklok.com", "https://www.laklok.com", ];
const origin = req.headers.origin; log("✈️ Preflight origin:", origin); // Check if the incoming origin is allowed if (allowedOrigins.includes(origin)) { reply.header("Access-Control-Allow-Origin", origin); } reply.header("Access-Control-Allow-Methods", "GET, POST, PUT, DELETE"); reply.send();});
const server = fastify.server;
const DEV_LOG_DIR = "/tmp/dev-logs/";const deviceLogFiles = new Map(); // Track which devices have log files
// Ensure log directory existsif (dev) { try { fs.mkdirSync(DEV_LOG_DIR, { recursive: true }); } catch (error) { console.error("Failed to create dev log directory:", error); }}
const info = { port: process.env.PORT, // 8889 in development via `package.json` name: process.env.SESSION_BACKEND_ID, service: process.env.JAMSOCKET_SERVICE,};
const codeChannels = {}; // Used to filter `code` updates from redis to// clients who explicitly have the channel set.const codeChannelState = {}; // Store last code sent to each channel for late joiners
// DAW channel for M4L device ↔ IDE communicationconst dawDevices = new Set(); // Connection IDs of /device instancesconst dawIDEs = new Set(); // Connection IDs of IDE instances in Ableton mode
// Unified client tracking: each client has handle, user, location, and connection typesconst clients = {}; // Map of connection ID to { handle, user, location, websocket: true/false, udp: true/false }
// Device naming for local dev (persisted to file)const DEVICE_NAMES_FILE = path.join(process.cwd(), "../.device-names.json");let deviceNames = {}; // Map of IP -> { name, group }function loadDeviceNames() { try { if (fs.existsSync(DEVICE_NAMES_FILE)) { deviceNames = JSON.parse(fs.readFileSync(DEVICE_NAMES_FILE, 'utf8')); log("📱 Loaded device names:", Object.keys(deviceNames).length); } } catch (e) { log("📱 Could not load device names:", e.message); }}function saveDeviceNames() { try { fs.writeFileSync(DEVICE_NAMES_FILE, JSON.stringify(deviceNames, null, 2)); } catch (e) { log("📱 Could not save device names:", e.message); }}if (dev) loadDeviceNames();
// Get the dev host machine nameimport os from "os";const DEV_HOST_NAME = os.hostname();const DEV_LAN_IP = (() => { // First, try to read from /tmp/host-lan-ip (written by entry.fish in devcontainer) try { const hostIpFile = '/tmp/host-lan-ip'; if (fs.existsSync(hostIpFile)) { const ip = fs.readFileSync(hostIpFile, 'utf-8').trim(); if (ip && ip.match(/^\d+\.\d+\.\d+\.\d+$/)) { console.log(`🖥️ Using host LAN IP from ${hostIpFile}: ${ip}`); return ip; } } } catch (e) { /* ignore */ } // Fallback: try to detect from network interfaces const interfaces = os.networkInterfaces(); for (const name of Object.keys(interfaces)) { for (const iface of interfaces[name]) { if (iface.family === 'IPv4' && !iface.internal && iface.address.startsWith('192.168.')) { return iface.address; } } } return null;})();console.log(`🖥️ Dev host: ${DEV_HOST_NAME}, LAN IP: ${DEV_LAN_IP || 'N/A'}`);
// Helper: Assign device letters (A, B, C...) based on connection orderfunction getDeviceLetter(connectionId) { // Get sorted list of connection IDs const sortedIds = Object.keys(connections) .map(id => parseInt(id)) .sort((a, b) => a - b); const index = sortedIds.indexOf(parseInt(connectionId)); if (index === -1) return '?'; // A=65, B=66, etc. Wrap around after Z return String.fromCharCode(65 + (index % 26));}
// Helper: Find connections by ID, IP, handle, or device letterfunction targetClients(target) { if (target === 'all') { return Object.entries(connections) .filter(([id, ws]) => ws?.readyState === WebSocket.OPEN) .map(([id, ws]) => ({ id: parseInt(id), ws })); } const results = []; for (const [id, ws] of Object.entries(connections)) { const client = clients[id]; const cleanTarget = target.replace('@', ''); const cleanIp = client?.ip?.replace('::ffff:', ''); const deviceLetter = getDeviceLetter(id); if ( String(id) === String(target) || cleanIp === target || client?.handle === `@${cleanTarget}` || client?.handle === cleanTarget || deviceNames[cleanIp]?.name?.toLowerCase() === target.toLowerCase() || deviceLetter.toLowerCase() === target.toLowerCase() // Match by letter (A, B, C...) ) { if (ws?.readyState === WebSocket.OPEN) { results.push({ id: parseInt(id), ws }); } } } return results;}
// *** Start up two `redis` clients. (One for subscribing, and for publishing)const redisEnabled = !!redisConnectionString;const sub = redisEnabled ? (!dev ? createClient({ url: redisConnectionString }) : createClient()) : null;if (sub) sub.on("error", (err) => { log("🔴 Redis subscriber client error!", err); logError('error', `Redis sub: ${err.message}`);});
const pub = redisEnabled ? (!dev ? createClient({ url: redisConnectionString }) : createClient()) : null;if (pub) pub.on("error", (err) => { log("🔴 Redis publisher client error!", err); logError('error', `Redis pub: ${err.message}`);});
try { if (sub && pub) { await sub.connect(); await pub.connect();
await sub.subscribe("code", (message) => { const parsed = JSON.parse(message); if (codeChannels[parsed.codeChannel]) { const msg = pack("code", message, "development"); subscribers(codeChannels[parsed.codeChannel], msg); } });
await sub.subscribe("scream", (message) => { everyone(pack("scream", message, "screamer")); // Socket back to everyone. }); } else { log("⚠️ Redis disabled — code/scream channels unavailable"); }} catch (err) { error("🔴 Could not connect to `redis` instance.");}
const secret = process.env.GITHUB_WEBHOOK_SECRET;
fastify.post("/update", (request, reply) => { const signature = request.headers["x-hub-signature"]; const hash = "sha1=" + crypto .createHmac("sha1", secret) .update(JSON.stringify(request.body)) .digest("hex");
if (hash !== signature) { reply.status(401).send({ error: "Invalid signature" }); return; }
// log("Path:", process.env.PATH);
// Restart service in production. // exec( // "cd /home/aesthetic-computer/session-server; pm2 stop all; git pull; npm install; pm2 start all", // (err, stdout, stderr) => { // if (err) { // error(`exec error: ${error}`); // return; // } // log(`stdout: ${stdout}`); // error(`stderr: ${stderr}`); // }, // );
reply.send({ status: "ok" });});
// *** Robots.txt - prevent crawling ***fastify.get("/robots.txt", async (req, reply) => { reply.type("text/plain"); return "User-agent: *\nDisallow: /";});
// *** Module HTTP endpoint - serve modules directly (bypasses Netlify proxy) ***// Used by boot.mjs on localhost when the main proxy is flakyfastify.get("/module/*", async (req, reply) => { const modulePath = req.params["*"]; const moduleData = getModuleHash(modulePath); if (moduleData) { reply .header("Content-Type", "application/javascript; charset=utf-8") .header("Access-Control-Allow-Origin", "*") .header("Cache-Control", "no-cache") .send(moduleData.content); } else { reply.status(404).send({ error: "Module not found", path: modulePath }); }});
// *** Build Stream - pipe terminal output to WebSocket clients ***// Available in both dev and production for build progress streamingfastify.post("/build-stream", async (req) => { const line = typeof req.body === 'string' ? req.body : req.body.line || ''; everyone(pack("build:log", { line, timestamp: Date.now() })); return { status: "ok" };});
fastify.post("/build-status", async (req) => { everyone(pack("build:status", { ...req.body, timestamp: Date.now() })); return { status: "ok" };});
// *** FF1 Art Computer Proxy ***// Proxies displayPlaylist commands to FF1 via direct connection or cloud relayconst FF1_RELAY_URL = "https://artwork-info.feral-file.workers.dev/api/cast";
// Load FF1 config from machines.jsonfunction getFF1Config() { try { const machinesPath = path.resolve(process.cwd(), "../aesthetic-computer-vault/machines.json"); const machines = JSON.parse(fs.readFileSync(machinesPath, "utf8")); return machines.machines?.["ff1-dvveklza"] || null; } catch (e) { log("⚠️ Could not load FF1 config from machines.json:", e.message); return null; }}
// Execute FF1 cast via SSH through MacBook (for devcontainer)async function castViaSSH(ff1Config, payload) { const { exec } = await import("child_process"); const { promisify } = await import("util"); const execAsync = promisify(exec); const ip = ff1Config.ip; const port = ff1Config.port || 1111; const payloadJson = JSON.stringify(payload).replace(/'/g, "'\\''"); // Escape for shell // SSH through MacBook to reach FF1 on local network const sshCmd = `ssh -o ConnectTimeout=5 jas@host.docker.internal "curl -s --connect-timeout 5 -X POST -H 'Content-Type: application/json' http://${ip}:${port}/api/cast -d '${payloadJson}'"`; log(`📡 FF1 cast via SSH: http://${ip}:${port}/api/cast`); const { stdout, stderr } = await execAsync(sshCmd, { timeout: 15000 }); if (stderr && !stdout) { throw new Error(stderr); } try { return JSON.parse(stdout); } catch { return { raw: stdout }; }}
fastify.post("/ff1/cast", async (req, reply) => { reply.header("Access-Control-Allow-Origin", "*"); reply.header("Access-Control-Allow-Methods", "POST, OPTIONS"); reply.header("Access-Control-Allow-Headers", "Content-Type"); const { topicID, apiKey, command, request, useDirect } = req.body || {}; const ff1Config = getFF1Config(); // Build the DP-1 payload const payload = { command: command || "displayPlaylist", request: request || {} }; // Strategy 1: Try direct connection via SSH tunnel (in dev mode) if (dev && ff1Config?.ip) { try { const result = await castViaSSH(ff1Config, payload); return { success: true, method: "direct-ssh", response: result }; } catch (sshErr) { log(`⚠️ FF1 SSH cast failed: ${sshErr.message}`); // Fall through to cloud relay } } // Strategy 2: Try direct connection (if useDirect or localhost tunnel is running) if (useDirect) { const deviceUrl = `http://localhost:1111/api/cast`; try { log(`📡 FF1 direct cast to ${deviceUrl}`); const directResponse = await fetch(deviceUrl, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), signal: AbortSignal.timeout(5000), // 5s timeout }); if (directResponse.ok) { const result = await directResponse.json(); return { success: true, method: "direct", response: result }; } log(`⚠️ FF1 direct cast failed: ${directResponse.status}`); } catch (directErr) { log(`⚠️ FF1 direct connection failed: ${directErr.message}`); } } // Strategy 3: Use cloud relay with topicID const relayTopicId = topicID || ff1Config?.topicId; if (!relayTopicId) { reply.status(400); return { success: false, error: "No topicID provided and no FF1 config found. Get topicID from your FF1 app settings." }; } const relayUrl = `${FF1_RELAY_URL}?topicID=${encodeURIComponent(relayTopicId)}`; try { log(`☁️ FF1 relay cast to ${relayUrl}`); const headers = { "Content-Type": "application/json" }; if (apiKey || ff1Config?.apiKey) { headers["API-KEY"] = apiKey || ff1Config?.apiKey; } const relayResponse = await fetch(relayUrl, { method: "POST", headers, body: JSON.stringify(payload), signal: AbortSignal.timeout(10000), // 10s timeout }); const responseText = await relayResponse.text(); let responseData; try { responseData = JSON.parse(responseText); } catch { responseData = { raw: responseText }; } if (!relayResponse.ok) { // Check if relay is down (404 or Cloudflare errors) if (relayResponse.status === 404 || responseText.includes("error code:")) { reply.status(503); return { success: false, error: "FF1 cloud relay is unavailable", hint: "The Feral File relay service appears to be down. Use ac-ff1 tunnel for local development.", details: responseData }; } reply.status(relayResponse.status); return { success: false, error: `FF1 relay error: ${relayResponse.status}`, details: responseData }; } return { success: true, method: "relay", response: responseData }; } catch (relayErr) { reply.status(500); return { success: false, error: relayErr.message }; }});
// FF1 CORS preflightfastify.options("/ff1/cast", async (req, reply) => { reply.header("Access-Control-Allow-Origin", "*"); reply.header("Access-Control-Allow-Methods", "POST, OPTIONS"); reply.header("Access-Control-Allow-Headers", "Content-Type"); return "";});
// *** Chat Log Endpoint (for system logs from other services) ***fastify.post("/chat/log", async (req, reply) => { const host = req.headers.host; // Determine which chat instance based on a header or default to chat-system const chatHost = req.headers["x-chat-instance"] || "chat-system.aesthetic.computer"; const instance = chatManager.getInstance(chatHost); if (!instance) { reply.status(404); return { status: "error", message: "Unknown chat instance" }; } const result = await chatManager.handleLog(instance, req.body, req.headers.authorization); reply.status(result.status); return result.body;});
// *** Chat Status Endpoint ***fastify.get("/chat/status", async (req) => { return chatManager.getStatus();});
const PROFILE_SECRET_CACHE_MS = 60 * 1000;let profileSecretCacheValue = null;let profileSecretCacheAt = 0;
function pickProfileStreamSecret(record) { if (!record || typeof record !== "object") return null; const candidates = [ record.secret, record.token, record.profileSecret, record.value, ]; for (const raw of candidates) { if (!raw) continue; const value = `${raw}`.trim(); if (value) return value; } return null;}
function profileSecretsMatch(expected, provided) { if (!expected || !provided) return false; const left = Buffer.from(expected); const right = Buffer.from(provided); if (left.length !== right.length) return false; try { return crypto.timingSafeEqual(left, right); } catch (_) { return false; }}
async function resolveProfileStreamSecret() { const now = Date.now(); if (profileSecretCacheAt && now - profileSecretCacheAt < PROFILE_SECRET_CACHE_MS) { return profileSecretCacheValue; }
let resolved = null; try { if (chatManager?.db) { const record = await chatManager.db .collection("secrets") .findOne({ _id: "profile-stream" }); resolved = pickProfileStreamSecret(record); } } catch (err) { error("👤 Could not load profile-stream secret from MongoDB:", err?.message || err); }
if (!resolved) { const envSecret = `${process.env.PROFILE_STREAM_SECRET || ""}`.trim(); resolved = envSecret || null; }
profileSecretCacheValue = resolved; profileSecretCacheAt = now; return profileSecretCacheValue;}
// *** Profile Stream Event Ingest ***// Accepts server-to-server profile events from Netlify functions.fastify.post("/profile-event", async (req, reply) => { try { const expectedSecret = await resolveProfileStreamSecret(); const providedSecret = `${req.headers["x-profile-secret"] || ""}`.trim() || null; if (expectedSecret && !profileSecretsMatch(expectedSecret, providedSecret)) { reply.status(401); return { ok: false, error: "Unauthorized" }; }
const body = req.body || {}; const handle = body.handle; const handleKey = normalizeProfileHandle(handle); if (!handleKey) { reply.status(400); return { ok: false, error: "Missing or invalid handle" }; }
if (body.event && typeof body.event === "object") { emitProfileActivity(handleKey, body.event); }
if (body.counts && typeof body.counts === "object") { broadcastProfileStream(handleKey, "counts:update", { handle: handleKey, counts: body.counts, }); }
if (body.countsDelta && typeof body.countsDelta === "object") { emitProfileCountDelta(handleKey, body.countsDelta); }
if (body.presence && typeof body.presence === "object") { broadcastProfileStream(handleKey, "presence:update", { handle: handleKey, reason: body.reason || "external", changed: Array.isArray(body.changed) ? body.changed : [], presence: body.presence, }); }
return { ok: true }; } catch (err) { error("👤 profile-event ingest failed:", err); reply.status(500); return { ok: false, error: err.message }; }});
// *** Live Reload of Pieces in Development ***if (dev) { fastify.post("/reload", async (req) => { everyone(pack("reload", req.body, "pieces")); return { msg: "Reload request sent!", body: req.body }; }); // Jump to a specific piece (navigate) fastify.post("/jump", async (req) => { const { piece } = req.body; // Broadcast to all browser clients everyone(pack("jump", { piece }, "pieces")); // Send direct message to VSCode extension clients vscodeClients.forEach(client => { if (client?.readyState === WebSocket.OPEN) { client.send(pack("vscode:jump", { piece }, "vscode")); } }); return { msg: "Jump request sent!", piece, vscodeConnected: vscodeClients.size > 0 }; }); // GET /devices - List all connected clients with metadata and names fastify.get("/devices", async () => { const clientList = getClientStatus(); // Enhance with device names and letters const enhanced = clientList.map((c, index) => ({ ...c, letter: getDeviceLetter(c.id), deviceName: deviceNames[c.ip]?.name || null, deviceGroup: deviceNames[c.ip]?.group || null, })); return { devices: enhanced, host: { name: DEV_HOST_NAME, ip: DEV_LAN_IP }, timestamp: Date.now() }; }); // GET /dev-info - Get dev host info for client overlay fastify.get("/dev-info", async (req, reply) => { // Add CORS headers for cross-origin requests from main site reply.header("Access-Control-Allow-Origin", "*"); reply.header("Access-Control-Allow-Methods", "GET"); return { host: DEV_HOST_NAME, ip: DEV_LAN_IP, mode: "LAN Dev", timestamp: Date.now() }; }); // POST /jump/:target - Targeted jump (by ID, IP, handle, or device name) fastify.post("/jump/:target", async (req) => { const { target } = req.params; const { piece, ahistorical, alias } = req.body; const targeted = targetClients(target); if (targeted.length === 0) { return { error: "No matching device", target }; } targeted.forEach(({ ws }) => { ws.send(pack("jump", { piece, ahistorical, alias }, "pieces")); }); return { msg: "Targeted jump sent", piece, count: targeted.length, targets: targeted.map(t => t.id) }; }); // POST /reload/:target - Targeted reload fastify.post("/reload/:target", async (req) => { const { target } = req.params; const targeted = targetClients(target); targeted.forEach(({ ws }) => { ws.send(pack("reload", req.body, "pieces")); }); return { msg: "Targeted reload sent", count: targeted.length }; }); // POST /piece-reload/:target - Targeted KidLisp reload fastify.post("/piece-reload/:target", async (req) => { const { target } = req.params; const { source, createCode, authToken } = req.body; const targeted = targetClients(target); targeted.forEach(({ ws }) => { ws.send(pack("piece-reload", { source, createCode, authToken }, "kidlisp")); }); return { msg: "Targeted piece-reload sent", count: targeted.length }; }); // POST /device/name - Set a friendly name for a device by IP fastify.post("/device/name", async (req) => { const { ip, name, group } = req.body; if (!ip) return { error: "IP required" }; const cleanIp = ip.replace('::ffff:', ''); if (name) { deviceNames[cleanIp] = { name, group: group || null, updatedAt: Date.now() }; } else { delete deviceNames[cleanIp]; } saveDeviceNames(); // Notify the device of its new name const targeted = targetClients(cleanIp); targeted.forEach(({ ws }) => { ws.send(pack("dev:identity", { name, host: DEV_HOST_NAME, hostIp: DEV_LAN_IP, mode: "LAN Dev" }, "dev")); }); return { msg: name ? "Device named" : "Device name cleared", ip: cleanIp, name, notified: targeted.length }; }); // GET /device/names - List all device names fastify.get("/device/names", async () => { return { names: deviceNames }; });}
// *** HTTP Server Initialization ***
// Track UDP channels manually (geckos.io doesn't expose this)const udpChannels = {};
// 🩰 Initialize geckos.io BEFORE server starts listening// Configure for devcontainer/Docker environment:// - iceServers: Use local TURN server for relay (required in Docker/devcontainer)// - portRange: constrain UDP to small range that can be exposed from container// - cors: allow from any origin in dev mode
// Detect external IP for TURN server (browsers need to reach TURN from outside container)// In devcontainer, we expose ports to the host, so use the host's LAN IP// Priority: TURN_HOST env var > DEV_LAN_IP > localhostconst getExternalTurnHost = () => { // Check for explicitly set TURN host if (process.env.TURN_HOST) return process.env.TURN_HOST; // Use the DEV_LAN_IP if available (detected earlier) if (DEV_LAN_IP) return DEV_LAN_IP; // Fallback to localhost (won't work for external clients but ok for local testing) return 'localhost';};
const turnHost = getExternalTurnHost();console.log("🩰 TURN server host for ICE:", turnHost);
const devIceServers = [ { urls: `stun:${turnHost}:3478` }, { urls: `turn:${turnHost}:3478`, username: 'aesthetic', credential: 'computer123' },];const prodIceServers = [ { urls: 'stun:stun.l.google.com:19302' }, // TODO: Add production TURN server];
const io = geckos({ iceServers: dev ? devIceServers : prodIceServers, // Force relay-only mode in dev to work through container networking // Direct UDP won't work from host browser to container internal IP iceTransportPolicy: dev ? 'relay' : 'all', portRange: { min: 10000, max: 10007, }, cors: { allowAuthorization: true, origin: dev ? "*" : (req) => { const allowed = ["https://aesthetic.computer", "https://notepat.com", "https://kidlisp.com", "https://pj.kidlisp.com", "https://laklok.com", "https://www.laklok.com"]; const reqOrigin = req.headers?.origin; return allowed.includes(reqOrigin) ? reqOrigin : allowed[0]; }, },});io.addServer(server); // Hook up to the HTTP Server - must be before listen()console.log("🩰 Geckos.io server attached to fastify server (UDP ports 10000-10007)");
const start = async () => { try { if (dev) { fastify.listen({ host: "0.0.0.0", // ip.address(), port: info.port, }); } else { fastify.listen({ host: "0.0.0.0", port: info.port }); } } catch (err) { fastify.log.error(err); process.exit(1); }};
await start();
// *** Status Page Data Collection ***
// Get unified client status - user-centric viewfunction getClientStatus() { const identityMap = new Map(); // Map by identity (handle or user or IP) // Helper to get identity key for a client const getIdentityKey = (client) => { // Priority: handle > user > IP (for grouping same person) if (client.handle) return `handle:${client.handle}`; if (client.user) return `user:${client.user}`; if (client.ip) return `ip:${client.ip}`; return null; }; // Process all WebSocket connections Object.keys(connections).forEach((id) => { const client = clients[id] || {}; const ws = connections[id]; const identityKey = getIdentityKey(client); if (!identityKey) return; // Skip if no identity info if (!identityMap.has(identityKey)) { identityMap.set(identityKey, { handle: client.handle || null, location: client.location || null, ip: client.ip || null, geo: client.geo || null, connectionIds: { websocket: [], udp: [] }, protocols: { websocket: false, udp: false }, connections: { websocket: [], udp: [] } }); } const identity = identityMap.get(identityKey); // Update with latest info if (client.handle && !identity.handle) identity.handle = client.handle; if (client.location) identity.location = client.location; if (client.ip && !identity.ip) identity.ip = client.ip; if (client.geo && !identity.geo) identity.geo = client.geo; identity.connectionIds.websocket.push(parseInt(id)); identity.protocols.websocket = true; identity.connections.websocket.push({ id: parseInt(id), alive: ws.isAlive || false, readyState: ws.readyState, ping: ws.lastPing || null, codeChannel: findCodeChannel(parseInt(id)), worlds: getWorldMemberships(parseInt(id)) }); }); // Process all UDP connections Object.keys(udpChannels).forEach((id) => { const client = clients[id] || {}; const udp = udpChannels[id]; const identityKey = getIdentityKey(client); if (!identityKey) return; // Skip if no identity info if (!identityMap.has(identityKey)) { identityMap.set(identityKey, { handle: client.handle || null, location: client.location || null, ip: client.ip || null, geo: client.geo || null, connectionIds: { websocket: [], udp: [] }, protocols: { websocket: false, udp: false }, connections: { websocket: [], udp: [] } }); } const identity = identityMap.get(identityKey); // Update with latest info if (client.handle && !identity.handle) identity.handle = client.handle; if (client.location) identity.location = client.location; if (client.ip && !identity.ip) identity.ip = client.ip; if (client.geo && !identity.geo) identity.geo = client.geo; identity.connectionIds.udp.push(id); identity.protocols.udp = true; identity.connections.udp.push({ id: id, connectedAt: udp.connectedAt, state: udp.state || 'unknown' }); }); // Convert to array and add summary info return Array.from(identityMap.values()).map(identity => { const wsCount = identity.connectionIds.websocket.length; const udpCount = identity.connectionIds.udp.length; const totalConnections = wsCount + udpCount; return { handle: identity.handle, location: identity.location, ip: identity.ip, geo: identity.geo, protocols: identity.protocols, connectionCount: { websocket: wsCount, udp: udpCount, total: totalConnections }, // Simplified connection info - just take first of each type for display websocket: identity.connections.websocket.length > 0 ? identity.connections.websocket[0] : null, udp: identity.connections.udp.length > 0 ? identity.connections.udp[0] : null, multipleTabs: totalConnections > 1 }; });}
function getWorldMemberships(connectionId) { const worlds = []; Object.keys(worldClients).forEach(piece => { if (worldClients[piece][connectionId]) { worlds.push({ piece, handle: worldClients[piece][connectionId].handle, showing: worldClients[piece][connectionId].showing, ghost: worldClients[piece][connectionId].ghost || false, }); } }); return worlds;}
function findCodeChannel(connectionId) { for (const [channel, subscribers] of Object.entries(codeChannels)) { if (subscribers.has(connectionId)) return channel; } return null;}
function getFullStatus() { const clientList = getClientStatus(); // Get chat status with recent messages const chatStatus = chatManager.getStatus(); const chatWithMessages = chatStatus.map(instance => { // Don't expose sotce chat messages — it's a paid subscriber network. const isSotce = instance.name === "chat-sotce"; const recentMessages = (!isSotce && instance.messages > 0) ? chatManager.getRecentMessages(instance.host, 5) : []; return { ...instance, recentMessages }; }); // Filter old errors const cutoff = Date.now() - ERROR_RETENTION_MS; const recentErrors = errorLog.filter(e => new Date(e.timestamp).getTime() > cutoff); return { timestamp: Date.now(), server: { uptime: process.uptime(), environment: dev ? "development" : "production", port: info.port, }, totals: { websocket: wss.clients.size, udp: Object.keys(udpChannels).length, unique_clients: clientList.length }, clients: clientList, chat: chatWithMessages, errors: recentErrors.slice(-20).reverse() // Most recent first };}
// *** Socket Server Initialization ***// #region socketlet wss;let connections = {}; // All active WebSocket connections.const worldClients = {}; // All connected 🧒 to a space like `field`.
let connectionId = 0; // TODO: Eventually replace with a username arrived at through// a client <-> server authentication function.
wss = new WebSocketServer({ server });log( `🤖 session.aesthetic.computer (${ dev ? "Development" : "Production" }) socket: wss://${ip.address()}:${info.port}`,);
// *** Status Page Routes (defined after wss initialization) ***// Status JSON endpointfastify.get("/status", async (request, reply) => { return getFullStatus();});
// Status dashboard HTML at rootfastify.get("/", async (request, reply) => { reply.type("text/html"); return `<!DOCTYPE html><html><head> <meta charset="utf-8"> <meta name="robots" content="noindex, nofollow"> <title>session-server</title> <style> * { margin: 0; padding: 0; box-sizing: border-box; } body { font-family: monospace; background: #000; color: #0f0; padding: 1.5rem; line-height: 1.5; } .header { border-bottom: 1px solid #333; padding-bottom: 1rem; margin-bottom: 1.5rem; } .header h1 { color: #0ff; font-size: 1.2rem; } .header .status { color: #888; font-size: 0.9rem; margin-top: 0.5rem; } .grid { display: grid; grid-template-columns: 1fr 1fr; gap: 1.5rem; } @media (max-width: 900px) { .grid { grid-template-columns: 1fr; } } .section { background: #0a0a0a; border: 1px solid #222; border-radius: 4px; padding: 1rem; } .section h2 { color: #0ff; font-size: 0.95rem; margin-bottom: 0.75rem; border-bottom: 1px solid #222; padding-bottom: 0.5rem; } .client { background: #111; border-left: 3px solid #0f0; padding: 0.75rem; margin-bottom: 0.75rem; } .name { color: #0ff; font-weight: bold; } .ping { color: yellow; } .detail { color: #888; margin-top: 0.2rem; font-size: 0.85rem; } .empty { color: #555; font-style: italic; } .chat-instance { background: #111; border-left: 3px solid #f0f; padding: 0.75rem; margin-bottom: 0.75rem; } .chat-instance.offline { border-left-color: #f00; opacity: 0.6; } .chat-instance .name { color: #f0f; } .chat-msg { background: #0a0a0a; padding: 0.4rem 0.6rem; margin-top: 0.4rem; font-size: 0.8rem; border-radius: 3px; } .chat-msg .from { color: #0ff; } .chat-msg .text { color: #aaa; } .chat-msg .time { color: #555; font-size: 0.75rem; } .error-log { background: #1a0000; border-left: 3px solid #f00; padding: 0.5rem; margin-bottom: 0.5rem; font-size: 0.8rem; } .error-log .time { color: #555; } .error-log .msg { color: #f66; } .warn-log { background: #1a1a00; border-left: 3px solid #ff0; } .warn-log .msg { color: #ff6; } .no-errors { color: #0f0; font-style: italic; } .tabs { display: flex; gap: 0.5rem; margin-bottom: 1rem; } .tab { padding: 0.4rem 0.8rem; background: #111; border: 1px solid #333; color: #888; cursor: pointer; border-radius: 3px; font-family: monospace; font-size: 0.85rem; } .tab.active { background: #0f0; color: #000; border-color: #0f0; } .tab-content { display: none; } .tab-content.active { display: block; } </style></head><body> <div class="header"> <h1>🧩 session-server</h1> <div class="status"> <span id="ws-status">🔴</span> | Uptime: <span id="uptime">--</span> | Online: <span id="client-count">0</span> | Chat: <span id="chat-count">0</span> </div> </div> <div class="tabs"> <button class="tab active" data-tab="overview">Overview</button> <button class="tab" data-tab="chat">💬 Chat</button> <button class="tab" data-tab="errors">⚠️ Errors</button> </div> <div id="overview" class="tab-content active"> <div class="grid"> <div class="section"> <h2>🧑💻 Connected Clients</h2> <div id="clients"></div> </div> <div class="section"> <h2>💬 Chat Instances</h2> <div id="chat-status"></div> </div> </div> </div> <div id="chat" class="tab-content"> <div class="grid"> <div class="section" id="chat-system-section"> <h2>💬 chat-system</h2> <div id="chat-system-messages"></div> </div> <div class="section" id="chat-clock-section"> <h2>🕐 chat-clock</h2> <div id="chat-clock-messages"></div> </div>
</div> </div> <div id="errors" class="tab-content"> <div class="section"> <h2>⚠️ Recent Errors & Warnings</h2> <div id="error-log"></div> </div> </div>
<script> // Tab switching document.querySelectorAll('.tab').forEach(tab => { tab.addEventListener('click', () => { document.querySelectorAll('.tab').forEach(t => t.classList.remove('active')); document.querySelectorAll('.tab-content').forEach(c => c.classList.remove('active')); tab.classList.add('active'); document.getElementById(tab.dataset.tab).classList.add('active'); }); }); const ws = new WebSocket(\`\${location.protocol === 'https:' ? 'wss:' : 'ws:'}//\${location.host}/status-stream\`); ws.onopen = () => { document.getElementById('ws-status').innerHTML = '🟢'; }; ws.onclose = () => { document.getElementById('ws-status').innerHTML = '🔴'; setTimeout(() => location.reload(), 2000); }; ws.onmessage = (event) => { const data = JSON.parse(event.data); if (data.type === 'status') update(data.data); }; function formatTime(dateStr) { if (!dateStr) return ''; const d = new Date(dateStr); return d.toLocaleTimeString('en-US', { hour: '2-digit', minute: '2-digit' }); } function escapeHtml(str) { if (!str) return ''; return str.replace(/&/g, '&').replace(/</g, '<').replace(/>/g, '>'); } function update(s) { const hrs = Math.floor(s.server.uptime / 3600); const min = Math.floor((s.server.uptime % 3600) / 60); document.getElementById('uptime').textContent = \`\${hrs}h \${min}m\`; document.getElementById('client-count').textContent = s.totals.unique_clients; // Chat instance count const totalChatters = s.chat ? s.chat.reduce((sum, c) => sum + c.connections, 0) : 0; document.getElementById('chat-count').textContent = totalChatters; // Clients section const clientsHtml = s.clients.length === 0 ? '<div class="empty">Nobody online</div>' : s.clients.map(c => { let out = '<div class="client">'; out += '<div class="name">'; out += escapeHtml(c.handle) || '(anonymous)'; if (c.multipleTabs && c.connectionCount.total > 1) out += \` (×\${c.connectionCount.total})\`; if (c.websocket?.ping) out += \` <span class="ping">(\${c.websocket.ping}ms)</span>\`; out += '</div>'; if (c.location && c.location !== '*keep-alive*') out += \`<div class="detail">📍 \${escapeHtml(c.location)}</div>\`; if (c.geo) { let geo = '🗺️ '; if (c.geo.city) geo += c.geo.city + ', '; if (c.geo.region) geo += c.geo.region + ', '; geo += c.geo.country; out += \`<div class="detail">\${geo}</div>\`; } else if (c.ip) { out += \`<div class="detail">🌐 \${c.ip}</div>\`; } if (c.websocket?.worlds?.length > 0) { const w = c.websocket.worlds[0]; out += \`<div class="detail">🌍 \${escapeHtml(w.piece)}\`; if (w.showing) out += \` (viewing \${escapeHtml(w.showing)})\`; if (w.ghost) out += ' 👻'; out += '</div>'; } const p = []; if (c.protocols.websocket) p.push(c.connectionCount.websocket > 1 ? \`ws×\${c.connectionCount.websocket}\` : 'ws'); if (c.protocols.udp) p.push(c.connectionCount.udp > 1 ? \`udp×\${c.connectionCount.udp}\` : 'udp'); if (p.length) out += \`<div class="detail" style="opacity:0.5">\${p.join(' + ')}</div>\`; out += '</div>'; return out; }).join(''); document.getElementById('clients').innerHTML = clientsHtml; // Chat status section (overview) if (s.chat) { const chatHtml = s.chat.map(c => { const isOnline = c.messages >= 0; return \`<div class="chat-instance \${isOnline ? '' : 'offline'}"> <div class="name">\${escapeHtml(c.name)} \${isOnline ? '🟢' : '🔴'}</div> <div class="detail">🧑🤝🧑 \${c.connections} connected</div> <div class="detail">💾 \${c.messages} messages loaded</div> </div>\`; }).join(''); document.getElementById('chat-status').innerHTML = chatHtml; } else { document.getElementById('chat-status').innerHTML = '<div class="empty">Chat not initialized</div>'; } // Chat messages (detailed view) if (s.chat) { s.chat.forEach(c => { const name = c.name.replace('chat-', ''); const el = document.getElementById(\`chat-\${name}-messages\`) || document.getElementById(\`chat-\${c.name}-messages\`); if (el && c.recentMessages) { const msgsHtml = c.recentMessages.length === 0 ? '<div class="empty">No recent messages</div>' : c.recentMessages.map(m => \`<div class="chat-msg"> <span class="from">\${escapeHtml(m.from)}</span> <span class="text">\${escapeHtml(m.text)}</span> <span class="time">\${formatTime(m.when)}</span> </div>\`).join(''); el.innerHTML = msgsHtml; } }); } // Error log if (s.errors && s.errors.length > 0) { const errHtml = s.errors.map(e => \`<div class="\${e.level === 'error' ? 'error-log' : 'warn-log error-log'}"> <span class="time">[\${formatTime(e.timestamp)}]</span> <span class="msg">\${escapeHtml(e.message)}</span> </div>\`).join(''); document.getElementById('error-log').innerHTML = errHtml; } else { document.getElementById('error-log').innerHTML = '<div class="no-errors">✅ No errors in the last hour</div>'; } } </script></body></html>`;});
// Pack messages into a simple object protocol of `{type, content}`.function pack(type, content, id) { return JSON.stringify({ type, content, id });}
// Enable ping-pong behavior to keep connections persistently tracked.// (In the future could just tie connections to logged in users or// persistent tokens to keep persistence.)const interval = setInterval(function ping() { wss.clients.forEach((client) => { if (client.isAlive === false) { return client.terminate(); } client.isAlive = false; client.pingStart = Date.now(); // Start ping timer client.ping(); });}, 15000); // 15 second pings from server before termination.
wss.on("close", function close() { clearInterval(interval); connections = {};});
// Construct the server.wss.on("connection", async (ws, req) => { const connectionInfo = { url: req.url, host: req.headers.host, origin: req.headers.origin, userAgent: req.headers['user-agent'], remoteAddress: req.socket.remoteAddress, }; log('🔌 WebSocket connection received:', JSON.stringify(connectionInfo, null, 2)); log('🔌 Total wss.clients.size:', wss.clients.size); log('🔌 Current connections count:', Object.keys(connections).length); // Route status dashboard WebSocket connections separately if (req.url === '/status-stream') { log('📊 Status dashboard viewer connected from:', req.socket.remoteAddress); statusClients.add(ws); // Mark as dashboard viewer (don't add to game clients) ws.isDashboardViewer = true; // Send initial state ws.send(JSON.stringify({ type: 'status', data: getFullStatus(), })); ws.on('close', () => { log('📊 Status dashboard viewer disconnected'); statusClients.delete(ws); }); ws.on('error', (err) => { error('📊 Status dashboard error:', err); statusClients.delete(ws); }); return; // Don't process as a game client }
// Route targeted profile stream connections if (req.url?.startsWith('/profile-stream')) { let requestedHandle = null; try { const parsedUrl = new URL(req.url, 'http://localhost'); requestedHandle = parsedUrl.searchParams.get('handle'); } catch (err) { error('👤 Invalid profile-stream URL:', err); }
const key = addProfileStreamClient(ws, requestedHandle); if (!key) { ws.send( JSON.stringify({ type: 'profile:error', data: { message: 'Missing or invalid handle query param.' }, }), ); try { ws.close(); } catch (_) {} return; }
log('👤 Profile stream viewer connected for:', key, 'from:', req.socket.remoteAddress);
ws.on('close', () => { removeProfileStreamClient(ws); log('👤 Profile stream viewer disconnected for:', key); });
ws.on('error', (err) => { error('👤 Profile stream error:', err); removeProfileStreamClient(ws); });
return; // Don't process as a game client } // Route chat connections to ChatManager based on host const host = req.headers.host; if (chatManager.isChatHost(host)) { log('💬 Chat client connection from:', host); chatManager.handleConnection(ws, req); return; // Don't process as a game client } // Route AC Machines connections — device monitoring & remote commands if (req.url.startsWith('/machines')) { const urlParams = new URL(req.url, 'http://localhost').searchParams; const role = urlParams.get('role') || 'device'; const token = urlParams.get('token') || ''; const machineId = urlParams.get('machineId') || '';
if (role === 'viewer') { // Browser dashboard viewer — verify AC auth token via Auth0 const authUser = await verifyACToken(token); if (!authUser?.sub) { ws.close(4001, 'Unauthorized'); return; } const userSub = authUser.sub; const userHandle = authUser.nickname || authUser.name || null;
log(`Machines viewer connected: ${userHandle || userSub}`);
if (!machinesViewers.has(userSub)) machinesViewers.set(userSub, new Set()); machinesViewers.get(userSub).add(ws);
// Send initial state: all online machines for this user const userMachines = []; for (const [mid, device] of machinesDevices) { if (device.user === userSub) { userMachines.push({ machineId: mid, ...device.info, status: "online", lastHeartbeat: device.lastHeartbeat, }); } } ws.send(JSON.stringify({ type: "machines-state", machines: userMachines }));
// Handle viewer → device commands ws.on('message', (data) => { try { const msg = JSON.parse(data.toString()); if (msg.type === "command" && msg.machineId) { const device = machinesDevices.get(msg.machineId); if (device && device.user === userSub && device.ws.readyState === WebSocket.OPEN) { const commandId = Date.now().toString(36) + Math.random().toString(36).slice(2, 6); device.ws.send(JSON.stringify({ type: "command", command: msg.cmd, commandId, target: msg.args?.target || msg.args?.piece || undefined, // Free-text payload for cmd:"prompt" — runs through the // device's prompt.mjs execute() exactly as if typed locally. text: typeof msg.args?.text === "string" ? msg.args.text : undefined, })); log(`Command '${msg.cmd}' → ${msg.machineId} (${commandId})`); } } // Swank eval: forward CL expression to device for evaluation if (msg.type === "swank:eval" && msg.machineId && msg.expr) { const device = machinesDevices.get(msg.machineId); if (device && device.user === userSub && device.ws.readyState === WebSocket.OPEN) { const evalId = Date.now().toString(36) + Math.random().toString(36).slice(2, 6); device.ws.send(JSON.stringify({ type: "swank:eval", expr: msg.expr, evalId, })); log(`🔮 Swank eval → ${msg.machineId}: ${msg.expr.slice(0, 60)}`); } } } catch (e) { error('🖥️ Machines viewer message error:', e); } });
ws.on('close', () => { log(`🖥️ Machines viewer disconnected: ${userHandle || userSub}`); const viewers = machinesViewers.get(userSub); if (viewers) { viewers.delete(ws); if (viewers.size === 0) machinesViewers.delete(userSub); } });
ws.on('error', (err) => { error('🖥️ Machines viewer error:', err); const viewers = machinesViewers.get(userSub); if (viewers) { viewers.delete(ws); if (viewers.size === 0) machinesViewers.delete(userSub); } });
} else { // Device connection const tokenPayload = await verifyMachineToken(token); const userSub = tokenPayload?.sub || null; const userHandle = tokenPayload?.handle || null; const linked = !!tokenPayload;
log(`📡 Machines device connected: ${machineId} (${linked ? userHandle : 'unlinked'})`);
machinesDevices.set(machineId, { ws, user: userSub, handle: userHandle, machineId, linked, info: {}, lastHeartbeat: Date.now(), });
if (userSub) { broadcastToMachineViewers(userSub, { type: "device-connected", machineId, linked }); }
ws.on('message', async (data) => { try { const msg = JSON.parse(data.toString()); const device = machinesDevices.get(machineId); if (!device) return;
switch (msg.type) { case "register": device.info = { version: msg.version, buildName: msg.buildName, gitHash: msg.gitHash, buildTs: msg.buildTs, hw: msg.hw, ip: msg.ip, wifiSSID: msg.wifiSSID, hostname: msg.hostname, label: msg.label, currentPiece: msg.currentPiece || "notepat", }; device.lastHeartbeat = Date.now(); try { await upsertMachine(userSub, machineId, device.info); } catch (e) { error("📡 upsert:", e.message); } if (userSub) broadcastToMachineViewers(userSub, { type: "machine-registered", machineId, ...device.info, status: "online" }); break;
case "heartbeat": device.lastHeartbeat = Date.now(); device.info.uptime = msg.uptime; device.info.currentPiece = msg.currentPiece || device.info.currentPiece; device.info.battery = msg.battery; device.info.charging = msg.charging; device.info.fps = msg.fps; try { await updateMachineHeartbeat(userSub, machineId, msg.uptime, device.info.currentPiece); } catch (e) { error("📡 heartbeat:", e.message); } if (userSub) broadcastToMachineViewers(userSub, { type: "heartbeat", machineId, uptime: msg.uptime, currentPiece: device.info.currentPiece, battery: msg.battery, charging: msg.charging, fps: msg.fps, timestamp: Date.now(), }); break;
case "log": try { await insertMachineLog(userSub, machineId, msg); } catch (e) { error("📡 log insert:", e.message); } if (userSub) { const logMessage = msg.message || (typeof msg.data === "string" ? msg.data : JSON.stringify(msg.data)); broadcastToMachineViewers(userSub, { type: "log", machineId, level: msg.logType === "crash" ? "error" : (msg.level || "info"), message: logMessage, logType: msg.logType || "log", data: msg.data || null, when: msg.when || new Date().toISOString(), }); } break;
case "command-ack": case "command-response": if (userSub) broadcastToMachineViewers(userSub, { type: msg.type, machineId, commandId: msg.commandId, command: msg.command, data: msg.data }); break;
case "swank:result": // Forward Swank eval result from device to viewer if (userSub) broadcastToMachineViewers(userSub, { type: "swank:result", machineId, evalId: msg.evalId, ok: msg.ok, result: msg.result, }); break; } } catch (e) { error('📡 Machines device message error:', e); } });
ws.on('close', async () => { log(`📡 Machines device disconnected: ${machineId}`); machinesDevices.delete(machineId); if (userSub) { broadcastToMachineViewers(userSub, { type: "status-change", machineId, status: "offline" }); try { await setMachineOffline(userSub, machineId); } catch (e) { error("📡 offline:", e.message); } } });
ws.on('error', (err) => { error(`📡 Machines device error (${machineId}):`, err); machinesDevices.delete(machineId); }); }
return; // Don't process as a game client }
// Route socklogs connections - devices sending logs and viewers subscribing if (req.url.startsWith('/socklogs')) { const urlParams = new URL(req.url, 'http://localhost').searchParams; const role = urlParams.get('role') || 'device'; // 'device' or 'viewer' const deviceId = urlParams.get('deviceId') || `device-${Date.now()}`; if (role === 'viewer') { // Viewer wants to see logs from devices log('👁️ SockLogs viewer connected'); socklogsViewers.add(ws); // Send current status ws.send(JSON.stringify({ type: 'status', ...socklogsStatus() })); ws.on('close', () => { log('👁️ SockLogs viewer disconnected'); socklogsViewers.delete(ws); }); ws.on('error', (err) => { error('👁️ SockLogs viewer error:', err); socklogsViewers.delete(ws); }); } else { // Device sending logs log(`📱 SockLogs device connected: ${deviceId}`); socklogsDevices.set(deviceId, { ws, logCount: 0, lastLog: null, connectedAt: Date.now() }); // Notify viewers of new device for (const viewer of socklogsViewers) { if (viewer.readyState === WebSocket.OPEN) { viewer.send(JSON.stringify({ type: 'device-connected', deviceId, status: socklogsStatus() })); } } ws.on('message', (data) => { try { const msg = JSON.parse(data.toString()); if (msg.type === 'log') { const device = socklogsDevices.get(deviceId); if (device) { device.logCount++; device.lastLog = Date.now(); } socklogsBroadcast(deviceId, msg); } } catch (e) { error('📱 SockLogs parse error:', e); } }); ws.on('close', () => { log(`📱 SockLogs device disconnected: ${deviceId}`); socklogsDevices.delete(deviceId); // Notify viewers for (const viewer of socklogsViewers) { if (viewer.readyState === WebSocket.OPEN) { viewer.send(JSON.stringify({ type: 'device-disconnected', deviceId, status: socklogsStatus() })); } } }); ws.on('error', (err) => { error(`📱 SockLogs device error (${deviceId}):`, err); socklogsDevices.delete(deviceId); }); } return; // Don't process as a game client } log('🎮 Game client connection detected, adding to connections'); // Regular game client connection handling below const ip = req.socket.remoteAddress || "localhost"; // beautify ip ws.isAlive = true; // For checking persistence between ping-pong messages. ws.pingStart = null; // Track ping timing ws.lastPing = null; // Store last measured ping
ws.on("pong", () => { ws.isAlive = true; if (ws.pingStart) { ws.lastPing = Date.now() - ws.pingStart; ws.pingStart = null; } }); // Receive a pong and stay alive!
// Assign the conection a unique id. connections[connectionId] = ws; const id = connectionId; let codeChannel; // Used to subscribe to incoming piece code. // Initialize client record with IP and geolocation if (!clients[id]) clients[id] = {}; clients[id].websocket = true; // Clean IP and get geolocation const cleanIp = ip.replace('::ffff:', ''); clients[id].ip = cleanIp; const geo = geoip.lookup(cleanIp); if (geo) { clients[id].geo = { country: geo.country, region: geo.region, city: geo.city, timezone: geo.timezone, ll: geo.ll // [latitude, longitude] }; log(`🌍 Geolocation for ${cleanIp}:`, geo.country, geo.region, geo.city); } else { log(`🌍 No geolocation data for ${cleanIp}`); }
log("🧏 Someone joined:", `${id}:${ip}`, "Online:", wss.clients.size, "🫂"); log("🎮 Added to connections. Total game clients:", Object.keys(connections).length);
const content = { id, playerCount: wss.clients.size };
// Send a message to all other clients except this one. function others(string) { wss.clients.forEach((c) => { if (c !== ws && c?.readyState === WebSocket.OPEN) c.send(string); }); }
// Send a self-connection message back to the client. ws.send( pack( "connected", JSON.stringify({ ip, playerCount: content.playerCount }), id, ), );
// In dev mode, send device identity info for LAN overlay if (dev) { const deviceName = deviceNames[cleanIp]?.name || null; const deviceLetter = getDeviceLetter(id); const identityPayload = { name: deviceName, letter: deviceLetter, host: DEV_HOST_NAME, hostIp: DEV_LAN_IP, mode: "LAN Dev", connectionId: id, }; console.log(`📱 Sending dev:identity to ${cleanIp}:`, identityPayload); ws.send(pack("dev:identity", identityPayload, "dev")); }
// Send a join message to everyone else. others( pack( "joined", JSON.stringify({ text: `${connectionId} has joined. Connections open: ${content.playerCount}`, }), id, ), );
connectionId += 1;
// Relay all incoming messages from this client to everyone else. ws.on("message", (data) => { // Parse incoming message and attach client identifier. let msg; try { msg = JSON.parse(data.toString()); } catch (error) { console.error("📚 Failed to parse JSON:", error); return; }
// 📦 Module streaming - handle module requests before other processing if (msg.type === "module:request") { const modulePath = msg.path; const withDeps = msg.withDeps === true; // Request all dependencies too const knownHashes = msg.knownHashes || {}; // Client's cached hashes if (withDeps) { // Recursively gather all dependencies const modules = {}; let skippedCount = 0; const gatherDeps = (p, fromPath = null) => { if (modules[p] || modules[p] === null) return; // Already gathered (or marked as cached) const data = getModuleHash(p); if (!data) { // Only warn for top-level not found, not for deps (which might be optional) if (!fromPath) log(`📦 Module not found: ${p}`); return; } // Check if client already has this hash cached if (knownHashes[p] === data.hash) { modules[p] = null; // Mark as "client has it" - don't send content skippedCount++; } else { modules[p] = { hash: data.hash, content: data.content }; } // Debug: show when gathering specific important modules if (p.includes('headers') || p.includes('kidlisp')) { log(`📦 Gathering ${p} (from ${fromPath || 'top'})${knownHashes[p] === data.hash ? ' [cached]' : ''}`); } // Parse static imports from content - match ES module import/export from statements // This regex only matches valid relative imports ending in .mjs or .js // Skip commented lines by checking each line doesn't start with // const staticImportRegex = /^(?!\s*\/\/).*?(?:import|export)\s+(?:[^;]*?\s+from\s+)?["'](\.{1,2}\/[^"'\s]+\.m?js)["']/gm; let match; while ((match = staticImportRegex.exec(data.content)) !== null) { const importPath = match[1]; // Skip invalid paths if (importPath.includes('...') || importPath.length > 200) continue; // Resolve relative path const dir = path.dirname(p); const resolved = path.normalize(path.join(dir, importPath)); log(`📦 Found dep: ${p} -> ${importPath} (resolved: ${resolved})`); gatherDeps(resolved, p); } // Parse dynamic imports - import("./path") or import('./path') or import(`./path`) // Skip commented lines const dynamicImportRegex = /^(?!\s*\/\/).*?import\s*\(\s*["'`](\.{1,2}\/[^"'`\s]+\.m?js)["'`]\s*\)/gm; while ((match = dynamicImportRegex.exec(data.content)) !== null) { const importPath = match[1]; // Skip invalid paths if (importPath.includes('...') || importPath.length > 200) continue; // Resolve relative path const dir = path.dirname(p); const resolved = path.normalize(path.join(dir, importPath)); gatherDeps(resolved, p); } }; gatherDeps(modulePath); // Filter out null entries (modules client already has) and count const modulesToSend = {}; const cachedPaths = []; for (const [p, data] of Object.entries(modules)) { if (data === null) { cachedPaths.push(p); } else { modulesToSend[p] = data; } } const totalModules = Object.keys(modules).length; const sentModules = Object.keys(modulesToSend).length; if (totalModules > 0) { // Log bundle stats if (skippedCount > 0) { log(`📦 Bundle for ${modulePath}: ${sentModules}/${totalModules} sent (${skippedCount} cached)`); } else { log(`📦 Bundle for ${modulePath}: ${sentModules} modules`); } ws.send(JSON.stringify({ type: "module:bundle", entry: modulePath, modules: modulesToSend, cached: cachedPaths // Tell client which paths to use from cache })); } else { ws.send(JSON.stringify({ type: "module:error", path: modulePath, error: "Module not found" })); } } else { // Single module request (original behavior) const moduleData = getModuleHash(modulePath); if (moduleData) { ws.send(JSON.stringify({ type: "module:response", path: modulePath, hash: moduleData.hash, content: moduleData.content })); log(`📦 Module sent: ${modulePath} (${moduleData.content.length} bytes)`); } else { ws.send(JSON.stringify({ type: "module:error", path: modulePath, error: "Module not found" })); log(`📦 Module not found: ${modulePath}`); } } return; } if (msg.type === "module:check") { const modulePath = msg.path; const clientHash = msg.hash; const moduleData = getModuleHash(modulePath); if (moduleData) { ws.send(JSON.stringify({ type: "module:status", path: modulePath, changed: moduleData.hash !== clientHash, hash: moduleData.hash })); } else { ws.send(JSON.stringify({ type: "module:status", path: modulePath, changed: true, hash: null, error: "Module not found" })); } return; } if (msg.type === "module:list") { // Return list of available modules (for prefetching) const modules = [ "lib/disk.mjs", "lib/graph.mjs", "lib/num.mjs", "lib/geo.mjs", "lib/parse.mjs", "lib/help.mjs", "lib/text.mjs", "bios.mjs" ]; const moduleInfo = modules.map(p => { const data = getModuleHash(p); return data ? { path: p, hash: data.hash, size: data.content.length } : null; }).filter(Boolean); ws.send(JSON.stringify({ type: "module:list", modules: moduleInfo })); return; }
// 🎹 DAW Channel - M4L device ↔ IDE communication if (msg.type === "daw:join") { // Device (kidlisp.com/device) joining to receive code dawDevices.add(id); log(`🎹 DAW device joined: ${id} (total: ${dawDevices.size})`); ws.send(JSON.stringify({ type: "daw:joined", id })); return; } if (msg.type === "daw:code") { // IDE sending code to all connected devices log(`🎹 DAW code broadcast from ${id} to ${dawDevices.size} devices`); const codeMsg = JSON.stringify({ type: "daw:code", content: msg.content, from: id }); // Broadcast to all DAW devices for (const deviceId of dawDevices) { const deviceWs = connections[deviceId]; if (deviceWs && deviceWs.readyState === WebSocket.OPEN) { deviceWs.send(codeMsg); log(`🎹 Sent code to device ${deviceId}`); } } return; }
if (msg.type === "notepat:midi:sources") { sendNotepatMidiSources(ws); return; }
if (msg.type === "notepat:midi:subscribe") { const filter = msg.content || {}; addNotepatMidiSubscriber(id, ws, filter); return; }
if (msg.type === "notepat:midi:unsubscribe") { removeNotepatMidiSubscriber(id); if (ws.readyState === WebSocket.OPEN) { ws.send(pack("notepat:midi:unsubscribed", true, "midi-relay")); } return; }
// 🎹 Browser-side notepat (notepat.com on web) publishes MIDI events // via WS. Native bare-metal sends the same payload via raw UDP at port // 10010 (see handleNotepatMidiUdpPacket). Both paths converge through // broadcastNotepatMidiEvent, fanning out to WS + UDP subscribers (M4L // notepat-remote devices listen on those channels). Heartbeats keep // the source visible in `notepat:midi:sources` between notes. if (msg.type === "notepat:midi:publish" || msg.type === "notepat:midi:heartbeat") { const payload = msg.content || {}; const now = Date.now(); const source = upsertNotepatMidiSource({ handle: payload.handle, machineId: payload.machineId, piece: payload.piece || "notepat", lastEvent: msg.type === "notepat:midi:heartbeat" ? "heartbeat" : payload.event, ts: now, }); if (msg.type === "notepat:midi:heartbeat") return; if (!source.handle && !source.machineId) return; const rawNote = Number(payload.note); const rawVelocity = Number(payload.velocity); const rawChannel = Number(payload.channel); if (!Number.isFinite(rawNote) || !Number.isFinite(rawVelocity) || !Number.isFinite(rawChannel)) return; let event = payload.event === "note_off" ? "note_off" : "note_on"; const note = Math.max(0, Math.min(127, Math.round(rawNote))); const velocity = Math.max(0, Math.min(127, Math.round(rawVelocity))); const channel = Math.max(0, Math.min(15, Math.round(rawChannel))); if (event === "note_on" && velocity === 0) event = "note_off"; broadcastNotepatMidiEvent({ type: "notepat:midi", event, note, velocity, channel, handle: source.handle, machineId: source.machineId, piece: source.piece || "notepat", ts: Number.isFinite(Number(payload.ts)) ? Number(payload.ts) : now, }); return; }
msg.id = id; // TODO: When sending a server generated message, use a special id.
// Extract user identity and handle from ANY message that contains it if (msg.content?.user?.sub) { if (!clients[id]) clients[id] = { websocket: true }; const userSub = msg.content.user.sub; const userChanged = !clients[id].user || clients[id].user !== userSub; if (userChanged) { clients[id].user = userSub; log("🔑 User identity from", msg.type + ":", userSub.substring(0, 20) + "...", "conn:", id); } // Extract handle from message if present (e.g., location:broadcast includes it) if (msg.content.handle && (!clients[id].handle || clients[id].handle !== msg.content.handle)) { clients[id].handle = msg.content.handle; log("✅ Handle from message:", msg.content.handle, "conn:", id); emitProfilePresence(msg.content.handle, "identify", ["handle"]); } }
if (msg.type === "scream") { // Alert all connected users via redis pub/sub to the scream. log("😱 About to scream..."); const out = filter(msg.content); pub .publish("scream", out) .then((result) => { log("😱 Scream succesfully published:", result);
let piece = ""; if (out.indexOf("pond") > -1) piece = "pond"; else if (out.indexOf("field") > -1) piece = "field";
if (chatManager?.db) { broadcastToTopic( chatManager.db, "scream", { title: "😱 Scream", body: out, urgent: true, // time-sensitive on iOS, Urgency: high on web ttl: 0, // don't store undelivered screams data: { piece }, }, log, ).catch((error) => { log("📵 Error sending notification:", error); }); } }) .catch((error) => { log("🙅♀️ Error publishing scream:", error); }); // Send a notification to all devices subscribed to the `scream` topic. } else if (msg.type === "code-channel:sub") { // Filter code-channel updates based on this user. codeChannel = msg.content; if (!codeChannels[codeChannel]) codeChannels[codeChannel] = new Set(); codeChannels[codeChannel].add(id); // Send current channel state to late joiners if (codeChannelState[codeChannel]) { // Note: codeChannelState stores the original msg.content object, // pack() will JSON.stringify it, so don't double-stringify here const stateMsg = pack("code", codeChannelState[codeChannel], id); send(stateMsg); log(`📥 Sent current state to late joiner on channel ${codeChannel}`); } } else if (msg.type === "code-channel:info") { // Return viewer count for a code channel const ch = msg.content; const count = codeChannels[ch]?.size || 0; send(pack("code-channel:info", { channel: ch, viewers: count }, id)); } else if (msg.type === "slide" && msg.content?.codeChannel) { // Handle slide broadcast (low-latency value updates, no state storage) const targetChannel = msg.content.codeChannel; // Don't store slide updates as state (they're transient) // Just broadcast immediately for low latency if (codeChannels[targetChannel]) { const slideMsg = pack("slide", msg.content, id); subscribers(codeChannels[targetChannel], slideMsg); } } else if (msg.type === "code" && msg.content?.codeChannel) { // Handle code broadcast to channel subscribers (for kidlisp.com pop-out sync) const targetChannel = msg.content.codeChannel; // Store the latest state for late joiners codeChannelState[targetChannel] = msg.content; if (codeChannels[targetChannel]) { // Note: msg.content is already an object, pack() will JSON.stringify it const codeMsg = pack("code", msg.content, id); subscribers(codeChannels[targetChannel], codeMsg); log(`📢 Broadcast code to channel ${targetChannel} (${codeChannels[targetChannel].size} subscribers)`); } } else if (msg.type === "login") { if (msg.content?.user?.sub) { if (!clients[id]) clients[id] = { websocket: true }; clients[id].user = msg.content.user.sub; // Fetch the user's handle from the API const userSub = msg.content.user.sub; log("🔑 Login attempt for user:", userSub.substring(0, 20) + "...", "connection:", id); fetch(`https://aesthetic.computer/handle/${encodeURIComponent(userSub)}`) .then(response => { log("📡 Handle API response status:", response.status, "for", userSub.substring(0, 20) + "..."); return response.json(); }) .then(data => { log("📦 Handle API data:", JSON.stringify(data), "for connection:", id); if (data.handle) { clients[id].handle = data.handle; log("✅ User logged in:", data.handle, `(${userSub.substring(0, 12)}...)`, "connection:", id); emitProfilePresence(data.handle, "login", ["handle", "online", "connections"]); } else { log("⚠️ User logged in (no handle in response):", userSub.substring(0, 12), "..., connection:", id); } }) .catch(err => { log("❌ Failed to fetch handle for:", userSub.substring(0, 20) + "...", "Error:", err.message); }); } } else if (msg.type === "identify") { // VSCode extension identifying itself if (msg.content?.type === "vscode") { vscodeClients.add(ws); log("✅ VSCode extension connected, conn:", id); // Send confirmation ws.send(pack("identified", { type: "vscode", id }, id)); } } else if (msg.type === "dev:log") { // 📡 Remote log forwarding from connected devices (LAN Dev mode) if (dev && msg.content) { const { level, args, deviceName, connectionId, time, queued } = msg.content; const client = clients[id]; const deviceLabel = deviceName || client?.ip || `conn:${connectionId}`; const levelEmoji = level === 'error' ? '🔴' : level === 'warn' ? '🟡' : '🔵'; const queuedTag = queued ? ' [Q]' : ''; // Format the log output const timestamp = new Date(time).toLocaleTimeString(); const argsStr = Array.isArray(args) ? args.join(' ') : String(args); console.log(`${levelEmoji} [${timestamp}] ${deviceLabel}${queuedTag}: ${argsStr}`); } } else if (msg.type === "location:broadcast") { /* sub .subscribe(`logout:broadcast:${msg.content.user.sub}`, () => { ws.send(pack(`logout:broadcast:${msg.content.user.sub}`, true, id)); }) .then(() => { log("🏃 Subscribed to logout updates from:", msg.content.user.sub); }) .catch((err) => error( "🏃 Could not unsubscribe from logout:broadcast for:", msg.content.user.sub, err, ), ); */ } else if (msg.type === "logout:broadcast:subscribe") { /* console.log("Logout broadcast:", msg.type, msg.content); pub .publish(`logout:broadcast:${msg.content.user.sub}`, "true") .then((result) => { console.log("🏃 Logout broadcast successful for:", msg.content); }) .catch((error) => { log("🙅♀️ Error publishing logout:", error); }); */ } else if (msg.type === "location:broadcast") { // Receive a slug location for this handle. if (msg.content.slug !== "*keep-alive*") { log("🗼 Location:", msg.content.slug, "Handle:", msg.content.handle, "ID:", id); } // Store handle and location for this client if (!clients[id]) clients[id] = { websocket: true }; const previousLocation = clients[id].location; // Extract user identity from message if (msg.content?.user?.sub) { clients[id].user = msg.content.user.sub; } // Extract handle directly from message if (msg.content.handle) { clients[id].handle = msg.content.handle; } // Extract and store location if (msg.content.slug) { // Don't overwrite location with keep-alive if (msg.content.slug !== "*keep-alive*") { clients[id].location = msg.content.slug; log(`📍 Location updated for ${clients[id].handle || id}: "${msg.content.slug}"`); if (previousLocation !== msg.content.slug) { emitProfileActivity(msg.content.handle || clients[id].handle, { type: "piece", when: Date.now(), label: `Piece ${msg.content.slug}`, ref: msg.content.slug, }); } } else { log(`💓 Keep-alive from ${clients[id].handle || id}, location unchanged`); } }
emitProfilePresence( msg.content.handle || clients[id].handle, "location:broadcast", ["online", "currentPiece", "connections"], );
// Publish to redis... pub .publish("slug:" + msg.content.handle, msg.content.slug) .then((result) => { if (msg.content.slug !== "*keep-alive*") { log( "🐛 Slug succesfully published for:", msg.content.handle, msg.content.slug, ); } }) .catch((error) => { log("🙅♀️ Error publishing slug:", error); });
// TODO: - [] When a user is ghosted, then subscribe to their location // updates. // - [] And stop subscribing when they are unghosted. } else if (msg.type === "dev-log" && dev) { // Create device-specific log files and only notify in terminal const timestamp = new Date().toISOString(); const deviceId = `client-${id}`; const logFileName = `${DEV_LOG_DIR}${deviceId}.log`; // Check if this is a new device if (!deviceLogFiles.has(deviceId)) { deviceLogFiles.set(deviceId, logFileName); console.log(`📱 New device logging: ${deviceId} -> ${logFileName}`); console.log(` tail -f ${logFileName}`); } // Write to device-specific log file const logEntry = `[${timestamp}] ${msg.content.level || 'LOG'}: ${msg.content.message}\n`; try { fs.appendFileSync(logFileName, logEntry); } catch (error) { console.error(`Failed to write to ${logFileName}:`, error); } } else { // 🗺️ World Messages // TODO: Should all messages be prefixed with their piece?
// Filter for `world:${piece}:${label}` type messages. if (msg.type.startsWith("world:")) { const parsed = msg.type.split(":"); const piece = parsed[1]; const label = parsed.pop(); const worldHandle = resolveProfileHandle(id, piece, msg.content?.handle);
// TODO: Store client position on disconnect, based on their handle.
if (label === "show") { // Store any existing show picture in clients list. worldClients[piece][id].showing = msg.content; emitProfileActivity(worldHandle, { type: "show", when: Date.now(), label: `Showing in ${piece}`, piece, ref: piece, }); emitProfilePresence(worldHandle, `world:${piece}:show`, ["world", "showing"]); }
if (label === "hide") { // Store any existing show picture in clients list. worldClients[piece][id].showing = null; emitProfileActivity(worldHandle, { type: "hide", when: Date.now(), label: `Hide in ${piece}`, piece, ref: piece, }); emitProfilePresence(worldHandle, `world:${piece}:hide`, ["world", "showing"]); }
// Intercept chats and filter them (skip for laer-klokken). if (label === "write") { if (piece !== "laer-klokken") msg.content = filter(msg.content); const chatText = typeof msg.content === "string" ? msg.content : msg.content?.text; if (chatText) { emitProfileActivity(worldHandle, { type: "chat", when: Date.now(), label: `Chat ${piece}: ${truncateProfileText(chatText, 80)}`, piece, ref: piece, text: chatText, }); emitProfileCountDelta(worldHandle, { chats: 1 }); } }
if (label === "join") { if (!worldClients[piece]) worldClients[piece] = {};
// Check to see if the client handle matches and a connection can // be reassociated.
let pickedUpConnection = false; keys(worldClients[piece]).forEach((clientID) => { // TODO: Break out of this loop early. const client = worldClients[piece][clientID]; if ( client["handle"].startsWith("@") && client["handle"] === msg.content.handle && client.ghosted ) { // log("👻 Ghosted?", client);
log( "👻 Unghosting:", msg.content.handle, "old id:", clientID, "new id:", id, ); pickedUpConnection = true;
client.ghosted = false;
sub .unsubscribe("slug:" + msg.content.handle) .then(() => { log("🐛 Unsubscribed from slug for:", msg.content.handle); }) .catch((err) => { error( "🐛 Could not unsubscribe from slug for:", msg.content.handle, err, ); });
delete worldClients[piece][clientID];
ws.send(pack(`world:${piece}:list`, worldClients[piece], id));
// Replace the old client with the new data. worldClients[piece][id] = { ...msg.content }; } });
if (!pickedUpConnection) ws.send(pack(`world:${piece}:list`, worldClients[piece], id));
// ❤️🔥 TODO: No need to send the current user back via `list` here. if (!pickedUpConnection) worldClients[piece][id] = { ...msg.content };
// ^ Send existing list to just this user.
others(JSON.stringify(msg)); // Alert everyone else about the join.
log("🧩 Clients in piece:", piece, worldClients[piece]); emitProfileActivity(worldHandle, { type: "join", when: Date.now(), label: `Joined ${piece}`, piece, ref: piece, }); emitProfilePresence(worldHandle, `world:${piece}:join`, ["world", "connections"]); return; } else if (label === "move") { // log("🚶♂️", piece, msg.content); if (typeof worldClients?.[piece]?.[id] === "object") worldClients[piece][id].pos = msg.content.pos; } else { log(`${label}:`, msg.content); }
if (label === "persist") { log("🧮 Persisting this client...", msg.content); }
// All world: messages are only broadcast to "others", with the // exception of "write" with relays the filtered message back: if (label === "write") { everyone(JSON.stringify(msg)); } else { others(JSON.stringify(msg)); } return; }
// 🎮 1v1 game position updates should only go to others (not back to sender) if (msg.type === "1v1:move") { // Log occasionally in production for debugging (1 in 100 messages) if (Math.random() < 0.01) { log(`🎮 1v1:move relay: ${msg.content?.handle || id} -> ${wss.clients.size - 1} others`); } others(JSON.stringify(msg)); return; } // 🎾 Squash game position updates — relay to others only (not back to sender) if (msg.type === "squash:move") { others(JSON.stringify(msg)); return; }
// 🔊 Audio data from kidlisp.com — relay only to code-channel subscribers if (msg.type === "audio" && msg.content?.codeChannel) { const ch = msg.content.codeChannel; if (codeChannels[ch]) { subscribers(codeChannels[ch], pack("audio", msg.content, id)); } return; }
// 🎮 1v1 join/state messages - log and relay to everyone if (msg.type === "1v1:join" || msg.type === "1v1:state") { log(`🎮 ${msg.type}: ${msg.content?.handle || id} -> all ${wss.clients.size} clients`); }
// 🎯 Duel messages — routed to DuelManager (server-authoritative) if (msg.type === "duel:join") { const handle = typeof msg.content === "string" ? JSON.parse(msg.content).handle : msg.content?.handle; if (handle) duelManager.playerJoin(handle, id); return; } if (msg.type === "duel:leave") { const handle = typeof msg.content === "string" ? JSON.parse(msg.content).handle : msg.content?.handle; if (handle) duelManager.playerLeave(handle); return; } if (msg.type === "duel:ping") { const parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; if (parsed?.handle) duelManager.handlePing(parsed.handle, parsed.ts, id); return; } if (msg.type === "duel:clientlog") { const parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; console.log(`🎯 [CLIENT ${parsed?.handle}] ${parsed?.msg}`, parsed?.bullets?.length > 0 ? JSON.stringify(parsed.bullets) : ""); return; } if (msg.type === "duel:input") { const parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; if (parsed?.handle) duelManager.receiveInput(parsed.handle, parsed); return; }
if (msg.type === "fight:join") { const parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; if (parsed?.handle) { if (!clients[id]) clients[id] = {}; clients[id].handle = parsed.handle; fightManager.join(parsed.handle, id); } return; } if (msg.type === "fight:leave") { const parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; if (parsed?.handle) fightManager.leave(parsed.handle, id); return; }
// 🌐 World messages (arena:*, land:*) — routed to the matching // WorldManager. Generic over the worldManagers map so new worlds // don't need new routing code. Consumes ALL messages with a known // world prefix (unknown verbs are dropped, not relayed to everyone). { const sep = msg.type.indexOf(":"); const wm = sep > 0 ? worldManagers[msg.type.slice(0, sep)] : undefined; if (wm) { const verb = msg.type.slice(sep + 1); let parsed; try { parsed = typeof msg.content === "string" ? JSON.parse(msg.content) : msg.content; } catch { return; } if (verb === "hello" && parsed?.handle) { // Ensure clients[id].handle is set so the WS-close handler can // look it up and call playerLeave. Without this, probes + any // client that hasn't sent a chat login message would leak forever. if (!clients[id]) clients[id] = {}; clients[id].handle = parsed.handle; wm.playerJoin(parsed.handle, id, { probe: !!parsed.probe }); } else if (verb === "bye" && parsed?.handle) { wm.playerLeave(parsed.handle); } else if (verb === "cmd" && parsed?.handle) { // WS fallback path; the fast path is the UDP channel.on("<world>:cmd", ...) handler. wm.receiveCmd(parsed.handle, parsed); } else if (verb === "ping" && parsed?.handle) { wm.handlePing(parsed.handle, parsed.ts, id); } return; } }
everyone(JSON.stringify(msg)); // Relay any other message to every user. } });
// More info: https://stackoverflow.com/a/49791634/8146077 ws.on("close", () => { log("🚪 Someone left:", id, "Online:", wss.clients.size, "🫂"); const rawDepartingHandle = clients?.[id]?.handle; const departingHandle = normalizeProfileHandle(rawDepartingHandle); if (departingHandle) duelManager.playerLeave(departingHandle); if (rawDepartingHandle) fightManager.leave(rawDepartingHandle, id); // Worlds use the raw handle (matches <world>:hello), not the @-normalized // form. Pass the closing wsId so a quick reload-race doesn't delete the // freshly-rebound player (the new hello will have set a different wsId). if (rawDepartingHandle) { for (const wm of Object.values(worldManagers)) wm.playerLeave(rawDepartingHandle, id); } removeNotepatMidiSubscriber(id);
// Remove from VSCode clients if present vscodeClients.delete(ws); // Remove from DAW devices if present if (dawDevices.has(id)) { dawDevices.delete(id); log(`🎹 DAW device disconnected: ${id} (remaining: ${dawDevices.size})`); } if (dawIDEs.has(id)) { dawIDEs.delete(id); log(`🎹 DAW IDE disconnected: ${id}`); }
// Delete the user from the worldClients pieces index. // keys(worldClients).forEach((piece) => { // delete worldClients[piece][id]; // if (keys(worldClients[piece]).length === 0) // delete worldClients[piece]; // });
if (clients[id]?.user) { const userSub = clients[id].user; sub .unsubscribe("logout:broadcast:" + userSub) .then(() => { log("🏃 Unsubscribed from logout:broadcast for:", userSub); }) .catch((err) => { error( "🏃 Could not unsubscribe from logout:broadcast for:", userSub, err, ); }); }
// Send a message to everyone else on the server that this client left.
let ghosted = false;
keys(worldClients).forEach((piece) => { if (worldClients[piece][id]) { // Turn this client into a ghost, unless it's the last one in the // world region. if ( worldClients[piece][id].handle.startsWith("@") && keys(worldClients[piece]).length > 1 ) { const handle = worldClients[piece][id].handle; log("👻 Ghosted:", handle); log("World clients after ghosting:", worldClients[piece]); worldClients[piece][id].ghost = true; ghosted = true;
function kick() { log("👢 Kicked:", handle, id); clearTimeout(kickTimer); sub .unsubscribe("slug:" + handle) .then(() => { log("🐛 Unsubscribed from slug for:", handle); }) .catch((err) => { error("🐛 Could not unsubscribe from slug for:", handle, err); }); // Delete the user from the worldClients pieces index. delete worldClients[piece][id]; if (keys(worldClients[piece]).length === 0) delete worldClients[piece]; everyone(pack(`world:${piece}:kick`, {}, id)); // Kick this ghost. }
let kickTimer = setTimeout(kick, 5000);
const worlds = ["field", "horizon"]; // Whitelist for worlds... // This could eventually be communicated based on a parameter in // the subscription? 24.03.09.15.05
// Subscribe to slug updates from redis. sub .subscribe("slug:" + handle, (slug) => { if (slug !== "*keep-alive*") { log(`🐛 ${handle} is now in:`, slug); if (!worlds.includes(slug)) everyone(pack(`world:${piece}:slug`, { handle, slug }, id)); }
if (worlds.includes(slug)) { kick(); } else { clearTimeout(kickTimer); kickTimer = setTimeout(kick, 5000); } // Whitelist slugs here }) .then(() => { log("🐛 Subscribed to slug updates from:", handle); }) .catch((err) => error("🐛 Could not subscribe to slug for:", handle, err), );
// Send a message to everyone on the server that this client is a ghost. everyone(pack(`world:${piece}:ghost`, {}, id)); } else { // Delete the user from the worldClients pieces index. delete worldClients[piece][id]; if (keys(worldClients[piece]).length === 0) delete worldClients[piece]; } } });
// Send a message to everyone else on the server that this client left. if (!ghosted) everyone(pack("left", { count: wss.clients.size }, id));
// Delete from the connection index. delete connections[id]; // Clean up client record if no longer connected via any protocol if (clients[id]) { clients[id].websocket = false; // If also not connected via UDP, delete the client record entirely if (!udpChannels[id]) { delete clients[id]; } }
// Clear out the codeChannel if the last user disconnects from it. if (codeChannel !== undefined) { codeChannels[codeChannel]?.delete(id); if (codeChannels[codeChannel]?.size === 0) { delete codeChannels[codeChannel]; delete codeChannelState[codeChannel]; // Clean up stored state too log(`🗑️ Cleaned up empty channel: ${codeChannel}`); } }
if (departingHandle) { emitProfilePresence(departingHandle, "disconnect", ["online", "connections"]); emitProfileActivity(departingHandle, { type: "presence", when: Date.now(), label: "Disconnected", }); } });});
// Sends a message to all connected clients.function everyone(string) { wss.clients.forEach((c) => { if (c?.readyState === WebSocket.OPEN) c.send(string); });}
// Sends a message to a particular set of client ids on// this instance that are part of the `subs` Set.function subscribers(subs, msg) { subs.forEach((connectionId) => { connections[connectionId]?.send(msg); });}
// 🎯 Wire DuelManager send functionsduelManager.setSendFunctions({ sendUDP: (channelId, event, data) => { const entry = udpChannels[channelId]; if (entry?.channel?.webrtcConnection?.state === "open") { try { entry.channel.emit(event, data); return true; } catch {} } return false; // Signal failure so caller can fall back to WS }, sendWS: (wsId, type, content) => { connections[wsId]?.send(pack(type, JSON.stringify(content), "duel")); }, broadcastWS: (type, content) => { everyone(pack(type, JSON.stringify(content), "duel")); }, resolveUdpForHandle: (handle) => { for (const [id, client] of Object.entries(clients)) { if (client.handle === handle && udpChannels[id]) return id; } return null; },});
fightManager.setSendFunction((wsId, type, content) => { connections[wsId]?.send(pack(type, JSON.stringify(content), "fight"));});
// 🌐 Wire WorldManager send functions (same shape as DuelManager; source tag// = the world prefix). Lifecycle broadcasts are member-scoped inside the// manager, so no broadcastWS/everyone hook here — that was the lobby leak.for (const wm of Object.values(worldManagers)) { wm.setSendFunctions({ sendUDP: (channelId, event, data) => { const entry = udpChannels[channelId]; if (entry?.channel?.webrtcConnection?.state === "open") { try { entry.channel.emit(event, data); return true; } catch {} } return false; }, sendWS: (wsId, type, content) => { connections[wsId]?.send(pack(type, JSON.stringify(content), wm.prefix)); }, resolveUdpForHandle: (handle) => { for (const [id, client] of Object.entries(clients)) { if (client.handle === handle && udpChannels[id]) return id; } return null; }, // Used to distinguish reconnect (old ws dead → silent refresh) from a // real takeover (old ws still alive → demote to spectator). isLive: (wsId) => connections[wsId]?.readyState === WebSocket.OPEN, });}// #endregion
// *** Status WebSocket Stream ***// Track status dashboard clients (separate from game clients)const statusClients = new Set();// Track targeted profile subscribers by normalized handle key (`@name`)const profileStreamClients = new Map();const profileLastSeen = new Map();
// *** VSCode Extension Clients ***// Track VSCode extension clients for direct jump message routingconst vscodeClients = new Set();
function normalizeProfileHandle(handle) { if (!handle) return null; const raw = `${handle}`.trim(); if (!raw) return null; return `@${raw.replace(/^@+/, "").toLowerCase()}`;}
function normalizeMidiHandle(handle) { const normalized = normalizeProfileHandle(handle); return normalized ? normalized.slice(1) : "";}
function notepatMidiSourceKey(handle, machineId) { const handleKey = normalizeProfileHandle(handle) || "@unknown"; const machineKey = `${machineId || "unknown"}`.trim() || "unknown"; return `${handleKey}:${machineKey}`;}
function listNotepatMidiSources() { return [...notepatMidiSources.values()] .sort((a, b) => (b.lastSeen || 0) - (a.lastSeen || 0)) .map((source) => ({ handle: source.handle || null, machineId: source.machineId, piece: source.piece || "notepat", lastSeen: source.lastSeen || 0, lastEvent: source.lastEvent || null, }));}
function sendNotepatMidiSources(ws) { if (!ws || ws.readyState !== WebSocket.OPEN) return; try { ws.send(pack("notepat:midi:sources", { sources: listNotepatMidiSources() }, "midi-relay")); } catch (err) { error("🎹 Failed to send notepat midi sources:", err); }}
function removeNotepatMidiSubscriber(id) { if (id === undefined || id === null) return; notepatMidiSubscribers.delete(id);}
function addNotepatMidiSubscriber(id, ws, filter = {}) { if (id === undefined || id === null || !ws) return;
notepatMidiSubscribers.set(id, { ws, all: filter.all === true, handle: normalizeMidiHandle(filter.handle), machineId: filter.machineId ? `${filter.machineId}`.trim() : "", });
if (ws.readyState === WebSocket.OPEN) { ws.send(pack("notepat:midi:subscribed", { all: filter.all === true, handle: normalizeMidiHandle(filter.handle) || null, machineId: filter.machineId ? `${filter.machineId}`.trim() : null, }, "midi-relay")); }
sendNotepatMidiSources(ws);}
function broadcastNotepatMidiSources() { for (const [id, sub] of notepatMidiSubscribers) { if (!sub?.ws || sub.ws.readyState !== WebSocket.OPEN) { notepatMidiSubscribers.delete(id); continue; } sendNotepatMidiSources(sub.ws); }}
function notepatMidiSubscriberMatches(sub, event) { if (!sub) return false; if (sub.all) return true;
const eventHandle = normalizeMidiHandle(event?.handle); const eventMachine = event?.machineId ? `${event.machineId}`.trim() : "";
if (sub.handle && sub.handle !== eventHandle) return false; if (sub.machineId && sub.machineId !== eventMachine) return false;
return !!(sub.handle || sub.machineId);}
function broadcastNotepatMidiEvent(event) { for (const [id, sub] of notepatMidiSubscribers) { if (!sub?.ws || sub.ws.readyState !== WebSocket.OPEN) { notepatMidiSubscribers.delete(id); continue; } if (!notepatMidiSubscriberMatches(sub, event)) continue; try { sub.ws.send(pack("notepat:midi", event, "midi-relay")); } catch (err) { error("🎹 Failed to fan out notepat midi event:", err); } } // UDP fan-out. Same filter model, emitted on the geckos channel. M4L // notepat-remote devices care about this path for sub-frame latency — // the WS path is ~5-15 ms slower end-to-end on a typical home network. for (const [id, sub] of notepatMidiUdpSubscribers) { if (!sub?.channel || sub.channel.webrtcConnection?.state !== "open") { notepatMidiUdpSubscribers.delete(id); continue; } if (!notepatMidiSubscriberMatches(sub, event)) continue; try { sub.channel.emit("notepat:midi", event); } catch (err) { error("🎹 UDP fan-out failed:", err); } }}
function upsertNotepatMidiSource({ handle, machineId, piece, lastEvent, ts, address, port }) { const cleanHandle = normalizeMidiHandle(handle); const cleanMachineId = `${machineId || "unknown"}`.trim() || "unknown"; const key = notepatMidiSourceKey(cleanHandle, cleanMachineId); const previous = notepatMidiSources.get(key); const next = { handle: cleanHandle || null, machineId: cleanMachineId, piece: piece || "notepat", lastSeen: ts || Date.now(), lastEvent: lastEvent || previous?.lastEvent || null, address: address || previous?.address || null, port: port || previous?.port || null, };
notepatMidiSources.set(key, next);
if (!previous) { log(`🎹 Notepat MIDI source online: ${next.handle ? "@" + next.handle : "@unknown"} ${next.machineId}`); }
if ( !previous || previous.handle !== next.handle || previous.machineId !== next.machineId || previous.piece !== next.piece ) { broadcastNotepatMidiSources(); }
return next;}
function compactProfileText(value) { return `${value || ""}`.replace(/\s+/g, " ").trim();}
function truncateProfileText(value, max = 100) { const text = compactProfileText(value); if (text.length <= max) return text; return `${text.slice(0, Math.max(0, max - 3))}...`;}
function getProfilePresence(handleKey) { if (!handleKey) return null; const clientsForStatus = getClientStatus(); const matched = clientsForStatus.find( (client) => normalizeProfileHandle(client?.handle) === handleKey, );
if (!matched) { return { online: false, currentPiece: null, worldPiece: null, showing: null, connections: { websocket: 0, udp: 0, total: 0 }, pingMs: null, lastSeenAt: profileLastSeen.get(handleKey) || null, }; }
const now = Date.now(); profileLastSeen.set(handleKey, now);
const world = matched?.websocket?.worlds?.[0] || null;
return { online: true, currentPiece: matched.location || null, worldPiece: world?.piece || null, showing: world?.showing || null, connections: matched.connectionCount || { websocket: 0, udp: 0, total: 0 }, pingMs: matched?.websocket?.ping || null, lastSeenAt: now, };}
function sendProfileStream(ws, type, data) { if (!ws || ws.readyState !== WebSocket.OPEN) return; try { ws.send(JSON.stringify({ type, data, timestamp: Date.now() })); } catch (err) { error("👤 Failed to send profile stream event:", err); }}
function broadcastProfileStream(handleKey, type, data) { const subs = profileStreamClients.get(handleKey); if (!subs || subs.size === 0) return;
const stale = []; subs.forEach((ws) => { if (ws.readyState !== WebSocket.OPEN) { stale.push(ws); return; } sendProfileStream(ws, type, data); });
stale.forEach((ws) => subs.delete(ws)); if (subs.size === 0) profileStreamClients.delete(handleKey);}
function addProfileStreamClient(ws, handle) { const handleKey = normalizeProfileHandle(handle); if (!handleKey) return null;
if (!profileStreamClients.has(handleKey)) { profileStreamClients.set(handleKey, new Set()); }
profileStreamClients.get(handleKey).add(ws); ws.profileHandleKey = handleKey;
const presence = getProfilePresence(handleKey); sendProfileStream(ws, "profile:snapshot", { handle: handleKey, presence, }); sendProfileStream(ws, "counts:update", { handle: handleKey, counts: { online: presence?.online ? 1 : 0, connections: presence?.connections?.total || 0, }, });
return handleKey;}
function removeProfileStreamClient(ws) { const handleKey = ws?.profileHandleKey; if (!handleKey) return;
const subs = profileStreamClients.get(handleKey); if (!subs) { ws.profileHandleKey = null; return; }
subs.delete(ws); if (subs.size === 0) profileStreamClients.delete(handleKey); ws.profileHandleKey = null;}
function emitProfilePresence(handle, reason = "update", changed = []) { const handleKey = normalizeProfileHandle(handle); if (!handleKey) return;
const presence = getProfilePresence(handleKey); broadcastProfileStream(handleKey, "presence:update", { handle: handleKey, reason, changed, presence, }); broadcastProfileStream(handleKey, "counts:update", { handle: handleKey, counts: { online: presence?.online ? 1 : 0, connections: presence?.connections?.total || 0, }, });}
function emitProfileCountDelta(handle, delta = {}) { const handleKey = normalizeProfileHandle(handle); if (!handleKey) return; if (!delta || typeof delta !== "object") return;
const cleanDelta = {}; for (const [key, value] of Object.entries(delta)) { const amount = Number(value); if (!Number.isFinite(amount) || amount === 0) continue; cleanDelta[key] = amount; } if (Object.keys(cleanDelta).length === 0) return;
broadcastProfileStream(handleKey, "counts:delta", { handle: handleKey, delta: cleanDelta, });}
function emitProfileActivity(handle, event = {}) { const handleKey = normalizeProfileHandle(handle); if (!handleKey) return;
const label = truncateProfileText( event.label || event.text || event.type || "event", 120, ); if (!label) return;
broadcastProfileStream(handleKey, "activity:append", { handle: handleKey, event: { type: event.type || "event", when: event.when || Date.now(), label, ref: event.ref || null, piece: event.piece || null, }, });}
function resolveProfileHandle(id, piece, fromMessage) { return ( normalizeProfileHandle(fromMessage) || normalizeProfileHandle(clients?.[id]?.handle) || normalizeProfileHandle(worldClients?.[piece]?.[id]?.handle) );}
chatManager.setActivityEmitter((payload = {}) => { const handle = payload.handle; if (payload.event) emitProfileActivity(handle, payload.event); if (payload.countsDelta) emitProfileCountDelta(handle, payload.countsDelta);});
// Broadcast status updates every 2 secondssetInterval(() => { if (statusClients.size > 0) { const status = getFullStatus(); statusClients.forEach(client => { if (client.readyState === WebSocket.OPEN) { try { client.send(JSON.stringify({ type: 'status', data: status })); } catch (err) { error('📊 Failed to send status update:', err); } } }); }}, 2000);
// Broadcast targeted profile heartbeat updates every 2 secondssetInterval(() => { if (profileStreamClients.size === 0) return;
for (const handleKey of profileStreamClients.keys()) { const presence = getProfilePresence(handleKey); broadcastProfileStream(handleKey, "presence:update", { handle: handleKey, reason: "heartbeat", changed: [], presence, }); }}, 2000);
// 🧚 UDP Server (using Twilio ICE servers)// #endregion udp
// Note: This currently works off of a monolith via `udp.aesthetic.computer`// as the ports are blocked on jamsocket.
// geckos.io is imported at top and initialized before server.listen()
io.onConnection((channel) => { // Track this UDP channel udpChannels[channel.id] = { connectedAt: Date.now(), state: channel.webrtcConnection.state, user: null, handle: null, channel: channel, // Store reference for targeted sends }; // Get IP address from channel const udpIp = channel.userData?.address || channel.remoteAddress || null; log(`🩰 UDP ${channel.id} connected from:`, udpIp || 'unknown'); // Initialize client record with IP if (!clients[channel.id]) clients[channel.id] = { udp: true }; if (udpIp) { const cleanIp = udpIp.replace('::ffff:', ''); clients[channel.id].ip = cleanIp; // Get geolocation for UDP client const geo = geoip.lookup(cleanIp); if (geo) { clients[channel.id].geo = { country: geo.country, region: geo.region, city: geo.city, timezone: geo.timezone, ll: geo.ll }; log(`🌍 UDP ${channel.id} geolocation:`, geo.city || geo.country); } } // Set a timeout to warn about missing identity setTimeout(() => { if (!clients[channel.id]?.user && !clients[channel.id]?.handle) { log(`⚠️ UDP ${channel.id} has been connected for 10s but hasn't sent identity message`); } }, 10000); // Handle identity message channel.on("udp:identity", (data) => { try { const identity = JSON.parse(data); log(`🩰 UDP ${channel.id} sent identity:`, JSON.stringify(identity).substring(0, 100)); // Initialize client record if needed if (!clients[channel.id]) clients[channel.id] = { udp: true }; // Extract user identity if (identity.user?.sub) { clients[channel.id].user = identity.user.sub; log(`🩰 UDP ${channel.id} user:`, identity.user.sub.substring(0, 20) + "..."); } // Extract handle directly from identity message if (identity.handle) { clients[channel.id].handle = identity.handle; log(`✅ UDP ${channel.id} handle: "${identity.handle}"`); // Resolve UDP channel for duel if this handle is in a duel duelManager.resolveUdpChannel(identity.handle, channel.id); // Resolve UDP channel for any world this handle is in for (const wm of Object.values(worldManagers)) { wm.resolveUdpChannel(identity.handle, channel.id); } } } catch (e) { error(`🩰 Failed to parse identity for ${channel.id}:`, e); } }); channel.onDisconnect(() => { log(`🩰 ${channel.id} got disconnected`); delete udpChannels[channel.id]; fairyThrottle.delete(channel.id); notepatMidiUdpSubscribers.delete(channel.id);
// Clean up client record if no longer connected via any protocol if (clients[channel.id]) { clients[channel.id].udp = false; // If also not connected via WebSocket, delete the client record entirely if (!connections[channel.id]) { delete clients[channel.id]; } }
channel.close(); });
// 🎹 Notepat MIDI relay over UDP. Same subscribe/unsubscribe model as the // WS path (cross-session, filter on handle/machineId or all:true). Events // fan out via notepatMidiUdpSubscribers in broadcastNotepatMidiEvent. channel.on("notepat:midi:subscribe", (data) => { let filter = {}; try { filter = typeof data === "string" ? JSON.parse(data) : (data || {}); } catch {} // Optional: wrap in `{ filter: {...} }` or pass fields directly — accept both. if (filter.filter) filter = filter.filter; notepatMidiUdpSubscribers.set(channel.id, { channel, all: filter.all === true, handle: normalizeMidiHandle(filter.handle), machineId: filter.machineId ? `${filter.machineId}`.trim() : "", }); try { channel.emit("notepat:midi:subscribed", { all: filter.all === true, handle: normalizeMidiHandle(filter.handle) || null, machineId: filter.machineId ? `${filter.machineId}`.trim() : null, }); } catch {} log(`🎹 UDP ${channel.id} subscribed to notepat:midi (all=${filter.all === true})`); });
channel.on("notepat:midi:unsubscribe", () => { notepatMidiUdpSubscribers.delete(channel.id); try { channel.emit("notepat:midi:unsubscribed", true); } catch {} });
// 🎹 Browser-side notepat (notepat.com on web) publishes MIDI events // via the geckos UDP channel for sub-frame latency. Mirror of the WS // `notepat:midi:publish` handler — same shape, same broadcaster. const handleUdpNotepatPublish = (data, isHeartbeat) => { let payload = {}; try { payload = typeof data === "string" ? JSON.parse(data) : (data || {}); } catch {} const now = Date.now(); const source = upsertNotepatMidiSource({ handle: payload.handle, machineId: payload.machineId, piece: payload.piece || "notepat", lastEvent: isHeartbeat ? "heartbeat" : payload.event, ts: now, }); if (isHeartbeat) return; if (!source.handle && !source.machineId) return; const rawNote = Number(payload.note); const rawVelocity = Number(payload.velocity); const rawChannel = Number(payload.channel); if (!Number.isFinite(rawNote) || !Number.isFinite(rawVelocity) || !Number.isFinite(rawChannel)) return; let event = payload.event === "note_off" ? "note_off" : "note_on"; const note = Math.max(0, Math.min(127, Math.round(rawNote))); const velocity = Math.max(0, Math.min(127, Math.round(rawVelocity))); const channel2 = Math.max(0, Math.min(15, Math.round(rawChannel))); if (event === "note_on" && velocity === 0) event = "note_off"; broadcastNotepatMidiEvent({ type: "notepat:midi", event, note, velocity, channel: channel2, handle: source.handle, machineId: source.machineId, piece: source.piece || "notepat", ts: Number.isFinite(Number(payload.ts)) ? Number(payload.ts) : now, }); }; channel.on("notepat:midi", (data) => handleUdpNotepatPublish(data, false)); channel.on("notepat:midi:heartbeat", (data) => handleUdpNotepatPublish(data, true));
// 💎 TODO: Make these channel names programmable somehow? 24.12.08.04.12
channel.on("tv", (data) => { if (channel.webrtcConnection.state === "open") { try { channel.room.emit("tv", data); } catch (err) { console.warn("Broadcast error:", err); } } else { console.log(channel.webrtcConnection.state); } });
// Just for testing via the aesthetic `udp` piece for now. channel.on("fairy:point", (data) => { // See docs here: https://github.com/geckosio/geckos.io#reliable-messages // TODO: - [] Learn about the differences between channels and rooms.
// log(`🩰 fairy:point - ${data}`); if (channel.webrtcConnection.state === "open") { try { channel.broadcast.emit("fairy:point", data); // ^ emit the to all channels in the same room except the sender
// Bridge to raw UDP clients (native bare-metal) try { const parsed = typeof data === "string" ? JSON.parse(data) : data; const x = parseFloat(parsed.x) || 0; const y = parseFloat(parsed.y) || 0; const handle = parsed.handle || ""; const hlen = Buffer.byteLength(handle, "utf8"); const pkt = Buffer.alloc(10 + hlen); pkt[0] = 0x02; // fairy broadcast pkt.writeFloatLE(x, 1); pkt.writeFloatLE(y, 5); pkt[9] = hlen; pkt.write(handle, 10, "utf8"); for (const [, client] of udpClients) { udpRelay.send(pkt, client.port, client.address); } } catch (e) { /* ignore bridge errors */ }
// Publish to Redis for silo firehose visualization (throttled ~10Hz) const now = Date.now(); const last = fairyThrottle.get(channel.id) || 0; if (now - last >= FAIRY_THROTTLE_MS) { fairyThrottle.set(channel.id, now); pub.publish("fairy:point", data).catch(() => {}); } } catch (err) { console.warn("Broadcast error:", err); } } else { console.log(channel.webrtcConnection.state); } });
// 🎮 1v1 FPS game position updates over UDP (low latency) channel.on("1v1:move", (data) => { if (channel.webrtcConnection.state === "open") { try { // Log occasionally for production debugging (1 in 100) if (Math.random() < 0.01) { const parsed = typeof data === 'string' ? JSON.parse(data) : data; log(`🩰 UDP 1v1:move: ${parsed?.handle || channel.id} broadcasting`); } // Broadcast position to all other players except sender channel.broadcast.emit("1v1:move", data); } catch (err) { console.warn("1v1:move broadcast error:", err); } } });
// 🎾 Squash game position updates over UDP (low latency) channel.on("squash:move", (data) => { if (channel.webrtcConnection.state === "open") { try { channel.broadcast.emit("squash:move", data); } catch (err) { console.warn("squash:move broadcast error:", err); } } });
// 🎯 Duel input over UDP (server-authoritative — NOT relayed, fed to DuelManager) channel.on("duel:input", (data) => { if (channel.webrtcConnection.state === "open") { try { const parsed = typeof data === "string" ? JSON.parse(data) : data; // Resolve handle from channel identity OR from message payload const handle = clients[channel.id]?.handle || parsed.handle; if (handle) { duelManager.receiveInput(handle, parsed); // Also resolve UDP channel if not yet linked if (!clients[channel.id]?.handle && parsed.handle) { duelManager.resolveUdpChannel(parsed.handle, channel.id); } } } catch (err) { console.warn("duel:input error:", err); } } });
// 🌐 World usercmds over UDP (fast path; WS is the fallback) for (const wm of Object.values(worldManagers)) { channel.on(`${wm.prefix}:cmd`, (data) => { if (channel.webrtcConnection.state !== "open") return; try { const parsed = typeof data === "string" ? JSON.parse(data) : data; const handle = clients[channel.id]?.handle || parsed.handle; if (!handle) return; wm.receiveCmd(handle, parsed); // Re-resolve on every UDP cmd: a working client→server channel is // the freshest promotion signal. The manager itself blocks // re-promotion while its stall watchdog has demoted this player. wm.resolveUdpChannel(handle, channel.id); } catch (err) { console.warn(`${wm.prefix}:cmd error:`, err); } }); }
// 🎚️ Slide mode: real-time code value updates via UDP (lowest latency) channel.on("slide:code", (data) => { if (channel.webrtcConnection.state === "open") { try { // Broadcast to all including sender (room.emit) for sync channel.room.emit("slide:code", data); } catch (err) { console.warn("slide:code broadcast error:", err); } } });
// 🔊 Audio: real-time audio analysis data via UDP (lowest latency) channel.on("udp:audio", (data) => { if (channel.webrtcConnection.state === "open") { try { channel.room.emit("udp:audio", data); } catch (err) { console.warn("udp:audio broadcast error:", err); } } });});
// #endregion
// ---------------------------------------------------------------------------// 🧚 Raw UDP fairy relay (port 10010) — for native bare-metal clients// Binary packet format:// [1 byte type] [4 float x LE] [4 float y LE] [1 handle_len] [N handle]// Type 0x01 = client→server, 0x02 = server→client broadcast// ---------------------------------------------------------------------------const UDP_FAIRY_PORT = 10010;
function handleNotepatMidiUdpPacket(payload, rinfo) { if (!payload || (payload.type !== "notepat:midi" && payload.type !== "notepat:midi:heartbeat")) { return false; }
const now = Date.now(); const source = upsertNotepatMidiSource({ handle: payload.handle, machineId: payload.machineId, piece: payload.piece || "notepat", lastEvent: payload.type === "notepat:midi" ? payload.event : "heartbeat", ts: now, address: rinfo.address, port: rinfo.port, });
if (!source.handle && !source.machineId) { return true; }
if (payload.type === "notepat:midi:heartbeat") { return true; }
const rawNote = Number(payload.note); const rawVelocity = Number(payload.velocity); const rawChannel = Number(payload.channel); if (!Number.isFinite(rawNote) || !Number.isFinite(rawVelocity) || !Number.isFinite(rawChannel)) { log("🎹 Invalid notepat midi UDP payload:", payload); return true; }
let event = payload.event === "note_off" ? "note_off" : "note_on"; const note = Math.max(0, Math.min(127, Math.round(rawNote))); const velocity = Math.max(0, Math.min(127, Math.round(rawVelocity))); const channel = Math.max(0, Math.min(15, Math.round(rawChannel))); if (event === "note_on" && velocity === 0) event = "note_off";
broadcastNotepatMidiEvent({ type: "notepat:midi", event, note, velocity, channel, handle: source.handle, machineId: source.machineId, piece: source.piece || "notepat", ts: Number.isFinite(Number(payload.ts)) ? Number(payload.ts) : now, });
return true;}
function pruneNotepatMidiSources() { const now = Date.now(); let changed = false;
for (const [key, source] of notepatMidiSources) { if (now - (source.lastSeen || 0) > UDP_MIDI_SOURCE_TTL_MS) { notepatMidiSources.delete(key); changed = true; } }
if (changed) broadcastNotepatMidiSources();}
udpRelay.on("message", (msg, rinfo) => { if (msg.length > 0 && msg[0] === 0x01 && msg.length >= 10) { const key = `${rinfo.address}:${rinfo.port}`; const x = msg.readFloatLE(1); const y = msg.readFloatLE(5); const hlen = msg[9]; const handle = msg.slice(10, 10 + hlen).toString("utf8");
udpClients.set(key, { address: rinfo.address, port: rinfo.port, handle, lastSeen: Date.now() });
// Build broadcast packet (type 0x02) const bcast = Buffer.alloc(msg.length); msg.copy(bcast); bcast[0] = 0x02;
// Broadcast to all other UDP clients for (const [k, client] of udpClients) { if (k !== key) { udpRelay.send(bcast, client.port, client.address); } }
// Also broadcast to Geckos.io WebRTC clients as fairy:point const fairyData = JSON.stringify({ x, y, handle }); try { // Emit to all geckos channels io.room().emit("fairy:point", fairyData); } catch (e) { /* ignore */ }
// Publish to Redis for silo firehose (throttled) const now = Date.now(); const lastFairy = fairyThrottle.get(key) || 0; if (now - lastFairy >= FAIRY_THROTTLE_MS) { fairyThrottle.set(key, now); pub.publish("fairy:point", fairyData).catch(() => {}); } return; }
if (msg.length > 0 && msg[0] === 0x7b) { try { const payload = JSON.parse(msg.toString("utf8")); if (handleNotepatMidiUdpPacket(payload, rinfo)) return; } catch (err) { log("🎹 Failed to parse UDP JSON packet:", err?.message || err); } }});
// Clean up stale UDP clients every 30ssetInterval(() => { const now = Date.now(); for (const [key, client] of udpClients) { if (now - client.lastSeen > 30000) udpClients.delete(key); } pruneNotepatMidiSources();}, 30000);
udpRelay.bind(UDP_FAIRY_PORT, () => { console.log(`🧚 Raw UDP fairy relay listening on port ${UDP_FAIRY_PORT}`);});
// Bridge: forward Geckos fairy:point to UDP clients// (patched into the existing fairy:point handler above via io.room().emit)// When a Geckos client sends fairy:point, also relay to UDP clients:const origFairyHandler = true; // marker — actual bridging done in channel.on("fairy:point") below
// #endregion UDP fairy relay
// 🚧 File Watching in Local Development Mode// File watching uses: https://github.com/paulmillr/chokidarif (dev) { // 1. Watch for local file changes in pieces. chokidar .watch("../system/public/aesthetic.computer/disks") .on("all", (event, path) => { if (event === "change") { const piece = path .split("/") .pop() .replace(/\.mjs|\.lisp$/, ""); everyone(pack("reload", { piece: piece || "*" }, "local")); } }); // 2. Watch base system files. chokidar .watch([ "../system/netlify/functions", "../system/public/privacy-policy.html", "../system/public/aesthetic-direct.html", "../system/public/aesthetic.computer/lib", "../system/public/aesthetic.computer/systems", // This doesn't need a full reload / could just reload the disk module? "../system/public/aesthetic.computer/boot.mjs", "../system/public/aesthetic.computer/bios.mjs", "../system/public/aesthetic.computer/style.css", "../system/public/kidlisp.com", "../system/public/l5.aesthetic.computer", "../system/public/gift.aesthetic.computer", "../system/public/give.aesthetic.computer", "../system/public/news.aesthetic.computer", ]) .on("all", (event, path) => { if (event === "change") everyone(pack("reload", { piece: "*refresh*" }, "local")); });
// 2b. Watch prompt files separately (piece reload instead of full refresh) chokidar .watch("../system/public/aesthetic.computer/prompts") .on("all", (event, path) => { if (event === "change") { const filename = path.split("/").pop(); console.log(`🎨 Prompt file changed: ${filename}`); everyone(pack("reload", { piece: "*piece-reload*" }, "local")); } });
// 3. Watch vscode extension chokidar.watch("../vscode-extension/out").on("all", (event, path) => { if (event === "change") everyone(pack("vscode-extension:reload", { reload: true }, "local")); });}
/*if (termkit) { term = termkit.terminal;
const doc = term.createDocument({ palette: new termkit.Palette(), });
// Create left (log) and right (client list) columns const leftColumn = new termkit.Container({ parent: doc, x: 0, width: "70%", height: "100%", });
const rightColumn = new termkit.Container({ parent: doc, x: "70%", width: "30%", height: "100%", });
term.grabInput();
console.log("grabbed input");
term.on("key", function (name, matches, data) { console.log("'key' event:", name);
// Detect CTRL-C and exit 'manually' if (name === "CTRL_C") { process.exit(); } });
term.on("mouse", function (name, data) { console.log("'mouse' event:", name, data); });
// Log box in the left column const logBox = new termkit.TextBox({ parent: leftColumn, content: "Your logs will appear here...\n", scrollable: true, vScrollBar: true, x: 0, y: 0, width: "100%", height: "100%", mouse: true, // to allow mouse interactions if needed });
// Static list box in the right column const clientList = new termkit.TextBox({ parent: rightColumn, content: "Client List:\n", x: 0, y: 0, width: "100%", height: "100%", });
// Example functions to update contents function addLog(message) { logBox.setContent(logBox.getContent() + message + "\n"); // logBox.scrollBottom(); doc.draw(); }
function updateClientList(clients) { clientList.setContent("Client List:\n" + clients.join("\n")); doc.draw(); }
// Example usage addLog("Server started..."); updateClientList(["Client1", "Client2"]);
// Handle input for graceful exit // term.grabInput(); // term.on("key", (key) => { // if (key === "CTRL_C") { // process.exit(); // } // });
// doc.draw();}*/
function log() { console.log(...arguments);}
function error() { console.error(...arguments);}