// TODO code inspired by https://github.com/webrtc/samples/blob/gh-pages/src/content/insertable-streams/endtoend-encryption/js/worker.js
import { EventEmitter } from 'events';
import type TypedEventEmitter from 'typed-emitter';
import { getErrorDescription } from '../../api/utils';
import {
  appendPacketTrailerToEncodedFrame,
  processPacketTrailer,
} from '../../frameMetadata/frameMetadata';
import type { FrameMetadataPublishOptions } from '../../frameMetadata/types';
import { hasFrameMetadataPublishOptions } from '../../frameMetadata/utils';
import { workerLogger } from '../../logger';
import { type VideoCodec, videoCodecs } from '../../room/track/options';
import { mimeTypeToVideoCodecString } from '../../room/track/utils';
import type { NonSharedUint8Array } from '../../type-polyfills/non-shared-typed-arrays';
import { ENCRYPTION_ALGORITHM, IV_LENGTH, UNENCRYPTED_BYTES } from '../constants';
import { CryptorError, CryptorErrorReason } from '../errors';
import { type CryptorCallbacks, CryptorEvent } from '../events';
import type {
  DecodeRatchetOptions,
  KeyProviderOptions,
  KeySet,
  PTMetadataFromE2EEMessage,
  RatchetResult,
} from '../types';
import { deriveKeys, isVideoFrame, needsRbspUnescaping, parseRbsp, writeRbsp } from '../utils';
import { ErrorRateLimiter } from './ErrorRateLimiter';
import type { ParticipantKeyHandler } from './ParticipantKeyHandler';
import { processNALUsForEncryption } from './naluUtils';
import { identifySifPayload } from './sifPayload';

export const encryptionEnabledMap: Map<string, boolean> = new Map();

export interface FrameCryptorConstructor {
  new (opts?: unknown): BaseFrameCryptor;
}

export interface TransformerInfo {
  readable: ReadableStream;
  writable: WritableStream;
  transformer: TransformStream;
  trackId: string;
  symbol: symbol;
}

export class BaseFrameCryptor extends (EventEmitter as new () => TypedEventEmitter<CryptorCallbacks>) {
  protected encodeFunction(
    encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame,
    controller: TransformStreamDefaultController,
  ): Promise<any> {
    throw Error('not implemented for subclass');
  }

  protected decodeFunction(
    encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame,
    controller: TransformStreamDefaultController,
  ): Promise<any> {
    throw Error('not implemented for subclass');
  }
}

/**
 * Cryptor is responsible for en-/decrypting media frames.
 * Each Cryptor instance is responsible for en-/decrypting a single mediaStreamTrack.
 */
export class FrameCryptor extends BaseFrameCryptor {
  private sendCounts: Map<number, number>;

  private participantIdentity: string | undefined;

  private trackId: string | undefined;

  private keys: ParticipantKeyHandler;

  private videoCodec?: VideoCodec;

  private rtpMap: Map<number, VideoCodec>;

  private keyProviderOptions: KeyProviderOptions;

  /**
   * used for detecting server injected unencrypted frames
   */
  private sifTrailer: NonSharedUint8Array;

  private detectedCodec?: VideoCodec;

  private currentTransform?: TransformerInfo;

  /**
   * The encoded streams this cryptor was last set up with, retained beyond the
   * lifetime of {@link currentTransform}.
   *
   * The main thread transfers a receiver's encoded streams to the worker once and
   * cannot transfer them again (a transferred stream is locked), and
   * `createEncodedStreams()` may only be called once per receiver. So when a
   * transceiver is reused, or when a pipe dies underneath us, re-establishing the
   * transform is only possible from here. See {@link ensureTransform}.
   */
  private retainedStreams?: {
    readable: ReadableStream<RTCEncodedVideoFrame | RTCEncodedAudioFrame>;
    writable: WritableStream<RTCEncodedVideoFrame | RTCEncodedAudioFrame>;
    operation: 'encode' | 'decode';
  };

  /**
   * Whether the subscribed track advertises packet trailer features.
   * When false, we skip the per-frame trailer extraction path entirely
   * on decode to avoid unnecessary work on tracks that don't use it.
   */
  private hasFrameMetadata: boolean = false;

  private frameMetadataOpts?: FrameMetadataPublishOptions;

  private frameMetadataFrameId = 0;

  private errorLimiter = new ErrorRateLimiter();

  private undecryptedTrackTimeout?: ReturnType<typeof setTimeout>;

  /** grace period for a teardown or a resubscribe to land before we report a stalled track */
  private readonly UNDECRYPTED_TRACK_GRACE_MS = 2000;

  /**
   * Tracks (participant, trackId, payloadType) tuples for which we've already logged a NALU
   * fallback, so a persistent bad state doesn't flood the console (Firefox doesn't filter debug).
   */
  private loggedNALUFallbacks: Set<string> = new Set();

  constructor(opts: {
    keys: ParticipantKeyHandler;
    participantIdentity: string;
    keyProviderOptions: KeyProviderOptions;
    sifTrailer?: NonSharedUint8Array;
  }) {
    super();
    this.sendCounts = new Map();
    this.keys = opts.keys;
    this.participantIdentity = opts.participantIdentity;
    this.rtpMap = new Map();
    this.keyProviderOptions = opts.keyProviderOptions;
    this.sifTrailer = opts.sifTrailer ?? Uint8Array.from([]);
  }

  private get logContext() {
    return {
      participant: this.participantIdentity,
      mediaTrackId: this.trackId,
      fallbackCodec: this.videoCodec,
    };
  }

  /**
   * Assign a different participant to the cryptor.
   * useful for transceiver re-use
   * @param id
   * @param keys
   */
  setParticipant(id: string, keys: ParticipantKeyHandler) {
    workerLogger.debug('setting new participant on cryptor', {
      ...this.logContext,
      newParticipant: id,
      hadPreviousParticipant: !!this.participantIdentity,
    });

    if (this.participantIdentity && this.participantIdentity !== id) {
      workerLogger.warn('cryptor has already a participant set, cleaning up before switching', {
        oldParticipant: this.participantIdentity,
        newParticipant: id,
        trackId: this.trackId,
      });
      // Clean up state from previous participant
      this.unsetParticipant();
    }

    this.participantIdentity = id;
    this.keys = keys;
  }

  unsetParticipant() {
    workerLogger.debug('unsetting participant', this.logContext);

    // NOTE: deliberately does not clear `currentTransform`. Nothing here cancels
    // the pipe, so the transform really is still running and clearing it would
    // desync our bookkeeping from reality -- which previously let a later
    // resubscribe skip setup while frames kept flowing through a cryptor with no
    // participant assigned.
    clearTimeout(this.undecryptedTrackTimeout);
    this.undecryptedTrackTimeout = undefined;
    this.participantIdentity = undefined;
    this.errorLimiter.reset();
  }

  isEnabled() {
    if (this.participantIdentity) {
      return encryptionEnabledMap.get(this.participantIdentity);
    } else {
      return undefined;
    }
  }

  getParticipantIdentity() {
    return this.participantIdentity;
  }

  getTrackId() {
    return this.trackId;
  }

  /**
   * Re-point this cryptor at a new trackId, keeping its (already transferred)
   * encoded streams. Used when a transceiver is reused for a new track.
   */
  setTrackId(trackId: string) {
    if (this.trackId === trackId) {
      return;
    }
    workerLogger.debug('re-pointing cryptor at new trackId', {
      ...this.logContext,
      newTrackId: trackId,
    });
    this.trackId = trackId;
  }

  hasActiveTransform() {
    return !!this.currentTransform;
  }

  /**
   * A track that is subscribed and known to be encrypted, but has no transform to
   * decrypt it, will never produce a decodable frame again. That state used to be
   * completely silent (a black tile and endless PLIs), so report it.
   *
   * Deferred, because on teardown the encoded streams routinely close before the
   * 'removeTransform' message arrives -- checking immediately would cry wolf on
   * every unsubscribe.
   *
   * Only meaningful while decoding: a sender's pipe closing on unpublish is
   * routine and has no 'removeTransform' equivalent to quiet it down.
   */
  private scheduleUndecryptedTrackWatchdog(operation: 'encode' | 'decode') {
    clearTimeout(this.undecryptedTrackTimeout);
    this.undecryptedTrackTimeout = undefined;
    if (operation !== 'decode') {
      return;
    }
    this.undecryptedTrackTimeout = setTimeout(() => {
      if (this.currentTransform || !this.isEnabled()) {
        return;
      }
      workerLogger.warn('encrypted track has no active decrypt transform', this.logContext);
      this.emitThrottledError(
        new CryptorError(
          `no active decrypt transform for encrypted track ${this.trackId}`,
          CryptorErrorReason.InternalError,
          this.participantIdentity,
        ),
      );
    }, this.UNDECRYPTED_TRACK_GRACE_MS);
  }

  /**
   * Re-establish the transform if it is gone while we still own the encoded
   * streams, e.g. after a pipe died on its own (the swallowed
   * 'Destination stream closed') or when a reused transceiver never got a new
   * pipeline. Without this the frames pile up in a readable nobody reads and the
   * decoder starves, which shows up as a permanently black tile.
   */
  ensureTransform() {
    if (this.currentTransform) {
      return true;
    }
    if (!this.retainedStreams || this.trackId === undefined) {
      workerLogger.warn('no streams retained, cannot re-establish transform', this.logContext);
      return false;
    }
    const { readable, writable, operation } = this.retainedStreams;
    workerLogger.info('re-establishing transform', { ...this.logContext, operation });
    return this.setupTransform(operation, readable, writable, this.trackId);
  }

  /**
   * Update the video codec used by the mediaStreamTrack
   * @param codec
   */
  setVideoCodec(codec: VideoCodec) {
    this.videoCodec = codec;
  }

  /**
   * rtp payload type map used for figuring out codec of payload type when encoding
   * @param map
   */
  setRtpMap(map: Map<number, VideoCodec>) {
    this.rtpMap = map;
  }

  /**
   * Sets whether the track associated with this cryptor carries packet
   * trailer data. When false, {@link decodeFunction} skips the per-frame
   * trailer extraction branch entirely.
   */
  setHasFrameMetadata(hasFrameMetadata: boolean) {
    this.hasFrameMetadata = hasFrameMetadata;
  }

  setFrameMetadataOpts(frameMetadata?: FrameMetadataPublishOptions) {
    this.frameMetadataOpts = frameMetadata;
    this.frameMetadataFrameId = 0;
  }

  setupTransform(
    operation: 'encode' | 'decode',
    readable: ReadableStream<RTCEncodedVideoFrame | RTCEncodedAudioFrame>,
    writable: WritableStream<RTCEncodedVideoFrame | RTCEncodedAudioFrame>,
    trackId: string,
    codec?: VideoCodec,
    frameMetadata?: FrameMetadataPublishOptions,
  ) {
    if (codec) {
      workerLogger.info('setting codec on cryptor to', { codec });
      this.videoCodec = codec;
    }
    if (operation === 'encode') {
      this.setFrameMetadataOpts(frameMetadata);
    }

    workerLogger.debug('Setting up frame cryptor transform', {
      operation,
      passedTrackId: trackId,
      codec,
      hasCurrentTransform: !!this.currentTransform,
      ...this.logContext,
    });

    this.trackId = trackId;

    // Retain the streams so we can rebuild the pipe later even if this cryptor
    // gets detached from its participant in between (see `ensureTransform`).
    this.retainedStreams = { readable, writable, operation };

    clearTimeout(this.undecryptedTrackTimeout);

    const symbol = Symbol('transform');

    const transformFn = operation === 'encode' ? this.encodeFunction : this.decodeFunction;
    const transformStream = new TransformStream({
      transform: transformFn.bind(this),
    });

    // Store transform info before starting the pipe
    this.currentTransform = {
      readable,
      writable,
      transformer: transformStream,
      trackId,
      symbol,
    };

    // pipeThrough/pipeTo throw synchronously on an already locked stream, which
    // would bypass the .catch below and leave us with no transform at all.
    try {
      readable
        .pipeThrough(transformStream)
        .pipeTo(writable)
        .catch((e) => {
          if (e instanceof TypeError && e.message === 'Destination stream closed') {
            // this can happen when subscriptions happen in quick successions, but doesn't influence functionality
            workerLogger.debug('destination stream closed');
          } else {
            workerLogger.warn('transform error', { error: e, ...this.logContext });
            this.emit(
              CryptorEvent.Error,
              e instanceof CryptorError
                ? e
                : new CryptorError(e.message, undefined, this.participantIdentity),
            );
          }
        })
        .finally(() => {
          // Only clear currentTransform if it's still the same one we started
          if (this.currentTransform?.symbol === symbol) {
            workerLogger.debug('transform completed', {
              ...this.logContext,
              trackId,
            });
            this.currentTransform = undefined;
          }
          // A pipe ending while the track is still assigned to a participant that
          // publishes encrypted media means we have silently stopped decrypting.
          this.scheduleUndecryptedTrackWatchdog(operation);
        });
    } catch (e: any) {
      if (this.currentTransform?.symbol === symbol) {
        this.currentTransform = undefined;
      }
      workerLogger.error('failed to set up transform', { error: e, ...this.logContext });
      this.emit(
        CryptorEvent.Error,
        e instanceof CryptorError
          ? e
          : new CryptorError(e.message, CryptorErrorReason.InternalError, this.participantIdentity),
      );
      return false;
    }

    return true;
  }

  setSifTrailer(trailer: NonSharedUint8Array) {
    workerLogger.debug('setting SIF trailer', { ...this.logContext, trailer });
    this.sifTrailer = trailer;
  }

  private emitThrottledError(error: CryptorError) {
    const errorKey = `${this.participantIdentity}-${error.reason}-decrypt`;
    const emit = this.errorLimiter.shouldEmit(errorKey, () => {
      workerLogger.warn(`Suppressing further decryption errors for ${this.participantIdentity}`, {
        ...this.logContext,
        errorKey,
      });
    });
    if (!emit) return;

    const count = this.errorLimiter.countFor(errorKey);
    if (count > 1) {
      workerLogger.debug(`Decryption error (${count} occurrences in window)`, {
        ...this.logContext,
        reason: CryptorErrorReason[error.reason],
      });
    }
    this.emit(CryptorEvent.Error, error);
  }

  /**
   * Function that will be injected in a stream and will encrypt the given encoded frames.
   *
   * @param {RTCEncodedVideoFrame|RTCEncodedAudioFrame} encodedFrame - Encoded video frame.
   * @param {TransformStreamDefaultController} controller - TransportStreamController.
   *
   * The VP8 payload descriptor described in
   * https://tools.ietf.org/html/rfc7741#section-4.2
   * is part of the RTP packet and not part of the frame and is not controllable by us.
   * This is fine as the SFU keeps having access to it for routing.
   *
   * The encrypted frame is formed as follows:
   * 1) Find unencrypted byte length, depending on the codec, frame type and kind.
   * 2) Form the GCM IV for the frame as described above.
   * 3) Encrypt the rest of the frame using AES-GCM.
   * 4) Allocate space for the encrypted frame.
   * 5) Copy the unencrypted bytes to the start of the encrypted frame.
   * 6) Append the ciphertext to the encrypted frame.
   * 7) Append the IV.
   * 8) Append a single byte for the key identifier.
   * 9) Enqueue the encrypted frame for sending.
   */
  protected async encodeFunction(
    encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame,
    controller: TransformStreamDefaultController,
  ) {
    // skip for encryption and packet trailer writes for empty dtx frames
    if (encodedFrame.data.byteLength === 0) {
      return controller.enqueue(encodedFrame);
    }

    if (!this.isEnabled()) {
      this.appendFrameMetadata(encodedFrame);
      return controller.enqueue(encodedFrame);
    }
    const keySet = this.keys.getKeySet();
    if (!keySet) {
      this.emitThrottledError(
        new CryptorError(
          `key set not found for ${
            this.participantIdentity
          } at index ${this.keys.getCurrentKeyIndex()}`,
          CryptorErrorReason.MissingKey,
          this.participantIdentity,
        ),
      );
      return;
    }
    const { encryptionKey } = keySet;
    const keyIndex = this.keys.getCurrentKeyIndex();

    if (encryptionKey) {
      const iv = this.makeIV(
        encodedFrame.getMetadata().synchronizationSource ?? -1,
        encodedFrame.timestamp,
      );
      let frameInfo = this.getUnencryptedBytes(encodedFrame);

      // Thіs is not encrypted and contains the VP8 payload descriptor or the Opus TOC byte.
      const frameHeader = new Uint8Array(encodedFrame.data, 0, frameInfo.unencryptedBytes);

      // Frame trailer contains the R|IV_LENGTH and key index
      const frameTrailer = new Uint8Array(2);

      frameTrailer[0] = IV_LENGTH;
      frameTrailer[1] = keyIndex;

      // Construct frame trailer. Similar to the frame header described in
      // https://tools.ietf.org/html/draft-omara-sframe-00#section-4.2
      // but we put it at the end.
      //
      // ---------+-------------------------+-+---------+----
      // payload  |IV...(length = IV_LENGTH)|R|IV_LENGTH|KID |
      // ---------+-------------------------+-+---------+----
      try {
        const cipherText = await crypto.subtle.encrypt(
          {
            name: ENCRYPTION_ALGORITHM,
            iv,
            additionalData: new Uint8Array(encodedFrame.data, 0, frameHeader.byteLength),
          },
          encryptionKey,
          new Uint8Array(encodedFrame.data, frameInfo.unencryptedBytes),
        );

        let newDataWithoutHeader: NonSharedUint8Array = new Uint8Array(
          cipherText.byteLength + iv.byteLength + frameTrailer.byteLength,
        );
        newDataWithoutHeader.set(new Uint8Array(cipherText)); // add ciphertext.
        newDataWithoutHeader.set(new Uint8Array(iv), cipherText.byteLength); // append IV.
        newDataWithoutHeader.set(frameTrailer, cipherText.byteLength + iv.byteLength); // append frame trailer.

        if (frameInfo.requiresNALUProcessing) {
          newDataWithoutHeader = writeRbsp(newDataWithoutHeader);
        }

        var newData = new Uint8Array(frameHeader.byteLength + newDataWithoutHeader.byteLength);
        newData.set(frameHeader);
        newData.set(newDataWithoutHeader, frameHeader.byteLength);

        encodedFrame.data = newData.buffer;
        this.appendFrameMetadata(encodedFrame);

        return controller.enqueue(encodedFrame);
      } catch (e: any) {
        // TODO: surface this to the app.
        workerLogger.error(`error while encrypting`, { ...this.logContext, error: e });
      }
    } else {
      workerLogger.debug('failed to encrypt, emitting error', this.logContext);
      this.emitThrottledError(
        new CryptorError(
          `encryption key missing for encoding`,
          CryptorErrorReason.MissingKey,
          this.participantIdentity,
        ),
      );
    }
  }

  private appendFrameMetadata(encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame) {
    if (!hasFrameMetadataPublishOptions(this.frameMetadataOpts) || !isVideoFrame(encodedFrame)) {
      return;
    }

    if (this.frameMetadataOpts?.frameId) {
      this.frameMetadataFrameId =
        this.frameMetadataFrameId === 0xffffffff ? 1 : this.frameMetadataFrameId + 1;
    }
    appendPacketTrailerToEncodedFrame(
      encodedFrame,
      this.frameMetadataOpts,
      this.frameMetadataFrameId,
    );
  }

  /**
   * Function that will be injected in a stream and will decrypt the given encoded frames.
   *
   * @param {RTCEncodedVideoFrame|RTCEncodedAudioFrame} encodedFrame - Encoded video frame.
   * @param {TransformStreamDefaultController} controller - TransportStreamController.
   */
  protected async decodeFunction(
    encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame,
    controller: TransformStreamDefaultController,
  ) {
    if (this.hasFrameMetadata && isVideoFrame(encodedFrame)) {
      try {
        const ptResult = processPacketTrailer(encodedFrame, this.trackId);
        if (ptResult.data) {
          encodedFrame.data = ptResult.data;
        }
        if (ptResult.payload && this.participantIdentity) {
          const msg: PTMetadataFromE2EEMessage = {
            kind: 'packetTrailerMetadata',
            data: ptResult.payload,
          };
          postMessage(msg);
        }
      } catch {
        // best-effort: never break the media pipeline if trailer parsing fails
      }
    }

    // skip decryption for empty dtx frames
    if (encodedFrame.data.byteLength === 0) {
      return controller.enqueue(encodedFrame);
    }

    const encryptionEnabled = this.isEnabled();

    if (encryptionEnabled === undefined) {
      // Either way we drop: forwarding would hand ciphertext straight to the
      // decoder, which renders as a black frame and logs nothing.
      if (this.participantIdentity === undefined) {
        // The pipe outlived its subscription. Expected on unsubscribe/disconnect:
        // we deliberately leave the pipe running so a reused transceiver can be
        // re-pointed at a new track, so in-flight frames still arrive here for a
        // moment. Not worth reporting -- the watchdog covers a track that stays
        // subscribed without a transform.
        workerLogger.debug('dropping frame for unassigned cryptor', this.logContext);
        return;
      }
      // A live subscription whose encryption state we were never told about.
      this.emitThrottledError(
        new CryptorError(
          `encryption state unknown for track ${this.trackId}, dropping frame`,
          CryptorErrorReason.InternalError,
          this.participantIdentity,
        ),
      );
      return;
    }

    if (!encryptionEnabled) {
      // participant is known to publish unencrypted media
      return controller.enqueue(encodedFrame);
    }

    if (isFrameServerInjected(encodedFrame.data, this.sifTrailer)) {
      encodedFrame.data = encodedFrame.data.slice(
        0,
        encodedFrame.data.byteLength - this.sifTrailer.byteLength,
      );
      if (await identifySifPayload(encodedFrame.data)) {
        workerLogger.debug('enqueue SIF', this.logContext);
        return controller.enqueue(encodedFrame);
      } else {
        workerLogger.warn('Unexpected SIF frame payload, dropping frame', this.logContext);
        return;
      }
    }
    const data = new Uint8Array(encodedFrame.data);
    const keyIndex = data[encodedFrame.data.byteLength - 1];

    if (this.keys.hasInvalidKeyAtIndex(keyIndex)) {
      // drop frame
      return;
    }

    if (this.keys.getKeySet(keyIndex)) {
      try {
        const decodedFrame = await this.decryptFrame(encodedFrame, keyIndex);
        this.keys.decryptionSuccess(keyIndex);
        if (decodedFrame) {
          return controller.enqueue(decodedFrame);
        }
      } catch (error) {
        if (error instanceof CryptorError && error.reason === CryptorErrorReason.InvalidKey) {
          // emit an error if the key handler thinks we have a valid key
          if (this.keys.hasValidKey) {
            this.emitThrottledError(error);
            this.keys.decryptionFailure(keyIndex);
          }
        } else {
          workerLogger.warn('decoding frame failed', { error });
        }
      }
    } else {
      // emit an error if the key index is out of bounds but the key handler thinks we still have a valid key
      workerLogger.warn(`skipping decryption due to missing key at index ${keyIndex}`);
      this.emitThrottledError(
        new CryptorError(
          `missing key at index ${keyIndex} for participant ${this.participantIdentity}`,
          CryptorErrorReason.MissingKey,
          this.participantIdentity,
        ),
      );
      this.keys.decryptionFailure(keyIndex);
    }
  }

  /**
   * Function that will decrypt the given encoded frame. If the decryption fails, it will
   * ratchet the key for up to RATCHET_WINDOW_SIZE times.
   */
  private async decryptFrame(
    encodedFrame: RTCEncodedVideoFrame | RTCEncodedAudioFrame,
    keyIndex: number,
    initialMaterial: KeySet | undefined = undefined,
    ratchetOpts: DecodeRatchetOptions = { ratchetCount: 0 },
  ): Promise<RTCEncodedVideoFrame | RTCEncodedAudioFrame | undefined> {
    const keySet = this.keys.getKeySet(keyIndex);
    if (!ratchetOpts.encryptionKey && !keySet) {
      throw new TypeError(`no encryption key found for decryption of ${this.participantIdentity}`);
    }
    let frameInfo = this.getUnencryptedBytes(encodedFrame);

    // Construct frame trailer. Similar to the frame header described in
    // https://tools.ietf.org/html/draft-omara-sframe-00#section-4.2
    // but we put it at the end.
    //
    // ---------+-------------------------+-+---------+----
    // payload  |IV...(length = IV_LENGTH)|R|IV_LENGTH|KID |
    // ---------+-------------------------+-+---------+----

    try {
      const frameHeader: NonSharedUint8Array = new Uint8Array(
        encodedFrame.data,
        0,
        frameInfo.unencryptedBytes,
      );
      var encryptedData: NonSharedUint8Array = new Uint8Array(
        encodedFrame.data,
        frameHeader.length,
        encodedFrame.data.byteLength - frameHeader.length,
      );
      if (frameInfo.requiresNALUProcessing && needsRbspUnescaping(encryptedData)) {
        encryptedData = parseRbsp(encryptedData);
        const newUint8 = new Uint8Array(frameHeader.byteLength + encryptedData.byteLength);
        newUint8.set(frameHeader);
        newUint8.set(encryptedData, frameHeader.byteLength);
        encodedFrame.data = newUint8.buffer;
      }

      const frameTrailer = new Uint8Array(encodedFrame.data, encodedFrame.data.byteLength - 2, 2);

      const ivLength = frameTrailer[0];
      const iv = new Uint8Array(
        encodedFrame.data,
        encodedFrame.data.byteLength - ivLength - frameTrailer.byteLength,
        ivLength,
      );

      const cipherTextStart = frameHeader.byteLength;
      const cipherTextLength =
        encodedFrame.data.byteLength -
        (frameHeader.byteLength + ivLength + frameTrailer.byteLength);

      const plainText = await crypto.subtle.decrypt(
        {
          name: ENCRYPTION_ALGORITHM,
          iv,
          additionalData: new Uint8Array(encodedFrame.data, 0, frameHeader.byteLength),
        },
        ratchetOpts.encryptionKey ?? keySet!.encryptionKey,
        new Uint8Array(encodedFrame.data, cipherTextStart, cipherTextLength),
      );

      const newData = new ArrayBuffer(frameHeader.byteLength + plainText.byteLength);
      const newUint8 = new Uint8Array(newData);

      newUint8.set(new Uint8Array(encodedFrame.data, 0, frameHeader.byteLength));
      newUint8.set(new Uint8Array(plainText), frameHeader.byteLength);

      encodedFrame.data = newData;

      return encodedFrame;
    } catch (error: any) {
      if (this.keyProviderOptions.ratchetWindowSize > 0) {
        if (ratchetOpts.ratchetCount < this.keyProviderOptions.ratchetWindowSize) {
          workerLogger.debug(
            `ratcheting key attempt ${ratchetOpts.ratchetCount} of ${
              this.keyProviderOptions.ratchetWindowSize
            }, for kind ${encodedFrame instanceof RTCEncodedAudioFrame ? 'audio' : 'video'}`,
          );

          let ratchetedKeySet: KeySet | undefined;
          let ratchetResult: RatchetResult | undefined;
          if ((initialMaterial ?? keySet) === this.keys.getKeySet(keyIndex)) {
            // only ratchet if the currently set key is still the same as the one used to decrypt this frame
            // if not, it might be that a different frame has already ratcheted and we try with that one first
            ratchetResult = await this.keys.ratchetKey(keyIndex, false);

            ratchetedKeySet = await deriveKeys(ratchetResult.cryptoKey, this.keyProviderOptions);
          }

          const frame = await this.decryptFrame(encodedFrame, keyIndex, initialMaterial || keySet, {
            ratchetCount: ratchetOpts.ratchetCount + 1,
            encryptionKey: ratchetedKeySet?.encryptionKey,
          });
          if (frame && ratchetedKeySet) {
            // before updating the keys, make sure that the keySet used for this frame is still the same as the currently set key
            // if it's not, a new key might have been set already, which we don't want to override
            if ((initialMaterial ?? keySet) === this.keys.getKeySet(keyIndex)) {
              this.keys.setKeySet(ratchetedKeySet, keyIndex, ratchetResult);
              // decryption was successful, set the new key index to reflect the ratcheted key set
              this.keys.setCurrentKeyIndex(keyIndex);
            }
          }
          return frame;
        } else {
          /**
           * Because we only set a new key once decryption has been successful,
           * we can be sure that we don't need to reset the key to the initial material at this point
           * as the key has not been updated on the keyHandler instance
           */

          workerLogger.warn('maximum ratchet attempts exceeded');
          throw new CryptorError(
            `valid key missing for participant ${this.participantIdentity}`,
            CryptorErrorReason.InvalidKey,
            this.participantIdentity,
          );
        }
      } else {
        throw new CryptorError(
          `Decryption failed: ${getErrorDescription(error, 'decryption')}`,
          CryptorErrorReason.InvalidKey,
          this.participantIdentity,
        );
      }
    }
  }

  /**
   * Construct the IV used for AES-GCM and sent (in plain) with the packet similar to
   * https://tools.ietf.org/html/rfc7714#section-8.1
   * It concatenates
   * - the 32 bit synchronization source (SSRC) given on the encoded frame,
   * - the 32 bit rtp timestamp given on the encoded frame,
   * - a send counter that is specific to the SSRC. Starts at a random number.
   * The send counter is essentially the pictureId but we currently have to implement this ourselves.
   * There is no XOR with a salt. Note that this IV leaks the SSRC to the receiver but since this is
   * randomly generated and SFUs may not rewrite this is considered acceptable.
   * The SSRC is used to allow demultiplexing multiple streams with the same key, as described in
   *   https://tools.ietf.org/html/rfc3711#section-4.1.1
   * The RTP timestamp is 32 bits and advances by the codec clock rate (90khz for video, 48khz for
   * opus audio) every second. For video it rolls over roughly every 13 hours.
   * The send counter will advance at the frame rate (30fps for video, 50fps for 20ms opus audio)
   * every second. It will take a long time to roll over.
   *
   * See also https://developer.mozilla.org/en-US/docs/Web/API/AesGcmParams
   */
  private makeIV(synchronizationSource: number, timestamp: number) {
    const iv = new ArrayBuffer(IV_LENGTH);
    const ivView = new DataView(iv);

    // having to keep our own send count (similar to a picture id) is not ideal.
    if (!this.sendCounts.has(synchronizationSource)) {
      // Initialize with a random offset, similar to the RTP sequence number.
      this.sendCounts.set(synchronizationSource, Math.floor(Math.random() * 0xffff));
    }

    const sendCount = this.sendCounts.get(synchronizationSource) ?? 0;

    ivView.setUint32(0, synchronizationSource);
    ivView.setUint32(4, timestamp);
    ivView.setUint32(8, timestamp - (sendCount % 0xffff));

    this.sendCounts.set(synchronizationSource, sendCount + 1);

    return iv;
  }

  private getUnencryptedBytes(frame: RTCEncodedVideoFrame | RTCEncodedAudioFrame): {
    unencryptedBytes: number;
    requiresNALUProcessing: boolean;
  } {
    // Handle audio frames
    if (!isVideoFrame(frame)) {
      return { unencryptedBytes: UNENCRYPTED_BYTES.audio, requiresNALUProcessing: false };
    }

    // Detect and track codec changes
    const detectedCodec = this.getVideoCodec(frame) ?? this.videoCodec;
    if (detectedCodec !== this.detectedCodec) {
      workerLogger.debug('detected different codec', {
        detectedCodec,
        oldCodec: this.detectedCodec,
        ...this.logContext,
      });
      this.detectedCodec = detectedCodec;
    }

    // Check for unsupported codecs
    if (detectedCodec === 'av1') {
      throw new Error(`${detectedCodec} is not yet supported for end to end encryption`);
    }

    // Handle VP8/VP9 codecs (no NALU processing needed)
    if (detectedCodec === 'vp8') {
      return { unencryptedBytes: UNENCRYPTED_BYTES[frame.type], requiresNALUProcessing: false };
    }
    if (detectedCodec === 'vp9') {
      return { unencryptedBytes: 0, requiresNALUProcessing: false };
    }

    // Try NALU processing for H.264/H.265 codecs
    const payloadType = frame.getMetadata().payloadType;
    const fallbackKey = `${this.participantIdentity}-${this.trackId}-${payloadType}`;
    try {
      const knownCodec =
        detectedCodec === 'h264' || detectedCodec === 'h265' ? detectedCodec : undefined;
      const naluResult = processNALUsForEncryption(new Uint8Array(frame.data), knownCodec);

      if (naluResult.requiresNALUProcessing) {
        // Recovered for this tuple, allow a future failure to log again.
        this.loggedNALUFallbacks.delete(fallbackKey);
        return {
          unencryptedBytes: naluResult.unencryptedBytes,
          requiresNALUProcessing: true,
        };
      }
    } catch (e) {
      this.logNALUFallbackOnce(fallbackKey, payloadType, e);
    }

    // Fallback to VP8 handling
    return { unencryptedBytes: UNENCRYPTED_BYTES[frame.type], requiresNALUProcessing: false };
  }

  /**
   * Logs a NALU processing fallback at most once per (participant, trackId, payloadType) tuple,
   * so a persistent bad state doesn't flood the console (Firefox doesn't filter debug).
   */
  private logNALUFallbackOnce(
    fallbackKey: string,
    payloadType: number | undefined,
    error: unknown,
  ) {
    if (this.loggedNALUFallbacks.has(fallbackKey)) {
      return;
    }
    this.loggedNALUFallbacks.add(fallbackKey);
    workerLogger.warn('NALU processing failed, falling back to VP8 handling', {
      error,
      payloadType,
      ...this.logContext,
    });
  }

  /**
   * inspects frame mimetype if available. falls back to payloadtype and maps it to the codec specified in rtpMap
   */
  private getVideoCodec(frame: RTCEncodedVideoFrame): VideoCodec | undefined {
    const metadata = frame.getMetadata();
    if (metadata.mimeType) {
      const maybeKnownCodec = mimeTypeToVideoCodecString(metadata.mimeType);
      if (videoCodecs.includes(maybeKnownCodec)) {
        return maybeKnownCodec;
      }
    }
    if (this.rtpMap.size === 0) {
      return undefined;
    }
    const payloadType = metadata.payloadType;
    const codec = payloadType ? this.rtpMap.get(payloadType) : undefined;
    return codec;
  }
}

/**
 * we use a magic frame trailer to detect whether a frame is injected
 * by the livekit server and thus to be treated as unencrypted
 * @internal
 */
export function isFrameServerInjected(
  frameData: ArrayBuffer,
  trailerBytes: NonSharedUint8Array,
): boolean {
  if (trailerBytes.byteLength === 0) {
    return false;
  }
  const frameTrailer = new Uint8Array(
    frameData.slice(frameData.byteLength - trailerBytes.byteLength),
  );
  return trailerBytes.every((value, index) => value === frameTrailer[index]);
}
