import {
  createConfig,
  createServer,
  defaultEndpointsFactory,
  DependsOnMethod,
  Routing,
  EventStreamFactory,
  Documentation,
} from "express-zod-api";
import {
  postCandlesBodySchema,
  getCandlesQuerySchema,
  getCandlesRangeSchema,
  dumpMessageSchema,
  candles,
} from "./schema";
import {
  upsertTicker,
  upsertTimeframe,
  upsertCandles,
  fetchCandles,
  fetchAllTickers,
  fetchAllTimeframes,
  fetchTimestampRange,
  getTickerId,
  getTimeframeId,
} from "./candles";
import { z } from "zod";
import ui from "swagger-ui-express";
import { createPostgresDrizzleStreamer } from "./extend_drizzle";
import { sql, db } from "./db";
import { eq, and, asc } from "drizzle-orm";
import express from "express";

const port = process.env.PORT ? parseInt(process.env.PORT) : 3000;

// Output schemas
const postCandlesOutput = z.object({
  status: z.string(),
  inserted: z.number(),
});
const getCandlesOutput = z.object({
  ticker: z.string(),
  timeframe: z.string(),
  headers: z.array(z.string()),
  candles: z.array(z.array(z.number())),
  pagination: z.object({
    page: z.number(),
    limit: z.number(),
    count: z.number(),
  }),
});
const getTickersOutput = z.object({ tickers: z.array(z.string()) });
const getTimeframesOutput = z.object({ timeframes: z.array(z.string()) });
const getCandlesRangeOutput = z.object({
  ticker: z.string(),
  timeframe: z.string(),
  range: z.object({
    min: z.number().nullable(),
    max: z.number().nullable(),
    minDate: z.string().nullable(),
    maxDate: z.string().nullable(),
    count: z.number().nullable(),
  }),
});

const postCandles = defaultEndpointsFactory.build({
  method: "post",
  input: postCandlesBodySchema,
  output: postCandlesOutput,
  handler: async ({ input }) => {
    const { ticker, timeframe, data } = input;
    const [tickerId, timeframeId] = await Promise.all([
      upsertTicker(ticker),
      upsertTimeframe(timeframe),
    ]);
    const inserted = await upsertCandles(tickerId, timeframeId, data);
    return { status: "ok", inserted };
  },
});

const getCandles = defaultEndpointsFactory.build({
  method: "get",
  input: getCandlesQuerySchema,
  output: getCandlesOutput,
  handler: async ({ input }) => {
    const { ticker, timeframe, from, to, limit, page } = input;
    const offset = (page - 1) * limit;
    const candlesArr = await fetchCandles(
      ticker,
      timeframe,
      from,
      to,
      limit,
      offset
    );
    const headers = ["timestamp", "open", "high", "close", "low", "volume"];
    return {
      ticker,
      timeframe,
      headers,
      candles: candlesArr,
      pagination: { page, limit, count: candlesArr.length },
    };
  },
});

const getTickers = defaultEndpointsFactory.build({
  method: "get",
  input: z.object({}),
  output: getTickersOutput,
  handler: async () => {
    const tickers = await fetchAllTickers();
    return { tickers };
  },
});

const getTimeframes = defaultEndpointsFactory.build({
  method: "get",
  input: z.object({}),
  output: getTimeframesOutput,
  handler: async () => {
    const timeframes = await fetchAllTimeframes();
    return { timeframes };
  },
});

const getCandlesRange = defaultEndpointsFactory.build({
  method: "get",
  input: getCandlesRangeSchema,
  output: getCandlesRangeOutput,
  handler: async ({ input }) => {
    const { ticker, timeframe } = input;
    const range = await fetchTimestampRange(ticker, timeframe);
    return {
      ticker,
      timeframe,
      range: {
        min: range.min,
        max: range.max,
        minDate: range.min ? new Date(range.min * 1000).toISOString() : null,
        maxDate: range.max ? new Date(range.max * 1000).toISOString() : null,
        count: range.count,
      },
    };
  },
});

// Streaming dump endpoint
const dumpCandlesStream = new EventStreamFactory({
  summary: z
    .object({
      min: z.number().nullable(),
      max: z.number().nullable(),
      count: z.number().nullable(),
    })
    .describe("Summary of timestamp range (min, max, count)"),
  candles: z
    .array(z.array(z.number()))
    .describe("Chunk of candles [timestamp, open, high, low, close, volume]"),
}).buildVoid({
  method: "get",
  input: z
    .object({
      ticker: z.string().describe("Ticker symbol"),
      timeframe: z.string().describe("Timeframe name"),
      chunkSize: z.coerce
        .number()
        .int()
        .min(100)
        .max(50_000)
        .default(500)
        .describe("Number of candles per chunk"),
    })
    .describe("Query parameters for streaming candles dump"),
  handler: async ({ input, options: { emit, isClosed }, logger }) => {
    const { ticker, timeframe, chunkSize } = input;
    logger.debug(`Starting dump for ${ticker} ${timeframe}`);
    // emit summary of available range before streaming data
    fetchTimestampRange(ticker, timeframe).then(({ min, max, count }) => {
      emit("summary", { min, max, count });
    });
    const [tickerId, timeframeId] = await Promise.all([
      getTickerId(ticker),
      getTimeframeId(timeframe),
    ]);
    const streamer = createPostgresDrizzleStreamer(sql);
    const query = db
      .select({
        timestamp: candles.timestamp,
        open: candles.open,
        high: candles.high,
        close: candles.close,
        low: candles.low,
        volume: candles.volume,
      })
      .from(candles)
      .$dynamic()
      .where(
        and(
          eq(candles.tickerId, tickerId),
          eq(candles.timeframeId, timeframeId)
        )
      )
      .orderBy(asc(candles.timestamp));
    for await (const batch of streamer(query, chunkSize)) {
      if (isClosed()) break;
      await emit(
        "candles",
        batch.map((row) => [
          row.timestamp,
          Number(row.open),
          Number(row.high),
          Number(row.low),
          Number(row.close),
          Number(row.volume),
        ])
      );
    }
  },
});

// Raw Streaming dump endpoint
const dumpCandlesRawStream = new EventStreamFactory({
  // Define events for raw data, likely just a chunk of raw rows
  // The client will need to know the structure/order
  raw_chunk: z
    .array(z.any())
    .describe("Chunk of raw candle data from the database"),
  // Optional: Still emit summary
  summary: z
    .object({
      min: z.number().nullable(),
      max: z.number().nullable(),
      count: z.number().nullable(),
    })
    .describe("Summary of timestamp range (min, max, count)"),
}).buildVoid({
  method: "get",
  input: z
    .object({
      ticker: z.string().describe("Ticker symbol"),
      timeframe: z.string().describe("Timeframe name"),
      chunkSize: z.coerce
        .number()
        .int()
        .min(100)
        .max(50_000)
        .default(500)
        .describe("Number of rows per chunk"),
    })
    .describe("Query parameters for raw streaming candles dump"),
  handler: async ({ input, options: { emit, isClosed }, logger }) => {
    const { ticker, timeframe, chunkSize } = input;
    logger.debug(`Starting RAW dump for ${ticker} ${timeframe}`);

    // Optional: Emit summary first
    fetchTimestampRange(ticker, timeframe)
      .then(({ min, max, count }) => {
        emit("summary", { min, max, count });
      })
      .catch((err) =>
        logger.error("Error fetching range for raw dump summary:", err)
      ); // Add error handling

    try {
      const [tickerId, timeframeId] = await Promise.all([
        getTickerId(ticker),
        getTimeframeId(timeframe),
      ]);

      if (!tickerId || !timeframeId) {
        logger.warn(
          `Ticker or Timeframe not found for raw dump: ${ticker}, ${timeframe}`
        );
        // Optionally emit an error event or just close
        return;
      }
      const query = db
        .select({
          timestamp: candles.timestamp,
          open: candles.open,
          high: candles.high,
          close: candles.close,
          low: candles.low,
          volume: candles.volume,
        })
        .from(candles)
        .$dynamic()
        .where(
          and(
            eq(candles.tickerId, tickerId),
            eq(candles.timeframeId, timeframeId)
          )
        )
        .orderBy(asc(candles.timestamp));

      const raw_query = query.toSQL();

      // Use postgres.js cursor directly
      const cursor = sql
        .unsafe<Record<string, any>[]>(
          raw_query.sql,
          raw_query.params as Parameters<typeof sql.unsafe>[1]
        )
        .cursor(chunkSize);

      for await (const chunk of cursor) {
        if (isClosed()) {
          logger.debug("Raw stream closed by client.");
          break;
        }
        // Emit the raw chunk directly
        // The client receives data as returned by the driver (e.g., numbers might be strings)
        await emit("raw_chunk", chunk);
      }
      logger.debug(`Finished RAW dump for ${ticker} ${timeframe}`);
    } catch (error) {
      logger.error(`Error during RAW dump for ${ticker} ${timeframe}:`, error);
      // Optionally emit an error event to the client
      // await emit('error', { message: 'Internal server error during streaming' });
    }
  },
});

const routing: Routing = {
  candles: new DependsOnMethod({
    post: postCandles,
    get: getCandles,
  }).nest({
    range: getCandlesRange,
    dump: dumpCandlesStream.nest({
      raw: dumpCandlesRawStream,
    }),
  }),
  tickers: getTickers,
  timeframes: getTimeframes,
};

const config = createConfig({
  http: { listen: port },
  cors: true,
  logger: { level: "debug" },
  beforeRouting: ({ app, getLogger }) => {
    // Serve Swagger UI at /docs
    const documentation = new Documentation({
      routing,
      config,
      version: "1.0.0",
      title: "MacroFinder API",
      serverUrl: `https://candles.macrofinder.flolep.fr`,
    });
    const spec = documentation.getSpec();
    app.use("/docs", ui.serve, ui.setup(spec));
  },
  jsonParser: express.json({
    limit: "50mb",
  }),
  rawParser: express.raw({
    type: "application/octet-stream",
    limit: "50mb",
  }),
  formParser: express.urlencoded({
    extended: true,
    limit: "50mb",
  }),
});

createServer(config, routing);
