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
1 change: 1 addition & 0 deletions .github/workflows/backend-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -624,6 +624,7 @@ jobs:
./internal/notifications/store
./internal/office/configsync
./internal/office/repository/sqlite
./internal/office/retention
./internal/orchestrator/messagequeue
./internal/persistence
./internal/persistence/storeconformance
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -269,17 +269,18 @@ func TestAggregator_FailedPausedPushIsRetriedOnUnchangedContribution(t *testing.
w.WriteHeader(http.StatusBadRequest)
return
}
modes <- body.Mode
if body.Mode == string(WorkspacePollModePaused) {
mu.Lock()
pausedCalls++
call := pausedCalls
mu.Unlock()
if call == 1 {
modes <- body.Mode
w.WriteHeader(http.StatusServiceUnavailable)
return
}
}
modes <- body.Mode
w.WriteHeader(http.StatusOK)
}))
t.Cleanup(srv.Close)
Expand Down
22 changes: 22 additions & 0 deletions apps/backend/internal/backendapp/helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ import (
notificationhandlers "github.com/kandev/kandev/internal/notifications/handlers"
officeagents "github.com/kandev/kandev/internal/office/agents"
officesqlite "github.com/kandev/kandev/internal/office/repository/sqlite"
"github.com/kandev/kandev/internal/office/retention"
officetestharness "github.com/kandev/kandev/internal/office/testharness"
"github.com/kandev/kandev/internal/orchestrator"
"github.com/kandev/kandev/internal/org"
Expand Down Expand Up @@ -1521,6 +1522,7 @@ func registerSecondaryRoutes(

registerHealthRoutes(p)
registerSystemRoutes(p)
registerRetentionRoutes(p)
if p.runtimeFlagsSvc != nil {
runtimeflags.RegisterRoutes(p.router, p.runtimeFlagsSvc)
}
Expand Down Expand Up @@ -1722,6 +1724,23 @@ func registerSystemRoutes(p routeParams) {
p.systemSvc.RegisterRoutes(p.router, p.log)
}

// registerRetentionRoutes mounts GET/PUT /api/v1/system/retention. It is a
// separate group from systemSvc's own /api/v1/system group (rather than a
// field on system.Service) because internal/office/retention cannot be
// imported by internal/system without inverting the existing system ->
// office dependency direction; gin allows two RouterGroups to share a path
// prefix as long as no route collides, and none does here. Read/admin
// split mirrors system.Service.RegisterRoutes: GET is member-readable,
// PUT requires the admin-scoped settings-manage permission.
func registerRetentionRoutes(p routeParams) {
if p.services == nil || p.services.Retention == nil {
return
}
read := p.router.Group("/api/v1/system")
admin := read.Group("", authz.RequireOrgScope(authz.ScopeOrgSettingsManage))
retention.RegisterRoutes(read, admin, p.services.Retention.Handler)
}

// registerHealthRoutes sets up the system health endpoint with all health checkers.
func registerHealthRoutes(p routeParams) {
var githubProvider health.GitHubStatusProvider
Expand All @@ -1747,6 +1766,9 @@ func registerHealthRoutes(p routeParams) {
if p.systemSvc != nil && p.systemSvc.StorageRuntime != nil {
checkers = append(checkers, p.systemSvc.StorageRuntime)
}
if p.services != nil && p.services.Retention != nil {
checkers = append(checkers, p.services.Retention.Checker)
}
healthSvc := health.NewService(p.log, checkers...)
health.RegisterRoutes(p.router, healthSvc, p.log)
}
Expand Down
11 changes: 11 additions & 0 deletions apps/backend/internal/backendapp/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ import (
officepause "github.com/kandev/kandev/internal/office/pause"
officeprojects "github.com/kandev/kandev/internal/office/projects"
officesqlite "github.com/kandev/kandev/internal/office/repository/sqlite"
"github.com/kandev/kandev/internal/office/retention"
officeroutines "github.com/kandev/kandev/internal/office/routines"
"github.com/kandev/kandev/internal/office/routing"
officescheduler "github.com/kandev/kandev/internal/office/scheduler"
Expand Down Expand Up @@ -1181,6 +1182,16 @@ func startGatewayAndServe(
})
systemSvc.Storage = storageComposition.handler
systemSvc.StorageRuntime = storageComposition.runtime

// Office run history retention: bounds office_routine_runs, runs, and
// their satellites on its own interval, separate from the 5s Office
// tick. Kept regardless of the Office feature flag — see Services.Retention.
services.Retention = retention.NewRuntime(dbPool, repos.SystemSettings,
func(message string, err error) { log.Error(message, zap.Error(err)) })
if err := services.Retention.Start(ctx); err != nil {
log.Warn("office run retention scheduler failed to start", zap.Error(err))
}
addCleanup(func() error { services.Retention.Stop(); return nil })
if systemSvc.LogBundles != nil {
systemSvc.LogBundles.SetNotifier(gateway.Hub)
systemSvc.LogBundles.SetSessionProvider(newDiagnosticSessionProvider(services.Task))
Expand Down
7 changes: 7 additions & 0 deletions apps/backend/internal/backendapp/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import (
notificationstore "github.com/kandev/kandev/internal/notifications/store"
office "github.com/kandev/kandev/internal/office"
officesqlite "github.com/kandev/kandev/internal/office/repository/sqlite"
"github.com/kandev/kandev/internal/office/retention"
officeservice "github.com/kandev/kandev/internal/office/service"
"github.com/kandev/kandev/internal/org"
"github.com/kandev/kandev/internal/orgunit"
Expand Down Expand Up @@ -112,6 +113,12 @@ type Services struct {
// WorktreeMgr is the worktree manager. Exposed here so the install-wide
// storage-maintenance composition can reach it for workspace cleanup.
WorktreeMgr *worktree.Manager
// Retention owns the office_routine_runs/runs history sweep scheduler,
// its HTTP surface, and its health checker. Kept regardless of the
// Office feature flag, matching every other required-schema owner: rows
// written while Office was enabled still need bounding after it is
// turned off.
Retention *retention.Runtime
// Terminal is the first-class user-terminal service (rename, park, etc.).
// Wired into the gateway once lifecycle.Manager is up so the PTY backend
// is available.
Expand Down
27 changes: 27 additions & 0 deletions apps/backend/internal/office/repository/sqlite/base_migrations.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,9 @@ func (r *Repository) runMigrations() error {
r.migrateBudgetPolicyRevision()
r.migrateWorkspacePauseSkipAttribution()
r.migrateLoopLivenessCausationID()
if err := r.migrateRetentionIndexes(); err != nil {
return err
}
if err := r.migrate.Err(); err != nil {
return err
}
Expand Down Expand Up @@ -141,6 +144,27 @@ func (r *Repository) migrateLoopLivenessCausationID() {
ON runs(causation_id) WHERE causation_id != ''`)
}

// migrateRetentionIndexes adds the two indexes the run-history retention
// sweep depends on (docs/specs/office/system-design/run-history-retention.md
// "Indexes to add"). Both are expression indexes over the same
// COALESCE(...) the sweep both filters and orders by; a plain-column index
// on the nullable completion column would serve neither the WHERE clause
// nor the ORDER BY the sweep actually issues, on either engine.
func (r *Repository) migrateRetentionIndexes() error {
if err := r.migrate.Apply(
"idx_office_routine_runs_retention",
`CREATE INDEX IF NOT EXISTS idx_office_routine_runs_retention
ON office_routine_runs(routine_id, status, (COALESCE(completed_at, created_at)) DESC, id DESC)`,
); err != nil {
return err
}
return r.migrate.Apply(
"idx_runs_retention",
`CREATE INDEX IF NOT EXISTS idx_runs_retention
ON runs(agent_profile_id, status, (COALESCE(finished_at, requested_at)) DESC, id DESC)`,
)
}

// migrateContinuationScope adds runs.continuation_scope for databases
// created before WO-16's claim-time scope persistence. Existing rows receive
// a scope from their stored context snapshot so queued or claimed taskless
Expand Down Expand Up @@ -418,6 +442,9 @@ func (r *Repository) migrateFailureColumns() error {
if _, err := r.db.Exec(`CREATE INDEX IF NOT EXISTS idx_office_agent_pause_recoveries_agent ON office_agent_pause_recoveries(agent_id)`); err != nil {
return fmt.Errorf("idx_office_agent_pause_recoveries_agent: %w", err)
}
if _, err := r.db.Exec(`CREATE INDEX IF NOT EXISTS idx_office_agent_pause_recoveries_failed_run ON office_agent_pause_recoveries(failed_run_id)`); err != nil {
return fmt.Errorf("idx_office_agent_pause_recoveries_failed_run: %w", err)
}
return nil
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
package sqlite_test

import (
"testing"

"github.com/kandev/kandev/internal/office/repository/sqlite"
taskrepo "github.com/kandev/kandev/internal/task/repository/sqlite"
"github.com/kandev/kandev/internal/testutil"
)

// TestPostgresRetentionIndexes_CreatedFreshAndReplaySafe is the PostgreSQL
// half of TestRetentionIndexes_CreatedFreshAndReplaySafe: the two
// expression indexes the retention sweep depends on must exist there too,
// with identical CREATE INDEX IF NOT EXISTS replay safety. Skips unless
// KANDEV_TEST_POSTGRES_DSN is set.
func TestPostgresRetentionIndexes_CreatedFreshAndReplaySafe(t *testing.T) {
dsn := testutil.PostgresDSNFromEnv(t)
conn := testutil.OpenIsolatedPostgres(t, dsn)

// tasks is created by the task repository's schema init, mirroring
// production boot order (see child_summaries_postgres_test.go).
if _, err := taskrepo.NewWithDB(conn, conn, nil); err != nil {
t.Fatalf("init task repo: %v", err)
}
if _, err := sqlite.NewWithDB(conn, conn, nil); err != nil {
t.Fatalf("fresh NewWithDB: %v", err)
}
assertPostgresIndexExists(t, conn, "idx_office_routine_runs_retention")
assertPostgresIndexExists(t, conn, "idx_runs_retention")
assertPostgresIndexExists(t, conn, "idx_office_agent_pause_recoveries_failed_run")

if _, err := sqlite.NewWithDB(conn, conn, nil); err != nil {
t.Fatalf("replay NewWithDB: %v", err)
}
assertPostgresIndexExists(t, conn, "idx_office_routine_runs_retention")
assertPostgresIndexExists(t, conn, "idx_runs_retention")
assertPostgresIndexExists(t, conn, "idx_office_agent_pause_recoveries_failed_run")
}

func assertPostgresIndexExists(t *testing.T, conn interface {
Get(dest interface{}, query string, args ...interface{}) error
}, name string) {
t.Helper()
var count int
if err := conn.Get(&count,
`SELECT COUNT(*) FROM pg_indexes WHERE indexname = $1`, name,
); err != nil {
t.Fatalf("query pg_indexes for %s: %v", name, err)
}
if count != 1 {
t.Fatalf("index %s: found %d, want 1", name, count)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package sqlite_test

import (
"testing"

"github.com/jmoiron/sqlx"
_ "github.com/mattn/go-sqlite3"

"github.com/kandev/kandev/internal/office/repository/sqlite"
)

// TestRetentionIndexes_CreatedFreshAndReplaySafe proves the retention indexes
// exist after a fresh boot and that re-running schema init against the same
// database (the upgrade-path replay) is a no-op, not an error.
func TestRetentionIndexes_CreatedFreshAndReplaySafe(t *testing.T) {
conn, err := sqlx.Open("sqlite3", ":memory:")
if err != nil {
t.Fatalf("open sqlite: %v", err)
}
conn.SetMaxOpenConns(1)
t.Cleanup(func() { _ = conn.Close() })

if _, err := sqlite.NewWithDB(conn, conn, nil); err != nil {
t.Fatalf("fresh NewWithDB: %v", err)
}
assertIndexExists(t, conn, "idx_office_routine_runs_retention")
assertIndexExists(t, conn, "idx_runs_retention")
assertIndexExists(t, conn, "idx_office_agent_pause_recoveries_failed_run")

// Replay: schema init against the same, already-initialized database.
if _, err := sqlite.NewWithDB(conn, conn, nil); err != nil {
t.Fatalf("replay NewWithDB: %v", err)
}
assertIndexExists(t, conn, "idx_office_routine_runs_retention")
assertIndexExists(t, conn, "idx_runs_retention")
assertIndexExists(t, conn, "idx_office_agent_pause_recoveries_failed_run")
}

func assertIndexExists(t *testing.T, conn *sqlx.DB, name string) {
t.Helper()
var count int
if err := conn.Get(&count,
`SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = ?`, name,
); err != nil {
t.Fatalf("query sqlite_master for %s: %v", name, err)
}
if count != 1 {
t.Fatalf("index %s: found %d, want 1", name, count)
}
}
Loading
Loading