From e299891b36b061219fd6f716d0714d80211623f3 Mon Sep 17 00:00:00 2001 From: Daniel Rosales <111561081+dnlrsls@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:10:15 -0500 Subject: [PATCH] feat(store): expose eligible pending mutation high-water --- internal/store/store.go | 16 +++++++++++ internal/store/store_test.go | 55 ++++++++++++++++++++++++++++++++++++ 2 files changed, 71 insertions(+) diff --git a/internal/store/store.go b/internal/store/store.go index 4d9b169db..7996292e6 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -6389,6 +6389,22 @@ func (s *Store) CloudSyncSummary() (CloudSyncSummary, error) { return summary, nil } +// MaxPendingSyncMutationSeq returns the highest sequence eligible for pending sync +// on targetKey, or zero when none exists. Like ListPendingSyncMutationsAfterSeq, +// only unacknowledged pending mutations for enrolled or global projects qualify; +// sync_state counters do not determine this bound. +func (s *Store) MaxPendingSyncMutationSeq(targetKey string) (int64, error) { + targetKey = normalizeSyncTargetKey(targetKey) + var seq int64 + err := s.db.QueryRow(` + SELECT COALESCE(MAX(sm.seq), 0) + FROM sync_mutations sm + LEFT JOIN sync_enrolled_projects sep ON sm.project = sep.project + WHERE sm.target_key = ? AND sm.acked_at IS NULL AND sm.disposition = 'pending' + AND (sm.project = '' OR sep.project IS NOT NULL)`, targetKey).Scan(&seq) + return seq, err +} + func (s *Store) ListPendingSyncMutationsAfterSeq(targetKey string, afterSeq int64, limit int) ([]SyncMutation, error) { targetKey = normalizeSyncTargetKey(targetKey) if limit <= 0 { diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 8c7ed35a8..571b6058d 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -11982,6 +11982,61 @@ func TestListPendingSyncMutationsIncludesProject(t *testing.T) { } } +func TestMaxPendingSyncMutationSeq(t *testing.T) { + s := newTestStore(t) + if err := s.EnrollProject("enrolled"); err != nil { + t.Fatalf("enroll: %v", err) + } + const target = "cloud:high-water-test" + for _, key := range []string{target, "cloud:other-high-water"} { + if _, err := s.db.Exec(`INSERT INTO sync_state (target_key, lifecycle, last_enqueued_seq, updated_at) VALUES (?, 'idle', 0, datetime('now'))`, key); err != nil { + t.Fatalf("insert sync state %s: %v", key, err) + } + } + insert := func(key, targetKey, project, ackedAt, disposition string) int64 { + t.Helper() + var acked any + if ackedAt != "" { + acked = ackedAt + } + result, err := s.db.Exec(`INSERT INTO sync_mutations + (target_key, entity, entity_key, op, payload, source, project, acked_at, disposition) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, targetKey, SyncEntityObservation, key, + SyncOpUpsert, `{}`, SyncSourceLocal, project, acked, disposition) + if err != nil { + t.Fatalf("insert %s: %v", key, err) + } + seq, err := result.LastInsertId() + if err != nil { + t.Fatalf("sequence %s: %v", key, err) + } + return seq + } + check := func(targetKey string, want int64) { + t.Helper() + got, err := s.MaxPendingSyncMutationSeq(targetKey) + if err != nil || got != want { + t.Fatalf("MaxPendingSyncMutationSeq(%q) = %d, %v; want %d", targetKey, got, err, want) + } + } + check(target, 0) + insert("other-target", "cloud:other-high-water", "", "", "pending") + global := insert("global", target, "", "", "pending") + insert("un-enrolled", target, "not-enrolled", "", "pending") + insert("acked", target, "enrolled", "2025-01-01T00:00:00Z", "pending") + insert("quarantined", target, "enrolled", "", "quarantined") + eligible := insert("enrolled", target, "enrolled", "", "pending") + if _, err := s.db.Exec(`UPDATE sync_state SET last_enqueued_seq = 0 WHERE target_key = ?`, target); err != nil { + t.Fatalf("stale sync state: %v", err) + } + check(target, eligible) + check("cloud:other-high-water", 1) + if _, err := s.db.Exec(`UPDATE sync_mutations SET acked_at = ? WHERE seq = ?`, "2025-01-01T00:00:00Z", eligible); err != nil { + t.Fatalf("ack eligible: %v", err) + } + check(target, global) +} + func TestCountPendingNonEnrolledSyncMutations(t *testing.T) { s := newTestStore(t)