Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions internal/store/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
55 changes: 55 additions & 0 deletions internal/store/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Loading