From f0809b37cf2e8bd4164c71e0752570d5e8bdea08 Mon Sep 17 00:00:00 2001 From: akshaydeo Date: Sat, 12 Sep 2026 00:10:13 +0530 Subject: [PATCH] adds clickhouse delte rows improvements --- .github/workflows/scripts/test-framework.sh | 21 +- .../deployment-guides/config-json/storage.mdx | 2 +- framework/changelog.md | 1 + framework/logstore/cleaner.go | 8 +- framework/logstore/clickhouse.go | 41 +- framework/logstore/clickhousemigrate.go | 98 +++- framework/logstore/clickhousestore.go | 225 +++++---- framework/logstore/clickhousestore_test.go | 464 +++++++++++++++++- tests/docker-compose.yml | 30 +- 9 files changed, 775 insertions(+), 115 deletions(-) diff --git a/.github/workflows/scripts/test-framework.sh b/.github/workflows/scripts/test-framework.sh index f39002ddd79..89575b81bdd 100755 --- a/.github/workflows/scripts/test-framework.sh +++ b/.github/workflows/scripts/test-framework.sh @@ -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 diff --git a/docs/deployment-guides/config-json/storage.mdx b/docs/deployment-guides/config-json/storage.mdx index 1d65906b71a..4f99306ce8b 100644 --- a/docs/deployment-guides/config-json/storage.mdx +++ b/docs/deployment-guides/config-json/storage.mdx @@ -279,7 +279,7 @@ ClickHouse is a **`logs_store`-only** backend. The `config_store` supports only ``` -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. ### Disabled diff --git a/framework/changelog.md b/framework/changelog.md index e69de29bb2d..ef31cc105a7 100644 --- a/framework/changelog.md +++ b/framework/changelog.md @@ -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) diff --git a/framework/logstore/cleaner.go b/framework/logstore/cleaner.go index e86a57e9e21..c1787820014 100644 --- a/framework/logstore/cleaner.go +++ b/framework/logstore/cleaner.go @@ -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 } } diff --git a/framework/logstore/clickhouse.go b/framework/logstore/clickhouse.go index 5080c7a2f8b..de44ec019b7 100644 --- a/framework/logstore/clickhouse.go +++ b/framework/logstore/clickhouse.go @@ -5,6 +5,7 @@ import ( "fmt" "net" "net/url" + "strconv" "strings" "time" @@ -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 { @@ -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) diff --git a/framework/logstore/clickhousemigrate.go b/framework/logstore/clickhousemigrate.go index e2cb7840a40..e8c99b7b10f 100644 --- a/framework/logstore/clickhousemigrate.go +++ b/framework/logstore/clickhousemigrate.go @@ -4,6 +4,8 @@ import ( "context" "fmt" "reflect" + "regexp" + "strconv" "strings" "time" @@ -223,8 +225,9 @@ 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 "" @@ -232,6 +235,77 @@ func chLogsTTL(retentionDays int) string { 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 @@ -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 @@ -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 @@ -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 @@ -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, diff --git a/framework/logstore/clickhousestore.go b/framework/logstore/clickhousestore.go index fe1108c01c3..094f010ebc9 100644 --- a/framework/logstore/clickhousestore.go +++ b/framework/logstore/clickhousestore.go @@ -23,10 +23,17 @@ import ( // `final = 1` setting, so a plain INSERT is correct), // - row updates (ClickHouse has no cheap UPDATE; we read-modify-write and // re-insert, letting the `ver` DEFAULT now64() column make the newest -// insert win on merge - see clickhousemigrate.go). +// insert win on merge - see clickhousemigrate.go), +// - every delete. The GORM ClickHouse driver rewrites Delete() into +// `ALTER TABLE ... DELETE`, a heavyweight mutation that rewrites every +// column of every part holding a matching row and leaves the old copy on +// disk for old_parts_lifetime (#7098). All deletes here go through +// chLightweightDelete, a raw `DELETE FROM ... WHERE` that only writes the +// _row_exists mask. Never call GORM Delete() on this store. // -// Deletes are left to the embedded methods: the GORM ClickHouse driver emits -// lightweight `DELETE ... WHERE`, and TTL is the primary retention mechanism. +// Table TTL (logs_store.retention_days) is the primary retention mechanism and +// is reconciled on every startup (clickhouseReconcileTTL); the LogsCleaner +// sweep is a single lightweight delete per run. type ClickHouseLogStore struct { *RDBLogStore // cluster is the optional ON CLUSTER name (empty = single-node). Retained @@ -483,110 +490,146 @@ func (s *ClickHouseLogStore) UpdateMCPToolLog(ctx context.Context, id string, en return s.chReinsert(ctx, &existing) } -// DeleteLogsBatch deletes logs older than cutoff in batches. Overridden -// because the GORM ClickHouse driver rewrites DELETE into an ALTER TABLE -// mutation whose driver result reports 0 rows affected - the inherited -// implementation would always return 0 and the LogsCleaner would treat every -// batch as empty and stop early. The ids are selected first, so their count -// is the deleted count once the (mutations_sync=1) delete returns. -func (s *ClickHouseLogStore) DeleteLogsBatch(ctx context.Context, cutoff time.Time, batchSize int) (int64, error) { - var ids []string - if err := s.db.WithContext(ctx). - Model(&Log{}). - Select("id"). - Where("created_at < ?", cutoff). - Order("created_at ASC"). - Limit(batchSize). - Pluck("id", &ids).Error; err != nil { +// chLightweightDelete issues a ClickHouse lightweight DELETE through raw Exec, +// bypassing the GORM driver's rewrite of Delete() into a heavyweight +// `ALTER TABLE ... DELETE`. ClickHouse records it as +// `UPDATE _row_exists = 0 WHERE ...`: on wide parts only the _row_exists mask +// is written and every other column file is hardlinked, so one call costs a +// mask per affected part instead of a full part rewrite per call (#7098). +// lightweight_deletes_sync = 1 waits for the current replica only, matching +// the mutations_sync=1 the DSN sets for the remaining heavyweight mutations +// (the default 2 would block on every replica of a cluster). Requires +// ClickHouse 24.4+, where the setting was introduced. +// +// The driver never reports rows affected for mutations, so callers that need +// a count select it first (chDeleteWhere). +func (s *ClickHouseLogStore) chLightweightDelete(ctx context.Context, table, where string, args ...any) error { + stmt := fmt.Sprintf("DELETE FROM `%s` WHERE %s SETTINGS lightweight_deletes_sync = 1", chEscapeIdentifier(table), where) + return s.db.WithContext(ctx).Exec(stmt, args...).Error +} + +// chCountWhere counts the logical rows matching where. The DSN-level final=1 +// collapses ReplacingMergeTree versions, so the count matches what the SQL +// stores report for the same predicate (see the delete_logs_batch parity test). +func (s *ClickHouseLogStore) chCountWhere(ctx context.Context, table, where string, args ...any) (int64, error) { + var count int64 + err := s.db.WithContext(ctx). + Raw(fmt.Sprintf("SELECT count() FROM `%s` WHERE %s", chEscapeIdentifier(table), where), args...). + Scan(&count).Error + return count, err +} + +// chExistsWhere reports whether any current row matches where. It runs under +// the connection-level final=1 so a superseded ReplacingMergeTree version (a +// log created as processing and later re-inserted as success) does not match; +// otherwise the minute sweep would issue a mutation on every run until the old +// version merged away, which on a large part can take hours. LIMIT 1 stops at +// the first hit. +func (s *ClickHouseLogStore) chExistsWhere(ctx context.Context, table, where string, args ...any) (bool, error) { + var hits []uint8 + err := s.db.WithContext(ctx). + Raw(fmt.Sprintf("SELECT 1 FROM `%s` WHERE %s LIMIT 1", chEscapeIdentifier(table), where), args...). + Scan(&hits).Error + return len(hits) > 0, err +} + +// chDeleteWhere counts the rows matching where and, when there are any, +// removes them with a single lightweight delete. It returns the count so the +// cleaners can log and pace on an accurate number; when nothing matches no +// mutation is issued at all. +func (s *ClickHouseLogStore) chDeleteWhere(ctx context.Context, table, where string, args ...any) (int64, error) { + count, err := s.chCountWhere(ctx, table, where, args...) + if err != nil || count == 0 { + return 0, err + } + if err := s.chLightweightDelete(ctx, table, where, args...); err != nil { return 0, err } + return count, nil +} + +// chFlushProcessing removes rows still marked processing that were created +// before since. It probes first so the once-a-minute sweep from the logging +// plugin issues no mutation on an idle table (the issue counted ~1,440 +// heavyweight mutations per table per day from this path alone). +func (s *ClickHouseLogStore) chFlushProcessing(ctx context.Context, table string, since time.Time) error { + const where = "status = 'processing' AND created_at < ?" + exists, err := s.chExistsWhere(ctx, table, where, since) + if err != nil || !exists { + return err + } + return s.chLightweightDelete(ctx, table, where, since) +} + +// DeleteLogsBatch deletes every log older than cutoff with one lightweight +// delete per call. batchSize is ignored: a ClickHouse mutation costs the same +// per affected part whether it matches 100 rows or all of them, so batching +// by id would rewrite the same part once per batch (#7098). The returned +// count is selected before the delete because mutations never report rows +// affected. The LogsCleaner loop stops after this call because the count +// differs from batchSize; when it happens to equal batchSize the next call +// finds nothing and returns 0 without issuing a mutation. +func (s *ClickHouseLogStore) DeleteLogsBatch(ctx context.Context, cutoff time.Time, _ int) (int64, error) { + return s.chDeleteWhere(ctx, "logs", "created_at < ?", cutoff) +} + +// DeleteLog deletes a log entry by id with a lightweight delete. +func (s *ClickHouseLogStore) DeleteLog(ctx context.Context, id string) error { + return s.chLightweightDelete(ctx, "logs", "id = ?", id) +} + +// DeleteLogs deletes multiple log entries by id with one lightweight delete. +func (s *ClickHouseLogStore) DeleteLogs(ctx context.Context, ids []string) error { if len(ids) == 0 { - return 0, nil + return nil } - if err := s.db.WithContext(ctx).Where("id IN ?", ids).Delete(&Log{}).Error; err != nil { - return 0, err + return s.chLightweightDelete(ctx, "logs", "id IN ?", ids) +} + +// DeleteMCPToolLogs deletes multiple MCP tool log entries by id with one +// lightweight delete. +func (s *ClickHouseLogStore) DeleteMCPToolLogs(ctx context.Context, ids []string) error { + if len(ids) == 0 { + return nil } - return int64(len(ids)), nil + return s.chLightweightDelete(ctx, "mcp_tool_logs", "id IN ?", ids) } -// DeleteExpiredAsyncJobs deletes async jobs whose expiry has passed. -// Overridden for the same reason as DeleteLogsBatch: mutation deletes report -// 0 rows affected, so ids are selected first and their count returned. -func (s *ClickHouseLogStore) DeleteExpiredAsyncJobs(ctx context.Context) (int64, error) { - now := time.Now().UTC() - const batchLimit = 100 - var total int64 - for { - var ids []string - if err := s.db.WithContext(ctx).Model(&AsyncJob{}).Select("id"). - Where("expires_at IS NOT NULL AND expires_at < ?", now). - Limit(batchLimit).Pluck("id", &ids).Error; err != nil { - return total, err - } - if len(ids) == 0 { - return total, nil - } - if err := s.db.WithContext(ctx).Where("id IN ?", ids).Delete(&AsyncJob{}).Error; err != nil { - return total, err - } - total += int64(len(ids)) - if len(ids) < batchLimit { - return total, nil - } +// Flush removes stale processing log rows. Overridden so the minute sweep is +// a probe plus at most one lightweight delete. The error text matches the +// SQL stores so the logging plugin's warnings are unchanged. +func (s *ClickHouseLogStore) Flush(ctx context.Context, since time.Time) error { + if err := s.chFlushProcessing(ctx, "logs", since); err != nil { + return fmt.Errorf("failed to cleanup old processing logs: %w", err) } + return nil +} + +// FlushMCPToolLogs removes stale processing MCP tool log rows. See Flush. +func (s *ClickHouseLogStore) FlushMCPToolLogs(ctx context.Context, since time.Time) error { + if err := s.chFlushProcessing(ctx, "mcp_tool_logs", since); err != nil { + return fmt.Errorf("failed to cleanup old processing MCP tool logs: %w", err) + } + return nil +} + +// DeleteExpiredAsyncJobs deletes async jobs whose expiry has passed with one +// lightweight delete; the count is selected first because mutations never +// report rows affected. +func (s *ClickHouseLogStore) DeleteExpiredAsyncJobs(ctx context.Context) (int64, error) { + return s.chDeleteWhere(ctx, "async_jobs", "expires_at IS NOT NULL AND expires_at < ?", time.Now().UTC()) } // DeleteStaleAsyncJobs deletes processing jobs created before staleSince. -// See DeleteExpiredAsyncJobs for why the count is derived from a prior select. +// See DeleteExpiredAsyncJobs. func (s *ClickHouseLogStore) DeleteStaleAsyncJobs(ctx context.Context, staleSince time.Time) (int64, error) { - const batchLimit = 100 - var total int64 - for { - var ids []string - if err := s.db.WithContext(ctx).Model(&AsyncJob{}).Select("id"). - Where("status = ? AND created_at < ?", "processing", staleSince). - Limit(batchLimit).Pluck("id", &ids).Error; err != nil { - return total, err - } - if len(ids) == 0 { - return total, nil - } - if err := s.db.WithContext(ctx).Where("id IN ?", ids).Delete(&AsyncJob{}).Error; err != nil { - return total, err - } - total += int64(len(ids)) - if len(ids) < batchLimit { - return total, nil - } - } + return s.chDeleteWhere(ctx, "async_jobs", "status = 'processing' AND created_at < ?", staleSince) } // DeleteExpiredWebhookDeliveries deletes webhook delivery history whose -// expiry has passed. Overridden for the same reason as DeleteLogsBatch: -// mutation deletes report 0 rows affected, so ids are selected first and -// their count returned. +// expiry has passed. See DeleteExpiredAsyncJobs. func (s *ClickHouseLogStore) DeleteExpiredWebhookDeliveries(ctx context.Context) (int64, error) { - now := time.Now().UTC() - const batchLimit = 100 - var total int64 - for { - var ids []string - if err := s.db.WithContext(ctx).Model(&WebhookDelivery{}).Select("id"). - Where("expires_at IS NOT NULL AND expires_at < ?", now). - Limit(batchLimit).Pluck("id", &ids).Error; err != nil { - return total, err - } - if len(ids) == 0 { - return total, nil - } - if err := s.db.WithContext(ctx).Where("id IN ?", ids).Delete(&WebhookDelivery{}).Error; err != nil { - return total, err - } - total += int64(len(ids)) - if len(ids) < batchLimit { - return total, nil - } - } + return s.chDeleteWhere(ctx, "webhook_deliveries", "expires_at IS NOT NULL AND expires_at < ?", time.Now().UTC()) } // UpdateAsyncJob applies a column->value map to an async job row via diff --git a/framework/logstore/clickhousestore_test.go b/framework/logstore/clickhousestore_test.go index da7b8c4e61a..af3827de211 100644 --- a/framework/logstore/clickhousestore_test.go +++ b/framework/logstore/clickhousestore_test.go @@ -3,6 +3,7 @@ package logstore import ( "context" "fmt" + "os" "reflect" "strings" "sync" @@ -17,9 +18,18 @@ import ( gormschema "gorm.io/gorm/schema" ) -// ClickHouse test connection matches the clickhouse service in +// ClickHouse test connection defaults match the clickhouse service in // framework/docker-compose.yml (native protocol on host port 9001; host 9000 -// is taken by Weaviate). +// is taken by Weaviate). Each value can be overridden through the +// BIFROST_TEST_CLICKHOUSE_* environment variables so the suite can target +// another local instance (for example one whose ports collide with 9001). +// The suite TRUNCATEs the log tables and rewrites their TTL, so the default +// "bifrost" database is accepted only for the stock docker-compose target: +// setting any override, even just the host or port, requires +// BIFROST_TEST_CLICKHOUSE_DB to name a database containing "test" +// (requireDedicatedClickHouseTestDB fails the test otherwise). For example: +// +// BIFROST_TEST_CLICKHOUSE_PORT=9011 BIFROST_TEST_CLICKHOUSE_DB=bifrost_test const ( clickhouseTestHost = "localhost" clickhouseTestPort = "9001" @@ -28,24 +38,70 @@ const ( clickhouseTestPassword = "bifrost_password" ) +// chTestEnv returns the environment override for key, or def when unset. +func chTestEnv(key, def string) string { + if v := os.Getenv(key); v != "" { + return v + } + return def +} + +// chTestOverridden reports whether any BIFROST_TEST_CLICKHOUSE_* override is +// set, meaning the suite is pointed away from the stock docker-compose target. +func chTestOverridden() bool { + for _, k := range []string{"HOST", "PORT", "DB", "USER", "PASSWORD"} { + if os.Getenv("BIFROST_TEST_CLICKHOUSE_"+k) != "" { + return true + } + } + return false +} + +// chTestTargetIsDedicated reports whether the suite may run its destructive +// setup (TRUNCATE of every log table, TTL rewrites) against cfg. The stock +// docker-compose target is dedicated by definition; an overridden target is +// accepted only when its database name marks it as a test database. +func chTestTargetIsDedicated(cfg *ClickHouseConfig, overridden bool) bool { + if !overridden { + return true + } + return strings.Contains(strings.ToLower(cfg.Database.GetValue()), "test") +} + +// requireDedicatedClickHouseTestDB fails the test before any connection is +// opened when the configured target is not safe to truncate. +func requireDedicatedClickHouseTestDB(t *testing.T, cfg *ClickHouseConfig) { + t.Helper() + if !chTestTargetIsDedicated(cfg, chTestOverridden()) { + t.Fatalf("refusing to run destructive ClickHouse tests against database %q: BIFROST_TEST_CLICKHOUSE_* overrides must point at a database whose name contains \"test\" (the suite truncates logs, mcp_tool_logs, async_jobs and webhook_deliveries and rewrites their TTL)", cfg.Database.GetValue()) + } +} + func clickhouseTestConfig() *ClickHouseConfig { return &ClickHouseConfig{ - Host: schemas.NewSecretVar(clickhouseTestHost), - Port: schemas.NewSecretVar(clickhouseTestPort), - Database: schemas.NewSecretVar(clickhouseTestDatabase), - Username: schemas.NewSecretVar(clickhouseTestUser), - Password: schemas.NewSecretVar(clickhouseTestPassword), + Host: schemas.NewSecretVar(chTestEnv("BIFROST_TEST_CLICKHOUSE_HOST", clickhouseTestHost)), + Port: schemas.NewSecretVar(chTestEnv("BIFROST_TEST_CLICKHOUSE_PORT", clickhouseTestPort)), + Database: schemas.NewSecretVar(chTestEnv("BIFROST_TEST_CLICKHOUSE_DB", clickhouseTestDatabase)), + Username: schemas.NewSecretVar(chTestEnv("BIFROST_TEST_CLICKHOUSE_USER", clickhouseTestUser)), + Password: schemas.NewSecretVar(chTestEnv("BIFROST_TEST_CLICKHOUSE_PASSWORD", clickhouseTestPassword)), } } // trySetupClickHouseStore connects to the docker-compose ClickHouse, runs // migrations, and truncates the log tables for a clean slate. Skips the test -// when ClickHouse is unavailable. +// when ClickHouse is unavailable locally; in CI (the CI env var is set by +// GitHub Actions and tests/docker-compose.yml provides the service) an +// unreachable ClickHouse is a failure, so the suite can never silently skip. func trySetupClickHouseStore(t *testing.T) *ClickHouseLogStore { t.Helper() ctx := context.Background() - store, err := newClickHouseLogStore(ctx, clickhouseTestConfig(), 0, testLogger{}) + cfg := clickhouseTestConfig() + requireDedicatedClickHouseTestDB(t, cfg) + store, err := newClickHouseLogStore(ctx, cfg, 0, testLogger{}) if err != nil { + if os.Getenv("CI") != "" { + t.Fatalf("ClickHouse not available in CI (is the clickhouse service in tests/docker-compose.yml up?): %v", err) + } t.Skipf("ClickHouse not available, skipping test: %v", err) } ch := store.(*ClickHouseLogStore) @@ -152,7 +208,15 @@ func TestBuildClickHouseDSN(t *testing.T) { }) t.Run("CredentialsPortAndDatabase", func(t *testing.T) { - dsn, err := buildClickHouseDSN(clickhouseTestConfig()) + // A literal config (not clickhouseTestConfig) so the assertion stays + // deterministic when BIFROST_TEST_CLICKHOUSE_* overrides are set. + dsn, err := buildClickHouseDSN(&ClickHouseConfig{ + Host: schemas.NewSecretVar("localhost"), + Port: schemas.NewSecretVar("9001"), + Database: schemas.NewSecretVar("bifrost"), + Username: schemas.NewSecretVar("bifrost"), + Password: schemas.NewSecretVar("bifrost_password"), + }) require.NoError(t, err) assert.Contains(t, dsn, "bifrost:bifrost_password@localhost:9001/bifrost") }) @@ -186,6 +250,104 @@ func chUnitSchemaDB() *gorm.DB { return &gorm.DB{Config: &gorm.Config{NamingStrategy: gormschema.NamingStrategy{}}} } +func TestChTestTargetIsDedicated(t *testing.T) { + cfg := func(db string) *ClickHouseConfig { return &ClickHouseConfig{Database: schemas.NewSecretVar(db)} } + assert.True(t, chTestTargetIsDedicated(cfg("bifrost"), false), "stock compose target is always allowed") + assert.True(t, chTestTargetIsDedicated(cfg("bifrost_test"), true)) + assert.True(t, chTestTargetIsDedicated(cfg("TestLogs"), true), "case-insensitive") + assert.False(t, chTestTargetIsDedicated(cfg("bifrost"), true), "an override onto a non-test database is refused") + assert.False(t, chTestTargetIsDedicated(cfg(""), true)) + + t.Setenv("BIFROST_TEST_CLICKHOUSE_PORT", "9011") + assert.True(t, chTestOverridden()) +} + +func TestChServerVersionSupported(t *testing.T) { + for _, tc := range []struct { + version string + ok bool + }{ + {"26.6.1.1193", true}, + {"24.8.14.39", true}, + {"24.4.1.1", true}, + {"24.3.9.5", false}, + {"23.8.2.7", false}, + {" 25.1 ", true}, + } { + ok, err := chServerVersionSupported(tc.version) + require.NoError(t, err, tc.version) + assert.Equal(t, tc.ok, ok, tc.version) + } + for _, bad := range []string{"", "26", "x.y.z", "24.four"} { + _, err := chServerVersionSupported(bad) + assert.Error(t, err, bad) + } +} + +func TestChTTLDays(t *testing.T) { + t.Run("EngineFull", func(t *testing.T) { + // Fixtures copied verbatim from system.tables.engine_full on ClickHouse 26.6. + days, ok := chTTLDaysFromEngineFull("ReplacingMergeTree(ver) ORDER BY id TTL toDateTime(created_at) + toIntervalDay(7) SETTINGS index_granularity = 8192") + assert.True(t, ok) + assert.Equal(t, 7, days) + + days, ok = chTTLDaysFromEngineFull("ReplacingMergeTree(ver) PARTITION BY toYYYYMM(timestamp) ORDER BY (timestamp, id) SETTINGS index_granularity = 8192") + assert.False(t, ok, "no TTL") + assert.Equal(t, 0, days) + + _, ok = chTTLDaysFromEngineFull("ReplacingMergeTree(ver) ORDER BY id TTL timestamp + toIntervalDay(3) SETTINGS index_granularity = 8192") + assert.False(t, ok, "a TTL Bifrost did not write is not managed") + }) + t.Run("Clause", func(t *testing.T) { + days, ok := chTTLDaysFromClause(chLogsTTL(3)) + assert.True(t, ok) + assert.Equal(t, 3, days) + + _, ok = chTTLDaysFromClause(chLogsTTL(0)) + assert.False(t, ok, "retention < 1 is unmanaged") + + _, ok = chTTLDaysFromClause("toDateTime(created_at) + INTERVAL 3 DAY DELETE WHERE status = 'x'") + assert.False(t, ok, "only the exact clause shape is managed") + }) +} + +// countingRetentionManager is a LogRetentionManager stub that returns a +// scripted count per call and records how often it was asked. +type countingRetentionManager struct { + counts []int64 + calls int +} + +func (m *countingRetentionManager) DeleteLogsBatch(_ context.Context, _ time.Time, _ int) (int64, error) { + m.calls++ + if m.calls > len(m.counts) { + return 0, nil + } + return m.counts[m.calls-1], nil +} + +// TestLogsCleanerStopsAfterOversizedBatch pins the loop contract ClickHouse +// relies on: a store that deletes the whole expired range in one statement +// returns a count above batchSize, and the cleaner must stop there instead of +// issuing the delete again. +func TestLogsCleanerStopsAfterOversizedBatch(t *testing.T) { + t.Run("OversizedCountEndsTheLoop", func(t *testing.T) { + m := &countingRetentionManager{counts: []int64{250}} + NewLogsCleaner(m, CleanerConfig{RetentionDays: 3}, testLogger{}).cleanupOldLogs(context.Background()) + assert.Equal(t, 1, m.calls, "a count above batchSize means the store already deleted everything") + }) + t.Run("FullBatchesKeepGoing", func(t *testing.T) { + m := &countingRetentionManager{counts: []int64{100, 100, 40}} + NewLogsCleaner(m, CleanerConfig{RetentionDays: 3}, testLogger{}).cleanupOldLogs(context.Background()) + assert.Equal(t, 3, m.calls, "SQL stores return exactly batchSize while rows remain") + }) + t.Run("ExactBatchThenEmpty", func(t *testing.T) { + m := &countingRetentionManager{counts: []int64{100, 0}} + NewLogsCleaner(m, CleanerConfig{RetentionDays: 3}, testLogger{}).cleanupOldLogs(context.Background()) + assert.Equal(t, 2, m.calls, "a full batch is followed by one more probe that finds nothing") + }) +} + func TestChApplyUpdateMapSkipsDedupKeys(t *testing.T) { ctx := context.Background() st, err := chParseSchema(chUnitSchemaDB(), &Log{}) @@ -541,6 +703,288 @@ func TestClickHouseDeleteLogsBatch(t *testing.T) { assert.NoError(t, err) } +// chMutationIDs snapshots the mutation ids ClickHouse has recorded for table. +// system.mutations keeps finished entries (and trims them in the background), +// so tests diff a before/after snapshot instead of asserting absolute counts. +func chMutationIDs(t *testing.T, db *gorm.DB, table string) map[string]struct{} { + t.Helper() + var ids []string + require.NoError(t, db.Raw("SELECT mutation_id FROM system.mutations WHERE database = currentDatabase() AND table = ?", table).Scan(&ids).Error) + set := make(map[string]struct{}, len(ids)) + for _, id := range ids { + set[id] = struct{}{} + } + return set +} + +// chNewMutationCommands returns the command text of every mutation recorded +// for table since the before snapshot, in creation order. +func chNewMutationCommands(t *testing.T, db *gorm.DB, table string, before map[string]struct{}) []string { + t.Helper() + type row struct { + MutationID string + Command string + } + var rows []row + require.NoError(t, db.Raw("SELECT mutation_id, command FROM system.mutations WHERE database = currentDatabase() AND table = ? ORDER BY create_time, mutation_id", table).Scan(&rows).Error) + var cmds []string + for _, r := range rows { + if _, seen := before[r.MutationID]; seen { + continue + } + cmds = append(cmds, r.Command) + } + return cmds +} + +// chLightweightDeletePrefix is how ClickHouse records a lightweight DELETE in +// system.mutations (the command column wraps each command in parentheses). +// A heavyweight ALTER TABLE ... DELETE is recorded as "(DELETE WHERE ...)" +// and rewrites every column of every affected part. +const chLightweightDeletePrefix = "UPDATE _row_exists = 0" + +func assertLightweightMutations(t *testing.T, table string, cmds []string) { + t.Helper() + for _, cmd := range cmds { + assert.True(t, strings.HasPrefix(strings.TrimLeft(cmd, "("), chLightweightDeletePrefix), "%s: expected a lightweight delete mutation, got %q", table, cmd) + assert.NotContains(t, cmd, "DELETE WHERE", "%s: heavyweight ALTER TABLE ... DELETE rewrites whole parts (#7098)", table) + } +} + +// TestClickHouseDeleteLogsBatchIsSingleLightweightMutation covers #7098: one +// retention sweep must cost one lightweight mutation, not one heavyweight +// part rewrite per 100 rows. It drives the same loop LogsCleaner runs. +func TestClickHouseDeleteLogsBatchIsSingleLightweightMutation(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + + old := time.Now().UTC().Add(-48 * time.Hour).Truncate(time.Millisecond) + entries := make([]*Log, 0, 250) + for i := range 250 { + entries = append(entries, chTestLog(fmt.Sprintf("ch-sweep-%03d", i), old.Add(time.Duration(i)*time.Millisecond))) + } + require.NoError(t, store.BatchCreateIfNotExists(ctx, entries)) + require.NoError(t, store.CreateIfNotExists(ctx, chTestLog("ch-sweep-fresh", time.Now().UTC().Truncate(time.Millisecond)))) + + before := chMutationIDs(t, store.db, "logs") + cutoff := time.Now().UTC().Add(-24 * time.Hour) + var total int64 + for { + deleted, err := store.DeleteLogsBatch(ctx, cutoff, batchSize) + require.NoError(t, err) + total += deleted + if deleted != int64(batchSize) { + break + } + } + assert.Equal(t, int64(250), total, "cleaner logging relies on an accurate deleted count") + + _, err := store.FindByID(ctx, "ch-sweep-fresh") + assert.NoError(t, err, "rows newer than the cutoff must survive") + _, err = store.FindByID(ctx, "ch-sweep-000") + assert.ErrorIs(t, err, ErrNotFound) + + cmds := chNewMutationCommands(t, store.db, "logs", before) + require.Len(t, cmds, 1, "one sweep must issue exactly one mutation, got %v", cmds) + assertLightweightMutations(t, "logs", cmds) +} + +// TestClickHouseFlushIsLightweightAndSkipsWhenEmpty covers the minute sweep +// from plugins/logging: Flush and FlushMCPToolLogs must use lightweight +// deletes and must not issue any mutation when nothing is left to flush. +func TestClickHouseFlushIsLightweightAndSkipsWhenEmpty(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + + old := time.Now().UTC().Add(-40 * time.Minute).Truncate(time.Millisecond) + stuck := chTestLog("ch-flush-stuck", old) + done := chTestLog("ch-flush-done", old) + done.Status = "success" + require.NoError(t, store.CreateIfNotExists(ctx, stuck)) + require.NoError(t, store.CreateIfNotExists(ctx, done)) + require.NoError(t, store.BatchCreateMCPToolLogsIfNotExists(ctx, []*MCPToolLog{chTestMCPToolLog("ch-flush-mcp", old)})) + + logsBefore := chMutationIDs(t, store.db, "logs") + mcpBefore := chMutationIDs(t, store.db, "mcp_tool_logs") + since := time.Now().UTC().Add(-30 * time.Minute) + require.NoError(t, store.Flush(ctx, since)) + require.NoError(t, store.FlushMCPToolLogs(ctx, since)) + + _, err := store.FindByID(ctx, "ch-flush-stuck") + assert.ErrorIs(t, err, ErrNotFound, "stale processing row must be flushed") + _, err = store.FindByID(ctx, "ch-flush-done") + assert.NoError(t, err, "completed rows must survive the flush") + _, err = store.FindMCPToolLog(ctx, "ch-flush-mcp") + assert.ErrorIs(t, err, ErrNotFound, "stale processing MCP row must be flushed") + + logsCmds := chNewMutationCommands(t, store.db, "logs", logsBefore) + require.Len(t, logsCmds, 1, "logs: one flush must issue exactly one mutation, got %v", logsCmds) + assertLightweightMutations(t, "logs", logsCmds) + mcpCmds := chNewMutationCommands(t, store.db, "mcp_tool_logs", mcpBefore) + require.Len(t, mcpCmds, 1, "mcp_tool_logs: one flush must issue exactly one mutation, got %v", mcpCmds) + assertLightweightMutations(t, "mcp_tool_logs", mcpCmds) + + // Nothing left to flush: the once-a-minute sweep must not touch the + // tables at all (the issue counted ~1,440 mutations per table per day). + logsBefore = chMutationIDs(t, store.db, "logs") + mcpBefore = chMutationIDs(t, store.db, "mcp_tool_logs") + require.NoError(t, store.Flush(ctx, since)) + require.NoError(t, store.FlushMCPToolLogs(ctx, since)) + assert.Empty(t, chNewMutationCommands(t, store.db, "logs", logsBefore), "empty flush must not issue a mutation") + assert.Empty(t, chNewMutationCommands(t, store.db, "mcp_tool_logs", mcpBefore), "empty flush must not issue a mutation") +} + +// TestClickHouseFlushSkipsSupersededProcessingRow: a log created as processing +// and then updated to success leaves the old version physically in place until +// merge. The minute sweep must not treat that superseded version as stale, or +// it would issue a mutation every run until the part merged. +func TestClickHouseFlushSkipsSupersededProcessingRow(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + old := time.Now().UTC().Add(-40 * time.Minute).Truncate(time.Millisecond) + require.NoError(t, store.CreateIfNotExists(ctx, chTestLog("ch-flush-superseded", old))) + require.NoError(t, store.Update(ctx, "ch-flush-superseded", map[string]interface{}{"status": "success"})) + require.NoError(t, store.BatchCreateMCPToolLogsIfNotExists(ctx, []*MCPToolLog{chTestMCPToolLog("ch-flush-superseded-mcp", old)})) + require.NoError(t, store.UpdateMCPToolLog(ctx, "ch-flush-superseded-mcp", map[string]interface{}{"status": "success"})) + + logsBefore := chMutationIDs(t, store.db, "logs") + mcpBefore := chMutationIDs(t, store.db, "mcp_tool_logs") + since := time.Now().UTC().Add(-30 * time.Minute) + require.NoError(t, store.Flush(ctx, since)) + require.NoError(t, store.FlushMCPToolLogs(ctx, since)) + assert.Empty(t, chNewMutationCommands(t, store.db, "logs", logsBefore), "superseded processing version must not trigger a flush mutation") + assert.Empty(t, chNewMutationCommands(t, store.db, "mcp_tool_logs", mcpBefore), "superseded processing version must not trigger a flush mutation") + + found, err := store.FindByID(ctx, "ch-flush-superseded") + require.NoError(t, err) + assert.Equal(t, "success", found.Status) + mcp, err := store.FindMCPToolLog(ctx, "ch-flush-superseded-mcp") + require.NoError(t, err) + assert.Equal(t, "success", mcp.Status) +} + +// TestClickHouseFinalAfterLightweightDelete proves the _row_exists mask left +// by a lightweight delete cannot resurrect an older ReplacingMergeTree version +// of the row under the connection-level final=1 setting, and that the +// user-triggered deletes are lightweight too. +func TestClickHouseFinalAfterLightweightDelete(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + ts := time.Now().UTC().Truncate(time.Millisecond) + + for _, id := range []string{"ch-final-1", "ch-final-2", "ch-final-3"} { + require.NoError(t, store.CreateIfNotExists(ctx, chTestLog(id, ts))) + // A second ReplacingMergeTree version of the same id via read-modify-write. + require.NoError(t, store.Update(ctx, id, map[string]interface{}{"status": "success"})) + // Two physical versions must exist (FINAL would hide the older one), + // otherwise the resurrection case below is not actually exercised. + require.EqualValues(t, 2, chCountIDsNoFinal(t, store, "logs", []string{id}), "expected both versions of %s on disk before the delete", id) + require.Equal(t, int64(1), chCountRows(t, store.db, "logs", id), "FINAL collapses them to one logical row") + } + require.NoError(t, store.BatchCreateMCPToolLogsIfNotExists(ctx, []*MCPToolLog{chTestMCPToolLog("ch-final-mcp", ts)})) + + logsBefore := chMutationIDs(t, store.db, "logs") + mcpBefore := chMutationIDs(t, store.db, "mcp_tool_logs") + require.NoError(t, store.DeleteLog(ctx, "ch-final-1")) + require.NoError(t, store.DeleteLogs(ctx, []string{"ch-final-2", "ch-final-3"})) + require.NoError(t, store.DeleteMCPToolLogs(ctx, []string{"ch-final-mcp"})) + + for _, id := range []string{"ch-final-1", "ch-final-2", "ch-final-3"} { + _, err := store.FindByID(ctx, id) + assert.ErrorIs(t, err, ErrNotFound, "log %s should be deleted", id) + assert.Equal(t, int64(0), chCountRows(t, store.db, "logs", id), "no version of %s may survive under FINAL", id) + assert.EqualValues(t, 0, chCountIDsNoFinal(t, store, "logs", []string{id}), "both physical versions of %s must be masked, not just the newest", id) + } + _, err := store.FindMCPToolLog(ctx, "ch-final-mcp") + assert.ErrorIs(t, err, ErrNotFound) + + logsCmds := chNewMutationCommands(t, store.db, "logs", logsBefore) + require.Len(t, logsCmds, 2, "DeleteLog + DeleteLogs must issue one mutation each, got %v", logsCmds) + assertLightweightMutations(t, "logs", logsCmds) + mcpCmds := chNewMutationCommands(t, store.db, "mcp_tool_logs", mcpBefore) + require.Len(t, mcpCmds, 1, "DeleteMCPToolLogs must issue one mutation, got %v", mcpCmds) + assertLightweightMutations(t, "mcp_tool_logs", mcpCmds) +} + +func chEngineFull(t *testing.T, db *gorm.DB, table string) string { + t.Helper() + var engineFull string + require.NoError(t, db.Raw("SELECT engine_full FROM system.tables WHERE database = currentDatabase() AND name = ?", table).Scan(&engineFull).Error) + return engineFull +} + +// TestClickHouseTTLReconciledOnExistingTables covers the second half of +// #7098: CREATE TABLE IF NOT EXISTS never updates the TTL of an existing +// table, so a changed logs_store.retention_days must be reconciled with +// MODIFY TTL / REMOVE TTL at startup. +func TestClickHouseTTLReconciledOnExistingTables(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + retained := []string{"logs", "mcp_tool_logs", "webhook_deliveries"} + // removeTTLs strips any TTL from the shared tables. REMOVE TTL errors on + // a table that has none (BAD_ARGUMENTS), so it is only issued when one is + // present. It runs before the initial-state check, because retention 0 + // deliberately preserves whatever an interrupted earlier run left behind, + // and again on cleanup so every other test still sees TTL-free tables. + removeTTLs := func() { + t.Helper() + for _, table := range retained { + if strings.Contains(chEngineFull(t, store.db, table), "TTL ") { + require.NoError(t, store.db.Exec("ALTER TABLE `"+table+"` REMOVE TTL").Error, "%s: reset TTL", table) + } + } + } + removeTTLs() + t.Cleanup(removeTTLs) + // managedDays asserts the table carries exactly one Bifrost-managed TTL of + // want days, using the production parser so an appended or leftover rule + // cannot satisfy a looser substring check. + managedDays := func(table string, want int) { + t.Helper() + engineFull := chEngineFull(t, store.db, table) + days, ok := chTTLDaysFromEngineFull(engineFull) + require.True(t, ok, "%s: expected a managed TTL, got %q", table, engineFull) + assert.Equal(t, want, days, "%s: %q", table, engineFull) + assert.Equal(t, 1, strings.Count(engineFull, "TTL "), "%s: exactly one TTL clause expected in %q", table, engineFull) + } + + for _, table := range retained { + require.NotContains(t, chEngineFull(t, store.db, table), "TTL", "%s: fixture tables are created without retention", table) + } + + withTTL, err := newClickHouseLogStore(ctx, clickhouseTestConfig(), 3, testLogger{}) + require.NoError(t, err) + require.NoError(t, withTTL.Close(ctx)) + for _, table := range retained { + managedDays(table, 3) + } + managedDays("async_jobs", 7) + + // A different retention replaces the TTL in place. + longer, err := newClickHouseLogStore(ctx, clickhouseTestConfig(), 5, testLogger{}) + require.NoError(t, err) + require.NoError(t, longer.Close(ctx)) + for _, table := range retained { + managedDays(table, 5) + } + // Snapshot the complete definitions so the retention-zero restart is + // proven to leave them byte-for-byte unchanged, not merely still matching. + before := map[string]string{} + for _, table := range append(retained, "async_jobs") { + before[table] = chEngineFull(t, store.db, table) + } + + // Retention 0 (the default when the field is omitted) means "not managed + // by Bifrost": an existing TTL, including one an operator applied by hand + // as the #7098 workaround, must survive a restart. + unmanaged, err := newClickHouseLogStore(ctx, clickhouseTestConfig(), 0, testLogger{}) + require.NoError(t, err) + require.NoError(t, unmanaged.Close(ctx)) + for table, want := range before { + assert.Equal(t, want, chEngineFull(t, store.db, table), "%s: retention_days=0 must leave the table definition untouched", table) + } +} + func chTestMCPToolLog(id string, ts time.Time) *MCPToolLog { return &MCPToolLog{ ID: id, diff --git a/tests/docker-compose.yml b/tests/docker-compose.yml index 41c89d4e867..8b364765058 100644 --- a/tests/docker-compose.yml +++ b/tests/docker-compose.yml @@ -216,6 +216,33 @@ services: # docker image CI test. ipv4_address: 172.28.0.16 + # ClickHouse instance for the framework logstore tests + # (framework/logstore/clickhousestore_test.go connects to host port 9001; + # host 9000 is taken by weaviate above). Without this service every + # TestClickHouse* test is skipped, or fails when CI is set. + clickhouse: + image: clickhouse/clickhouse-server:24.8-alpine + environment: + CLICKHOUSE_DB: bifrost + CLICKHOUSE_USER: bifrost + CLICKHOUSE_PASSWORD: bifrost_password + ports: + - "9001:9000" + volumes: + - clickhouse_data:/var/lib/clickhouse + ulimits: + nofile: + soft: 262144 + hard: 262144 + healthcheck: + test: ["CMD", "wget", "--spider", "-q", "http://127.0.0.1:8123/ping"] + interval: 5s + timeout: 5s + retries: 10 + networks: + bifrost_network: + ipv4_address: 172.28.0.20 + networks: bifrost_network: driver: bridge @@ -227,4 +254,5 @@ networks: volumes: weaviate_data: redis_data: - qdrant_data: + qdrant_data: + clickhouse_data: