import { unlink } from "node:fs/promises";
import { FailureStore } from "./failure-store.js";
import { HttpError, sleep } from "./http.js";
import { GraphClient } from "./graph.js";
import { isRetryableTransferError, SalesforceClient, SourceValidationError } from "./salesforce.js";
import type { ArchiveAudit, ArchiveWorkItem } from "./types.js";

export interface WorkerOptions {
  workerId: string;
  concurrency: number;
  scanIntervalMs: number;
  deleteAfterArchive: boolean;
}

function log(level: "info" | "error", message: string, fields: Record<string, unknown> = {}): void {
  const output = JSON.stringify({ timestamp: new Date().toISOString(), level, message, ...fields });
  (level === "error" ? console.error : console.log)(output);
}

function errorMessage(error: unknown): string {
  if (error instanceof HttpError) return `${error.message} ${error.responseBody}`.trim().slice(0, 4000);
  return error instanceof Error ? error.message.slice(0, 4000) : String(error).slice(0, 4000);
}

export class ArchiveWorker {
  private stopping = false;

  constructor(
    private readonly salesforce: SalesforceClient,
    private readonly graph: GraphClient,
    private readonly failures: FailureStore,
    private readonly options: WorkerOptions,
  ) {}

  stop(): void {
    this.stopping = true;
  }

  async run(): Promise<void> {
    await this.failures.load();
    log("info", "SharePoint archive worker started.", {
      workerId: this.options.workerId,
      concurrency: this.options.concurrency,
      deleteAfterArchive: this.options.deleteAfterArchive,
    });
    while (!this.stopping) {
      try {
        await this.salesforce.verifyArchiveSchema();
        const work = (await this.salesforce.discoverWork())
          .filter((item) => this.failures.canRun(this.failureKey(item)));
        log("info", "Salesforce archive scan completed.", { discoveredSources: work.length });
        await this.processPool(work);
      } catch (error) {
        log("error", "Salesforce archive scan failed.", { error: errorMessage(error) });
      }
      if (!this.stopping) await sleep(this.options.scanIntervalMs);
    }
    log("info", "SharePoint archive worker stopped.", { workerId: this.options.workerId });
  }

  private async processPool(items: ArchiveWorkItem[]): Promise<void> {
    let nextIndex = 0;
    const slots = Array.from({ length: this.options.concurrency }, async (_, slotIndex) => {
      while (!this.stopping) {
        const item = items[nextIndex];
        nextIndex += 1;
        if (!item) return;
        await this.process(item, slotIndex + 1);
      }
    });
    await Promise.all(slots);
  }

  private async process(item: ArchiveWorkItem, slot: number): Promise<void> {
    const failureKey = this.failureKey(item);
    let tempPath: string | undefined;
    try {
      const existingAudits = await this.salesforce.successfulAudits(
        item.targets.map((target) => target.archiveKey),
      );
      const missingTargets = item.targets.filter((target) => !existingAudits.has(target.archiveKey));
      if (missingTargets.length > 0) {
        const downloaded = await this.salesforce.download(item.source);
        tempPath = downloaded.path;
        if (downloaded.bytes !== item.source.fileSizeBytes) {
          throw new SourceValidationError(
            `Downloaded ${downloaded.bytes} bytes but Salesforce expected ${item.source.fileSizeBytes}.`,
          );
        }
        for (const target of missingTargets) {
          const driveItem = await this.graph.upload(tempPath, downloaded.bytes, target.sharePointPath);
          if (Number(driveItem.size) !== item.source.fileSizeBytes) {
            throw new SourceValidationError(
              `SharePoint reported ${driveItem.size} bytes but Salesforce expected ${item.source.fileSizeBytes}.`,
            );
          }
          await this.salesforce.upsertSuccessAudit(
            item.source,
            target,
            driveItem.id,
            driveItem.webUrl,
          );
          existingAudits.set(target.archiveKey, {
            Archive_Key__c: target.archiveKey,
            Archive_Status__c: "Success",
            SharePoint_File_Url__c: driveItem.webUrl,
            Graph_Item_Id__c: driveItem.id,
          });
          log("info", "SharePoint destination archived.", {
            slot,
            sourceType: item.source.sourceType,
            sourceRecordId: item.source.sourceRecordId,
            parentRecordId: target.parentRecordId,
            bytes: item.source.fileSizeBytes,
          });
        }
      }

      const safeToDelete = await this.readyForDeletion(item, existingAudits);
      if (safeToDelete && this.options.deleteAfterArchive) {
        await this.salesforce.deleteSource(item.source);
        log("info", "Verified Salesforce source deleted.", {
          slot,
          sourceType: item.source.sourceType,
          sourceRecordId: item.source.sourceRecordId,
        });
      }
      await this.failures.clear(failureKey);
    } catch (error) {
      const retryable = isRetryableTransferError(error);
      await this.failures.recordFailure(failureKey, errorMessage(error), retryable);
      log("error", "Archive source processing failed; source was retained.", {
        slot,
        sourceType: item.source.sourceType,
        sourceRecordId: item.source.sourceRecordId,
        retryable,
        error: errorMessage(error),
      });
    } finally {
      if (tempPath) {
        await unlink(tempPath).catch((error: unknown) => {
          log("error", "Unable to remove temporary archive file.", {
            sourceRecordId: item.source.sourceRecordId,
            tempPath,
            error: errorMessage(error),
          });
        });
      }
    }
  }

  private async readyForDeletion(
    item: ArchiveWorkItem,
    knownAudits: Map<string, ArchiveAudit>,
  ): Promise<boolean> {
    if (item.source.sourceType === "Attachment") {
      return item.targets.every((target) => knownAudits.has(target.archiveKey));
    }
    if (!(await this.salesforce.isLatestContentVersion(item.source))) return false;
    const currentTargets = await this.salesforce.currentContentTargets(item.source);
    const currentAudits = await this.salesforce.successfulAudits(
      currentTargets.map((target) => target.archiveKey),
    );
    return currentTargets.length > 0
      && currentTargets.every((target) => currentAudits.has(target.archiveKey));
  }

  private failureKey(item: ArchiveWorkItem): string {
    return `${item.source.sourceType}:${item.source.sourceRecordId}`;
  }
}
