# Data Retention Implementation Plan > **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. **Goal:** Implement per-organization data retention with a 1-year default, replacing the current global 90-day TimescaleDB retention. **Architecture:** Add `data_retention_days` field to organizations (default 365). An Oban cron worker runs nightly, groups organizations by retention period, and batch-deletes expired data from all time-series tables. TimescaleDB global retention policies are bumped to 730 days as a safety net for orphaned data. **Tech Stack:** Elixir/Ecto migrations, Oban workers, TimescaleDB retention policies, PostgreSQL DELETE with subqueries. --- ## Current State - **TimescaleDB hypertables** (90-day global retention): `monitoring_checks`, `snmp_sensor_readings`, `snmp_interface_stats`, `wireless_client_readings` - **Hypertable with NO retention**: `check_results` - **Regular tables (no retention)**: `snmp_processor_readings`, `snmp_storage_readings` - **Continuous aggregates**: hourly (1yr), daily (5yr) — these stay as-is ## Data Relationship Map (for DELETE queries) Tables with direct `organization_id`: - `check_results`, `wireless_client_readings`, `alerts` Tables requiring joins: - `monitoring_checks` → `device_id` → `devices.organization_id` - `snmp_sensor_readings` → `sensor_id` → `snmp_sensors.snmp_device_id` → `snmp_devices.device_id` → `devices.organization_id` - `snmp_interface_stats` → `interface_id` → `snmp_interfaces.snmp_device_id` → `snmp_devices.device_id` → `devices.organization_id` - `snmp_processor_readings` → `processor_id` → `snmp_processors.snmp_device_id` → `snmp_devices.device_id` → `devices.organization_id` - `snmp_storage_readings` → `storage_id` → `snmp_storage.snmp_device_id` → `snmp_devices.device_id` → `devices.organization_id` --- ### Task 1: Migration — Add `data_retention_days` to organizations **Files:** - Create: `priv/repo/migrations/TIMESTAMP_add_data_retention_days_to_organizations.exs` **Step 1: Generate migration** ```bash cd /Users/graham/dev/towerops/towerops-web mix ecto.gen.migration add_data_retention_days_to_organizations ``` **Step 2: Write migration** ```elixir defmodule Towerops.Repo.Migrations.AddDataRetentionDaysToOrganizations do use Ecto.Migration def change do alter table(:organizations) do add :data_retention_days, :integer, default: 365, null: false end end end ``` **Step 3: Run migration** ```bash mix ecto.migrate ``` Expected: Migration succeeds, all existing orgs get `data_retention_days = 365`. --- ### Task 2: Update Organization Schema and Changeset **Files:** - Modify: `lib/towerops/organizations/organization.ex` - Test: `test/towerops/workers/data_retention_worker_test.exs` (created in Task 5) **Step 1: Add field to schema** In `lib/towerops/organizations/organization.ex`, add to the schema block: ```elixir field :data_retention_days, :integer, default: 365 ``` **Step 2: Add to typespec** Add to the `@type t` block: ```elixir data_retention_days: integer(), ``` **Step 3: Add to changeset** Add `:data_retention_days` to the `cast/2` list in the main `changeset/2` function, and add validation: ```elixir |> validate_number(:data_retention_days, greater_than_or_equal_to: 30, less_than_or_equal_to: 730, message: "must be between 30 and 730 days" ) ``` Min 30 days (no sense retaining less), max 730 days (matches the TimescaleDB safety net). **Step 4: Verify compilation** ```bash mix compile --warnings-as-errors ``` **Step 5: Commit** ```bash git add lib/towerops/organizations/organization.ex priv/repo/migrations/*_add_data_retention_days_to_organizations.exs git commit -m "feat: add data_retention_days to organizations (default 365)" ``` --- ### Task 3: Migration — Update TimescaleDB Retention Policies **Files:** - Create: `priv/repo/migrations/TIMESTAMP_update_timescaledb_retention_to_730_days.exs` This migration: 1. Changes existing 90-day retention to 730-day on all hypertables (safety net) 2. Adds missing retention policy for `check_results` **Step 1: Generate migration** ```bash mix ecto.gen.migration update_timescaledb_retention_to_730_days ``` **Step 2: Write migration** ```elixir defmodule Towerops.Repo.Migrations.UpdateTimescaledbRetentionTo730Days do @moduledoc """ Updates TimescaleDB retention policies from 90 days to 730 days. Per-organization retention is now handled by DataRetentionWorker (default 365 days). The 730-day TimescaleDB policy acts as a safety net for orphaned data from deleted organizations or data that the worker misses. """ use Ecto.Migration @disable_ddl_transaction true @disable_migration_lock true @hypertables ~w(monitoring_checks snmp_sensor_readings snmp_interface_stats wireless_client_readings) def up do if timescaledb_available?() do # Update existing retention policies from 90 days to 730 days for table <- @hypertables do execute("SELECT remove_retention_policy('#{table}', if_exists => true)") execute("SELECT add_retention_policy('#{table}', INTERVAL '730 days', if_not_exists => true)") end # Add missing retention policy for check_results execute("SELECT add_retention_policy('check_results', INTERVAL '730 days', if_not_exists => true)") end end def down do if timescaledb_available?() do # Revert to original 90-day retention for table <- @hypertables do execute("SELECT remove_retention_policy('#{table}', if_exists => true)") execute("SELECT add_retention_policy('#{table}', INTERVAL '90 days', if_not_exists => true)") end # Remove check_results retention (it had none before) execute("SELECT remove_retention_policy('check_results', if_exists => true)") end end defp timescaledb_available? do case repo().query("SELECT 1 FROM pg_available_extensions WHERE name = 'timescaledb'") do {:ok, %{num_rows: n}} when n > 0 -> true _ -> false end end end ``` **Step 3: Run migration** ```bash mix ecto.migrate ``` **Step 4: Commit** ```bash git add priv/repo/migrations/*_update_timescaledb_retention_to_730_days.exs git commit -m "feat: update TimescaleDB retention to 730 days as safety net" ``` --- ### Task 4: Write DataRetentionWorker Tests **Files:** - Create: `test/towerops/workers/data_retention_worker_test.exs` **Step 1: Write the test file** Tests must cover: 1. Worker deletes data older than org's retention period 2. Worker respects per-org retention overrides 3. Worker handles tables with direct `organization_id` (check_results, wireless_client_readings) 4. Worker handles tables requiring joins (monitoring_checks, sensor_readings, etc.) 5. Worker does NOT delete data within retention window 6. Worker returns summary counts Use existing test factories/fixtures. Key patterns: - Create an org with a short retention (e.g., 1 day) - Insert old data (2 days ago) and recent data (now) - Run worker - Assert old data deleted, recent data preserved ```elixir defmodule Towerops.Workers.DataRetentionWorkerTest do use Towerops.DataCase, async: false import Towerops.AccountsFixtures import Towerops.OrganizationsFixtures alias Towerops.Organizations alias Towerops.Repo alias Towerops.Workers.DataRetentionWorker describe "perform/1" do setup do user = user_fixture() org = organization_fixture(user) # Set short retention for testing {:ok, org} = Organizations.update_organization(org, %{data_retention_days: 1}) device = device_fixture(org) %{org: org, device: device, user: user} end test "deletes monitoring_checks older than org retention", %{device: device, org: org} do old_time = DateTime.add(DateTime.utc_now(), -2, :day) recent_time = DateTime.utc_now() # Insert old and recent monitoring checks Repo.insert_all("monitoring_checks", [ %{ id: Ecto.UUID.generate(), device_id: device.id, status: "success", checked_at: old_time, inserted_at: old_time }, %{ id: Ecto.UUID.generate(), device_id: device.id, status: "success", checked_at: recent_time, inserted_at: recent_time } ]) assert {:ok, result} = DataRetentionWorker.perform(%Oban.Job{args: %{}}) # Old data deleted, recent data preserved checks = Repo.all(from(m in "monitoring_checks", where: m.device_id == ^device.id, select: m.id)) assert length(checks) == 1 end test "deletes check_results older than org retention", %{org: org} do check = check_fixture(org) old_time = DateTime.add(DateTime.utc_now(), -2, :day) recent_time = DateTime.utc_now() Repo.insert_all("check_results", [ %{ id: Ecto.UUID.generate(), check_id: check.id, organization_id: org.id, status: 0, checked_at: old_time }, %{ id: Ecto.UUID.generate(), check_id: check.id, organization_id: org.id, status: 0, checked_at: recent_time } ]) assert {:ok, _result} = DataRetentionWorker.perform(%Oban.Job{args: %{}}) results = Repo.all(from(r in "check_results", where: r.organization_id == ^org.id, select: r.id)) assert length(results) == 1 end test "respects different retention periods per org" do user1 = user_fixture() user2 = user_fixture() org_short = organization_fixture(user1) org_long = organization_fixture(user2) {:ok, org_short} = Organizations.update_organization(org_short, %{data_retention_days: 1}) {:ok, org_long} = Organizations.update_organization(org_long, %{data_retention_days: 60}) device_short = device_fixture(org_short) device_long = device_fixture(org_long) # Data from 3 days ago — within org_long's window, outside org_short's three_days_ago = DateTime.add(DateTime.utc_now(), -3, :day) Repo.insert_all("monitoring_checks", [ %{ id: Ecto.UUID.generate(), device_id: device_short.id, status: "success", checked_at: three_days_ago, inserted_at: three_days_ago }, %{ id: Ecto.UUID.generate(), device_id: device_long.id, status: "success", checked_at: three_days_ago, inserted_at: three_days_ago } ]) assert {:ok, _result} = DataRetentionWorker.perform(%Oban.Job{args: %{}}) short_checks = Repo.all(from(m in "monitoring_checks", where: m.device_id == ^device_short.id, select: m.id)) long_checks = Repo.all(from(m in "monitoring_checks", where: m.device_id == ^device_long.id, select: m.id)) assert length(short_checks) == 0 assert length(long_checks) == 1 end test "does not delete data within retention window", %{device: device} do recent_time = DateTime.utc_now() Repo.insert_all("monitoring_checks", [ %{ id: Ecto.UUID.generate(), device_id: device.id, status: "success", checked_at: recent_time, inserted_at: recent_time } ]) assert {:ok, _result} = DataRetentionWorker.perform(%Oban.Job{args: %{}}) checks = Repo.all(from(m in "monitoring_checks", where: m.device_id == ^device.id, select: m.id)) assert length(checks) == 1 end end end ``` Note: These tests use `Repo.insert_all` directly into table names (strings) to insert into hypertables that may not have standard Ecto schema primary keys. Adjust fixtures as needed — check what `device_fixture/1` and `check_fixture/1` return in the existing test helpers. **Step 2: Run tests to verify they fail** ```bash mix test test/towerops/workers/data_retention_worker_test.exs ``` Expected: Compilation error — `DataRetentionWorker` module doesn't exist yet. --- ### Task 5: Implement DataRetentionWorker **Files:** - Create: `lib/towerops/workers/data_retention_worker.ex` - Modify: `config/runtime.exs` (add to Oban crontab) **Step 1: Create the worker** ```elixir defmodule Towerops.Workers.DataRetentionWorker do @moduledoc """ Oban worker for per-organization data retention. Runs nightly to delete time-series data older than each organization's `data_retention_days` setting (default: 365 days). Tables cleaned: - monitoring_checks (via device_id → devices.organization_id) - snmp_sensor_readings (via sensor → snmp_device → device → org) - snmp_interface_stats (via interface → snmp_device → device → org) - snmp_processor_readings (via processor → snmp_device → device → org) - snmp_storage_readings (via storage → snmp_device → device → org) - check_results (direct organization_id) - wireless_client_readings (direct organization_id) - alerts (direct organization_id, only resolved alerts) TimescaleDB global retention (730 days) acts as a safety net. """ use Oban.Worker, queue: :maintenance require Logger alias Towerops.Repo @impl Oban.Worker def perform(_job) do # Group orgs by retention_days for efficient batch deletion retention_groups = fetch_retention_groups() results = Enum.flat_map(retention_groups, fn {retention_days, org_ids} -> cutoff = DateTime.add(DateTime.utc_now(), -retention_days, :day) purge_for_orgs(org_ids, cutoff) end) total = Enum.reduce(results, 0, fn {_table, count}, acc -> acc + count end) if total > 0 do summary = Enum.reject(results, fn {_table, count} -> count == 0 end) Logger.info("Data retention: deleted #{total} total rows — #{inspect(summary)}") end {:ok, %{total_deleted: total, details: Map.new(results)}} end defp fetch_retention_groups do Repo.query!(""" SELECT data_retention_days, array_agg(id) AS org_ids FROM organizations GROUP BY data_retention_days """) |> Map.get(:rows) |> Enum.map(fn [days, ids] -> {days, Enum.map(ids, &Ecto.UUID.cast!/1)} end) end defp purge_for_orgs(org_ids, cutoff) do # Tables with direct organization_id direct_tables = [ {"check_results", :checked_at}, {"wireless_client_readings", :checked_at} ] direct_results = Enum.map(direct_tables, fn {table, time_col} -> {count, _} = Repo.query!( "DELETE FROM #{table} WHERE organization_id = ANY($1) AND #{time_col} < $2", [org_ids, cutoff] ) |> then(fn %{num_rows: n} -> {n, nil} end) {table, count} end) # Resolved alerts (keep unresolved regardless of age) {alert_count, _} = Repo.query!( "DELETE FROM alerts WHERE organization_id = ANY($1) AND inserted_at < $2 AND resolved_at IS NOT NULL", [org_ids, cutoff] ) |> then(fn %{num_rows: n} -> {n, nil} end) # monitoring_checks: device_id → devices.organization_id mc_result = Repo.query!( """ DELETE FROM monitoring_checks WHERE checked_at < $2 AND device_id IN (SELECT id FROM devices WHERE organization_id = ANY($1)) """, [org_ids, cutoff] ) # snmp_sensor_readings: sensor_id → sensors.snmp_device_id → snmp_devices.device_id → devices.organization_id sr_result = Repo.query!( """ DELETE FROM snmp_sensor_readings WHERE checked_at < $2 AND sensor_id IN ( SELECT s.id FROM snmp_sensors s JOIN snmp_devices sd ON sd.id = s.snmp_device_id JOIN devices d ON d.id = sd.device_id WHERE d.organization_id = ANY($1) ) """, [org_ids, cutoff] ) # snmp_interface_stats: interface_id → interfaces.snmp_device_id → ... is_result = Repo.query!( """ DELETE FROM snmp_interface_stats WHERE checked_at < $2 AND interface_id IN ( SELECT i.id FROM snmp_interfaces i JOIN snmp_devices sd ON sd.id = i.snmp_device_id JOIN devices d ON d.id = sd.device_id WHERE d.organization_id = ANY($1) ) """, [org_ids, cutoff] ) # snmp_processor_readings pr_result = Repo.query!( """ DELETE FROM snmp_processor_readings WHERE checked_at < $2 AND processor_id IN ( SELECT p.id FROM snmp_processors p JOIN snmp_devices sd ON sd.id = p.snmp_device_id JOIN devices d ON d.id = sd.device_id WHERE d.organization_id = ANY($1) ) """, [org_ids, cutoff] ) # snmp_storage_readings str_result = Repo.query!( """ DELETE FROM snmp_storage_readings WHERE checked_at < $2 AND storage_id IN ( SELECT st.id FROM snmp_storage st JOIN snmp_devices sd ON sd.id = st.snmp_device_id JOIN devices d ON d.id = sd.device_id WHERE d.organization_id = ANY($1) ) """, [org_ids, cutoff] ) direct_results ++ [ {"alerts_resolved", alert_count}, {"monitoring_checks", mc_result.num_rows}, {"snmp_sensor_readings", sr_result.num_rows}, {"snmp_interface_stats", is_result.num_rows}, {"snmp_processor_readings", pr_result.num_rows}, {"snmp_storage_readings", str_result.num_rows} ] end end ``` **Step 2: Add to Oban crontab in `config/runtime.exs`** Add after the billing sync entry (around line 255): ```elixir # Data retention cleanup nightly at 1 AM {"0 1 * * *", Towerops.Workers.DataRetentionWorker}, ``` **Step 3: Run the tests** ```bash mix test test/towerops/workers/data_retention_worker_test.exs ``` Expected: All tests pass. **Step 4: Run full test suite** ```bash mix test ``` **Step 5: Run quality checks** ```bash mix credo --strict mix dialyzer ``` **Step 6: Commit** ```bash git add lib/towerops/workers/data_retention_worker.ex test/towerops/workers/data_retention_worker_test.exs config/runtime.exs git commit -m "feat: add per-org data retention worker (default 1 year)" ``` --- ### Task 6: Final Verification **Step 1: Run precommit** ```bash mix precommit ``` **Step 2: Verify migrations run cleanly on fresh DB** ```bash mix ecto.reset ``` **Step 3: Run full test suite with coverage** ```bash mix test --cover ``` --- ## Notes - **Continuous aggregates** (hourly 1yr, daily 5yr) are NOT touched by this worker — they have their own TimescaleDB retention policies - **Unresolved alerts** are never deleted regardless of age - **Login history, sessions, weather** have their own cleanup workers and are unaffected - **Future UI**: The `data_retention_days` field is ready for an org settings page — just needs a form field in the settings LiveView - **Performance**: The worker runs at 1 AM to minimize load. For very large datasets, consider adding `LIMIT` batching per table