import { createHash, randomUUID } from "node:crypto";
import { createWriteStream } from "node:fs";
import { mkdir } from "node:fs/promises";
import { join } from "node:path";
import { Readable, Transform } from "node:stream";
import { pipeline } from "node:stream/promises";
import { fetchWithRetry, HttpError } from "./http.js";
import type {
  ArchiveAudit,
  ArchiveSource,
  ArchiveTarget,
  ArchiveWorkItem,
  SalesforceAccessToken,
  TokenProvider,
} from "./types.js";

type Fetch = typeof fetch;

interface DiscoveryOptions {
  apiVersion: string;
  allowedObjects: string[];
  contentVersionRetentionMonths: number;
  attachmentRetentionMonths: number;
}

interface QueryResponse<T> {
  records: T[];
  done: boolean;
  nextRecordsUrl?: string;
}

interface ContentVersionRow {
  Id: string;
  Title: string;
  FileExtension: string | null;
  ContentSize: number;
  ContentDocumentId: string;
  IsLatest?: boolean;
}

interface ContentVersionStateRow {
  Id: string;
  ContentSize: number;
  IsLatest: boolean;
}

interface AttachmentRow {
  Id: string;
  Name: string;
  BodyLength: number;
  ParentId: string;
}

interface ContentDocumentLinkRow {
  ContentDocumentId: string;
  LinkedEntityId: string;
}

interface EmailMessageRow {
  Id: string;
  ParentId: string | null;
  RelatedToId: string | null;
}

interface GlobalDescribeResponse {
  sobjects: Array<{ name: string; keyPrefix: string | null }>;
}

export interface DownloadResult {
  path: string;
  bytes: number;
}

function escapeSoql(value: string): string {
  return value.replace(/\\/g, "\\\\").replace(/'/g, "\\'");
}

function soqlValues(values: string[]): string {
  return values.map((value) => `'${escapeSoql(value)}'`).join(",");
}

function cutoffDate(months: number): string {
  const cutoff = new Date();
  cutoff.setUTCMonth(cutoff.getUTCMonth() - months);
  return cutoff.toISOString();
}

function chunks<T>(values: T[], size: number): T[][] {
  const output: T[][] = [];
  for (let offset = 0; offset < values.length; offset += size) output.push(values.slice(offset, offset + size));
  return output;
}

function truncate(value: string, maxLength: number): string {
  return value.length <= maxLength ? value : value.slice(0, maxLength);
}

function sanitizeFileName(rawName: string): string {
  const cleaned = (rawName.trim() || "unnamed_file")
    .replace(/[\\/:*?"<>|#%&{}~]/g, "_")
    .replace(/[. ]+$/g, "_");
  return truncate(cleaned, 255);
}

function sanitizeFolderName(rawName: string): string {
  const cleaned = (rawName.trim() || "UnknownObject")
    .replace(/[\\/:*?"<>|#%&{}~ ]/g, "_")
    .replace(/[.]+$/g, "_");
  return truncate(cleaned, 100);
}

function contentVersionFileName(title: string, extension: string | null): string {
  const cleanTitle = sanitizeFileName(title);
  const cleanExtension = extension ? sanitizeFileName(extension).toLowerCase() : "";
  const suffix = cleanExtension ? `.${cleanExtension}` : "";
  if (suffix && cleanTitle.toLowerCase().endsWith(suffix.toLowerCase())) return truncate(cleanTitle, 255);
  return `${truncate(cleanTitle, 255 - suffix.length)}${suffix}`;
}

export function archiveKey(sourceType: string, sourceRecordId: string, parentRecordId: string): string {
  return createHash("sha256").update(`${sourceType}|${sourceRecordId}|${parentRecordId}`).digest("hex");
}

export class SalesforceClient {
  private keyPrefixToType?: Map<string, string>;

  constructor(
    private readonly tokens: TokenProvider<SalesforceAccessToken>,
    private readonly tempDir: string,
    private readonly options: DiscoveryOptions,
    private readonly fetchImpl: Fetch = fetch,
  ) {}

  async discoverWork(): Promise<ArchiveWorkItem[]> {
    const [contentVersions, attachments] = await Promise.all([
      this.discoverContentVersions(),
      this.discoverAttachments(),
    ]);
    return [...contentVersions, ...attachments];
  }

  async verifyArchiveSchema(): Promise<void> {
    await this.queryAll<ArchiveAudit>(
      "SELECT Archive_Key__c,Archive_Status__c,Archived_On__c,"
      + "Source_Type__c,Source_Record_Id__c,Content_Document_Id__c,"
      + "Parent_Record_Id__c,Parent_Object_API_Name__c,File_Name__c,"
      + "File_Size_Bytes__c,Uploaded_Size_Bytes__c,SharePoint_Path__c,"
      + "SharePoint_File_Url__c,Graph_Item_Id__c,Error_Message__c "
      + "FROM SharePoint_Archived_File__c LIMIT 1",
    );
  }

  async discoverContentVersions(): Promise<ArchiveWorkItem[]> {
    const allowed = soqlValues(this.options.allowedObjects);
    const versions = await this.queryAll<ContentVersionRow>(
      `SELECT Id,Title,FileExtension,ContentSize,ContentDocumentId,IsLatest FROM ContentVersion `
      + `WHERE IsLatest=true AND ContentLocation='S' `
      + `AND CreatedDate < ${cutoffDate(this.options.contentVersionRetentionMonths)} `
      + `AND ContentDocumentId IN (SELECT ContentDocumentId FROM ContentDocumentLink `
      + `WHERE LinkedEntity.Type IN (${allowed}))`,
    );
    if (versions.length === 0) return [];
    const links = await this.contentLinks(versions.map((row) => row.ContentDocumentId));
    const linksByDocument = new Map<string, ContentDocumentLinkRow[]>();
    for (const link of links) {
      const grouped = linksByDocument.get(link.ContentDocumentId) ?? [];
      grouped.push(link);
      linksByDocument.set(link.ContentDocumentId, grouped);
    }
    const prefixMap = await this.objectTypesByPrefix();
    return versions.flatMap((version) => {
      const source = this.contentSource(version);
      const targets = (linksByDocument.get(version.ContentDocumentId) ?? [])
        .map((link) => this.targetFor(source, link.LinkedEntityId, prefixMap.get(link.LinkedEntityId.slice(0, 3))))
        .filter((target): target is ArchiveTarget => target !== null);
      return targets.length ? [{ source, targets }] : [];
    });
  }

  async discoverAttachments(): Promise<ArchiveWorkItem[]> {
    const allowed = soqlValues(this.options.allowedObjects);
    const attachments = await this.queryAll<AttachmentRow>(
      `SELECT Id,Name,BodyLength,ParentId FROM Attachment `
      + `WHERE CreatedDate < ${cutoffDate(this.options.attachmentRetentionMonths)} `
      + `AND (Parent.Type IN (${allowed}) OR Parent.Type='EmailMessage')`,
    );
    if (attachments.length === 0) return [];
    const prefixMap = await this.objectTypesByPrefix();
    const emailIds = attachments
      .filter((row) => prefixMap.get(row.ParentId.slice(0, 3)) === "EmailMessage")
      .map((row) => row.ParentId);
    const emails = new Map<string, EmailMessageRow>();
    for (const group of chunks(emailIds, 100)) {
      for (const email of await this.queryAll<EmailMessageRow>(
        `SELECT Id,ParentId,RelatedToId FROM EmailMessage WHERE Id IN (${soqlValues(group)})`,
      )) emails.set(email.Id, email);
    }

    return attachments.flatMap((attachment) => {
      let parentId: string | null = attachment.ParentId;
      let parentType = prefixMap.get(parentId.slice(0, 3));
      if (parentType === "EmailMessage") {
        const email = emails.get(parentId);
        const candidates = [email?.ParentId, email?.RelatedToId].filter((id): id is string => Boolean(id));
        parentId = candidates.find((id) => this.options.allowedObjects.includes(prefixMap.get(id.slice(0, 3)) ?? "")) ?? null;
        parentType = parentId ? prefixMap.get(parentId.slice(0, 3)) : undefined;
      }
      if (!parentId || !parentType || !this.options.allowedObjects.includes(parentType)) return [];
      const source: ArchiveSource = {
        sourceType: "Attachment",
        sourceRecordId: attachment.Id,
        contentDocumentId: null,
        fileName: sanitizeFileName(attachment.Name),
        fileSizeBytes: Number(attachment.BodyLength),
        downloadPath: this.sobjectBlobPath("Attachment", attachment.Id, "Body"),
      };
      const target = this.targetFor(source, parentId, parentType);
      return target ? [{ source, targets: [target] }] : [];
    });
  }

  async currentContentTargets(source: ArchiveSource): Promise<ArchiveTarget[]> {
    if (!source.contentDocumentId) return [];
    const prefixMap = await this.objectTypesByPrefix();
    return (await this.contentLinks([source.contentDocumentId]))
      .map((link) => this.targetFor(source, link.LinkedEntityId, prefixMap.get(link.LinkedEntityId.slice(0, 3))))
      .filter((target): target is ArchiveTarget => target !== null);
  }

  async isLatestContentVersion(source: ArchiveSource): Promise<boolean> {
    const rows = await this.queryAll<ContentVersionStateRow>(
      `SELECT Id,ContentSize,IsLatest FROM ContentVersion WHERE Id='${escapeSoql(source.sourceRecordId)}' LIMIT 1`,
    );
    const row = rows[0];
    return Boolean(row && row.IsLatest === true && Number(row.ContentSize) === source.fileSizeBytes);
  }

  async successfulAudits(keys: string[]): Promise<Map<string, ArchiveAudit>> {
    const audits = new Map<string, ArchiveAudit>();
    for (const group of chunks([...new Set(keys)], 100)) {
      if (group.length === 0) continue;
      const records = await this.queryAll<ArchiveAudit>(
        `SELECT Archive_Key__c,Archive_Status__c,SharePoint_File_Url__c,Graph_Item_Id__c `
        + `FROM SharePoint_Archived_File__c WHERE Archive_Status__c='Success' `
        + `AND Archive_Key__c IN (${soqlValues(group)})`,
      );
      for (const record of records) audits.set(record.Archive_Key__c, record);
    }
    return audits;
  }

  async upsertSuccessAudit(
    source: ArchiveSource,
    target: ArchiveTarget,
    graphItemId: string,
    webUrl: string,
  ): Promise<void> {
    await this.authorizedRequest(
      `${this.dataPath()}/sobjects/SharePoint_Archived_File__c/Archive_Key__c/${encodeURIComponent(target.archiveKey)}`,
      () => ({
        method: "PATCH",
        headers: { "Content-Type": "application/json" },
        body: JSON.stringify({
          Archive_Status__c: "Success",
          Archived_On__c: new Date().toISOString(),
          Source_Type__c: source.sourceType,
          Source_Record_Id__c: source.sourceRecordId,
          Content_Document_Id__c: source.contentDocumentId,
          Parent_Record_Id__c: target.parentRecordId,
          Parent_Object_API_Name__c: target.parentObjectApiName,
          File_Name__c: source.fileName,
          File_Size_Bytes__c: source.fileSizeBytes,
          Uploaded_Size_Bytes__c: source.fileSizeBytes,
          SharePoint_Path__c: target.sharePointPath,
          SharePoint_File_Url__c: webUrl,
          Graph_Item_Id__c: graphItemId,
          Error_Message__c: null,
        }),
      }),
    );
  }

  async download(source: ArchiveSource): Promise<DownloadResult> {
    await mkdir(this.tempDir, { recursive: true, mode: 0o700 });
    const tempPath = join(this.tempDir, `${source.sourceRecordId}-${randomUUID()}.part`);
    const response = await this.authorizedRequest(source.downloadPath, () => ({ method: "GET" }));
    if (!response.body) throw new Error("Salesforce download response had no body stream.");
    let bytes = 0;
    const counter = new Transform({
      transform(chunk: Buffer, _encoding, callback) {
        bytes += chunk.length;
        callback(null, chunk);
      },
    });
    await pipeline(
      Readable.fromWeb(response.body as never),
      counter,
      createWriteStream(tempPath, { flags: "wx", mode: 0o600 }),
    );
    return { path: tempPath, bytes };
  }

  async deleteSource(source: ArchiveSource): Promise<void> {
    const objectName = source.sourceType === "ContentVersion" ? "ContentDocument" : "Attachment";
    const recordId = source.sourceType === "ContentVersion" ? source.contentDocumentId : source.sourceRecordId;
    if (!recordId) throw new SourceValidationError("ContentDocument id is missing from archive source metadata.");
    await this.authorizedRequest(
      `${this.dataPath()}/sobjects/${objectName}/${recordId}`,
      () => ({ method: "DELETE" }),
    );
  }

  private contentSource(version: ContentVersionRow): ArchiveSource {
    return {
      sourceType: "ContentVersion",
      sourceRecordId: version.Id,
      contentDocumentId: version.ContentDocumentId,
      fileName: contentVersionFileName(version.Title, version.FileExtension),
      fileSizeBytes: Number(version.ContentSize),
      downloadPath: this.sobjectBlobPath("ContentVersion", version.Id, "VersionData"),
    };
  }

  private targetFor(
    source: ArchiveSource,
    parentRecordId: string,
    parentObjectApiName: string | undefined,
  ): ArchiveTarget | null {
    if (!parentObjectApiName || !this.options.allowedObjects.includes(parentObjectApiName)) return null;
    const key = archiveKey(source.sourceType, source.sourceRecordId, parentRecordId);
    const sharePointPath = [
      sanitizeFolderName(parentObjectApiName),
      parentRecordId,
      source.sourceRecordId,
      source.fileName,
    ].join("/");
    return { archiveKey: key, parentRecordId, parentObjectApiName, sharePointPath };
  }

  private async contentLinks(contentDocumentIds: string[]): Promise<ContentDocumentLinkRow[]> {
    const allowed = soqlValues(this.options.allowedObjects);
    const links: ContentDocumentLinkRow[] = [];
    for (const group of chunks([...new Set(contentDocumentIds)], 100)) {
      links.push(...await this.queryAll<ContentDocumentLinkRow>(
        `SELECT ContentDocumentId,LinkedEntityId FROM ContentDocumentLink `
        + `WHERE ContentDocumentId IN (${soqlValues(group)}) AND LinkedEntity.Type IN (${allowed})`,
      ));
    }
    return links;
  }

  private async objectTypesByPrefix(): Promise<Map<string, string>> {
    if (this.keyPrefixToType) return this.keyPrefixToType;
    const response = await this.authorizedRequest(`${this.dataPath()}/sobjects`, () => ({ method: "GET" }));
    const describe = (await response.json()) as GlobalDescribeResponse;
    this.keyPrefixToType = new Map(
      describe.sobjects
        .filter((item): item is { name: string; keyPrefix: string } => Boolean(item.keyPrefix))
        .map((item) => [item.keyPrefix, item.name]),
    );
    return this.keyPrefixToType;
  }

  private async queryAll<T>(soql: string): Promise<T[]> {
    const output: T[] = [];
    let path: string | undefined = `${this.dataPath()}/query?q=${encodeURIComponent(soql)}`;
    while (path) {
      const response = await this.authorizedRequest(path, () => ({ method: "GET" }));
      const page = (await response.json()) as QueryResponse<T>;
      output.push(...page.records);
      path = page.done ? undefined : page.nextRecordsUrl;
    }
    return output;
  }

  private sobjectBlobPath(objectName: string, recordId: string, fieldName: string): string {
    return `${this.dataPath()}/sobjects/${objectName}/${recordId}/${fieldName}`;
  }

  private dataPath(): string {
    return `/services/data/v${this.options.apiVersion}`;
  }

  private async authorizedRequest(path: string, init: () => RequestInit): Promise<Response> {
    return fetchWithRetry(async (attempt) => {
      const token = await this.tokens.getToken(attempt > 0);
      const request = init();
      const headers = new Headers(request.headers);
      headers.set("Authorization", `Bearer ${token.accessToken}`);
      return this.fetchImpl(`${token.instanceUrl}${path}`, { ...request, headers });
    });
  }
}

export function isRetryableTransferError(error: unknown): boolean {
  if (error instanceof HttpError) return error.retryable;
  if (error instanceof SyntaxError || error instanceof TypeError) return true;
  return !(error instanceof SourceValidationError);
}

export class SourceValidationError extends Error {
  constructor(message: string) {
    super(message);
    this.name = "SourceValidationError";
  }
}
