Building an Async RAG API in Go - No Langchain

Published on: 2026-03-02 · #go

What and why?

A technical walkthrough of a Retrieval Augmented Generation backend built in Go, with worker pools, buffered channels, semantic caching, and probably too much Redis.

Every RAG implementation is now done through python and Langchain. Langchain is now slowly becaoming commercialized and I dont want to be vendor locked. This application aims to avoid third part libraries and be flexible enough for any LLM / embedding or vector DB vendor.


What This Thing Actually Is

This is a RESTful Go API built primarily for RAG, Retrieval Augmented Generation. You send it a question, it finds relevant documents in a vector database, feeds them to an LLM along with your question, and gives you an answer. Not exactly a novel concept, but the devil’s in the implementation.

It also ingests documents. PDFs, DOCX, plain text, it chunks them up, embeds them, and stores them in Qdrant. There’s nascent MCP functionality too, but let’s not talk about that yet.

The stack:

  • Go 1.24 with Chi router (can work without Chi too)
  • Redis for job state and conversation history
  • Qdrant for vector storage and semantic caching, swappable for Pinecone, Weaviate, or whatever vector DB you prefer
  • Google Gemini for both embeddings and LLM generation, swappable for OpenAI, Anthropic, Cohere, or any provider that speaks the same interface
  • Prometheus for metrics
  • Docker for deployment (because running things locally is for people with patience)

Four endpoints. That’s it.

MethodPathPurpose
POST/chatStart or continue a conversation
GET/status/{id}Poll for job result
POST/ingestUpload a document for ingestion
POST/mcpMCP (in development)

Why Polling?

The application this backs is browser based and used in areas with low or unreliable internet. It’s also data heavy, it needs to sync its own operational data, handle panic alerts, process queued transactions and all of which take priority over a chatbot response when bandwidth is tight. The app can go offline entirely and keep working, and when it reconnects, the chatbot is not at the top of the priority list.

Polling fits that constraint. The client decides when it has bandwidth to spare. It sends a request, gets a job ID, and checks back when it can. The chatbot doesn’t get to demand a persistent connection or fight for bandwidth against more important traffic. That tradeoff is the whole point.

With a WebSocket or SSE approach, you need a live connection, a reconnect strategy, session state, and a way to handle all the things that go wrong when connectivity is spotty. With polling, you need none of that. You ask, you get an ID, you check back later. Simple to implement on the client side, and the server doesn’t care whether the client polls in 2 seconds or 2 minutes.

Every POST returns 202 Accepted with a job ID and a status URL. The client polls GET /status/{id} whenever it feels like it. That status endpoint reads directly from Redis — no channels, no workers, no queue. Just a key lookup. It’s fast. Redis handles atomicity, so the data is never stale. And we don’t block the processing channel. Also, side steps injection concerns.


The Architecture

The system is asynchronous by design. Here’s the flow:

  1. A request hits an HTTP handler
  2. The handler creates a Job, stamps it as QUEUED, and drops it into a buffered channel (capacity: 100)
  3. The handler immediately returns 202 Accepted with the job ID
  4. A worker picks the job off the channel
  5. The worker runs the RAG pipeline (embed → cache check → vector search → LLM → answer)
  6. The worker saves the result to Redis
  7. The client polls /status/{id} and eventually gets the answer

The buffered channel is the backbone. If 100 jobs are already queued and you send another, the handler blocks. Its an intentional backpressure mechanism for DDOS defence.

The Composition Root Pattern

Everything gets wired together in main.go using the composition root pattern. Channels, stores, external service clients, the RAG service, the worker pool, they are all instantiated in one place, passed everywhere else. If Redis is offline at startup, the system falls back to in memory stores. If Qdrant or the embedding service fails to initialize, the app shuts down. Could probably handle that more gracefully, but for now it just refuses to start half broken.


Worker Pool

The worker pool is where the actual work happens. It’s a dynamic pool that scales based on demand and shrinks when things are quiet.

On startup, the dispatcher creates exactly one worker. Every 10 incoming requests, the handler signals the dispatcher to spawn another one. Every ingest job also triggers a signal since ingestion is resource heavy. Workers cap out at 10. When a worker sits idle for a minute with nothing to do, it retires itself with the quiet dignity of someone who knows their time has come. The minimum is always 1, someone has to keep the lights on. It reminds me of my unstoppable march towards death. Wow, that was bleak. ’;)’

Each worker sits in a select loop with three cases: pick up a job, receive a stop signal, or time out from inactivity. Worker count is tracked with atomic.Int64 because mutexes are for people who enjoy typing.

Job execution gives each job a 60 sec timeout(configurable). State transitions (QUEUED → RUNNING → COMPLETE/ERROR) are saved to Redis at every step. For query jobs, conversation history is appended after a successful answer. Prometheus captures the duration and status of everything.


The RAG Pipeline

The RAG service uses Go’s “opaque interface” pattern. There’s a public Service interface and a private service struct. The worker interacts with the interface. It has no idea what’s behind it. If it quacks like a duck, etc.

The query pipeline has five steps:

  1. Embed — Convert the question into a 1536 dimensional vector using Google’s embedding model (dimensions can also be configured)
  2. Cache check — Search the semantic cache in Qdrant for a sufficiently similar past question (≥0.97 cosine similarity)
  3. Vector search — If no cache hit, search the main collection for the top 3 relevant document chunks
  4. LLM generation — Pass the question, context chunks, and conversation history to Gemini or LLM
  5. Cache save — Fire and forget goroutine that saves the new Q&A pair to the cache (yeet the ring into the mountain bruh)

Each step updates the job’s CurrentStep field, starts a timer, and defers a Prometheus metrics capture. If anything fails, the job is marked as error with a generic “Internal Server Error” message for the client and a detailed trace tagged log on the server side. All errors are flagged as retryable, whether that’s always appropriate is debatable, but it keeps the client from giving up on transient failures.

The semantic cache cutoff at 0.97 is deliberately tight. “What is Go?” and “What’s Go?” will probably hit the cache. “Tell me about Go” probably won’t. I think its better to burn LLM tokens regenerating an answer than to serve a cached response to the wrong question.

Again, this is a for a workplace application that’s in the utility workspace. So tight semantic cache cut off is neccessary.

Document Ingestion

Ingestion follows its own pipeline: (This is buggy at times, with large files. I will keep the structure but refactor for reliability in the coming weeks)

I am also thinking that document ingestion could be ported to python because of its large support libraries and OCR.

  1. Ensure the Qdrant collection exists
  2. Detect file type (PDF, DOCX, TXT, RTF)
  3. Extract text from the document
  4. Split into chunks (max 1000 characters, 150 character overlap)
  5. Batch embed and upsert to Qdrant (100 chunks per batch)
  6. Delete the temp file

To be honest, the recursive chunking algorithm was suggested by Claude, when I was trying to come up with a good way to process data.

The text splitter is a recursive chunking algorithm that tries separators in order of semantic quality: paragraph breaks first, then line breaks, then sentences, then words, then individual characters (for when all hope is lost). The 150 character overlap maintains semantic continuity across chunk boundaries, so if a sentence gets split between chunks, the next chunk starts with enough context to not be completely confused.

PDF extraction runs page by page with a 10sec timeout per page. If a corrupted page hangs the extraction goroutine, the timeout fires and the page is skipped. One bad page does not hold the entire ingestion hostage.

For truly massive documents (>1M chunks), the embedding step switches from synchronous API calls to Google’s asynchronous batch embeddings API. It submits the job, then polls every 30 minutes. Because sometimes patience really is a virtue.


Pluggable Providers

The three external dependencies —> embedding, LLM, and vector database, they are all behind interfaces:

Go
UTF-8|16 Lines|
type Embedder interface {
    GetEmbedding(ctx context.Context, query string) ([]float32, error)
    BatchEmbedding(ctx context.Context, chunks []string, isHugeDataSet bool) ([][]float32, error)
}

type Provider interface {
    Generate(ctx context.Context, query string, matches []string, messageHistory []string) (string, error)
}

type DataProcessor interface {
    Search(ctx context.Context, vectorVal []float32) ([]string, []string, error)
    GetCachedAnswer(ctx context.Context, queryVector []float32) (string, bool, error)
    SaveToCache(ctx context.Context, id string, vector []float32, answer string) error
    CreateCollection(ctx context.Context, collectionName string) error
    UpsertBatch(ctx context.Context, collectionName string, chunks []commonModels.DocChunk, vectors [][]float32) error
}

The point is to avoid vendor lock in. Swap Gemini for OpenAI. Replace Qdrant with Pinecone or Weaviate. Mix and match, Gemini for LLM, OpenAI for embeddings. The rest of the system doesn’t care, as long as the vector dimensions stay consistent. In theory, a config change is all it takes. In practice, each provider has quirks that the interface can’t fully abstract away, so “just swap it” is slightly optimistic. But the plumbing is there.

Currently, the implementations are:

  • Embeddings: Google’s gemini-embedding-001 via the GenAI SDK, singleton with sync.Once
  • LLM: Gemini (gemini-2.5-flash-lite-preview-09-2025), temperature 0.7, with a system instruction that politely tells the model to not get jailbroken.
  • Vector DB: Qdrant via gRPC, with both a main collection and a semantic-cache collection

Data Layer

Redis does the heavy lifting for state management. Two separate Redis databases are used because Redis generously provides 16 of them and most people use exactly one.

  • DB 0 — Job store. Jobs are JSON serialized and stored with a 24 hour TTL. The /status/{id} endpoint reads from here. It’s a straight GET — no workers, no queues. Basically the fastest thing in the entire system.

  • DB 1 — Message store. Conversation history lives in Redis Lists. Each chat ID maps to a list of serialized JobPayload objects. When the RAG pipeline needs context, it fetches the last 5 messages using LRANGE with negative offsets. No need to read the whole list, Redis handles the windowing.

The Redis wrapper uses a singleton DB pattern with double checked locking, fast path with RLock, slow path with Lock. Each store instance cleans up via a goroutine that waits for the service context to cancel.

If Redis is unreachable at startup, the system falls back to in memory stores backed by map + RWMutex. Same interface, same behavior, just won’t survive a restart. Fine for dev. Not great for production, but it avoids a hard crash.


Middleware

Every request passes through a middleware pipeline before reaching a handler. The pipeline runs three checks:

  1. Trace injection — Every request gets a trace ID. If the caller provides X-Trace-Id, we use it. Otherwise, a UUID is generated. This ID propagates through the entire system, logs, Redis operations, external API calls.

  2. Authentication — Bearer token validation with subtle.ConstantTimeCompare to prevent timing attacks. There’s a NoAuthBypass flag for local development, because typing bearer tokens into curl gets old fast. It logs a loud warning when active, as it should (for debug only).

  3. Rate limiting — Per-IP token bucket using Go’s x/time/rate. Default: 2 requests/second with a burst allowance of 5. Each IP gets its own limiter stored in a map behind an RWMutex. Currently commented out in the pipeline, but the implementation is complete and waiting. There’s a TODO to move the IP to limiter map to Redis when the user count justifies it.

The middleware also wraps the response writer in a status recorder so Prometheus can track response codes without the handlers having to care about metrics.


Observability

Metrics

Prometheus metrics cover the system end-to-end:

  • http_requests_total — Counter, labeled by path and status code
  • count_jobs_in_queue — Gauge for current queue depth
  • active_worker_count — Gauge for live workers
  • dispatcher_signal_count — Gauge tracking dispatcher triggers
  • process_request_duration_seconds — Histogram for total job time (buckets: 0.1s to 30s)
  • dependency_latency_seconds — Histogram per external service: embedding, cache lookup, vector search, LLM generation, ingestion (buckets: 0.05s to 10s)

All exposed at GET /metrics. Scrape it with Prometheus, visualize it with Grafana, stare at it at 3 AM wondering why the P99 just spiked.

Logging

The logger wraps Go’s standard slog package. In development, it outputs readable text at DEBUG level. In production, it switches to JSON at INFO level, it’s structured, parseable, ready for whatever log aggregator you’ve committed to.

Every logger instance is tagged with a component name (“Server”, “WorkerPool”, “RAG Service”). Additional context, trace IDs, job IDs, is chained via .With(). Error logs include the source file and line number via runtime.Callers. Could be more sophisticated, but it gets the job done.


Configuration

All configuration lives in a single Go file as constants. No YAML parsing, no .env file loaders. Just constants. This works fine at the current scale, it would probably need to change if the config got more complex. Sensitive values (API keys, auth tokens) are gitignored and kept separate.

Key settings worth knowing:

SettingValueWhy
Buffer limit100Jobs in the channel before backpressure kicks in
Max workers10Upper bound on concurrent processing
Idle worker timeout1 minuteWorkers that do nothing eventually cease to exist
Embedding dimensions1536Must match the embedding model
Cache similarity cutoff0.97Very tight — prefers re-generation over wrong cached answers
Redis job TTL24 hoursJobs expire after a day
Redis message TTL24 hoursConversations expire after a day
Server read timeout5 seconds
Server write timeout10 seconds
Job execution timeout60 seconds

Environment variables REDIS_ADDR, QDRANT_HOST, and QDRANT_PORT override defaults at runtime. This is how Docker Compose tells the containers where to find each other.


Testing

The test suite covers the critical paths without the ceremony of an enterprise testing framework.

RAG service tests are table driven, covering the happy path and every failure mode: full pipeline success, cache hit short circuit, embedding failure, vector search failure, and LLM failure. Each scenario wires up mock implementations using a callback pattern, set the behavior you want, or don’t and get sensible defaults.

Worker pool tests verify that workers spawn on dispatcher signals, process jobs from the channel, shut down cleanly on stop signals, and retire after idle timeouts.

Ingestion tests cover file type detection (PDF→PDF, DOCX→DOCX, TXT→DOCX, PNG→error), text chunking with overlap verification, batch processing boundaries (150 chunks → 2 batches of 100 and 50), and chunk metadata integrity.

Redis store tests use miniredis — an in memory Redis implementation so tests don’t need a real Redis server. They cover save/get roundtrips, missing key lookups, deletes, and a race condition test that fires 50 concurrent goroutines at the same job to verify the RWMutex is doing its job.

k6 performance tests simulate 20 virtual users for 5 seconds, each running the full workflow: POST a chat request → poll for status up to 10 times → verify the response. Rate limiting and error paths are exercised.


Deployment

The API ships as a Docker container. The Dockerfile is a multi stage build: first stage runs on golang:1.24.1-alpine, downloads dependencies, generates Swagger docs via swag init, and compiles the binary with CGO_ENABLED=0. Second stage copies just the binary onto a bare alpine:latest image. The result is a container that’s small enough to deploy basically anywhere.

Docker Compose brings up three services:

  • api — The Go binary, exposed on port 3000, with environment variables pointing to the other containers
  • redis — redis:alpine on port 6379, with a health check (redis-cli ping) so the API doesn’t start before Redis is ready
  • qdrant — qdrant/qdrant on ports 6333 (REST) and 6334 (gRPC)
plaintext
UTF-8|1 Line|
docker-compose up --build

Three containers, one command. The API waits for Redis to be healthy, connects to Qdrant, initializes the collections, spins up a worker, and starts listening on port 3000.

Multiple instances can run simultaneously as long as they all talk to the same Redis and Qdrant. The worker pool inside each pod handles local concurrency. Redis and Qdrant handle shared state. Whether this scales well under real production load is still being validated — the architecture supports it, but the proof is in the deployment.


Project Structure

plaintext
UTF-8|25 Lines|
cmd/api/main.go           → entry point, composition root
internal/
  adapter/                → domain-to-API model conversion
  api/                    → request/response DTOs
  config/                 → constants and environment config
  customHttpClient/       → HTTP transport with connection pooling
  data/
    redisStore/           → generic Redis wrapper (singleton per DB)
    store/                → job + message stores (Redis and in-memory)
  domain/
    commonModels/         → Document, DocChunk
    jobModel/             → Job, status enums, store interfaces
  handlers/               → HTTP handlers, job creation
  job/                    → job service (channels + stores)
  metrics/                → Prometheus definitions
  middleware/             → trace, auth, rate limiting
  rag/                    → RAG pipeline orchestration
    embedding/            → Embedder interface + Google implementation
    ingest/               → document ingestion (extract, chunk, embed, upsert)
    llm/                  → LLM Provider interface + Gemini implementation
    vectorDB/             → DataProcessor interface + Qdrant implementation
  server/                 → HTTP server + graceful shutdown
  worker/                 → worker pool + job execution
pkg/logger_i/             → slog wrapper
k6/script.js              → load test

The internal/ folder is Go’s built in access control. Nothing outside this module can import from it. The pkg/ folder is the opposite, it’s explicitly public. One folder says “mine,” the other says “yours.”


Design Decisions

Nil pointers for optional JSON fields. The API response uses pointer fields for Error and RAGResponse. When there’s no error or no answer yet, those fields are nil and get omitted from the JSON. Standard Go idiom: “communicate empty state with the absence of the object.”

Adapter layer between domain and API. Domain models have no JSON tags. API contracts have no business logic. The adapter handles translation between them. It’s an extra layer of indirection, adds some boilerplate, but keeps the two concerns from leaking into each other.

Singleton external clients with sync.Once. The Qdrant client, embedding client, and LLM client are all singletons initialized exactly once. Each registers a cleanup goroutine on context cancellation. A reasonable approach for now, though it does make testing with different configurations within the same process harder.

Constant-time auth comparison. The bearer token check uses subtle.ConstantTimeCompare instead of ==. Prevents timing-based token guessing.

Graceful shutdown. When SIGINT or SIGTERM arrives, the server stops accepting new connections, workers finish their current jobs, the service context gets cancelled (cleaning up external clients), and then the process exits. There are edge cases this probably doesn’t handle perfectly, long running ingestion jobs, for example but it covers the common case.