mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-09-11 14:00:29 +00:00
Fix legacy resource_changes migration
This commit is contained in:
@@ -84,6 +84,10 @@ registry rebuilds and supplemental ingest into `ResourceChange` records, and
|
||||
`internal/unifiedresources/store.go` persists those changes so `RecentChanges`
|
||||
can round-trip through the SQLite-backed resource store instead of living only
|
||||
in memory or adapter-local state.
|
||||
That store also now migrates legacy `resource_changes` tables that still carry
|
||||
the pre-v6 `timestamp` column by backfilling canonical `observed_at` values,
|
||||
adding the newer `occurred_at` field, and preserving the legacy timestamp on
|
||||
write when the target database still requires it.
|
||||
`internal/api/resources.go` now exposes that same history through dedicated
|
||||
`/api/resources/{id}/timeline` reads, while `/api/resources/{id}/capabilities`
|
||||
and `/api/resources/{id}/relationships` expose the current graph facets as
|
||||
|
||||
@@ -59,9 +59,10 @@ type ResourceExclusion struct {
|
||||
|
||||
// SQLiteResourceStore stores overrides in SQLite.
|
||||
type SQLiteResourceStore struct {
|
||||
db *sql.DB
|
||||
dbPath string
|
||||
mu sync.Mutex
|
||||
db *sql.DB
|
||||
dbPath string
|
||||
resourceChangesHasTimestamp bool
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
const (
|
||||
@@ -303,7 +304,6 @@ func (s *SQLiteResourceStore) initSchema() error {
|
||||
related_resources TEXT,
|
||||
metadata_json TEXT
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_resource_changes_canonical_time ON resource_changes(canonical_id, observed_at DESC);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS action_audits (
|
||||
id TEXT PRIMARY KEY,
|
||||
@@ -361,7 +361,16 @@ func (s *SQLiteResourceStore) migrateResourceChangesSchema() error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, ok := columns["timestamp"]; ok {
|
||||
s.resourceChangesHasTimestamp = true
|
||||
}
|
||||
|
||||
if err := s.addResourceChangesColumnIfMissing(columns, "observed_at", "DATETIME"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.addResourceChangesColumnIfMissing(columns, "occurred_at", "DATETIME"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.addResourceChangesColumnIfMissing(columns, "source_type", "TEXT NOT NULL DEFAULT 'pulse_diff'"); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -378,6 +387,9 @@ func (s *SQLiteResourceStore) migrateResourceChangesSchema() error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := s.backfillLegacyResourceChangeObservedAt(columns); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.normalizeResourceChangeRows(columns); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -466,6 +478,23 @@ func (s *SQLiteResourceStore) normalizeResourceChangeRows(columns map[string]str
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SQLiteResourceStore) backfillLegacyResourceChangeObservedAt(columns map[string]struct{}) error {
|
||||
if _, ok := columns["observed_at"]; !ok {
|
||||
return nil
|
||||
}
|
||||
|
||||
expressions := []string{"observed_at = COALESCE(observed_at, CURRENT_TIMESTAMP)"}
|
||||
if _, ok := columns["timestamp"]; ok {
|
||||
expressions[0] = "observed_at = COALESCE(observed_at, timestamp, CURRENT_TIMESTAMP)"
|
||||
}
|
||||
|
||||
query := `UPDATE resource_changes SET ` + expressions[0] + ` WHERE observed_at IS NULL OR TRIM(COALESCE(observed_at, '')) = ''`
|
||||
if _, err := s.db.Exec(query); err != nil {
|
||||
return fmt.Errorf("backfill resource_changes.observed_at: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SQLiteResourceStore) AddLink(link ResourceLink) error {
|
||||
if link.ResourceA == "" || link.ResourceB == "" {
|
||||
return fmt.Errorf("resource IDs required")
|
||||
@@ -603,10 +632,48 @@ func (s *SQLiteResourceStore) RecordChange(change ResourceChange) error {
|
||||
relJSON, _ := json.Marshal(change.RelatedResources)
|
||||
metaJSON, _ := json.Marshal(change.Metadata)
|
||||
|
||||
_, err := s.db.Exec(`
|
||||
INSERT INTO resource_changes (id, canonical_id, observed_at, occurred_at, kind, from_state, to_state, source_type, source_adapter, actor, confidence, reason, related_resources, metadata_json)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
`, change.ID, CanonicalResourceID(change.ResourceID), change.ObservedAt, change.OccurredAt, string(change.Kind), change.From, change.To, change.SourceType, change.SourceAdapter, change.Actor, string(change.Confidence), change.Reason, string(relJSON), string(metaJSON))
|
||||
columns := []string{"id", "canonical_id", "observed_at"}
|
||||
values := []any{change.ID, CanonicalResourceID(change.ResourceID), change.ObservedAt}
|
||||
if s.resourceChangesHasTimestamp {
|
||||
columns = append(columns, "timestamp")
|
||||
values = append(values, change.ObservedAt)
|
||||
}
|
||||
columns = append(columns,
|
||||
"occurred_at",
|
||||
"kind",
|
||||
"from_state",
|
||||
"to_state",
|
||||
"source_type",
|
||||
"source_adapter",
|
||||
"actor",
|
||||
"confidence",
|
||||
"reason",
|
||||
"related_resources",
|
||||
"metadata_json",
|
||||
)
|
||||
values = append(values,
|
||||
change.OccurredAt,
|
||||
string(change.Kind),
|
||||
change.From,
|
||||
change.To,
|
||||
change.SourceType,
|
||||
change.SourceAdapter,
|
||||
change.Actor,
|
||||
string(change.Confidence),
|
||||
change.Reason,
|
||||
string(relJSON),
|
||||
string(metaJSON),
|
||||
)
|
||||
|
||||
placeholders := make([]string, len(columns))
|
||||
for i := range placeholders {
|
||||
placeholders[i] = "?"
|
||||
}
|
||||
|
||||
_, err := s.db.Exec(
|
||||
`INSERT INTO resource_changes (`+strings.Join(columns, ", ")+`) VALUES (`+strings.Join(placeholders, ", ")+`)`,
|
||||
values...,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("insert resource change: %w", err)
|
||||
}
|
||||
|
||||
@@ -141,8 +141,7 @@ func TestNewSQLiteResourceStore_MigratesLegacyResourceChangesTable(t *testing.T)
|
||||
CREATE TABLE resource_changes (
|
||||
id TEXT PRIMARY KEY,
|
||||
canonical_id TEXT NOT NULL,
|
||||
observed_at DATETIME NOT NULL,
|
||||
occurred_at DATETIME,
|
||||
timestamp DATETIME NOT NULL,
|
||||
kind TEXT NOT NULL,
|
||||
from_state TEXT,
|
||||
to_state TEXT,
|
||||
@@ -156,9 +155,9 @@ func TestNewSQLiteResourceStore_MigratesLegacyResourceChangesTable(t *testing.T)
|
||||
}
|
||||
if _, err := db.Exec(`
|
||||
INSERT INTO resource_changes (
|
||||
id, canonical_id, observed_at, occurred_at, kind, from_state, to_state, source, confidence, reason
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
`, "chg-legacy", "vm:legacy", time.Date(2026, 3, 18, 12, 0, 0, 0, time.UTC), nil, string(ChangeStateTransition), "offline", "online", "proxmox", string(ConfidenceHigh), "legacy row"); err != nil {
|
||||
id, canonical_id, timestamp, kind, from_state, to_state, source, confidence, reason
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
`, "chg-legacy", "vm:legacy", time.Date(2026, 3, 18, 12, 0, 0, 0, time.UTC), string(ChangeStateTransition), "offline", "online", "proxmox", string(ConfidenceHigh), "legacy row"); err != nil {
|
||||
_ = db.Close()
|
||||
t.Fatalf("insert legacy resource change failed: %v", err)
|
||||
}
|
||||
@@ -188,6 +187,12 @@ func TestNewSQLiteResourceStore_MigratesLegacyResourceChangesTable(t *testing.T)
|
||||
if results[0].SourceAdapter != ChangeSourceAdapter("proxmox") {
|
||||
t.Fatalf("legacy source adapter = %q, want proxmox", results[0].SourceAdapter)
|
||||
}
|
||||
if !results[0].ObservedAt.Equal(time.Date(2026, 3, 18, 12, 0, 0, 0, time.UTC)) {
|
||||
t.Fatalf("legacy observed_at = %v, want 2026-03-18T12:00:00Z", results[0].ObservedAt)
|
||||
}
|
||||
if results[0].OccurredAt != nil {
|
||||
t.Fatalf("legacy occurred_at = %v, want nil", results[0].OccurredAt)
|
||||
}
|
||||
|
||||
if err := store.RecordChange(ResourceChange{
|
||||
ID: "chg-new",
|
||||
@@ -209,6 +214,16 @@ func TestNewSQLiteResourceStore_MigratesLegacyResourceChangesTable(t *testing.T)
|
||||
t.Fatalf("GetRecentChanges after migration write returned %d rows, want 2", len(results))
|
||||
}
|
||||
|
||||
columns, err := resourceChangeColumns(store.db)
|
||||
if err != nil {
|
||||
t.Fatalf("resourceChangeColumns: %v", err)
|
||||
}
|
||||
for _, want := range []string{"observed_at", "occurred_at", "source_type", "source_adapter", "actor", "related_resources", "metadata_json"} {
|
||||
if _, ok := columns[want]; !ok {
|
||||
t.Fatalf("expected migrated resource_changes column %q, got %#v", want, columns)
|
||||
}
|
||||
}
|
||||
|
||||
indexes, err := resourceChangesIndexes(store.db)
|
||||
if err != nil {
|
||||
t.Fatalf("resourceChangesIndexes: %v", err)
|
||||
@@ -225,6 +240,34 @@ func TestNewSQLiteResourceStore_MigratesLegacyResourceChangesTable(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func resourceChangeColumns(db *sql.DB) (map[string]struct{}, error) {
|
||||
rows, err := db.Query(`PRAGMA table_info(resource_changes)`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
columns := make(map[string]struct{})
|
||||
for rows.Next() {
|
||||
var (
|
||||
cid int
|
||||
name string
|
||||
typ string
|
||||
notNull int
|
||||
dflt sql.NullString
|
||||
pk int
|
||||
)
|
||||
if err := rows.Scan(&cid, &name, &typ, ¬Null, &dflt, &pk); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
columns[name] = struct{}{}
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return columns, nil
|
||||
}
|
||||
|
||||
func resourceChangesIndexes(db *sql.DB) (map[string]struct{}, error) {
|
||||
rows, err := db.Query(`PRAGMA index_list(resource_changes)`)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user