diff --git a/internal/storage/processing.go b/internal/storage/processing.go new file mode 100644 index 0000000..ecfaddf --- /dev/null +++ b/internal/storage/processing.go @@ -0,0 +1,1070 @@ +package storage + +import ( + "context" + "database/sql" + "errors" + "fmt" + "strings" + "time" +) + +const ( + ScheduleRunStatusRunning = "running" + ScheduleRunStatusSucceeded = "succeeded" + ScheduleRunStatusPartial = "partial" + ScheduleRunStatusFailed = "failed" + + AlertEventStatusOpen = "open" + AlertEventStatusResolved = "resolved" + + CandidateSuggestionStatusPendingReview = "pending_review" + CandidateSuggestionStatusAccepted = "accepted" + CandidateSuggestionStatusDismissed = "dismissed" +) + +type ScheduleRun struct { + ID int64 + DatasetID string + ScheduleJobID string + SourceID string + Status string + StartedAt time.Time + FinishedAt *time.Time + ErrorMessage string + FetchRunsStarted int + FetchRunsSucceeded int + FetchRunsFailed int +} + +type StartScheduleRunParams struct { + DatasetID string + ScheduleJobID string + StartedAt time.Time +} + +type FinalizeScheduleRunSuccessParams struct { + ID int64 + FinishedAt time.Time + FetchRunsStarted int + FetchRunsSucceeded int + FetchRunsFailed int +} + +type FinalizeScheduleRunPartialParams struct { + ID int64 + FinishedAt time.Time + ErrorMessage string + FetchRunsStarted int + FetchRunsSucceeded int + FetchRunsFailed int +} + +type FinalizeScheduleRunFailureParams struct { + ID int64 + FinishedAt time.Time + ErrorMessage string + FetchRunsStarted int + FetchRunsSucceeded int + FetchRunsFailed int +} + +type AlertEvent struct { + ID int64 + DatasetID string + AlertRuleID string + EntityID string + SourceID string + Severity string + Status string + Summary string + PayloadJSON string + ObservedAt time.Time + TriggeredAt time.Time + ResolvedAt *time.Time +} + +type InsertAlertEventParams struct { + DatasetID string + AlertRuleID string + EntityID string + Summary string + PayloadJSON string + ObservedAt time.Time + TriggeredAt time.Time +} + +type ResolveAlertEventParams struct { + ID int64 + ResolvedAt time.Time +} + +type CandidateSuggestion struct { + ID int64 + DatasetID string + AlertRuleID string + CandidateKey string + DisplayName string + DiscoveredFromEntityID string + SourceID string + MentionCount int + FirstSeenAt time.Time + LastSeenAt time.Time + Status string + EvidenceJSON string +} + +type UpsertCandidateSuggestionParams struct { + DatasetID string + AlertRuleID string + CandidateKey string + DisplayName string + DiscoveredFromEntityID string + MentionCount int + FirstSeenAt time.Time + LastSeenAt time.Time + EvidenceJSON string +} + +type UpdateCandidateSuggestionStatusParams struct { + ID int64 + Status string +} + +type scheduleJobContext struct { + DatasetID string + ScheduleJobID string + SourceID string +} + +type alertRuleContext struct { + DatasetID string + AlertRuleID string + SourceID string + Severity string +} + +type rowScanner interface { + Scan(dest ...any) error +} + +func StartScheduleRun(ctx context.Context, db *sql.DB, params StartScheduleRunParams) (ScheduleRun, error) { + if db == nil { + return ScheduleRun{}, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID := strings.TrimSpace(params.DatasetID) + if datasetID == "" { + return ScheduleRun{}, fmt.Errorf("dataset_id is required") + } + + scheduleJobID := strings.TrimSpace(params.ScheduleJobID) + if scheduleJobID == "" { + return ScheduleRun{}, fmt.Errorf("schedule_job_id is required") + } + + jobContext, err := lookupScheduleJobContext(ctx, db, datasetID, scheduleJobID) + if err != nil { + return ScheduleRun{}, err + } + + startedAt := normalizeTimestamp(params.StartedAt) + + result, err := db.ExecContext( + ctx, + `INSERT INTO schedule_runs ( + dataset_id, schedule_job_id, source_id, status, started_at, finished_at, + error_message, fetch_runs_started, fetch_runs_succeeded, fetch_runs_failed + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + jobContext.DatasetID, + jobContext.ScheduleJobID, + jobContext.SourceID, + ScheduleRunStatusRunning, + formatTimestamp(startedAt), + nil, + "", + 0, + 0, + 0, + ) + if err != nil { + return ScheduleRun{}, fmt.Errorf("insert schedule run: %w", err) + } + + id, err := result.LastInsertId() + if err != nil { + return ScheduleRun{}, fmt.Errorf("read schedule run id: %w", err) + } + + return ScheduleRun{ + ID: id, + DatasetID: jobContext.DatasetID, + ScheduleJobID: jobContext.ScheduleJobID, + SourceID: jobContext.SourceID, + Status: ScheduleRunStatusRunning, + StartedAt: startedAt, + FinishedAt: nil, + ErrorMessage: "", + FetchRunsStarted: 0, + FetchRunsSucceeded: 0, + FetchRunsFailed: 0, + }, nil +} + +func FinishScheduleRunSuccess(ctx context.Context, db *sql.DB, params FinalizeScheduleRunSuccessParams) error { + if err := validateTerminalScheduleRunCounts( + ScheduleRunStatusSucceeded, + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ); err != nil { + return err + } + + return finalizeScheduleRun( + ctx, + db, + params.ID, + ScheduleRunStatusSucceeded, + normalizeTimestamp(params.FinishedAt), + "", + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ) +} + +func FinishScheduleRunPartial(ctx context.Context, db *sql.DB, params FinalizeScheduleRunPartialParams) error { + if err := validateTerminalScheduleRunCounts( + ScheduleRunStatusPartial, + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ); err != nil { + return err + } + + return finalizeScheduleRun( + ctx, + db, + params.ID, + ScheduleRunStatusPartial, + normalizeTimestamp(params.FinishedAt), + strings.TrimSpace(params.ErrorMessage), + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ) +} + +func FinishScheduleRunFailure(ctx context.Context, db *sql.DB, params FinalizeScheduleRunFailureParams) error { + if err := validateTerminalScheduleRunCounts( + ScheduleRunStatusFailed, + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ); err != nil { + return err + } + + return finalizeScheduleRun( + ctx, + db, + params.ID, + ScheduleRunStatusFailed, + normalizeTimestamp(params.FinishedAt), + strings.TrimSpace(params.ErrorMessage), + params.FetchRunsStarted, + params.FetchRunsSucceeded, + params.FetchRunsFailed, + ) +} + +func LookupLatestTerminalScheduleRun( + ctx context.Context, + db *sql.DB, + datasetID string, + scheduleJobID string, +) (ScheduleRun, bool, error) { + if db == nil { + return ScheduleRun{}, false, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID = strings.TrimSpace(datasetID) + if datasetID == "" { + return ScheduleRun{}, false, fmt.Errorf("dataset_id is required") + } + + scheduleJobID = strings.TrimSpace(scheduleJobID) + if scheduleJobID == "" { + return ScheduleRun{}, false, fmt.Errorf("schedule_job_id is required") + } + + row, err := scanScheduleRun( + db.QueryRowContext( + ctx, + `SELECT + id, + dataset_id, + schedule_job_id, + source_id, + status, + started_at, + finished_at, + error_message, + fetch_runs_started, + fetch_runs_succeeded, + fetch_runs_failed + FROM schedule_runs + WHERE dataset_id = ? + AND schedule_job_id = ? + AND status IN (?, ?, ?) + ORDER BY started_at DESC, id DESC + LIMIT 1`, + datasetID, + scheduleJobID, + ScheduleRunStatusSucceeded, + ScheduleRunStatusPartial, + ScheduleRunStatusFailed, + ), + ) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return ScheduleRun{}, false, nil + } + return ScheduleRun{}, false, fmt.Errorf("lookup latest terminal schedule run: %w", err) + } + + return row, true, nil +} + +func InsertAlertEvent(ctx context.Context, db *sql.DB, params InsertAlertEventParams) (AlertEvent, error) { + if db == nil { + return AlertEvent{}, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID := strings.TrimSpace(params.DatasetID) + if datasetID == "" { + return AlertEvent{}, fmt.Errorf("dataset_id is required") + } + + alertRuleID := strings.TrimSpace(params.AlertRuleID) + if alertRuleID == "" { + return AlertEvent{}, fmt.Errorf("alert_rule_id is required") + } + + entityID := strings.TrimSpace(params.EntityID) + if entityID == "" { + return AlertEvent{}, fmt.Errorf("entity_id is required") + } + + summary := strings.TrimSpace(params.Summary) + if summary == "" { + return AlertEvent{}, fmt.Errorf("summary is required") + } + + payloadJSON, err := normalizedJSONText(params.PayloadJSON) + if err != nil { + return AlertEvent{}, fmt.Errorf("normalize payload_json: %w", err) + } + + observedAt := normalizeTimestamp(params.ObservedAt) + triggeredAt := normalizeTimestamp(params.TriggeredAt) + if observedAt.After(triggeredAt) { + return AlertEvent{}, fmt.Errorf("observed_at must be less than or equal to triggered_at") + } + + ruleContext, err := lookupAlertRuleContext(ctx, db, datasetID, alertRuleID) + if err != nil { + return AlertEvent{}, err + } + + result, err := db.ExecContext( + ctx, + `INSERT INTO alert_events ( + dataset_id, alert_rule_id, entity_id, source_id, severity, status, + summary, payload_json, observed_at, triggered_at, resolved_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + ruleContext.DatasetID, + ruleContext.AlertRuleID, + entityID, + ruleContext.SourceID, + ruleContext.Severity, + AlertEventStatusOpen, + summary, + payloadJSON, + formatTimestamp(observedAt), + formatTimestamp(triggeredAt), + nil, + ) + if err != nil { + return AlertEvent{}, fmt.Errorf("insert alert event: %w", err) + } + + id, err := result.LastInsertId() + if err != nil { + return AlertEvent{}, fmt.Errorf("read alert event id: %w", err) + } + + return AlertEvent{ + ID: id, + DatasetID: ruleContext.DatasetID, + AlertRuleID: ruleContext.AlertRuleID, + EntityID: entityID, + SourceID: ruleContext.SourceID, + Severity: ruleContext.Severity, + Status: AlertEventStatusOpen, + Summary: summary, + PayloadJSON: payloadJSON, + ObservedAt: observedAt, + TriggeredAt: triggeredAt, + ResolvedAt: nil, + }, nil +} + +func ResolveAlertEvent(ctx context.Context, db *sql.DB, params ResolveAlertEventParams) error { + if db == nil { + return fmt.Errorf("database handle is required") + } + if params.ID <= 0 { + return fmt.Errorf("alert_event_id must be greater than zero") + } + + ctx = normalizeContext(ctx) + + result, err := db.ExecContext( + ctx, + `UPDATE alert_events + SET status = ?, resolved_at = ? + WHERE id = ? AND status = ?`, + AlertEventStatusResolved, + formatTimestamp(normalizeTimestamp(params.ResolvedAt)), + params.ID, + AlertEventStatusOpen, + ) + if err != nil { + return fmt.Errorf("resolve alert event %d: %w", params.ID, err) + } + + rowsAffected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("read resolve alert event result for %d: %w", params.ID, err) + } + if rowsAffected != 1 { + return fmt.Errorf("no open alert event found for id %d", params.ID) + } + + return nil +} + +func ListOpenAlertEventsByRuleAndEntity( + ctx context.Context, + db *sql.DB, + datasetID string, + alertRuleID string, + entityID string, +) ([]AlertEvent, error) { + if db == nil { + return nil, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID = strings.TrimSpace(datasetID) + if datasetID == "" { + return nil, fmt.Errorf("dataset_id is required") + } + + alertRuleID = strings.TrimSpace(alertRuleID) + if alertRuleID == "" { + return nil, fmt.Errorf("alert_rule_id is required") + } + + entityID = strings.TrimSpace(entityID) + if entityID == "" { + return nil, fmt.Errorf("entity_id is required") + } + + rows, err := db.QueryContext( + ctx, + `SELECT + id, + dataset_id, + alert_rule_id, + entity_id, + source_id, + severity, + status, + summary, + payload_json, + observed_at, + triggered_at, + resolved_at + FROM alert_events + WHERE dataset_id = ? + AND alert_rule_id = ? + AND entity_id = ? + AND status = ? + ORDER BY triggered_at ASC, id ASC`, + datasetID, + alertRuleID, + entityID, + AlertEventStatusOpen, + ) + if err != nil { + return nil, fmt.Errorf("query open alert events: %w", err) + } + defer rows.Close() + + result := make([]AlertEvent, 0) + for rows.Next() { + row, err := scanAlertEvent(rows) + if err != nil { + return nil, fmt.Errorf("scan open alert event: %w", err) + } + result = append(result, row) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate open alert events: %w", err) + } + + return result, nil +} + +func UpsertCandidateSuggestion( + ctx context.Context, + db *sql.DB, + params UpsertCandidateSuggestionParams, +) (CandidateSuggestion, error) { + if db == nil { + return CandidateSuggestion{}, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID := strings.TrimSpace(params.DatasetID) + if datasetID == "" { + return CandidateSuggestion{}, fmt.Errorf("dataset_id is required") + } + + alertRuleID := strings.TrimSpace(params.AlertRuleID) + if alertRuleID == "" { + return CandidateSuggestion{}, fmt.Errorf("alert_rule_id is required") + } + + candidateKey := strings.TrimSpace(params.CandidateKey) + if candidateKey == "" { + return CandidateSuggestion{}, fmt.Errorf("candidate_key is required") + } + + displayName := strings.TrimSpace(params.DisplayName) + if displayName == "" { + return CandidateSuggestion{}, fmt.Errorf("display_name is required") + } + + discoveredFromEntityID := strings.TrimSpace(params.DiscoveredFromEntityID) + if discoveredFromEntityID == "" { + return CandidateSuggestion{}, fmt.Errorf("discovered_from_entity_id is required") + } + + if params.MentionCount <= 0 { + return CandidateSuggestion{}, fmt.Errorf("mention_count must be greater than zero") + } + + firstSeenAt := normalizeTimestamp(params.FirstSeenAt) + lastSeenAt := normalizeTimestamp(params.LastSeenAt) + if firstSeenAt.After(lastSeenAt) { + return CandidateSuggestion{}, fmt.Errorf("first_seen_at must be less than or equal to last_seen_at") + } + + evidenceJSON, err := normalizedJSONText(params.EvidenceJSON) + if err != nil { + return CandidateSuggestion{}, fmt.Errorf("normalize evidence_json: %w", err) + } + + ruleContext, err := lookupAlertRuleContext(ctx, db, datasetID, alertRuleID) + if err != nil { + return CandidateSuggestion{}, err + } + + if _, err := db.ExecContext( + ctx, + `INSERT INTO candidate_suggestions ( + dataset_id, + alert_rule_id, + candidate_key, + display_name, + discovered_from_entity_id, + source_id, + mention_count, + first_seen_at, + last_seen_at, + status, + evidence_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(dataset_id, discovered_from_entity_id, source_id, candidate_key) DO UPDATE SET + mention_count = CASE + WHEN excluded.mention_count > candidate_suggestions.mention_count + THEN excluded.mention_count + ELSE candidate_suggestions.mention_count + END, + last_seen_at = CASE + WHEN excluded.last_seen_at > candidate_suggestions.last_seen_at + THEN excluded.last_seen_at + ELSE candidate_suggestions.last_seen_at + END, + evidence_json = CASE + WHEN excluded.mention_count > candidate_suggestions.mention_count + OR excluded.last_seen_at > candidate_suggestions.last_seen_at + THEN excluded.evidence_json + ELSE candidate_suggestions.evidence_json + END`, + ruleContext.DatasetID, + ruleContext.AlertRuleID, + candidateKey, + displayName, + discoveredFromEntityID, + ruleContext.SourceID, + params.MentionCount, + formatTimestamp(firstSeenAt), + formatTimestamp(lastSeenAt), + CandidateSuggestionStatusPendingReview, + evidenceJSON, + ); err != nil { + return CandidateSuggestion{}, fmt.Errorf("upsert candidate suggestion: %w", err) + } + + suggestion, found, err := LookupCandidateSuggestion( + ctx, + db, + ruleContext.DatasetID, + discoveredFromEntityID, + ruleContext.SourceID, + candidateKey, + ) + if err != nil { + return CandidateSuggestion{}, err + } + if !found { + return CandidateSuggestion{}, fmt.Errorf( + "candidate suggestion %q for entity %q not found after upsert", + candidateKey, + discoveredFromEntityID, + ) + } + + return suggestion, nil +} + +func UpdateCandidateSuggestionStatus( + ctx context.Context, + db *sql.DB, + params UpdateCandidateSuggestionStatusParams, +) error { + if db == nil { + return fmt.Errorf("database handle is required") + } + if params.ID <= 0 { + return fmt.Errorf("candidate_suggestion_id must be greater than zero") + } + + ctx = normalizeContext(ctx) + + status, err := normalizeCandidateSuggestionStatus(params.Status) + if err != nil { + return err + } + + result, err := db.ExecContext( + ctx, + `UPDATE candidate_suggestions + SET status = ? + WHERE id = ?`, + status, + params.ID, + ) + if err != nil { + return fmt.Errorf("update candidate suggestion status for %d: %w", params.ID, err) + } + + rowsAffected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("read candidate suggestion status update result for %d: %w", params.ID, err) + } + if rowsAffected != 1 { + return fmt.Errorf("candidate suggestion %d not found", params.ID) + } + + return nil +} + +func LookupCandidateSuggestion( + ctx context.Context, + db *sql.DB, + datasetID string, + discoveredFromEntityID string, + sourceID string, + candidateKey string, +) (CandidateSuggestion, bool, error) { + if db == nil { + return CandidateSuggestion{}, false, fmt.Errorf("database handle is required") + } + + ctx = normalizeContext(ctx) + + datasetID = strings.TrimSpace(datasetID) + if datasetID == "" { + return CandidateSuggestion{}, false, fmt.Errorf("dataset_id is required") + } + + discoveredFromEntityID = strings.TrimSpace(discoveredFromEntityID) + if discoveredFromEntityID == "" { + return CandidateSuggestion{}, false, fmt.Errorf("discovered_from_entity_id is required") + } + + sourceID = strings.TrimSpace(sourceID) + if sourceID == "" { + return CandidateSuggestion{}, false, fmt.Errorf("source_id is required") + } + + candidateKey = strings.TrimSpace(candidateKey) + if candidateKey == "" { + return CandidateSuggestion{}, false, fmt.Errorf("candidate_key is required") + } + + row, err := scanCandidateSuggestion( + db.QueryRowContext( + ctx, + `SELECT + id, + dataset_id, + alert_rule_id, + candidate_key, + display_name, + discovered_from_entity_id, + source_id, + mention_count, + first_seen_at, + last_seen_at, + status, + evidence_json + FROM candidate_suggestions + WHERE dataset_id = ? + AND discovered_from_entity_id = ? + AND source_id = ? + AND candidate_key = ?`, + datasetID, + discoveredFromEntityID, + sourceID, + candidateKey, + ), + ) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return CandidateSuggestion{}, false, nil + } + return CandidateSuggestion{}, false, fmt.Errorf("lookup candidate suggestion: %w", err) + } + + return row, true, nil +} + +func finalizeScheduleRun( + ctx context.Context, + db *sql.DB, + scheduleRunID int64, + status string, + finishedAt time.Time, + errorMessage string, + fetchRunsStarted int, + fetchRunsSucceeded int, + fetchRunsFailed int, +) error { + if db == nil { + return fmt.Errorf("database handle is required") + } + if scheduleRunID <= 0 { + return fmt.Errorf("schedule_run_id must be greater than zero") + } + + ctx = normalizeContext(ctx) + + result, err := db.ExecContext( + ctx, + `UPDATE schedule_runs + SET status = ?, + finished_at = ?, + error_message = ?, + fetch_runs_started = ?, + fetch_runs_succeeded = ?, + fetch_runs_failed = ? + WHERE id = ? AND status = ?`, + status, + formatTimestamp(finishedAt), + strings.TrimSpace(errorMessage), + fetchRunsStarted, + fetchRunsSucceeded, + fetchRunsFailed, + scheduleRunID, + ScheduleRunStatusRunning, + ) + if err != nil { + return fmt.Errorf("finalize schedule run %d: %w", scheduleRunID, err) + } + + rowsAffected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("read finalize schedule run result for %d: %w", scheduleRunID, err) + } + if rowsAffected != 1 { + return fmt.Errorf("no running schedule run found for id %d", scheduleRunID) + } + + return nil +} + +func validateTerminalScheduleRunCounts(status string, started, succeeded, failed int) error { + switch { + case started < 0: + return fmt.Errorf("fetch_runs_started must be greater than or equal to zero") + case succeeded < 0: + return fmt.Errorf("fetch_runs_succeeded must be greater than or equal to zero") + case failed < 0: + return fmt.Errorf("fetch_runs_failed must be greater than or equal to zero") + case succeeded+failed != started: + return fmt.Errorf("fetch run counts must satisfy fetch_runs_started = fetch_runs_succeeded + fetch_runs_failed") + } + + switch status { + case ScheduleRunStatusSucceeded: + if failed != 0 { + return fmt.Errorf("succeeded schedule runs cannot have failed fetch runs") + } + case ScheduleRunStatusPartial: + if succeeded == 0 || failed == 0 { + return fmt.Errorf("partial schedule runs require both succeeded and failed fetch runs") + } + case ScheduleRunStatusFailed: + if succeeded != 0 { + return fmt.Errorf("failed schedule runs cannot have successful fetch runs") + } + default: + return fmt.Errorf("unsupported schedule run status %q", status) + } + + return nil +} + +func lookupScheduleJobContext(ctx context.Context, db *sql.DB, datasetID, scheduleJobID string) (scheduleJobContext, error) { + var context scheduleJobContext + + err := db.QueryRowContext( + ctx, + `SELECT dataset_id, id, source_id + FROM schedule_jobs + WHERE dataset_id = ? AND id = ?`, + datasetID, + scheduleJobID, + ).Scan( + &context.DatasetID, + &context.ScheduleJobID, + &context.SourceID, + ) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return scheduleJobContext{}, fmt.Errorf("schedule job %q not found in dataset %q", scheduleJobID, datasetID) + } + return scheduleJobContext{}, fmt.Errorf("lookup schedule job context: %w", err) + } + + return context, nil +} + +func lookupAlertRuleContext(ctx context.Context, db *sql.DB, datasetID, alertRuleID string) (alertRuleContext, error) { + var context alertRuleContext + + err := db.QueryRowContext( + ctx, + `SELECT dataset_id, id, source_id, severity + FROM alert_rules + WHERE dataset_id = ? AND id = ?`, + datasetID, + alertRuleID, + ).Scan( + &context.DatasetID, + &context.AlertRuleID, + &context.SourceID, + &context.Severity, + ) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return alertRuleContext{}, fmt.Errorf("alert rule %q not found in dataset %q", alertRuleID, datasetID) + } + return alertRuleContext{}, fmt.Errorf("lookup alert rule context: %w", err) + } + + return context, nil +} + +func scanScheduleRun(scanner rowScanner) (ScheduleRun, error) { + var row ScheduleRun + var startedAt string + var finishedAt sql.NullString + + if err := scanner.Scan( + &row.ID, + &row.DatasetID, + &row.ScheduleJobID, + &row.SourceID, + &row.Status, + &startedAt, + &finishedAt, + &row.ErrorMessage, + &row.FetchRunsStarted, + &row.FetchRunsSucceeded, + &row.FetchRunsFailed, + ); err != nil { + return ScheduleRun{}, err + } + + parsedStartedAt, err := parseStoredTimestamp("schedule_runs.started_at", startedAt) + if err != nil { + return ScheduleRun{}, err + } + parsedFinishedAt, err := parseStoredOptionalTimestamp("schedule_runs.finished_at", finishedAt) + if err != nil { + return ScheduleRun{}, err + } + + row.StartedAt = parsedStartedAt + row.FinishedAt = parsedFinishedAt + + return row, nil +} + +func scanAlertEvent(scanner rowScanner) (AlertEvent, error) { + var row AlertEvent + var observedAt string + var triggeredAt string + var resolvedAt sql.NullString + + if err := scanner.Scan( + &row.ID, + &row.DatasetID, + &row.AlertRuleID, + &row.EntityID, + &row.SourceID, + &row.Severity, + &row.Status, + &row.Summary, + &row.PayloadJSON, + &observedAt, + &triggeredAt, + &resolvedAt, + ); err != nil { + return AlertEvent{}, err + } + + parsedObservedAt, err := parseStoredTimestamp("alert_events.observed_at", observedAt) + if err != nil { + return AlertEvent{}, err + } + parsedTriggeredAt, err := parseStoredTimestamp("alert_events.triggered_at", triggeredAt) + if err != nil { + return AlertEvent{}, err + } + parsedResolvedAt, err := parseStoredOptionalTimestamp("alert_events.resolved_at", resolvedAt) + if err != nil { + return AlertEvent{}, err + } + + row.ObservedAt = parsedObservedAt + row.TriggeredAt = parsedTriggeredAt + row.ResolvedAt = parsedResolvedAt + + return row, nil +} + +func scanCandidateSuggestion(scanner rowScanner) (CandidateSuggestion, error) { + var row CandidateSuggestion + var firstSeenAt string + var lastSeenAt string + + if err := scanner.Scan( + &row.ID, + &row.DatasetID, + &row.AlertRuleID, + &row.CandidateKey, + &row.DisplayName, + &row.DiscoveredFromEntityID, + &row.SourceID, + &row.MentionCount, + &firstSeenAt, + &lastSeenAt, + &row.Status, + &row.EvidenceJSON, + ); err != nil { + return CandidateSuggestion{}, err + } + + parsedFirstSeenAt, err := parseStoredTimestamp("candidate_suggestions.first_seen_at", firstSeenAt) + if err != nil { + return CandidateSuggestion{}, err + } + parsedLastSeenAt, err := parseStoredTimestamp("candidate_suggestions.last_seen_at", lastSeenAt) + if err != nil { + return CandidateSuggestion{}, err + } + + row.FirstSeenAt = parsedFirstSeenAt + row.LastSeenAt = parsedLastSeenAt + + return row, nil +} + +func parseStoredTimestamp(fieldName, value string) (time.Time, error) { + parsed, err := time.Parse(time.RFC3339, strings.TrimSpace(value)) + if err != nil { + return time.Time{}, fmt.Errorf("%s parse error: %w", fieldName, err) + } + return parsed.UTC(), nil +} + +func parseStoredOptionalTimestamp(fieldName string, value sql.NullString) (*time.Time, error) { + if !value.Valid || strings.TrimSpace(value.String) == "" { + return nil, nil + } + + parsed, err := parseStoredTimestamp(fieldName, value.String) + if err != nil { + return nil, err + } + + return &parsed, nil +} + +func normalizeCandidateSuggestionStatus(value string) (string, error) { + status := strings.TrimSpace(strings.ToLower(value)) + + switch status { + case CandidateSuggestionStatusPendingReview, + CandidateSuggestionStatusAccepted, + CandidateSuggestionStatusDismissed: + return status, nil + default: + return "", fmt.Errorf( + "unsupported candidate suggestion status %q (supported: %q, %q, %q)", + value, + CandidateSuggestionStatusPendingReview, + CandidateSuggestionStatusAccepted, + CandidateSuggestionStatusDismissed, + ) + } +} diff --git a/internal/storage/processing_test.go b/internal/storage/processing_test.go new file mode 100644 index 0000000..b0db1df --- /dev/null +++ b/internal/storage/processing_test.go @@ -0,0 +1,907 @@ +package storage + +import ( + "context" + "database/sql" + "path/filepath" + "testing" + "time" + + "github.com/CrisSTEM/signalscope/internal/config" + repomigrations "github.com/CrisSTEM/signalscope/migrations" +) + +func TestBootstrapCreatesProcessingTables(t *testing.T) { + t.Parallel() + + ctx := context.Background() + + db, err := Open(filepath.Join(t.TempDir(), "signalscope.db")) + if err != nil { + t.Fatalf("Open() error = %v", err) + } + defer db.Close() + + if err := Bootstrap(ctx, db); err != nil { + t.Fatalf("Bootstrap() error = %v", err) + } + + for _, table := range []string{"schedule_runs", "alert_events", "candidate_suggestions"} { + exists, err := processingTableExists(ctx, db, table) + if err != nil { + t.Fatalf("processingTableExists(%q) error = %v", table, err) + } + if !exists { + t.Fatalf("processing table %q does not exist after Bootstrap()", table) + } + } + + version3Count, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM schema_migrations WHERE version = 3`, + ) + if err != nil { + t.Fatalf("processingQueryCount(schema_migrations version 3) error = %v", err) + } + if version3Count != 1 { + t.Fatalf("version3Count = %d, want %d", version3Count, 1) + } +} + +func TestBootstrapUpgradesV030DatabaseWithoutClearingIngestionHistory(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-upgrade-history" + + db, err := Open(filepath.Join(t.TempDir(), "signalscope.db")) + if err != nil { + t.Fatalf("Open() error = %v", err) + } + defer db.Close() + + if err := ensureMigrationTable(ctx, db); err != nil { + t.Fatalf("ensureMigrationTable() error = %v", err) + } + + processingApplyEmbeddedMigration(t, ctx, db, "001_initial.sql") + processingApplyEmbeddedMigration(t, ctx, db, "002_ingestion.sql") + + if err := SeedPack(ctx, db, processingFixturePack(datasetID)); err != nil { + t.Fatalf("SeedPack() error = %v", err) + } + + startedAt := time.Date(2026, 4, 25, 10, 0, 0, 0, time.UTC) + run, err := StartFetchRun(ctx, db, StartFetchRunParams{ + DatasetID: datasetID, + BindingID: "news-binding-1", + StartedAt: startedAt, + }) + if err != nil { + t.Fatalf("StartFetchRun() error = %v", err) + } + + metricValue := 2.0 + metricsInserted, contentItemsInserted, err := InsertMetricSnapshotsAndContentItems( + ctx, + db, + run.ID, + []MetricSnapshotInput{ + { + MetricKey: "news.article_count", + MetricValueNum: &metricValue, + Unit: "count", + WindowKey: "point_in_time", + CapturedAt: startedAt, + }, + }, + []ContentItemInput{ + { + ItemType: "news_article", + Title: "Upgrade path fixture", + Summary: "preserve existing ingestion history", + URL: "https://example.com/news/upgrade-path", + DiscoveredAt: startedAt, + MetadataJSON: `{"fixture":"upgrade-path"}`, + }, + }, + ) + if err != nil { + t.Fatalf("InsertMetricSnapshotsAndContentItems() error = %v", err) + } + if metricsInserted != 1 { + t.Fatalf("metricsInserted = %d, want %d", metricsInserted, 1) + } + if contentItemsInserted != 1 { + t.Fatalf("contentItemsInserted = %d, want %d", contentItemsInserted, 1) + } + + if err := FinishFetchRunSuccess(ctx, db, FinalizeFetchRunSuccessParams{ + ID: run.ID, + FinishedAt: startedAt.Add(time.Second), + RecordsWritten: 1, + MetricsWritten: 1, + ContentItemsWritten: 1, + }); err != nil { + t.Fatalf("FinishFetchRunSuccess() error = %v", err) + } + + if err := Bootstrap(ctx, db); err != nil { + t.Fatalf("Bootstrap() upgrade error = %v", err) + } + + for _, table := range []string{"schedule_runs", "alert_events", "candidate_suggestions"} { + exists, err := processingTableExists(ctx, db, table) + if err != nil { + t.Fatalf("processingTableExists(%q) error = %v", table, err) + } + if !exists { + t.Fatalf("processing table %q does not exist after upgrade Bootstrap()", table) + } + } + + version3Count, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM schema_migrations WHERE version = 3`, + ) + if err != nil { + t.Fatalf("processingQueryCount(version 3) error = %v", err) + } + if version3Count != 1 { + t.Fatalf("version3Count = %d, want %d", version3Count, 1) + } + + fetchRunCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM fetch_runs WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(fetch_runs) error = %v", err) + } + if fetchRunCount != 1 { + t.Fatalf("fetchRunCount = %d, want %d", fetchRunCount, 1) + } + + metricSnapshotCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM metric_snapshots WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(metric_snapshots) error = %v", err) + } + if metricSnapshotCount != 1 { + t.Fatalf("metricSnapshotCount = %d, want %d", metricSnapshotCount, 1) + } + + contentItemCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM content_items WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(content_items) error = %v", err) + } + if contentItemCount != 1 { + t.Fatalf("contentItemCount = %d, want %d", contentItemCount, 1) + } +} + +func TestScheduleRunLifecycleAndLookupLatestTerminal(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-schedule-run" + + db := processingOpenSeededDB(t, datasetID) + defer db.Close() + + startedAtOne := time.Date(2026, 4, 25, 11, 0, 0, 0, time.UTC) + runOne, err := StartScheduleRun(ctx, db, StartScheduleRunParams{ + DatasetID: datasetID, + ScheduleJobID: "job-github", + StartedAt: startedAtOne, + }) + if err != nil { + t.Fatalf("StartScheduleRun(runOne) error = %v", err) + } + if runOne.SourceID != "github" { + t.Fatalf("runOne.SourceID = %q, want %q", runOne.SourceID, "github") + } + if runOne.Status != ScheduleRunStatusRunning { + t.Fatalf("runOne.Status = %q, want %q", runOne.Status, ScheduleRunStatusRunning) + } + + if err := FinishScheduleRunSuccess(ctx, db, FinalizeScheduleRunSuccessParams{ + ID: runOne.ID, + FinishedAt: startedAtOne.Add(2 * time.Second), + FetchRunsStarted: 1, + FetchRunsSucceeded: 1, + FetchRunsFailed: 0, + }); err != nil { + t.Fatalf("FinishScheduleRunSuccess(runOne) error = %v", err) + } + + startedAtTwo := time.Date(2026, 4, 25, 12, 0, 0, 0, time.UTC) + runTwo, err := StartScheduleRun(ctx, db, StartScheduleRunParams{ + DatasetID: datasetID, + ScheduleJobID: "job-github", + StartedAt: startedAtTwo, + }) + if err != nil { + t.Fatalf("StartScheduleRun(runTwo) error = %v", err) + } + + latestBeforeFinalizingRunTwo, found, err := LookupLatestTerminalScheduleRun(ctx, db, datasetID, "job-github") + if err != nil { + t.Fatalf("LookupLatestTerminalScheduleRun(before finalizing runTwo) error = %v", err) + } + if !found { + t.Fatal("LookupLatestTerminalScheduleRun(before finalizing runTwo) found = false, want true") + } + if latestBeforeFinalizingRunTwo.ID != runOne.ID { + t.Fatalf( + "latestBeforeFinalizingRunTwo.ID = %d, want %d", + latestBeforeFinalizingRunTwo.ID, + runOne.ID, + ) + } + + if err := FinishScheduleRunPartial(ctx, db, FinalizeScheduleRunPartialParams{ + ID: runTwo.ID, + FinishedAt: startedAtTwo.Add(3 * time.Second), + ErrorMessage: "one downstream binding failed", + FetchRunsStarted: 2, + FetchRunsSucceeded: 1, + FetchRunsFailed: 1, + }); err != nil { + t.Fatalf("FinishScheduleRunPartial(runTwo) error = %v", err) + } + + latestAfterFinalizingRunTwo, found, err := LookupLatestTerminalScheduleRun(ctx, db, datasetID, "job-github") + if err != nil { + t.Fatalf("LookupLatestTerminalScheduleRun(after finalizing runTwo) error = %v", err) + } + if !found { + t.Fatal("LookupLatestTerminalScheduleRun(after finalizing runTwo) found = false, want true") + } + if latestAfterFinalizingRunTwo.ID != runTwo.ID { + t.Fatalf("latestAfterFinalizingRunTwo.ID = %d, want %d", latestAfterFinalizingRunTwo.ID, runTwo.ID) + } + if latestAfterFinalizingRunTwo.Status != ScheduleRunStatusPartial { + t.Fatalf( + "latestAfterFinalizingRunTwo.Status = %q, want %q", + latestAfterFinalizingRunTwo.Status, + ScheduleRunStatusPartial, + ) + } + if latestAfterFinalizingRunTwo.ErrorMessage != "one downstream binding failed" { + t.Fatalf( + "latestAfterFinalizingRunTwo.ErrorMessage = %q, want %q", + latestAfterFinalizingRunTwo.ErrorMessage, + "one downstream binding failed", + ) + } + + storedRunTwo, err := processingReadScheduleRunByID(ctx, db, runTwo.ID) + if err != nil { + t.Fatalf("processingReadScheduleRunByID(runTwo) error = %v", err) + } + if storedRunTwo.FetchRunsStarted != 2 { + t.Fatalf("storedRunTwo.FetchRunsStarted = %d, want %d", storedRunTwo.FetchRunsStarted, 2) + } + if storedRunTwo.FetchRunsSucceeded != 1 { + t.Fatalf("storedRunTwo.FetchRunsSucceeded = %d, want %d", storedRunTwo.FetchRunsSucceeded, 1) + } + if storedRunTwo.FetchRunsFailed != 1 { + t.Fatalf("storedRunTwo.FetchRunsFailed = %d, want %d", storedRunTwo.FetchRunsFailed, 1) + } +} + +func TestScheduleRunFailureFinalization(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-schedule-failure" + + db := processingOpenSeededDB(t, datasetID) + defer db.Close() + + startedAt := time.Date(2026, 4, 25, 13, 0, 0, 0, time.UTC) + run, err := StartScheduleRun(ctx, db, StartScheduleRunParams{ + DatasetID: datasetID, + ScheduleJobID: "job-github", + StartedAt: startedAt, + }) + if err != nil { + t.Fatalf("StartScheduleRun() error = %v", err) + } + + if err := FinishScheduleRunFailure(ctx, db, FinalizeScheduleRunFailureParams{ + ID: run.ID, + FinishedAt: startedAt.Add(time.Second), + ErrorMessage: "scheduler timed out before downstream fetch work started", + FetchRunsStarted: 0, + FetchRunsSucceeded: 0, + FetchRunsFailed: 0, + }); err != nil { + t.Fatalf("FinishScheduleRunFailure() error = %v", err) + } + + storedRun, err := processingReadScheduleRunByID(ctx, db, run.ID) + if err != nil { + t.Fatalf("processingReadScheduleRunByID() error = %v", err) + } + if storedRun.Status != ScheduleRunStatusFailed { + t.Fatalf("storedRun.Status = %q, want %q", storedRun.Status, ScheduleRunStatusFailed) + } + if storedRun.ErrorMessage != "scheduler timed out before downstream fetch work started" { + t.Fatalf( + "storedRun.ErrorMessage = %q, want %q", + storedRun.ErrorMessage, + "scheduler timed out before downstream fetch work started", + ) + } +} + +func TestInsertAlertEventCopiesRuleContextAndResolve(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-alert-events" + + db := processingOpenSeededDB(t, datasetID) + defer db.Close() + + observedAt := time.Date(2026, 4, 25, 14, 0, 0, 0, time.UTC) + triggeredAt := observedAt.Add(5 * time.Minute) + + event, err := InsertAlertEvent(ctx, db, InsertAlertEventParams{ + DatasetID: datasetID, + AlertRuleID: "rule-release", + EntityID: "org-1", + Summary: "Release published after a silence window", + PayloadJSON: `{"event_key":"release.published","release_tag":"v1.0.0"}`, + ObservedAt: observedAt, + TriggeredAt: triggeredAt, + }) + if err != nil { + t.Fatalf("InsertAlertEvent() error = %v", err) + } + + if event.SourceID != "github" { + t.Fatalf("event.SourceID = %q, want %q", event.SourceID, "github") + } + if event.Severity != "medium" { + t.Fatalf("event.Severity = %q, want %q", event.Severity, "medium") + } + if event.Status != AlertEventStatusOpen { + t.Fatalf("event.Status = %q, want %q", event.Status, AlertEventStatusOpen) + } + + openEvents, err := ListOpenAlertEventsByRuleAndEntity(ctx, db, datasetID, "rule-release", "org-1") + if err != nil { + t.Fatalf("ListOpenAlertEventsByRuleAndEntity() error = %v", err) + } + if len(openEvents) != 1 { + t.Fatalf("len(openEvents) = %d, want %d", len(openEvents), 1) + } + if openEvents[0].ID != event.ID { + t.Fatalf("openEvents[0].ID = %d, want %d", openEvents[0].ID, event.ID) + } + + resolvedAt := triggeredAt.Add(10 * time.Minute) + if err := ResolveAlertEvent(ctx, db, ResolveAlertEventParams{ + ID: event.ID, + ResolvedAt: resolvedAt, + }); err != nil { + t.Fatalf("ResolveAlertEvent() error = %v", err) + } + + openEventsAfterResolve, err := ListOpenAlertEventsByRuleAndEntity(ctx, db, datasetID, "rule-release", "org-1") + if err != nil { + t.Fatalf("ListOpenAlertEventsByRuleAndEntity(after resolve) error = %v", err) + } + if len(openEventsAfterResolve) != 0 { + t.Fatalf("len(openEventsAfterResolve) = %d, want %d", len(openEventsAfterResolve), 0) + } + + storedEvent, err := processingReadAlertEventByID(ctx, db, event.ID) + if err != nil { + t.Fatalf("processingReadAlertEventByID() error = %v", err) + } + if storedEvent.Status != AlertEventStatusResolved { + t.Fatalf("storedEvent.Status = %q, want %q", storedEvent.Status, AlertEventStatusResolved) + } + if storedEvent.ResolvedAt == nil { + t.Fatal("storedEvent.ResolvedAt = nil, want value") + } + if !storedEvent.ResolvedAt.Equal(resolvedAt) { + t.Fatalf("storedEvent.ResolvedAt = %v, want %v", storedEvent.ResolvedAt, resolvedAt) + } +} + +func TestUpsertCandidateSuggestionPreservesCanonicalRowAcrossReruns(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-candidate-suggestions" + + db := processingOpenSeededDB(t, datasetID) + defer db.Close() + + firstSeenAt := time.Date(2026, 4, 25, 15, 0, 0, 0, time.UTC) + + firstSuggestion, err := UpsertCandidateSuggestion(ctx, db, UpsertCandidateSuggestionParams{ + DatasetID: datasetID, + AlertRuleID: "rule-candidate", + CandidateKey: "candidate:acme-wallet", + DisplayName: "Acme Wallet", + DiscoveredFromEntityID: "org-1", + MentionCount: 2, + FirstSeenAt: firstSeenAt, + LastSeenAt: firstSeenAt, + EvidenceJSON: `{"items":["news-1","news-2"]}`, + }) + if err != nil { + t.Fatalf("UpsertCandidateSuggestion(first) error = %v", err) + } + + if firstSuggestion.SourceID != "news-rss" { + t.Fatalf("firstSuggestion.SourceID = %q, want %q", firstSuggestion.SourceID, "news-rss") + } + if firstSuggestion.Status != CandidateSuggestionStatusPendingReview { + t.Fatalf( + "firstSuggestion.Status = %q, want %q", + firstSuggestion.Status, + CandidateSuggestionStatusPendingReview, + ) + } + + sameEvidenceSuggestion, err := UpsertCandidateSuggestion(ctx, db, UpsertCandidateSuggestionParams{ + DatasetID: datasetID, + AlertRuleID: "rule-candidate", + CandidateKey: "candidate:acme-wallet", + DisplayName: "Acme Wallet Updated Name Should Not Replace Canonical Row", + DiscoveredFromEntityID: "org-1", + MentionCount: 2, + FirstSeenAt: firstSeenAt, + LastSeenAt: firstSeenAt, + EvidenceJSON: `{"items":["news-1","news-2"],"rerun":true}`, + }) + if err != nil { + t.Fatalf("UpsertCandidateSuggestion(same evidence) error = %v", err) + } + + if sameEvidenceSuggestion.ID != firstSuggestion.ID { + t.Fatalf("sameEvidenceSuggestion.ID = %d, want %d", sameEvidenceSuggestion.ID, firstSuggestion.ID) + } + if sameEvidenceSuggestion.MentionCount != 2 { + t.Fatalf("sameEvidenceSuggestion.MentionCount = %d, want %d", sameEvidenceSuggestion.MentionCount, 2) + } + if sameEvidenceSuggestion.LastSeenAt != firstSeenAt { + t.Fatalf("sameEvidenceSuggestion.LastSeenAt = %v, want %v", sameEvidenceSuggestion.LastSeenAt, firstSeenAt) + } + if sameEvidenceSuggestion.EvidenceJSON != `{"items":["news-1","news-2"]}` { + t.Fatalf( + "sameEvidenceSuggestion.EvidenceJSON = %q, want %q", + sameEvidenceSuggestion.EvidenceJSON, + `{"items":["news-1","news-2"]}`, + ) + } + + if err := UpdateCandidateSuggestionStatus(ctx, db, UpdateCandidateSuggestionStatusParams{ + ID: firstSuggestion.ID, + Status: CandidateSuggestionStatusAccepted, + }); err != nil { + t.Fatalf("UpdateCandidateSuggestionStatus() error = %v", err) + } + + secondSeenAt := firstSeenAt.Add(24 * time.Hour) + updatedSuggestion, err := UpsertCandidateSuggestion(ctx, db, UpsertCandidateSuggestionParams{ + DatasetID: datasetID, + AlertRuleID: "rule-candidate", + CandidateKey: "candidate:acme-wallet", + DisplayName: "Acme Wallet", + DiscoveredFromEntityID: "org-1", + MentionCount: 3, + FirstSeenAt: firstSeenAt, + LastSeenAt: secondSeenAt, + EvidenceJSON: `{"items":["news-1","news-2","news-3"]}`, + }) + if err != nil { + t.Fatalf("UpsertCandidateSuggestion(updated evidence) error = %v", err) + } + + if updatedSuggestion.ID != firstSuggestion.ID { + t.Fatalf("updatedSuggestion.ID = %d, want %d", updatedSuggestion.ID, firstSuggestion.ID) + } + if updatedSuggestion.Status != CandidateSuggestionStatusAccepted { + t.Fatalf( + "updatedSuggestion.Status = %q, want %q", + updatedSuggestion.Status, + CandidateSuggestionStatusAccepted, + ) + } + if updatedSuggestion.MentionCount != 3 { + t.Fatalf("updatedSuggestion.MentionCount = %d, want %d", updatedSuggestion.MentionCount, 3) + } + if !updatedSuggestion.FirstSeenAt.Equal(firstSeenAt) { + t.Fatalf("updatedSuggestion.FirstSeenAt = %v, want %v", updatedSuggestion.FirstSeenAt, firstSeenAt) + } + if !updatedSuggestion.LastSeenAt.Equal(secondSeenAt) { + t.Fatalf("updatedSuggestion.LastSeenAt = %v, want %v", updatedSuggestion.LastSeenAt, secondSeenAt) + } + if updatedSuggestion.EvidenceJSON != `{"items":["news-1","news-2","news-3"]}` { + t.Fatalf( + "updatedSuggestion.EvidenceJSON = %q, want %q", + updatedSuggestion.EvidenceJSON, + `{"items":["news-1","news-2","news-3"]}`, + ) + } + + lookedUp, found, err := LookupCandidateSuggestion( + ctx, + db, + datasetID, + "org-1", + "news-rss", + "candidate:acme-wallet", + ) + if err != nil { + t.Fatalf("LookupCandidateSuggestion() error = %v", err) + } + if !found { + t.Fatal("LookupCandidateSuggestion() found = false, want true") + } + if lookedUp.ID != firstSuggestion.ID { + t.Fatalf("lookedUp.ID = %d, want %d", lookedUp.ID, firstSuggestion.ID) + } +} + +func TestSyncPackPreservesProcessingHistory(t *testing.T) { + t.Parallel() + + ctx := context.Background() + datasetID := "processing-sync-preserves-history" + + db := processingOpenSeededDB(t, datasetID) + defer db.Close() + + startedAt := time.Date(2026, 4, 25, 16, 0, 0, 0, time.UTC) + scheduleRun, err := StartScheduleRun(ctx, db, StartScheduleRunParams{ + DatasetID: datasetID, + ScheduleJobID: "job-github", + StartedAt: startedAt, + }) + if err != nil { + t.Fatalf("StartScheduleRun() error = %v", err) + } + + if err := FinishScheduleRunSuccess(ctx, db, FinalizeScheduleRunSuccessParams{ + ID: scheduleRun.ID, + FinishedAt: startedAt.Add(time.Second), + FetchRunsStarted: 1, + FetchRunsSucceeded: 1, + FetchRunsFailed: 0, + }); err != nil { + t.Fatalf("FinishScheduleRunSuccess() error = %v", err) + } + + if _, err := InsertAlertEvent(ctx, db, InsertAlertEventParams{ + DatasetID: datasetID, + AlertRuleID: "rule-release", + EntityID: "org-1", + Summary: "Preserve alert history through SyncPack", + PayloadJSON: `{"fixture":"sync-pack"}`, + ObservedAt: startedAt, + TriggeredAt: startedAt.Add(2 * time.Minute), + }); err != nil { + t.Fatalf("InsertAlertEvent() error = %v", err) + } + + if _, err := UpsertCandidateSuggestion(ctx, db, UpsertCandidateSuggestionParams{ + DatasetID: datasetID, + AlertRuleID: "rule-candidate", + CandidateKey: "candidate:sync-pack", + DisplayName: "Sync Pack Candidate", + DiscoveredFromEntityID: "org-1", + MentionCount: 2, + FirstSeenAt: startedAt, + LastSeenAt: startedAt, + EvidenceJSON: `{"fixture":"sync-pack"}`, + }); err != nil { + t.Fatalf("UpsertCandidateSuggestion() error = %v", err) + } + + updatedPack := processingFixturePack(datasetID) + updatedPack.Sources.Bindings[1].Notes = "updated news binding note" + updatedPack.Entities.Entities = append(updatedPack.Entities.Entities, config.Entity{ + ID: "org-2", + Slug: "org-2", + Kind: "organization", + Name: "Org 2", + }) + updatedPack.Sources.Bindings = append(updatedPack.Sources.Bindings, config.Binding{ + ID: "github-binding-2", + EntityID: "org-2", + SourceID: "github", + Enabled: true, + Scope: config.JSONMap{ + "repos": []string{"example/repo-2"}, + }, + Notes: "second github binding", + }) + + if err := SyncPack(ctx, db, updatedPack); err != nil { + t.Fatalf("SyncPack() error = %v", err) + } + + scheduleRunCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM schedule_runs WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(schedule_runs) error = %v", err) + } + if scheduleRunCount != 1 { + t.Fatalf("scheduleRunCount = %d, want %d", scheduleRunCount, 1) + } + + alertEventCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM alert_events WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(alert_events) error = %v", err) + } + if alertEventCount != 1 { + t.Fatalf("alertEventCount = %d, want %d", alertEventCount, 1) + } + + candidateSuggestionCount, err := processingQueryCount( + ctx, + db, + `SELECT COUNT(*) FROM candidate_suggestions WHERE dataset_id = ?`, + datasetID, + ) + if err != nil { + t.Fatalf("processingQueryCount(candidate_suggestions) error = %v", err) + } + if candidateSuggestionCount != 1 { + t.Fatalf("candidateSuggestionCount = %d, want %d", candidateSuggestionCount, 1) + } +} + +func processingOpenSeededDB(t *testing.T, datasetID string) *sql.DB { + t.Helper() + + ctx := context.Background() + + db, err := Open(filepath.Join(t.TempDir(), "signalscope.db")) + if err != nil { + t.Fatalf("Open() error = %v", err) + } + + if err := Bootstrap(ctx, db); err != nil { + _ = db.Close() + t.Fatalf("Bootstrap() error = %v", err) + } + + if err := SeedPack(ctx, db, processingFixturePack(datasetID)); err != nil { + _ = db.Close() + t.Fatalf("SeedPack() error = %v", err) + } + + return db +} + +func processingFixturePack(datasetID string) config.Pack { + return config.Pack{ + Entities: config.EntitiesFile{ + Version: 1, + DatasetID: datasetID, + DatasetName: "Processing Fixture Dataset", + Entities: []config.Entity{ + { + ID: "org-1", + Slug: "org-1", + Kind: "organization", + Name: "Org 1", + }, + }, + Relationships: []config.Relationship{}, + }, + Sources: config.SourcesFile{ + Version: 1, + DatasetID: datasetID, + Sources: []config.Source{ + { + ID: "github", + Kind: "github", + Enabled: true, + Defaults: config.JSONMap{}, + }, + { + ID: "news-rss", + Kind: "news_rss", + Enabled: true, + Defaults: config.JSONMap{}, + }, + }, + Bindings: []config.Binding{ + { + ID: "github-binding-1", + EntityID: "org-1", + SourceID: "github", + Enabled: true, + Scope: config.JSONMap{ + "repos": []string{"example/repo-1"}, + }, + Notes: "github binding 1", + }, + { + ID: "news-binding-1", + EntityID: "org-1", + SourceID: "news-rss", + Enabled: true, + Scope: config.JSONMap{ + "query": "Example Org", + }, + Notes: "news binding 1", + }, + }, + }, + Alerts: config.AlertsFile{ + Version: 1, + DatasetID: datasetID, + Rules: []config.AlertRule{ + { + ID: "rule-release", + Name: "Release Silence Alert", + Enabled: true, + SourceID: "github", + AppliesToKinds: []string{"organization"}, + EventKey: "release.published", + Window: &config.TimeRange{ + Value: 7, + Unit: "day", + }, + Condition: config.JSONMap{ + "type": "silence_then_event", + "silence_days": 7, + }, + Severity: "medium", + }, + { + ID: "rule-candidate", + Name: "Candidate Repeat Alert", + Enabled: true, + SourceID: "news-rss", + AppliesToKinds: []string{"organization"}, + EventKey: "news.co_mention", + Window: &config.TimeRange{ + Value: 7, + Unit: "day", + }, + Condition: config.JSONMap{ + "type": "candidate_repeat_gte", + "value": 2, + }, + Severity: "low", + }, + }, + }, + Schedules: config.SchedulesFile{ + Version: 1, + DatasetID: datasetID, + Jobs: []config.ScheduleJob{ + { + ID: "job-github", + SourceID: "github", + Enabled: true, + Cadence: "24h", + TimeoutSeconds: 30, + JitterSeconds: 0, + Notes: "github schedule job", + }, + }, + }, + } +} + +func processingApplyEmbeddedMigration(t *testing.T, ctx context.Context, db *sql.DB, name string) { + t.Helper() + + sqlText, err := repomigrations.Files.ReadFile(name) + if err != nil { + t.Fatalf("ReadFile(%q) error = %v", name, err) + } + + version, err := parseMigrationVersion(name) + if err != nil { + t.Fatalf("parseMigrationVersion(%q) error = %v", name, err) + } + + if err := applyMigration(ctx, db, version, name, string(sqlText)); err != nil { + t.Fatalf("applyMigration(%q) error = %v", name, err) + } +} + +func processingTableExists(ctx context.Context, db *sql.DB, tableName string) (bool, error) { + var count int + if err := db.QueryRowContext( + ctx, + `SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?`, + tableName, + ).Scan(&count); err != nil { + return false, err + } + + return count == 1, nil +} + +func processingReadScheduleRunByID(ctx context.Context, db *sql.DB, id int64) (ScheduleRun, error) { + return scanScheduleRun( + db.QueryRowContext( + ctx, + `SELECT + id, + dataset_id, + schedule_job_id, + source_id, + status, + started_at, + finished_at, + error_message, + fetch_runs_started, + fetch_runs_succeeded, + fetch_runs_failed + FROM schedule_runs + WHERE id = ?`, + id, + ), + ) +} + +func processingReadAlertEventByID(ctx context.Context, db *sql.DB, id int64) (AlertEvent, error) { + return scanAlertEvent( + db.QueryRowContext( + ctx, + `SELECT + id, + dataset_id, + alert_rule_id, + entity_id, + source_id, + severity, + status, + summary, + payload_json, + observed_at, + triggered_at, + resolved_at + FROM alert_events + WHERE id = ?`, + id, + ), + ) +} + +func processingQueryCount(ctx context.Context, db *sql.DB, query string, args ...any) (int, error) { + var count int + if err := db.QueryRowContext(ctx, query, args...).Scan(&count); err != nil { + return 0, err + } + return count, nil +} diff --git a/migrations/003_processing.sql b/migrations/003_processing.sql new file mode 100644 index 0000000..145896b --- /dev/null +++ b/migrations/003_processing.sql @@ -0,0 +1,107 @@ +PRAGMA foreign_keys = ON; + +CREATE TABLE IF NOT EXISTS schedule_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + dataset_id TEXT NOT NULL, + schedule_job_id TEXT NOT NULL, + source_id TEXT NOT NULL, + status TEXT NOT NULL, + started_at TEXT NOT NULL, + finished_at TEXT, + error_message TEXT NOT NULL DEFAULT '', + fetch_runs_started INTEGER NOT NULL DEFAULT 0, + fetch_runs_succeeded INTEGER NOT NULL DEFAULT 0, + fetch_runs_failed INTEGER NOT NULL DEFAULT 0, + FOREIGN KEY (dataset_id) REFERENCES datasets(id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, schedule_job_id) REFERENCES schedule_jobs(dataset_id, id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, source_id) REFERENCES sources(dataset_id, id) ON DELETE CASCADE, + CHECK (status IN ('running', 'succeeded', 'partial', 'failed')), + CHECK (started_at <> ''), + CHECK ( + (status = 'running' AND finished_at IS NULL) OR + (status IN ('succeeded', 'partial', 'failed') AND finished_at IS NOT NULL) + ), + CHECK (fetch_runs_started >= 0), + CHECK (fetch_runs_succeeded >= 0), + CHECK (fetch_runs_failed >= 0) +); + +CREATE TABLE IF NOT EXISTS alert_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + dataset_id TEXT NOT NULL, + alert_rule_id TEXT NOT NULL, + entity_id TEXT NOT NULL, + source_id TEXT NOT NULL, + severity TEXT NOT NULL, + status TEXT NOT NULL, + summary TEXT NOT NULL, + payload_json TEXT NOT NULL DEFAULT '{}', + observed_at TEXT NOT NULL, + triggered_at TEXT NOT NULL, + resolved_at TEXT, + FOREIGN KEY (dataset_id) REFERENCES datasets(id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, alert_rule_id) REFERENCES alert_rules(dataset_id, id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, entity_id) REFERENCES entities(dataset_id, id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, source_id) REFERENCES sources(dataset_id, id) ON DELETE CASCADE, + CHECK (severity IN ('low', 'medium', 'high')), + CHECK (status IN ('open', 'resolved')), + CHECK (summary <> ''), + CHECK (payload_json <> ''), + CHECK (observed_at <> ''), + CHECK (triggered_at <> ''), + CHECK (observed_at <= triggered_at), + CHECK ( + (status = 'open' AND resolved_at IS NULL) OR + (status = 'resolved' AND resolved_at IS NOT NULL) + ), + CHECK (resolved_at IS NULL OR resolved_at >= triggered_at) +); + +CREATE TABLE IF NOT EXISTS candidate_suggestions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + dataset_id TEXT NOT NULL, + alert_rule_id TEXT NOT NULL, + candidate_key TEXT NOT NULL, + display_name TEXT NOT NULL, + discovered_from_entity_id TEXT NOT NULL, + source_id TEXT NOT NULL, + mention_count INTEGER NOT NULL, + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + status TEXT NOT NULL, + evidence_json TEXT NOT NULL DEFAULT '{}', + FOREIGN KEY (dataset_id) REFERENCES datasets(id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, alert_rule_id) REFERENCES alert_rules(dataset_id, id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, discovered_from_entity_id) REFERENCES entities(dataset_id, id) ON DELETE CASCADE, + FOREIGN KEY (dataset_id, source_id) REFERENCES sources(dataset_id, id) ON DELETE CASCADE, + CHECK (candidate_key <> ''), + CHECK (display_name <> ''), + CHECK (mention_count > 0), + CHECK (first_seen_at <> ''), + CHECK (last_seen_at <> ''), + CHECK (first_seen_at <= last_seen_at), + CHECK (status IN ('pending_review', 'accepted', 'dismissed')), + CHECK (evidence_json <> ''), + UNIQUE (dataset_id, discovered_from_entity_id, source_id, candidate_key) +); + +CREATE INDEX IF NOT EXISTS idx_schedule_runs_dataset_job_started + ON schedule_runs(dataset_id, schedule_job_id, started_at DESC); + +CREATE INDEX IF NOT EXISTS idx_schedule_runs_dataset_source_status_started + ON schedule_runs(dataset_id, source_id, status, started_at DESC); + +CREATE INDEX IF NOT EXISTS idx_alert_events_dataset_rule_entity_status + ON alert_events(dataset_id, alert_rule_id, entity_id, status); + +CREATE INDEX IF NOT EXISTS idx_alert_events_dataset_entity_observed + ON alert_events(dataset_id, entity_id, observed_at DESC); + +CREATE INDEX IF NOT EXISTS idx_alert_events_dataset_source_triggered + ON alert_events(dataset_id, source_id, triggered_at DESC); + +CREATE INDEX IF NOT EXISTS idx_candidate_suggestions_dataset_entity_status_last_seen + ON candidate_suggestions(dataset_id, discovered_from_entity_id, status, last_seen_at DESC); + +CREATE INDEX IF NOT EXISTS idx_candidate_suggestions_dataset_rule_last_seen + ON candidate_suggestions(dataset_id, alert_rule_id, last_seen_at DESC);