import { synchronize } from '@nozbe/watermelondb/sync';
import { apiFetch } from '@/modules/core/http/client';
import { getDatabase } from './db';
import { Outbox } from './db/models/Outbox';
import { syncProgress } from './syncProgress';

type PullResponse = {
    changes: any;
    timestamp: number;
    // Dump INITIAL paginé : le client boucle tant que `has_more`, en repassant `next_cursor`. `total` = nb de
    // conversations (pour la barre de progression). Un DELTA renvoie has_more=false, next_cursor=null, total=null.
    has_more?: boolean;
    next_cursor?: string | null;
    total?: number | null;
};

let inFlight: Promise<void> | null = null;

/**
 * Un pull = potentiellement PLUSIEURS pages (dump initial du 1ᵉʳ login, paginé par conversations côté serveur). On
 * BOUCLE `synchronize()` en repassant `cursor` (closure) et en ÉPINGLANT le timestamp de la page 1 : chaque page
 * commit sa propre `database.write` → les conversations apparaissent au fur et à mesure dans la liste réactive, et le
 * watermark final reste celui de la page 1 (le delta suivant repart proprement). Un DELTA fait une seule itération.
 */
async function runSync(): Promise<void> {
    let cursor: string | null = null;
    let pinnedTs: number | null = null;
    let hasMore = true;
    let done = 0;

    while (hasMore) {
        try {
            await synchronize({
                database: getDatabase(),
                // `sendCreatedAsUpdated` : on APPLIQUE les enregistrements en UPSERT (créer OU mettre à jour) au lieu
                // d'échouer sur « already exists ». Indispensable : `created_at` en précision SECONDE → un renvoi
                // (message d'une même seconde, ou recouvrement dû au watermark épinglé) est inoffensif.
                pullChanges: async ({ lastPulledAt }) => {
                    const params: string[] = [];
                    if (lastPulledAt != null) params.push(`last_pulled_at=${lastPulledAt}`);
                    if (cursor) params.push(`cursor=${encodeURIComponent(cursor)}`);
                    const qs = params.length ? `?${params.join('&')}` : '';
                    const res = await apiFetch<PullResponse>(`/api/v1/chat-sync/pull${qs}`, { method: 'GET' });
                    hasMore = !!res.has_more;
                    cursor = res.next_cursor ?? null;
                    if (pinnedTs == null) pinnedTs = res.timestamp;
                    // Progression : barre UNIQUEMENT pour un dump initial MULTI-pages (page 1 avec has_more, ou déjà
                    // en cours). Delta / page unique → `total` absent ou has_more=false d'emblée → pas de barre.
                    const total = res.total ?? 0;
                    const pageConv = res?.changes?.conversations?.updated?.length ?? 0;
                    done += pageConv;
                    // DIAGNOSTIC (dev) : ce que renvoie réellement le pull, avant application. Permet de distinguer
                    // « 200 vide » (counts 0) de « 200 rempli mais non appliqué » (counts > 0 → l'erreur vient de
                    // synchronize, capturée dans le catch ci-dessous). À retirer une fois la cause corrigée.
                    if (__DEV__) {
                        const c: any = res?.changes ?? {};
                        const cnt = (t: string) => ((c?.[t]?.created?.length ?? 0) + (c?.[t]?.updated?.length ?? 0) + (c?.[t]?.deleted?.length ?? 0));
                        console.log('[chat-sync] pull', {
                            total, has_more: res.has_more, next_cursor: res.next_cursor ? '…' : null,
                            tables: Object.keys(c), conversations: cnt('conversations'), messages: cnt('messages'), calls: cnt('calls'),
                        });
                    }
                    if (total > 0 && (hasMore || syncProgress.get().syncing)) {
                        syncProgress.set({ syncing: true, done: Math.min(done, total), total });
                    }
                    // Timestamp épinglé sur toutes les pages → watermark final = page 1.
                    return { changes: res.changes, timestamp: pinnedTs ?? res.timestamp };
                },
                pushChanges: async ({ changes }) => {
                    await apiFetch(`/api/v1/chat-sync/push`, {
                        method: 'POST',
                        body: JSON.stringify(changes),
                    });
                },
                sendCreatedAsUpdated: true,
            });
        } catch (e) {
            // Point AVEUGLE historique : l'erreur d'APPLICATION des changes (WatermelonDB) était totalement
            // silencieuse. On la logue (dev) puis on relance — la stratégie de récupération reste dans syncChat().
            if (__DEV__) console.warn('[chat-sync] échec application des changes (synchronize)', e);
            throw e;
        }
    }
}

/**
 * Synchronise le chat local (WatermelonDB) avec le canal dédié `chat-sync` (pull scopé à l'utilisateur).
 * - timestamp serveur opaque pour WatermelonDB (renvoyé tel quel au pull suivant).
 * - `pushChanges` : no-op backend en P1 ; l'écriture offline arrive en P3.
 *
 * Anti-concurrence : un seul `synchronize()` à la fois. Sans réseau, `apiFetch` jette `ApiError(0)` → la
 * sync échoue silencieusement (la base locale reste la source de vérité).
 *
 * Auto-réparation : si l'état de sync local devient incohérent (curseur perdu → le serveur renvoie en
 * `created` des enregistrements déjà présents), on réinitialise la base locale puis on re-pull tout.
 * L'outbox (écritures hors-ligne en attente) est PRÉSERVÉ via snapshot → reset → restore (sinon le reset
 * effaçait des messages en attente : perte silencieuse).
 */
export async function syncChat(): Promise<void> {
    if (inFlight) return inFlight;
    inFlight = (async () => {
        try {
            await runSync();
        } catch (e: any) {
            const msg = String(e?.message ?? e);
            if (__DEV__) console.warn('[chat-sync] syncChat a échoué —', msg);
            // Récupération réservée à la vraie corruption de l'état de sync : « already exists » = curseur perdu, le
            // serveur renvoie en `created` des enregistrements déjà présents. ATTENTION : unsafeResetDatabase() EFFACE
            // TOUTE la base ; on ne le déclenche donc PAS sur une erreur d'application quelconque — une « Diagnostic
            // error » de forme/colonne wiperait la base à chaque pull sans corriger la cause, et masquerait le vrai bug
            // (c'était le cas des motifs `cannot create|diagnostic`, retirés ici).
            if (/already exists/i.test(msg)) {
                if (__DEV__) console.warn('[chat-sync] curseur perdu → réinitialisation base + re-pull complet');
                const db = getDatabase();
                // Repartir d'une base propre + re-pull complet — SANS perdre les écritures hors-ligne en attente.
                // On snapshot l'outbox, on reset, puis on le restaure avant de re-synchroniser.
                const snapshot = (await db.get<Outbox>('outbox').query().fetch()).map((o) => ({
                    op: o.op,
                    conversationId: o.conversationId,
                    clientId: o.clientId,
                    kind: o.kind,
                    text: o.text,
                    replyToId: o.replyToId,
                    replyToMedia: o.replyToMedia,
                    fileUri: o.fileUri,
                    fileName: o.fileName,
                    fileType: o.fileType,
                    duration: o.duration,
                    emoji: o.emoji,
                    mentions: o.mentions,
                    status: o.status,
                    queuedAt: o.queuedAt,
                }));
                await db.write(async () => {
                    await db.unsafeResetDatabase();
                });
                if (snapshot.length) {
                    const outbox = db.get<Outbox>('outbox');
                    await db.write(async () => {
                        await Promise.all(
                            snapshot.map((it) =>
                                outbox.create((o) => {
                                    o.op = it.op;
                                    o.conversationId = it.conversationId;
                                    o.clientId = it.clientId;
                                    o.kind = it.kind;
                                    o.text = it.text;
                                    o.replyToId = it.replyToId;
                                    o.replyToMedia = it.replyToMedia;
                                    o.fileUri = it.fileUri;
                                    o.fileName = it.fileName;
                                    o.fileType = it.fileType;
                                    o.duration = it.duration;
                                    o.emoji = it.emoji;
                                    o.mentions = it.mentions;
                                    o.status = it.status;
                                    o.queuedAt = it.queuedAt;
                                }),
                            ),
                        );
                    });
                }
                await runSync();
            } else {
                throw e;
            }
        }
    })();
    try {
        await inFlight;
    } finally {
        inFlight = null;
        syncProgress.finish(); // masque la barre à la fin (succès OU échec) ; no-op si aucune n'était affichée
    }
}

// --- Demande de sync COALESCÉE ---------------------------------------------------------------------------
// Une même action utilisateur déclenche plusieurs demandes de pull (chemin d'écriture + écho temps réel reçu
// par le fil + écho reçu par la liste). On les fond en UN SEUL pull par rafale : debounce, puis si des
// demandes arrivent PENDANT un pull en cours, on en relance exactement un à la fin (aucun delta manqué).
let reqTimer: ReturnType<typeof setTimeout> | null = null;
let running = false;
let pendingAgain = false;

async function runRequested(): Promise<void> {
    if (running) {
        pendingAgain = true;
        return;
    }
    running = true;
    try {
        await syncChat();
    } catch {
        // best-effort (hors-ligne / serveur injoignable) — la base locale reste la source de vérité
    } finally {
        running = false;
    }
    if (pendingAgain) {
        pendingAgain = false;
        runRequested();
    }
}

/**
 * Déclencheur de sync à utiliser pour les changements pilotés par le temps réel / les actions : coalesce les
 * rafales en un seul pull incrémental. Évite « une sync par déclencheur » tout en gardant la persistance
 * fiable (y compris quand le temps réel est coupé : le chemin d'écriture demande quand même un pull).
 */
export function requestChatSync(delay = 300): void {
    if (reqTimer) clearTimeout(reqTimer);
    reqTimer = setTimeout(() => {
        reqTimer = null;
        runRequested();
    }, delay);
}
