← 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(),
};
};