WordloopWordloop
LearnPlatform ServicesML (Python)

ML Pipelines

Transcription, diarization, embedding, RAG, and streaming pipeline patterns for wordloop-ml.

ML Pipelines

TL;DR

Each pipeline is a sequence of composable steps with explicit data contracts at every boundary. Provider outputs are validated before domain logic consumes them. Confidence thresholds are configured, not hardcoded. Every significant step emits a trace span.


The Pipeline Contract

Every ML pipeline step follows the same structure:

  1. Receive a typed domain input — no SDK types, no raw dicts
  2. Call a provider through its gateway interface
  3. Validate the output — confidence, shape, required fields — before it crosses the boundary
  4. Return a typed domain output safe for downstream consumption

No step owns both inference and persistence. No step reaches into another step's domain. This composability is what makes each step independently testable and replaceable.


Transcription Pipeline

The transcription pipeline converts audio into a structured Transcript domain object.

Data Contract

Every transcript crossing the provider boundary conforms to this structure:

from datetime import datetime
from pydantic import BaseModel, Field

class TranscriptSegment(BaseModel):
    id: int
    speaker: str
    start: float        # seconds — never store as "HH:MM:SS" at domain boundaries
    end: float          # seconds
    text: str
    confidence: float = Field(ge=0.0, le=1.0)

class TranscriptMetadata(BaseModel):
    duration: float     # seconds
    language: str       # BCP-47
    speakers: list[str]
    created_at: datetime

class Transcript(BaseModel):
    segments: list[TranscriptSegment]
    metadata: TranscriptMetadata

Timestamps are always float seconds at domain boundaries. String formats ("00:03:12") are a presentation concern handled at the entrypoint layer before returning to callers.

Confidence Validation

Filter or flag low-confidence segments before they enter domain logic. The threshold is configuration, not a literal:

def validate_segments(
    segments: list[TranscriptSegment],
    threshold: float,
) -> list[TranscriptSegment]:
    valid = [s for s in segments if s.confidence >= threshold]
    if len(valid) < len(segments):
        logger.warning(
            "transcription.low_confidence_segments_filtered",
            total=len(segments),
            filtered=len(segments) - len(valid),
            threshold=threshold,
        )
    return valid

Provider Mapping

The provider maps SDK responses to domain types at the boundary. Domain logic never imports assemblyai, whisperx, or any other transcription SDK:

class AssemblyAIProvider:
    async def transcribe(self, audio_uri: str) -> Transcript:
        raw = await self._client.transcribe(audio_uri)
        return self._to_domain(raw)

    def _to_domain(self, raw: aai.Transcript) -> Transcript:
        return Transcript(
            segments=[
                TranscriptSegment(
                    id=i,
                    speaker=u.speaker or "unknown",
                    start=u.start / 1000,   # ms → seconds
                    end=u.end / 1000,
                    text=u.text,
                    confidence=u.confidence or 0.0,
                )
                for i, u in enumerate(raw.utterances or [])
            ],
            metadata=TranscriptMetadata(
                duration=raw.audio_duration or 0.0,
                language=raw.language_code or "en",
                speakers=list({u.speaker for u in raw.utterances or [] if u.speaker}),
                created_at=datetime.utcnow(),
            ),
        )

Diarization Pipeline

Diarization assigns speaker identities to transcript segments. It runs after transcription — consuming a Transcript and returning enriched speaker labels.

Diarization confidence is independent of transcription confidence. A segment can have high transcription confidence and low diarization confidence. Validate both independently.

class DiarisedSegment(BaseModel):
    segment_id: int
    speaker_id: str
    speaker_confidence: float = Field(ge=0.0, le=1.0)

When diarization is unavailable or fails, return the transcript without speaker attribution rather than surfacing an error to the caller. The transcript is still useful without speaker labels. Define the degradation boundary explicitly in the service layer, not the provider (see Resilience — Graceful Degradation).


Embedding Pipeline

Embeddings convert text into dense vectors for semantic search and similarity matching.

Batching

Never embed one segment at a time. Batch requests reduce API cost and latency significantly:

class EmbeddingProvider:
    async def embed_batch(self, texts: list[str]) -> list[EmbeddingVector]:
        raw = await self._client.embeddings.create(
            input=texts,
            model=self._model,
        )
        return [
            EmbeddingVector(
                vector=item.embedding,
                model=self._model,
                dimensions=len(item.embedding),
            )
            for item in raw.data
        ]

Versioning

Store the model name and version alongside every vector. Queries against stale embeddings — generated by a previous model version — produce incorrect similarity scores. Validate dimensions at the boundary: an unexpected dimension count signals a model change that needs addressing.

Caching

Cache on the combination of text content, model name, and model version. Embedding the same transcript segment twice is a pure cost with no benefit.


RAG Pipeline

Retrieval-Augmented Generation combines vector search with LLM generation. The pipeline has three composable steps, each independently testable.

1. Retrieve

Embed the query and search the vector store for the top-k candidates:

async def retrieve(self, query: str, top_k: int = 10) -> list[RetrievalCandidate]:
    query_vector = await self.embedder.embed(query)
    return await self.vector_store.search(query_vector, top_k=top_k)

2. Rerank

Score candidates for relevance and discard those below the threshold before passing to generation. This step prevents low-quality context from degrading the LLM output:

def rerank(
    self,
    candidates: list[RetrievalCandidate],
    query: str,
    threshold: float,
) -> list[RetrievalCandidate]:
    scored = self.reranker.score(candidates, query)
    return [c for c in scored if c.relevance_score >= threshold]

3. Generate

Pass reranked context to the LLM. Prompts are version-controlled Python constants — not runtime configuration, not database records:

INSIGHT_PROMPT = """
You are analysing a meeting transcript.

Context:
{context}

Question: {question}

Provide a concise, factual answer based only on the context above.
If the context does not contain enough information, say so explicitly.
Do not infer or extrapolate beyond what the context states.
""".strip()

async def generate(self, question: str, context: list[RetrievalCandidate]) -> str:
    prompt = INSIGHT_PROMPT.format(
        context="\n\n".join(c.text for c in context),
        question=question,
    )
    return await self.llm_gateway.complete(prompt)

A prompt edited outside version control and without a corresponding eval run is an untested code change. Prompts go through PR review and are covered by the eval suite the same way code changes are.


Streaming Pipeline

Real-time audio processing uses semantic endpointing: buffer audio until a natural speech boundary — a pause, a sentence end — then process the completed chunk. Sending partial sentences to the LLM produces incoherent outputs.

class SemanticEndpointer:
    MIN_CHUNK_DURATION_S = 3.0
    MIN_SILENCE_MS = 500

    def should_process(self, buffer: AudioBuffer) -> bool:
        if buffer.duration < self.MIN_CHUNK_DURATION_S:
            return False
        if buffer.silence_since_ms < self.MIN_SILENCE_MS:
            return False
        return True

Each chunk is an independent pipeline invocation: transcribe → extract insights → publish to Pub/Sub. Results stream back to the client as they complete. Chunk boundaries, processing latency, and insight extraction timing are all emitted as trace spans so the pipeline is observable end-to-end.

No-Audio Detection: Header Caching + Batched Decode

StreamingSessionService (src/wordloop/core/services/streaming_session_service.py) scores every forwarded audio frame for silence, so it can report no_audio_detected when a full window (ML_NO_AUDIO_WINDOW_SECONDS) has been silent and recovered when audio returns. Raw PCM16 frames (the browser's AudioWorklet capture path) are scored inline with a numpy RMS check — no codec involved.

Containerised frames (the browser's encoding: "webm" capture path) need decoding first, and that decode step is deliberately not per-frame:

  1. Header caching. A MediaRecorder-style WebM stream sends the initialisation segment (EBML header + Segment/Tracks) only once, in the first timeslice; every later frame is a bare Cluster with no header, which AudioCodecGateway.decode_pcm16 cannot decode on its own (None, not an error). split_container_header() splits the stream's first chunk into (header, first cluster) once; the header is cached for the life of the session and prepended to every later batch before decoding.
  2. Batched decode. Frames accumulate in a pending buffer; once ML_NO_AUDIO_BATCH_FRAMES (default 10, ~1s of audio at the browser's 100ms timeslice) have queued, header + joined_clusters is decoded once on a dedicated ThreadPoolExecutor (thread_name_prefix="ml-audio-score") — never on the drain thread that forwards audio to the transcription provider. Decoding per-frame would mean one ffmpeg subprocess spawn per 100ms frame on that thread; batching cuts that to roughly once a second, and moving it off the drain thread means a slow decode can never add latency to the live forwarding path.
  3. Frames that arrive while a batch is decoding keep accumulating (bounded, oldest dropped past a small multiple of the batch size) so a stream's last batch is never stranded when audio stops mid-batch.

The decoded PCM16 feeds the same silent-window state machine PCM16 streams use directly. This design is specific to no-audio detection — the composed audio Core stores and later batch-transcribes is unaffected; the codec call here only produces a same-signal RMS score, not a stored artifact.


Anti-Patterns

Raw SDK responses in domain logic. assemblyai.Transcript in a service method is an architecture violation. Map at the provider boundary; pass domain types across.

Hardcoded thresholds. Confidence thresholds, relevance scores, and chunk sizes belong in configuration. They need to be tunable per deployment without a code change.

Embedding one item at a time. Calling the embedding API per segment in a loop is 10–100× more expensive than batching. Always batch.

Embedding without versioning. Vectors stored without a model identifier become unqueryable the moment the embedding model changes.

Prompts in runtime config. A prompt in a database or environment variable is an unreviewed, untested code change. Prompts are code.

Time-based audio chunking. Fixed-duration chunks (e.g. every 15 seconds) split sentences mid-word. Semantic endpointing based on silence detection produces coherent chunks that the LLM can reason about correctly.

On this page