import { useAuth } from '@clerk/nextjs';
import { noop } from 'lodash-es';
import { useCallback, useEffect, useMemo, useRef, useState } from 'react';

import queryClient from '@/components/studio/queryClient';
import { components } from '@/lib/gen';

type DownbeatsSchema = components['schemas']['DownbeatsSchema'] & {
  downbeats?: [number, number][] | null;
  raw_downbeats?: [number, number][] | null;
  raw_onset_downbeats?: [number, number][] | null;
  seconds_missing_from_start?: number | null;
};

type StreamState = {
  callbacks: Set<
    ({
      downbeats,
      final,
      secondsMissingFromStart,
    }: {
      downbeats: [number, number][];
      final: boolean;
      secondsMissingFromStart: number;
    }) => void
  >;
  secondsMissingFromStart: number;
  accumulatedDownbeats: [number, number][];
  final: boolean;
};

const getDownbeatsCacheKey = (clipId: string) => {
  return ['sep-23-2025', 'downbeats_streaming_v2', clipId];
};

// Helper functions
const cacheFinalResult = (clipId: string, streamState: StreamState) => {
  if (streamState.final && streamState.accumulatedDownbeats.length > 0) {
    queryClient.setQueryData(getDownbeatsCacheKey(clipId), {
      downbeats: streamState.accumulatedDownbeats,
      final: true,
      secondsMissingFromStart: streamState.secondsMissingFromStart,
    });
  }
};

const notifyCallbacks = (streamState: StreamState, finalOverride?: boolean) => {
  streamState.callbacks.forEach((callback) => {
    try {
      callback({
        downbeats: streamState.accumulatedDownbeats,
        secondsMissingFromStart: streamState.secondsMissingFromStart,
        final: finalOverride !== undefined ? finalOverride : streamState.final,
      });
    } catch (error) {
      console.warn('Error calling downbeats callback', error);
    }
  });
};

const processDataChunk = (
  clipId: string,
  streamState: StreamState,
  data: DownbeatsSchema,
  finalOverride?: boolean
) => {
  if (data.downbeats) streamState.accumulatedDownbeats = data.downbeats;
  streamState.final =
    finalOverride !== undefined ? finalOverride : !!data.final;
  if (data.seconds_missing_from_start !== undefined) {
    streamState.secondsMissingFromStart = data.seconds_missing_from_start ?? 0;
  }

  cacheFinalResult(clipId, streamState);
  notifyCallbacks(streamState, finalOverride);
};

// Global Map to share active streams across all hook instances
const activeStreams = new Map<string, StreamState>();

export const useStreamingDownbeats = () => {
  const { getToken } = useAuth();

  const startStream = useCallback(
    async (clipId: string, streamState: StreamState) => {
      let reader: ReadableStreamDefaultReader<Uint8Array> | null = null;

      try {
        const token = await getToken();

        const response = await fetch(
          `${process.env.NEXT_PUBLIC_API_BASE_ECS}/api/gen/${clipId}/downbeats_streaming/v2`,
          {
            method: 'POST',
            headers: {
              'Content-Type': 'application/json',
              Authorization: `Bearer ${token}`,
            },
            body: JSON.stringify({}),
          }
        );

        if (!response.ok) {
          throw new Error(`HTTP error streaming downbeats: ${response.status}`);
        }

        const contentType = response.headers.get('content-type');
        if (!contentType) {
          throw new Error(`No content type for downbeats stream`);
        }

        if (contentType.includes('application/json')) {
          const data = (await response.json()) as DownbeatsSchema;
          const final = !!data.final || data.state === 'complete';
          processDataChunk(clipId, streamState, data, final);
          return;
        }

        if (!contentType.includes('text/event-stream')) {
          throw new Error(
            `Unexpected content type for downbeats stream: ${contentType}`
          );
        }

        reader = response.body?.getReader() || null;
        const decoder = new TextDecoder('utf-8');

        if (!reader) {
          throw new Error(
            'No response body reader available for streaming downbeats'
          );
        }

        let accum = '';

        while (true) {
          const { done, value } = await reader.read();

          if (done) break;

          if (value) {
            const chunk = decoder.decode(value);
            accum += chunk;
            const splitLines = accum.split('\n\n');

            if (splitLines.length > 1) {
              accum = splitLines.pop() as string;
              for (const line of splitLines) {
                const trimmedLine = line.trim();
                if (trimmedLine.startsWith('data: ')) {
                  try {
                    const data = JSON.parse(
                      trimmedLine.slice(6)
                    ) as DownbeatsSchema;
                    processDataChunk(clipId, streamState, data);
                  } catch (parseError) {
                    console.warn(
                      'Error parsing streaming downbeats data frame',
                      parseError
                    );
                    continue;
                  }
                }
              }
            }
          }
        }

        if (accum.trim()) {
          const trimmedAccum = accum.trim();
          if (trimmedAccum.startsWith('data: ')) {
            try {
              const data = JSON.parse(trimmedAccum.slice(6)) as DownbeatsSchema;
              processDataChunk(clipId, streamState, data);
            } catch (parseError) {
              console.warn(
                'Error parsing final streaming downbeats data frame',
                parseError
              );
            }
          }
        }
      } catch (error) {
        console.error('Error in downbeats stream', error);
      } finally {
        if (reader) {
          try {
            reader.releaseLock();
          } catch (releaseError) {
            console.warn(
              'Error releasing reader lock for streaming downbeats',
              releaseError
            );
          }
        }
        // Clean up the stream state when done
        activeStreams.delete(clipId);
      }
    },
    [getToken]
  );

  const streamDownbeatsForClip = useCallback(
    (
      clipId: string,
      callbackFunction?: ({
        downbeats,
        final,
        secondsMissingFromStart,
      }: {
        downbeats: [number, number][];
        final: boolean;
        secondsMissingFromStart: number;
      }) => void
    ): (() => void) => {
      // Check for cached final result first
      const cachedData = queryClient.getQueryData<{
        downbeats: [number, number][];
        final: boolean;
        secondsMissingFromStart: number;
      }>(getDownbeatsCacheKey(clipId));

      if (cachedData?.final && callbackFunction) {
        // Return cached data immediately
        try {
          callbackFunction(cachedData);
        } catch (error) {
          console.warn(
            'Error calling downbeats callback with cached data',
            error
          );
        }
        return noop;
      }

      const existingStream = activeStreams.get(clipId);

      if (existingStream) {
        // Stream already exists, just register the callback and send existing data
        if (callbackFunction) {
          existingStream.callbacks.add(callbackFunction);
          // Send all existing data to the new callback
          try {
            callbackFunction({
              downbeats: existingStream.accumulatedDownbeats,
              secondsMissingFromStart: existingStream.secondsMissingFromStart,
              final: existingStream.final,
            });
          } catch (error) {
            console.warn(
              'Error calling downbeats callback with existing data',
              error
            );
          }

          // Return cleanup function
          return () => {
            const stream = activeStreams.get(clipId);
            if (stream) {
              stream.callbacks.delete(callbackFunction);
            }
          };
        }
        return noop;
      }

      // Create new stream state
      const streamState: StreamState = {
        callbacks: new Set(callbackFunction ? [callbackFunction] : []),
        accumulatedDownbeats: [],
        final: false,
        secondsMissingFromStart: 0,
      };

      activeStreams.set(clipId, streamState);

      // Start the stream asynchronously
      startStream(clipId, streamState).catch((error) => {
        console.error('Error starting downbeats stream', error);
        activeStreams.delete(clipId);
      });

      // Return cleanup function if callback was provided
      if (callbackFunction) {
        return () => {
          const stream = activeStreams.get(clipId);
          if (stream) {
            stream.callbacks.delete(callbackFunction);
          }
        };
      } else {
        return noop;
      }
    },
    [startStream]
  );

  return {
    streamDownbeatsForClip,
  };
};

export const useClipDownbeats = (
  clipId?: string | null,
  finalOnly: boolean = false
) => {
  const { streamDownbeatsForClip } = useStreamingDownbeats();
  const currentClipId = useRef(clipId);
  currentClipId.current = clipId;

  const [resultAndClipId, setResultAndClipId] = useState<{
    result: {
      downbeats: [number, number][];
      secondsMissingFromStart: number;
      final: boolean;
    };
    clipId: string;
  } | null>(null);

  useEffect(() => {
    if (!clipId) return;
    return streamDownbeatsForClip(
      clipId,
      ({ downbeats, secondsMissingFromStart, final }) => {
        if (currentClipId.current !== clipId) return;
        setResultAndClipId({
          result: { downbeats, secondsMissingFromStart, final },
          clipId,
        });
      }
    );
  }, [streamDownbeatsForClip, clipId]);

  const result = useMemo(() => {
    if (!resultAndClipId) return null;
    if (clipId !== resultAndClipId.clipId) return null;
    if (finalOnly && !resultAndClipId.result.final) return null;
    return resultAndClipId.result;
  }, [resultAndClipId, clipId, finalOnly]);

  return result;
};
