-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathserver.js
More file actions
819 lines (741 loc) · 28.3 KB
/
Copy pathserver.js
File metadata and controls
819 lines (741 loc) · 28.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
/* eslint-disable @typescript-eslint/no-require-imports */
const { loadAppEnv } = require('./server/load-app-env');
loadAppEnv(process.cwd());
const dev = process.env.NODE_ENV !== 'production';
const { assertProductionAuthSecret } = require('./app/lib/security/auth-secret');
assertProductionAuthSecret();
if (typeof globalThis.AsyncLocalStorage === 'undefined') {
const { AsyncLocalStorage } = require('node:async_hooks');
globalThis.AsyncLocalStorage = AsyncLocalStorage;
}
// The custom Node server is server-side code, but it imports some Next app
// modules directly for WebSocket runtime handling. Next aliases `server-only`
// during its own bundling; plain Node would otherwise execute the package's
// throwing stub.
const Module = require('module');
const path = require('path');
const originalLoad = Module._load;
const originalResolveFilename = Module._resolveFilename;
function getExportTarget(exportValue) {
if (typeof exportValue === 'string') {
return exportValue;
}
if (!exportValue || typeof exportValue !== 'object') {
return null;
}
return exportValue.import || exportValue.default || exportValue.require || null;
}
function addEsmOnlyPackageAliases(packageName, aliases) {
const packageRoot = path.resolve(process.cwd(), 'node_modules', packageName);
const packageJsonPath = path.join(packageRoot, 'package.json');
const packageJson = require(packageJsonPath);
const exportsMap = packageJson.exports && typeof packageJson.exports === 'object'
? packageJson.exports
: { '.': packageJson.main || './dist/index.js' };
for (const [exportPath, exportValue] of Object.entries(exportsMap)) {
const target = getExportTarget(exportValue);
if (!target) {
continue;
}
const request = exportPath === '.'
? packageName
: `${packageName}/${exportPath.replace(/^\.\//, '')}`;
aliases.set(request, path.resolve(packageRoot, target));
}
}
const esmOnlyPackageAliases = new Map();
addEsmOnlyPackageAliases('@earendil-works/pi-ai', esmOnlyPackageAliases);
addEsmOnlyPackageAliases('@earendil-works/pi-agent-core', esmOnlyPackageAliases);
addEsmOnlyPackageAliases('@earendil-works/pi-telemetry', esmOnlyPackageAliases);
Module._resolveFilename = function resolveWithEsmPackageAliases(request, parent, isMain, options) {
const aliasedPath = esmOnlyPackageAliases.get(request);
if (aliasedPath) {
return aliasedPath;
}
return originalResolveFilename.call(this, request, parent, isMain, options);
};
Module._load = function loadWithServerOnlyMarker(request, parent, isMain) {
if (request === 'server-only') {
return {};
}
return originalLoad.call(this, request, parent, isMain);
};
const http = require('http');
const fs = require('fs');
const next = require('next');
// Terminal service now runs as separate process via Unix Socket
// See server/terminal-service.ts
const {
resolveSkillsDataDir,
} = require('./app/lib/runtime-data-paths');
const port = parseInt(process.env.PORT || '3000', 10);
const hostname = process.env.HOSTNAME || 'localhost';
const useWebpackDev = dev && process.env.CANVAS_DEV_BUNDLER === 'webpack';
if (dev) {
console.log(`[Startup] Next.js dev bundler: ${useWebpackDev ? 'webpack' : 'turbopack'}`);
}
const app = next({
dev,
hostname,
port,
...(useWebpackDev ? { webpack: true, turbopack: false } : {}),
});
const handle = app.getRequestHandler();
let authInstance = null;
function getAuth() {
if (!authInstance) {
authInstance = require('./app/lib/auth').auth;
}
return authInstance;
}
// Helper to get session from Node.js request using better-auth
async function getAuthSession(req) {
try {
const auth = getAuth();
const webHeaders = new Headers();
for (const [key, value] of Object.entries(req.headers)) {
if (typeof value === 'string') {
webHeaders.append(key, value);
} else if (Array.isArray(value)) {
for (const v of value) {
webHeaders.append(key, v);
}
}
}
return await auth.api.getSession({ headers: webHeaders });
} catch (e) {
console.error('[Auth] Error verifying session:', e);
return null;
}
}
const DATA = process.env.DATA || path.resolve(process.cwd(), 'data');
const MEDIA_ROOT = path.join(DATA, 'workspace');
const MEDIA_TYPES = {
pdf: 'application/pdf',
png: 'image/png',
jpg: 'image/jpeg',
jpeg: 'image/jpeg',
gif: 'image/gif',
webp: 'image/webp',
svg: 'image/svg+xml',
mp4: 'video/mp4',
webm: 'video/webm',
ogv: 'video/ogg',
mov: 'video/quicktime',
wav: 'audio/wav',
mp3: 'audio/mpeg',
m4a: 'audio/mp4',
aac: 'audio/aac',
ogg: 'audio/ogg',
opus: 'audio/opus',
flac: 'audio/flac',
};
function setNoIndexHeader(res) {
res.setHeader('X-Robots-Tag', 'noindex, nofollow, noarchive, nosnippet, noimageindex, notranslate');
}
function resolveMediaPath(requestPath) {
const basePath = path.resolve(MEDIA_ROOT);
const normalized = path.resolve(basePath, requestPath);
if (normalized === basePath || normalized.startsWith(`${basePath}${path.sep}`)) {
return normalized;
}
return null;
}
function getContentType(filePath) {
const ext = path.extname(filePath).slice(1).toLowerCase();
return MEDIA_TYPES[ext] || 'application/octet-stream';
}
const RETIRED_SEED_SKILLS = [
{
name: 'browser-tools',
marker: 'author: canvas-studios',
},
];
function cleanupRetiredSeedSkills(skillsDir) {
for (const skill of RETIRED_SEED_SKILLS) {
const skillDir = path.join(skillsDir, skill.name);
const skillMdPath = path.join(skillDir, 'SKILL.md');
if (!fs.existsSync(skillMdPath)) {
continue;
}
const skillMd = fs.readFileSync(skillMdPath, 'utf8');
if (!skillMd.includes(`name: ${skill.name}`) || !skillMd.includes(skill.marker)) {
console.warn(`[Startup] Retired seed skill "${skill.name}" exists but did not match Canvas seed markers; preserving it.`);
continue;
}
fs.rmSync(skillDir, { recursive: true, force: true });
console.log(`[Startup] Removed retired seed skill: ${skill.name}`);
}
}
function syncSeedSkills(repoSkillsDir, skillsDir) {
const copyOptions = {
recursive: true,
force: true,
verbatimSymlinks: true,
};
const skillsRoot = path.resolve(skillsDir);
const maxAttempts = 5;
for (let attempt = 1; attempt <= maxAttempts; attempt += 1) {
try {
fs.cpSync(repoSkillsDir, skillsDir, copyOptions);
return;
} catch (error) {
const targetPath = typeof error?.path === 'string' ? path.resolve(error.path) : null;
const canRetry =
error?.code === 'EEXIST' &&
targetPath &&
targetPath !== skillsRoot &&
targetPath.startsWith(`${skillsRoot}${path.sep}`);
if (!canRetry || attempt === maxAttempts) {
throw error;
}
console.warn(`[Startup] Replacing existing seed skill path before retry: ${targetPath}`);
fs.rmSync(targetPath, { recursive: true, force: true });
}
}
}
function ensureSkillsDirectory() {
const skillsDir = resolveSkillsDataDir(process.cwd());
const repoSkillsDir = path.resolve(process.cwd(), 'seed_skills');
try {
if (!fs.existsSync(skillsDir)) {
fs.mkdirSync(skillsDir, { recursive: true });
console.log(`[Startup] Created skills directory: ${skillsDir}`);
}
if (fs.existsSync(repoSkillsDir)) {
syncSeedSkills(repoSkillsDir, skillsDir);
console.log(`[Startup] Synced seed skills to ${skillsDir}`);
}
cleanupRetiredSeedSkills(skillsDir);
} catch (error) {
console.error('[Startup] Failed to sync skills directory:', error);
}
}
function ensureRuntimeDirectories() {
try {
fs.mkdirSync(path.resolve(MEDIA_ROOT), { recursive: true });
console.log(`[Startup] Ensured workspace directory exists: ${MEDIA_ROOT}`);
} catch (error) {
console.error(`[Startup] Failed to create WORKSPACE_DIR at ${MEDIA_ROOT}:`, error);
throw error;
}
try {
fs.mkdirSync(DATA, { recursive: true });
console.log(`[Startup] Ensured data directory exists: ${DATA}`);
} catch (error) {
console.error(`[Startup] Failed to create data directory at ${DATA}:`, error);
throw error;
}
// Ensure secrets directory exists for Canvas-Integrations.env
const secretsDir = path.join(DATA, 'secrets');
try {
fs.mkdirSync(secretsDir, { recursive: true });
console.log(`[Startup] Ensured secrets directory exists: ${secretsDir}`);
} catch (error) {
console.error(`[Startup] Failed to create secrets directory at ${secretsDir}:`, error);
throw error;
}
}
function serveMedia(req, res) {
setNoIndexHeader(res);
if (req.method !== 'GET' && req.method !== 'HEAD') {
res.statusCode = 405;
res.setHeader('Allow', 'GET, HEAD');
res.end();
return;
}
const url = new URL(req.url, 'http://localhost');
const rawPath = url.pathname.replace(/^\/media\/?/, '');
if (!rawPath) {
res.statusCode = 404;
res.end('Not found');
return;
}
let decodedPath;
try {
decodedPath = decodeURIComponent(rawPath);
} catch {
res.statusCode = 400;
res.end('Bad request');
return;
}
const filePath = resolveMediaPath(decodedPath);
console.log(`[Media Debug] Request: ${decodedPath} -> Resolved: ${filePath} | MEDIA_ROOT: ${MEDIA_ROOT} | DATA: ${DATA}`);
if (!filePath) {
console.log(`[Media Debug] Forbidden: Path resolved to null`);
res.statusCode = 403;
res.end('Forbidden');
return;
}
fs.stat(filePath, (statErr, stats) => {
if (statErr || !stats.isFile()) {
console.log(`[Media Debug] 404: File not found at ${filePath} | Error: ${statErr?.message || 'Not a file'}`);
res.statusCode = 404;
res.end('Not found');
return;
}
const totalSize = stats.size;
const range = req.headers.range;
const contentType = getContentType(filePath);
const ext = path.extname(filePath).slice(1).toLowerCase();
const isImage = ['png', 'jpg', 'jpeg', 'gif', 'webp', 'svg'].includes(ext);
const isMedia = ['mp4', 'webm', 'ogv', 'mov', 'wav', 'mp3', 'm4a', 'aac', 'ogg', 'opus', 'flac'].includes(ext);
const cacheControl = isImage
? 'private, max-age=300'
: isMedia
? 'private, max-age=60'
: 'no-store, max-age=0';
res.setHeader('Content-Type', contentType);
res.setHeader('Content-Disposition', `inline; filename="${path.basename(filePath)}"`);
res.setHeader('Accept-Ranges', 'bytes');
res.setHeader('Cache-Control', cacheControl);
res.setHeader('X-Accel-Buffering', 'no');
if (!range) {
res.statusCode = 200;
res.setHeader('Content-Length', totalSize);
if (req.method === 'HEAD') {
res.end();
return;
}
const stream = fs.createReadStream(filePath, { highWaterMark: 1024 * 1024 });
stream.on('error', () => {
res.destroy();
});
stream.pipe(res);
return;
}
const match = /bytes=(\d*)-(\d*)/i.exec(range);
let start = 0;
let end = totalSize - 1;
if (match) {
if (match[1]) start = Number(match[1]);
if (match[2]) end = Number(match[2]);
if (!match[1] && match[2]) {
const suffixLength = Number(match[2]);
if (Number.isFinite(suffixLength)) {
start = Math.max(totalSize - suffixLength, 0);
end = totalSize - 1;
}
}
}
if (!Number.isFinite(start) || !Number.isFinite(end) || start > end || start >= totalSize) {
res.statusCode = 416;
res.setHeader('Content-Range', `bytes */${totalSize}`);
res.end();
return;
}
end = Math.min(end, totalSize - 1);
res.statusCode = 206;
res.setHeader('Content-Range', `bytes ${start}-${end}/${totalSize}`);
res.setHeader('Content-Length', end - start + 1);
if (req.method === 'HEAD') {
res.end();
return;
}
const stream = fs.createReadStream(filePath, {
start,
end,
highWaterMark: 1024 * 1024,
});
stream.on('error', () => {
res.destroy();
});
stream.pipe(res);
});
}
async function runStartupDatabaseMigrations() {
if (process.env.CANVAS_DATABASE_MIGRATIONS_COMPLETED === 'true') {
console.log('[Startup] Database migrations already completed by the container entrypoint');
return;
}
const { runStartupDatabaseMigrations: migrateDatabase } = require('./app/lib/db/startup-migrations');
await migrateDatabase();
// Runtime modules can be loaded later by Next.js on demand. They must not
// reopen the schema migration path while long-lived SQLite connections are
// already serving requests.
process.env.CANVAS_DATABASE_MIGRATIONS_COMPLETED = 'true';
}
// Ensure all runtime directories and tokens are set up before starting the server
console.log('[Startup] Starting runtime setup...');
try {
console.log('[Startup] Calling ensureRuntimeDirectories()...');
ensureRuntimeDirectories();
console.log('[Startup] ensureRuntimeDirectories() completed');
} catch (error) {
console.error('[Startup] CRITICAL ERROR in ensureRuntimeDirectories():', error.message);
console.error('[Startup] Stack trace:', error.stack);
// Continue anyway - don't block server startup
}
try {
console.log('[Startup] Calling ensureSkillsDirectory()...');
ensureSkillsDirectory();
console.log('[Startup] ensureSkillsDirectory() completed');
} catch (error) {
console.error('[Startup] ERROR in ensureSkillsDirectory():', error.message);
console.error('[Startup] Stack trace:', error.stack);
}
console.log('[Startup] Runtime setup complete');
function runOrphanedAssetsCleanup() {
try {
console.log('[Startup] Running orphaned-assets cleanup...');
const { cleanupOrphanedStudioAssets } = require('./app/lib/cleanup/orphaned-assets');
cleanupOrphanedStudioAssets().then((result) => {
console.log(`[Startup] Orphaned-assets cleanup: ${result.deleted} files deleted, ${result.errors.length} errors`);
}).catch((err) => {
console.warn('[Startup] Orphaned-assets cleanup failed:', err.message);
});
} catch (err) {
console.warn('[Startup] Orphaned-assets cleanup could not be loaded:', err.message);
}
}
function runStudioPresetSeeding() {
try {
console.log('[Startup] Seeding studio preset assets...');
const { ensureDefaultStudioPresetsSeeded } = require('./app/lib/integrations/studio-preset-defaults');
const { ensureStudioAssetsWorkspace } = require('./app/lib/integrations/studio-workspace');
ensureStudioAssetsWorkspace().then(() => {
return ensureDefaultStudioPresetsSeeded();
}).then((result) => {
console.log(`[Startup] Studio preset seeding: ${result.total} presets (${result.inserted} inserted, ${result.updated} updated)`);
}).catch((err) => {
console.warn('[Startup] Studio preset seeding failed:', err.message);
});
} catch (err) {
console.warn('[Startup] Studio preset seeding could not be loaded:', err.message);
}
}
function scheduleExpiredSessionCleanup() {
try {
const { openDb, getDatabaseProvider } = require('./app/lib/db/index');
const isPostgres = getDatabaseProvider() === 'postgres';
const CLEANUP_INTERVAL_MS = 15 * 60 * 1000;
const cleanupQuery = isPostgres
? "DELETE FROM session WHERE expires_at < floor(extract(epoch from now()))::bigint"
: "DELETE FROM session WHERE expires_at < unixepoch()";
async function purgeExpiredSessions() {
try {
const dbConn = await openDb();
const result = dbConn.run(cleanupQuery);
if (result.changes > 0) {
console.log(`[Session Cleanup] Deleted ${result.changes} expired session(s)`);
}
if (!isPostgres) {
dbConn.run("PRAGMA optimize");
}
dbConn.close();
} catch (err) {
console.warn('[Session Cleanup] Failed:', err.message);
}
}
purgeExpiredSessions();
setInterval(purgeExpiredSessions, CLEANUP_INTERVAL_MS).unref?.();
console.log('[Startup] Expired session cleanup scheduled (every 15min)');
} catch (err) {
console.warn('[Startup] Session cleanup could not be initialized:', err.message);
}
}
function scheduleBackgroundMaintenance() {
const timer = setTimeout(() => {
console.log('[Startup] Starting background maintenance...');
runOrphanedAssetsCleanup();
runStudioPresetSeeding();
scheduleExpiredSessionCleanup();
import('./app/lib/license/refresh.ts')
.then((refreshModule) => {
resolveImportedServerModule(
refreshModule,
['initializeCommunityLicenseRefreshRuntime'],
'Community license refresh',
).initializeCommunityLicenseRefreshRuntime();
})
.catch((err) => {
console.warn('[Startup] Community license refresh could not be initialized:', err.message);
});
import('./app/lib/license/team-license-lifecycle.ts')
.then((lifecycleModule) => {
resolveImportedServerModule(
lifecycleModule,
['initializeTeamLicenseLifecycleRuntime'],
'Team license lifecycle',
).initializeTeamLicenseLifecycleRuntime();
})
.catch((err) => {
console.warn('[Startup] Team license lifecycle could not be initialized:', err.message);
});
import('./app/lib/memory/review-worker.ts')
.then((memoryReviewModule) => {
resolveImportedServerModule(
memoryReviewModule,
['initializeMemoryReviewWorkerRuntime'],
'Memory review worker',
).initializeMemoryReviewWorkerRuntime();
})
.catch((err) => {
console.warn('[Startup] Memory review worker could not be initialized:', err.message);
});
}, 1500);
timer.unref?.();
console.log('[Startup] Background maintenance scheduled');
}
function recoverStaleAutomationRuns() {
try {
console.log('[Startup] Running automation stale-run recovery...');
const { markStaleAutomationRunsFailed } = require('./app/lib/automations/store');
markStaleAutomationRunsFailed().then((count) => {
if (count > 0) {
console.log(`[Startup] Recovered ${count} stale automation run(s)`);
} else {
console.log('[Startup] No stale automation runs to recover');
}
}).catch((err) => {
console.warn('[Startup] Failed to recover stale automation runs:', err.message);
});
} catch (err) {
console.warn('[Startup] Could not load automation recovery module:', err.message);
}
}
const server = http.createServer((req, res) => {
const url = new URL(req.url, 'http://localhost');
if (url.pathname.startsWith('/media/')) {
getAuthSession(req)
.then((sessionData) => {
if (!sessionData || !sessionData.user) {
res.statusCode = 401;
res.setHeader('Content-Type', 'application/json');
setNoIndexHeader(res);
res.end(JSON.stringify({ success: false, error: 'Unauthorized' }));
return;
}
serveMedia(req, res);
})
.catch(() => {
res.statusCode = 401;
res.setHeader('Content-Type', 'application/json');
setNoIndexHeader(res);
res.end(JSON.stringify({ success: false, error: 'Unauthorized' }));
});
return;
}
// Terminal kill endpoint is now handled by Next.js API routes
// See app/api/terminal/kill/route.ts
handle(req, res);
});
let shutdownInProgress = false;
let flushCollaborationDocuments = async () => {};
let flushExcalidrawCollaborationDocuments = async () => {};
let closeChatWebSocketServer = async () => {};
function exitCodeForSignal(signal) {
if (signal === 'SIGINT') return 130;
if (signal === 'SIGTERM') return 143;
return 0;
}
async function shutdownServer(signal) {
if (shutdownInProgress) {
return;
}
shutdownInProgress = true;
console.log(`[Startup] Received ${signal}; closing HTTP server...`);
const forceExitTimer = setTimeout(() => {
console.warn(`[Startup] Forced exit after ${signal}; HTTP server did not close in time`);
process.exit(exitCodeForSignal(signal));
}, 10_000);
forceExitTimer.unref?.();
try {
await Promise.all([
flushCollaborationDocuments(),
flushExcalidrawCollaborationDocuments(),
closeChatWebSocketServer(),
]);
} catch (error) {
console.error('[Startup] Error while flushing collaboration documents:', error);
}
server.close((error) => {
if (error) {
console.error(`[Startup] Error while closing HTTP server after ${signal}:`, error);
process.exit(1);
}
console.log(`[Startup] HTTP server closed after ${signal}`);
process.exit(exitCodeForSignal(signal));
});
}
process.on('SIGTERM', () => shutdownServer('SIGTERM'));
process.on('SIGINT', () => shutdownServer('SIGINT'));
process.on('uncaughtException', (error) => {
console.error('[Startup] Uncaught exception:', error);
process.exit(1);
});
process.on('unhandledRejection', (reason) => {
console.error('[Startup] Unhandled rejection:', reason);
process.exit(1);
});
// Register the chat WebSocket handler before Next attaches its own upgrade
// listeners. After app.prepare() we wrap Next's listeners so they never see
// /ws/chat sockets; otherwise Next can still corrupt or close the upgraded
// connection after our ws server has accepted it.
let isCanvasWebSocketRequest = () => false;
function guardNonChatUpgradeListener(listener) {
if (typeof listener !== 'function' || listener.__canvasUpgradeGuarded) {
return listener;
}
const guardedListener = function guardedUpgradeListener(request, socket, head) {
if (isCanvasWebSocketRequest(request.url)) {
return;
}
return listener.call(this, request, socket, head);
};
guardedListener.__canvasUpgradeGuarded = true;
guardedListener.__canvasOriginalListener = listener;
return guardedListener;
}
function installChatUpgradeGuard(targetServer) {
const originalOn = targetServer.on.bind(targetServer);
const originalAddListener = targetServer.addListener.bind(targetServer);
const originalPrependListener = targetServer.prependListener.bind(targetServer);
const originalOnce = targetServer.once.bind(targetServer);
const originalPrependOnceListener = targetServer.prependOnceListener.bind(targetServer);
targetServer.on = function guardedOn(eventName, listener) {
return originalOn(eventName, eventName === 'upgrade' ? guardNonChatUpgradeListener(listener) : listener);
};
targetServer.addListener = function guardedAddListener(eventName, listener) {
return originalAddListener(eventName, eventName === 'upgrade' ? guardNonChatUpgradeListener(listener) : listener);
};
targetServer.prependListener = function guardedPrependListener(eventName, listener) {
return originalPrependListener(eventName, eventName === 'upgrade' ? guardNonChatUpgradeListener(listener) : listener);
};
targetServer.once = function guardedOnce(eventName, listener) {
return originalOnce(eventName, eventName === 'upgrade' ? guardNonChatUpgradeListener(listener) : listener);
};
targetServer.prependOnceListener = function guardedPrependOnceListener(eventName, listener) {
return originalPrependOnceListener(eventName, eventName === 'upgrade' ? guardNonChatUpgradeListener(listener) : listener);
};
}
function resolveImportedServerModule(importedModule, requiredFunctions, label) {
const candidates = [
importedModule,
importedModule?.default,
importedModule?.['module.exports'],
];
const resolved = candidates.find((candidate) => (
candidate
&& requiredFunctions.every((name) => typeof candidate[name] === 'function')
));
if (!resolved) {
throw new Error(`${label} module did not expose expected functions`);
}
return resolved;
}
async function startServer() {
try {
await runStartupDatabaseMigrations();
} catch (error) {
console.error('[Startup] CRITICAL ERROR in database migrations:', error.message);
console.error('[Startup] Stack trace:', error.stack);
throw error;
}
try {
const { assertDirectMcpStartupReady } = require('./app/lib/mcp/server/readiness');
await assertDirectMcpStartupReady();
console.log('[Startup] Direct MCP readiness check completed');
} catch (error) {
console.error('[Startup] CRITICAL ERROR in Direct MCP readiness:', error.message);
throw error;
}
let agentRuntimeWarmupPromise = null;
let managedCatalogWarmupPromise = null;
console.log('[Startup] Initializing WebSocket Server...');
try {
const websocketModule = await import('./server/websocket-server.ts');
const websocketServer = resolveImportedServerModule(
websocketModule,
['createWebSocketServer', 'isChatWebSocketRequest'],
'WebSocket server',
);
const chatWebSocketServer = websocketServer.createWebSocketServer(server);
closeChatWebSocketServer = () => websocketServer.closeWebSocketServer(chatWebSocketServer);
console.log('[Startup] WebSocket Server ready on ws://localhost:' + port + '/ws/chat');
const browserViewModule = await import('./server/browser-view-server.ts');
const browserViewServer = resolveImportedServerModule(
browserViewModule,
['createBrowserViewServer', 'isBrowserViewWebSocketRequest'],
'Browser view server',
);
browserViewServer.createBrowserViewServer(server);
console.log('[Startup] Browser View WebSocket ready on ws://localhost:' + port + '/ws/browser');
// Keep the custom server and Next server externals on the CommonJS Yjs
// entry so the long-lived Node process has exactly one constructor set.
const collaborationModule = require('./server/collaboration-server.ts');
collaborationModule.createCollaborationServer(server);
flushCollaborationDocuments = collaborationModule.flushCollaborationDocuments;
const excalidrawCollaborationModule = require('./server/excalidraw-collaboration/server.ts');
excalidrawCollaborationModule.createExcalidrawCollaborationServer(server);
flushExcalidrawCollaborationDocuments = excalidrawCollaborationModule.flushExcalidrawCollaborationDocuments;
isCanvasWebSocketRequest = (requestUrl) => (
websocketServer.isChatWebSocketRequest(requestUrl)
|| browserViewServer.isBrowserViewWebSocketRequest(requestUrl)
|| collaborationModule.isCollaborationWebSocketRequest(requestUrl)
|| excalidrawCollaborationModule.isExcalidrawCollaborationWebSocketRequest(requestUrl)
);
console.log('[Startup] Collaboration WebSocket ready on ws://localhost:' + port + '/ws/collaboration');
console.log('[Startup] Excalidraw Collaboration WebSocket ready on ws://localhost:' + port + '/ws/collaboration/excalidraw');
const agentRuntimeLoaderModule = await import('./server/agent-runtime-loader.ts');
const agentRuntimeLoader = resolveImportedServerModule(
agentRuntimeLoaderModule,
['preloadAgentRuntimeModules'],
'Agent runtime loader',
);
const { preloadAgentRuntimeModules } = agentRuntimeLoader;
agentRuntimeWarmupPromise = preloadAgentRuntimeModules().then((result) => {
console.log('[Startup] Agent runtime modules preloaded', result);
return result;
});
const managedCatalogModule = await import('./app/lib/managed/control-plane-models.ts');
const managedCatalog = resolveImportedServerModule(
managedCatalogModule,
['primeCanvasControlPlaneCatalog'],
'Managed Control Plane catalog',
);
const { primeCanvasControlPlaneCatalog } = managedCatalog;
managedCatalogWarmupPromise = primeCanvasControlPlaneCatalog().then((catalog) => {
console.log('[Startup] Managed model catalog warmup finished', {
status: catalog.status,
errorCode: catalog.errorCode,
modelCount: catalog.models.length,
});
return catalog;
});
} catch (error) {
console.error('[Startup] ERROR initializing WebSocket Server:', error.message);
console.error('[Startup] Stack trace:', error.stack);
isCanvasWebSocketRequest = () => false;
if (process.env.CANVAS_ALLOW_HTTP_WITHOUT_CHAT_WS !== 'true') {
throw error;
}
}
installChatUpgradeGuard(server);
// Channel Manager start (Telegram polling etc.)
try {
const { getChannelManager } = require('./app/lib/channels/manager.ts');
const manager = getChannelManager();
await manager.start();
console.log('[Startup] Channel Manager started');
} catch (error) {
console.error('[Startup] Channel Manager failed:', error.message);
}
console.log('[Startup] Preparing Next.js app...');
await Promise.all([
app.prepare(),
agentRuntimeWarmupPromise,
managedCatalogWarmupPromise,
]);
console.log('[Startup] Next.js app prepared');
server.listen(port, hostname, (err) => {
if (err) throw err;
console.log(`> Ready on http://localhost:${port}`);
scheduleBackgroundMaintenance();
recoverStaleAutomationRuns();
});
}
startServer().catch((error) => {
console.error('Failed to start server', error);
process.exit(1);
});