package store import ( "context" "database/sql" "errors" "fmt" "strings" "time" ) const orphanReportRetention = 7 * 24 * time.Hour func (s *Postgres) AcquireOperationalLease( ctx context.Context, name, owner, token string, _ time.Time, duration time.Duration, ) (bool, error) { if strings.TrimSpace(name) == "" || strings.TrimSpace(owner) == "" || strings.TrimSpace(token) == "" || duration <= 0 { return false, errors.New("invalid operational lease") } var acquired int err := s.db.QueryRowContext(ctx, `INSERT INTO sense.operational_leases( lease_name, owner_id, fencing_token, lease_until, updated_at ) VALUES ( $1, $2, $3, clock_timestamp() + ($4 * interval '1 second'), clock_timestamp() ) ON CONFLICT (lease_name) DO UPDATE SET owner_id = EXCLUDED.owner_id, fencing_token = EXCLUDED.fencing_token, lease_until = EXCLUDED.lease_until, updated_at = EXCLUDED.updated_at WHERE sense.operational_leases.lease_until <= clock_timestamp() OR (sense.operational_leases.owner_id = $2 AND sense.operational_leases.fencing_token = $3) RETURNING 1`, name, owner, token, duration.Seconds()).Scan(&acquired) if errors.Is(err, sql.ErrNoRows) { return false, nil } if err != nil { return false, errors.New("acquire postgres operational lease") } return acquired == 1, nil } func (s *Postgres) ReleaseOperationalLease( ctx context.Context, name, owner, token string, _ time.Time, ) error { _, err := s.db.ExecContext(ctx, `UPDATE sense.operational_leases SET lease_until = clock_timestamp(), updated_at = clock_timestamp() WHERE lease_name = $1 AND owner_id = $2 AND fencing_token = $3`, name, owner, token) if err != nil { return errors.New("release postgres operational lease") } return nil } func (s *Postgres) ListMediaPathOwnership(ctx context.Context) ([]MediaPathOwnership, error) { rows, err := s.db.QueryContext(ctx, `SELECT o.path_name, o.device_id, EXISTS ( SELECT 1 FROM sense.devices d WHERE d.id = o.device_id AND d.path_name = o.path_name AND EXISTS ( SELECT 1 FROM sense.device_capabilities c WHERE c.device_id = d.id AND c.capability = 'video_capture' ) ) AS current_claim FROM sense.media_path_ownership o ORDER BY o.path_name`) if err != nil { return nil, errors.New("list postgres MediaMTX path ownership") } defer rows.Close() values := make([]MediaPathOwnership, 0) for rows.Next() { var value MediaPathOwnership if err := rows.Scan(&value.PathName, &value.DeviceID, &value.CurrentClaim); err != nil { return nil, errors.New("scan postgres MediaMTX path ownership") } values = append(values, value) } if err := rows.Err(); err != nil { return nil, errors.New("iterate postgres MediaMTX path ownership") } return values, nil } func (s *Postgres) SaveOrphanScan( ctx context.Context, scan OrphanScan, owner, token string, ) error { if err := validateOrphanScan(scan); err != nil { return err } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin postgres orphan scan save") } defer tx.Rollback() var lease int err = tx.QueryRowContext(ctx, `SELECT 1 FROM sense.operational_leases WHERE lease_name = $1 AND owner_id = $2 AND fencing_token = $3 AND lease_until > clock_timestamp() FOR UPDATE`, OperationalLeaseOrphanScan, owner, token).Scan(&lease) if errors.Is(err, sql.ErrNoRows) { return ErrOperationalLeaseLost } if err != nil { return errors.New("verify postgres orphan scan lease") } if _, err := tx.ExecContext(ctx, `INSERT INTO sense.orphan_scan_runs( id, instance_id, observed_count, owned_stale_count, unowned_count, safety_allowed, safety_reason, completed_at, expires_at ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9)`, scan.ID, scan.InstanceID, scan.ObservedCount, scan.OwnedStaleCount, scan.UnownedCount, scan.SafetyAllowed, scan.SafetyReason, scan.CompletedAt, scan.ExpiresAt, ); err != nil { return errors.New("insert postgres orphan scan") } for _, finding := range scan.Findings { var deviceID any if finding.DeviceID != "" { deviceID = finding.DeviceID } if _, err := tx.ExecContext(ctx, `INSERT INTO sense.orphan_scan_findings( scan_id, path_name, classification, device_id ) VALUES ($1,$2,$3,$4)`, scan.ID, finding.PathName, finding.Classification, deviceID); err != nil { return errors.New("insert postgres orphan scan finding") } } if _, err := tx.ExecContext(ctx, `UPDATE sense.operational_leases SET lease_until = clock_timestamp(), updated_at = clock_timestamp() WHERE lease_name = $1 AND owner_id = $2 AND fencing_token = $3`, OperationalLeaseOrphanScan, owner, token); err != nil { return errors.New("release postgres orphan scan lease") } if _, err := tx.ExecContext(ctx, `DELETE FROM sense.orphan_scan_runs r WHERE r.completed_at < $1 AND NOT EXISTS ( SELECT 1 FROM sense.orphan_cleanup_actions a WHERE a.scan_id = r.id )`, scan.CompletedAt.Add(-orphanReportRetention)); err != nil { return errors.New("expire postgres orphan scan reports") } if err := tx.Commit(); err != nil { return errors.New("commit postgres orphan scan") } return nil } func validateOrphanScan(scan OrphanScan) error { if strings.TrimSpace(scan.ID) == "" || strings.TrimSpace(scan.InstanceID) == "" || scan.ObservedCount < 0 || scan.OwnedStaleCount < 0 || scan.UnownedCount < 0 || scan.OwnedStaleCount+scan.UnownedCount > scan.ObservedCount || strings.TrimSpace(scan.SafetyReason) == "" || !scan.ExpiresAt.After(scan.CompletedAt) { return errors.New("invalid orphan scan") } seen := make(map[string]struct{}, len(scan.Findings)) staleCount, unownedCount := 0, 0 for _, finding := range scan.Findings { if strings.TrimSpace(finding.PathName) == "" { return errors.New("invalid orphan finding path") } if _, duplicate := seen[finding.PathName]; duplicate { return errors.New("duplicate orphan finding path") } seen[finding.PathName] = struct{}{} switch finding.Classification { case OrphanOwnedStale: staleCount++ if strings.TrimSpace(finding.DeviceID) == "" { return errors.New("owned stale finding lacks device") } case OrphanUnowned: unownedCount++ if finding.DeviceID != "" { return errors.New("unowned finding has device") } default: return fmt.Errorf("invalid orphan finding classification %q", finding.Classification) } } if staleCount != scan.OwnedStaleCount || unownedCount != scan.UnownedCount { return errors.New("orphan scan counts do not match findings") } return nil } func (s *Postgres) GetOrphanScan(ctx context.Context, id string) (OrphanScan, error) { var scan OrphanScan err := s.db.QueryRowContext(ctx, `SELECT id, instance_id, observed_count, owned_stale_count, unowned_count, safety_allowed, safety_reason, completed_at, expires_at FROM sense.orphan_scan_runs WHERE id = $1`, id).Scan( &scan.ID, &scan.InstanceID, &scan.ObservedCount, &scan.OwnedStaleCount, &scan.UnownedCount, &scan.SafetyAllowed, &scan.SafetyReason, &scan.CompletedAt, &scan.ExpiresAt, ) if errors.Is(err, sql.ErrNoRows) { return OrphanScan{}, ErrOrphanScanNotFound } if err != nil { return OrphanScan{}, errors.New("read postgres orphan scan") } rows, err := s.db.QueryContext(ctx, `SELECT f.path_name, f.classification, COALESCE(f.device_id, ''), COALESCE(a.status = 'deleted', false) FROM sense.orphan_scan_findings f LEFT JOIN sense.orphan_cleanup_actions a ON a.scan_id = f.scan_id AND a.path_name = f.path_name WHERE f.scan_id = $1 ORDER BY f.path_name`, id) if err != nil { return OrphanScan{}, errors.New("list postgres orphan scan findings") } defer rows.Close() scan.Findings = make([]OrphanFinding, 0) for rows.Next() { var finding OrphanFinding if err := rows.Scan( &finding.PathName, &finding.Classification, &finding.DeviceID, &finding.Deleted, ); err != nil { return OrphanScan{}, errors.New("scan postgres orphan finding") } scan.Findings = append(scan.Findings, finding) } if err := rows.Err(); err != nil { return OrphanScan{}, errors.New("iterate postgres orphan findings") } return scan, nil } func (s *Postgres) RecordOrphanCleanup( ctx context.Context, scanID, pathName, actorID, status, errorCode string, now time.Time, ) error { if status != "deleted" && status != "failed" { return errors.New("invalid orphan cleanup status") } var storedError any if status == "failed" { if strings.TrimSpace(errorCode) == "" { return errors.New("failed orphan cleanup requires an error code") } storedError = errorCode } _, err := s.db.ExecContext(ctx, `INSERT INTO sense.orphan_cleanup_actions( scan_id, path_name, classification, actor_id, status, error_code, attempted_at ) VALUES ($1,$2,'owned_stale',$3,$4,$5,$6) ON CONFLICT (scan_id, path_name) DO UPDATE SET actor_id = EXCLUDED.actor_id, status = EXCLUDED.status, error_code = EXCLUDED.error_code, attempted_at = EXCLUDED.attempted_at WHERE sense.orphan_cleanup_actions.status <> 'deleted'`, scanID, pathName, actorID, status, storedError, now) if err != nil { return errors.New("record postgres orphan cleanup result") } return nil }