← 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;