package reconciliation import ( "context" "fmt" "time" "archivmail/internal/audit" ) // bucketKey identifies one reconciliation bucket in memory. Nil tenant/source // IDs are encoded as -1 so they can be used as map keys. type bucketKey struct { tenant int64 sourceType string sourceID int64 } func keyOf(tenant, sourceID *int64, sourceType string) bucketKey { t := int64(-1) if tenant != nil { t = *tenant } sid := int64(-1) if sourceID != nil { sid = *sourceID } return bucketKey{tenant: t, sourceType: sourceType, sourceID: sid} } func ptr(v int64) *int64 { return &v } func nilIfNeg(v int64) *int64 { if v < 0 { return nil } return ptr(v) } // imapExpected holds the IMAP soll/ist snapshot for one account. type imapExpected struct { tenant *int64 expected int64 // sum of per-folder last_uid high-water marks (source proxy) cumulArch int64 // cumulative archived mails for this account } // Anomaly describes a detected significant drop for one source/day. type Anomaly struct { Date time.Time TenantID *int64 SourceKey string Archived int64 Average float64 ThresholdPct int } // ComputeForDate reconciles a single calendar day (UTC) and upserts one row per // known source bucket. Rows are only written after every read query has // succeeded, so a mid-job DB failure leaves the day WITHOUT a report row // (dashboard shows "data missing") instead of a misleading all-zero report. // // After persisting, it evaluates the trailing 7-day average per source and // writes a `reconciliation_anomaly` audit entry when today's archived count has // dropped more than thresholdPct percent below that average. Sources with fewer // than 7 prior daily records are skipped ("not enough data yet"). // // Returns the anomalies detected (also useful for the CLI summary/tests). func (s *Store) ComputeForDate(ctx context.Context, day time.Time, thresholdPct int) ([]Anomaly, error) { dayStart := time.Date(day.Year(), day.Month(), day.Day(), 0, 0, 0, 0, time.UTC) dayEnd := dayStart.Add(24 * time.Hour) // ── Read phase (all-or-nothing) ──────────────────────────────────────── archived, err := s.archivedForDay(ctx, dayStart, dayEnd) if err != nil { return nil, err } known, err := s.knownBuckets(ctx) if err != nil { return nil, err } imap, err := s.imapExpectedSnapshot(ctx) if err != nil { return nil, err } // ── Build rows ───────────────────────────────────────────────────────── rows := make([]Report, 0, len(known)) for k := range known { r := Report{ Date: dayStart, TenantID: nilIfNeg(k.tenant), SourceType: k.sourceType, SourceID: nilIfNeg(k.sourceID), ArchivedCount: archived[k], // 0 when the source had no mail that day } // IMAP soll/ist: expected = source mailbox size proxy (sum of last_uid), // delta = cumulative archived for the account − expected. delta going // negative because the user emptied the source mailbox is a "good" // direction and never alerts (alerting keys off archived_count only). if k.sourceType == "imap" && k.sourceID >= 0 { if ie, ok := imap[k.sourceID]; ok { r.ExpectedCount = ptr(ie.expected) r.Delta = ptr(ie.cumulArch - ie.expected) } } rows = append(rows, r) } // ── Write phase ──────────────────────────────────────────────────────── if err := s.upsertRows(ctx, rows); err != nil { return nil, err } // ── Alert phase ──────────────────────────────────────────────────────── var anomalies []Anomaly for _, r := range rows { avg, n, err := s.trailingAverage(ctx, r, dayStart) if err != nil { s.logger.Warn("reconciliation: trailing average failed", "source", SourceKey(r.SourceType, r.SourceID), "err", err) continue } if n < 7 { continue // not enough history yet } limit := avg * (1 - float64(thresholdPct)/100.0) if float64(r.ArchivedCount) < limit { a := Anomaly{ Date: dayStart, TenantID: r.TenantID, SourceKey: SourceKey(r.SourceType, r.SourceID), Archived: r.ArchivedCount, Average: avg, ThresholdPct: thresholdPct, } anomalies = append(anomalies, a) s.logAnomaly(a) } } return anomalies, nil } // archivedForDay returns the count of newly archived mails per source bucket for // the given day window. Mails with NULL source_type (legacy) bucket as 'import'. func (s *Store) archivedForDay(ctx context.Context, start, end time.Time) (map[bucketKey]int64, error) { rows, err := s.pool.Query(ctx, ` SELECT tenant_id, COALESCE(source_type, 'import') AS st, source_id, COUNT(*) FROM emails WHERE received_at >= $1 AND received_at < $2 GROUP BY tenant_id, st, source_id `, start, end) if err != nil { return nil, fmt.Errorf("reconciliation: archived-for-day query: %w", err) } defer rows.Close() out := make(map[bucketKey]int64) for rows.Next() { var tenant, sourceID *int64 var st string var cnt int64 if err := rows.Scan(&tenant, &st, &sourceID, &cnt); err != nil { return nil, fmt.Errorf("reconciliation: archived-for-day scan: %w", err) } out[keyOf(tenant, sourceID, st)] = cnt } return out, rows.Err() } // knownBuckets returns every (tenant, source_type, source_id) combination that // has ever produced an archived mail. These are the buckets for which a report // row is written every day — including explicit 0 on inactive days so gaps in // the cron run are distinguishable from genuine zero-activity days. func (s *Store) knownBuckets(ctx context.Context) (map[bucketKey]struct{}, error) { rows, err := s.pool.Query(ctx, ` SELECT DISTINCT tenant_id, COALESCE(source_type, 'import') AS st, source_id FROM emails `) if err != nil { return nil, fmt.Errorf("reconciliation: known-buckets query: %w", err) } defer rows.Close() out := make(map[bucketKey]struct{}) for rows.Next() { var tenant, sourceID *int64 var st string if err := rows.Scan(&tenant, &st, &sourceID); err != nil { return nil, fmt.Errorf("reconciliation: known-buckets scan: %w", err) } out[keyOf(tenant, sourceID, st)] = struct{}{} } return out, rows.Err() } // imapExpectedSnapshot returns the IMAP soll/ist snapshot keyed by account ID. // expected reuses the per-folder UID high-water marks from imap_folder_state // (PROJ-45) — no additional IMAP login. cumulArch is the cumulative number of // archived mails attributed to the account. func (s *Store) imapExpectedSnapshot(ctx context.Context) (map[int64]imapExpected, error) { out := make(map[int64]imapExpected) // Expected proxy + tenant per account. LEFT JOIN so accounts without any // synced folder yet still appear (expected 0). rows, err := s.pool.Query(ctx, ` SELECT a.id, a.tenant_id, COALESCE(SUM(fs.last_uid), 0) FROM imap_accounts a LEFT JOIN imap_folder_state fs ON fs.account_id = a.id GROUP BY a.id, a.tenant_id `) if err != nil { return nil, fmt.Errorf("reconciliation: imap expected query: %w", err) } defer rows.Close() for rows.Next() { var id int64 var tenant *int64 var expected int64 if err := rows.Scan(&id, &tenant, &expected); err != nil { return nil, fmt.Errorf("reconciliation: imap expected scan: %w", err) } out[id] = imapExpected{tenant: tenant, expected: expected} } if err := rows.Err(); err != nil { return nil, err } // Cumulative archived count per IMAP account. crows, err := s.pool.Query(ctx, ` SELECT source_id, COUNT(*) FROM emails WHERE source_type = 'imap' AND source_id IS NOT NULL GROUP BY source_id `) if err != nil { return nil, fmt.Errorf("reconciliation: imap archived query: %w", err) } defer crows.Close() for crows.Next() { var id, cnt int64 if err := crows.Scan(&id, &cnt); err != nil { return nil, fmt.Errorf("reconciliation: imap archived scan: %w", err) } ie := out[id] ie.cumulArch = cnt out[id] = ie } return out, crows.Err() } // upsertRows writes all report rows in a single transaction. On any error the // transaction is rolled back so no partial day is persisted. func (s *Store) upsertRows(ctx context.Context, rows []Report) error { tx, err := s.pool.Begin(ctx) if err != nil { return fmt.Errorf("reconciliation: begin tx: %w", err) } defer tx.Rollback(ctx) for _, r := range rows { _, err := tx.Exec(ctx, ` INSERT INTO reconciliation_reports (date, tenant_id, source_type, source_id, expected_count, archived_count, delta) VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT (date, COALESCE(tenant_id, -1), source_type, COALESCE(source_id, -1)) DO UPDATE SET expected_count = EXCLUDED.expected_count, archived_count = EXCLUDED.archived_count, delta = EXCLUDED.delta, created_at = NOW() `, r.Date, r.TenantID, r.SourceType, r.SourceID, r.ExpectedCount, r.ArchivedCount, r.Delta) if err != nil { return fmt.Errorf("reconciliation: upsert row: %w", err) } } if err := tx.Commit(ctx); err != nil { return fmt.Errorf("reconciliation: commit: %w", err) } return nil } // trailingAverage returns the mean archived_count of the up-to-7 report rows // immediately preceding the given day for the same source bucket, plus how many // prior day rows were found. NULL-safe matching on tenant_id / source_id. func (s *Store) trailingAverage(ctx context.Context, r Report, day time.Time) (float64, int, error) { rows, err := s.pool.Query(ctx, ` SELECT archived_count FROM reconciliation_reports WHERE source_type = $1 AND tenant_id IS NOT DISTINCT FROM $2 AND source_id IS NOT DISTINCT FROM $3 AND date < $4 ORDER BY date DESC LIMIT 7 `, r.SourceType, r.TenantID, r.SourceID, day) if err != nil { return 0, 0, fmt.Errorf("reconciliation: trailing average query: %w", err) } defer rows.Close() var sum int64 var n int for rows.Next() { var c int64 if err := rows.Scan(&c); err != nil { return 0, 0, err } sum += c n++ } if err := rows.Err(); err != nil { return 0, 0, err } if n == 0 { return 0, 0, nil } return float64(sum) / float64(n), n, nil } // logAnomaly emits a structured log line and, when wired, a tenant-visible // audit entry for a detected drop. func (s *Store) logAnomaly(a Anomaly) { detail := fmt.Sprintf("source=%s date=%s archived=%d avg_7d=%.1f threshold=%d%%", a.SourceKey, a.Date.Format("2006-01-02"), a.Archived, a.Average, a.ThresholdPct) s.logger.Warn("reconciliation: anomaly detected", "detail", detail) if s.audlog != nil { s.audlog.Log(audit.Entry{ EventType: audit.EventReconciliationAnomaly, Username: "system", TenantID: a.TenantID, Success: false, Detail: detail, }) } }