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: 17 additions & 4 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,6 @@
#
# Mirrors `make test`: build tag duckdb_arrow is mandatory (Arc fails closed
# without it, #598) and the suite runs under -race.
#
# gofmt is deliberately NOT gated yet: a handful of pre-existing files
# (compaction/reconciliation/scheduler) are gofmt-dirty on main. Gate it once
# that debt is cleared.
name: CI

on:
Expand Down Expand Up @@ -36,6 +32,23 @@ jobs:
go-version-file: go.mod
cache: true

# Gated as of #875. The debt this used to wait on — eight files in
# compaction, reconciliation and scheduler — is cleared, and the project
# instructions have always told contributors this must be clean, so
# leaving it ungated trained people to ignore the output.
#
# Whole repo, not just ./internal ./cmd: pkg/ and scripts/ are Go too,
# they are clean today, so covering them is free now and closes the hole
# permanently. It also matches `make fmt`, which formats the whole tree.
- name: Format
run: |
unformatted=$(gofmt -l .)
if [ -n "$unformatted" ]; then
echo "These files are not gofmt-clean. Run: gofmt -w ."
echo "$unformatted"
exit 1
fi

- name: Build
run: go build -tags=duckdb_arrow ./...

Expand Down
8 changes: 4 additions & 4 deletions internal/compaction/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,10 +42,10 @@ type Scheduler struct {
// paths and by existing tests that predate Phase 4.
clusterGate ClusterGate

cron *cron.Cron
running bool
roleGated bool // true when Start found CanCompact=false; used by Status
stopCh chan struct{}
cron *cron.Cron
running bool
roleGated bool // true when Start found CanCompact=false; used by Status
stopCh chan struct{}

logger zerolog.Logger
mu sync.Mutex
Expand Down
13 changes: 7 additions & 6 deletions internal/compaction/watcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,8 @@ var ErrNotLeader = errors.New("compaction bridge: not the Raft leader")
// The bridge converts CompactedFile to raft.FileEntry inside the cluster
// package, keeping the compaction package free of any raft.* imports.
type CompactedFile struct {
Path string // storage-relative path
SHA256 string // hex-encoded 64 chars
Path string // storage-relative path
SHA256 string // hex-encoded 64 chars
SizeBytes int64
Database string
Measurement string
Expand Down Expand Up @@ -120,10 +120,11 @@ type CompletionWatcherConfig struct {
// pending manifests via the ManifestBridge. Run one per compactor node.
//
// Lifecycle:
// w := NewCompletionWatcher(cfg)
// w.Start(ctx) // kicks off the poll loop in a background goroutine
// ...
// w.Stop() // signals stop, waits for the loop to drain one final poll
//
// w := NewCompletionWatcher(cfg)
// w.Start(ctx) // kicks off the poll loop in a background goroutine
// ...
// w.Stop() // signals stop, waits for the loop to drain one final poll
//
// The watcher is safe to Start/Stop multiple times. It is NOT safe for
// concurrent Start calls from different goroutines.
Expand Down
2 changes: 1 addition & 1 deletion internal/reconciliation/diff_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ func TestComputeDiff_BothKindsMixed(t *testing.T) {
storage := []objectRecord{
{path: "db/m/a.parquet", lastModified: now.Add(-2 * time.Hour)},
{path: "db/m/c.parquet", lastModified: now.Add(-2 * time.Hour)},
{path: "db/m/orphan.parquet", lastModified: now.Add(-48 * time.Hour)}, // missing from manifest
{path: "db/m/orphan.parquet", lastModified: now.Add(-48 * time.Hour)}, // missing from manifest
{path: "db/m/young.parquet", lastModified: now.Add(-30 * time.Minute)}, // grace skip
}
d := computeDiff(manifest, storage, now, 24*time.Hour)
Expand Down
14 changes: 7 additions & 7 deletions internal/reconciliation/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -282,12 +282,12 @@ const (
// Run is a single reconcile-cycle summary. Held in the reconciler's ring
// buffer of recent runs and surfaced via Status().
type Run struct {
ID string `json:"id"`
StartedAt time.Time `json:"started_at"`
FinishedAt time.Time `json:"finished_at"`
DryRun bool `json:"dry_run"`
ID string `json:"id"`
StartedAt time.Time `json:"started_at"`
FinishedAt time.Time `json:"finished_at"`
DryRun bool `json:"dry_run"`
BackendKind BackendKind `json:"backend_kind"`
Role string `json:"role"`
Role string `json:"role"`

// Counts
ManifestFileCount int `json:"manifest_file_count"`
Expand Down Expand Up @@ -499,8 +499,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, dryRun bool) (*Run, error) {
defer cancel()

run := &Run{
ID: uuid.NewString(),
StartedAt: time.Now().UTC(),
ID: uuid.NewString(),
StartedAt: time.Now().UTC(),
// Honor the caller's dryRun exactly. The cron path
// (scheduler.tick) already passes cfg.ManifestOnlyDryRun, so the
// safety policy is enforced upstream. OR'ing it back in here
Expand Down
16 changes: 8 additions & 8 deletions internal/reconciliation/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -998,14 +998,14 @@ func TestLooksLikeManagedPath(t *testing.T) {
path string
want bool
}{
{"db/m/2026/04/27/12/file.parquet", true}, // canonical 7-segment
{"db/m/2026/04/27/12/sub/file.parquet", true}, // 8 segments — fine
{"db/m/file.parquet", false}, // too short
{"db/m/2026/04/27/file.parquet", false}, // 6 segments — too short
{"db/m/2026/04/27/12/", false}, // empty trailing segment
{"../etc/passwd/2026/04/27/12/x.parquet", false}, // .. segment
{"db/m/./27/04/12/x.parquet", false}, // . segment
{"//m/2026/04/27/12/file.parquet", false}, // empty leading segment
{"db/m/2026/04/27/12/file.parquet", true}, // canonical 7-segment
{"db/m/2026/04/27/12/sub/file.parquet", true}, // 8 segments — fine
{"db/m/file.parquet", false}, // too short
{"db/m/2026/04/27/file.parquet", false}, // 6 segments — too short
{"db/m/2026/04/27/12/", false}, // empty trailing segment
{"../etc/passwd/2026/04/27/12/x.parquet", false}, // .. segment
{"db/m/./27/04/12/x.parquet", false}, // . segment
{"//m/2026/04/27/12/file.parquet", false}, // empty leading segment
}
for _, c := range cases {
got := looksLikeManagedPath(c.path)
Expand Down
20 changes: 10 additions & 10 deletions internal/reconciliation/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,17 +180,17 @@ func (s *Scheduler) Status() map[string]interface{} {

cfg := s.reconciler.cfg
out := map[string]interface{}{
"enabled": cfg.Enabled,
"running": running,
"schedule": s.schedule,
"backend_kind": string(cfg.BackendKind),
"grace_window": cfg.GraceWindow.String(),
"clock_skew_allowance": cfg.ClockSkewAllowance.String(),
"max_run_duration": cfg.MaxRunDuration.String(),
"max_manifest_size": cfg.MaxManifestSize,
"max_deletes_per_run": cfg.MaxDeletesPerRun,
"enabled": cfg.Enabled,
"running": running,
"schedule": s.schedule,
"backend_kind": string(cfg.BackendKind),
"grace_window": cfg.GraceWindow.String(),
"clock_skew_allowance": cfg.ClockSkewAllowance.String(),
"max_run_duration": cfg.MaxRunDuration.String(),
"max_manifest_size": cfg.MaxManifestSize,
"max_deletes_per_run": cfg.MaxDeletesPerRun,
"manifest_only_dry_run": cfg.ManifestOnlyDryRun,
"in_flight": s.reconciler.IsRunning(),
"in_flight": s.reconciler.IsRunning(),
}
if running {
out["next_run"] = s.nextRun().UTC().Format(time.RFC3339)
Expand Down
1 change: 0 additions & 1 deletion internal/scheduler/cq_scheduler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -307,4 +307,3 @@ func TestCQScheduler_ClusterGate_FailoverTransition(t *testing.T) {
t.Error("after failover gate should return IsPrimaryWriter=true without restart")
}
}

8 changes: 4 additions & 4 deletions internal/scheduler/retention_scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,10 @@ type RetentionScheduler struct {
clusterGate RetentionClusterGate
schedule string // Cron schedule (e.g., "0 3 * * *" = 3am daily)
cron *cron.Cron
running bool
runningJob bool // true while a retention cycle is in progress; prevents overlap
mu sync.Mutex
logger zerolog.Logger
running bool
runningJob bool // true while a retention cycle is in progress; prevents overlap
mu sync.Mutex
logger zerolog.Logger
}

// RetentionSchedulerConfig holds configuration for the retention scheduler
Expand Down
Loading