2026-03-14 00:35:31 -07:00
// Telegram Alerter v2 — Multi-tier alerts, semantic dedup, two-way bot commands
// USP feature: Crucix becomes a conversational intelligence agent via Telegram
import { createHash } from 'crypto' ;
2026-03-12 23:45:46 -07:00
const TELEGRAM _API = 'https://api.telegram.org' ;
2026-03-17 14:04:32 +01:00
/** Telegram Bot API limit for sendMessage text (bytes/characters). */
const TELEGRAM _MAX _TEXT = 4096 ;
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
// ─── Alert Tiers ────────────────────────────────────────────────────────────
// FLASH: Immediate action required — market-moving, time-critical (e.g. war escalation, flash crash)
// PRIORITY: Important signal cluster — act within hours (e.g. rate surprise, major OSINT shift)
// ROUTINE: Noteworthy change — FYI, no urgency (e.g. trend continuation, moderate delta)
const TIER _CONFIG = {
FLASH : { emoji : '🔴' , label : 'FLASH' , cooldownMs : 5 * 60 * 1000 , maxPerHour : 6 } ,
PRIORITY : { emoji : '🟡' , label : 'PRIORITY' , cooldownMs : 30 * 60 * 1000 , maxPerHour : 4 } ,
ROUTINE : { emoji : '🔵' , label : 'ROUTINE' , cooldownMs : 60 * 60 * 1000 , maxPerHour : 2 } ,
} ;
// ─── Bot Commands ───────────────────────────────────────────────────────────
const COMMANDS = {
'/status' : 'Get current system health, last sweep time, source status' ,
'/sweep' : 'Trigger a manual sweep cycle' ,
'/brief' : 'Get a compact text summary of the latest intelligence' ,
'/portfolio' : 'Show current positions and P&L (if Alpaca connected)' ,
'/alerts' : 'Show recent alert history' ,
'/mute' : 'Mute alerts for 1h (or /mute 2h, /mute 4h)' ,
'/unmute' : 'Resume alerts' ,
'/help' : 'Show available commands' ,
} ;
2026-03-12 23:45:46 -07:00
export class TelegramAlerter {
constructor ( { botToken , chatId } ) {
this . botToken = botToken ;
this . chatId = chatId ;
2026-03-14 00:35:31 -07:00
this . _alertHistory = [ ] ; // Recent alerts for rate limiting
this . _contentHashes = { } ; // Semantic dedup: hash → timestamp
this . _muteUntil = null ; // Mute timestamp
this . _lastUpdateId = 0 ; // For polling bot commands
this . _commandHandlers = { } ; // Registered command callbacks
this . _pollingInterval = null ;
2026-03-17 13:43:55 +01:00
this . _botUsername = null ;
2026-03-12 23:45:46 -07:00
}
get isConfigured ( ) {
return ! ! ( this . botToken && this . chatId ) ;
}
2026-03-14 00:35:31 -07:00
// ─── Core Messaging ─────────────────────────────────────────────────────
2026-03-12 23:45:46 -07:00
/ * *
2026-03-17 14:04:32 +01:00
* Send a message via Telegram Bot API . Splits at TELEGRAM _MAX _TEXT so long messages
* ( e . g . / brief ) are sent in multiple messages instead of being truncated or failing .
2026-03-12 23:45:46 -07:00
* @ param { string } message - markdown - formatted message
2026-03-17 14:04:32 +01:00
* @ param { object } opts - optional : { parseMode , disablePreview , replyToMessageId , chatId }
2026-03-14 00:35:31 -07:00
* @ returns { Promise < { ok : boolean , messageId ? : number } > }
2026-03-12 23:45:46 -07:00
* /
2026-03-14 00:35:31 -07:00
async sendMessage ( message , opts = { } ) {
if ( ! this . isConfigured ) return { ok : false } ;
2026-03-17 14:04:32 +01:00
const chatId = opts . chatId ? ? this . chatId ;
const parseMode = opts . parseMode || 'Markdown' ;
const chunks = this . _chunkText ( message , TELEGRAM _MAX _TEXT ) ;
2026-03-12 23:45:46 -07:00
try {
2026-03-17 14:04:32 +01:00
let lastResult = { ok : false , messageId : undefined } ;
for ( let i = 0 ; i < chunks . length ; i ++ ) {
const res = await fetch ( ` ${ TELEGRAM _API } /bot ${ this . botToken } /sendMessage ` , {
method : 'POST' ,
headers : { 'Content-Type' : 'application/json' } ,
body : JSON . stringify ( {
chat _id : chatId ,
text : chunks [ i ] ,
parse _mode : parseMode ,
disable _web _page _preview : opts . disablePreview !== false ,
... ( opts . replyToMessageId && i === 0 ? { reply _to _message _id : opts . replyToMessageId } : { } ) ,
} ) ,
signal : AbortSignal . timeout ( 15000 ) ,
} ) ;
2026-03-12 23:45:46 -07:00
2026-03-17 14:04:32 +01:00
if ( ! res . ok ) {
const err = await res . text ( ) . catch ( ( ) => '' ) ;
console . error ( ` [Telegram] Send failed ( ${ res . status } ): ${ err . substring ( 0 , 200 ) } ` ) ;
return lastResult ;
}
2026-03-12 23:45:46 -07:00
2026-03-17 14:04:32 +01:00
const data = await res . json ( ) ;
lastResult = { ok : true , messageId : data . result ? . message _id } ;
}
return lastResult ;
2026-03-12 23:45:46 -07:00
} catch ( err ) {
console . error ( '[Telegram] Send error:' , err . message ) ;
2026-03-14 00:35:31 -07:00
return { ok : false } ;
2026-03-12 23:45:46 -07:00
}
}
2026-03-17 14:04:32 +01:00
/ * *
* Split text into chunks of at most maxLen . Prefer breaking at newlines to avoid
* splitting mid - Markdown .
* /
_chunkText ( text , maxLen = TELEGRAM _MAX _TEXT ) {
if ( ! text || text . length <= maxLen ) return text ? [ text ] : [ ] ;
const chunks = [ ] ;
let start = 0 ;
while ( start < text . length ) {
let end = Math . min ( start + maxLen , text . length ) ;
if ( end < text . length ) {
const lastNewline = text . lastIndexOf ( '\n' , end - 1 ) ;
if ( lastNewline > start ) end = lastNewline + 1 ;
}
chunks . push ( text . slice ( start , end ) ) ;
start = end ;
}
return chunks ;
}
2026-03-14 00:35:31 -07:00
// Backward-compatible alias
async sendAlert ( message ) {
const result = await this . sendMessage ( message ) ;
return result . ok ;
}
// ─── Multi-Tier Alert Evaluation ────────────────────────────────────────
2026-03-12 23:45:46 -07:00
/ * *
2026-03-14 00:35:31 -07:00
* Evaluate delta signals with LLM and send tiered alert if warranted .
* Uses semantic dedup , rate limiting , and a much richer evaluation prompt .
2026-03-12 23:45:46 -07:00
* /
async evaluateAndAlert ( llmProvider , delta , memory ) {
2026-03-14 00:35:31 -07:00
if ( ! this . isConfigured ) return false ;
if ( ! delta ? . summary ? . totalChanges ) return false ;
if ( this . _isMuted ( ) ) {
console . log ( '[Telegram] Alerts muted until' , new Date ( this . _muteUntil ) . toLocaleTimeString ( ) ) ;
return false ;
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
// 1. Gather new signals — filter already-alerted AND semantically duplicate
const allSignals = [
2026-03-12 23:45:46 -07:00
... ( delta . signals ? . new || [ ] ) ,
... ( delta . signals ? . escalated || [ ] ) ,
2026-03-14 00:35:31 -07:00
] ;
const newSignals = allSignals . filter ( s => {
const key = this . _signalKey ( s ) ;
// Check decay-based suppression (if memory supports it)
if ( typeof memory . isSignalSuppressed === 'function' ) {
if ( memory . isSignalSuppressed ( key ) ) return false ;
} else {
// Legacy: check flat alerted map
const alerted = memory . getAlertedSignals ( ) ;
if ( alerted [ key ] ) return false ;
}
// Check semantic/content hash dedup
if ( this . _isSemanticDuplicate ( s ) ) return false ;
return true ;
2026-03-12 23:45:46 -07:00
} ) ;
if ( newSignals . length === 0 ) return false ;
2026-03-14 00:35:31 -07:00
// 2. Try LLM evaluation first, fall back to rule-based if unavailable
let evaluation = null ;
if ( llmProvider ? . isConfigured ) {
try {
const systemPrompt = this . _buildEvaluationPrompt ( ) ;
const userMessage = this . _buildSignalContext ( newSignals , delta ) ;
const result = await llmProvider . complete ( systemPrompt , userMessage , {
maxTokens : 800 ,
timeout : 30000 ,
} ) ;
evaluation = parseJSON ( result . text ) ;
} catch ( err ) {
console . warn ( '[Telegram] LLM evaluation failed, falling back to rules:' , err . message ) ;
// Fall through to rule-based evaluation
}
}
// Rule-based fallback: fires when LLM is unavailable or returns garbage
if ( ! evaluation || typeof evaluation . shouldAlert !== 'boolean' ) {
evaluation = this . _ruleBasedEvaluation ( newSignals , delta ) ;
if ( evaluation ) evaluation . _source = 'rules' ;
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
if ( ! evaluation ? . shouldAlert ) {
console . log ( '[Telegram] No alert —' , evaluation ? . reason || 'no qualifying signals' ) ;
return false ;
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
// 3. Validate tier and check rate limits
const tier = TIER _CONFIG [ evaluation . tier ] ? evaluation . tier : 'ROUTINE' ;
if ( ! this . _checkRateLimit ( tier ) ) {
console . log ( ` [Telegram] Rate limited for tier ${ tier } ` ) ;
return false ;
}
// 4. Format and send tiered alert
const message = this . _formatTieredAlert ( evaluation , delta , tier ) ;
const sent = await this . sendAlert ( message ) ;
if ( sent ) {
// Mark signals as alerted with content hashing
for ( const s of newSignals ) {
const key = this . _signalKey ( s ) ;
memory . markAsAlerted ( key , new Date ( ) . toISOString ( ) ) ;
this . _recordContentHash ( s ) ;
}
this . _recordAlert ( tier ) ;
console . log ( ` [Telegram] ${ tier } alert sent ( ${ evaluation . _source || 'llm' } ): ${ evaluation . headline } ` ) ;
}
return sent ;
}
// ─── Rule-Based Alert Fallback ────────────────────────────────────────
/ * *
* Deterministic alert evaluation when LLM is unavailable .
* Uses signal counts , severity , and cross - domain correlation .
* /
_ruleBasedEvaluation ( signals , delta ) {
const criticals = signals . filter ( s => s . severity === 'critical' ) ;
const highs = signals . filter ( s => s . severity === 'high' ) ;
const nukeSignal = signals . find ( s => s . key === 'nuke_anomaly' ) ;
const osintNew = signals . filter ( s => s . key ? . startsWith ( 'tg_urgent' ) ) ;
const marketSignals = signals . filter ( s => [ 'vix' , 'hy_spread' , 'wti' , 'brent' , '10y2y' ] . includes ( s . key ) ) ;
const conflictSignals = signals . filter ( s => [ 'conflict_events' , 'conflict_fatalities' , 'thermal_total' ] . includes ( s . key ) ) ;
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
// FLASH: nuclear anomaly, or ≥3 critical signals across domains
if ( nukeSignal ) {
return {
shouldAlert : true , tier : 'FLASH' , confidence : 'HIGH' ,
headline : 'Nuclear Anomaly Detected' ,
reason : 'Safecast radiation monitors have flagged an anomaly. This requires immediate attention.' ,
actionable : 'Check dashboard for affected sites. Monitor confirmation from secondary sources.' ,
signals : [ 'nuke_anomaly' ] ,
crossCorrelation : 'radiation monitors' ,
} ;
}
// FLASH: ≥2 critical signals AND they span multiple domains
const hasCriticalMarket = criticals . some ( s => marketSignals . includes ( s ) ) ;
const hasCriticalConflict = criticals . some ( s => conflictSignals . includes ( s ) || osintNew . includes ( s ) ) ;
if ( criticals . length >= 2 && hasCriticalMarket && hasCriticalConflict ) {
return {
shouldAlert : true , tier : 'FLASH' , confidence : 'HIGH' ,
headline : ` ${ criticals . length } Critical Cross-Domain Signals ` ,
reason : ` ${ criticals . length } critical signals detected across market and conflict domains. Multi-domain correlation suggests systemic event. ` ,
actionable : 'Review dashboard immediately. Assess portfolio exposure.' ,
signals : criticals . map ( s => s . label || s . key ) . slice ( 0 , 5 ) ,
crossCorrelation : 'market + conflict' ,
} ;
}
// PRIORITY: ≥2 high/critical signals in same direction
const escalatedHighs = [ ... criticals , ... highs ] . filter ( s => s . direction === 'up' ) ;
if ( escalatedHighs . length >= 2 ) {
return {
shouldAlert : true , tier : 'PRIORITY' , confidence : 'MEDIUM' ,
headline : ` ${ escalatedHighs . length } Escalating Signals ` ,
reason : ` Multiple indicators escalating simultaneously: ${ escalatedHighs . map ( s => s . label || s . key ) . slice ( 0 , 3 ) . join ( ', ' ) } . ` ,
actionable : 'Monitor for continuation. Check if trend persists in next sweep.' ,
signals : escalatedHighs . map ( s => s . label || s . key ) . slice ( 0 , 5 ) ,
crossCorrelation : 'multi-indicator' ,
} ;
}
// PRIORITY: ≥5 new OSINT posts (surge in conflict reporting)
if ( osintNew . length >= 5 ) {
return {
shouldAlert : true , tier : 'PRIORITY' , confidence : 'MEDIUM' ,
headline : ` OSINT Surge: ${ osintNew . length } New Urgent Posts ` ,
reason : ` ${ osintNew . length } new urgent OSINT signals detected. Elevated conflict reporting tempo. ` ,
actionable : 'Review OSINT stream for pattern. Cross-check with satellite and ACLED data.' ,
2026-03-21 12:59:30 -04:00
signals : osintNew . map ( s => s . text || s . label || s . key ) . slice ( 0 , 5 ) ,
2026-03-14 00:35:31 -07:00
crossCorrelation : 'telegram OSINT' ,
} ;
}
// ROUTINE: any critical signal OR ≥3 high signals
if ( criticals . length >= 1 || highs . length >= 3 ) {
const topSignal = criticals [ 0 ] || highs [ 0 ] ;
return {
shouldAlert : true , tier : 'ROUTINE' , confidence : 'LOW' ,
headline : topSignal . label || topSignal . reason || 'Signal Change Detected' ,
reason : ` ${ criticals . length } critical, ${ highs . length } high-severity signals. ${ delta . summary . direction } bias. ` ,
actionable : 'Monitor' ,
signals : [ ... criticals , ... highs ] . map ( s => s . label || s . key ) . slice ( 0 , 4 ) ,
crossCorrelation : 'single-domain' ,
} ;
}
// No alert
return {
shouldAlert : false ,
reason : ` ${ signals . length } signals, but none meet alert threshold ( ${ criticals . length } critical, ${ highs . length } high). ` ,
} ;
}
// ─── Two-Way Bot Commands ───────────────────────────────────────────────
/ * *
* Register command handlers that the bot can respond to .
* @ param { string } command - e . g . '/status'
* @ param { Function } handler - async ( args , messageId ) => responseText
* /
onCommand ( command , handler ) {
this . _commandHandlers [ command . toLowerCase ( ) ] = handler ;
}
/ * *
* Start polling for incoming messages / commands .
* Call this once during server startup .
* @ param { number } intervalMs - polling interval ( default 5000 ms )
* /
startPolling ( intervalMs = 5000 ) {
if ( ! this . isConfigured ) return ;
if ( this . _pollingInterval ) return ; // Already polling
console . log ( '[Telegram] Bot command polling started' ) ;
2026-03-17 13:43:55 +01:00
this . _initializeBotCommands ( ) . catch ( ( err ) => {
console . error ( '[Telegram] Command initialization failed:' , err . message ) ;
} ) ;
2026-03-14 00:35:31 -07:00
this . _pollingInterval = setInterval ( ( ) => this . _pollUpdates ( ) , intervalMs ) ;
// Initial poll
this . _pollUpdates ( ) ;
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
/ * *
* Stop polling for incoming messages .
* /
stopPolling ( ) {
if ( this . _pollingInterval ) {
clearInterval ( this . _pollingInterval ) ;
this . _pollingInterval = null ;
console . log ( '[Telegram] Bot command polling stopped' ) ;
}
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
async _pollUpdates ( ) {
2026-03-12 23:45:46 -07:00
try {
2026-03-14 00:35:31 -07:00
const params = new URLSearchParams ( {
offset : String ( this . _lastUpdateId + 1 ) ,
timeout : '0' ,
limit : '10' ,
allowed _updates : JSON . stringify ( [ 'message' ] ) ,
} ) ;
const res = await fetch ( ` ${ TELEGRAM _API } /bot ${ this . botToken } /getUpdates? ${ params } ` , {
signal : AbortSignal . timeout ( 10000 ) ,
} ) ;
if ( ! res . ok ) return ;
const data = await res . json ( ) ;
if ( ! data . ok || ! Array . isArray ( data . result ) ) return ;
for ( const update of data . result ) {
this . _lastUpdateId = Math . max ( this . _lastUpdateId , update . update _id ) ;
const msg = update . message ;
if ( ! msg ? . text ) continue ;
const chatId = String ( msg . chat ? . id ) ;
2026-03-18 08:20:32 +01:00
// Restrict command execution to the configured chat/group only.
if ( chatId !== String ( this . chatId ) ) continue ;
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
await this . _handleMessage ( msg ) ;
2026-03-12 23:45:46 -07:00
}
2026-03-14 00:35:31 -07:00
} catch ( err ) {
// Silent — polling failures are non-fatal
if ( ! err . message ? . includes ( 'aborted' ) ) {
console . error ( '[Telegram] Poll error:' , err . message ) ;
}
}
}
async _handleMessage ( msg ) {
const text = msg . text . trim ( ) ;
const parts = text . split ( /\s+/ ) ;
2026-03-17 13:43:55 +01:00
const rawCommand = parts [ 0 ] . toLowerCase ( ) ;
const command = this . _normalizeCommand ( rawCommand ) ;
if ( ! command ) return ;
2026-03-14 00:35:31 -07:00
const args = parts . slice ( 1 ) . join ( ' ' ) ;
2026-03-17 13:43:55 +01:00
const replyChatId = msg . chat ? . id ;
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
// Built-in commands
if ( command === '/help' ) {
const helpText = Object . entries ( COMMANDS )
. map ( ( [ cmd , desc ] ) => ` ${ cmd } — ${ desc } ` )
. join ( '\n' ) ;
await this . sendMessage (
` 🤖 *CRUCIX BOT COMMANDS* \n \n ${ helpText } \n \n _Tip: Commands are case-insensitive_ ` ,
2026-03-17 13:43:55 +01:00
{ chatId : replyChatId , replyToMessageId : msg . message _id }
2026-03-14 00:35:31 -07:00
) ;
return ;
}
if ( command === '/mute' ) {
const hours = parseFloat ( args ) || 1 ;
this . _muteUntil = Date . now ( ) + hours * 60 * 60 * 1000 ;
await this . sendMessage (
` 🔇 Alerts muted for ${ hours } h — until ${ new Date ( this . _muteUntil ) . toLocaleTimeString ( ) } UTC \n Use /unmute to resume. ` ,
2026-03-17 13:43:55 +01:00
{ chatId : replyChatId , replyToMessageId : msg . message _id }
2026-03-14 00:35:31 -07:00
) ;
return ;
}
if ( command === '/unmute' ) {
this . _muteUntil = null ;
await this . sendMessage (
` 🔔 Alerts resumed. You'll receive the next signal evaluation. ` ,
2026-03-17 13:43:55 +01:00
{ chatId : replyChatId , replyToMessageId : msg . message _id }
2026-03-14 00:35:31 -07:00
) ;
return ;
}
2026-03-12 23:45:46 -07:00
2026-03-14 00:35:31 -07:00
if ( command === '/alerts' ) {
const recent = this . _alertHistory . slice ( - 10 ) ;
if ( recent . length === 0 ) {
2026-03-17 13:43:55 +01:00
await this . sendMessage ( 'No recent alerts.' , { chatId : replyChatId , replyToMessageId : msg . message _id } ) ;
2026-03-14 00:35:31 -07:00
return ;
}
const lines = recent . map ( a =>
` ${ TIER _CONFIG [ a . tier ] ? . emoji || '⚪' } ${ a . tier } — ${ new Date ( a . timestamp ) . toLocaleTimeString ( ) } `
) ;
await this . sendMessage (
` 📋 *Recent Alerts (last ${ recent . length } )* \n \n ${ lines . join ( '\n' ) } ` ,
2026-03-17 13:43:55 +01:00
{ chatId : replyChatId , replyToMessageId : msg . message _id }
2026-03-14 00:35:31 -07:00
) ;
return ;
}
// Delegate to registered handlers
const handler = this . _commandHandlers [ command ] ;
if ( handler ) {
try {
const response = await handler ( args , msg . message _id ) ;
if ( response ) {
2026-03-17 13:43:55 +01:00
await this . sendMessage ( response , { chatId : replyChatId , replyToMessageId : msg . message _id } ) ;
2026-03-12 23:45:46 -07:00
}
2026-03-14 00:35:31 -07:00
} catch ( err ) {
console . error ( ` [Telegram] Command ${ command } error: ` , err . message ) ;
await this . sendMessage (
` ❌ Command failed: ${ err . message } ` ,
2026-03-17 13:43:55 +01:00
{ chatId : replyChatId , replyToMessageId : msg . message _id }
2026-03-14 00:35:31 -07:00
) ;
2026-03-12 23:45:46 -07:00
}
2026-03-14 00:35:31 -07:00
}
// Unknown commands are silently ignored to avoid spamming
}
2026-03-12 23:45:46 -07:00
2026-03-17 13:43:55 +01:00
async _initializeBotCommands ( ) {
await this . _loadBotIdentity ( ) ;
const botCommands = Object . entries ( COMMANDS ) . map ( ( [ command , description ] ) => ( {
command : command . replace ( '/' , '' ) ,
description : description . substring ( 0 , 256 ) ,
} ) ) ;
2026-03-18 08:20:32 +01:00
// Register commands only for the configured chat to avoid global discovery.
await this . _setMyCommands ( botCommands , this . _buildConfiguredChatScope ( ) ) ;
2026-03-17 13:43:55 +01:00
}
async _loadBotIdentity ( ) {
const res = await fetch ( ` ${ TELEGRAM _API } /bot ${ this . botToken } /getMe ` , {
signal : AbortSignal . timeout ( 10000 ) ,
} ) ;
if ( ! res . ok ) {
const err = await res . text ( ) . catch ( ( ) => '' ) ;
throw new Error ( ` getMe failed ( ${ res . status } ): ${ err . substring ( 0 , 200 ) } ` ) ;
}
const data = await res . json ( ) ;
if ( ! data . ok || ! data . result ? . username ) {
throw new Error ( 'getMe returned invalid bot profile' ) ;
}
this . _botUsername = String ( data . result . username ) . toLowerCase ( ) ;
}
async _setMyCommands ( commands , scope = null ) {
const body = { commands } ;
if ( scope ) body . scope = scope ;
const res = await fetch ( ` ${ TELEGRAM _API } /bot ${ this . botToken } /setMyCommands ` , {
method : 'POST' ,
headers : { 'Content-Type' : 'application/json' } ,
body : JSON . stringify ( body ) ,
signal : AbortSignal . timeout ( 10000 ) ,
} ) ;
if ( ! res . ok ) {
const err = await res . text ( ) . catch ( ( ) => '' ) ;
throw new Error ( ` setMyCommands failed ( ${ res . status } ): ${ err . substring ( 0 , 200 ) } ` ) ;
}
const data = await res . json ( ) ;
if ( ! data . ok ) {
throw new Error ( ` setMyCommands rejected: ${ JSON . stringify ( data ) . substring ( 0 , 200 ) } ` ) ;
}
}
2026-03-18 08:20:32 +01:00
_buildConfiguredChatScope ( ) {
const chatId = Number ( this . chatId ) ;
if ( ! Number . isSafeInteger ( chatId ) ) {
throw new Error ( ` TELEGRAM_CHAT_ID must be a numeric chat id, got: ${ this . chatId } ` ) ;
}
return { type : 'chat' , chat _id : chatId } ;
}
2026-03-17 13:43:55 +01:00
_normalizeCommand ( rawCommand ) {
if ( ! rawCommand . startsWith ( '/' ) ) return null ;
const atIdx = rawCommand . indexOf ( '@' ) ;
if ( atIdx === - 1 ) return rawCommand ;
const command = rawCommand . substring ( 0 , atIdx ) ;
const mentionedBot = rawCommand . substring ( atIdx + 1 ) . toLowerCase ( ) ;
if ( ! this . _botUsername || mentionedBot === this . _botUsername ) return command ;
return null ;
}
2026-03-14 00:35:31 -07:00
// ─── Semantic Dedup ─────────────────────────────────────────────────────
/ * *
* Generate a content - based hash for a signal to detect near - duplicates .
* Uses normalized text + key metrics rather than raw text prefix matching .
* /
_contentHash ( signal ) {
// Normalize: lowercase, strip numbers that change frequently (timestamps, exact values)
let content = '' ;
if ( signal . text ) {
content = signal . text . toLowerCase ( )
. replace ( /\d{1,2}:\d{2}/g , '' ) // strip times
. replace ( /\d+\.\d+%?/g , 'NUM' ) // normalize numbers
. replace ( /\s+/g , ' ' )
. trim ( )
. substring ( 0 , 120 ) ;
} else if ( signal . label ) {
// For metric signals, hash the label + direction (not exact values)
content = ` ${ signal . label } : ${ signal . direction || 'none' } ` ;
} else {
content = signal . key || JSON . stringify ( signal ) . substring ( 0 , 80 ) ;
}
return createHash ( 'sha256' ) . update ( content ) . digest ( 'hex' ) . substring ( 0 , 16 ) ;
}
_isSemanticDuplicate ( signal ) {
const hash = this . _contentHash ( signal ) ;
const lastSeen = this . _contentHashes [ hash ] ;
if ( ! lastSeen ) return false ;
// Consider duplicate if seen within last 4 hours
const fourHoursAgo = Date . now ( ) - 4 * 60 * 60 * 1000 ;
return new Date ( lastSeen ) . getTime ( ) > fourHoursAgo ;
}
_recordContentHash ( signal ) {
const hash = this . _contentHash ( signal ) ;
this . _contentHashes [ hash ] = new Date ( ) . toISOString ( ) ;
// Prune hashes older than 24h
const cutoff = Date . now ( ) - 24 * 60 * 60 * 1000 ;
for ( const [ h , ts ] of Object . entries ( this . _contentHashes ) ) {
if ( new Date ( ts ) . getTime ( ) < cutoff ) delete this . _contentHashes [ h ] ;
}
}
_signalKey ( signal ) {
// Improved key generation — use content hash for text signals, structured key for metrics
if ( signal . text ) return ` tg: ${ this . _contentHash ( signal ) } ` ;
return signal . key || signal . label || JSON . stringify ( signal ) . substring ( 0 , 60 ) ;
}
// ─── Rate Limiting ──────────────────────────────────────────────────────
_checkRateLimit ( tier ) {
const config = TIER _CONFIG [ tier ] ;
if ( ! config ) return true ;
const now = Date . now ( ) ;
const oneHourAgo = now - 60 * 60 * 1000 ;
// Check cooldown since last alert of same or lower tier
const lastSameTier = this . _alertHistory
. filter ( a => a . tier === tier )
. pop ( ) ;
if ( lastSameTier && ( now - lastSameTier . timestamp ) < config . cooldownMs ) {
2026-03-12 23:45:46 -07:00
return false ;
}
2026-03-14 00:35:31 -07:00
// Check hourly cap
const recentCount = this . _alertHistory
. filter ( a => a . tier === tier && a . timestamp > oneHourAgo )
. length ;
if ( recentCount >= config . maxPerHour ) {
return false ;
}
return true ;
}
_recordAlert ( tier ) {
this . _alertHistory . push ( { tier , timestamp : Date . now ( ) } ) ;
// Keep only last 50 alerts
if ( this . _alertHistory . length > 50 ) {
this . _alertHistory = this . _alertHistory . slice ( - 50 ) ;
}
}
_isMuted ( ) {
if ( ! this . _muteUntil ) return false ;
if ( Date . now ( ) > this . _muteUntil ) {
this . _muteUntil = null ;
return false ;
}
return true ;
}
// ─── Prompt Engineering ─────────────────────────────────────────────────
_buildEvaluationPrompt ( ) {
return ` You are Crucix, an elite intelligence alert evaluator for a personal OSINT monitoring system. You analyze signal deltas from a 25-source intelligence sweep and decide if the user needs to be alerted via Telegram.
# # Your Decision Framework
You must classify each evaluation into one of four outcomes :
# # # NO ALERT — suppress if :
- Routine scheduled data ( NFP , CPI , FOMC minutes on expected dates ) UNLESS the deviation from consensus is extreme ( > 2 σ )
- Continuation of existing trends already flagged in prior sweeps
- Low - confidence signals from single sources without corroboration
- Social media noise without hard - data confirmation ( Telegram chatter alone is NOT enough )
# # # 🔴 FLASH — immediate , life - of - portfolio risk :
- Active military escalation between nuclear powers or NATO - involved states
- Flash crash indicators ( VIX spike > 40 % , major index down > 3 % intraday )
- Central bank emergency action ( unscheduled rate decision , emergency lending facility )
- Nuclear / radiological anomaly confirmed by multiple monitors
- Sanctions against major economy announced without warning
FLASH requires : ≥ 2 corroborating sources across different domains ( e . g . OSINT + market data + satellite )
# # # 🟡 PRIORITY — act within hours :
- Significant market dislocation ( VIX > 25 AND credit spreads widening )
- Geopolitical escalation with clear energy / commodity transmission ( conflict + oil move > 3 % )
- Unexpected economic data ( > 1.5 σ miss on major indicator )
- New conflict front or ceasefire collapse confirmed by ACLED + Telegram
PRIORITY requires : ≥ 2 signals moving in same direction , at least 1 from hard data
# # # 🔵 ROUTINE — informational , no urgency :
- Notable trend shifts or reversals worth tracking
- Single - source signals of moderate importance
- Cumulative drift ( multiple small moves in same direction over several sweeps )
# # Output Format
Respond with ONLY valid JSON :
{
"shouldAlert" : true / false ,
"tier" : "FLASH" | "PRIORITY" | "ROUTINE" ,
"headline" : "10-word max headline" ,
"reason" : "2-3 sentences. What happened, why it matters, what to watch next." ,
"actionable" : "Specific action the user could take (or 'Monitor' if just informational)" ,
"signals" : [ "signal1" , "signal2" ] ,
"confidence" : "HIGH" | "MEDIUM" | "LOW" ,
"crossCorrelation" : "Which domains are confirming each other (e.g. 'conflict + energy + satellite')"
} ` ;
}
_buildSignalContext ( signals , delta ) {
const sections = [ ] ;
// Categorize signals
const marketSignals = signals . filter ( s => [ 'vix' , 'hy_spread' , 'wti' , 'brent' , 'natgas' , '10y2y' , 'fed_funds' , '10y_yield' , 'usd_index' ] . includes ( s . key ) ) ;
const osintSignals = signals . filter ( s => s . key === 'tg_urgent' || s . item ? . channel ) ;
const conflictSignals = signals . filter ( s => [ 'conflict_events' , 'conflict_fatalities' , 'thermal_total' ] . includes ( s . key ) ) ;
const otherSignals = signals . filter ( s => ! marketSignals . includes ( s ) && ! osintSignals . includes ( s ) && ! conflictSignals . includes ( s ) ) ;
if ( marketSignals . length > 0 ) {
sections . push ( '📊 MARKET SIGNALS:\n' + marketSignals . map ( s =>
` ${ s . label } : ${ s . from } → ${ s . to } ( ${ s . pctChange > 0 ? '+' : '' } ${ s . pctChange ? . toFixed ( 1 ) || s . change } ${ s . pctChange !== undefined ? '%' : '' } ) `
) . join ( '\n' ) ) ;
}
if ( osintSignals . length > 0 ) {
sections . push ( '📡 OSINT SIGNALS:\n' + osintSignals . map ( s => {
const post = s . item || s ;
2026-03-20 16:49:58 -04:00
return ` [ ${ post . channel || 'UNKNOWN' } ] ${ post . text || s . reason || '' } ` ;
2026-03-14 00:35:31 -07:00
} ) . join ( '\n' ) ) ;
}
if ( conflictSignals . length > 0 ) {
sections . push ( '⚔️ CONFLICT INDICATORS:\n' + conflictSignals . map ( s =>
` ${ s . label } : ${ s . from } → ${ s . to } ( ${ s . direction } ) `
) . join ( '\n' ) ) ;
}
if ( otherSignals . length > 0 ) {
sections . push ( '📌 OTHER:\n' + otherSignals . map ( s =>
` ${ s . label || s . key || s . reason } : ${ s . from !== undefined ? ` ${ s . from } → ${ s . to } ` : 'new signal' } `
) . join ( '\n' ) ) ;
}
sections . push ( ` \n 📈 SWEEP DELTA: direction= ${ delta . summary . direction } , total= ${ delta . summary . totalChanges } , critical= ${ delta . summary . criticalChanges } ` ) ;
return sections . join ( '\n\n' ) ;
}
// ─── Message Formatting ─────────────────────────────────────────────────
_formatTieredAlert ( evaluation , delta , tier ) {
const tc = TIER _CONFIG [ tier ] ;
const confidenceEmoji = { HIGH : '🟢' , MEDIUM : '🟡' , LOW : '⚪' } [ evaluation . confidence ] || '⚪' ;
const lines = [
` ${ tc . emoji } *CRUCIX ${ tc . label } * ` ,
` ` ,
` * ${ evaluation . headline } * ` ,
` ` ,
evaluation . reason ,
` ` ,
` Confidence: ${ confidenceEmoji } ${ evaluation . confidence || 'MEDIUM' } ` ,
` Direction: ${ delta . summary . direction . toUpperCase ( ) } ` ,
] ;
if ( evaluation . crossCorrelation ) {
lines . push ( ` Cross-correlation: ${ evaluation . crossCorrelation } ` ) ;
}
if ( evaluation . actionable && evaluation . actionable !== 'Monitor' ) {
lines . push ( ` ` , ` 💡 *Action:* ${ evaluation . actionable } ` ) ;
}
if ( evaluation . signals ? . length ) {
2026-03-21 12:59:30 -04:00
lines . push ( '' , ` *Signals:* ` ) ;
for ( const sig of evaluation . signals ) {
lines . push ( ` • ${ sig } ` ) ;
}
2026-03-14 00:35:31 -07:00
}
lines . push ( '' , ` _ ${ new Date ( ) . toISOString ( ) . replace ( 'T' , ' ' ) . substring ( 0 , 19 ) } UTC_ ` ) ;
return lines . join ( '\n' ) ;
2026-03-12 23:45:46 -07:00
}
}
2026-03-14 00:35:31 -07:00
// ─── Helpers ──────────────────────────────────────────────────────────────
function parseJSON ( text ) {
2026-03-12 23:45:46 -07:00
if ( ! text ) return null ;
let cleaned = text . trim ( ) ;
if ( cleaned . startsWith ( '```' ) ) {
cleaned = cleaned . replace ( /^```(?:json)?\n?/ , '' ) . replace ( /\n?```$/ , '' ) ;
}
try {
return JSON . parse ( cleaned ) ;
} catch {
const match = cleaned . match ( /\{[\s\S]*\}/ ) ;
if ( match ) {
try { return JSON . parse ( match [ 0 ] ) ; } catch { /* give up */ }
}
return null ;
}
}