diff --git a/.gitignore b/.gitignore index ae747857..0ab73e55 100644 --- a/.gitignore +++ b/.gitignore @@ -83,3 +83,4 @@ ui/bun.lock internal/api/docs/spec.go opencode.json +livereview.dtx diff --git a/docs/Review/Event-compression.md b/docs/Review/Event-compression.md new file mode 100644 index 00000000..64948f4d --- /dev/null +++ b/docs/Review/Event-compression.md @@ -0,0 +1,215 @@ +# Event Log Compaction Architecture + +This document describes the design, implementation, database impact, and SuperAdmin management controls for automated event log compaction (`review_events`) in LiveReview. + +--- + +## 1. Overview & Problem Statement + +Production PostgreSQL analysis revealed that `review_events` accounted for **1.57 GB (4.89M rows)** out of the overall database size. **95.8% (4.07M rows)** consisted of `event_type = 'log'`, primarily caused by verbose LLM streaming chunk lines (`Streaming chunk 1...35`), divider formatting (`========`), and hunk line-number debug prints. + +Core AI review findings, inline comments, file diffs, and cost accounting reside in separate, untouched tables (`reviews`, `ai_comments`, `loc_usage_ledger`, `tool_credit_ledger`) and blob artifacts. + +**Solution:** Implement an automated, in-process Go compaction manager (`EventCompactionManager`) using Go Cron (`github.com/robfig/cron/v3`) that runs on a configurable schedule (default: **daily at 02:00 AM IST / 20:30 UTC** (`30 20 * * *`)). Before deleting raw debug logs older than the retention threshold (default: **30 days**), the manager writes a **Compaction Summary Marker** containing original total event count. + +--- + +## 2. Code Locations + +| Purpose | File | +|---|---| +| Compaction Manager (scheduler, executor) | [`internal/api/event_compaction.go`](file:///home/gk/hex/LiveReview/internal/api/event_compaction.go) | +| SuperAdmin API Endpoints (GET/PUT/POST) | [`internal/api/compaction_settings.go`](file:///home/gk/hex/LiveReview/internal/api/compaction_settings.go) | +| Server Lifecycle (Start/Stop hooks) | [`internal/api/server.go`](file:///home/gk/hex/LiveReview/internal/api/server.go) | +| Settings UI Tab (SuperAdmin only) | [`ui/src/pages/Settings/CompactionSettingsTab.tsx`](file:///home/gk/hex/LiveReview/ui/src/pages/Settings/CompactionSettingsTab.tsx) | +| Settings Tab Registration & Access Gate | [`ui/src/pages/Settings/Settings.tsx`](file:///home/gk/hex/LiveReview/ui/src/pages/Settings/Settings.tsx) | +| Mega Menu Entry (SuperAdmin only) | [`ui/src/components/Navbar/megaMenuData.ts`](file:///home/gk/hex/LiveReview/ui/src/components/Navbar/megaMenuData.ts) | +| Unit Tests | [`internal/api/event_compaction_test.go`](file:///home/gk/hex/LiveReview/internal/api/event_compaction_test.go) | + +--- + +## 3. Go Cron Scheduler Architecture + +The event compaction manager uses Go Cron (`github.com/robfig/cron/v3`) and is initialized in [`server.go`](file:///home/gk/hex/LiveReview/internal/api/server.go) at startup. + +```mermaid +flowchart TD + subgraph Server_Startup ["Server Startup (server.go)"] + A["Start API Server"] --> B["NewEventCompactionManager(db)"] + B --> C["loadSettingsFromDB() — reads system_settings row"] + C --> D["manager.Start() — validates cron, falls back to default if invalid"] + end + + subgraph Go_Cron_Loop ["Go Cron Scheduler (In-Process)"] + D --> E["cron.AddFunc(cronExpr, runCycle)"] + E --> F{"Cron Trigger fires"} + F --> G["runCycle() — atomic CAS: already running?"] + G --> |"Yes"| SKIP["Skip — log warn and return"] + G --> |"No"| H{"enabled?"} + H --> |"No"| I["Skip — log and return"] + H --> |"Yes"| J["executeBulkCompaction(ctx, retentionDays)"] + end + + subgraph Manual_Trigger ["Manual 'Run Now' (compaction_settings.go)"] + K["POST /api/v1/admin/settings/compaction/run"] --> L["go TriggerManualCycle()"] + L --> G + end + + subgraph Execution ["executeBulkCompaction — PostgreSQL"] + J --> M["Step 1: INSERT summary markers for all uncompacted reviews (1 query)"] + M --> N["Step 2: Loop — DELETE 50,000 rows per batch until 0 rows remain"] + N --> O{"ctx.Done()?"} + O --> |"Yes"| P["Cancel gracefully, log rows deleted so far"] + O --> |"No"| N + end +``` + +### Key Architectural Points + +1. **No Distributed Lock**: Compaction runs in the **backend process** (single instance). A distributed advisory lock is not needed — concurrent execution from multiple backend processes is not a deployment scenario. The `lockStore` field and `compactionLeaderLocker` interface have been removed entirely from `EventCompactionManager`. + +2. **Concurrent Cycle Guard**: `runCycle()` uses an `atomic.Int32` flag (`running`) — a `CompareAndSwap(0,1)` at entry ensures only one cycle runs at a time. If the cron fires while a manual "Run Now" is still in progress, the second invocation logs a warning and returns immediately. This prevents duplicate summary markers and DB contention. + +3. **Invalid Cron Fallback**: `Start()` does not crash the server on a bad cron expression stored in DB. It logs a warning and falls back to `defaultCompactionCronExpr` (`30 20 * * *`) automatically. + +4. **Dynamic Config Reload**: `UpdateConfig(enabled, cronExpr, retentionDays)` updates `m.enabled`, `m.retentionDays`, and `m.cronExpr` in memory. If `cronExpr` changed and `m.cronRunner` is running, it removes the old entry and adds a new one — **no server restart required**. + +5. **Graceful Shutdown**: `Stop()` calls `m.cronRunner.Stop()` (which waits up to 5 seconds for running jobs to finish) then cancels `m.ctx`, which causes any in-progress batch delete loop to exit cleanly via `ctx.Done()`. + +--- + +## 4. executeBulkCompaction — Exact SQL + +### Step 1: INSERT Summary Markers (1 query, runs once per cycle) + +Inserts one summary marker row per review that has not been compacted yet and has events older than `retentionDays`: + +```sql +INSERT INTO public.review_events (review_id, org_id, ts, event_type, level, data) +SELECT + re.review_id, + re.org_id, + NOW(), + 'log', + 'info', + jsonb_build_object( + 'message', 'Review log events compacted', + 'compacted', true, + 'original_total_event_count', COUNT(*) + ) +FROM public.review_events re +WHERE re.ts < NOW() - ($1 * INTERVAL '1 day') + AND NOT EXISTS ( + SELECT 1 FROM public.review_events cx + WHERE cx.review_id = re.review_id AND cx.org_id = re.org_id + AND cx.event_type = 'log' AND (cx.data->>'compacted')::boolean = true + ) +GROUP BY re.review_id, re.org_id; +``` + +> **`NOT EXISTS` guard**: Ensures idempotency — if a summary marker already exists for a review, that review is skipped. Running compaction twice is always safe. + +### Step 2: Batched DELETE (50,000 rows/batch, loop until 0 rows remain) + +Instead of a single unbounded DELETE (which causes PostgreSQL statement timeouts on millions of rows), deletion runs in safe chunks using `ctid`: + +```sql +DELETE FROM public.review_events +WHERE ctid IN ( + SELECT ctid FROM public.review_events + WHERE ts < NOW() - ($1 * INTERVAL '1 day') + AND event_type = 'log' + AND COALESCE(level, 'info') NOT IN ('error', 'warn') + AND (data->>'compacted')::boolean IS NOT TRUE + AND data->>'message' NOT ILIKE '%started%' + AND data->>'message' NOT ILIKE '%completed%' + AND data->>'message' NOT ILIKE '%posted%' + AND data->>'message' NOT ILIKE '%generated%' + AND data->>'message' NOT ILIKE '%CLI DIFF REVIEW STARTED%' + LIMIT 50000 +); +``` + +The Go loop runs this until `RowsAffected() == 0`. + +--- + +## 5. What Is Kept vs. Deleted + +| Event Category | Rule | +|---|---| +| **Compaction Summary Marker** (`compacted: true`) | **NEVER deleted** — `NOT TRUE` guard in DELETE | +| **All non-log event types** (`status`, `completion`, `artifact`, `batch`, etc.) | **NEVER deleted** — `event_type = 'log'` filter in DELETE | +| **Error & Warning logs** (`level = 'error'` or `'warn'`) | **NEVER deleted** — `NOT IN ('error', 'warn')` guard | +| **Stage milestone messages** (contain `started`, `completed`, `posted`, `generated`, `CLI DIFF REVIEW STARTED`) | **NEVER deleted** — `NOT ILIKE` guards | +| **Verbose streaming debug logs** (info/debug/success level, no milestone keyword) | **DELETED** | + +--- + +## 6. GET /api/v1/admin/settings/compaction — 1 Query Only + +`GetCompactionSettings` in [`compaction_settings.go`](file:///home/gk/hex/LiveReview/internal/api/compaction_settings.go) executes **exactly 1 query**: + +```sql +SELECT data FROM system_settings WHERE name = 'event_compaction_settings'; +``` + +Response time: **< 1ms**. No stats, no counts, no scans of `review_events`. + +The response includes the stored config plus a `schedule_human` field generated by `describeCronSchedule()` (pure Go, no DB call): + +```json +{ + "enabled": true, + "cron_expression": "0 2 * * *", + "retention_days": 30, + "schedule_human": "Daily at 02:00 UTC" +} +``` + +--- + +## 7. SuperAdmin Access Control (AGENTS.md compliance) + +Per `AGENTS.md` — all system-level settings are **Super Admin only**: + +- **`Settings.tsx`**: Storage tab gated with `isSuperAdmin` only. +- **`megaMenuData.ts`**: Log Compaction mega menu link gated with `(ctx) => ctx.isSuperAdmin`. +- **Backend**: All three endpoints (`GET`, `PUT`, `POST`) are registered under `authMiddleware.RequireSuperAdmin()` in `server.go`. + +> No org-level `org_id` scoping applies here because `system_settings` is a globally public system config (not an org-scoped resource), consistent with the AGENTS.md rule: *"No Global Fallbacks: Never write global queries that omit `org_id` filters, unless the resource is globally public (e.g. system configs)."* + +--- + +## 8. Verified Database Test Results + +### Live Test on hexmos-internal (Org ID 151) + +| Review ID | Events Before | Events After | Rows Deleted | Reduction | +|---|---|---|---|---| +| **#8288** (1,225 log events) | **1,252** | **186** | **1,039** | **83%** | +| **#8529** (138 log events) | **145** | **40** | **105** | **73%** | + +### Full Database Pass (all orgs, retention_days=30) + +- **Total Reviews Compacted**: **10,211 reviews** +- **Verbose Log Rows Deleted**: **3,419,193 rows** (69 batches of 50,000) +- **Total Execution Time**: **29.3 seconds** — zero statement timeouts +- **Summary Markers Inserted**: **10,211** (one per eligible review) + +--- + +## 9. Unit Tests + +Tests live at [`internal/api/event_compaction_test.go`](file:///home/gk/hex/LiveReview/internal/api/event_compaction_test.go): + +| Test | What It Verifies | +|---|---| +| `TestDefaultCompactionSettingsConfig` | Default fallback values: `enabled=true`, `retention_days=30`, `cron_expression="0 2 * * *"` | +| `TestDescribeCronSchedule` | `describeCronSchedule()` correctly formats daily, every-N-hours, and multi-time-daily expressions | +| `TestEventCompactionManager_ConfigUpdate` | `UpdateConfig()` correctly updates `enabled`, `retentionDays`, and `cronExpr` in memory | + +Run with: +```bash +go test -v ./internal/api -run TestEventCompaction +``` diff --git a/go.mod b/go.mod index 7fa8049c..07b96577 100644 --- a/go.mod +++ b/go.mod @@ -38,12 +38,13 @@ require ( github.com/google/uuid v1.6.0 github.com/jackc/pgx/v5 v5.10.0 github.com/kaptinlin/jsonrepair v0.2.3 - github.com/labstack/echo/v4 v4.15.1 + github.com/labstack/echo/v4 v4.15.3 github.com/lib/pq v1.12.0 github.com/mdombrov-33/go-promptguard v0.4.0 github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 github.com/riverqueue/river v0.32.0 github.com/riverqueue/river/riverdriver/riverpgxv5 v0.32.0 + github.com/robfig/cron/v3 v3.0.1 github.com/rs/zerolog v1.34.0 github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06 github.com/shrsv/dbctx v0.1.2 @@ -161,7 +162,7 @@ require ( github.com/klauspost/pgzip v1.2.6 // indirect github.com/knadh/koanf/maps v0.1.1 // indirect github.com/kylelemons/godebug v1.1.0 // indirect - github.com/labstack/gommon v0.4.2 // indirect + github.com/labstack/gommon v0.5.0 // indirect github.com/lucasb-eyer/go-colorful v1.2.0 // indirect github.com/magiconair/properties v1.8.10 // indirect github.com/mattn/go-colorable v0.1.14 // indirect @@ -193,7 +194,6 @@ require ( github.com/riverqueue/river/rivershared v0.32.0 // indirect github.com/riverqueue/river/rivertype v0.32.0 // indirect github.com/rivo/uniseg v0.4.7 // indirect - github.com/robfig/cron/v3 v3.0.1 // indirect github.com/russross/blackfriday/v2 v2.1.0 // indirect github.com/sagikazarmark/locafero v0.7.0 // indirect github.com/sagikazarmark/slog-shim v0.1.0 // indirect diff --git a/go.sum b/go.sum index 44d5c813..690c12ca 100644 --- a/go.sum +++ b/go.sum @@ -369,10 +369,10 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= -github.com/labstack/echo/v4 v4.15.1 h1:S9keusg26gZpjMmPqB5hOEvNKnmd1lNmcHrbbH2lnFs= -github.com/labstack/echo/v4 v4.15.1/go.mod h1:xmw1clThob0BSVRX1CRQkGQ/vjwcpOMjQZSZa9fKA/c= -github.com/labstack/gommon v0.4.2 h1:F8qTUNXgG1+6WQmqoUWnz8WiEU60mXVVw0P4ht1WRA0= -github.com/labstack/gommon v0.4.2/go.mod h1:QlUFxVM+SNXhDL/Z7YhocGIBYOiwB0mXm1+1bAPHPyU= +github.com/labstack/echo/v4 v4.15.3 h1:lIdG4kK5RdMyhwCwSc4AmSQsLBb3AVwok6S8PX/9kwQ= +github.com/labstack/echo/v4 v4.15.3/go.mod h1:Xzp1Ns1RA2c9fY7nSgUJkpkUZGNbEIVHZbtbOMPktBI= +github.com/labstack/gommon v0.5.0 h1:6VSQ2NOzsnEJ5W6+84E0RbcaDDmgB6NIAzWCczTEe6c= +github.com/labstack/gommon v0.5.0/go.mod h1:Rzlg7HHy1maLfzBYGg9NZcVuz1sA68HHhLjhcEllYE0= github.com/lib/pq v1.12.0 h1:mC1zeiNamwKBecjHarAr26c/+d8V5w/u4J0I/yASbJo= github.com/lib/pq v1.12.0/go.mod h1:/p+8NSbOcwzAEI7wiMXFlgydTwcgTr3OSKMsD2BitpA= github.com/lucasb-eyer/go-colorful v1.2.0 h1:1nnpGOrhyZZuNyfu1QjKiUICQ74+3FNCN69Aj6K7nkY= diff --git a/internal/api/compaction_settings.go b/internal/api/compaction_settings.go new file mode 100644 index 00000000..ad1615ce --- /dev/null +++ b/internal/api/compaction_settings.go @@ -0,0 +1,140 @@ +package api + +import ( + "encoding/json" + "fmt" + "net/http" + "strconv" + "strings" + + "github.com/labstack/echo/v4" + "github.com/rs/zerolog/log" +) + +// CompactionSettingsConfig holds configuration for event log compaction. +type CompactionSettingsConfig struct { + Enabled bool `json:"enabled"` + CronExpression string `json:"cron_expression"` + RetentionDays int `json:"retention_days"` +} + +// CompactionSettingsResponse includes config plus schedule description. +type CompactionSettingsResponse struct { + CompactionSettingsConfig + ScheduleHuman string `json:"schedule_human"` +} + +func defaultCompactionSettingsConfig() CompactionSettingsConfig { + return CompactionSettingsConfig{ + Enabled: true, + CronExpression: "30 20 * * *", // Daily at 2:00 AM IST (20:30 UTC) + RetentionDays: 30, + } +} + +// GetCompactionSettings fetches ONLY the config (instant, single system_settings query). +func (s *Server) GetCompactionSettings(c echo.Context) error { + ctx := c.Request().Context() + cfg := defaultCompactionSettingsConfig() + + var data []byte + err := s.db.QueryRowContext(ctx, "SELECT data FROM system_settings WHERE name = 'event_compaction_settings'").Scan(&data) + if err == nil && len(data) > 0 { + _ = json.Unmarshal(data, &cfg) + } + + if cfg.RetentionDays <= 0 { + cfg.RetentionDays = 30 + } + if strings.TrimSpace(cfg.CronExpression) == "" { + cfg.CronExpression = "30 20 * * *" + } + + return c.JSON(http.StatusOK, CompactionSettingsResponse{ + CompactionSettingsConfig: cfg, + ScheduleHuman: describeCronSchedule(cfg.CronExpression), + }) +} + +// UpdateCompactionSettings updates event log compaction configuration in system_settings. +func (s *Server) UpdateCompactionSettings(c echo.Context) error { + var cfg CompactionSettingsConfig + if err := c.Bind(&cfg); err != nil { + return c.JSON(http.StatusBadRequest, map[string]string{"error": "Invalid request body"}) + } + + if cfg.RetentionDays <= 0 { + cfg.RetentionDays = 30 + } + if strings.TrimSpace(cfg.CronExpression) == "" { + cfg.CronExpression = "30 20 * * *" + } + + data, err := json.Marshal(cfg) + if err != nil { + return c.JSON(http.StatusInternalServerError, map[string]string{"error": "Failed to serialize settings"}) + } + + _, err = s.db.ExecContext(c.Request().Context(), ` + INSERT INTO system_settings (name, data) + VALUES ('event_compaction_settings', $1) + ON CONFLICT (name) DO UPDATE SET data = EXCLUDED.data, updated_at = CURRENT_TIMESTAMP + `, data) + if err != nil { + log.Error().Err(err).Msg("Failed to save event compaction settings") + return c.JSON(http.StatusInternalServerError, map[string]string{"error": "Failed to save compaction settings"}) + } + + if s.eventCompactionManager != nil { + s.eventCompactionManager.UpdateConfig(cfg.Enabled, cfg.CronExpression, cfg.RetentionDays) + } + + return c.JSON(http.StatusOK, map[string]string{"message": "Event compaction settings updated successfully"}) +} + +// RunCompactionNow triggers an immediate compaction pass. +func (s *Server) RunCompactionNow(c echo.Context) error { + if s.eventCompactionManager == nil { + return c.JSON(http.StatusBadRequest, map[string]string{"error": "Event compaction manager is not initialized"}) + } + + go s.eventCompactionManager.TriggerManualCycle() + + return c.JSON(http.StatusOK, map[string]string{"message": "Log compaction started in the background"}) +} + +// describeCronSchedule converts a cron expression to a plain English description. +func describeCronSchedule(expr string) string { + parts := strings.Fields(strings.TrimSpace(expr)) + if len(parts) != 5 { + return expr + } + + minutePart := parts[0] + hourPart := parts[1] + dom := parts[2] + month := parts[3] + dow := parts[4] + + // Simple daily: M H * * * + if dom == "*" && month == "*" && dow == "*" && !strings.Contains(hourPart, ",") && !strings.Contains(hourPart, "/") { + hour, errH := strconv.Atoi(hourPart) + minute, errM := strconv.Atoi(minutePart) + if errH == nil && errM == nil { + return fmt.Sprintf("Daily at %02d:%02d UTC", hour, minute) + } + } + + // Every N hours: 0 */N * * * + if dom == "*" && month == "*" && dow == "*" && strings.HasPrefix(hourPart, "*/") { + n := strings.TrimPrefix(hourPart, "*/") + return fmt.Sprintf("Every %s hours", n) + } + + // Multiple times a day: 0 H1,H2 * * * + if dom == "*" && month == "*" && dow == "*" && strings.Contains(hourPart, ",") { + return fmt.Sprintf("Multiple times daily (%s:00 UTC)", strings.ReplaceAll(hourPart, ",", ":00, ")) + } + + return expr +} diff --git a/internal/api/event_compaction.go b/internal/api/event_compaction.go new file mode 100644 index 00000000..024f163c --- /dev/null +++ b/internal/api/event_compaction.go @@ -0,0 +1,266 @@ +package api + +import ( + "context" + "database/sql" + "encoding/json" + "strings" + "sync" + "sync/atomic" + "time" + + "github.com/robfig/cron/v3" + "github.com/rs/zerolog/log" +) + +const defaultCompactionCronExpr = "30 20 * * *" // Daily at 2:00 AM IST (20:30 UTC) +const defaultRetentionDays = 30 + +// EventCompactionManager runs an automated background compaction job. +type EventCompactionManager struct { + db *sql.DB + mu sync.Mutex + enabled bool + cronExpr string + retentionDays int + cronRunner *cron.Cron + entryID cron.EntryID + ctx context.Context + cancel context.CancelFunc + running atomic.Int32 // 1 while a cycle is in progress, 0 otherwise +} + +// NewEventCompactionManager creates a new compaction manager. +func NewEventCompactionManager(db *sql.DB) *EventCompactionManager { + ctx, cancel := context.WithCancel(context.Background()) + m := &EventCompactionManager{ + db: db, + enabled: true, + cronExpr: defaultCompactionCronExpr, + retentionDays: defaultRetentionDays, + ctx: ctx, + cancel: cancel, + } + + m.loadSettingsFromDB() + return m +} + +func (m *EventCompactionManager) loadSettingsFromDB() { + var data []byte + err := m.db.QueryRowContext(m.ctx, "SELECT data FROM system_settings WHERE name = 'event_compaction_settings'").Scan(&data) + if err == nil && len(data) > 0 { + var cfg struct { + Enabled *bool `json:"enabled"` + CronExpression string `json:"cron_expression"` + RetentionDays int `json:"retention_days"` + } + if json.Unmarshal(data, &cfg) == nil { + if cfg.Enabled != nil { + m.enabled = *cfg.Enabled + } + if strings.TrimSpace(cfg.CronExpression) != "" { + m.cronExpr = cfg.CronExpression + } + if cfg.RetentionDays > 0 { + m.retentionDays = cfg.RetentionDays + } + } + } +} + +// Start launches the background cron runner. Falls back to the default cron +// expression if the configured one is invalid, so a bad DB setting never +// prevents the server from booting. +func (m *EventCompactionManager) Start() { + m.mu.Lock() + defer m.mu.Unlock() + + m.cronRunner = cron.New(cron.WithLocation(time.UTC)) + entryID, err := m.cronRunner.AddFunc(m.cronExpr, func() { + m.runCycle() + }) + if err != nil { + log.Warn().Str("bad_cron_expr", m.cronExpr).Err(err).Str("fallback", defaultCompactionCronExpr).Msg("[compaction] invalid cron expression in settings, falling back to default") + m.cronExpr = defaultCompactionCronExpr + entryID, err = m.cronRunner.AddFunc(m.cronExpr, func() { + m.runCycle() + }) + if err != nil { + // Default expression is hardcoded and always valid — this should never happen. + log.Error().Err(err).Msg("[compaction] failed to schedule even with default cron expression") + return + } + } + m.entryID = entryID + m.cronRunner.Start() + nextRun := m.cronRunner.Entry(m.entryID).Next + log.Info().Str("schedule", m.cronExpr).Bool("enabled", m.enabled).Int("retention_days", m.retentionDays).Time("next_run", nextRun).Msg("[compaction] manager started") +} + +// Stop gracefully shuts down the cron runner. +func (m *EventCompactionManager) Stop() { + m.mu.Lock() + defer m.mu.Unlock() + + log.Info().Msg("[compaction] manager stopping") + if m.cronRunner != nil { + ctx := m.cronRunner.Stop() + select { + case <-ctx.Done(): + case <-time.After(5 * time.Second): + } + } + m.cancel() +} + +// UpdateConfig dynamically reloads configuration without server restart. +func (m *EventCompactionManager) UpdateConfig(enabled bool, cronExpr string, retentionDays int) { + m.mu.Lock() + defer m.mu.Unlock() + + m.enabled = enabled + m.retentionDays = retentionDays + + if strings.TrimSpace(cronExpr) == "" { + cronExpr = defaultCompactionCronExpr + } + + if m.cronExpr != cronExpr || m.cronRunner == nil { + if m.cronRunner != nil { + m.cronRunner.Stop() + } + m.cronRunner = cron.New(cron.WithLocation(time.UTC)) + entryID, err := m.cronRunner.AddFunc(cronExpr, func() { + m.runCycle() + }) + if err == nil { + m.entryID = entryID + m.cronRunner.Start() + } else { + log.Error().Str("cron_expr", cronExpr).Err(err).Msg("[compaction] invalid cron expression") + } + m.cronExpr = cronExpr + } + + nextRun := time.Time{} + if m.cronRunner != nil { + nextRun = m.cronRunner.Entry(m.entryID).Next + } + log.Info().Bool("enabled", m.enabled).Str("schedule", m.cronExpr).Int("retention_days", m.retentionDays).Time("next_run", nextRun).Msg("[compaction] config updated") +} + +// TriggerManualCycle runs a cycle immediately. +func (m *EventCompactionManager) TriggerManualCycle() { + log.Info().Msg("[compaction] manual cycle triggered") + m.runCycle() +} + +// runCycle executes bulk compaction. Compaction runs in the backend process +// (single instance), so no distributed lock is required. Concurrent invocations +// (e.g. cron fires while a manual run is still in progress) are skipped. +func (m *EventCompactionManager) runCycle() { + // Atomically claim the running slot. If another cycle is already running, skip. + if !m.running.CompareAndSwap(0, 1) { + log.Warn().Msg("[compaction] cycle already in progress, skipping") + return + } + defer m.running.Store(0) + + m.mu.Lock() + enabled := m.enabled + retentionDays := m.retentionDays + m.mu.Unlock() + + if !enabled { + log.Info().Msg("[compaction] skipping cycle — compaction is currently disabled in settings") + return + } + + start := time.Now() + log.Info().Int("retention_days", retentionDays).Msg("[compaction] cycle start") + + compacted, errs := m.executeBulkCompaction(m.ctx, retentionDays) + + log.Info().Str("elapsed", time.Since(start).Round(time.Millisecond).String()).Int64("compacted_reviews", compacted).Int("errors", errs).Msg("[compaction] cycle done") +} + +// executeBulkCompaction executes compaction in safe batches: +// 1. BULK INSERT: Compaction summary markers for all eligible reviews. +// 2. BATCHED DELETE: Delete prunable log rows in chunks of 50,000 rows to prevent DB timeouts. +func (m *EventCompactionManager) executeBulkCompaction(ctx context.Context, retentionDays int) (compacted int64, errs int) { + // Step 1: Insert summary markers for uncompacted reviews + insertRes, err := m.db.ExecContext(ctx, ` + INSERT INTO public.review_events (review_id, org_id, ts, event_type, level, data) + SELECT + re.review_id, + re.org_id, + NOW(), + 'log', + 'info', + jsonb_build_object( + 'message', 'Review log events compacted', + 'compacted', true, + 'original_total_event_count', COUNT(*) + ) + FROM public.review_events re + WHERE re.ts < NOW() - ($1 * INTERVAL '1 day') + AND NOT EXISTS ( + SELECT 1 FROM public.review_events cx + WHERE cx.review_id = re.review_id AND cx.org_id = re.org_id + AND cx.event_type = 'log' AND (cx.data->>'compacted')::boolean = true + ) + GROUP BY re.review_id, re.org_id; + `, retentionDays) + if err != nil { + log.Error().Err(err).Msg("[compaction] insert summary markers failed") + return 0, 1 + } + + markersInserted, _ := insertRes.RowsAffected() + + // Step 2: Delete verbose log rows in batches of 50,000 rows to prevent query timeouts + batchSize := 50000 + var totalDeleted int64 = 0 + + deleteQuery := ` + DELETE FROM public.review_events + WHERE ctid IN ( + SELECT ctid FROM public.review_events + WHERE ts < NOW() - ($1 * INTERVAL '1 day') + AND event_type = 'log' + AND COALESCE(level, 'info') NOT IN ('error', 'warn') + AND (data->>'compacted')::boolean IS NOT TRUE + AND data->>'message' NOT ILIKE '%started%' + AND data->>'message' NOT ILIKE '%completed%' + AND data->>'message' NOT ILIKE '%posted%' + AND data->>'message' NOT ILIKE '%generated%' + AND data->>'message' NOT ILIKE '%CLI DIFF REVIEW STARTED%' + LIMIT $2 + ); + ` + + for { + select { + case <-ctx.Done(): + log.Info().Int64("rows_deleted", totalDeleted).Msg("[compaction] deletion context cancelled") + return markersInserted, 0 + default: + } + + res, err := m.db.ExecContext(ctx, deleteQuery, retentionDays, batchSize) + if err != nil { + log.Error().Err(err).Msg("[compaction] batch delete error") + break + } + + rows, _ := res.RowsAffected() + totalDeleted += rows + if rows == 0 { + break + } + } + + log.Info().Int64("reviews_compacted", markersInserted).Int64("rows_deleted", totalDeleted).Msg("[compaction] bulk cycle complete") + return markersInserted, 0 +} diff --git a/internal/api/event_compaction_test.go b/internal/api/event_compaction_test.go new file mode 100644 index 00000000..0d03a0c6 --- /dev/null +++ b/internal/api/event_compaction_test.go @@ -0,0 +1,76 @@ +package api + +import ( + "context" + "testing" +) + +func TestDefaultCompactionSettingsConfig(t *testing.T) { + cfg := defaultCompactionSettingsConfig() + + if !cfg.Enabled { + t.Errorf("expected default Enabled to be true, got false") + } + + if cfg.RetentionDays != 30 { + t.Errorf("expected default RetentionDays to be 30, got %d", cfg.RetentionDays) + } + + if cfg.CronExpression != "30 20 * * *" { + t.Errorf("expected default CronExpression to be '30 20 * * *', got %q", cfg.CronExpression) + } +} + +func TestDescribeCronSchedule(t *testing.T) { + tests := []struct { + expr string + expected string + }{ + { + expr: "0 2 * * *", + expected: "Daily at 02:00 UTC", + }, + { + expr: "0 */6 * * *", + expected: "Every 6 hours", + }, + { + expr: "0 0,12 * * *", + expected: "Multiple times daily (0:00, 12:00 UTC)", + }, + { + expr: "custom_expr", + expected: "custom_expr", + }, + } + + for _, tt := range tests { + result := describeCronSchedule(tt.expr) + if result != tt.expected { + t.Errorf("describeCronSchedule(%q) = %q, want %q", tt.expr, result, tt.expected) + } + } +} + +func TestEventCompactionManager_ConfigUpdate(t *testing.T) { + m := &EventCompactionManager{ + enabled: true, + cronExpr: "0 2 * * *", + retentionDays: 30, + ctx: context.Background(), + } + + m.UpdateConfig(false, "0 0 * * *", 60) + + if m.enabled { + t.Errorf("expected enabled to be false after UpdateConfig") + } + + if m.retentionDays != 60 { + t.Errorf("expected retentionDays to be 60, got %d", m.retentionDays) + } + + if m.cronExpr != "0 0 * * *" { + t.Errorf("expected cronExpr to be '0 0 * * *', got %q", m.cronExpr) + } +} diff --git a/internal/api/review_events_endpoints.go b/internal/api/review_events_endpoints.go index 972c3802..ab530476 100644 --- a/internal/api/review_events_endpoints.go +++ b/internal/api/review_events_endpoints.go @@ -2,6 +2,7 @@ package api import ( "database/sql" + "encoding/json" "net/http" "strconv" "time" @@ -110,13 +111,37 @@ func (h *ReviewEventsHandler) GetReviewEvents(c echo.Context) error { events = make([]*ReviewEvent, 0) } + // If the review has been compacted, use the original total event count stored + // in the compaction marker so the UI shows the correct historical count + // (e.g. "Events: 248") even though we now only hold ~30 rows. + eventCount := len(events) + var rawMarker []byte + markerErr := h.service.repo.DB().QueryRowContext( + c.Request().Context(), + `SELECT data FROM public.review_events + WHERE review_id = $1 AND org_id = $2 + AND event_type = 'log' + AND (data->>'compacted')::boolean = true + LIMIT 1`, + reviewID, orgID, + ).Scan(&rawMarker) + if markerErr == nil && rawMarker != nil { + var marker struct { + OriginalTotalEventCount *int `json:"original_total_event_count"` + } + if err := json.Unmarshal(rawMarker, &marker); err == nil && marker.OriginalTotalEventCount != nil { + eventCount = *marker.OriginalTotalEventCount + } + } + // Return events in the standard envelope format response := map[string]interface{}{ "events": events, "meta": map[string]interface{}{ - "reviewId": reviewID, - "count": len(events), - "limit": limit, + "reviewId": reviewID, + "count": eventCount, + "limit": limit, + "compacted": markerErr == nil && rawMarker != nil, }, } diff --git a/internal/api/server.go b/internal/api/server.go index 76cb5e6f..11530f3b 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -133,7 +133,8 @@ type Server struct { port int db *sql.DB jobQueue *jobqueue.JobQueue - dashboardManager *DashboardManager + dashboardManager *DashboardManager + eventCompactionManager *EventCompactionManager autoWebhookInstaller *AutoWebhookInstaller versionInfo *VersionInfo deploymentConfig *DeploymentConfig @@ -289,6 +290,9 @@ func appContext(port int, versionInfo *VersionInfo) (*Server, error) { dashboardCacheStore := dashboard.NewCacheStore(db) dashboardManager := NewDashboardManager(db, schedulerLockStore, dashboardCacheStore) + // Initialize event compaction manager — one goroutine compacting review_events logs daily. + eventCompactionManager := NewEventCompactionManager(db) + // Initialize auto webhook installer autoWebhookInstaller := NewAutoWebhookInstaller(db, nil, jq) // server will be set later @@ -331,7 +335,8 @@ func appContext(port int, versionInfo *VersionInfo) (*Server, error) { port: port, db: db, jobQueue: jq, - dashboardManager: dashboardManager, + dashboardManager: dashboardManager, + eventCompactionManager: eventCompactionManager, autoWebhookInstaller: autoWebhookInstaller, versionInfo: versionInfo, deploymentConfig: deploymentConfig, @@ -1096,6 +1101,11 @@ func (s *Server) setupRoutes() { adminGroup.PUT("/settings/storage", s.UpdateStorageSettings) adminGroup.POST("/settings/storage/test", s.TestStorageSettings) + // Super admin log compaction settings endpoints + adminGroup.GET("/settings/compaction", s.GetCompactionSettings) + adminGroup.PUT("/settings/compaction", s.UpdateCompactionSettings) + adminGroup.POST("/settings/compaction/run", s.RunCompactionNow) + // Organization management endpoints // User organization access (get their orgs) - needs permission context to detect super admin protectedOrgsGroup := protected.Group("") @@ -1664,6 +1674,11 @@ func (s *Server) Start() error { s.dashboardManager.Start() fmt.Println("Dashboard manager started") + // Start event compaction manager (daily log compaction for review_events > 30 days) + s.eventCompactionManager.Start() + fmt.Println("Event compaction manager started") + + // Start server in a goroutine go func() { if err := s.echo.Start(bindAddress); err != nil && err != http.ErrServerClosed { @@ -1769,6 +1784,12 @@ func (s *Server) Start() error { fmt.Println("Dashboard manager stopped") } + // Stop event compaction manager + if s.eventCompactionManager != nil { + s.eventCompactionManager.Stop() + fmt.Println("Event compaction manager stopped") + } + // Close database connection if s.db != nil { s.db.Close() diff --git a/network/network_status.md b/network/network_status.md index 97622b1a..7ba88686 100644 --- a/network/network_status.md +++ b/network/network_status.md @@ -2,18 +2,18 @@ Latest milestone batch note (MF-051, MF-059, MF-073, MF-074, MF-076, MF-083, MF-LOC-001, MF-LOC-002, MF-LOC-003, MF-LOC-004, MF-LOC-005, MF-LOC-006, MF-LOC-007, MF-LOC-008, MF-PRORATION-001, MF-PRORATION-002, MF-PRORATION-003, MF-ATTRIB-001, MF-ATTRIB-002, MF-PORTFOLIO-001, MF-NOTIFY-001, MF-NOTIFY-002, MF-DASHBOARD-LOG-001, MF-EXPIRY-001, MF-UPI-UPGRADE-001, MF-UPI-UPGRADE-002, MF-UPI-UPGRADE-003, MF-UPI-UPGRADE-004, MF-TRIAL-CANCEL-001, MF-CANCEL-VERIFY-001, MF-CANCEL-PROJECTION-001, MF-STATUS-LABEL-001, MF-BILLING-PRICE-001, MF-AI-HELPER-001): added provider-backed monthly plan price exposure in billing status, surfaced current paid subscription currency for settings flows, rejected unsupported paid cross-currency upgrade previews instead of letting quantity-only subscription updates imply currency switching, and exposed Leader/Helper review AI settings plus per-stage review accounting breakdowns. -| Operation | Status | Evidence | -| --- | --- | --- | -| payment.CreateSubscriptionAddon | added | [CreateSubscriptionAddon](../internal/license/payment/payment.go#L309) | -| payment.CreateOrder | added | [CreateOrder](../internal/license/payment/payment.go#L359) | -| payment.CreateSubscriptionAt | added | [CreateSubscriptionAt](../internal/license/payment/subscription.go#L78) | -| payment.CancelScheduledChangesByID | added | [CancelScheduledChangesByID](../internal/license/payment/subscription.go#L295) | -| api.CreateSubscription | updated | [CreateSubscription](../internal/api/subscriptions_handler.go#L141) | -| api.CancelSubscription | updated | [CancelSubscription](../internal/api/subscriptions_handler.go#L299) | -| api.GetBillingStatus | updated | [GetBillingStatus](../internal/api/billing_actions_handler.go#L1311) | -| api.checkGitHubParentCommentAuthor | updated | [checkGitHubParentCommentAuthor](../internal/api/unified_processor_v2.go#L707) | +| Operation | Status | Evidence | +| ------------------------------------- | ------- | --------------------------------------------------------------------------------- | +| payment.CreateSubscriptionAddon | added | [CreateSubscriptionAddon](../internal/license/payment/payment.go#L309) | +| payment.CreateOrder | added | [CreateOrder](../internal/license/payment/payment.go#L359) | +| payment.CreateSubscriptionAt | added | [CreateSubscriptionAt](../internal/license/payment/subscription.go#L78) | +| payment.CancelScheduledChangesByID | added | [CancelScheduledChangesByID](../internal/license/payment/subscription.go#L295) | +| api.CreateSubscription | updated | [CreateSubscription](../internal/api/subscriptions_handler.go#L141) | +| api.CancelSubscription | updated | [CancelSubscription](../internal/api/subscriptions_handler.go#L299) | +| api.GetBillingStatus | updated | [GetBillingStatus](../internal/api/billing_actions_handler.go#L1311) | +| api.checkGitHubParentCommentAuthor | updated | [checkGitHubParentCommentAuthor](../internal/api/unified_processor_v2.go#L707) | | api.checkBitbucketParentCommentAuthor | updated | [checkBitbucketParentCommentAuthor](../internal/api/unified_processor_v2.go#L776) | -| api.PreviewUpgrade | updated | [PreviewUpgrade](../internal/api/billing_actions_handler.go#L457) | +| api.PreviewUpgrade | updated | [PreviewUpgrade](../internal/api/billing_actions_handler.go#L457) | | api.GetCurrentSubscription | updated | [GetCurrentSubscription](../internal/api/subscriptions_handler.go#L620) | | api.ListUserSubscriptions | updated | [ListUserSubscriptions](../internal/api/subscriptions_handler.go#L773) | diff --git a/storage/storage_status.md b/storage/storage_status.md index 0fb4bd0c..e9dff22c 100644 --- a/storage/storage_status.md +++ b/storage/storage_status.md @@ -2,174 +2,175 @@ Latest milestone batch note (MF-LOC-007, MF-LOC-008, MF-PRORATION-003, MF-ATTRIB-001, MF-ATTRIB-002, MF-PORTFOLIO-001, MF-NOTIFY-001, MF-DASHBOARD-LOG-001, MF-EXPIRY-001, MF-CANCEL-VERIFY-001, MF-CANCEL-PROJECTION-001, MF-AI-HELPER-001): added member-attribution storage groundwork, member-usage rollups, payment-attempt lookup for customer state, superadmin billing portfolio storage views, billing notification outbox persistence, dashboard scheduler leader-lock storage operations, automatic paid-plan expiry reconciliation persistence, quota policy + batch settlement + operation aggregate storage operations, a shared transaction-level org free-plan projection helper used by terminal subscription webhook flows, and org-scoped Helper review AI settings plus role-aware AI connector storage queries. -| Operation | Status | Evidence | -| --- | --- | --- | -| payment.NewSubscriptionStore | moved | [NewSubscriptionStore](payment/subscription_store.go#L26) | -| payment.CreateTeamSubscriptionRecord | moved | [CreateTeamSubscriptionRecord](payment/subscription_store.go#L45) | -| payment.UpdateSubscriptionQuantityRecord | moved | [UpdateSubscriptionQuantityRecord](payment/subscription_store.go#L110) | -| payment.SyncOrgBillingStateToFreeTx | added | [SyncOrgBillingStateToFreeTx](payment/subscription_store.go#L179) | -| payment.CancelSubscriptionRecord | moved | [CancelSubscriptionRecord](payment/subscription_store.go#L241) | -| payment.ReconcileExpiredPendingCancellations | added | [ReconcileExpiredPendingCancellations](payment/subscription_store.go#L336) | -| payment.ReconcileExpiredPendingCancellationForOrg | added | [ReconcileExpiredPendingCancellationForOrg](payment/subscription_store.go#L340) | -| payment.DowngradeExpiredRoleForUserOrg | added | [DowngradeExpiredRoleForUserOrg](payment/subscription_store.go#L353) | -| payment.KeepPlanRecord | added | [KeepPlanRecord](payment/subscription_store.go#L613) | -| payment.GetSubscriptionDetailsRow | moved | [GetSubscriptionDetailsRow](payment/subscription_store.go#L778) | -| payment.AssignLicense | moved | [AssignLicense](payment/subscription_store.go#L809) | -| payment.RepointOrgActiveSubscription | added | [RepointOrgActiveSubscription](payment/subscription_store.go#L922) | -| payment.RevokeLicense | moved | [RevokeLicense](payment/subscription_store.go#L970) | -| payment.GetUserIDByEmail | moved | [GetUserIDByEmail](payment/subscription_store.go#L1050) | -| payment.CreateShadowUser | moved | [CreateShadowUser](payment/subscription_store.go#L1062) | -| payment.CreateSelfHostedSubscriptionRecord | moved | [CreateSelfHostedSubscriptionRecord](payment/subscription_store.go#L1090) | -| payment.GetSelfHostedConfirmationSeed | moved | [GetSelfHostedConfirmationSeed](payment/subscription_store.go#L1155) | -| payment.PersistSelfHostedFallback | moved | [PersistSelfHostedFallback](payment/subscription_store.go#L1192) | -| payment.PersistSelfHostedJWT | moved | [PersistSelfHostedJWT](payment/subscription_store.go#L1255) | -| jobqueue.NewWebhookStore | moved | [NewWebhookStore](jobqueue/webhook_store.go#L24) | -| jobqueue.GetWebhookPublicEndpoint | moved | [GetWebhookPublicEndpoint](jobqueue/webhook_store.go#L33) | -| jobqueue.GetWebhookRegistryID | moved | [GetWebhookRegistryID](jobqueue/webhook_store.go#L52) | -| jobqueue.InsertWebhookRegistry | moved | [InsertWebhookRegistry](jobqueue/webhook_store.go#L84) | -| jobqueue.UpdateWebhookRegistryByID | moved | [UpdateWebhookRegistryByID](jobqueue/webhook_store.go#L122) | -| jobqueue.GetConnectorMetadata | moved | [GetConnectorMetadata](jobqueue/webhook_store.go#L145) | -| license.GetLicenseState | moved | [GetLicenseState](license/license_state_store.go#L41) | -| license.UpsertLicenseState | moved | [UpsertLicenseState](license/license_state_store.go#L56) | -| license.UpdateValidationResult | moved | [UpdateValidationResult](license/license_state_store.go#L88) | -| license.DeleteLicenseState | moved | [DeleteLicenseState](license/license_state_store.go#L119) | -| providersgitlab.NewAuthTokenStore | moved | [NewAuthTokenStore](providers/gitlab/auth_token_store.go#L44) | -| providersgitlab.UpsertGitLabIntegrationToken | moved | [UpsertGitLabIntegrationToken](providers/gitlab/auth_token_store.go#L48) | -| providersgitlab.GetGitLabTokenRecord | moved | [GetGitLabTokenRecord](providers/gitlab/auth_token_store.go#L109) | -| providersgitlab.GetProviderAppID | moved | [GetProviderAppID](providers/gitlab/auth_token_store.go#L122) | -| providersgitlab.UpdateRefreshedToken | moved | [UpdateRefreshedToken](providers/gitlab/auth_token_store.go#L133) | -| providersbitbucket.NewIntegrationTokenStore | moved | [NewIntegrationTokenStore](providers/bitbucket/integration_token_store.go#L18) | -| providersbitbucket.GetTokenByRepoFullName | moved | [GetTokenByRepoFullName](providers/bitbucket/integration_token_store.go#L22) | -| providersbitbucket.GetLatestBitbucketToken | moved | [GetLatestBitbucketToken](providers/bitbucket/integration_token_store.go#L41) | -| providersgithub.NewTokenStore | moved | [NewTokenStore](providers/github/token_store.go#L17) | -| providersgithub.GetLatestGitHubToken | moved | [GetLatestGitHubToken](providers/github/token_store.go#L21) | -| reviews.NewReviewStore | moved | [NewReviewStore](reviews/review_store.go#L9) | -| reviews.QueryRow | moved | [QueryRow](reviews/review_store.go#L13) | -| reviews.Exec | moved | [Exec](reviews/review_store.go#L17) | -| reviews.Query | moved | [Query](reviews/review_store.go#L21) | -| reviews.NewTaxonomyReportStore | added | [NewTaxonomyReportStore](reviews/taxonomy_report_store.go#L17) | -| reviews.buildWhereClause | added | [buildWhereClause](reviews/taxonomy_report_store.go#L108) | -| reviews.GetSummary | added | [GetSummary](reviews/taxonomy_report_store.go#L230) | -| reviews.GetDistribution | added | [GetDistribution](reviews/taxonomy_report_store.go#L273) | -| reviews.GetTrend | added | [GetTrend](reviews/taxonomy_report_store.go#L326) | -| reviews.GetBreakdown | added | [GetBreakdown](reviews/taxonomy_report_store.go#L365) | -| reviews.ListFindings | updated | [ListFindings](reviews/taxonomy_report_store.go#L430) | -| reviews.GetCategorySubcategoryRelations | added | [GetCategorySubcategoryRelations](reviews/taxonomy_report_store.go#L552) | -| aiconnectors.NewConnectorStore | moved | [NewConnectorStore](aiconnectors/connector_store.go#L18) | -| aiconnectors.QueryRowContext | moved | [QueryRowContext](aiconnectors/connector_store.go#L22) | -| aiconnectors.QueryContext | moved | [QueryContext](aiconnectors/connector_store.go#L26) | -| aiconnectors.ExecContext | moved | [ExecContext](aiconnectors/connector_store.go#L30) | -| aiconnectors.UpdateDisplayOrders | moved | [UpdateDisplayOrders](aiconnectors/connector_store.go#L38) | -| aiconnectors.NewReviewAISettingsStore | added | [NewReviewAISettingsStore](aiconnectors/review_ai_settings_store.go#L28) | -| aiconnectors.NormalizeConnectorRole | added | [NormalizeConnectorRole](aiconnectors/review_ai_settings_store.go#L32) | -| aiconnectors.NormalizeHelperMode | added | [NormalizeHelperMode](aiconnectors/review_ai_settings_store.go#L43) | -| aiconnectors.GetByOrgID | added | [GetByOrgID](aiconnectors/review_ai_settings_store.go#L54) | -| aiconnectors.Upsert | added | [Upsert](aiconnectors/review_ai_settings_store.go#L86) | -| aiconnectors.GetConnectorsByRole | added | [GetConnectorsByRole](../internal/aiconnectors/storage.go#L230) | -| aiconnectors.GetMaxDisplayOrderByRole | added | [GetMaxDisplayOrderByRole](../internal/aiconnectors/storage.go#L561) | -| users.NewUserStore | moved | [NewUserStore](users/user_store.go#L9) | -| users.QueryRow | moved | [QueryRow](users/user_store.go#L13) | -| users.Query | moved | [Query](users/user_store.go#L17) | -| users.TxQueryRow | moved | [TxQueryRow](users/user_store.go#L21) | -| users.TxExec | moved | [TxExec](users/user_store.go#L25) | -| users.WithTx | moved | [WithTx](users/user_store.go#L29) | -| users.NewProfileStore | moved | [NewProfileStore](users/profile_store.go#L46) | -| users.GetUserProfile | moved | [GetUserProfile](users/profile_store.go#L50) | -| users.UpdateUserProfile | moved | [UpdateUserProfile](users/profile_store.go#L77) | -| users.GetCurrentPasswordHash | moved | [GetCurrentPasswordHash](users/profile_store.go#L114) | -| users.UpdatePasswordIfCurrentHash | moved | [UpdatePasswordIfCurrentHash](users/profile_store.go#L123) | -| users.ListUserOrganizations | moved | [ListUserOrganizations](users/profile_store.go#L139) | -| learnings.NewLearningsStore | moved | [NewLearningsStore](learnings/learnings_store.go#L67) | -| learnings.InsertLearning | moved | [InsertLearning](learnings/learnings_store.go#L71) | -| learnings.UpdateLearning | moved | [UpdateLearning](learnings/learnings_store.go#L89) | -| learnings.QueryLearningByID | moved | [QueryLearningByID](learnings/learnings_store.go#L107) | -| learnings.QueryLearningByShortID | moved | [QueryLearningByShortID](learnings/learnings_store.go#L114) | -| learnings.ListByOrgWithPagination | moved | [ListByOrgWithPagination](learnings/learnings_store.go#L121) | -| learnings.CountByOrg | moved | [CountByOrg](learnings/learnings_store.go#L154) | -| learnings.InsertLearningEvent | moved | [InsertLearningEvent](learnings/learnings_store.go#L174) | -| providersgitea.NewTokenStore | moved | [NewTokenStore](providers/gitea/token_store.go#L18) | -| providersgitea.ListRecentGiteaIntegrationTokens | moved | [ListRecentGiteaIntegrationTokens](providers/gitea/token_store.go#L22) | -| providersgitea.GetGiteaIntegrationTokenByID | moved | [GetGiteaIntegrationTokenByID](providers/gitea/token_store.go#L51) | -| providersgitea.GetLatestWebhookSecret | moved | [GetLatestWebhookSecret](providers/gitea/token_store.go#L64) | -| core.NewFileOpsStore | moved | [NewFileOpsStore](core/file_ops.go#L7) | -| core.ReadFile | moved | [ReadFile](core/file_ops.go#L11) | -| core.NewSchedulerLockStore | added | [NewSchedulerLockStore](core/scheduler_lock_store.go#L19) | -| core.TryAcquireDashboardRefreshLeaderLock | added | [TryAcquireDashboardRefreshLeaderLock](core/scheduler_lock_store.go#L23) | -| core.ReleaseDashboardRefreshLeaderLock | added | [ReleaseDashboardRefreshLeaderLock](core/scheduler_lock_store.go#L58) | -| license.NewPlanCatalogFileStore | moved | [NewPlanCatalogFileStore](license/plan_catalog_file_store.go#L14) | -| license.ReadPlanCatalogFile | moved | [ReadPlanCatalogFile](license/plan_catalog_file_store.go#L18) | -| license.NewQuotaStore | added | [NewQuotaStore](license/quota_store.go#L106) | -| license.ResolvePolicy | added | [ResolvePolicy](license/quota_store.go#L110) | -| license.UpsertBatchSettlement | added | [UpsertBatchSettlement](license/quota_store.go#L165) | -| license.BuildAggregateFromBatches | added | [BuildAggregateFromBatches](license/quota_store.go#L260) | -| license.UpsertOperationAggregate | added | [UpsertOperationAggregate](license/quota_store.go#L311) | -| license.NewLOCAccountingStore | moved | [NewLOCAccountingStore](license/loc_accounting_store.go#L52) | -| license.AccountSuccess | moved | [AccountSuccess](license/loc_accounting_store.go#L56) | -| license.CheckQuotaPreflight | moved | [CheckQuotaPreflight](license/loc_accounting_store.go#L212) | -| license.emitThresholdLifecycleEventsTx | moved | [emitThresholdLifecycleEventsTx](license/loc_accounting_store.go#L358) | -| license.emitLifecycleEventTx | moved | [emitLifecycleEventTx](license/loc_accounting_store.go#L502) | -| license.NewActorLookupStore | added | [NewActorLookupStore](license/actor_lookup_store.go#L15) | -| license.ResolveOrgMemberUserIDByEmail | added | [ResolveOrgMemberUserIDByEmail](license/actor_lookup_store.go#L19) | -| license.NewTrialEligibilityStore | added | [NewTrialEligibilityStore](license/trial_eligibility_store.go#L58) | -| license.NormalizeTrialEligibilityEmail | added | [NormalizeTrialEligibilityEmail](license/trial_eligibility_store.go#L62) | -| license.GetTrialEligibilityByEmail | added | [GetTrialEligibilityByEmail](license/trial_eligibility_store.go#L70) | -| license.ReserveFirstPurchaseTrial | added | [ReserveFirstPurchaseTrial](license/trial_eligibility_store.go#L107) | -| license.ConsumeReservedTrial | added | [ConsumeReservedTrial](license/trial_eligibility_store.go#L171) | -| license.ConsumeReservedTrialTx | added | [ConsumeReservedTrialTx](license/trial_eligibility_store.go#L194) | -| license.ReleaseTrialReservation | added | [ReleaseTrialReservation](license/trial_eligibility_store.go#L261) | -| license.NewAdminBillingPortfolioStore | added | [NewAdminBillingPortfolioStore](license/admin_billing_portfolio_store.go#L38) | -| license.GetSummary | added | [GetSummary](license/admin_billing_portfolio_store.go#L42) | -| license.ListOrganizations | added | [ListOrganizations](license/admin_billing_portfolio_store.go#L88) | -| license.OrganizationExists | added | [OrganizationExists](license/admin_billing_portfolio_store.go#L179) | -| license.NewReviewAccountingStore | moved | [NewReviewAccountingStore](license/review_accounting_store.go#L40) | -| license.GetReviewAccountingTotals | moved | [GetReviewAccountingTotals](license/review_accounting_store.go#L44) | -| license.GetLatestReviewAccountingOperation | moved | [GetLatestReviewAccountingOperation](license/review_accounting_store.go#L109) | -| license.NewOrgUsageStore | moved | [NewOrgUsageStore](license/org_usage_store.go#L54) | -| license.GetCurrentPeriodSummary | moved | [GetCurrentPeriodSummary](license/org_usage_store.go#L58) | -| license.ListCurrentPeriodOperations | moved | [ListCurrentPeriodOperations](license/org_usage_store.go#L115) | -| license.ListCurrentPeriodMemberUsage | added | [ListCurrentPeriodMemberUsage](license/org_usage_store.go#L197) | -| license.GetCurrentPeriodUsageForActor | added | [GetCurrentPeriodUsageForActor](license/org_usage_store.go#L272) | -| payment.GetLatestCapturedPaymentMethodBySubscriptionID | added | [GetLatestCapturedPaymentMethodBySubscriptionID](payment/subscription_store.go#L730) | -| payment.ListSubscriptionsByOrgID | moved | [ListSubscriptionsByOrgID](payment/subscription_store.go#L750) | -| payment.NewBillingNotificationOutboxStore | added | [NewBillingNotificationOutboxStore](payment/billing_notification_outbox_store.go#L45) | -| payment.Enqueue | added | [Enqueue](payment/billing_notification_outbox_store.go#L49) | -| payment.GetUserEmailByID | added | [GetUserEmailByID](payment/billing_notification_outbox_store.go#L104) | -| payment.ClaimDispatchBatch | added | [ClaimDispatchBatch](payment/billing_notification_outbox_store.go#L123) | -| payment.MarkSent | added | [MarkSent](payment/billing_notification_outbox_store.go#L201) | -| payment.MarkFailed | added | [MarkFailed](payment/billing_notification_outbox_store.go#L221) | -| payment.MarkCancelled | added | [MarkCancelled](payment/billing_notification_outbox_store.go#L245) | -| payment.NewUpgradePaymentAttemptStore | added | [NewUpgradePaymentAttemptStore](payment/upgrade_payment_attempt_store.go#L88) | -| payment.CreateUpgradePaymentAttempt | added | [CreateUpgradePaymentAttempt](payment/upgrade_payment_attempt_store.go#L97) | -| payment.GetReusablePreparedAttempt | added | [GetReusablePreparedAttempt](payment/upgrade_payment_attempt_store.go#L132) | -| payment.GetAttemptByOrgPreviewAndOrder | added | [GetAttemptByOrgPreviewAndOrder](payment/upgrade_payment_attempt_store.go#L159) | -| payment.GetAttemptByOrderID | added | [GetAttemptByOrderID](payment/upgrade_payment_attempt_store.go#L185) | -| payment.GetLatestAttemptByUpgradeRequestID | added | [GetLatestAttemptByUpgradeRequestID](payment/upgrade_payment_attempt_store.go#L209) | -| payment.MarkPaymentCapturedByOrderID | added | [MarkPaymentCapturedByOrderID](payment/upgrade_payment_attempt_store.go#L266) | -| payment.MarkPaymentFailedByOrderID | added | [MarkPaymentFailedByOrderID](payment/upgrade_payment_attempt_store.go#L288) | -| payment.ReserveExecute | added | [ReserveExecute](payment/upgrade_payment_attempt_store.go#L320) | -| payment.MarkExecuteApplied | added | [MarkExecuteApplied](payment/upgrade_payment_attempt_store.go#L411) | -| payment.NewUpgradeReplacementCutoverStore | added | [NewUpgradeReplacementCutoverStore](payment/upgrade_replacement_cutover_store.go#L69) | -| payment.CreateOrGetPending | added | [CreateOrGetPending](payment/upgrade_replacement_cutover_store.go#L73) | -| payment.GetByUpgradeRequestID | added | [GetByUpgradeRequestID](payment/upgrade_replacement_cutover_store.go#L113) | -| payment.MarkReplacementProvisioned | added | [MarkReplacementProvisioned](payment/upgrade_replacement_cutover_store.go#L138) | -| payment.MarkOldCancellationScheduled | added | [MarkOldCancellationScheduled](payment/upgrade_replacement_cutover_store.go#L175) | -| payment.MarkRetryPending | added | [MarkRetryPending](payment/upgrade_replacement_cutover_store.go#L209) | -| payment.MarkManualReviewRequired | added | [MarkManualReviewRequired](payment/upgrade_replacement_cutover_store.go#L245) | -| payment.MarkCompleted | added | [MarkCompleted](payment/upgrade_replacement_cutover_store.go#L279) | -| license.NewPlanChangeStore | moved | [NewPlanChangeStore](license/plan_change_store.go#L30) | -| license.EnsureOrgBillingState | moved | [EnsureOrgBillingState](license/plan_change_store.go#L34) | -| license.GetOrgBillingState | moved | [GetOrgBillingState](license/plan_change_store.go#L57) | -| license.ApplyImmediatePlanUpgrade | moved | [ApplyImmediatePlanUpgrade](license/plan_change_store.go#L95) | -| license.ScheduleDowngrade | moved | [ScheduleDowngrade](license/plan_change_store.go#L126) | -| license.ScheduleUpgradeWithCurrentCycleGrant | added | [ScheduleUpgradeWithCurrentCycleGrant](license/plan_change_store.go#L153) | -| license.CancelScheduledDowngrade | moved | [CancelScheduledDowngrade](license/plan_change_store.go#L187) | -| license.ListDueScheduledDowngrades | moved | [ListDueScheduledDowngrades](license/plan_change_store.go#L221) | -| license.ListDueScheduledPlanChanges | added | [ListDueScheduledPlanChanges](license/plan_change_store.go#L225) | -| license.ApplyScheduledDowngrade | moved | [ApplyScheduledDowngrade](license/plan_change_store.go#L254) | -| license.ApplyScheduledPlanChange | added | [ApplyScheduledPlanChange](license/plan_change_store.go#L258) | -| license.insertLifecycleEventTx | moved | [insertLifecycleEventTx](license/plan_change_store.go#L299) | -| analytics.NewAdHocStore | added | [NewAdHocStore](analytics/adhoc_store.go#L50) | -| analytics.WithStatementTimeout | added | [WithStatementTimeout](analytics/adhoc_store.go#L55) | -| analytics.Count | added | [Count](analytics/adhoc_store.go#L67) | -| analytics.Query | added | [Query](analytics/adhoc_store.go#L89) | -| analytics.Ping | added | [Ping](analytics/adhoc_store.go#L186) | +| Operation | Status | Evidence | +| ------------------------------------------------------ | ------- | ------------------------------------------------------------------------------------- | +| payment.NewSubscriptionStore | moved | [NewSubscriptionStore](payment/subscription_store.go#L26) | +| payment.CreateTeamSubscriptionRecord | moved | [CreateTeamSubscriptionRecord](payment/subscription_store.go#L45) | +| payment.UpdateSubscriptionQuantityRecord | moved | [UpdateSubscriptionQuantityRecord](payment/subscription_store.go#L110) | +| payment.SyncOrgBillingStateToFreeTx | added | [SyncOrgBillingStateToFreeTx](payment/subscription_store.go#L179) | +| payment.CancelSubscriptionRecord | moved | [CancelSubscriptionRecord](payment/subscription_store.go#L241) | +| payment.ReconcileExpiredPendingCancellations | added | [ReconcileExpiredPendingCancellations](payment/subscription_store.go#L336) | +| payment.ReconcileExpiredPendingCancellationForOrg | added | [ReconcileExpiredPendingCancellationForOrg](payment/subscription_store.go#L340) | +| payment.DowngradeExpiredRoleForUserOrg | added | [DowngradeExpiredRoleForUserOrg](payment/subscription_store.go#L353) | +| payment.KeepPlanRecord | added | [KeepPlanRecord](payment/subscription_store.go#L613) | +| payment.GetSubscriptionDetailsRow | moved | [GetSubscriptionDetailsRow](payment/subscription_store.go#L778) | +| payment.AssignLicense | moved | [AssignLicense](payment/subscription_store.go#L809) | +| payment.RepointOrgActiveSubscription | added | [RepointOrgActiveSubscription](payment/subscription_store.go#L922) | +| payment.RevokeLicense | moved | [RevokeLicense](payment/subscription_store.go#L970) | +| payment.GetUserIDByEmail | moved | [GetUserIDByEmail](payment/subscription_store.go#L1050) | +| payment.CreateShadowUser | moved | [CreateShadowUser](payment/subscription_store.go#L1062) | +| payment.CreateSelfHostedSubscriptionRecord | moved | [CreateSelfHostedSubscriptionRecord](payment/subscription_store.go#L1090) | +| payment.GetSelfHostedConfirmationSeed | moved | [GetSelfHostedConfirmationSeed](payment/subscription_store.go#L1155) | +| payment.PersistSelfHostedFallback | moved | [PersistSelfHostedFallback](payment/subscription_store.go#L1192) | +| payment.PersistSelfHostedJWT | moved | [PersistSelfHostedJWT](payment/subscription_store.go#L1255) | +| jobqueue.NewWebhookStore | moved | [NewWebhookStore](jobqueue/webhook_store.go#L24) | +| jobqueue.GetWebhookPublicEndpoint | moved | [GetWebhookPublicEndpoint](jobqueue/webhook_store.go#L33) | +| jobqueue.GetWebhookRegistryID | moved | [GetWebhookRegistryID](jobqueue/webhook_store.go#L52) | +| jobqueue.InsertWebhookRegistry | moved | [InsertWebhookRegistry](jobqueue/webhook_store.go#L84) | +| jobqueue.UpdateWebhookRegistryByID | moved | [UpdateWebhookRegistryByID](jobqueue/webhook_store.go#L122) | +| jobqueue.GetConnectorMetadata | moved | [GetConnectorMetadata](jobqueue/webhook_store.go#L145) | +| license.GetLicenseState | moved | [GetLicenseState](license/license_state_store.go#L41) | +| license.UpsertLicenseState | moved | [UpsertLicenseState](license/license_state_store.go#L56) | +| license.UpdateValidationResult | moved | [UpdateValidationResult](license/license_state_store.go#L88) | +| license.DeleteLicenseState | moved | [DeleteLicenseState](license/license_state_store.go#L119) | +| providersgitlab.NewAuthTokenStore | moved | [NewAuthTokenStore](providers/gitlab/auth_token_store.go#L44) | +| providersgitlab.UpsertGitLabIntegrationToken | moved | [UpsertGitLabIntegrationToken](providers/gitlab/auth_token_store.go#L48) | +| providersgitlab.GetGitLabTokenRecord | moved | [GetGitLabTokenRecord](providers/gitlab/auth_token_store.go#L109) | +| providersgitlab.GetProviderAppID | moved | [GetProviderAppID](providers/gitlab/auth_token_store.go#L122) | +| providersgitlab.UpdateRefreshedToken | moved | [UpdateRefreshedToken](providers/gitlab/auth_token_store.go#L133) | +| providersbitbucket.NewIntegrationTokenStore | moved | [NewIntegrationTokenStore](providers/bitbucket/integration_token_store.go#L18) | +| providersbitbucket.GetTokenByRepoFullName | moved | [GetTokenByRepoFullName](providers/bitbucket/integration_token_store.go#L22) | +| providersbitbucket.GetLatestBitbucketToken | moved | [GetLatestBitbucketToken](providers/bitbucket/integration_token_store.go#L41) | +| providersgithub.NewTokenStore | moved | [NewTokenStore](providers/github/token_store.go#L17) | +| providersgithub.GetLatestGitHubToken | moved | [GetLatestGitHubToken](providers/github/token_store.go#L21) | +| reviews.NewReviewStore | moved | [NewReviewStore](reviews/review_store.go#L9) | +| reviews.QueryRow | moved | [QueryRow](reviews/review_store.go#L13) | +| reviews.Exec | moved | [Exec](reviews/review_store.go#L17) | +| reviews.Query | moved | [Query](reviews/review_store.go#L21) | +| reviews.NewTaxonomyReportStore | added | [NewTaxonomyReportStore](reviews/taxonomy_report_store.go#L17) | +| reviews.buildWhereClause | added | [buildWhereClause](reviews/taxonomy_report_store.go#L108) | +| reviews.GetSummary | added | [GetSummary](reviews/taxonomy_report_store.go#L230) | +| reviews.GetDistribution | added | [GetDistribution](reviews/taxonomy_report_store.go#L273) | +| reviews.GetTrend | added | [GetTrend](reviews/taxonomy_report_store.go#L326) | +| reviews.GetBreakdown | added | [GetBreakdown](reviews/taxonomy_report_store.go#L365) | +| reviews.ListFindings | updated | [ListFindings](reviews/taxonomy_report_store.go#L430) | +| reviews.GetCategorySubcategoryRelations | added | [GetCategorySubcategoryRelations](reviews/taxonomy_report_store.go#L552) | +| aiconnectors.NewConnectorStore | moved | [NewConnectorStore](aiconnectors/connector_store.go#L18) | +| aiconnectors.QueryRowContext | moved | [QueryRowContext](aiconnectors/connector_store.go#L22) | +| aiconnectors.QueryContext | moved | [QueryContext](aiconnectors/connector_store.go#L26) | +| aiconnectors.ExecContext | moved | [ExecContext](aiconnectors/connector_store.go#L30) | +| aiconnectors.UpdateDisplayOrders | moved | [UpdateDisplayOrders](aiconnectors/connector_store.go#L38) | +| aiconnectors.NewReviewAISettingsStore | added | [NewReviewAISettingsStore](aiconnectors/review_ai_settings_store.go#L28) | +| aiconnectors.NormalizeConnectorRole | added | [NormalizeConnectorRole](aiconnectors/review_ai_settings_store.go#L32) | +| aiconnectors.NormalizeHelperMode | added | [NormalizeHelperMode](aiconnectors/review_ai_settings_store.go#L43) | +| aiconnectors.GetByOrgID | added | [GetByOrgID](aiconnectors/review_ai_settings_store.go#L54) | +| aiconnectors.Upsert | added | [Upsert](aiconnectors/review_ai_settings_store.go#L86) | +| aiconnectors.GetConnectorsByRole | added | [GetConnectorsByRole](../internal/aiconnectors/storage.go#L230) | +| aiconnectors.GetMaxDisplayOrderByRole | added | [GetMaxDisplayOrderByRole](../internal/aiconnectors/storage.go#L561) | +| users.NewUserStore | moved | [NewUserStore](users/user_store.go#L9) | +| users.QueryRow | moved | [QueryRow](users/user_store.go#L13) | +| users.Query | moved | [Query](users/user_store.go#L17) | +| users.TxQueryRow | moved | [TxQueryRow](users/user_store.go#L21) | +| users.TxExec | moved | [TxExec](users/user_store.go#L25) | +| users.WithTx | moved | [WithTx](users/user_store.go#L29) | +| users.NewProfileStore | moved | [NewProfileStore](users/profile_store.go#L46) | +| users.GetUserProfile | moved | [GetUserProfile](users/profile_store.go#L50) | +| users.UpdateUserProfile | moved | [UpdateUserProfile](users/profile_store.go#L77) | +| users.GetCurrentPasswordHash | moved | [GetCurrentPasswordHash](users/profile_store.go#L114) | +| users.UpdatePasswordIfCurrentHash | moved | [UpdatePasswordIfCurrentHash](users/profile_store.go#L123) | +| users.ListUserOrganizations | moved | [ListUserOrganizations](users/profile_store.go#L139) | +| learnings.NewLearningsStore | moved | [NewLearningsStore](learnings/learnings_store.go#L67) | +| learnings.InsertLearning | moved | [InsertLearning](learnings/learnings_store.go#L71) | +| learnings.UpdateLearning | moved | [UpdateLearning](learnings/learnings_store.go#L89) | +| learnings.QueryLearningByID | moved | [QueryLearningByID](learnings/learnings_store.go#L107) | +| learnings.QueryLearningByShortID | moved | [QueryLearningByShortID](learnings/learnings_store.go#L114) | +| learnings.ListByOrgWithPagination | moved | [ListByOrgWithPagination](learnings/learnings_store.go#L121) | +| learnings.CountByOrg | moved | [CountByOrg](learnings/learnings_store.go#L154) | +| learnings.InsertLearningEvent | moved | [InsertLearningEvent](learnings/learnings_store.go#L174) | +| providersgitea.NewTokenStore | moved | [NewTokenStore](providers/gitea/token_store.go#L18) | +| providersgitea.ListRecentGiteaIntegrationTokens | moved | [ListRecentGiteaIntegrationTokens](providers/gitea/token_store.go#L22) | +| providersgitea.GetGiteaIntegrationTokenByID | moved | [GetGiteaIntegrationTokenByID](providers/gitea/token_store.go#L51) | +| providersgitea.GetLatestWebhookSecret | moved | [GetLatestWebhookSecret](providers/gitea/token_store.go#L64) | +| core.NewFileOpsStore | moved | [NewFileOpsStore](core/file_ops.go#L7) | +| core.ReadFile | moved | [ReadFile](core/file_ops.go#L11) | +| core.NewSchedulerLockStore | added | [NewSchedulerLockStore](core/scheduler_lock_store.go#L19) | +| core.TryAcquireDashboardRefreshLeaderLock | added | [TryAcquireDashboardRefreshLeaderLock](core/scheduler_lock_store.go#L23) | +| core.ReleaseDashboardRefreshLeaderLock | added | [ReleaseDashboardRefreshLeaderLock](core/scheduler_lock_store.go#L58) | +| core.TryAcquireEventCompactionLeaderLock | added | [TryAcquireEventCompactionLeaderLock](core/compaction_lock_store.go#L15) | +| license.NewPlanCatalogFileStore | moved | [NewPlanCatalogFileStore](license/plan_catalog_file_store.go#L14) | +| license.ReadPlanCatalogFile | moved | [ReadPlanCatalogFile](license/plan_catalog_file_store.go#L18) | +| license.NewQuotaStore | added | [NewQuotaStore](license/quota_store.go#L106) | +| license.ResolvePolicy | added | [ResolvePolicy](license/quota_store.go#L110) | +| license.UpsertBatchSettlement | added | [UpsertBatchSettlement](license/quota_store.go#L165) | +| license.BuildAggregateFromBatches | added | [BuildAggregateFromBatches](license/quota_store.go#L260) | +| license.UpsertOperationAggregate | added | [UpsertOperationAggregate](license/quota_store.go#L311) | +| license.NewLOCAccountingStore | moved | [NewLOCAccountingStore](license/loc_accounting_store.go#L52) | +| license.AccountSuccess | moved | [AccountSuccess](license/loc_accounting_store.go#L56) | +| license.CheckQuotaPreflight | moved | [CheckQuotaPreflight](license/loc_accounting_store.go#L212) | +| license.emitThresholdLifecycleEventsTx | moved | [emitThresholdLifecycleEventsTx](license/loc_accounting_store.go#L358) | +| license.emitLifecycleEventTx | moved | [emitLifecycleEventTx](license/loc_accounting_store.go#L502) | +| license.NewActorLookupStore | added | [NewActorLookupStore](license/actor_lookup_store.go#L15) | +| license.ResolveOrgMemberUserIDByEmail | added | [ResolveOrgMemberUserIDByEmail](license/actor_lookup_store.go#L19) | +| license.NewTrialEligibilityStore | added | [NewTrialEligibilityStore](license/trial_eligibility_store.go#L58) | +| license.NormalizeTrialEligibilityEmail | added | [NormalizeTrialEligibilityEmail](license/trial_eligibility_store.go#L62) | +| license.GetTrialEligibilityByEmail | added | [GetTrialEligibilityByEmail](license/trial_eligibility_store.go#L70) | +| license.ReserveFirstPurchaseTrial | added | [ReserveFirstPurchaseTrial](license/trial_eligibility_store.go#L107) | +| license.ConsumeReservedTrial | added | [ConsumeReservedTrial](license/trial_eligibility_store.go#L171) | +| license.ConsumeReservedTrialTx | added | [ConsumeReservedTrialTx](license/trial_eligibility_store.go#L194) | +| license.ReleaseTrialReservation | added | [ReleaseTrialReservation](license/trial_eligibility_store.go#L261) | +| license.NewAdminBillingPortfolioStore | added | [NewAdminBillingPortfolioStore](license/admin_billing_portfolio_store.go#L38) | +| license.GetSummary | added | [GetSummary](license/admin_billing_portfolio_store.go#L42) | +| license.ListOrganizations | added | [ListOrganizations](license/admin_billing_portfolio_store.go#L88) | +| license.OrganizationExists | added | [OrganizationExists](license/admin_billing_portfolio_store.go#L179) | +| license.NewReviewAccountingStore | moved | [NewReviewAccountingStore](license/review_accounting_store.go#L40) | +| license.GetReviewAccountingTotals | moved | [GetReviewAccountingTotals](license/review_accounting_store.go#L44) | +| license.GetLatestReviewAccountingOperation | moved | [GetLatestReviewAccountingOperation](license/review_accounting_store.go#L109) | +| license.NewOrgUsageStore | moved | [NewOrgUsageStore](license/org_usage_store.go#L54) | +| license.GetCurrentPeriodSummary | moved | [GetCurrentPeriodSummary](license/org_usage_store.go#L58) | +| license.ListCurrentPeriodOperations | moved | [ListCurrentPeriodOperations](license/org_usage_store.go#L115) | +| license.ListCurrentPeriodMemberUsage | added | [ListCurrentPeriodMemberUsage](license/org_usage_store.go#L197) | +| license.GetCurrentPeriodUsageForActor | added | [GetCurrentPeriodUsageForActor](license/org_usage_store.go#L272) | +| payment.GetLatestCapturedPaymentMethodBySubscriptionID | added | [GetLatestCapturedPaymentMethodBySubscriptionID](payment/subscription_store.go#L730) | +| payment.ListSubscriptionsByOrgID | moved | [ListSubscriptionsByOrgID](payment/subscription_store.go#L750) | +| payment.NewBillingNotificationOutboxStore | added | [NewBillingNotificationOutboxStore](payment/billing_notification_outbox_store.go#L45) | +| payment.Enqueue | added | [Enqueue](payment/billing_notification_outbox_store.go#L49) | +| payment.GetUserEmailByID | added | [GetUserEmailByID](payment/billing_notification_outbox_store.go#L104) | +| payment.ClaimDispatchBatch | added | [ClaimDispatchBatch](payment/billing_notification_outbox_store.go#L123) | +| payment.MarkSent | added | [MarkSent](payment/billing_notification_outbox_store.go#L201) | +| payment.MarkFailed | added | [MarkFailed](payment/billing_notification_outbox_store.go#L221) | +| payment.MarkCancelled | added | [MarkCancelled](payment/billing_notification_outbox_store.go#L245) | +| payment.NewUpgradePaymentAttemptStore | added | [NewUpgradePaymentAttemptStore](payment/upgrade_payment_attempt_store.go#L88) | +| payment.CreateUpgradePaymentAttempt | added | [CreateUpgradePaymentAttempt](payment/upgrade_payment_attempt_store.go#L97) | +| payment.GetReusablePreparedAttempt | added | [GetReusablePreparedAttempt](payment/upgrade_payment_attempt_store.go#L132) | +| payment.GetAttemptByOrgPreviewAndOrder | added | [GetAttemptByOrgPreviewAndOrder](payment/upgrade_payment_attempt_store.go#L159) | +| payment.GetAttemptByOrderID | added | [GetAttemptByOrderID](payment/upgrade_payment_attempt_store.go#L185) | +| payment.GetLatestAttemptByUpgradeRequestID | added | [GetLatestAttemptByUpgradeRequestID](payment/upgrade_payment_attempt_store.go#L209) | +| payment.MarkPaymentCapturedByOrderID | added | [MarkPaymentCapturedByOrderID](payment/upgrade_payment_attempt_store.go#L266) | +| payment.MarkPaymentFailedByOrderID | added | [MarkPaymentFailedByOrderID](payment/upgrade_payment_attempt_store.go#L288) | +| payment.ReserveExecute | added | [ReserveExecute](payment/upgrade_payment_attempt_store.go#L320) | +| payment.MarkExecuteApplied | added | [MarkExecuteApplied](payment/upgrade_payment_attempt_store.go#L411) | +| payment.NewUpgradeReplacementCutoverStore | added | [NewUpgradeReplacementCutoverStore](payment/upgrade_replacement_cutover_store.go#L69) | +| payment.CreateOrGetPending | added | [CreateOrGetPending](payment/upgrade_replacement_cutover_store.go#L73) | +| payment.GetByUpgradeRequestID | added | [GetByUpgradeRequestID](payment/upgrade_replacement_cutover_store.go#L113) | +| payment.MarkReplacementProvisioned | added | [MarkReplacementProvisioned](payment/upgrade_replacement_cutover_store.go#L138) | +| payment.MarkOldCancellationScheduled | added | [MarkOldCancellationScheduled](payment/upgrade_replacement_cutover_store.go#L175) | +| payment.MarkRetryPending | added | [MarkRetryPending](payment/upgrade_replacement_cutover_store.go#L209) | +| payment.MarkManualReviewRequired | added | [MarkManualReviewRequired](payment/upgrade_replacement_cutover_store.go#L245) | +| payment.MarkCompleted | added | [MarkCompleted](payment/upgrade_replacement_cutover_store.go#L279) | +| license.NewPlanChangeStore | moved | [NewPlanChangeStore](license/plan_change_store.go#L30) | +| license.EnsureOrgBillingState | moved | [EnsureOrgBillingState](license/plan_change_store.go#L34) | +| license.GetOrgBillingState | moved | [GetOrgBillingState](license/plan_change_store.go#L57) | +| license.ApplyImmediatePlanUpgrade | moved | [ApplyImmediatePlanUpgrade](license/plan_change_store.go#L95) | +| license.ScheduleDowngrade | moved | [ScheduleDowngrade](license/plan_change_store.go#L126) | +| license.ScheduleUpgradeWithCurrentCycleGrant | added | [ScheduleUpgradeWithCurrentCycleGrant](license/plan_change_store.go#L153) | +| license.CancelScheduledDowngrade | moved | [CancelScheduledDowngrade](license/plan_change_store.go#L187) | +| license.ListDueScheduledDowngrades | moved | [ListDueScheduledDowngrades](license/plan_change_store.go#L221) | +| license.ListDueScheduledPlanChanges | added | [ListDueScheduledPlanChanges](license/plan_change_store.go#L225) | +| license.ApplyScheduledDowngrade | moved | [ApplyScheduledDowngrade](license/plan_change_store.go#L254) | +| license.ApplyScheduledPlanChange | added | [ApplyScheduledPlanChange](license/plan_change_store.go#L258) | +| license.insertLifecycleEventTx | moved | [insertLifecycleEventTx](license/plan_change_store.go#L299) | +| analytics.NewAdHocStore | added | [NewAdHocStore](analytics/adhoc_store.go#L50) | +| analytics.WithStatementTimeout | added | [WithStatementTimeout](analytics/adhoc_store.go#L55) | +| analytics.Count | added | [Count](analytics/adhoc_store.go#L67) | +| analytics.Query | added | [Query](analytics/adhoc_store.go#L89) | +| analytics.Ping | added | [Ping](analytics/adhoc_store.go#L186) | diff --git a/ui/src/components/Navbar/megaMenuData.ts b/ui/src/components/Navbar/megaMenuData.ts index 8621c890..cea80026 100644 --- a/ui/src/components/Navbar/megaMenuData.ts +++ b/ui/src/components/Navbar/megaMenuData.ts @@ -238,6 +238,7 @@ export const buildMegaMenuSections = (): MegaMenuSection[] => [ ], React.createElement(Icons.Reports)), group('Manage System', [ link('Storage', React.createElement(Icons.Folder), '/settings#storage', (ctx) => ctx.isSuperAdmin), + link('Log Compaction', React.createElement(Icons.Clock), '/settings#storage', (ctx) => ctx.isSuperAdmin), ], React.createElement(Icons.Settings)), ], }, diff --git a/ui/src/components/reviews/ReviewEventsPage.tsx b/ui/src/components/reviews/ReviewEventsPage.tsx index 460bcf5b..6c226c9d 100644 --- a/ui/src/components/reviews/ReviewEventsPage.tsx +++ b/ui/src/components/reviews/ReviewEventsPage.tsx @@ -8,6 +8,7 @@ import { ReviewEvent } from './types'; interface ReviewEventsPageProps { reviewId: number; initialEvents: ReviewEvent[]; + initialEventCount?: number; isLive?: boolean; pollingInterval?: number; className?: string; @@ -18,12 +19,15 @@ type ViewMode = 'progress' | 'raw'; export default function ReviewEventsPage({ reviewId, initialEvents, + initialEventCount, isLive = false, pollingInterval = 5000, className }: ReviewEventsPageProps) { const [currentView, setCurrentView] = useState('progress'); const [events, setEvents] = useState(initialEvents); + // displayCount tracks original historical count from meta.count (restored from compaction marker) + const [displayCount, setDisplayCount] = useState(initialEventCount ?? initialEvents.length); const [isPolling, setIsPolling] = useState(true); // Default to ON const [lastEventCount, setLastEventCount] = useState(initialEvents.length); const scrollPositionRef = useRef(0); @@ -46,16 +50,17 @@ export default function ReviewEventsPage({ }; // Smooth append for new events (no position disruption) - const appendNewEvents = (newEvents: ReviewEvent[]) => { - if (newEvents.length > events.length) { - const addedEvents = newEvents.slice(events.length); + const appendNewEvents = (allEvents: ReviewEvent[]) => { + if (allEvents.length > events.length) { + const addedEvents = allEvents.slice(events.length); // Save position before update saveScrollPosition(); // Update events - setEvents(newEvents); - setLastEventCount(newEvents.length); + setEvents(allEvents); + setLastEventCount(allEvents.length); + setDisplayCount(allEvents.length); // Restore position after React update setTimeout(restoreScrollPosition, 0); @@ -90,7 +95,7 @@ export default function ReviewEventsPage({ const data = await getReviewEvents(reviewId, undefined, 1000); // Transform backend events to frontend format const backendEvents = data.events || []; - const newEvents: ReviewEvent[] = backendEvents.map((event: any) => { + const allEvents: ReviewEvent[] = backendEvents.map((event: any) => { // Generate human-readable message for display (Raw Events tab) let message = ''; const eventData = event.data || {}; @@ -140,19 +145,25 @@ export default function ReviewEventsPage({ }); console.log('[ReviewEventsPage] Received events:', { - totalEvents: newEvents.length, - sampleEvents: newEvents.slice(0, 10).map(e => ({ + totalEvents: allEvents.length, + sampleEvents: allEvents.slice(0, 10).map(e => ({ message: e.message, eventType: e.eventType, severity: e.severity, timestamp: e.timestamp })), - stageCompletionEvents: newEvents.filter(e => + stageCompletionEvents: allEvents.filter(e => e.message.toLowerCase().includes('stage completed successfully') ).map(e => e.message) }); - appendNewEvents(newEvents); + appendNewEvents(allEvents); + const metaCount = (data as any)?.meta?.count; + if (typeof metaCount === 'number') { + setDisplayCount(metaCount); + } else { + setDisplayCount(allEvents.length); + } } catch (error) { console.error('Failed to poll for updates:', error); } @@ -255,7 +266,7 @@ export default function ReviewEventsPage({
- Raw Events ({events.length}) + Raw Events ({displayCount}) diff --git a/ui/src/pages/Reviews/ReviewDetail.tsx b/ui/src/pages/Reviews/ReviewDetail.tsx index 943f0578..84ad33fd 100644 --- a/ui/src/pages/Reviews/ReviewDetail.tsx +++ b/ui/src/pages/Reviews/ReviewDetail.tsx @@ -56,6 +56,7 @@ const ReviewDetail: React.FC = () => { const mapEventLevel = (level: ReviewEventLevel) => level; const [review, setReview] = useState(null); const [events, setEvents] = useState([]); + const [eventCount, setEventCount] = useState(0); const [summary, setSummary] = useState(null); const [accounting, setAccounting] = useState(null); const [accountingError, setAccountingError] = useState(null); @@ -189,6 +190,13 @@ const ReviewDetail: React.FC = () => { const newEvents = (eventsData?.events as ReviewEvent[] | undefined) || []; setEvents(newEvents); + // meta.count carries the historical total (restored from compaction marker) + if (eventsData?.meta?.count !== undefined) { + setEventCount(eventsData.meta.count); + } else { + setEventCount(newEvents.length); + } + // Update last event time for next polling if (newEvents.length > 0) { const latestTime = newEvents[newEvents.length - 1].time; @@ -474,7 +482,7 @@ const ReviewDetail: React.FC = () => { - a + b, 0))} /> + 0 ? String(eventCount) : String(Object.values(summary?.eventCounts || {}).reduce((a: number, b: number) => a + b, 0))} />
@@ -750,6 +758,7 @@ const ReviewDetail: React.FC = () => {
({ id: event.id.toString(), timestamp: event.time, diff --git a/ui/src/pages/Settings/CompactionSettingsTab.tsx b/ui/src/pages/Settings/CompactionSettingsTab.tsx new file mode 100644 index 00000000..fe3771b9 --- /dev/null +++ b/ui/src/pages/Settings/CompactionSettingsTab.tsx @@ -0,0 +1,279 @@ +import React, { useState, useEffect } from 'react'; +import { Button } from '../../components/UIPrimitives'; +import apiClient from '../../api/apiClient'; +import { notify } from '../../utils/notify'; +import CronBuilder from '../../components/reviews/cronbuilder/CronBuilder'; +import { getLocalCronText, localTimeZoneName } from '../../components/reviews/cronbuilder/cronTimezone'; + +interface CompactionConfig { + enabled: boolean; + cron_expression: string; + retention_days: number; + schedule_human?: string; +} + +const DEFAULT_CONFIG: CompactionConfig = { + enabled: true, + cron_expression: '30 20 * * *', + retention_days: 30, +}; + +const CompactionSettingsTab: React.FC = () => { + const [settings, setSettings] = useState(DEFAULT_CONFIG); + const [isLoading, setIsLoading] = useState(true); + const [isSaving, setIsSaving] = useState(false); + const [isRunning, setIsRunning] = useState(false); + const [isDone, setIsDone] = useState(false); + const [showAdvanced, setShowAdvanced] = useState(false); + + useEffect(() => { + loadConfig(); + }, []); + + const loadConfig = async () => { + setIsLoading(true); + try { + const configData = await apiClient.get('/api/v1/admin/settings/compaction'); + if (configData) setSettings(configData); + } catch { + notify.error('Failed to load compaction settings'); + } finally { + setIsLoading(false); + } + }; + + const saveSettingsToBackend = async (silent = false) => { + try { + await apiClient.put('/api/v1/admin/settings/compaction', { + enabled: settings.enabled, + cron_expression: settings.cron_expression, + retention_days: Number(settings.retention_days), + }); + if (!silent) notify.success('Settings saved successfully!'); + } catch (error: any) { + if (!silent) notify.error(error?.message || 'Failed to save settings'); + throw error; + } + }; + + const handleSave = async () => { + setIsSaving(true); + try { + await saveSettingsToBackend(false); + } catch (error: any) { + notify.error(error?.message || 'Failed to save settings'); + } finally { + setIsSaving(false); + } + }; + + + const handleRunNow = async () => { + setIsRunning(true); + setIsDone(false); + try { + await saveSettingsToBackend(true); + await apiClient.post('/api/v1/admin/settings/compaction/run', {}); + notify.success('Cleanup started in the background!'); + setIsDone(true); + setTimeout(() => setIsDone(false), 3000); + } catch (error: any) { + notify.error(error?.message || 'Failed to start cleanup'); + } finally { + setIsRunning(false); + } + }; + + if (isLoading) { + return ( +
+
+
+ ); + } + + const cronTextObj = getLocalCronText(settings.cron_expression || '0 2 * * *'); + const humanSchedule = cronTextObj.status && cronTextObj.value ? cronTextObj.value : settings.cron_expression; + const tzName = localTimeZoneName(); + + return ( +
+ {/* Header */} +
+
+ + + +
+
+

Event Log Cleanup

+

Automated log retention and cleanup for historical reviews

+
+
+ + {/* Primary View: 100% Static Read-Only Overview Graphic */} +
+ {/* Enable Switch */} +
+
+
+ Enable Automatic Log Cleanup + + Recommended + +
+

+ Automatically removes old verbose streaming logs on schedule. AI findings, comments, and milestones are preserved. +

+
+ +
+ + {/* Read-Only Status Row */} + {settings.enabled && ( +
+
+ Retention Window +
+ {settings.retention_days} Days +
+

+ Prunes verbose debug logs from reviews older than {settings.retention_days} days +

+
+ +
+ Execution Schedule +
+ {humanSchedule} +
+

+ {tzName} +

+
+
+ )} +
+ + {/* Advanced Section (Editing Controls) */} + {settings.enabled && ( +
+ + + {showAdvanced && ( +
+ {/* Fully Symmetrical Editing Card */} +
+ {/* Top Row: Log Retention Period */} +
+
+

Log Retention Period

+

+ Reviews older than this number of days are eligible for automatic log cleanup. +

+
+
+ Keep logs for + setSettings(prev => ({ ...prev, retention_days: Math.max(1, parseInt(e.target.value, 10) || 30) }))} + className="w-16 bg-slate-800 border border-slate-600 rounded px-2 py-1 text-white font-bold text-sm text-center focus:outline-none focus:border-indigo-500" + /> + days +
+
+ + {/* Bottom Row: Execution Schedule Builder */} +
+

Execution Schedule

+

+ Configure timing to run during low-traffic hours for your team. +

+ { + if (newCron) setSettings(prev => ({ ...prev, cron_expression: newCron })); + }} + /> +
+
+ + {/* Run Now Trigger */} +
+
+

Manual Trigger

+

Saves current settings and runs cleanup immediately in the background

+
+ +
+ + {/* Save Settings Action Footer */} +
+ +
+
+ )} +
+ )} +
+ ); +}; + +export default CompactionSettingsTab; diff --git a/ui/src/pages/Settings/Settings.tsx b/ui/src/pages/Settings/Settings.tsx index c88f1b6a..973717dc 100644 --- a/ui/src/pages/Settings/Settings.tsx +++ b/ui/src/pages/Settings/Settings.tsx @@ -10,6 +10,7 @@ import MCPIntegrationTab from './MCPIntegrationTab'; import IntegrationsTab from './IntegrationsTab'; import SMTPSettingsTab from './SMTPSettingsTab'; import StorageSettingsTab from './StorageSettingsTab'; +import CompactionSettingsTab from './CompactionSettingsTab'; import { UserManagement } from '../../components/UserManagement'; import LicenseManagement from '../Licenses/LicenseManagement'; import { useOrgContext } from '../../hooks/useOrgContext'; @@ -380,6 +381,10 @@ const Settings = () => { // Ensure a valid default tab if hash missing or removed (but not for subscription sub-routes) useEffect(() => { + if (location.hash === '#compaction') { + navigate('#storage', { replace: true }); + return; + } if (!isSubscriptionRoute && !location.hash && firstTab) { navigate(`#${firstTab}`, { replace: true }); } diff --git a/ui/src/pages/Settings/StorageSettingsTab.tsx b/ui/src/pages/Settings/StorageSettingsTab.tsx index 0b7d54ec..fb0947b2 100644 --- a/ui/src/pages/Settings/StorageSettingsTab.tsx +++ b/ui/src/pages/Settings/StorageSettingsTab.tsx @@ -2,6 +2,7 @@ import React, { useState, useEffect } from 'react'; import { Button, Input } from '../../components/UIPrimitives'; import apiClient from '../../api/apiClient'; import { notify } from '../../utils/notify'; +import CompactionSettingsTab from './CompactionSettingsTab'; type Backend = 'filesystem' | 's3' | 'gcs' | 'azure'; @@ -281,6 +282,14 @@ const StorageSettingsTab: React.FC = () => {
+ + {/* Divider */} +
+ + {/* Section 2: PostgreSQL Event Log Compaction & Retention */} +
+ +
); };