← back to Ken
src/app/api/markets/route.js
109 lines
import { NextResponse } from 'next/server';
const { query } = require('../../../lib/db');
const { getMarkets, isWeatherMarket, parseWeatherEvent } = require('../../../lib/polymarket');
export const dynamic = 'force-dynamic';
export async function GET() {
try {
const res = await query(`
SELECT m.*,
(SELECT COUNT(*) FROM predictions p WHERE p.condition_id = m.condition_id) as prediction_count
FROM markets m
ORDER BY m.created_at DESC
LIMIT 100
`);
return NextResponse.json({ markets: res.rows });
} catch (err) {
return NextResponse.json({ error: err.message }, { status: 500 });
}
}
export async function POST(req) {
const body = await req.json();
if (body.action === 'test_api') {
try {
const markets = await getMarkets({ limit: 5 });
return NextResponse.json({ success: true, market_count: Array.isArray(markets) ? markets.length : 0 });
} catch (err) {
return NextResponse.json({ success: false, error: err.message });
}
}
if (body.action === 'ingest') {
try {
let allMarkets = [];
let offset = 0;
const limit = 100;
// Paginate through Gamma API
for (let page = 0; page < 10; page++) {
const batch = await getMarkets({ limit, offset: offset + page * limit });
if (!Array.isArray(batch) || batch.length === 0) break;
allMarkets = allMarkets.concat(batch);
if (batch.length < limit) break;
await new Promise(r => setTimeout(r, 600)); // rate limit
}
// Filter weather markets
const weatherMarkets = allMarkets.filter(isWeatherMarket);
let ingested = 0;
for (const m of weatherMarkets) {
const parsed = parseWeatherEvent(m);
try {
await query(`
INSERT INTO markets (condition_id, slug, question, description, resolution_source, end_ts, category, liquidity, parsed_event, raw_gamma)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
ON CONFLICT (condition_id) DO UPDATE SET
liquidity = EXCLUDED.liquidity, raw_gamma = EXCLUDED.raw_gamma,
parsed_event = EXCLUDED.parsed_event, updated_at = NOW()
`, [
m.conditionId || m.id, m.slug, m.question, m.description,
m.resolutionSource, m.endDate, 'weather',
m.liquidity || 0, parsed, m
]);
// Insert tokens if available
if (m.tokens && Array.isArray(m.tokens)) {
for (const t of m.tokens) {
await query(`
INSERT INTO tokens (asset_id, condition_id, outcome, price, raw_gamma)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (asset_id) DO UPDATE SET price = EXCLUDED.price, raw_gamma = EXCLUDED.raw_gamma, updated_at = NOW()
`, [t.token_id || t.asset_id || `${m.conditionId}_${t.outcome}`, m.conditionId || m.id, t.outcome, t.price || 0, t]);
}
}
ingested++;
} catch (err) {
console.error(`[Markets] Ingest error for ${m.conditionId}:`, err.message);
}
}
// Log alert
await query(`INSERT INTO alerts_log (alert_type, severity, message, details) VALUES ($1, $2, $3, $4)`,
['market_ingest', 'info', `Ingested ${ingested} weather markets from ${allMarkets.length} total`, { total_scanned: allMarkets.length, weather_found: weatherMarkets.length, ingested }]);
return NextResponse.json({ ingested, total_scanned: allMarkets.length, weather_candidates: weatherMarkets.length });
} catch (err) {
return NextResponse.json({ error: err.message }, { status: 500 });
}
}
if (body.action === 'parse') {
try {
const res = await query('SELECT * FROM markets WHERE condition_id = $1', [body.condition_id]);
if (!res.rows[0]) return NextResponse.json({ error: 'Market not found' }, { status: 404 });
const market = res.rows[0];
const parsed = parseWeatherEvent(market.raw_gamma || market);
await query('UPDATE markets SET parsed_event = $1, updated_at = NOW() WHERE condition_id = $2', [parsed, body.condition_id]);
return NextResponse.json({ parsed, condition_id: body.condition_id });
} catch (err) {
return NextResponse.json({ error: err.message }, { status: 500 });
}
}
return NextResponse.json({ error: 'Unknown action' }, { status: 400 });
}