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
21 changes: 19 additions & 2 deletions .github/workflows/scripts/test-framework.sh
Original file line number Diff line number Diff line change
Expand Up @@ -26,15 +26,32 @@ trap cleanup_docker EXIT
echo "🔧 Starting dependencies of framework tests..."
# Use docker compose (v2) if available, fallback to docker-compose (v1)
if command -v docker-compose >/dev/null 2>&1; then
docker-compose -f tests/docker-compose.yml up -d
COMPOSE="docker-compose"
elif docker compose version >/dev/null 2>&1; then
docker compose -f tests/docker-compose.yml up -d
COMPOSE="docker compose"
else
echo "❌ Neither docker-compose nor docker compose is available"
exit 1
fi
$COMPOSE -f tests/docker-compose.yml up -d
sleep 20

# The framework logstore tests fail (not skip) in CI when ClickHouse is
# unreachable, and `up -d` does not wait for health, so gate on its /ping.
echo "⏳ Waiting for ClickHouse to become ready..."
for attempt in $(seq 1 60); do
if $COMPOSE -f tests/docker-compose.yml exec -T clickhouse wget --spider -q http://127.0.0.1:8123/ping 2>/dev/null; then
echo "✅ ClickHouse is ready"
break
fi
if [ "$attempt" -eq 60 ]; then
echo "❌ ClickHouse did not become ready within 120s"
$COMPOSE -f tests/docker-compose.yml logs --tail=50 clickhouse || true
exit 1
fi
sleep 2
done

# Validate framework build
echo "🔨 Validating framework build..."
cd framework
Expand Down
2 changes: 1 addition & 1 deletion docs/deployment-guides/config-json/storage.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -279,7 +279,7 @@ ClickHouse is a **`logs_store`-only** backend. The `config_store` supports only
```

<Note>
With ClickHouse, `retention_days` is enforced by a native table **TTL** rather than a background delete job. Setting it to `0` (or omitting it) leaves the TTL unset, so ClickHouse itself never expires rows. Background cleanup is controlled separately by `client_config.log_retention_days`, so set that to your desired horizon as well. `matview_refresh_interval` and `matview_refresh_timeout` do not apply — they are PostgreSQL-only settings for materialized views.
With ClickHouse, `retention_days` is enforced by a native table **TTL** rather than a background delete job. The TTL is reconciled on every startup, so changing `retention_days` later updates the existing `logs`, `mcp_tool_logs` and `webhook_deliveries` tables (metadata only; expired rows are dropped by ClickHouse's regular TTL merges). Setting it to `0` (or omitting it) leaves any existing TTL untouched and creates new tables without one, so ClickHouse itself never expires rows unless you manage the TTL yourself. Background cleanup is controlled separately by `client_config.log_retention_days`; on ClickHouse each run is a single lightweight `DELETE FROM ... WHERE created_at < cutoff`, as are the periodic sweeps of stale `processing` rows, so setting both is safe. Lightweight deletes require ClickHouse 24.4 or newer. `matview_refresh_interval` and `matview_refresh_timeout` do not apply - they are PostgreSQL-only settings for materialized views.
</Note>

### Disabled
Expand Down
1 change: 1 addition & 0 deletions framework/changelog.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- fix: ClickHouse log store deletes no longer run as heavyweight `ALTER TABLE ... DELETE` mutations. The retention cleaner issued one such mutation per 100 rows, each rewriting the whole current-month part, and the once-a-minute stale-`processing` sweeps issued one per table unconditionally, filling replica disks in minutes. Every delete on the ClickHouse store (retention sweep, `Flush`/`FlushMCPToolLogs`, UI log deletes, async job and webhook delivery expiry) is now a single lightweight `DELETE FROM ... WHERE` per run, skipped entirely when nothing matches. The table TTL derived from `logs_store.retention_days` is now reconciled on every startup with `MODIFY TTL` (metadata only), so changing the value reaches existing tables; `0` leaves an existing TTL untouched (#7098)
8 changes: 6 additions & 2 deletions framework/logstore/cleaner.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,8 +141,12 @@ func (c *LogsCleaner) cleanupOldLogs(ctx context.Context) {
batchCount++
c.logger.Debug("deleted batch %d: %d logs", batchCount, deleted)

// If we deleted fewer than the batch size, we're done
if deleted < int64(batchSize) {
// A full batch means more rows may remain; anything else means the
// store is done. The SQL stores return at most batchSize. ClickHouse
// deletes the whole expired range in one lightweight statement and
// returns that count, so a count above batchSize must end the loop
// too instead of re-issuing the delete (#7098).
if deleted != int64(batchSize) {
break
}
}
Expand Down
41 changes: 39 additions & 2 deletions framework/logstore/clickhouse.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"fmt"
"net"
"net/url"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -144,9 +145,33 @@ func buildClickHouseDSN(config *ClickHouseConfig) (string, error) {
return u.String(), nil
}

// chMinServerVersion is the oldest ClickHouse release the store supports.
// chLightweightDelete relies on the lightweight_deletes_sync setting, which
// ClickHouse introduced in 24.4; older servers reject every delete with
// UNKNOWN_SETTING, so they are refused at startup instead of on the first
// sweep.
const chMinServerVersion = "24.4"

// chServerVersionSupported parses a ClickHouse version() string such as
// "26.6.1.1193" and reports whether it is at least chMinServerVersion.
func chServerVersionSupported(version string) (bool, error) {
parts := strings.Split(strings.TrimSpace(version), ".")
if len(parts) < 2 {
return false, fmt.Errorf("clickhouse: unrecognised server version %q", version)
}
major, err1 := strconv.Atoi(parts[0])
minor, err2 := strconv.Atoi(parts[1])
if err1 != nil || err2 != nil {
return false, fmt.Errorf("clickhouse: unrecognised server version %q", version)
}
return major > 24 || (major == 24 && minor >= 4), nil
}

// newClickHouseLogStore creates a new ClickHouse log store. retentionDays drives
// the table TTL; values < 1 leave TTL unset (the LogsCleaner still prunes via
// DeleteLogsBatch).
// the table TTL, which is reconciled on every start so a changed value reaches
// existing tables; values < 1 leave any TTL untouched. Independently, the
// LogsCleaner (client_config.log_retention_days) prunes with a single
// lightweight delete per run.
func newClickHouseLogStore(ctx context.Context, config *ClickHouseConfig, retentionDays int, logger schemas.Logger) (LogStore, error) {
dsn, err := buildClickHouseDSN(config)
if err != nil {
Expand Down Expand Up @@ -181,6 +206,18 @@ func newClickHouseLogStore(ctx context.Context, config *ClickHouseConfig, retent
return nil, fmt.Errorf("clickhouse ping failed: %w", err)
}

var serverVersion string
if err := db.WithContext(ctx).Raw("SELECT version()").Scan(&serverVersion).Error; err != nil {
logger.Error("logstore: failed to read clickhouse server version: %v", err)
return nil, fmt.Errorf("clickhouse: read server version: %w", err)
}
if ok, err := chServerVersionSupported(serverVersion); err != nil {
return nil, err
} else if !ok {
return nil, fmt.Errorf("clickhouse: server version %s is not supported; Bifrost requires ClickHouse %s or newer (lightweight DELETE with lightweight_deletes_sync)", serverVersion, chMinServerVersion)
}
logger.Info("logstore: clickhouse server version %s", serverVersion)

logger.Info("logstore: running clickhouse schema migrations")
if err := triggerClickHouseMigrations(ctx, db, config.Cluster, retentionDays, logger); err != nil {
logger.Error("logstore: clickhouse schema migrations failed: %v", err)
Expand Down
98 changes: 92 additions & 6 deletions framework/logstore/clickhousemigrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ import (
"context"
"fmt"
"reflect"
"regexp"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -223,15 +225,87 @@ func clickhouseReconcileColumns(ctx context.Context, db *gorm.DB, model any, tab
}

// chLogsTTL derives the logs/mcp_tool_logs TTL clause from the configured
// retention. Values < 1 leave TTL unset (the LogsCleaner still prunes via
// DeleteLogsBatch).
// retention. Values < 1 return "" which means the TTL is not managed by
// Bifrost: new tables get none and existing tables keep whatever they have
// (the LogsCleaner still prunes with one lightweight delete per run).
func chLogsTTL(retentionDays int) string {
if retentionDays < 1 {
return ""
}
return fmt.Sprintf("toDateTime(created_at) + INTERVAL %d DAY", retentionDays)
}

// chTTLClauseDaysRe matches the clause form chLogsTTL and the fixed backstops
// emit for CREATE TABLE / MODIFY TTL.
var chTTLClauseDaysRe = regexp.MustCompile(`^toDateTime\(created_at\) \+ INTERVAL (\d+) DAY$`)

// chTTLEngineFullDaysRe matches the normalized form ClickHouse stores in
// system.tables.engine_full: `INTERVAL n DAY` becomes `toIntervalDay(n)`.
var chTTLEngineFullDaysRe = regexp.MustCompile(`TTL toDateTime\(created_at\) \+ toIntervalDay\((\d+)\)`)

func chTTLDaysFromMatch(re *regexp.Regexp, s string) (int, bool) {
m := re.FindStringSubmatch(s)
if m == nil {
return 0, false
}
days, err := strconv.Atoi(m[1])
if err != nil {
return 0, false
}
return days, true
}

// chTTLDaysFromClause returns the retention days a TTL clause built by
// chLogsTTL (or a fixed backstop) encodes. ok is false for "" and for any
// expression Bifrost does not manage.
func chTTLDaysFromClause(ttl string) (days int, ok bool) {
return chTTLDaysFromMatch(chTTLClauseDaysRe, ttl)
}

// chTTLDaysFromEngineFull returns the retention days currently applied to a
// table, read from system.tables.engine_full. ok is false when the table has
// no TTL or one Bifrost did not write.
func chTTLDaysFromEngineFull(engineFull string) (days int, ok bool) {
return chTTLDaysFromMatch(chTTLEngineFullDaysRe, engineFull)
}

// clickhouseReconcileTTL brings an existing table's TTL in line with wantTTL.
// CREATE TABLE IF NOT EXISTS never updates the TTL of a table that already
// exists, so before this a changed logs_store.retention_days silently did
// nothing (#7098). An empty wantTTL means "unmanaged" and leaves the table
// alone, so a TTL an operator applied by hand survives restarts.
//
// materialize_ttl_after_modify = 0 keeps the ALTER metadata-only. The default
// would submit a MATERIALIZE TTL mutation that rewrites every part holding
// expired rows, on every pod's first boot, which is the exact heavyweight
// rewrite this code path exists to avoid. Expired rows are dropped by the
// regular TTL merges instead.
func clickhouseReconcileTTL(ctx context.Context, db *gorm.DB, table, cluster, wantTTL string, logger schemas.Logger) error {
wantDays, managed := chTTLDaysFromClause(wantTTL)
if !managed {
return nil
}
var engineFull string
if err := db.WithContext(ctx).
Raw("SELECT engine_full FROM system.tables WHERE database = currentDatabase() AND name = ?", table).
Scan(&engineFull).Error; err != nil {
return fmt.Errorf("clickhouse: read %s engine definition: %w", table, err)
}
if haveDays, ok := chTTLDaysFromEngineFull(engineFull); ok && haveDays == wantDays {
return nil
}
onCluster := ""
if cluster != "" {
onCluster = fmt.Sprintf(" ON CLUSTER `%s`", chEscapeIdentifier(cluster))
}
stmt := fmt.Sprintf("ALTER TABLE `%s`%s MODIFY TTL %s SETTINGS materialize_ttl_after_modify = 0", table, onCluster, wantTTL)
logger.Info("[logstore] clickhouse: setting %s TTL to %d days", table, wantDays)
if err := db.WithContext(ctx).Exec(stmt).Error; err != nil {
return fmt.Errorf("clickhouse: modify %s TTL: %w", table, err)
}
return nil
}

// clickhouseMigrationStep is one per-table migration: create the table if
// missing, then reconcile any columns added to its model since.
type clickhouseMigrationStep func(ctx context.Context, db *gorm.DB, cluster string, retentionDays int, logger schemas.Logger) error
Expand All @@ -257,7 +331,10 @@ func migrationClickHouseLogsTable(ctx context.Context, db *gorm.DB, cluster stri
}, cluster); err != nil {
return fmt.Errorf("clickhouse: create logs table: %w", err)
}
return clickhouseReconcileColumns(ctx, db, &Log{}, "logs", cluster, logger)
if err := clickhouseReconcileColumns(ctx, db, &Log{}, "logs", cluster, logger); err != nil {
return err
}
return clickhouseReconcileTTL(ctx, db, "logs", cluster, chLogsTTL(retentionDays), logger)
}

// migrationClickHouseMCPToolLogsTable creates the mcp_tool_logs table and
Expand All @@ -277,7 +354,10 @@ func migrationClickHouseMCPToolLogsTable(ctx context.Context, db *gorm.DB, clust
}, cluster); err != nil {
return fmt.Errorf("clickhouse: create mcp_tool_logs table: %w", err)
}
return clickhouseReconcileColumns(ctx, db, &MCPToolLog{}, "mcp_tool_logs", cluster, logger)
if err := clickhouseReconcileColumns(ctx, db, &MCPToolLog{}, "mcp_tool_logs", cluster, logger); err != nil {
return err
}
return clickhouseReconcileTTL(ctx, db, "mcp_tool_logs", cluster, chLogsTTL(retentionDays), logger)
}

// migrationClickHouseAsyncJobsTable creates the async_jobs table and reconciles
Expand All @@ -294,7 +374,10 @@ func migrationClickHouseAsyncJobsTable(ctx context.Context, db *gorm.DB, cluster
}, cluster); err != nil {
return fmt.Errorf("clickhouse: create async_jobs table: %w", err)
}
return clickhouseReconcileColumns(ctx, db, &AsyncJob{}, "async_jobs", cluster, logger)
if err := clickhouseReconcileColumns(ctx, db, &AsyncJob{}, "async_jobs", cluster, logger); err != nil {
return err
}
return clickhouseReconcileTTL(ctx, db, "async_jobs", cluster, "toDateTime(created_at) + INTERVAL 7 DAY", logger)
}

// migrationClickHouseWebhookDeliveriesTable creates the webhook_deliveries
Expand All @@ -311,7 +394,10 @@ func migrationClickHouseWebhookDeliveriesTable(ctx context.Context, db *gorm.DB,
}, cluster); err != nil {
return fmt.Errorf("clickhouse: create webhook_deliveries table: %w", err)
}
return clickhouseReconcileColumns(ctx, db, &WebhookDelivery{}, "webhook_deliveries", cluster, logger)
if err := clickhouseReconcileColumns(ctx, db, &WebhookDelivery{}, "webhook_deliveries", cluster, logger); err != nil {
return err
}
return clickhouseReconcileTTL(ctx, db, "webhook_deliveries", cluster, chLogsTTL(retentionDays), logger)
}

// clickhouseMigrationSteps lists the per-table migrations in execution order,
Expand Down
Loading
Loading