← Files Biological Sequence & Alignment ViewerARCHIVED FILE

src/persistent/streaming-source-replacement.ts

16.2 KB · Sep 30, 2026 · 23:01 UTC

↓ Download file

import {
  MAX_SEQUENCE_SOURCE_EDIT_CHUNK_BYTES,
  MAX_SEQUENCE_SOURCE_RANGE_BYTES,
  type ScientificSequenceDataClient,
} from "./scientific-data-client";
import type { SequenceDocument } from "../sequence/types";

const MAX_SOURCE_REPLACEMENT_INTERVALS = 4_096;
const MAX_SOURCE_REPLACEMENT_BYTES = 1024 * 1024;
const MAX_SOURCE_HEADER_BYTES = 64 * 1024;
const SUPPORTED_RESIDUE = /^[A-Za-z*.?-]$/u;

export type SequenceOriginalSourceMutation = Readonly<{
  recordNumber: number;
  recordId: string;
  edits: ReadonlyArray<
    Readonly<{
      expected: string;
      offset0: number;
      replacement: string;
    }>
  >;
}>;

export type SequenceOriginalSourceReplacementPlan = Readonly<{
  format: "fasta" | "fastq";
  mutations: ReadonlyArray<SequenceOriginalSourceMutation>;
  sourceRevision: string;
}>;

type SequenceOriginalSourceRangeClient = Pick<
  ScientificSequenceDataClient,
  "readOriginalSourceRange" | "session"
>;

/**
 * Limit original-source edits to equal-length residue substitutions in the
 * authenticated, currently visible record previews. Unloaded records, wrapped
 * lines, headers, qualities, and every other original byte remain untouched.
 */
export function createSequenceOriginalSourceReplacementPlan({
  current,
  initial,
  sourceRevision,
}: {
  current: SequenceDocument;
  initial: SequenceDocument;
  sourceRevision: string;
}): SequenceOriginalSourceReplacementPlan {
  if (/^(?:bgzf|gzip):/iu.test(sourceRevision)) {
    throw new Error(
      "Compressed original sequences require atomic recompression and index updates.",
    );
  }
  const fileName = current.fileName ?? "";
  const supportedPlainFile =
    current.format === "fasta"
      ? /\.(?:fa|faa|fas|fasta|ffn|fna|frn)$/iu.test(fileName)
      : current.format === "fastq" && /\.(?:fastq|fq)$/iu.test(fileName);
  if (!supportedPlainFile || initial.format !== current.format) {
    throw new Error(
      "Only uncompressed plaintext FASTA and FASTQ originals can be replaced.",
    );
  }

  const inventory = current.recordInventory;
  const initialInventory = initial.recordInventory;
  if (
    inventory == null ||
    initialInventory == null ||
    inventory.materializedCount !== current.records.length ||
    initialInventory.materializedCount !== initial.records.length ||
    inventory.totalCount !== initialInventory.totalCount ||
    inventory.truncated !== initialInventory.truncated ||
    current.records.length !== initial.records.length ||
    current.records.length === 0
  ) {
    throw new Error(
      "The original Sequence record inventory no longer matches its authenticated previews.",
    );
  }

  const mutations: SequenceOriginalSourceMutation[] = [];
  let intervalCount = 0;
  let changedBytes = 0;
  for (
    let recordIndex = 0;
    recordIndex < current.records.length;
    recordIndex++
  ) {
    const record = current.records[recordIndex];
    const original = initial.records[recordIndex];
    if (
      record == null ||
      original == null ||
      record.id !== original.id ||
      record.sourceLabel !== original.sourceLabel ||
      record.description !== original.description ||
      record.length !== original.length ||
      record.sequence.length !== original.sequence.length ||
      original.sequence.length !== original.length ||
      JSON.stringify(record.features) !== JSON.stringify(original.features) ||
      JSON.stringify(record.metadata) !== JSON.stringify(original.metadata) ||
      record.molecule !== original.molecule ||
      record.topology !== original.topology ||
      record.quality?.ascii !== original.quality?.ascii ||
      (current.format === "fastq" &&
        (record.quality == null ||
          record.quality.ascii.length !== record.sequence.length))
    ) {
      throw new Error(
        "Original-source replacement only supports equal-length residue edits without record, annotation, header, or quality changes.",
      );
    }

    const edits: Array<{
      expected: string;
      offset0: number;
      replacement: string;
    }> = [];
    for (let offset = 0; offset < record.sequence.length; ) {
      if (record.sequence[offset] === original.sequence[offset]) {
        offset += 1;
        continue;
      }
      const start = offset;
      while (
        offset < record.sequence.length &&
        record.sequence[offset] !== original.sequence[offset]
      ) {
        const expected = original.sequence[offset];
        const replacement = record.sequence[offset];
        if (
          expected == null ||
          replacement == null ||
          !SUPPORTED_RESIDUE.test(expected) ||
          !SUPPORTED_RESIDUE.test(replacement)
        ) {
          throw new Error(
            "Original-source replacement requires single-byte biological residues.",
          );
        }
        changedBytes += 1;
        if (changedBytes > MAX_SOURCE_REPLACEMENT_BYTES) {
          throw new Error(
            "The original Sequence edit exceeds its bounded residue budget.",
          );
        }
        offset += 1;
      }
      intervalCount += 1;
      if (intervalCount > MAX_SOURCE_REPLACEMENT_INTERVALS) {
        throw new Error(
          "The original Sequence edit exceeds its bounded interval budget.",
        );
      }
      edits.push({
        expected: original.sequence.slice(start, offset),
        offset0: start,
        replacement: record.sequence.slice(start, offset),
      });
    }
    if (edits.length > 0) {
      mutations.push({
        edits,
        recordId: original.id,
        recordNumber: recordIndex + 1,
      });
    }
  }
  if (mutations.length === 0) {
    throw new Error(
      "Only deliberate, equal-length residue substitutions can replace the original source.",
    );
  }
  return {
    format: current.format === "fasta" ? "fasta" : "fastq",
    mutations,
    sourceRevision,
  };
}

/** Stream original source bytes and patch only authenticated residue offsets. */
export async function* streamSequenceOriginalSourceReplacement({
  client,
  plan,
  signal,
}: {
  client: SequenceOriginalSourceRangeClient;
  plan: SequenceOriginalSourceReplacementPlan;
  signal?: AbortSignal;
}): AsyncGenerator<Uint8Array> {
  if (
    client.session.sourceRevision !== plan.sourceRevision ||
    /^(?:bgzf|gzip):/iu.test(plan.sourceRevision)
  ) {
    throw new Error(
      "The original Sequence source revision is compressed or has changed.",
    );
  }
  const parser = new ByteExactSequenceSourceParser(plan);
  let sourceOffset = 0n;
  let firstRange = true;
  for (;;) {
    signal?.throwIfAborted();
    if (client.session.sourceRevision !== plan.sourceRevision) {
      throw new Error("The original Sequence source revision has changed.");
    }
    const range = await client.readOriginalSourceRange({
      length: MAX_SEQUENCE_SOURCE_RANGE_BYTES,
      offsetDecimal: sourceOffset.toString(),
      signal,
    });
    if (
      firstRange &&
      range.bytes.byteLength >= 2 &&
      range.bytes[0] === 0x1f &&
      range.bytes[1] === 0x8b
    ) {
      throw new Error(
        "Compressed original sequences require atomic recompression and index updates.",
      );
    }
    firstRange = false;
    if (range.bytes.byteLength > 0) {
      const bytes = parser.consume(range.bytes);
      sourceOffset += BigInt(bytes.byteLength);
      yield bytes;
    }
    if (range.eof) {
      parser.finish();
      return;
    }
  }
}

/** Coalesce authenticated source fragments into bounded, exact-offset writes. */
export async function* sequenceOriginalSourceUploadChunks({
  maxChunkBytes,
  signal,
  source,
}: {
  maxChunkBytes: number;
  signal?: AbortSignal;
  source: AsyncIterable<Uint8Array>;
}): AsyncGenerator<{ bytes: Uint8Array; offsetDecimal: string }> {
  if (
    !Number.isSafeInteger(maxChunkBytes) ||
    maxChunkBytes <= 0 ||
    maxChunkBytes > MAX_SEQUENCE_SOURCE_EDIT_CHUNK_BYTES
  ) {
    throw new Error("The original Sequence upload chunk limit is not safe.");
  }
  let offset = 0n;
  let buffered = new Uint8Array(maxChunkBytes);
  let bufferedLength = 0;
  for await (const fragment of source) {
    signal?.throwIfAborted();
    if (
      !ArrayBuffer.isView(fragment) ||
      Object.prototype.toString.call(fragment) !== "[object Uint8Array]" ||
      fragment.byteLength === 0 ||
      fragment.byteLength > MAX_SEQUENCE_SOURCE_RANGE_BYTES
    ) {
      throw new Error("The original Sequence source fragment is not bounded.");
    }
    let fragmentOffset = 0;
    while (fragmentOffset < fragment.byteLength) {
      signal?.throwIfAborted();
      const length = Math.min(
        buffered.byteLength - bufferedLength,
        fragment.byteLength - fragmentOffset,
      );
      buffered.set(
        fragment.subarray(fragmentOffset, fragmentOffset + length),
        bufferedLength,
      );
      bufferedLength += length;
      fragmentOffset += length;
      if (bufferedLength === buffered.byteLength) {
        yield { bytes: buffered, offsetDecimal: offset.toString() };
        offset += BigInt(bufferedLength);
        buffered = new Uint8Array(maxChunkBytes);
        bufferedLength = 0;
      }
    }
  }
  if (bufferedLength > 0) {
    yield {
      bytes: buffered.slice(0, bufferedLength),
      offsetDecimal: offset.toString(),
    };
  }
}

type SourceParserState =
  | "fasta-expect-header"
  | "fasta-header"
  | "fasta-sequence"
  | "fastq-expect-header"
  | "fastq-header"
  | "fastq-sequence"
  | "fastq-plus"
  | "fastq-quality"
  | "fastq-quality-complete";

class ByteExactSequenceSourceParser {
  private state: SourceParserState;
  private lineStart = true;
  private recordNumber = 0;
  private sequenceOffset = 0n;
  private qualityOffset = 0n;
  private readonly header: number[] = [];
  private mutationIndex = 0;
  private editIndex = 0;
  private editCharacter = 0;
  private passthrough = false;

  constructor(private readonly plan: SequenceOriginalSourceReplacementPlan) {
    this.state =
      plan.format === "fasta" ? "fasta-expect-header" : "fastq-expect-header";
  }

  consume(input: Uint8Array): Uint8Array {
    if (this.passthrough) return input;
    const output = input.slice();
    for (let index = 0; index < output.byteLength; index++) {
      if (this.passthrough) break;
      this.consumeByte(output, index);
    }
    return output;
  }

  finish(): void {
    if (this.passthrough) return;
    if (this.state === "fastq-quality-complete") {
      this.finishFastqRecord();
    }
    if (this.mutationIndex !== this.plan.mutations.length) {
      throw new Error(
        "The original Sequence source ended before every residue edit was verified.",
      );
    }
    if (
      this.state === "fastq-header" ||
      this.state === "fastq-sequence" ||
      this.state === "fastq-plus" ||
      this.state === "fastq-quality"
    ) {
      throw new Error("The original FASTQ record has incomplete quality data.");
    }
  }

  private consumeByte(output: Uint8Array, index: number): void {
    const value = output[index];
    if (value == null) return;
    switch (this.state) {
      case "fasta-expect-header":
      case "fastq-expect-header": {
        if (value === 0x0a || value === 0x0d) return;
        const marker = this.state === "fasta-expect-header" ? 0x3e : 0x40;
        if (value !== marker) {
          throw new Error(
            "The original Sequence source has an invalid record header.",
          );
        }
        this.beginHeader();
        return;
      }
      case "fasta-header":
      case "fastq-header": {
        if (value === 0x0a) {
          this.finishHeader();
          return;
        }
        this.header.push(value);
        if (this.header.length > MAX_SOURCE_HEADER_BYTES) {
          throw new Error(
            "The original Sequence record header exceeds its bounded size.",
          );
        }
        return;
      }
      case "fasta-sequence": {
        if (value === 0x0a) {
          this.lineStart = true;
          return;
        }
        if (value === 0x0d) return;
        if (this.lineStart && value === 0x3e) {
          this.assertCurrentRecordComplete();
          this.beginHeader();
          return;
        }
        this.lineStart = false;
        this.patchResidue(output, index);
        if (
          this.mutationIndex === this.plan.mutations.length &&
          this.editIndex === 0
        ) {
          this.passthrough = true;
        }
        return;
      }
      case "fastq-sequence": {
        if (value === 0x0a) {
          this.lineStart = true;
          return;
        }
        if (value === 0x0d) return;
        if (this.lineStart && value === 0x2b) {
          this.assertCurrentRecordComplete();
          this.state = "fastq-plus";
          this.lineStart = false;
          return;
        }
        this.lineStart = false;
        this.patchResidue(output, index);
        return;
      }
      case "fastq-plus": {
        if (value === 0x0a) {
          this.state =
            this.sequenceOffset === 0n
              ? "fastq-quality-complete"
              : "fastq-quality";
          this.qualityOffset = 0n;
          this.lineStart = true;
        }
        return;
      }
      case "fastq-quality": {
        if (value === 0x0a || value === 0x0d) return;
        this.qualityOffset += 1n;
        if (this.qualityOffset > this.sequenceOffset) {
          throw new Error(
            "The original FASTQ quality exceeds its sequence length.",
          );
        }
        if (this.qualityOffset === this.sequenceOffset) {
          this.state = "fastq-quality-complete";
        }
        return;
      }
      case "fastq-quality-complete": {
        if (value === 0x0d) return;
        if (value !== 0x0a) {
          throw new Error(
            "The original FASTQ quality exceeds its sequence length.",
          );
        }
        this.finishFastqRecord();
        return;
      }
    }
  }

  private beginHeader(): void {
    this.recordNumber += 1;
    this.sequenceOffset = 0n;
    this.qualityOffset = 0n;
    this.header.length = 0;
    this.editIndex = 0;
    this.editCharacter = 0;
    this.lineStart = false;
    this.state = this.plan.format === "fasta" ? "fasta-header" : "fastq-header";
  }

  private finishHeader(): void {
    if (this.header.at(-1) === 0x0d) this.header.pop();
    const header = new TextDecoder("utf-8", { fatal: true }).decode(
      new Uint8Array(this.header),
    );
    const match = /^\s*(\S+)(?:\s+(.*))?$/u.exec(header);
    if (match == null) {
      throw new Error(
        "The original Sequence source has an invalid record header.",
      );
    }
    const expected = this.plan.mutations[this.mutationIndex];
    if (
      expected != null &&
      expected.recordNumber === this.recordNumber &&
      expected.recordId !== match[1]
    ) {
      throw new Error("The original Sequence record identity has changed.");
    }
    this.lineStart = true;
    this.state =
      this.plan.format === "fasta" ? "fasta-sequence" : "fastq-sequence";
  }

  private patchResidue(output: Uint8Array, index: number): void {
    const mutation = this.plan.mutations[this.mutationIndex];
    if (mutation?.recordNumber === this.recordNumber) {
      const edit = mutation.edits[this.editIndex];
      if (
        edit != null &&
        this.sequenceOffset === BigInt(edit.offset0 + this.editCharacter)
      ) {
        const expected = edit.expected.charCodeAt(this.editCharacter);
        if (output[index] !== expected) {
          throw new Error(
            "The original Sequence residue changed before replacement.",
          );
        }
        output[index] = edit.replacement.charCodeAt(this.editCharacter);
        this.editCharacter += 1;
        if (this.editCharacter === edit.expected.length) {
          this.editIndex += 1;
          this.editCharacter = 0;
          if (this.editIndex === mutation.edits.length) {
            this.mutationIndex += 1;
            this.editIndex = 0;
          }
        }
      }
    }
    this.sequenceOffset += 1n;
  }

  private assertCurrentRecordComplete(): void {
    const next = this.plan.mutations[this.mutationIndex];
    if (next?.recordNumber === this.recordNumber) {
      throw new Error(
        "The original Sequence record ended before all expected residues were found.",
      );
    }
  }

  private finishFastqRecord(): void {
    this.assertCurrentRecordComplete();
    this.state = "fastq-expect-header";
    this.lineStart = true;
    if (this.mutationIndex === this.plan.mutations.length) {
      this.passthrough = true;
    }
  }
}

SHA-256: cfa02b5ae14e9c0e812067d1272bd4bd5c92ea7d24688b24643763706b1bcda1