Go Reverse Proxy: Прекинувачи на коло, отфрлање на оптоварување и структурирано логирање во пракса
Обратниот прокси обично откажува многу пред да престане да прифаќа врски. Бавен upstream ги троши сите достапни слотови за барања, повторните обиди го зголемуваат притисокот, логовите стануваат непребарлив тек, а технички „здравиот“ прокси се претвора во многу ефикасен дистрибутер на прекини.
Овој туторијал изработува Go прокси ориентиран кон продукциска употреба што се однесува поинаку. Тој ја ограничува конкурентноста, ги изолира backend-системите што откажуваат со circuit breakers за секој backend, извршува активни здравствени проверки и емитува структурирани access логови. Имплементацијата ја користи само стандардната библиотека на Go и намерно избегнува автоматски повторни обиди, зачувувајќи безбедно однесување за неидемпотентни барања.
Предуслови и архитектура
Потребни ви се Go 1.22 или понов и два HTTP backend-системи достапни од хостот на проксито. Примерот ги користи следниве адреси:
127.0.0.1:8080: обратен прокси127.0.0.1:9001и127.0.0.1:9002: upstream сервиси/healthz: задолжителна upstream здравствена крајна точка/livezи/readyz: крајни точки за животниот циклус на проксито
Барањата прво поминуваат низ access логирање и неблокирачка порта за конкурентност. Потоа проксито избира здрав backend чиј circuit breaker дозволува сообраќај. Транспортните грешки и upstream 5xx одговорите се бројат како неуспеси; откажувањата од клиентот не се.
Постојат две намерни компромисни решенија. Изборот е round-robin наместо да е свесен за латентноста, што го прави однесувањето разбирливо под притисок. Барањата исто така се обидуваат само еднаш. Повторен обид за POST по двосмислен мрежен неуспех може да дуплира несакан ефект, па повторните обиди припаѓаат во изрично идемпотентен слој.
Структура на проектот
edgeproxy/
├── go.mod
├── cmd/
│ ├── edgeproxy/main.go
│ └── testbackend/main.go
└── deploy/
└── edgeproxy.service
Создајте ги директориумите со вашиот уредувач или mkdir -p edgeproxy/cmd/{edgeproxy,testbackend} edgeproxy/deploy, а потоа додајте ја следнава модулна датотека:
module example.com/edgeproxy
go 1.22
Имплементирајте го проксито
Circuit breaker-от користи дозволи нумерирани по генерации. Кога breaker-от ја менува состојбата, доцните резултати од претходната генерација се игнорираат. Ова спречува старо барање во тек случајно да затвори поново коло.
package main
import (
"context"
"crypto/rand"
"encoding/hex"
"flag"
"io"
"log/slog"
"net"
"net/http"
"net/http/httputil"
"net/url"
"os"
"os/signal"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"
)
type circuitState uint8
const (
closed circuitState = iota
open
halfOpen
)
type breaker struct {
mu sync.Mutex
state circuitState
failures int
threshold int
cooldown time.Duration
openedAt time.Time
probeInFlight bool
generation uint64
}
func (b *breaker) allow(now time.Time) (uint64, bool) {
b.mu.Lock()
defer b.mu.Unlock()
switch b.state {
case closed:
return b.generation, true
case open:
if now.Sub(b.openedAt) < b.cooldown {
return 0, false
}
b.state = halfOpen
b.probeInFlight = true
b.generation++
return b.generation, true
case halfOpen:
if b.probeInFlight {
return 0, false
}
b.probeInFlight = true
return b.generation, true
default:
return 0, false
}
}
func (b *breaker) done(ticket uint64, success bool) {
b.mu.Lock()
defer b.mu.Unlock()
if ticket != b.generation {
return
}
switch b.state {
case closed:
if success {
b.failures = 0
return
}
b.failures++
if b.failures >= b.threshold {
b.state = open
b.openedAt = time.Now()
b.generation++
}
case halfOpen:
b.probeInFlight = false
b.generation++
if success {
b.state = closed
b.failures = 0
} else {
b.state = open
b.openedAt = time.Now()
}
}
}
func (b *breaker) cancel(ticket uint64) {
b.mu.Lock()
defer b.mu.Unlock()
if ticket == b.generation && b.state == halfOpen {
b.probeInFlight = false
b.state = open
b.openedAt = time.Now()
b.generation++
}
}
func (b *breaker) ready(now time.Time) bool {
b.mu.Lock()
defer b.mu.Unlock()
return b.state == closed ||
(b.state == open && now.Sub(b.openedAt) >= b.cooldown)
}
type backend struct {
url *url.URL
healthy atomic.Bool
breaker breaker
proxy *httputil.ReverseProxy
}
type permit struct {
backend *backend
ticket uint64
}
type contextKey uint8
const (
permitKey contextKey = iota
metaKey
)
type requestMeta struct {
backend string
}
type selector struct {
backends []*backend
next atomic.Uint64
}
func (s *selector) choose() (*backend, uint64, bool) {
n := len(s.backends)
start := int(s.next.Add(1)-1) % n
for i := 0; i < n; i++ {
b := s.backends[(start+i)%n]
if !b.healthy.Load() {
continue
}
if ticket, ok := b.breaker.allow(time.Now()); ok {
return b, ticket, true
}
}
return nil, 0, false
}
type responseRecorder struct {
http.ResponseWriter
status int
bytes int
}
func (w *responseRecorder) WriteHeader(status int) {
if w.status == 0 {
w.status = status
w.ResponseWriter.WriteHeader(status)
}
}
func (w *responseRecorder) Write(p []byte) (int, error) {
if w.status == 0 {
w.WriteHeader(http.StatusOK)
}
n, err := w.ResponseWriter.Write(p)
w.bytes += n
return n, err
}
func (w *responseRecorder) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}
func newRequestID() string {
var value [16]byte
if _, err := rand.Read(value[:]); err == nil {
return hex.EncodeToString(value[:])
}
return hex.EncodeToString([]byte(time.Now().UTC().Format(time.RFC3339Nano)))
}
func clientAddress(remote string) string {
host, _, err := net.SplitHostPort(remote)
if err == nil {
return host
}
return remote
}
func accessLog(log *slog.Logger, next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
started := time.Now()
id := newRequestID()
meta := &requestMeta{}
r = r.WithContext(context.WithValue(r.Context(), metaKey, meta))
r.Header.Set("X-Request-ID", id)
rec := &responseRecorder{ResponseWriter: w}
next.ServeHTTP(rec, r)
status := rec.status
if status == 0 {
status = http.StatusOK
}
log.Info("access",
"request_id", id,
"method", r.Method,
"path", r.URL.Path,
"status", status,
"bytes", rec.bytes,
"duration_ms", time.Since(started).Milliseconds(),
quot;client_ip", clientAddress(r.RemoteAddr),
"backend", meta.backend,
)
})
}
func shed(limit int, next http.Handler) http.Handler {
slots := make(chan struct{}, limit)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
select {
case slots <- struct{}{}:
defer func() { <-slots }()
next.ServeHTTP(w, r)
default:
w.Header().Set("Retry-After", "1")
http.Error(w, "proxy capacity exhausted", http.StatusServiceUnavailable)
}
})
}
func checkBackend(ctx context.Context, client *http.Client, b *backend) bool {
u := *b.url
u.Path = "/healthz"
u.RawQuery = ""
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
if err != nil {
return false
}
resp, err := client.Do(req)
if err != nil {
return false
}
defer resp.Body.Close()
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 1024))
return resp.StatusCode >= 200 && resp.StatusCode < 300
}
func healthLoop(
ctx context.Context,
log *slog.Logger,
client *http.Client,
backends []*backend,
) {
run := func() {
for _, b := range backends {
ok := checkBackend(ctx, client, b)
previous := b.healthy.Swap(ok)
if previous != ok {
log.Info("backend_health_changed",
quot;backend", b.url.String(),
"healthy", ok,
)
}
}
}
run()
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
run()
}
}
}
func main() {
listen := flag.String("listen", "127.0.0.1:8080", "proxy listen address")
csv := flag.String(
"backends",
"http://127.0.0.1:9001,http://127.0.0.1:9002",
"comma-separated backend URLs",
)
maxInFlight := flag.Int("max-in-flight", 256, "concurrent request limit")
flag.Parse()
log := slog.New(slog.NewJSONHandler(os.Stdout, nil))
if *maxInFlight < 1 {
log.Error("max-in-flight must be positive")
os.Exit(2)
}
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{Timeout: 2 * time.Second}).DialContext,
ForceAttemptHTTP2: true,
MaxIdleConns: 256,
MaxIdleConnsPerHost: 64,
IdleConnTimeout: 60 * time.Second,
TLSHandshakeTimeout: 3 * time.Second,
ResponseHeaderTimeout: 5 * time.Second,
ExpectContinueTimeout: 1 * time.Second,
}
var backends []*backend
for _, raw := range strings.Split(*csv, ",") {
u, err := url.Parse(strings.TrimSpace(raw))
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") {
log.Error("invalid backend URL", "value", raw)
os.Exit(2)
}
b := &backend{
url: u,
breaker: breaker{
threshold: 5,
cooldown: 15 * time.Second,
},
}
b.proxy = &httputil.ReverseProxy{
Transport: transport,
Rewrite: func(pr *httputil.ProxyRequest) {
pr.SetURL(b.url)
pr.Out.Host = b.url.Host
pr.SetXForwarded()
pr.Out.Header.Set("X-Request-ID", pr.In.Header.Get("X-Request-ID"))
},
ModifyResponse: func(resp *http.Response) error {
p := resp.Request.Context().Value(permitKey).(permit)
p.backend.breaker.done(p.ticket, resp.StatusCode < 500)
return nil
},
ErrorHandler: func(w http.ResponseWriter, r *http.Request, err error) {
p := r.Context().Value(permitKey).(permit)
if r.Context().Err() != nil {
p.backend.breaker.cancel(p.ticket)
w.WriteHeader(499)
return
}
p.backend.breaker.done(p.ticket, false)
log.Warn("upstream_error", "backend", b.url.String(), "error", err)
http.Error(w, "bad gateway", http.StatusBadGateway)
},
ErrorLog: slog.NewLogLogger(log.Handler(), slog.LevelError),
}
backends = append(backends, b)
}
sel := &selector{backends: backends}
proxyHandler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
b, ticket, ok := sel.choose()
if !ok {
w.Header().Set("Retry-After", "1")
http.Error(w, "no backend available", http.StatusServiceUnavailable)
return
}
if meta, ok := r.Context().Value(metaKey).(*requestMeta); ok {
meta.backend = b.url.Host
}
ctx := context.WithValue(r.Context(), permitKey, permit{b, ticket})
b.proxy.ServeHTTP(w, r.WithContext(ctx))
})
healthCtx, stopHealth := context.WithCancel(context.Background())
healthClient := &http.Client{Timeout: time.Second}
go healthLoop(healthCtx, log, healthClient, backends)
mux := http.NewServeMux()
mux.HandleFunc("/livez", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
})
mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) {
for _, b := range backends {
if b.healthy.Load() && b.breaker.ready(time.Now()) {
w.WriteHeader(http.StatusNoContent)
return
}
}
http.Error(w, "no backend ready", http.StatusServiceUnavailable)
})
mux.Handle("/", accessLog(log, shed(*maxInFlight, proxyHandler)))
server := &http.Server{
Addr: *listen,
Handler: mux,
ReadHeaderTimeout: 5 * time.Second,
ReadTimeout: 30 * time.Second,
WriteTimeout: 30 * time.Second,
IdleTimeout: 60 * time.Second,
MaxHeaderBytes: 1 << 20,
}
signalCtx, stopSignals := signal.NotifyContext(
context.Background(), syscall.SIGINT, syscall.SIGTERM,
)
defer stopSignals()
errs := make(chan error, 1)
go func() {
log.Info("proxy_started", "listen", *listen)
errs <- server.ListenAndServe()
}()
select {
case <-signalCtx.Done():
log.Info("shutdown_started")
case err := <-errs:
if err != nil && err != http.ErrServerClosed {
log.Error("server_failed", "error", err)
}
}
stopHealth()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := server.Shutdown(shutdownCtx); err != nil {
log.Error("shutdown_incomplete", "error", err)
}
transport.CloseIdleConnections()
}
Резултатот од активната здравствена проверка и состојбата на circuit breaker-от остануваат одделни. Backend-системот може да одговара на здравствени проверки додека неговите апликациски барања откажуваат; circuit breaker-от сепак привремено го отстранува. Обратно, неуспешна здравствена проверка веднаш спречува избор без да ја наруши историјата на колото.
Создајте контролирани тест backend-системи
Овој мал сервер обезбедува здрава крајна точка и опционален, детерминистички 503 одговор. Зачувајте го како cmd/testbackend/main.go.
package main
import (
"flag"
"fmt"
"log"
"net/http"
"sync/atomic"
)
func main() {
listen := flag.String("listen", "127.0.0.1:9001", "listen address")
name := flag.String("name", "backend-1", "response name")
failEvery := flag.Uint64("fail-every", 0, "return 503 every N requests")
flag.Parse()
var requests atomic.Uint64
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNoContent)
})
mux.HandleFunc("/", func(w http.ResponseWriter, _ *http.Request) {
n := requests.Add(1)
if *failEvery > 0 && n%*failEvery == 0 {
http.Error(w, "simulated failure", http.StatusServiceUnavailable)
return
}
fmt.Fprintf(w, "%s request=%d\n", *name, n)
})
log.Printf("test backend listening on %s", *listen)
log.Fatal(http.ListenAndServe(*listen, mux))
}
Од три терминали, извршете:
go run ./cmd/testbackend -listen=127.0.0.1:9001 -name=backend-1
go run ./cmd/testbackend -listen=127.0.0.1:9002 -name=backend-2 -fail-every=3
go run ./cmd/edgeproxy -listen=127.0.0.1:8080 \
-backends=http://127.0.0.1:9001,http://127.0.0.1:9002 \
-max-in-flight=64
Проверете го однесувањето при неуспех
Прво потврдете ги крајните точки за животниот циклус и round-robin насочувањето:
curl --fail --silent --show-error -o /dev/null \
http://127.0.0.1:8080/readyz
for n in 1 2 3 4 5 6 7 8; do
curl --silent --show-error -i http://127.0.0.1:8080/demo
done
Запрете еден тест backend. Неговата здравствена состојба се менува по следната проверка, додека барањата продолжуваат преку преостанатиот backend. Рестартирајте го и дозволете еден здравствен интервал за повторен прием.
За да го испробате отфрлањето на товар, намалете го -max-in-flight и насочете го проксито кон backend обработувач што намерно блокира. Вишокот барања треба веднаш да добијат 503 со Retry-After: 1; не смеат да чекаат во неограничена внатрешнопроцесна редица.
Breaker-от се отвора по пет релевантни неуспеси во една генерација. Останува отворен 15 секунди, дозволува една half-open проба и се затвора само ако таа проба успее. Бидејќи секој автоматски повторен обид троши дополнителен backend капацитет, клиентите треба да применуваат ограничени повторни обиди со jitter и да повторуваат само операции за кои знаат дека се безбедни.
Безбедност, перформанси и набљудливост
Чувајте го проксито на приватна адреса кога TLS завршува на доверен frontend. Ако мора директно да прифаќа јавен сообраќај, додајте TLS, ограничете ги дозволените методи и големините на телата и поставете ги административните крајни точки на посебен заштитен listener. Firewall на хостот треба да ја дозволи јавната frontend порта, а backend портите да ги ограничи на хостот на проксито или приватната мрежа.
ProxyRequest.SetXForwarded е важен: Go ги отстранува недоверливите forwarding заглавија пред да се изврши функцијата за препишување, а потоа генерира контролирани вредности. Не копирајте произволни влезни X-Forwarded-For заглавија сами.
Транспортот одделно ги ограничува воспоставувањето врска, TLS преговарањето и чекањето на заглавието на одговорот. Timeout за поврзување не ограничува бавен одговор. Буџетите за читање и пишување на серверот се 30 секунди, па оваа конфигурација е наменета за вообичаен API сообраќај. Долготрајните текови и големите отпраќања бараат буџети специфични за рута, наместо глобално оневозможување на заштитата.
JSON access записите содржат ID на барањето, избран backend, статус, бајти во одговорот, адреса на клиентот и траење. Здравствените премини и транспортните неуспеси се одделни настани. Во поголем систем, изведете бројачи и хистограми за латентност од овие настани или инструментирајте ги истите граници со вашиот metrics stack. Избегнувајте логирање на заглавија за авторизација, колачиња, тајни во query-параметри или тела.
Распоредување со systemd
Изградете како непривилегиран корисник, проверете ја дестинацијата пред да замените постоечка бинарна датотека, а потоа инсталирајте го изричниот артефакт:
mkdir -p ./bin
go test ./...
go build -trimpath -o ./bin/edgeproxy ./cmd/edgeproxy
sudo install -m 0755 ./bin/edgeproxy /usr/local/bin/edgeproxy
Зачувајте ја оваа единица како /etc/systemd/system/edgeproxy.service. Создавањето на таа датотека и управувањето со сервисот бараат root привилегии.
[Unit]
Description=Go edge reverse proxy
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=edgeproxy
Group=edgeproxy
ExecStart=/usr/local/bin/edgeproxy -listen=127.0.0.1:8080 -backends=http://127.0.0.1:9001,http://127.0.0.1:9002 -max-in-flight=256
Restart=on-failure
RestartSec=2
NoNewPrivileges=true
PrivateTmp=true
ProtectSystem=strict
ProtectHome=true
ProtectKernelTunables=true
ProtectKernelModules=true
ProtectControlGroups=true
RestrictSUIDSGID=true
LockPersonality=true
MemoryDenyWriteExecute=true
RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6
CapabilityBoundingSet=
AmbientCapabilities=
[Install]
WantedBy=multi-user.target
Создајте го наменскиот системски корисник според локалната политика за сметки, потврдете ја единицата со systemd-analyze verify /etc/systemd/system/edgeproxy.service, а потоа користете systemctl daemon-reload, systemctl enable --now edgeproxy и journalctl -u edgeproxy. Ако се потребни upstream DNS, откривање сервиси или Unix сокети, намерно прилагодете ги претпоставките за адресното семејство и мрежата.
Вообичаени продукциски неуспеси
- Секој backend започнува недостапен: ова е намерно сè додека не успее првата активна здравствена проверка. Нека readiness probe-овите го толерираат времето на стартување.
- Здравите крајни точки маскираат расипани апликациски патеки: нека
/healthzги проверува само зависностите потребни за опслужување сообраќај, додека се потпирате на исходите од breaker-от за неуспеси на патеките на барањата. - 503 одговори се појавуваат при скокови на сообраќајот: утврдете дали дошле од отфрлање на товар, избор на backend или upstream одговор. Пораката и backend полето во логот ги разликуваат.
- Streaming одговорите завршуваат по 30 секунди: глобалниот timeout за пишување е несоодветен за WebSockets, server-sent настани и неограничени преземања. Дајте му на таквиот сообраќај посебен сервер или изрична политика.
- Гасењето надминува десет секунди: прифатено барање сè уште е блокирано. Усогласете ги ограничувањата на upstream одговорите, timeout-ите на серверот, времето за празнење на frontend и systemd буџетот за запирање.
Конечна контролна листа за верификација
- Двете upstream здравствени URL-адреси враќаат ограничен 2xx одговор.
- Нормалните барања се менуваат меѓу здравите backend-системи.
- Запрен backend исчезнува без да го запре проксито.
- Повторените 5xx или транспортни неуспеси го отвораат колото само на тој backend.
- Исцрпувањето на капацитетот враќа непосредни 503 одговори.
- Логовите содржат ID на барања, траења, статуси и backend идентитети без тајни.
- SIGTERM ги испразнува прифатените барања во рамките на буџетот за гасење.
- Само наменетиот frontend може да достигне proxy и backend порти.
Отпорното прокси не се дефинира со тоа колку барања може да прифати. Се дефинира со тоа колку внимателно одбива работа, колку тесно го ограничува неуспехот и колку јасно ја објаснува секоја одлука потоа. Тие својства претвораат тенка мрежна компонента во доверлива продукциска граница.