Implemented platform

I built an OpenAI/Anthropic-compatible API gateway in Go to connect LLMs to business pipelines. It handles RAG requests and publishes events to NATS JetStream using fire-and-forget. Dagster daemon sensors pull events and execute jobs. Vector (Rust) subscribes to NATS and sends telemetry to Prometheus/Loki.

The gateway publishes events; each middleware consumes them independently. This removes the gateway’s Dagster dependency and allows separate deployments.

The initial Rust (axum) prototype was migrated to Go. Goroutines and channels reduced the code needed for concurrent SSE relay, NATS subscription and PG/pgvector writes.

Components

  • Language: Go 1.25, Gin framework
  • Messaging: NATS 2.11 (JetStream enabled)
  • Orchestration: Dagster (Python) — daemon + sensor + asset materialization
  • Data store: PostgreSQL 18 + pgvector (JIT enabled)
  • Telemetry: Vector 0.45 (Rust) → Prometheus + Loki + Grafana
  • Inference backends: vLLM, llama.cpp, LM Studio
  • Embedding + Reranker: multi-bert-inference — Rust (Axum) + ONNX Runtime (INT8, ~17MB, ~10ms). ColBERT 256d embed + 64d token MaxSim rerank. DB-free, gRPC 256d vec direct transmission (~1KB/candidate)
  • Containers: Podman rootless + podman-compose
  • Host topology: 3 hosts (storage server / desktop mac / compute server)

Design Decisions

Layered Architecture

The implementation separates Transport, Domain and Infrastructure, using Clean Architecture dependency rules.

  Client (opencode CLI, API consumer)
    |
    v
Transport Layer (Gin) [internal/transport/http/]
    |-- middleware: RequestContext -> Logger -> Recovery
    |-- v1/ handlers (OpenAI/Anthropic compatible endpoints)
    |-- Request parsing and validation
    +-- Presentation DTO conversion
    |
    |  Presentation Layer [pkg/openai/, pkg/anthropic/]
    |    +-- OpenAI/Anthropic compatible request/response DTOs
    |
    v
Domain Layer [internal/domain/]
    |-- knowledge/  RAG orchestration
    |-- llm/        Multi-backend LLM routing
    |-- vectorstore/ Vector search (pgvector ANN)
    +-- pipeline/   Pipeline integration (NATS publish, fire-and-forget)
    |
    v
Infrastructure Layer [internal/infra/]
    |-- vllm/       vLLM client
    |-- llamacpp/   llama.cpp client
    |-- lmstudio/   LM Studio client
    |-- postgres/   PostgreSQL + pgvector
    |-- nats/       NATS JetStream publish
    |-- vector/     Vector telemetry publish
    +-- reranker/   multi-bert-inference gRPC client
  

The Domain layer has no external dependencies. Changes to LLM backends or data stores stay in Infrastructure. This also separates business logic from integrations in business application development.

Gateway Is Publish-Only

The gateway only publishes to NATS. No subscribing, no triggering. Consumers (Dagster sensor, Vector) independently pull/subscribe.

The intended effects are:

  • Separate deployments without a Dagster dependency in gateway code.
  • Add or change consumers without changing the gateway.
  • Isolate gateway and pipeline failures.

Dagster’s role

Dagster daemon acts as a durable consumer, pulling from JetStream. Sensor → RunRequest → direct job execution. The GraphQL trigger bridge was eliminated in favor of Dagster’s native sensor mechanism.

Why Dagster (Python)

Dagster manages execution order and state. Heavy computation runs outside Python.

Actual ComputationExecution EnginePython/Dagster Role
LLM inferencevLLM / llama.cpp (C++/CUDA)API call only
Vector searchPostgreSQL + pgvector (C)SQL query only
Embedding + Rerankingmulti-bert-inference (Rust/ONNX Runtime)gRPC call only
Telemetry conversionVector (Rust)YAML config only

I chose Dagster for lineage tracking, sensors and asset materialization. I had not found a Rust replacement. For this I/O-wait and state-management workload, those features mattered more than Rust’s CPU-bound performance and memory-safety benefits.

LLM Backend Selection

Backend is automatically selected based on the requested model name.

BackendDefault URLExample Models
vLLMhttp://$COMPUTE_HOST:8000qwen3-next:80b, qwen3-coder:30b
llama.cpphttp://$COMPUTE_HOST:8081nemotron-3-nano:30b
LM Studiohttp://$COMPUTE_HOST:1234lfm2.5-1.2b-instruct-mlx

Gateway goroutines handle LLM routing. Choosing a backend based on response content requires stateful control, so Vector does not perform that role.

I continue to evaluate model performance and quality. Verified models are archived in GGUF, NVFP4, ONNX and other formats on storage-server cold storage.

Model archive on storage server
Verified model archive on storage server cold storage. Accumulated in GGUF, NVFP4, ONNX and other formats

Rust to Go Migration

The initial Rust (axum) prototype used its type system to define the design, but required more code for concurrent asynchronous operations.

  • Managing async contexts for SSE relay + NATS subscription + PG/pgvector writes simultaneously was verbose in Rust
  • goroutine + channel naturally fits this pattern
  • Runtime overhead difference is negligible for this use case

All design principles were preserved after migration: OpenAI-compatible endpoints, NATS Pub/Sub event-driven architecture, trace_id on all logs, idempotent design. The Rust design documents served as specifications for the Go implementation.

Implementation

Endpoints

MethodPathDescription
POST/v1/chat/completionsChat completions (OpenAI compatible)
POST/v1/messagesAnthropic Messages API
POST/v1/responsesResponses API
POST/v1/embeddingsEmbeddings
GET/v1/modelsModel listing
GET/healthzHealth check

Middleware Chain

Gin middleware: RequestContext → Logger → Recovery

  • RequestContext — Extracts request tracking ID from X-Correlation-ID or X-Request-ID header. Auto-generates a timestamp-based ID when absent, setting it on both context and response header
  • Logger — Records method, path, status, duration per HTTP request via slog
  • Recovery — Gin standard panic recovery

RAG Data Flow

  Client
  |
  v
Chat Completion Request
  |
  v
Knowledge Service (RAG orchestration)
  |-- 1. User query extraction
  |-- 2. Embedding generation (multi-bert-inference -> 256d ColBERT)
  |-- 3. Vector search (pgvector HNSW ANN) + reranking (ColBERT MaxSim)
  |-- 4. Context injection
  |-- 5. Send to LLM backend
  |-- 6. NATS publish (pipeline.* + telemetry.*)  <- fire-and-forget
  +-- 7. Response to client
  

System Topology

                                 NATS JetStream :4222
                              +----------------------+
                              |  pipeline.*           |
                              |  telemetry.*          |
                              +---+----------+--------+
                    publish ------+          | pull subscribe
                    (fire & forget)         | (durable consumer)
                                           |
+-- agent-gateway :8080 ------+    +-- Dagster daemon -----------------------+
|                             |    |                                         |
|  Client Request             |    |  sensor: nats_dagster_chat_persist      |
|    |                        |    |  sensor: nats_dagster_embedding         |
|    v                        |    |  sensor: nats_dagster_retrieve          |
|  Knowledge Service          |    |  sensor: nats_dagster_compose           |
|    |-- embed (multi-bert)   |    |  sensor: nats_dagster_flow_lineage     |
|    |-- vector search (pg)   |    |  sensor: nats_dagster_tool_call        |
|    |-- rerank (ColBERT)     |    |         |                               |
|    |-- LLM inference -------|-->-|  vLLM / llama.cpp / LM Studio          |
|    |                        |    |         |                               |
|    |-- NATS pub: pipeline.* |    |         v                               |
|    |-- NATS pub: telemetry.*|    |  RunRequest -> Dagster job execution   |
|    +-- Response             |    |    |-- chat persist -> pg               |
|                             |    |    |-- embedding event -> pg            |
|  LLM Backends               |    |    |-- retrieve event -> pg             |
|    |-- vLLM :8000           |    |    |-- compose event -> pg              |
|    |-- llama.cpp :8081      |    |    |-- flow lineage -> pg               |
|    +-- LM Studio :1234     |    |    +-- tool call -> pg                  |
|                             |    |                                         |
+-----------------------------+    +-----------------------------------------+
                                           |
                              +------------+
                              | pull subscribe
                              v
                    +-- Vector (Rust) ----------------------+
                    |  source: nats (telemetry.*)           |
                    |    |-- sink: Prometheus exporter      |
                    |    +-- sink: Loki (CorrelationID)    |
                    +--------------------------------------+
  

NATS Event Design

Two independent paths share the same CorrelationID.

Pipeline (Dagster sensor → job execution):

TopicDagster Job
pipeline.knowledge.chat.persistknowledge_chat_persist + knowledge_chat_pair_materialize
pipeline.knowledge.embeddingknowledge_embedding
pipeline.knowledge.retrieveknowledge_retrieve
pipeline.knowledge.composeknowledge_compose
pipeline.knowledge.flow.lineageknowledge_flow_lineage
pipeline.knowledge.tool_callknowledge_tool_call

Telemetry (Vector → Prometheus / Loki):

TopicContent
telemetry.knowledge.chatModel, backend, token usage, streaming flag
telemetry.knowledge.embeddingModel, input count, dimensions
telemetry.knowledge.retrieveTopK, hit count, scores
telemetry.knowledge.tool_callFunction name, model

JetStream stream configuration:

  • PIPELINE — subjects: pipeline.>, retention: limits, max-age: 72h, storage: file
  • TELEMETRY — subjects: telemetry.>, retention: limits, max-age: 24h, storage: file

Dagster Assets and Jobs

The Dagster side is structured with asset materialization + sensor patterns.

Dagster Assets catalog — knowledge group asset listing
Dagster Assets catalog. chat_pair, flow_lineage, tool_call are Materialized; others show Never materialized

Each sensor pull-subscribes from NATS JetStream, detects events, and issues RunRequests for direct job execution. No GraphQL API involved.

Dagster tool_call_record asset detail
tool_call_record asset metadata. event_id, trace_id, function_name, model recorded at materialization time

Example tool_call_record: function_name: tokei, model: minimax-m2.5 — records which tool was called by which model as an event, enabling lineage tracking in Dagster.

Dagster knowledge_chat_pair_materialize job
chat_pair_record → chat_pair_dataset asset graph. Prompt-response pair persistence job

The knowledge_chat_pair_materialize job materializes two assets: chat_pair_record and chat_pair_dataset. It persists prompt-response pairs to PostgreSQL as raw material for future fine-tuning dataset construction.

devstack Container Layout

All middleware runs as rootless containers via podman-compose.

Podman Desktop — devstack container listing
Running devstack containers: postgres, reranker, vector, nats, dagster (user-code, webserver, daemon), grafana, node-exporter, promtail
ServiceImagePortsRole
NATSnats:2.11-alpine4222, 8222JetStream messaging
PostgreSQLpostgres:18-jit-vector5432pgvector + JIT, agent_gateway / dagster DB
DagsterCustom build3000Webserver + Daemon + User-code gRPC
Vectortimberio/vector:0.45.0-alpine8686, 9598NATS telemetry → Prometheus/Loki
multi-bert-inferenceRust (Axum) + ONNX Runtime3001 (HTTP), gRPCColBERT embed (256d) + rerank (64d MaxSim)

Lakehouse profile (for Phase 2, explicit startup):

  • Nessie (19120) — Iceberg metadata catalog (git-like branching)
  • Trino (8081) — SQL engine with Iceberg catalog
  • dbt-fusion — dbt CLI container

3-Host Network Topology

NATS co-located with the gateway, keeping all paths within 1 hop.

  +-- storage server (observability) --+  +-- desktop mac (gateway) ----------+  +-- compute server (inference) ------+
|                                    |  |                                   |  |                                     |
|  Prometheus <-- scrape --------------|-- :9598 (Vector exporter)          |  |  vLLM :8000 (GPU inference)         |
|  Grafana                           |  |                                   |  |  llama.cpp :8081                    |
|  Loki <-- push -----------------------  Vector (nats source)             |  |                                     |
|  Vector (aggregator)               |  |       |                          |  |  Dagster :3000                       |
|    |-- nats source --------------------  NATS JetStream :4222 <---------|--|---- daemon sensor (pull subscribe)  |
|    |-- -> Loki sink                |  |       ^ publish                  |  |    |-- webserver                     |
|    +-- -> Prometheus exporter      |  |       | (fire-and-forget)        |  |    |-- daemon (sensor -> job)        |
|  MinIO :9000                       |  |       |                          |  |    +-- user-code gRPC               |
|                                    |  |  agent-gateway :8080             |  |                                     |
|                                    |  |    |-- proxy API                 |  |  PostgreSQL 18 :5432                |
|                                    |  |    |   (OpenAI/Anthropic/        |  |    +-- pgvector (ANN + full-text)   |
|                                    |  |    |    Responses)               |  |                                     |
|                                    |  |    |-- knowledge orchestrator    |  |  multi-bert-inference :3001         |
|                                    |  |    +-- NATS pub ----------------|--->  pipeline.* + telemetry.*          |
|                                    |  |                                   |  |                                     |
|                                    |  |  LM Studio :1234 (CPU inference) |  |  Trino :8080                        |
|                                    |  |                                   |  |  Nessie :19120                       |
+------------------------------------+  +-----------------------------------+  +-------------------------------------+
  
HostRoleKey ServicesUptime
storage serverObservability + object storagePrometheus, Grafana, Loki, Vector, MinIO24/7
desktop macGateway + messagingagent-gateway, NATS JetStream, LM StudioWorking hours
compute serverGPU inference + data platformvLLM, llama.cpp, Dagster, PostgreSQL, TrinoWorking hours

NATS shares desktop mac with the gateway, the largest publisher, avoiding a host-to-host publish hop. Desktop and compute start and stop together. Storage runs continuously, and Vector reconnects after a disconnection.

Data Storage Strategy

Storage is split by data characteristics between PostgreSQL (pgvector) and Iceberg (Parquet).

                      PostgreSQL (pgvector)             Iceberg (Parquet via Trino/Nessie)
                    --------------------             -----------------------------------
Characteristics     Structured + vector search        File-oriented + bulk accumulation
Queries             WHERE + ANN (nearest neighbor)    SQL JOIN + time-travel
Update pattern      UPSERT (row-level)                append-only (immutable Parquet)
Versioning          None (self-managed history)       Nessie branches + Iceberg snapshots
  

In Phase 1 (real-time), all data is written directly to PostgreSQL. From Phase 2 onward, Dagster batch processes will convert PG data to Parquet → Iceberg append, with dbt-fusion building analytical datasets.

PostgreSQL Schema

  • document_chunks — 256-dimension ColBERT embeddings + HNSW vector index
  • chat_history — Conversation history by correlation_id
  • rerank_scores — Reranking scores per query x document x model
  • api_responses — Responses API persistence

Caveats

  • When desktop mac / compute server are stopped, NATS goes down too, but JetStream persistent data is retained on desktop mac disk and resumes from unconsumed messages on next startup
  • Vector detects disconnection and auto-reconnects, but telemetry events during NATS downtime are lost (pipeline side is protected by JetStream durable consumers)
  • PostgreSQL 18 JIT requires shm_size: 4gb configuration. Insufficient allocation risks OOM kills

Verification

  • Confirmed asset materialization status in Dagster UI: chat_pair, flow_lineage, tool_call successfully Materialized
  • Monitored stream state and consumer lag via NATS monitoring (:8222)
  • Verified telemetry metrics delivery through Vector exporter (:9598) → Prometheus
  • Integration tests: go test -tags=integration ./internal/infra/integration ./internal/domain/pipeline

Next Actions

Phase 1, the real-time platform, was operational at the time of writing. The following phases remain planned.

PhaseContentStatus
1gateway → NATS → Dagster/pg, Vector → Prometheus/LokiRunning
2pg → Parquet → Iceberg batch, dbt-fusion processingdevstack ready
3MCP server tool definitions, Trino query tool, hypothesis testing automationNot started
4Fine-tuning pipeline, MOE expert pruning, synthetic data augmentationNot started