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