Skip to content
Merged
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
4 changes: 4 additions & 0 deletions config.example.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,10 @@ database:
# Example: "postgres://user:password@localhost:5432/proxy?sslmode=disable"
url: ""

# How often cache hit counts and last-access times are written, batched in
# one transaction. "0" writes each hit as it happens.
hit_flush_interval: "1s"

# Logging configuration
log:
# Minimum log level: "debug", "info", "warn", "error"
Expand Down
15 changes: 15 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,21 @@ database:
|--------|-------------|------|-------------|
| `database.url` | `PROXY_DATABASE_URL` | `-database-url` | PostgreSQL connection URL |

### Hit counts

Every cache hit updates the artifact's hit count and last-access time, which the stats pages and LRU eviction use. Rather than writing each hit as its own transaction, the proxy counts hits in memory and writes them together every `hit_flush_interval`. SQLite allows one writer at a time, and the proxy uses a single SQLite connection, so with a write per hit, concurrent downloads queue behind each other.

```yaml
database:
hit_flush_interval: "1s"
```

| Config | Environment | Flag | Description |
|--------|-------------|------|-------------|
| `database.hit_flush_interval` | `PROXY_DATABASE_HIT_FLUSH_INTERVAL` | - | How often batched hits are written (default `1s`). `0` writes each hit as it happens. |

Hits not yet written are lost if the process is killed; a normal shutdown writes them.

## Logging

```yaml
Expand Down
42 changes: 42 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -418,6 +418,12 @@ type DatabaseConfig struct {

// URL is the PostgreSQL connection string.
URL string `json:"url" yaml:"url"`

// HitFlushInterval is how often cache hit counts and last-access times
// are written, batched in one transaction. Uses Go duration syntax
// (e.g. "1s", "10s"). Default: "1s". Set to "0" to write each hit as it
// happens.
HitFlushInterval string `json:"hit_flush_interval" yaml:"hit_flush_interval"`
}

// String returns a human-readable description of the configured database
Expand Down Expand Up @@ -896,6 +902,7 @@ func (c *Config) LoadFromEnv() {
setEnvString(&c.Database.Driver, "PROXY_DATABASE_DRIVER")
setEnvString(&c.Database.Path, "PROXY_DATABASE_PATH")
setEnvString(&c.Database.URL, "PROXY_DATABASE_URL")
setEnvString(&c.Database.HitFlushInterval, "PROXY_DATABASE_HIT_FLUSH_INTERVAL")
setEnvString(&c.Log.Level, "PROXY_LOG_LEVEL")
setEnvString(&c.Log.Format, "PROXY_LOG_FORMAT")
setEnvString(&c.AccessLog.Path, "PROXY_ACCESS_LOG_PATH")
Expand Down Expand Up @@ -1038,6 +1045,10 @@ func (c *Config) Validate() error {
return err
}

if err := validateHitFlushInterval(c.Database.HitFlushInterval); err != nil {
return err
}

return c.validateComponents()
}

Expand Down Expand Up @@ -1120,6 +1131,7 @@ const (
defaultMetadataTTL = 5 * time.Minute //nolint:mnd // sensible default
defaultDirectServeTTL = 15 * time.Minute //nolint:mnd // sensible default
defaultHTTPTimeout = 30 * time.Second //nolint:mnd // sensible default
defaultHitFlushInterval = time.Second
defaultMetadataMaxSize = 100 << 20
defaultGradleBuildCacheMaxUploadSize = 100 << 20
defaultGradleBuildCacheSweepInterval = 10 * time.Minute
Expand Down Expand Up @@ -1198,6 +1210,36 @@ func (c *Config) ParseHTTPTimeout() time.Duration {
return d
}

func validateHitFlushInterval(s string) error {
if s == "" || s == "0" {
return nil
}
d, err := time.ParseDuration(s)
if err != nil {
return fmt.Errorf("invalid database.hit_flush_interval %q: %w", s, err)
}
if d < 0 {
return fmt.Errorf("invalid database.hit_flush_interval %q: must be non-negative", s)
}
return nil
}

// ParseHitFlushInterval returns how often batched cache hits are written.
// Returns 1 second if unset or invalid, 0 if explicitly disabled.
func (c *Config) ParseHitFlushInterval() time.Duration {
if c.Database.HitFlushInterval == "" {
return defaultHitFlushInterval
}
if c.Database.HitFlushInterval == "0" {
return 0
}
d, err := time.ParseDuration(c.Database.HitFlushInterval)
if err != nil || d < 0 {
return defaultHitFlushInterval
}
return d
}

// ParseMetadataTTL returns the metadata TTL duration.
// Returns 5 minutes if unset, 0 if explicitly disabled.
func (c *Config) ParseMetadataTTL() time.Duration {
Expand Down
55 changes: 55 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,61 @@ func TestLoadHTTPTimeoutFromEnv(t *testing.T) {
}
}

func TestParseHitFlushInterval(t *testing.T) {
tests := []struct {
name string
interval string
want time.Duration
}{
{"empty defaults to 1s", "", time.Second},
{"explicit zero disables", "0", 0},
{"10 seconds", "10s", 10 * time.Second},
{"invalid defaults to 1s", "not-a-duration", time.Second},
{"negative defaults to 1s", "-5s", time.Second},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
cfg := Default()
cfg.Database.HitFlushInterval = tt.interval
got := cfg.ParseHitFlushInterval()
if got != tt.want {
t.Errorf("ParseHitFlushInterval() = %v, want %v", got, tt.want)
}
})
}
}

func TestValidateHitFlushInterval(t *testing.T) {
cfg := Default()
cfg.Database.HitFlushInterval = "not-a-duration"
if err := cfg.Validate(); err == nil {
t.Error("expected validation error for invalid hit_flush_interval")
}

cfg.Database.HitFlushInterval = "-5s"
if err := cfg.Validate(); err == nil {
t.Error("expected validation error for negative hit_flush_interval")
}

for _, ok := range []string{"10s", "0", ""} {
cfg.Database.HitFlushInterval = ok
if err := cfg.Validate(); err != nil {
t.Errorf("unexpected error for hit_flush_interval %q: %v", ok, err)
}
}
}

func TestLoadHitFlushIntervalFromEnv(t *testing.T) {
cfg := Default()
t.Setenv("PROXY_DATABASE_HIT_FLUSH_INTERVAL", "10s")
cfg.LoadFromEnv()

if cfg.Database.HitFlushInterval != "10s" {
t.Errorf("HitFlushInterval = %q, want %q", cfg.Database.HitFlushInterval, "10s")
}
}

func TestLoadMetadataTTLFromEnv(t *testing.T) {
cfg := Default()
t.Setenv("PROXY_METADATA_TTL", "10m")
Expand Down
1 change: 1 addition & 0 deletions internal/database/database.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ type DB struct {
*sqlx.DB
dialect Dialect
path string
hits *hitBatch
}

func (db *DB) Dialect() Dialect {
Expand Down
146 changes: 146 additions & 0 deletions internal/database/hits.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
package database

import (
"cmp"
"log/slog"
"slices"
"sync"
"time"
)

// Recording each cache hit as its own UPDATE makes every cached download a
// write transaction, and SQLite runs on a single connection, so concurrent
// downloads queue behind those writes. BatchHits counts hits in memory and
// writes them in one transaction per interval instead.

type hitKey struct{ versionPURL, filename string }

type hitEntry struct {
count int64
last time.Time
}

type hitBatch struct {
mu sync.Mutex
pending map[hitKey]hitEntry
stop chan struct{}
stopOnce sync.Once
done chan struct{}
}

func (b *hitBatch) add(k hitKey, e hitEntry) {
b.mu.Lock()
defer b.mu.Unlock()
cur := b.pending[k]
cur.count += e.count
if e.last.After(cur.last) {
cur.last = e.last
}
b.pending[k] = cur
}

func (b *hitBatch) take() map[hitKey]hitEntry {
b.mu.Lock()
defer b.mu.Unlock()
pending := b.pending
b.pending = map[hitKey]hitEntry{}
return pending
}

// BatchHits makes RecordArtifactHit buffer hits and write them every
// interval. An interval of zero or less leaves every hit written immediately.
// Close writes whatever is still pending. Failed writes are logged to logger.
func (db *DB) BatchHits(interval time.Duration, logger *slog.Logger) {
if interval <= 0 || db.hits != nil {
return
}
if logger == nil {
logger = slog.New(slog.DiscardHandler)
}
b := &hitBatch{pending: map[hitKey]hitEntry{}, stop: make(chan struct{}), done: make(chan struct{})}
db.hits = b
go func() {
defer close(b.done)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
if err := db.flushHits(); err != nil {
logger.Warn("failed to write cache hits, retrying next flush", "error", err)
}
case <-b.stop:
if err := db.flushHits(); err != nil {
logger.Error("failed to write cache hits on close, dropping them", "error", err)
}
return
}
}
}()
}

// flushHits writes the pending hits in one transaction. If the transaction
// fails they are put back for the next flush, so a failed write delays
// counts rather than losing them.
func (db *DB) flushHits() error {
pending := db.hits.take()
if len(pending) == 0 {
return nil
}

err := db.writeHits(pending)
if err != nil {
for k, e := range pending {
db.hits.add(k, e)
}
}
return err
}

func (db *DB) writeHits(pending map[hitKey]hitEntry) error {
tx, err := db.Beginx()
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()

// Timestamps only move forward: with proxies sharing a database, an older
// batch can flush after a newer hit is already written.
stmt, err := tx.Preparex(db.Rebind(`
UPDATE artifacts
SET hit_count = hit_count + ?,
last_accessed_at = CASE WHEN last_accessed_at IS NULL OR last_accessed_at < ? THEN ? ELSE last_accessed_at END,
updated_at = CASE WHEN updated_at IS NULL OR updated_at < ? THEN ? ELSE updated_at END
WHERE version_purl = ? AND filename = ?
`))
if err != nil {
return err
}
defer func() { _ = stmt.Close() }()

// A fixed order keeps proxies sharing a Postgres database from locking
// the same rows in opposite orders and deadlocking.
keys := make([]hitKey, 0, len(pending))
for k := range pending {
keys = append(keys, k)
}
slices.SortFunc(keys, func(a, b hitKey) int {
return cmp.Or(cmp.Compare(a.versionPURL, b.versionPURL), cmp.Compare(a.filename, b.filename))
})
for _, k := range keys {
e := pending[k]
if _, err := stmt.Exec(e.count, e.last, e.last, e.last, e.last, k.versionPURL, k.filename); err != nil {
return err
}
}
return tx.Commit()
}

// Close writes any batched hits, then closes the database.
func (db *DB) Close() error {
if b := db.hits; b != nil {
b.stopOnce.Do(func() { close(b.stop) })
<-b.done
}
return db.DB.Close()
}
Loading
Loading