import type { StandardSchemaV1 } from "@standard-schema/spec";
import { runWithResolvedTracingContext } from "../tracing/execution.js";
import {
  parseTraceCarrier,
  type TraceCarrier,
  type TracingPort,
} from "../tracing/index.js";
import {
  isEventPayloadParsed,
  isEventPayloadTransportStable,
  markEventPayloadTransportStable,
} from "./payload-state.js";
import {
  EventTransportError,
  type EventTransportValue,
  eventTransportValuesEqual,
  toEventTransportValue,
} from "./transport.js";

export {
  EventTransportError,
  type EventTransportErrorReason,
  type EventTransportValue,
} from "./transport.js";

/**
 * Any Standard Schema compatible validator.
 */
export type StandardSchema = StandardSchemaV1<unknown, unknown>;

/**
 * Value or promise of that value.
 */
export type MaybePromise<T> = T | Promise<T>;

const DEFAULT_LISTENER_READY_TIMEOUT_MS = 10_000;
const MAX_TIMER_MS = 2_147_483_647;

/**
 * Infer the parsed output type from a Standard Schema.
 */
export type InferSchemaOutput<T extends StandardSchemaV1> =
  StandardSchemaV1.InferOutput<T>;

/**
 * Minimal event definition shape accepted by event bus helpers.
 */
export interface EventPayloadDef<
  Name extends string = string,
  Payload extends StandardSchema = StandardSchema,
> {
  /**
   * Stable event name.
   */
  readonly name: Name;
  /**
   * Standard Schema payload validator.
   */
  readonly payload: Payload;
  /**
   * Optional human-readable description for docs and tooling.
   */
  readonly description?: string;
}

/**
 * Event definition created by `defineEvent(...)`.
 */
export interface EventDef<
  Name extends string = string,
  Payload extends StandardSchema = StandardSchema,
> extends EventPayloadDef<Name, Payload> {
  /**
   * Discriminator for event definitions.
   */
  readonly kind: "event";
}

/**
 * Infer the parsed payload type for an event definition.
 */
export type InferEventPayload<E extends EventPayloadDef> =
  E["payload"] extends StandardSchemaV1<unknown, infer Output> ? Output : never;

/** Metadata propagated with an event delivery. */
export interface EventPublishOptions {
  /** Versioned trace context captured by the event producer. */
  trace?: TraceCarrier;
}

/**
 * Lifecycle handle for one event subscription or a composed listener
 * registration.
 *
 * `ready` proves initial transport readiness. It does not represent ongoing
 * connectivity, durability, replay, or handler success after startup.
 */
export interface EventSubscription {
  /** Resolves when the subscription can receive events. */
  readonly ready: Promise<void>;
  /** Stop local delivery and await transport cleanup. Idempotent. */
  unsubscribe(): Promise<void>;
}

/** Error used when a subscription closes before initial readiness. */
export class EventSubscriptionClosedError extends Error {
  constructor(message = "Event subscription closed before it became ready.") {
    super(message);
    this.name = "EventSubscriptionClosedError";
  }
}

/** Error thrown when a listener registry misses its readiness deadline. */
export class ListenerRegistrationTimeoutError extends Error {
  /** Configured readiness deadline in milliseconds. */
  readonly timeoutMs: number;
  /** Listener names that were part of the registration. */
  readonly listenerNames: readonly string[];

  constructor(args: {
    timeoutMs: number;
    listenerNames: readonly string[];
  }) {
    const suffix = args.listenerNames.length
      ? `: ${args.listenerNames.join(", ")}`
      : "";
    super(
      `Listeners did not become ready within ${args.timeoutMs}ms${suffix}.`,
    );
    this.name = "ListenerRegistrationTimeoutError";
    this.timeoutMs = args.timeoutMs;
    this.listenerNames = [...args.listenerNames];
  }
}

/** Error thrown when listener rollback misses the registration deadline. */
export class ListenerRegistrationCleanupTimeoutError extends Error {
  /** Configured registration deadline in milliseconds. */
  readonly timeoutMs: number;
  /** Listener names that were part of the registration. */
  readonly listenerNames: readonly string[];

  constructor(args: {
    timeoutMs: number;
    listenerNames: readonly string[];
  }) {
    const suffix = args.listenerNames.length
      ? `: ${args.listenerNames.join(", ")}`
      : "";
    super(
      `Listener cleanup did not finish before the ${args.timeoutMs}ms registration deadline${suffix}.`,
    );
    this.name = "ListenerRegistrationCleanupTimeoutError";
    this.timeoutMs = args.timeoutMs;
    this.listenerNames = [...args.listenerNames];
  }
}

/**
 * Options for `defineEvent(...)`.
 */
export interface DefineEventOptions<Payload extends StandardSchema> {
  /**
   * Standard Schema payload validator.
   */
  payload: Payload;
  /**
   * Optional human-readable description for docs and tooling.
   */
  description?: string;
}

/**
 * Arguments passed to a listener handler.
 */
export interface ListenerHandleArgs<E extends EventDef, Ctx> {
  /**
   * Event definition being handled.
   */
  event: E;
  /**
   * Parsed event payload.
   */
  payload: InferEventPayload<E>;
  /**
   * Listener context.
   */
  ctx: Ctx;
}

/**
 * Listener definition created by `defineListener(...)`.
 */
export interface ListenerDef<
  E extends EventDef = EventDef,
  Ctx = unknown,
  Name extends string = string,
> {
  /**
   * Discriminator for listener definitions.
   */
  readonly kind: "listener";
  /**
   * Stable listener name.
   */
  readonly name: Name;
  /**
   * Event this listener handles.
   */
  readonly event: E;
  /**
   * Handle a parsed event payload.
   */
  handle(args: ListenerHandleArgs<E, Ctx>): MaybePromise<void>;
}

/**
 * Options for `defineListener(...)`.
 */
export interface DefineListenerOptions<E extends EventDef, Ctx> {
  /**
   * Event this listener handles.
   */
  event: E;
  /**
   * Handle a parsed event payload.
   */
  handle(args: ListenerHandleArgs<E, Ctx>): MaybePromise<void>;
}

/**
 * Event bus shape required by Beignet listener registration helpers.
 */
export interface EventBusLike {
  /**
   * Publish an event payload.
   */
  publish<E extends EventPayloadDef>(
    event: E,
    payload: InferEventPayload<E>,
    options?: EventPublishOptions,
  ): MaybePromise<void>;
  /**
   * Subscribe to an event and return its readiness and cleanup handle.
   */
  subscribe<E extends EventPayloadDef>(
    event: E,
    handler: (
      payload: InferEventPayload<E>,
      options?: EventPublishOptions,
    ) => MaybePromise<void>,
  ): EventSubscription;
}

/**
 * Options for `registerListeners(...)`.
 */
export interface RegisterListenersOptions<Ctx> {
  /**
   * Static listener context or factory evaluated for each delivered event.
   */
  ctx?: Ctx | (() => MaybePromise<Ctx>);
  /**
   * Runtime tracing port used to start the listener span before a lazy context
   * factory runs.
   */
  tracing?: TracingPort;
  /**
   * Called when a listener fails. When omitted, listener errors are rethrown to
   * the event bus subscription callback.
   */
  onError?: (error: unknown, listener: ListenerDef<EventDef, Ctx>) => void;
  /**
   * Maximum time for the complete listener registry to become ready. The same
   * registration deadline bounds automatic rollback after startup failure.
   * Defaults to 10 seconds.
   */
  readyTimeoutMs?: number;
}

/**
 * Context-bound listener helper factory.
 */
export interface Listeners<Ctx> {
  /**
   * Define a listener with the bound context type.
   */
  defineListener<Name extends string, E extends EventDef>(
    name: Name,
    options: DefineListenerOptions<E, Ctx>,
  ): ListenerDef<E, Ctx, Name>;
}

/**
 * Error thrown when event payload validation fails.
 */
export class EventValidationError extends Error {
  /**
   * Raw Standard Schema validation issues.
   */
  readonly issues: readonly StandardSchemaV1.Issue[];

  constructor(args: {
    name: string;
    issues: readonly StandardSchemaV1.Issue[];
  }) {
    super(
      `Event "${args.name}" payload validation failed: ${formatIssues(args.issues)}`,
    );
    this.name = "EventValidationError";
    this.issues = args.issues;
  }
}

/**
 * Parsed runtime value and canonical JSON value for one event publication.
 */
export interface PreparedEventPayload<E extends EventPayloadDef> {
  /** Parsed Standard Schema output delivered to in-process listeners. */
  readonly payload: InferEventPayload<E>;
  /** Canonical JSON value written by serialized transports. */
  readonly transportValue: EventTransportValue;
  /** Complete publish metadata to forward to in-process subscribers. */
  readonly publishOptions: EventPublishOptions;
}

function formatPath(path: StandardSchemaV1.Issue["path"]): string {
  if (!path?.length) return "";

  return path
    .map((segment) =>
      typeof segment === "object" && segment !== null && "key" in segment
        ? String(segment.key)
        : String(segment),
    )
    .join(".");
}

function formatIssues(issues: readonly StandardSchemaV1.Issue[]): string {
  return issues
    .map((issue) => {
      const path = formatPath(issue.path);
      return path ? `${path}: ${issue.message}` : issue.message;
    })
    .join("; ");
}

async function parsePayload<Schema extends StandardSchemaV1>(
  schema: Schema,
  input: unknown,
  args: { name: string },
): Promise<InferSchemaOutput<Schema>> {
  const result = await schema["~standard"].validate(input);

  if (result.issues?.length) {
    throw new EventValidationError({
      name: args.name,
      issues: result.issues,
    });
  }

  if ("value" in result) {
    return result.value as InferSchemaOutput<Schema>;
  }

  throw new Error("Invalid Standard Schema result: missing value");
}

/**
 * Define a typed event.
 *
 * Event payloads are validated before publishing through `publishEvent(...)`
 * and before registered listeners run. Producer helpers also require parsed
 * output to be plain JSON that remains unchanged when validated again after a
 * transport round trip.
 */
export function defineEvent<
  Name extends string,
  Payload extends StandardSchema,
>(name: Name, options: DefineEventOptions<Payload>): EventDef<Name, Payload> {
  return {
    kind: "event",
    name,
    payload: options.payload,
    description: options.description,
  };
}

function defineListenerImpl<
  Name extends string = string,
  E extends EventDef = EventDef,
  Ctx = unknown,
>(
  name: Name,
  options: DefineListenerOptions<E, Ctx>,
): ListenerDef<E, Ctx, Name> {
  return {
    kind: "listener",
    name,
    event: options.event,
    handle: options.handle,
  };
}

/**
 * Validate and parse an event payload with the event's Standard Schema.
 */
export async function parseEventPayload<E extends EventPayloadDef>(
  event: E,
  payload: unknown,
): Promise<InferEventPayload<E>> {
  return (await parsePayload(event.payload, payload, {
    name: event.name,
  })) as InferEventPayload<E>;
}

/**
 * Parse an event payload and prove that its canonical JSON survives transport.
 *
 * Custom event-bus providers must call this before publishing. The returned
 * runtime payload is suitable for in-process listeners; `transportValue` is
 * the exact JSON-safe value to encode for a serialized transport.
 */
export async function prepareEventPayloadForTransport<
  E extends EventPayloadDef,
>(
  event: E,
  payload: unknown,
  options?: EventPublishOptions,
): Promise<PreparedEventPayload<E>> {
  const parsed = isEventPayloadParsed(event, payload, options)
    ? (payload as InferEventPayload<E>)
    : await parseEventPayload(event, payload);
  const transportValue = toEventTransportValue(event.name, parsed);
  let canonicalPayload = parsed;
  let transportStateIsCurrent = false;

  try {
    transportStateIsCurrent = isEventPayloadTransportStable(
      event,
      parsed,
      transportValue,
      options,
    );
  } catch {
    // Revalidate if a mutated or proxied payload can no longer be compared.
  }

  if (!transportStateIsCurrent) {
    let reparsed: InferEventPayload<E>;
    try {
      reparsed = await parseEventPayload(event, transportValue);
    } catch (error) {
      throw new EventTransportError({
        eventName: event.name,
        reason: "not-stable",
        message:
          "is not transport-stable: its canonical JSON failed validation when decoded again. Event schemas must accept their own canonical JSON output.",
        cause: error,
      });
    }

    let reparsedTransportValue: EventTransportValue;
    try {
      reparsedTransportValue = toEventTransportValue(event.name, reparsed);
    } catch (error) {
      throw new EventTransportError({
        eventName: event.name,
        reason: "not-stable",
        path: error instanceof EventTransportError ? error.path : undefined,
        message:
          "is not transport-stable: validating its canonical JSON produced a value that is not JSON-safe.",
        cause: error,
      });
    }

    if (!eventTransportValuesEqual(transportValue, reparsedTransportValue)) {
      throw new EventTransportError({
        eventName: event.name,
        reason: "not-stable",
        message:
          "is not transport-stable: its canonical JSON changed when validated again. Event schema transforms must be idempotent.",
      });
    }
    canonicalPayload = reparsed;
  }

  try {
    return {
      payload: canonicalPayload,
      transportValue,
      publishOptions: markEventPayloadTransportStable(
        event,
        canonicalPayload,
        transportValue,
        options,
      ),
    };
  } catch (error) {
    throw new EventTransportError({
      eventName: event.name,
      reason: "not-stable",
      message:
        "is not transport-stable: its canonical output could not be inspected consistently.",
      cause: error,
    });
  }
}

/**
 * Validate an event payload, prove transport stability, and publish it through
 * an event bus.
 */
export async function publishEvent<E extends EventPayloadDef>(
  eventBus: EventBusLike,
  event: E,
  payload: InferEventPayload<E>,
  options?: EventPublishOptions,
): Promise<void> {
  const prepared = await prepareEventPayloadForTransport(
    event,
    payload,
    options,
  );
  await eventBus.publish(event, prepared.payload, prepared.publishOptions);
}

function assertReadyTimeoutMs(value: number): void {
  if (Number.isInteger(value) && value >= 1 && value <= MAX_TIMER_MS) return;
  throw new RangeError(
    `readyTimeoutMs must be an integer between 1 and ${MAX_TIMER_MS} milliseconds.`,
  );
}

function errorsFromCleanup(error: unknown): unknown[] {
  return error instanceof AggregateError ? [...error.errors] : [error];
}

function unsubscribeAllSettled(
  subscriptions: readonly EventSubscription[],
): Promise<PromiseSettledResult<void>[]> {
  const attempts: Promise<void>[] = [];
  for (const subscription of [...subscriptions].reverse()) {
    try {
      attempts.push(Promise.resolve(subscription.unsubscribe()));
    } catch (error) {
      attempts.push(Promise.reject(error));
    }
  }

  return Promise.allSettled(attempts);
}

function throwSubscriptionCleanupErrors(
  results: readonly PromiseSettledResult<void>[],
): void {
  const errors = results.flatMap((result) =>
    result.status === "rejected" ? errorsFromCleanup(result.reason) : [],
  );
  if (errors.length > 0) {
    throw new AggregateError(errors, "Event subscription cleanup failed");
  }
}

function withListenerDeadline<T>(
  operation: Promise<T>,
  args: {
    deadlineAt: number;
    timeoutError: Error;
  },
): Promise<T> {
  let timeout: ReturnType<typeof setTimeout> | undefined;
  const timedOut = new Promise<T>((_, reject) => {
    timeout = setTimeout(
      () => {
        reject(args.timeoutError);
      },
      Math.max(0, args.deadlineAt - performance.now()),
    );
  });

  return Promise.race([operation, timedOut]).finally(() => {
    if (timeout !== undefined) clearTimeout(timeout);
  });
}

/**
 * Register listeners against an event bus and return a composite lifecycle
 * handle.
 *
 * Payloads are validated before listener handlers run. Listener context is
 * resolved per delivery when `options.ctx` is a factory. Initial registration
 * starts every child cleanup after a synchronous subscribe failure, rejected
 * readiness promise, or readiness timeout. Cleanup that cannot finish inside
 * the registration deadline is reported without extending startup forever.
 */
export function registerListeners<Ctx>(
  eventBus: EventBusLike,
  listeners: readonly ListenerDef<EventDef, Ctx>[],
  options: RegisterListenersOptions<Ctx> = {},
): EventSubscription {
  const readyTimeoutMs =
    options.readyTimeoutMs ?? DEFAULT_LISTENER_READY_TIMEOUT_MS;
  assertReadyTimeoutMs(readyTimeoutMs);
  const registrationDeadlineAt = performance.now() + readyTimeoutMs;
  const listenerNames = listeners.map(({ name }) => name);

  const subscriptions: EventSubscription[] = [];
  let registrationError: unknown;

  for (const listener of listeners) {
    try {
      subscriptions.push(
        eventBus.subscribe(
          listener.event,
          async (rawPayload, publishOptions) => {
            try {
              const prepared = await prepareEventPayloadForTransport(
                listener.event,
                rawPayload,
                publishOptions,
              );
              const traceAttributes = {
                "beignet.listener.name": listener.name,
                "beignet.event.name": listener.event.name,
              } as const;
              await runWithResolvedTracingContext({
                tracing: options.tracing,
                ctx: options.ctx as Ctx | (() => MaybePromise<Ctx>),
                operation: {
                  name: `beignet.listener ${listener.name}`,
                  type: "listener",
                  kind: "consumer",
                  parent: parseTraceCarrier(prepared.publishOptions.trace),
                  attributes: traceAttributes,
                  metricAttributes: traceAttributes,
                },
                run: (ctx) =>
                  listener.handle({
                    event: listener.event,
                    payload: prepared.payload,
                    ctx,
                  }),
              });
            } catch (error) {
              options.onError?.(error, listener);
              if (!options.onError) throw error;
            }
          },
        ),
      );
    } catch (error) {
      registrationError = error;
      break;
    }
  }

  let rejectClosed!: (error: unknown) => void;
  const closedBeforeReady = new Promise<void>((_, reject) => {
    rejectClosed = reject;
  });
  // A consumer may only await cleanup. Keep an explicit observer on the
  // cancellation signal so that use does not create an unhandled rejection.
  void closedBeforeReady.catch(() => undefined);

  let cleanupPromise: Promise<void> | undefined;
  const cleanup = (deadlineAt?: number): Promise<void> => {
    if (cleanupPromise) return cleanupPromise;
    const operation = unsubscribeAllSettled(subscriptions);
    cleanupPromise =
      deadlineAt === undefined
        ? operation.then(throwSubscriptionCleanupErrors)
        : withListenerDeadline(operation, {
            deadlineAt,
            timeoutError: new ListenerRegistrationCleanupTimeoutError({
              timeoutMs: readyTimeoutMs,
              listenerNames,
            }),
          }).then(throwSubscriptionCleanupErrors);
    return cleanupPromise;
  };

  const initialReadiness = registrationError
    ? Promise.reject(registrationError)
    : Promise.all(subscriptions.map(({ ready }) => ready)).then(
        () => undefined,
      );
  let readySettled = false;
  const ready = withListenerDeadline(
    Promise.race([initialReadiness, closedBeforeReady]),
    {
      deadlineAt: registrationDeadlineAt,
      timeoutError: new ListenerRegistrationTimeoutError({
        timeoutMs: readyTimeoutMs,
        listenerNames,
      }),
    },
  )
    .catch(async (primaryError: unknown) => {
      try {
        await cleanup(registrationDeadlineAt);
      } catch (cleanupError) {
        throw new AggregateError(
          [primaryError, ...errorsFromCleanup(cleanupError)],
          "Listener registration failed and cleanup failed",
        );
      }
      throw primaryError;
    })
    .finally(() => {
      readySettled = true;
    });
  // EventSubscription allows callers that only own cleanup. Preserve the
  // rejection for awaiters while preventing process-level unhandled noise.
  void ready.catch(() => undefined);

  return {
    ready,
    unsubscribe() {
      if (!readySettled) {
        rejectClosed(new EventSubscriptionClosedError());
      }
      return cleanup(readySettled ? undefined : registrationDeadlineAt);
    },
  };
}

/**
 * Create listener helper methods bound to an application context type.
 *
 * Call it once in `lib/listeners.ts`:
 *
 * ```ts
 * export const { defineListener } = createListeners<AppContext>();
 * ```
 */
export function createListeners<Ctx>(): Listeners<Ctx> {
  return {
    defineListener<Name extends string, E extends EventDef>(
      name: Name,
      options: DefineListenerOptions<E, Ctx>,
    ): ListenerDef<E, Ctx, Name> {
      return defineListenerImpl(name, options);
    },
  };
}
