← back to Norma

agents/reddit-agent/skills/monitor.js

243 lines

/**
 * Reddit Monitor Skill
 *
 * Monitors configured subreddits for relevant advocacy discussions.
 * Fetches new/hot posts, scores relevance, and pushes to Pulse.
 *
 * Rate limit: 60/min (Reddit JSON API)
 */

const { RateLimiter } = require('../../shared/rate-limiter');
const { logAction } = require('../../shared/audit-logger');
const { pushSocialFeed, pushTrending } = require('../../shared/pulse-reporter');
const { getHotPosts, getNewPosts, scoreRelevance, SUBREDDITS, KEYWORDS } = require('../lib/reddit');
const { query: dbQuery } = require('../../shared/db');

const AGENT = 'reddit-agent';
const PLATFORM = 'reddit';

const rateLimiter = new RateLimiter(AGENT, PLATFORM);

/**
 * @param {Object} body - Request body
 * @param {string[]} [body.subreddits] - Override subreddits to monitor
 * @param {string[]} [body.keywords] - Additional keywords to look for
 * @param {number} [body.limit=25] - Posts per subreddit
 * @param {string} [body.feed='hot'] - Feed type: 'hot' or 'new'
 * @param {boolean} [body.push_to_pulse=true] - Push results to Pulse
 * @returns {Promise<Object>}
 */
module.exports = async function monitor(body) {
  const limit = Math.min(body.limit || 25, 100);
  const feed = body.feed || 'hot';
  const pushToPulse = body.push_to_pulse !== false;

  // Check rate limit
  const rateCheck = await rateLimiter.checkLimit('search');
  if (!rateCheck.allowed) {
    await logAction({
      agent: AGENT,
      actionType: 'monitor',
      platform: PLATFORM,
      status: 'rate_limited',
      errorMessage: `Rate limited. Retry after ${Math.ceil(rateCheck.retryAfterMs / 1000)}s`,
    });

    return {
      rate_limited: true,
      retry_after_seconds: Math.ceil(rateCheck.retryAfterMs / 1000),
      current: rateCheck.current,
      max: rateCheck.max,
    };
  }

  // Get active monitors from database
  let dbMonitors = [];
  try {
    const { rows } = await dbQuery(
      `SELECT id, target, keywords, monitor_type
       FROM sdcc_agent_monitors
       WHERE agent = $1 AND platform = $2 AND is_active = true`,
      [AGENT, PLATFORM]
    );
    dbMonitors = rows;
  } catch (err) {
    console.error(`[${AGENT}] Failed to fetch monitors:`, err.message);
  }

  // Build subreddit list from monitors + body + defaults
  const subredditsToMonitor = new Set();

  // From database monitors
  for (const m of dbMonitors) {
    if (m.monitor_type === 'subreddit' && m.target) {
      subredditsToMonitor.add(m.target.replace(/^r\//, ''));
    }
  }

  // From body params
  if (body.subreddits && body.subreddits.length > 0) {
    body.subreddits.forEach((s) => subredditsToMonitor.add(s));
  }

  // Default subreddits if none configured
  if (subredditsToMonitor.size === 0) {
    SUBREDDITS.forEach((s) => subredditsToMonitor.add(s));
  }

  // Additional keywords from body + DB monitors
  const extraKeywords = [...(body.keywords || [])];
  for (const m of dbMonitors) {
    if (m.keywords && Array.isArray(m.keywords)) {
      extraKeywords.push(...m.keywords);
    }
  }

  // Fetch posts from each subreddit
  const allResults = [];
  const fetchFn = feed === 'new' ? getNewPosts : getHotPosts;

  for (const subreddit of subredditsToMonitor) {
    try {
      const posts = await fetchFn(subreddit, limit);
      await rateLimiter.recordAction('search');

      // Score each post for relevance (combining default + extra keywords)
      const scored = posts.map((post) => {
        let relevance = scoreRelevance(post);

        // Bonus for extra keyword matches
        if (extraKeywords.length > 0) {
          const text = `${post.title} ${post.selftext}`.toLowerCase();
          for (const kw of extraKeywords) {
            if (text.includes(kw.toLowerCase())) {
              relevance = Math.min(100, relevance + 5);
            }
          }
        }

        return { ...post, relevance_score: relevance };
      });

      // Filter to relevant posts only
      const relevant = scored.filter((p) => p.relevance_score > 10);
      relevant.sort((a, b) => b.relevance_score - a.relevance_score);

      allResults.push({
        subreddit,
        feed,
        total_fetched: posts.length,
        relevant_count: relevant.length,
        posts: relevant,
      });

      // Update DB monitor stats if applicable
      const dbMonitor = dbMonitors.find(
        (m) => m.target === subreddit || m.target === `r/${subreddit}`
      );
      if (dbMonitor) {
        try {
          await dbQuery(
            `UPDATE sdcc_agent_monitors
             SET last_checked_at = NOW(), items_found = items_found + $1, updated_at = NOW()
             WHERE id = $2`,
            [relevant.length, dbMonitor.id]
          );
        } catch (err) {
          console.error(`[${AGENT}] Failed to update monitor ${dbMonitor.id}:`, err.message);
        }
      }
    } catch (err) {
      console.error(`[${AGENT}] Monitor failed for r/${subreddit}:`, err.message);
      allResults.push({
        subreddit,
        feed,
        total_fetched: 0,
        relevant_count: 0,
        posts: [],
        error: err.message,
      });
    }
  }

  // Aggregate trending keywords across all results
  const keywordCounts = {};
  allResults.flatMap((r) => r.posts).forEach((post) => {
    const text = `${post.title} ${post.selftext}`.toLowerCase();
    for (const kw of KEYWORDS) {
      if (text.includes(kw.toLowerCase())) {
        keywordCounts[kw] = (keywordCounts[kw] || 0) + 1;
      }
    }
  });

  const trending = Object.entries(keywordCounts)
    .sort(([, a], [, b]) => b - a)
    .slice(0, 15)
    .map(([topic, count]) => ({
      topic,
      platform: PLATFORM,
      volume: count,
      sentiment: 'neutral',
    }));

  // Push to Pulse
  if (pushToPulse) {
    if (trending.length > 0) {
      await pushTrending(trending);
    }

    // Push top-relevance posts
    const topPosts = allResults
      .flatMap((r) => r.posts)
      .sort((a, b) => b.relevance_score - a.relevance_score)
      .slice(0, 5);

    for (const post of topPosts) {
      await pushSocialFeed({
        agent: AGENT,
        platform: PLATFORM,
        type: 'discussion',
        title: `r/${post.subreddit}: ${post.title.substring(0, 100)}`,
        url: post.permalink,
        content: post.selftext.substring(0, 300),
        metadata: {
          post_id: post.id,
          score: post.score,
          num_comments: post.num_comments,
          relevance_score: post.relevance_score,
          subreddit: post.subreddit,
        },
      });
    }
  }

  const totalFetched = allResults.reduce((sum, r) => sum + r.total_fetched, 0);
  const totalRelevant = allResults.reduce((sum, r) => sum + r.relevant_count, 0);

  // Log the action
  await logAction({
    agent: AGENT,
    actionType: 'monitor',
    platform: PLATFORM,
    responseData: {
      subreddits_monitored: subredditsToMonitor.size,
      total_fetched: totalFetched,
      total_relevant: totalRelevant,
      trending_keywords: trending.length,
      db_monitors_checked: dbMonitors.length,
      feed,
    },
    status: 'success',
  });

  return {
    subreddits_monitored: subredditsToMonitor.size,
    total_fetched: totalFetched,
    total_relevant: totalRelevant,
    results: allResults,
    trending,
    monitored_at: new Date().toISOString(),
  };
};