← back to Norma
agents/twitch-agent/skills/monitor.js
198 lines
/**
* Twitch Monitor Skill
*
* Monitors configured Twitch channels for chat messages matching
* platform-related keywords. Reads monitor targets from the
* sdcc_agent_monitors table and joins channels via TMI.js.
*
* Rate limit: 30/min
*/
const { RateLimiter } = require('../../shared/rate-limiter');
const { logAction } = require('../../shared/audit-logger');
const { query: dbQuery } = require('../../shared/db');
const {
initChat,
joinChannel,
getRecentMessages,
getJoinedChannels,
getChatStatus,
getStreamsByCategory,
MONITORED_CATEGORIES,
KEYWORDS,
} = require('../lib/twitch');
const AGENT = 'twitch-agent';
const PLATFORM = 'twitch';
const rateLimiter = new RateLimiter(AGENT, PLATFORM);
/**
* @param {Object} body - Request body
* @param {string[]} [body.channels] - Additional channels to join/monitor
* @param {string[]} [body.categories] - Category names to scan for relevant streams
* @param {boolean} [body.check_db=true] - Check sdcc_agent_monitors table for targets
* @param {number} [body.recent_limit=50] - How many recent keyword messages to return
* @returns {Promise<Object>}
*/
module.exports = async function monitor(body) {
// Check rate limit
const limit = await rateLimiter.checkLimit('monitor');
if (!limit.allowed) {
await logAction({
agent: AGENT,
actionType: 'monitor',
platform: PLATFORM,
status: 'rate_limited',
errorMessage: `Rate limited. Retry after ${Math.ceil(limit.retryAfterMs / 1000)}s`,
});
return {
rate_limited: true,
retry_after_seconds: Math.ceil(limit.retryAfterMs / 1000),
current: limit.current,
max: limit.max,
};
}
const checkDb = body.check_db !== false;
const recentLimit = Math.min(body.recent_limit || 50, MAX_RECENT);
// Initialize chat client
await initChat();
const channelsToJoin = new Set();
const monitorEntries = [];
// 1. Load monitor targets from database
if (checkDb) {
try {
const { rows } = await dbQuery(
`SELECT id, target, keywords, is_active, last_checked_at, items_found
FROM sdcc_agent_monitors
WHERE agent = $1 AND platform = $2 AND is_active = true
ORDER BY created_at ASC`,
[AGENT, PLATFORM]
);
for (const row of rows) {
channelsToJoin.add(row.target.toLowerCase().replace('#', ''));
monitorEntries.push({
id: row.id,
channel: row.target,
keywords: row.keywords || KEYWORDS,
lastChecked: row.last_checked_at,
itemsFound: row.items_found || 0,
});
}
console.log(`[${AGENT}] Loaded ${rows.length} monitor targets from database`);
} catch (err) {
console.error(`[${AGENT}] Failed to load monitor targets:`, err.message);
}
}
// 2. Add channels from request body
if (body.channels && Array.isArray(body.channels)) {
for (const ch of body.channels) {
channelsToJoin.add(ch.toLowerCase().replace('#', ''));
}
}
// 3. Join all channels
const joinResults = [];
for (const ch of channelsToJoin) {
const joined = await joinChannel(ch);
joinResults.push({ channel: ch, joined });
}
// 4. Optionally scan categories for relevant live streams to auto-join
const categoryScans = [];
const categories = body.categories || [];
for (const catName of categories) {
const catId = MONITORED_CATEGORIES[catName];
if (!catId) {
categoryScans.push({ category: catName, error: 'Unknown category ID' });
continue;
}
try {
const streams = await getStreamsByCategory(catId, 10);
const relevant = streams.filter((s) => s.relevance_score > 5);
for (const stream of relevant.slice(0, 3)) {
const login = stream.user_login;
if (!channelsToJoin.has(login)) {
channelsToJoin.add(login);
const joined = await joinChannel(login);
joinResults.push({ channel: login, joined, source: catName, relevance: stream.relevance_score });
}
}
categoryScans.push({
category: catName,
totalStreams: streams.length,
relevantStreams: relevant.length,
autoJoined: relevant.slice(0, 3).map((s) => s.user_login),
});
} catch (err) {
categoryScans.push({ category: catName, error: err.message });
}
}
// 5. Get recent keyword-matching messages
const allRecent = getRecentMessages();
const recentMessages = allRecent.slice(-recentLimit);
// 6. Update last_checked_at in database for all monitor entries
for (const entry of monitorEntries) {
try {
const channelMessages = allRecent.filter(
(m) => m.channel === entry.channel.toLowerCase().replace('#', '')
);
await dbQuery(
`UPDATE sdcc_agent_monitors
SET last_checked_at = NOW(), items_found = $1, updated_at = NOW()
WHERE id = $2`,
[channelMessages.length, entry.id]
);
} catch (err) {
console.error(`[${AGENT}] Failed to update monitor entry ${entry.id}:`, err.message);
}
}
// Get chat client status
const chatStatus = getChatStatus();
// Record the action
await rateLimiter.recordAction('monitor');
await logAction({
agent: AGENT,
actionType: 'monitor',
platform: PLATFORM,
content: `Monitoring ${channelsToJoin.size} channels`,
responseData: {
channelsMonitored: channelsToJoin.size,
dbEntries: monitorEntries.length,
recentKeywordMessages: recentMessages.length,
categoriesScanned: categoryScans.length,
},
status: 'success',
});
return {
status: 'monitoring',
chat: chatStatus,
channels_joined: joinResults,
db_monitors: monitorEntries,
category_scans: categoryScans,
recent_keyword_messages: recentMessages,
keywords: KEYWORDS,
checked_at: new Date().toISOString(),
};
};
const MAX_RECENT = 200;