PidokuInfra

Calling LLMs

Expert Advanced 1h Difficulty 3/5 Topic 06 of 09

Prerequisites IV.03, V.05

The idea in one minute#

Calling a language model is an HTTP request with a JSON body and, usually, a streamed response. The mechanics are simple. Doing it well is four things: stream so users see output immediately and you can measure time to first token; bound time with a context; retry the right failures with backoff and never the wrong ones; and treat the model’s output as untrusted input that must be parsed and validated — ideally by asking for JSON that matches a schema.

Nearly every model server speaks the same wire format (the “OpenAI-compatible” chat API), so a client written once works against hosted APIs, vLLM, Ollama and llama.cpp by changing a URL.

An analogy#

Ordering by phone from a busy kitchen. You state the order precisely (the request). You want to hear items read back as they are plated rather than silence for ten minutes (streaming). If the line is engaged you call back after a pause, a little longer each time — but if they tell you the dish does not exist, calling again will not help (retry policy). And you check the bag before you leave (validation).

A picture#

flowchart TB
  REQ["Build request<br/>model, messages, max_tokens, response_format"] --> SEND["POST /v1/chat/completions<br/>with ctx deadline"]
  SEND --> ST{"status?"}
  ST -->|"200"| SSE["Read SSE events<br/>data: {delta} ... data: [DONE]"]
  ST -->|"429 or 5xx"| BACK["Wait: Retry-After, or<br/>exponential backoff + jitter"]
  BACK -->|"attempts left and ctx alive"| SEND
  BACK -->|"exhausted"| FAIL["return error"]
  ST -->|"other 4xx"| FAIL
  SSE --> TOK["onToken(text)<br/>record time to first token"]
  SSE --> USE["usage: prompt and completion tokens<br/>finish_reason"]
  USE --> VAL{"structured output?"}
  VAL -->|"yes"| PARSE["Decode strictly into a struct<br/>validate fields"]
  VAL -->|"no"| DONE["done"]
  PARSE --> DONE
  class REQ,DONE neutral
  class SEND,SSE io
  class ST,VAL queue
  class BACK,FAIL warn
  class TOK,USE,PARSE compute

How it really works#

The request#

JSON
{
  "model": "some-model",
  "messages": [
    {"role": "system", "content": "You are concise."},
    {"role": "user", "content": "Summarize this ticket."}
  ],
  "max_tokens": 400,
  "temperature": 0.2,
  "stream": true
}
  • Messages carry roles: system (instructions), user, assistant (earlier replies), and tool (results of tool calls — lesson 09). The server applies the model’s chat template (lesson 04).
  • Always set max_tokens. It bounds cost and latency, and it is the only limit a runaway generation respects.
  • The API is stateless: every call sends the whole conversation. Keeping the beginning of the prompt byte-for-byte identical between calls lets the server reuse its prefix cache, which is the largest cost lever you control from the client.

Providers differ in details — field names for token limits, how system prompts and tool calls are expressed, authentication headers. The official Go SDKs (openai-go, anthropic-sdk-go, google.golang.org/genai) hide those differences per provider; the shape underneath is what this lesson builds.

Streaming#

With "stream": true the response is Server-Sent Events (V.05): each data: line is a JSON chunk carrying a few characters in choices[0].delta.content; the last chunks carry a finish_reason and, if requested, token usage; then data: [DONE].

Why stream even when you do not show text to a person:

  • Time to first token is measurable, and is the best early signal that a request is healthy.
  • You can cancel as soon as you have what you need, or when the user leaves.
  • A per-chunk idle timeout detects a stalled stream in seconds, where a non-streaming call would wait out its full deadline.

Check finish_reason: stop is normal; length means the output hit max_tokens and is truncated; tool_calls means the model wants a function run; content_filter means it was blocked.

Timeouts#

Three clocks, three purposes:

ClockMechanismTypical
The whole call, including retriesContext deadlineYour caller’s budget
Time to first tokenA timer you start at send and stop on the first chunkSeconds
Idle gap between chunksA timer you reset per chunkSeconds

Do not use http.Client.Timeout for streaming: it would cut a long, healthy stream off.

Retries#

OutcomeRetry?
Connection error, timeout before any byteYes
429 Too Many RequestsYes — honour Retry-After if present
500, 502, 503, 529Yes
400, 401, 403, 404, 422No: the request is wrong; it will be wrong again
Context cancelled or deadline passedNo
Failure after tokens have streamedNot transparently: the caller has already seen partial output

Wait base × 2^attempt plus random jitter, capped, with a small maximum number of attempts. Without jitter, every client that failed together retries together and knocks the server over again. Retries also multiply load during an outage: bound them with a retry budget (for example, retries may be at most 10% of requests) or a circuit breaker.

Rate limits and cost#

Providers limit requests per minute and tokens per minute. Limit yourself first: a token bucket sized to your quota (IV.06), with WaitN(ctx, estimatedTokens). Record usage from every response — it is the bill, per request, in the two units you pay for — and export it as metrics (Observability V.04).

Structured output#

Free text is hard to use in a program. Ask for JSON that matches a JSON Schema: the server constrains generation so that the output parses and has the required fields. Then decode strictly and validate the values anyway:

  • Structured output guarantees shape, not truth. A confidence of 0.93 is still the model’s claim.
  • Decode with unknown fields disallowed; check ranges and enums in code.
  • If the output was truncated (finish_reason: length) the JSON is incomplete — treat it as a failure, not as something to repair.
  • Where a server does not support schemas, ask for JSON in the prompt, parse, and on failure retry once with the parse error appended.

One client, shared#

Create one client with one http.Client and use it from every goroutine. It is safe for concurrent use and holds the connection pool. Raise MaxIdleConnsPerHost when you send many concurrent requests to a single model server (V.05).

Security notes#

  • Keep API keys out of code and logs; read them from the environment or a secret store.
  • Model output is untrusted: never pass it to a shell, a SQL string or html/template as trusted markup.
  • Text retrieved from documents or tools can contain instructions aimed at the model (“prompt injection”). Do not give a model that reads untrusted text the ability to take consequential actions without a check (lesson 09).

Code#

A complete client against a fake OpenAI-compatible server in the same process, so it runs offline. Point BaseURL at any real compatible endpoint and it works unchanged.

Go
// llmclient.go — a streaming chat client: retries with backoff, SSE parsing, structured output.
package main

import (
	"bufio"
	"bytes"
	"context"
	"encoding/json"
	"errors"
	"fmt"
	"io"
	"math/rand"
	"net/http"
	"net/http/httptest"
	"strconv"
	"strings"
	"sync/atomic"
	"time"
)

// ---- the wire format shared by OpenAI-compatible servers (vLLM, Ollama, llama.cpp, hosted APIs)

type Message struct {
	Role    string `json:"role"`
	Content string `json:"content"`
}

type ChatRequest struct {
	Model          string          `json:"model"`
	Messages       []Message       `json:"messages"`
	Stream         bool            `json:"stream,omitempty"`
	MaxTokens      int             `json:"max_tokens,omitempty"`
	Temperature    float64         `json:"temperature"`
	ResponseFormat json.RawMessage `json:"response_format,omitempty"`
}

type Usage struct {
	PromptTokens     int `json:"prompt_tokens"`
	CompletionTokens int `json:"completion_tokens"`
}

type chunk struct {
	Choices []struct {
		Delta        struct{ Content string } `json:"delta"`
		Message      Message                  `json:"message"`
		FinishReason string                   `json:"finish_reason"`
	} `json:"choices"`
	Usage *Usage `json:"usage"`
}

type APIError struct {
	Status     int
	RetryAfter time.Duration
	Body       string
}

func (e *APIError) Error() string { return fmt.Sprintf("api status %d: %s", e.Status, e.Body) }

func (e *APIError) retryable() bool { return e.Status == 429 || e.Status >= 500 }

// ---- client

type Client struct {
	BaseURL, APIKey string
	HTTP            *http.Client // shared: it owns the connection pool
	MaxRetries      int
}

// do sends one request, retrying retryable failures with exponential backoff and jitter.
func (c *Client) do(ctx context.Context, req ChatRequest) (*http.Response, error) {
	body, err := json.Marshal(req)
	if err != nil {
		return nil, err
	}
	backoff := 20 * time.Millisecond
	for attempt := 0; ; attempt++ {
		hr, err := http.NewRequestWithContext(ctx, "POST", c.BaseURL+"/v1/chat/completions", bytes.NewReader(body))
		if err != nil {
			return nil, err
		}
		hr.Header.Set("Content-Type", "application/json")
		hr.Header.Set("Authorization", "Bearer "+c.APIKey)
		resp, err := c.HTTP.Do(hr)
		if err == nil && resp.StatusCode == http.StatusOK {
			return resp, nil
		}
		var apiErr *APIError
		if err == nil {
			msg, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
			resp.Body.Close() // always: otherwise the connection cannot be reused
			apiErr = &APIError{Status: resp.StatusCode, Body: strings.TrimSpace(string(msg))}
			if s, convErr := strconv.Atoi(resp.Header.Get("Retry-After")); convErr == nil {
				apiErr.RetryAfter = time.Duration(s) * time.Second
			}
			err = apiErr
		}
		if ctx.Err() != nil || attempt >= c.MaxRetries || (apiErr != nil && !apiErr.retryable()) {
			return nil, fmt.Errorf("after %d attempt(s): %w", attempt+1, err)
		}
		wait := backoff + time.Duration(rand.Int63n(int64(backoff))) // jitter: do not retry in step with everyone else
		if apiErr != nil && apiErr.RetryAfter > 0 {
			wait = apiErr.RetryAfter
		}
		fmt.Printf("  attempt %d failed (%v); retrying in %v\n", attempt+1, err, wait.Round(time.Millisecond))
		select {
		case <-time.After(wait):
		case <-ctx.Done():
			return nil, ctx.Err()
		}
		backoff *= 2
	}
}

// Stream calls onToken for each piece of text as it arrives and returns the usage.
func (c *Client) Stream(ctx context.Context, req ChatRequest, onToken func(string)) (Usage, string, error) {
	req.Stream = true
	resp, err := c.do(ctx, req)
	if err != nil {
		return Usage{}, "", err
	}
	defer resp.Body.Close()

	var usage Usage
	finish := ""
	sc := bufio.NewScanner(resp.Body)
	sc.Buffer(make([]byte, 0, 64<<10), 1<<20) // one event can exceed the 64 kB default
	for sc.Scan() {
		data, ok := strings.CutPrefix(sc.Text(), "data: ")
		if !ok {
			continue // blank separators and comments
		}
		if data == "[DONE]" {
			return usage, finish, nil
		}
		var ch chunk
		if err := json.Unmarshal([]byte(data), &ch); err != nil {
			return usage, finish, fmt.Errorf("bad event %q: %w", data, err)
		}
		if ch.Usage != nil {
			usage = *ch.Usage
		}
		for _, choice := range ch.Choices {
			if choice.Delta.Content != "" {
				onToken(choice.Delta.Content)
			}
			if choice.FinishReason != "" {
				finish = choice.FinishReason
			}
		}
	}
	if err := sc.Err(); err != nil {
		return usage, finish, err
	}
	return usage, finish, errors.New("stream ended without [DONE]")
}

// Structured asks for JSON matching a schema and decodes it strictly into out.
func (c *Client) Structured(ctx context.Context, req ChatRequest, schema string, out any) error {
	req.ResponseFormat = json.RawMessage(`{"type":"json_schema","json_schema":{"name":"result","strict":true,"schema":` + schema + `}}`)
	resp, err := c.do(ctx, req)
	if err != nil {
		return err
	}
	defer resp.Body.Close()
	var ch chunk
	if err := json.NewDecoder(resp.Body).Decode(&ch); err != nil {
		return err
	}
	if len(ch.Choices) == 0 {
		return errors.New("no choices in response")
	}
	dec := json.NewDecoder(strings.NewReader(ch.Choices[0].Message.Content))
	dec.DisallowUnknownFields() // the model's output is untrusted input: validate it
	return dec.Decode(out)
}

// ---- a fake OpenAI-compatible server, so the program runs offline

func fakeServer() *httptest.Server {
	var calls atomic.Int32
	return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
		var req ChatRequest
		if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
			http.Error(w, "bad json", http.StatusBadRequest)
			return
		}
		if len(req.Messages) == 0 {
			http.Error(w, "messages is required", http.StatusBadRequest)
			return
		}
		if calls.Add(1) == 1 { // the first call is rate limited
			http.Error(w, "rate limit exceeded", http.StatusTooManyRequests)
			return
		}
		if !req.Stream {
			content := `{"sentiment":"positive","confidence":0.93}`
			json.NewEncoder(w).Encode(map[string]any{"choices": []any{map[string]any{
				"message": Message{"assistant", content}, "finish_reason": "stop"}}})
			return
		}
		w.Header().Set("Content-Type", "text/event-stream")
		rc := http.NewResponseController(w)
		words := strings.Fields("Go handles the streaming and the retries ; the model does the thinking .")
		for _, word := range words {
			time.Sleep(8 * time.Millisecond)
			ev, _ := json.Marshal(map[string]any{"choices": []any{map[string]any{"delta": map[string]string{"content": word + " "}}}})
			fmt.Fprintf(w, "data: %s\n\n", ev)
			rc.Flush()
		}
		fmt.Fprintf(w, "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}],\"usage\":{\"prompt_tokens\":21,\"completion_tokens\":%d}}\n\n", len(words))
		fmt.Fprint(w, "data: [DONE]\n\n")
	}))
}

func main() {
	srv := fakeServer()
	defer srv.Close()
	c := &Client{BaseURL: srv.URL, APIKey: "test", HTTP: srv.Client(), MaxRetries: 3}

	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) // total budget for the call
	defer cancel()

	fmt.Println("streaming call:")
	start := time.Now()
	var first time.Duration
	fmt.Print("  ")
	usage, finish, err := c.Stream(ctx, ChatRequest{
		Model:    "local-model",
		Messages: []Message{{"system", "You are concise."}, {"user", "Who does what?"}},
	}, func(tok string) {
		if first == 0 {
			first = time.Since(start)
			fmt.Print("  ")
		}
		fmt.Print(tok)
	})
	fmt.Printf("\n  finish=%s  prompt=%d  completion=%d tokens  ttft=%dms  total=%dms  err=%v\n",
		finish, usage.PromptTokens, usage.CompletionTokens, first.Milliseconds(), time.Since(start).Milliseconds(), err)

	fmt.Println("structured call:")
	var result struct {
		Sentiment  string  `json:"sentiment"`
		Confidence float64 `json:"confidence"`
	}
	schema := `{"type":"object","properties":{"sentiment":{"type":"string","enum":["positive","negative","neutral"]},"confidence":{"type":"number"}},"required":["sentiment","confidence"],"additionalProperties":false}`
	err = c.Structured(ctx, ChatRequest{Model: "local-model", Messages: []Message{{"user", "Classify: 'this is great'"}}}, schema, &result)
	fmt.Printf("  decoded: %+v  err=%v\n", result, err)

	// A client error is not retried: sending the same bad request again cannot help.
	_, _, err = c.Stream(ctx, ChatRequest{Model: "local-model"}, func(string) {})
	fmt.Println("bad request:", err)
}

Remember this#

  • An LLM call is JSON over HTTP with an SSE response; the same client works for most servers.
  • Stream; bound total time with a context and stalls with an idle timer; always set max_tokens.
  • Retry 429 and 5xx with backoff and jitter; never retry other 4xx; respect Retry-After.
  • Ask for schema-constrained JSON, decode strictly, validate, and check finish_reason.
  • Record token usage on every call.

Try it#

  1. Run llmclient.go. Point it at a local Ollama or vLLM server if you have one; what had to change?
  2. Add a per-chunk idle timeout of 100 ms and make the fake server stall mid-stream. Confirm the client gives up quickly with a clear error.
  3. Add a token-bucket limiter of 1,000 tokens per second and send 50 concurrent requests. Plot when each one started.

Check yourself#

  1. Why should streaming be used even for machine consumers?
  2. Which HTTP statuses should be retried, and why is jitter necessary?
  3. What does schema-constrained output guarantee, and what does it not?

↑↓ navigate↵ openesc close