2026-06-29 17:40:42 -04:00
// Investor Flow — backend entry (DESIGN.md §2.1: one process exposes /api/trpc/*).
// Adapted from Bun.serve to node:http + @trpc/server fetch adapter (runtime glue only).
import { createServer , type IncomingMessage , type ServerResponse } from 'node:http' ;
import { fetchRequestHandler } from '@trpc/server/adapters/fetch' ;
import { db } from './db/client.ts' ;
import { createCacheRepository , type SourceKind } from './cache/CacheRepository.ts' ;
import { YFinanceAdapter } from './adapters/YFinanceAdapter.ts' ;
2026-07-25 12:59:43 -04:00
import { NasdaqAdapter } from './adapters/NasdaqAdapter.ts' ;
2026-07-25 13:03:49 -04:00
import { FinraBulkAdapter } from './adapters/FinraBulkAdapter.ts' ;
2026-07-12 12:55:59 -04:00
import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts' ;
2026-07-12 20:19:12 -04:00
import { SecLintAdapter } from './adapters/SecLintAdapter.ts' ;
2026-07-23 18:02:24 -04:00
import { XCookieAdapter } from './adapters/XCookieAdapter.ts' ;
2026-07-12 12:55:59 -04:00
import type { SourceFetch } from './adapters/SourceAdapter.ts' ;
2026-07-23 18:02:24 -04:00
import cryptoMod from './lib/crypto.ts' ;
2026-06-29 17:40:42 -04:00
import { AdapterQueue } from './queue/AdapterQueue.ts' ;
import { makeCreateContext } from './trpc/context.ts' ;
import { appRouter } from './trpc/router.ts' ;
const PORT = Number ( process . env . PORT ?? 3001 );
const database = db ();
2026-07-12 20:19:12 -04:00
const adapters = new Map < SourceKind , SourceFetch >([
[ 'yfinance' as const , new YFinanceAdapter () as unknown as SourceFetch ],
2026-07-25 12:59:43 -04:00
[ 'nasdaq' as const , new NasdaqAdapter () as unknown as SourceFetch ],
2026-07-25 13:03:49 -04:00
[ 'finra-bulk' as const , new FinraBulkAdapter ( database ) as unknown as SourceFetch ],
2026-07-12 20:19:12 -04:00
[ 'sec-fetch' as const , new SecFetchAdapter ( database ) as unknown as SourceFetch ],
[ 'sec-lint-holders' as const , new SecLintAdapter (() => database , 'sec-lint-holders' ) as unknown as SourceFetch ],
[ 'sec-lint-insiders' as const , new SecLintAdapter (() => database , 'sec-lint-insiders' ) as unknown as SourceFetch ],
2026-07-12 12:55:59 -04:00
]);
2026-07-23 18:02:24 -04:00
// Load X credentials at startup and register XCookieAdapter if available.
let xAdapter : XCookieAdapter | null = null ;
( function initXAdapter() {
const row = database . prepare ( 'SELECT ct0_enc, auth_token_enc FROM x_credentials WHERE id=?' ). get ( 'singleton' ) as { ct0_enc? : string ; auth_token_enc? : string } | undefined ;
if ( ! row ? . ct0_enc || ! row ? . auth_token_enc ) return ;
let creds : { ct0 : string ; auth_token : string };
try { creds = { ct0 : cryptoMod.decrypt ( row . ct0_enc ), auth_token : cryptoMod.decrypt ( row . auth_token_enc ) }; }
catch { return ; }
const updateHealth = ( status : string , err? : string | null ) => {
database . prepare (
`INSERT INTO x_credentials (id, healthy, last_error, updated_at) VALUES ('singleton', ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET healthy=excluded.healthy, last_error=excluded.last_error, updated_at=excluded.updated_at`
). run ( status === 'healthy' ? 1 : 0 , err ?? null , new Date (). toISOString ());
};
xAdapter = new XCookieAdapter ( creds , ( health ) => updateHealth ( health . sourceStatus , health . lastError ), database );
adapters . set ( 'x' as const , xAdapter as unknown as SourceFetch );
console . log ( '[investor-flow] X adapter registered (credentials configured)' );
})();
2026-06-29 17:40:42 -04:00
const queue = new AdapterQueue ({ db : database , adapters });
const cache = createCacheRepository ({ db : database , scheduler : queue });
queue . cache = cache ; // break the cache<->scheduler cycle
2026-07-12 12:55:59 -04:00
// Seed default schedules (noop if already seeded)
queue . seedDefaultSchedules ();
// Startup recovery: any job left 'in_flight' was interrupted by a restart/crash.
// Reset to 'pending' so the drain loop reprocesses it.
const recovered = database . prepare ( "UPDATE adapter_queue SET status='pending', error=NULL, retry_count=0 WHERE status='in_flight'" ). run ();
if ( Number ( recovered . changes ) > 0 ) console . log ( `[investor-flow] recovered ${ recovered . changes } interrupted in_flight jobs` );
2026-07-23 18:02:24 -04:00
const createContext = makeCreateContext ({ db : database , cache , queue , xAdapter });
2026-06-29 17:40:42 -04:00
// Background drain: stale-while-revalidate refreshes are queued by CacheRepository.get;
// this loop drains them (fetch via adapter -> write to cache), deduped + backed off.
const DRAIN_MS = Number ( process . env . IFLOW_DRAIN_MS ?? 2000 );
const drainTimer = setInterval (() => { queue . drain (). catch (( e ) => console . error ( '[drain error]' , e )); }, DRAIN_MS );
drainTimer . unref ();
2026-07-12 12:55:59 -04:00
// Auto-scheduler: every 30s, enqueue refreshes for due schedules
const SCHEDULE_MS = 30 _000 ;
const scheduleTimer = setInterval (() => { queue . enqueueDueSchedules (). catch (( e ) => console . error ( '[schedule error]' , e )); }, SCHEDULE_MS );
scheduleTimer . unref ();
2026-07-23 20:51:47 -04:00
// Per-fetch alert check: run alongside the schedule cycle for critical alert types.
const perFetchAlertTimer = setInterval ( async () => {
try {
const alerts = await runProducers ( database , 'per-fetch' );
if ( alerts . length > 0 ) {
console . log ( `[alert] ${ alerts . length } per-fetch alert(s) created` );
const { sendAlertEmail } = await import ( './services/emailAlertService.ts' );
alerts . forEach (( a ) => sendAlertEmail ( database , a ). catch (() => {}));
}
} catch ( e ) {
console . error ( '[alert] per-fetch check failed:' , e );
}
}, SCHEDULE_MS );
perFetchAlertTimer . unref ();
// Register alert producers.
import { registerProducer , runProducers } from './alerts/producers/index.ts' ;
import { informedBuyProducer , informedSellProducer } from './alerts/producers/insiderProducer.ts' ;
import { new13daProducer } from './alerts/producers/new13daProducer.ts' ;
registerProducer ( informedBuyProducer );
registerProducer ( informedSellProducer );
registerProducer ( new13daProducer );
// Batched alert check: run every 10 minutes for non-critical producers.
const ALERT_BATCH_MS = 10 * 60 * 1000 ;
const alertBatchTimer = setInterval ( async () => {
try {
const alerts = await runProducers ( database , 'batched' );
if ( alerts . length > 0 ) console . log ( `[alert] ${ alerts . length } batch alert(s) created` );
// Send email for each new alert.
const { sendAlertEmail } = await import ( './services/emailAlertService.ts' );
for ( const alert of alerts ) {
sendAlertEmail ( database , alert ). catch (( e ) => console . error ( '[alert:email] send error:' , e ));
}
} catch ( e ) {
console . error ( '[alert] batch check failed:' , e );
}
}, ALERT_BATCH_MS );
alertBatchTimer . unref ();
2026-06-29 17:40:42 -04:00
function readBody ( req : IncomingMessage ) : Promise < string > {
return new Promise (( resolve , reject ) => {
let data = '' ;
req . on ( 'data' , ( c : Buffer ) => { data += c . toString (); });
req . on ( 'end' , () => resolve ( data ));
req . on ( 'error' , reject );
});
}
const server = createServer ( async ( req , res ) => {
const url = new URL ( req . url ?? '/' , `http://localhost: ${ PORT } ` );
if ( url . pathname === '/health' ) {
res . writeHead ( 200 , { 'content-type' : 'application/json' });
res . end ( JSON . stringify ({ ok : true , queue : queue.health () }));
return ;
}
if ( url . pathname . startsWith ( '/api/trpc' )) {
const headers = new Headers ();
for ( const [ k , v ] of Object . entries ( req . headers )) if ( v != null ) headers . set ( k , Array . isArray ( v ) ? v . join ( ', ' ) : String ( v ));
const method = req . method ?? 'GET' ;
const body = method === 'GET' || method === 'HEAD' ? undefined : await readBody ( req );
2026-07-12 12:55:59 -04:00
const requestUrl = `http://localhost: ${ PORT }${ url . pathname }${ url . search } ` ;
const init : RequestInit = { method , headers };
if ( body ) {
init . body = new Blob ([ body ], { type : 'application/json' });
}
const request = new Request ( requestUrl , init );
2026-06-29 17:40:42 -04:00
try {
const response = await fetchRequestHandler ({ router : appRouter , createContext , endpoint : '/api/trpc' , req : request });
const buf = Buffer . from ( await response . arrayBuffer ());
res . writeHead ( response . status , Object . fromEntries ( response . headers . entries ()));
res . end ( buf );
} catch ( e ) {
console . error ( '[trpc error]' , e );
res . writeHead ( 500 , { 'content-type' : 'application/json' });
res . end ( JSON . stringify ({ error : 'tRPC handler error' }));
}
return ;
}
res . writeHead ( 404 , { 'content-type' : 'application/json' });
res . end ( JSON . stringify ({ error : 'not found' }));
});
2026-07-12 12:55:59 -04:00
const HOST = process . env . HOST ?? '0.0.0.0' ;
server . listen ( PORT , HOST , () => {
console . log ( `[investor-flow] backend on http:// ${ HOST } : ${ PORT } (tRPC /api/trpc, health /health, drain every ${ DRAIN_MS } ms)` );
2026-06-29 17:40:42 -04:00
});