← back to Norma

agents/datasource-agent/skills/discover-apis.js

342 lines

/**
 * Discover APIs Skill
 *
 * Searches known API catalogs (Data.gov, Socrata, ProPublica, CFPB, FRED, etc.)
 * for data APIs related to education, nonprofit advocacy, and consumer finance.
 * Uses Gemini 2.0 Flash to score relevance of each discovered API.
 */

const fetch = require('node-fetch');
const { logAction } = require('../../shared/audit-logger');
const apiCatalogs = require('../lib/api-registry');

const AGENT = 'datasource-agent';
const PLATFORM = 'data-apis';

const GEMINI_URL =
  'https://generativelanguage.googleapis.com/v1beta/models/gemini-2.0-flash:generateContent?key=${GOOGLE_API_KEY}';

const DB_URL = process.env.DATABASE_URL || 'postgresql://dw_admin@127.0.0.1:5432/sdcc';

// Lazy-load pg pool
let _pool;
function getPool() {
  if (!_pool) {
    const { Pool } = require('pg');
    _pool = new Pool({ connectionString: DB_URL, max: 5 });
  }
  return _pool;
}

/**
 * Score a data source with Gemini for relevance to nonprofit advocacy topics.
 *
 * @param {string} name - API/dataset name
 * @param {string} description - Description text
 * @returns {Promise<{ score: number, topics: string[], reasoning: string }>}
 */
async function geminiScoreAPI(name, description) {
  try {
    const prompt = `You are a data analyst evaluating public data sources for a nonprofit civic advocacy organization.

Rate the following data source on relevance (0-100) to these topics:
- Nonprofit advocacy, civic engagement, social impact
- Consumer finance, CFPB complaints, economic policy
- Higher education enrollment, graduation rates, earnings
- Economic inequality, wage stagnation, cost of living
- Education policy, federal aid, public benefits

Data source name: ${name}
Description: ${(description || 'No description available').substring(0, 500)}

Respond in EXACTLY this JSON format (no markdown, no code fences):
{"score": 75, "topics": ["nonprofit advocacy", "higher education"], "reasoning": "Brief explanation"}`;

    const res = await fetch(GEMINI_URL, {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify({
        contents: [{ parts: [{ text: prompt }] }],
        generationConfig: { temperature: 0.2, maxOutputTokens: 500 },
      }),
    });

    const data = await res.json();
    const text = data?.candidates?.[0]?.content?.parts?.[0]?.text || '';

    // Try to parse JSON from the response
    const jsonMatch = text.match(/\{[\s\S]*\}/);
    if (jsonMatch) {
      const parsed = JSON.parse(jsonMatch[0]);
      return {
        score: Math.min(100, Math.max(0, Number(parsed.score) || 0)),
        topics: Array.isArray(parsed.topics) ? parsed.topics : [],
        reasoning: parsed.reasoning || '',
      };
    }

    return { score: 0, topics: [], reasoning: 'Failed to parse Gemini response' };
  } catch (err) {
    console.error(`[${AGENT}] Gemini scoring error for "${name}":`, err.message);
    return { score: 0, topics: [], reasoning: 'Gemini API error: ' + err.message };
  }
}

/**
 * Search a CKAN catalog (Data.gov).
 *
 * @param {Object} catalog - Catalog entry from api-registry
 * @returns {Promise<Array>} Array of discovered APIs
 */
async function searchCKAN(catalog) {
  const results = [];
  try {
    const url = new URL(catalog.url);
    Object.entries(catalog.params).forEach(([k, v]) => url.searchParams.set(k, v));

    const res = await fetch(url.toString(), { timeout: 15000 });
    if (!res.ok) return results;

    const data = await res.json();
    const packages = data?.result?.results || [];

    for (const pkg of packages) {
      results.push({
        name: pkg.title || pkg.name || 'Unknown',
        description: (pkg.notes || pkg.description || '').substring(0, 2000),
        api_url: pkg.url || `https://catalog.data.gov/dataset/${pkg.name || pkg.id}`,
        documentation_url: pkg.url || null,
        category: 'government',
        data_type: pkg.type || 'dataset',
        auth_required: false,
      });
    }
  } catch (err) {
    console.error(`[${AGENT}] CKAN search error (${catalog.name}):`, err.message);
  }
  return results;
}

/**
 * Search Socrata Discovery API.
 *
 * @param {Object} catalog - Catalog entry from api-registry
 * @returns {Promise<Array>} Array of discovered APIs
 */
async function searchSocrata(catalog) {
  const results = [];
  try {
    const url = new URL(catalog.url);
    Object.entries(catalog.params).forEach(([k, v]) => url.searchParams.set(k, v));

    const res = await fetch(url.toString(), { timeout: 15000 });
    if (!res.ok) return results;

    const data = await res.json();
    const datasets = data?.results || [];

    for (const ds of datasets) {
      const resource = ds.resource || {};
      results.push({
        name: resource.name || ds.name || 'Unknown Socrata Dataset',
        description: (resource.description || ds.description || '').substring(0, 2000),
        api_url: ds.link || resource.download_count ? `https://${ds.metadata?.domain || 'data.gov'}/resource/${resource.id}.json` : '',
        documentation_url: ds.link || null,
        category: 'open-data',
        data_type: resource.type || 'dataset',
        auth_required: false,
      });
    }
  } catch (err) {
    console.error(`[${AGENT}] Socrata search error:`, err.message);
  }
  return results;
}

/**
 * Search ProPublica Nonprofit API.
 *
 * @param {Object} catalog - Catalog entry from api-registry
 * @returns {Promise<Array>} Array of discovered APIs
 */
async function searchProPublica(catalog) {
  const results = [];
  try {
    const url = new URL(catalog.url);
    Object.entries(catalog.params).forEach(([k, v]) => url.searchParams.set(k, v));

    const res = await fetch(url.toString(), { timeout: 15000 });
    if (!res.ok) return results;

    const data = await res.json();
    const orgs = data?.organizations || [];

    for (const org of orgs.slice(0, 20)) {
      results.push({
        name: `ProPublica: ${org.name || 'Unknown Org'}`,
        description: `Nonprofit organization - EIN: ${org.ein || 'N/A'}, State: ${org.state || 'N/A'}`,
        api_url: `https://projects.propublica.org/nonprofits/api/v2/organizations/${org.ein}.json`,
        documentation_url: 'https://projects.propublica.org/nonprofits/api',
        category: 'nonprofit',
        data_type: 'organization',
        auth_required: false,
      });
    }
  } catch (err) {
    console.error(`[${AGENT}] ProPublica search error:`, err.message);
  }
  return results;
}

/**
 * Generic search fallback for catalogs without a specific parser.
 *
 * @param {Object} catalog - Catalog entry
 * @returns {Promise<Array>} Array with one entry for the catalog itself
 */
async function searchGeneric(catalog) {
  return [{
    name: `${catalog.name} API`,
    description: `${catalog.name} data catalog endpoint (type: ${catalog.type})`,
    api_url: catalog.url,
    documentation_url: catalog.url,
    category: catalog.type,
    data_type: 'api',
    auth_required: catalog.type === 'fred' || catalog.type === 'fec',
  }];
}

/**
 * Insert a discovered API into the database.
 *
 * @param {Object} api - API data
 * @param {{ score: number, topics: string[], reasoning: string }} scoring - Gemini scoring result
 * @returns {Promise<Object|null>}
 */
async function insertAPI(api, scoring) {
  const pool = getPool();
  try {
    const { rows } = await pool.query(
      `INSERT INTO discovered_apis
        (name, description, api_url, documentation_url, category, data_type, auth_required, relevance_score, match_topics, status, discovered_by)
       VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
       ON CONFLICT DO NOTHING
       RETURNING id, name, relevance_score`,
      [
        api.name,
        api.description,
        api.api_url,
        api.documentation_url,
        api.category,
        api.data_type,
        api.auth_required,
        scoring.score,
        scoring.topics,
        scoring.score >= 50 ? 'relevant' : 'discovered',
        AGENT,
      ]
    );
    return rows[0] || null;
  } catch (err) {
    console.error(`[${AGENT}] DB insert error for "${api.name}":`, err.message);
    return null;
  }
}

/**
 * Main discover skill handler.
 *
 * @param {Object} body - Request body
 * @param {boolean} [body.score=true] - Whether to AI-score results
 * @param {number} [body.max_per_catalog=20] - Max results per catalog
 * @returns {Promise<Object>} Discovery results
 */
module.exports = async function discoverApis(body = {}) {
  const shouldScore = body.score !== false;
  const maxPerCatalog = body.max_per_catalog || 20;

  let totalDiscovered = 0;
  let totalScored = 0;
  const topApis = [];
  const errors = [];

  for (const catalog of apiCatalogs) {
    console.log(`[${AGENT}] Searching ${catalog.name} (${catalog.type})...`);

    let found = [];
    try {
      switch (catalog.type) {
        case 'ckan':
          found = await searchCKAN(catalog);
          break;
        case 'socrata':
          found = await searchSocrata(catalog);
          break;
        case 'propublica':
          found = await searchProPublica(catalog);
          break;
        default:
          found = await searchGeneric(catalog);
          break;
      }
    } catch (err) {
      errors.push({ catalog: catalog.name, error: err.message });
      continue;
    }

    found = found.slice(0, maxPerCatalog);
    console.log(`[${AGENT}] Found ${found.length} from ${catalog.name}`);

    for (const api of found) {
      let scoring = { score: 0, topics: [], reasoning: 'Not scored' };

      if (shouldScore) {
        scoring = await geminiScoreAPI(api.name, api.description);
        totalScored++;

        // Small delay to avoid Gemini rate limits
        await new Promise((r) => setTimeout(r, 300));
      }

      const inserted = await insertAPI(api, scoring);
      if (inserted) {
        totalDiscovered++;
        if (scoring.score >= 60) {
          topApis.push({
            name: inserted.name,
            score: scoring.score,
            topics: scoring.topics,
          });
        }
      }
    }
  }

  // Sort top APIs by score descending
  topApis.sort((a, b) => b.score - a.score);

  await logAction({
    agent: AGENT,
    actionType: 'discover',
    platform: PLATFORM,
    content: `Searched ${apiCatalogs.length} catalogs`,
    responseData: {
      catalogs_searched: apiCatalogs.length,
      discovered: totalDiscovered,
      scored: totalScored,
      top_count: topApis.length,
      errors: errors.length,
    },
    status: 'success',
  });

  return {
    discovered: totalDiscovered,
    scored: totalScored,
    top_apis: topApis.slice(0, 10),
    errors,
    catalogs_searched: apiCatalogs.length,
    fetched_at: new Date().toISOString(),
  };
};