Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 34 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ curl http://localhost:8080/<your-bucket>/<your-object>
- Optional `Access-Control-Allow-Origin` header (`-cors-origin`) for simple CORS use cases
- Structured logging via `log/slog` with text or JSON output (`-log-format`)
- `/_health` endpoint for liveness/readiness probes
- Graceful shutdown on SIGTERM/SIGINT with configurable drain windows (`-shutdown-delay`, `-shutdown-idle-grace-period`, `-shutdown-timeout`) for zero-downtime deploys

## Installation

Expand Down Expand Up @@ -93,6 +94,12 @@ Usage of gcsproxy:
Minimum log level: debug, info, warn, or error. (default "info")
-not-found string
Object served with HTTP 404 for unmatched routes.
-shutdown-delay duration
Delay after SIGTERM/SIGINT during which requests are served normally but responses carry Connection: close.
-shutdown-idle-grace-period duration
After the listener closes, keep idle connections open for up to this long; ends early once no connections remain.
-shutdown-timeout duration
Max wait for in-flight requests after the idle grace period; 0 waits indefinitely. (default 30s)
-spa
SPA fallback: serve -i from the bucket root with HTTP 200 for unmatched routes.
-v Show access log.
Expand Down Expand Up @@ -232,7 +239,30 @@ Edge cases:

### Health check

`/_health` returns `200 OK` with the body `OK`. It does not call GCS and is safe to use as a Kubernetes/Cloud Run liveness or readiness probe.
`/_health` returns `200 OK` with the body `OK`. It does not call GCS and is safe to use as a Kubernetes/Cloud Run liveness or readiness probe. It keeps returning `200` during graceful shutdown, so a liveness probe won't kill the process mid-drain.

### Graceful shutdown

On SIGTERM or SIGINT, gcsproxy shuts down in four phases designed to eliminate 502/503s and connection resets during rolling deploys:

1. **Drain** (`-shutdown-delay`): nothing stops — new connections are accepted and requests are served normally, but every response carries `Connection: close`, telling keep-alive clients (e.g. an nginx `upstream keepalive` pool) to retire the connection after use.
2. **Stop accepting**: the listener closes; new TCP connections are refused. Existing connections keep working.
3. **Idle grace** (`-shutdown-idle-grace-period`): existing idle connections stay open, because a client may send a request at the exact moment the server would close one (that race is what causes resets). Requests arriving on them are still served, each response again marked `Connection: close`. This phase ends as soon as no connections remain open, so it never waits longer than necessary.
4. **Close** (`-shutdown-timeout`): remaining idle connections are closed and gcsproxy waits for in-flight requests to finish — at most `-shutdown-timeout` (default `30s`; `0` waits indefinitely) — then exits 0.

`-shutdown-delay` and `-shutdown-idle-grace-period` default to `0s`, which still gives a basic graceful shutdown (in-flight requests complete before exit). Recommended production values:

```
gcsproxy -shutdown-delay 5s -shutdown-idle-grace-period 5s -shutdown-timeout 30s
```

Sizing guidance:

- Set `-shutdown-idle-grace-period` above your clients' idle-connection reuse timeout (nginx upstream `keepalive_timeout`, Go's `Transport.IdleConnTimeout`, etc.) so every pooled connection is either used once more (and cleanly closed) or closed by the client before the server closes it.
- Make sure your supervisor's kill timeout (systemd `TimeoutStopSec`, Kubernetes `terminationGracePeriodSeconds`) exceeds `shutdown-delay + shutdown-idle-grace-period + shutdown-timeout`, otherwise the process is SIGKILLed mid-drain.
- A second SIGTERM/SIGINT terminates the process immediately.

Note that even with the default values, SIGTERM now waits for in-flight requests (up to `-shutdown-timeout`) instead of exiting immediately; supervisors backstop this with SIGKILL after their kill timeout.

### Authentication

Expand Down Expand Up @@ -275,8 +305,10 @@ Wants=network-online.target

[Service]
Type=simple
ExecStart=/opt/gcsproxy/gcsproxy -v
ExecStart=/opt/gcsproxy/gcsproxy -v -shutdown-delay 5s -shutdown-idle-grace-period 5s -shutdown-timeout 30s
Restart=on-failure
# Must exceed shutdown-delay + shutdown-idle-grace-period + shutdown-timeout.
TimeoutStopSec=45

[Install]
WantedBy=multi-user.target
Expand Down
4 changes: 3 additions & 1 deletion gcsproxy.service
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,10 @@ Description=gcsproxy

[Service]
Type=simple
ExecStart=/opt/gcsproxy/gcsproxy -v
ExecStart=/opt/gcsproxy/gcsproxy -v -shutdown-timeout 30s
ExecStop=/bin/kill -SIGTERM $MAINPID
# Must exceed shutdown-delay + shutdown-idle-grace-period + shutdown-timeout.
TimeoutStopSec=45

[Install]
WantedBy = multi-user.target
197 changes: 171 additions & 26 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,16 @@ import (
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/url"
"os"
"os/signal"
"strconv"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"

"cloud.google.com/go/auth/credentials"
Expand All @@ -21,16 +26,20 @@ import (
)

type Server struct {
addr string
client *storage.Client
defaultIndex string
walkUpIndex bool
sourceBucket string
spa bool
notFoundPath string
contentLength bool
corsOrigin string
verbose bool
addr string
client *storage.Client
defaultIndex string
walkUpIndex bool
sourceBucket string
spa bool
notFoundPath string
contentLength bool
corsOrigin string
verbose bool
shutdownDelay time.Duration
idleGracePeriod time.Duration
shutdownTimeout time.Duration
draining atomic.Bool
}

func main() {
Expand All @@ -47,6 +56,9 @@ func main() {
logLevel = flag.String("log-level", "info", "Minimum log level: debug, info, warn, or error.")
contentLength = flag.Bool("content-length", false, "Send the Content-Length header (disables chunked transfer).")
corsOrigin = flag.String("cors-origin", "", "Value for the Access-Control-Allow-Origin header.")
shutdownDelay = flag.Duration("shutdown-delay", 0, "Delay after SIGTERM/SIGINT during which requests are served normally but responses carry Connection: close.")
idleGracePeriod = flag.Duration("shutdown-idle-grace-period", 0, "After the listener closes, keep idle connections open for up to this long; ends early once no connections remain.")
shutdownTimeout = flag.Duration("shutdown-timeout", 30*time.Second, "Max wait for in-flight requests after the idle grace period; 0 waits indefinitely.")
)
flag.Parse()

Expand Down Expand Up @@ -74,6 +86,9 @@ func main() {
if *spa && *notFoundPath != "" {
fatal("-spa and -not-found are mutually exclusive")
}
if *shutdownDelay < 0 || *idleGracePeriod < 0 || *shutdownTimeout < 0 {
fatal("-shutdown-delay, -shutdown-idle-grace-period and -shutdown-timeout must not be negative")
}

ctx := context.Background()
var opts []option.ClientOption
Expand All @@ -91,28 +106,158 @@ func main() {
if err != nil {
fatal("failed to create client", "err", err)
}
defer client.Close()

s := &Server{
addr: *bind,
client: client,
defaultIndex: *defaultIndex,
walkUpIndex: *walkUpIndex,
sourceBucket: *sourceBucket,
spa: *spa,
notFoundPath: *notFoundPath,
contentLength: *contentLength,
corsOrigin: *corsOrigin,
verbose: *verbose,
}

if err := s.ListenAndServe(); err != nil {
addr: *bind,
client: client,
defaultIndex: *defaultIndex,
walkUpIndex: *walkUpIndex,
sourceBucket: *sourceBucket,
spa: *spa,
notFoundPath: *notFoundPath,
contentLength: *contentLength,
corsOrigin: *corsOrigin,
verbose: *verbose,
shutdownDelay: *shutdownDelay,
idleGracePeriod: *idleGracePeriod,
shutdownTimeout: *shutdownTimeout,
}

sigCtx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stop()
// Unregister after the first signal so a second SIGTERM/SIGINT kills the
// process immediately via the default disposition.
context.AfterFunc(sigCtx, stop)

if err := s.ListenAndServe(sigCtx); err != nil {
fatal("server exited", "err", err)
}
}

func (s *Server) ListenAndServe() error {
slog.Info("listening", "addr", s.addr)
return http.ListenAndServe(s.addr, s.handler())
func (s *Server) ListenAndServe(ctx context.Context) error {
ln, err := net.Listen("tcp", s.addr)
if err != nil {
return err
}
slog.Info("listening", "addr", ln.Addr().String())
return s.serve(ctx, ln, s.handler())
}

// serve runs h on ln until ctx is canceled, then executes the phased graceful
// shutdown: drain (Connection: close) -> close listener -> idle grace ->
// Shutdown. Idle connections stay open until the idle grace period elapses so
// clients never race a new request against a server-side close.
func (s *Server) serve(ctx context.Context, ln net.Listener, h http.Handler) error {
tracker := &connTracker{}
srv := &http.Server{Handler: s.drainingHandler(h), ConnState: tracker.connState}

errc := make(chan error, 1)
go func() { errc <- srv.Serve(ln) }()

select {
case err := <-errc:
return err
case <-ctx.Done():
}

// Phase 1: keep serving; every new response tells clients to stop reusing.
s.draining.Store(true)
slog.Info("shutdown: draining", "delay", s.shutdownDelay)
time.Sleep(s.shutdownDelay)

// Phase 2: stop accepting new connections. Existing conns keep working.
slog.Info("shutdown: closing listener")
ln.Close()
// Wait for Serve to return so it untracks its listener; otherwise Shutdown
// below can report a spurious double-close error.
if err := <-errc; err != nil && !errors.Is(err, net.ErrClosed) {
slog.Warn("shutdown: serve loop exited with error", "err", err)
}

// Phase 3: leave idle connections open; requests on them are still served.
// Skip the wait when no connections remain, and end it early once the
// last one closes.
slog.Info("shutdown: waiting before closing idle connections", "grace", s.idleGracePeriod, "open", tracker.open())
select {
case <-tracker.noneOpen():
slog.Info("shutdown: no connections remain")
case <-time.After(s.idleGracePeriod):
}

// Phase 4: close idle connections, wait for in-flight requests.
slog.Info("shutdown: closing idle connections", "timeout", s.shutdownTimeout)
sctx := context.Background()
if s.shutdownTimeout > 0 {
var cancel context.CancelFunc
sctx, cancel = context.WithTimeout(sctx, s.shutdownTimeout)
defer cancel()
}
if err := srv.Shutdown(sctx); err != nil {
srv.Close()
return fmt.Errorf("graceful shutdown incomplete: %w", err)
}
slog.Info("shutdown: complete")
return nil
}

// connTracker counts open server connections via http.Server.ConnState so
// the shutdown sequence can stop waiting as soon as none remain. Active
// connections count too: a response whose headers predate the drain flag
// carries no Connection: close, so its connection can still go idle and be
// reused by the client.
type connTracker struct {
mu sync.Mutex
count int
waiter chan struct{}
}

func (ct *connTracker) connState(_ net.Conn, state http.ConnState) {
ct.mu.Lock()
defer ct.mu.Unlock()
switch state {
case http.StateNew:
ct.count++
case http.StateHijacked, http.StateClosed:
ct.count--
if ct.count == 0 && ct.waiter != nil {
close(ct.waiter)
ct.waiter = nil
}
}
}

func (ct *connTracker) open() int {
ct.mu.Lock()
defer ct.mu.Unlock()
return ct.count
}

// noneOpen returns a channel that is closed once no connections remain open.
// Only valid after the listener stopped accepting new connections, since the
// count never rises again from zero.
func (ct *connTracker) noneOpen() <-chan struct{} {
ct.mu.Lock()
defer ct.mu.Unlock()
ch := make(chan struct{})
if ct.count == 0 {
close(ch)
} else {
ct.waiter = ch
}
return ch
}

// drainingHandler marks every response written after shutdown begins with
// Connection: close, so keep-alive clients stop reusing the connection.
// net/http then closes the connection after the response completes.
func (s *Server) drainingHandler(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if s.draining.Load() {
w.Header().Set("Connection", "close")
}
next.ServeHTTP(w, r)
})
}

func (s *Server) handler() http.Handler {
Expand Down
Loading