From d18c32d8db4d26f38546305ec2ab16ff17c90858 Mon Sep 17 00:00:00 2001 From: Daniel Rosales <111561081+dnlrsls@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:31:23 -0500 Subject: [PATCH] feat(autosync): preflight keyed prompt authority before push --- cmd/engram/autosync_e2e_test.go | 27 +++ docs/codebase/prompt-inbox-provenance.md | 6 +- internal/cloud/autosync/manager.go | 148 ++++++++++++++-- internal/cloud/autosync/manager_test.go | 205 ++++++++++++++++++++--- 4 files changed, 345 insertions(+), 41 deletions(-) diff --git a/cmd/engram/autosync_e2e_test.go b/cmd/engram/autosync_e2e_test.go index 12680a4b6..e4fa5b718 100644 --- a/cmd/engram/autosync_e2e_test.go +++ b/cmd/engram/autosync_e2e_test.go @@ -349,6 +349,33 @@ func (s *autosyncFakeStore) ListPendingSyncMutations(_ string, limit int) ([]sto return result, nil } +func (s *autosyncFakeStore) MaxPendingSyncMutationSeq(string) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + var max int64 + for _, mutation := range s.mutations { + max = mutation.seq + } + return max, nil +} + +func (s *autosyncFakeStore) ListPendingSyncMutationsAfterSeq(target string, after int64, limit int) ([]store.SyncMutation, error) { + pending, err := s.ListPendingSyncMutations(target, len(s.mutations)) + if err != nil { + return nil, err + } + var page []store.SyncMutation + for _, mutation := range pending { + if mutation.Seq > after { + page = append(page, mutation) + if len(page) == limit { + break + } + } + } + return page, nil +} + func (s *autosyncFakeStore) CountPendingNonEnrolledSyncMutations(_ string) ([]store.PendingSyncMutationProjectCount, error) { return nil, nil } diff --git a/docs/codebase/prompt-inbox-provenance.md b/docs/codebase/prompt-inbox-provenance.md index 807a4cb27..72ad5ee07 100644 --- a/docs/codebase/prompt-inbox-provenance.md +++ b/docs/codebase/prompt-inbox-provenance.md @@ -1,6 +1,6 @@ # RFC: authenticated prompt inbox provenance before remote deletion -**Status: proposed end-to-end contract; cloud registration, pair-claim storage and admission routes exist; local origin marking is implemented, but no client handshake or verified delete gate exists.** This RFC defines the minimum non-cryptographic authority needed before #1464 can enter the merge queue. It does not change current sync behavior. The baseline is tracker `313f5269`; the review thread on #1464 records the hold. #1240 is separate. +**Status: proposed end-to-end contract; cloud registration and pair-claim routes exist, local origin marking and eligible-local-origin autosync preflight are implemented, but no verified cloud delete gate or old/imported source reauthorization exists.** This RFC defines the minimum non-cryptographic authority needed before #1464 can enter the merge queue. The RFC alone does not establish end-to-end enforcement. The baseline is tracker `313f5269`; the review thread on #1464 records the hold. #1240 is separate. ## Decision in one minute @@ -16,7 +16,7 @@ This is a server-enforced authorization contract, not a new inference rule based | Claim prompt pair | For a beta prompt referencing an alpha session, the claimant must hold authorization for **both** alpha and beta at claim time. Registration must already be verified. An authorized claim durably binds `(session_id, source_inbox_id)` to `(sync_id, beta)`; matching replay is idempotent and any competing sync ID or project is a conflict. Claims cannot bootstrap registration or be inferred from a beta upsert. Same-project claims still require session-owner and prompt-project authority; no special weaker bootstrap. | | Apply verified delete | A delete may reserve or replay a remote pair only against the existing verified binding, with beta authorization for the mutation. It need not require renewed alpha authorization: alpha consent was checked at the original claim. Validate sync ID, session ID, inbox ID and beta project against that binding. A beta-only writer cannot create or change the claim, even by sending an upsert followed by a delete; an alpha-only writer cannot claim/delete beta. | -Registration and claim are explicit authenticated server operations (`POST /sync/session-authorities` and `POST /sync/prompt-pair-claims`); the verified-delete gate and client handshake remain pending. Atomic uniqueness and conflict checks must survive concurrent retries. Authentication means server-verified principal and project grants, not fields supplied in the payload; no cryptographic offline capability is assumed. Revocation after a verified claim does not erase that historical binding, but current beta mutation authorization remains necessary. Idless legacy prompts remain pair-less and must not acquire a source inbox binding through inference. +Registration and claim are explicit authenticated server operations (`POST /sync/session-authorities` and `POST /sync/prompt-pair-claims`); autosync now performs the handshake only for independently eligible local keyed prompt mutations. The verified-delete gate and explicit old/imported source reauthorization remain pending. Atomic uniqueness and conflict checks must survive concurrent retries. Authentication means server-verified principal and project grants, not fields supplied in the payload; no cryptographic offline capability is assumed. Revocation after a verified claim does not erase that historical binding, but current beta mutation authorization remains necessary. Idless legacy prompts remain pair-less and must not acquire a source inbox binding through inference. ## Unverified deletes and compatibility @@ -44,4 +44,4 @@ Each PR is at most **400 authored diff lines**, includes its own tests/docs with ## Limitations -Pending retry UX and any cryptographic offline issuer still require design and tests; none may turn chunk/import history into authority. Session registration and its storage schema exist. SQLite records `local_creation_project` only on newly inserted CreateSession/StartSession rows. Newly inserted keyed `AddPromptWithResult` rows also record their original session ID, inbox ID and prompt project; existing/replayed, idless, imported and pulled rows do not acquire this marker. The live-row read requires exact agreement with all three original values, including beta prompt ownership under an alpha session. This marker is local creation evidence, not cloud authority. Explicit local `DeletePrompt` and `DeleteSession` retain the verified original triple in separate nullable tombstone columns in the same transaction; migration defaults to NULL. The exact-sync-ID lookup uses that triple only after no live row exists, and only if it agrees with the tombstone identity. Imported/pulled tombstones and tombstone replay cannot create or promote local origin. Ambiguous sync IDs, idless prompts and mismatched identities remain ineligible. NULL means unknown/unverified for pre-migration, imported, rescued and remotely created sessions. A pulled upsert into a locally created session preserves its existing marker. Eligibility requires the current owner to match that stored creation owner exactly; identity repair preserves the nullable value. A store read returns current owner and local eligibility, not cloud authority. Cloud storage supports explicit immutable prompt pair claims against registered sessions, without deriving claims from chunks or prompt history. `ClaimPromptPair` assumes its caller has authenticated the actor and checked both project grants; it does not perform authorization. Authenticated registration and pair-claim routes exist; no verified delete admission or client handshake is implemented yet. The current cloud behavior does **not** enforce this RFC. +Pending retry UX and any cryptographic offline issuer still require design and tests; none may turn chunk/import history into authority. Session registration and its storage schema exist. SQLite records `local_creation_project` only on newly inserted CreateSession/StartSession rows. Newly inserted keyed `AddPromptWithResult` rows also record their original session ID, inbox ID and prompt project; existing/replayed, idless, imported and pulled rows do not acquire this marker. The live-row read requires exact agreement with all three original values, including beta prompt ownership under an alpha session. This marker is local creation evidence, not cloud authority. Explicit local `DeletePrompt` and `DeleteSession` retain the verified original triple in separate nullable tombstone columns in the same transaction; migration defaults to NULL. The exact-sync-ID lookup uses that triple only after no live row exists, and only if it agrees with the tombstone identity. Imported/pulled tombstones and tombstone replay cannot create or promote local origin. Ambiguous sync IDs, idless prompts and mismatched identities remain ineligible. NULL means unknown/unverified for pre-migration, imported, rescued and remotely created sessions. A pulled upsert into a locally created session preserves its existing marker. Eligibility requires the current owner to match that stored creation owner exactly; identity repair preserves the nullable value. A store read returns current owner and local eligibility, not cloud authority. Cloud storage supports explicit immutable prompt pair claims against registered sessions, without deriving claims from chunks or prompt history. `ClaimPromptPair` assumes its caller has authenticated the actor and checked both project grants; it does not perform authorization. Autosync checks the exact journal sync ID, session, inbox and prompt project against independent locally recorded prompt origin and the live session's local creation owner; it registers that owner, then claims the pair before pushing a keyed prompt mutation. An unavailable capability, missing or mismatched origin, or failed registration/claim leaves the mutation pending and unpushed; other eligible entries in the same project batch may still push and ack independently. Idless prompts do not claim pairs. A locally deleted session has no live session provenance lookup, so its prompt tombstone remains pending even if its prompt origin is marked; T4b3b or a separately reviewed locally proven session-tombstone mechanism must address it, never project inference. This is a client preflight only: old/imported sources lack explicit reauthorization, and the cloud still does **not** enforce verified delete admission. diff --git a/internal/cloud/autosync/manager.go b/internal/cloud/autosync/manager.go index af3f3f867..3b74c82c2 100644 --- a/internal/cloud/autosync/manager.go +++ b/internal/cloud/autosync/manager.go @@ -120,6 +120,23 @@ type CloudTransport interface { PullMutations(sinceSeq int64, limit int) (*PullMutationsResponse, error) } +// Provenance capabilities are optional for legacy transports/stores, but keyed +// prompts fail closed when either capability is unavailable. +type pendingMutationPager interface { + ListPendingSyncMutationsAfterSeq(targetKey string, afterSeq int64, limit int) ([]store.SyncMutation, error) + MaxPendingSyncMutationSeq(targetKey string) (int64, error) +} + +type localPromptProvenance interface { + LocalSessionProvenance(id string) (owner string, eligible bool, err error) + LocalPromptCreationIdentity(syncID string) (session, inbox, project string, eligible bool, err error) +} + +type promptAuthorityTransport interface { + RegisterSessionAuthority(sessionID, ownerProject string) error + ClaimPromptPair(sessionID, sourceInboxID, syncID, ownerProject, promptProject string) error +} + // transportStatusError is an optional interface that transport errors may implement. // BW5: Allows Manager to detect 401 (auth_required) vs 403 (policy_forbidden) // vs generic transport failures without importing the remote package. @@ -618,11 +635,47 @@ func (m *Manager) push(ctx context.Context) error { } } - pending, err := m.store.ListPendingSyncMutations(m.cfg.TargetKey, m.cfg.PushBatchSize) + pager, ok := m.store.(pendingMutationPager) + if !ok { + return fmt.Errorf("bounded pending mutation pagination unavailable") + } + // Snapshot the eligible journal after repair: new enqueues belong to a later cycle. + highWater, err := pager.MaxPendingSyncMutationSeq(m.cfg.TargetKey) if err != nil { - return fmt.Errorf("list pending: %w", err) + return fmt.Errorf("read push high-water: %w", err) } - if len(pending) == 0 { + var failures []error + var afterSeq int64 + seen := false + for { + if err := ctx.Err(); err != nil { + return errors.Join(append(failures, err)...) + } + pending, err := pager.ListPendingSyncMutationsAfterSeq(m.cfg.TargetKey, afterSeq, m.cfg.PushBatchSize) + if err != nil { + return errors.Join(append(failures, fmt.Errorf("list pending: %w", err))...) + } + page := make([]store.SyncMutation, 0, len(pending)) + for _, mut := range pending { + if mut.Seq <= afterSeq { + return errors.Join(append(failures, fmt.Errorf("pending pagination did not advance"))...) + } + if mut.Seq > highWater { + break + } + page = append(page, mut) + afterSeq = mut.Seq + } + if len(page) == 0 { + break + } + seen = true + failures = append(failures, m.pushPage(ctx, page)...) + if len(pending) < m.cfg.PushBatchSize || afterSeq >= highWater { + break + } + } + if !seen { counts, err := m.store.CountPendingNonEnrolledSyncMutations(m.cfg.TargetKey) if err != nil { return fmt.Errorf("count pending non-enrolled mutations: %w", err) @@ -633,6 +686,10 @@ func (m *Manager) push(ctx context.Context) error { return nil } + return errors.Join(failures...) +} + +func (m *Manager) pushPage(ctx context.Context, pending []store.SyncMutation) []error { // Group by project (preserve order). Empty or padded project values are invalid // for cloud transport: never send them, but continue with healthy project groups. groups := make(map[string][]store.SyncMutation) @@ -653,22 +710,43 @@ func (m *Manager) push(ctx context.Context) error { for _, project := range order { if err := ctx.Err(); err != nil { failures = append(failures, err) - return errors.Join(failures...) + return failures } batch := groups[project] - entries := make([]MutationEntry, len(batch)) - seqs := make([]int64, len(batch)) - for i, mut := range batch { - entries[i] = MutationEntry{ - Project: mut.Project, - Entity: mut.Entity, - EntityKey: mut.EntityKey, - Op: mut.Op, - Payload: json.RawMessage(mut.Payload), + entries := make([]MutationEntry, 0, len(batch)) + seqs := make([]int64, 0, len(batch)) + for _, mut := range batch { + if mut.Entity == store.SyncEntityPrompt { + var identity struct { + SyncID string `json:"sync_id"` + Session string `json:"session_id"` + Inbox string `json:"source_inbox_id"` + Project string `json:"project"` + } + if err := json.Unmarshal([]byte(mut.Payload), &identity); err != nil { + failures = append(failures, fmt.Errorf("prompt seq %d: invalid identity: %w", mut.Seq, err)) + continue + } + if identity.SyncID == "" || identity.SyncID != mut.EntityKey || identity.Project != project || identity.Session == "" { + failures = append(failures, fmt.Errorf("prompt seq %d: invalid journal identity", mut.Seq)) + continue + } + if identity.Inbox != "" { + if err := m.preflightPrompt(mut, identity.SyncID, identity.Session, identity.Inbox, identity.Project); err != nil { + failures = append(failures, fmt.Errorf("prompt seq %d: %w", mut.Seq, err)) + continue + } + } } - seqs[i] = mut.Seq + entries = append(entries, MutationEntry{ + Project: mut.Project, Entity: mut.Entity, EntityKey: mut.EntityKey, + Op: mut.Op, Payload: json.RawMessage(mut.Payload), + }) + seqs = append(seqs, mut.Seq) + } + if len(entries) == 0 { + continue } - result, err := m.transport.PushMutations(entries) if err != nil { failures = append(failures, &projectTransportFailure{project: project, err: err}) @@ -684,11 +762,47 @@ func (m *Manager) push(ctx context.Context) error { } if err := m.store.AckSyncMutationSeqs(m.cfg.TargetKey, seqs); err != nil { failures = append(failures, fmt.Errorf("ack project %q: %w", project, err)) - return errors.Join(failures...) + return failures } } - return errors.Join(failures...) + return failures +} + +func (m *Manager) preflightPrompt(mut store.SyncMutation, syncID, session, inbox, project string) error { + if syncID == "" || session == "" || inbox == "" || project == "" || + syncID != mut.EntityKey || project != mut.Project { + return fmt.Errorf("unverified keyed prompt mutation identity") + } + local, ok := m.store.(localPromptProvenance) + if !ok { + return fmt.Errorf("local prompt provenance unavailable") + } + remote, ok := m.transport.(promptAuthorityTransport) + if !ok { + return fmt.Errorf("remote prompt authority unavailable") + } + originalSession, originalInbox, originalProject, eligible, err := local.LocalPromptCreationIdentity(syncID) + if err != nil { + return fmt.Errorf("read local prompt origin: %w", err) + } + if !eligible || originalSession != session || originalInbox != inbox || originalProject != project { + return fmt.Errorf("unverified keyed prompt origin") + } + owner, eligible, err := local.LocalSessionProvenance(session) + if err != nil { + return fmt.Errorf("read local session origin: %w", err) + } + if !eligible || owner == "" { + return fmt.Errorf("unverified local session owner") + } + if err := remote.RegisterSessionAuthority(session, owner); err != nil { + return fmt.Errorf("register session authority: %w", err) + } + if err := remote.ClaimPromptPair(session, inbox, syncID, owner, project); err != nil { + return fmt.Errorf("claim prompt pair: %w", err) + } + return nil } // ─── Pull ──────────────────────────────────────────────────────────────────── diff --git a/internal/cloud/autosync/manager_test.go b/internal/cloud/autosync/manager_test.go index 0ed7174d0..8d2d9e230 100644 --- a/internal/cloud/autosync/manager_test.go +++ b/internal/cloud/autosync/manager_test.go @@ -19,29 +19,30 @@ import ( // ─── Fakes ─────────────────────────────────────────────────────────────────── type fakeLocalStore struct { - mu sync.Mutex - mutations []store.SyncMutation - syncState *store.SyncState - leaseOwner string - leaseCalls int - pushErr error - pullErr error - failureMessage string - failureReason string - blockedReason string - blockedMessage string - appliedMuts []store.SyncMutation - acquireGranted bool - ackedSeqs []int64 - ackErr error - healthyCalls int + mu sync.Mutex + mutations []store.SyncMutation + syncState *store.SyncState + leaseOwner string + leaseCalls int + pushErr error + pullErr error + failureMessage string + failureReason string + blockedReason string + blockedMessage string + appliedMuts []store.SyncMutation + acquireGranted bool + ackedSeqs []int64 + ackErr error + healthyCalls int + staleHighWater bool blockedAfterSuccessCalls int blockedAfterSuccessErr error nonEnrolledCounts []store.PendingSyncMutationProjectCount - deferredProjects []string - listDeferredErr error - listedTargets []string - replayedScopes []string + deferredProjects []string + listDeferredErr error + listedTargets []string + replayedScopes []string } func newFakeLocalStore() *fakeLocalStore { @@ -61,7 +62,16 @@ func (s *fakeLocalStore) GetSyncState(_ string) (*store.SyncState, error) { if s.pullErr != nil { return nil, s.pullErr } - return s.syncState, nil + state := *s.syncState + if s.staleHighWater { + return &state, nil + } + for _, mutation := range s.mutations { + if mutation.Seq > state.LastEnqueuedSeq { + state.LastEnqueuedSeq = mutation.Seq + } + } + return &state, nil } func (s *fakeLocalStore) ListPendingSyncMutations(_ string, limit int) ([]store.SyncMutation, error) { @@ -80,6 +90,34 @@ func (s *fakeLocalStore) ListPendingSyncMutations(_ string, limit int) ([]store. return s.mutations[:n], nil } +func (s *fakeLocalStore) MaxPendingSyncMutationSeq(string) (int64, error) { + var max int64 + for _, mutation := range s.mutations { + if mutation.Seq > max { + max = mutation.Seq + } + } + return max, nil +} + +func (s *fakeLocalStore) ListPendingSyncMutationsAfterSeq(_ string, afterSeq int64, limit int) ([]store.SyncMutation, error) { + s.mu.Lock() + defer s.mu.Unlock() + if s.pushErr != nil { + return nil, s.pushErr + } + var page []store.SyncMutation + for _, mutation := range s.mutations { + if mutation.Seq > afterSeq { + page = append(page, mutation) + if len(page) == limit { + break + } + } + } + return page, nil +} + func (s *fakeLocalStore) CountPendingNonEnrolledSyncMutations(_ string) ([]store.PendingSyncMutationProjectCount, error) { s.mu.Lock() defer s.mu.Unlock() @@ -328,6 +366,131 @@ func attemptedProjects(t *fakeCloudTransport) []string { return projects } +type provenanceLocalStore struct { + *fakeLocalStore + owner, session, inbox, project string + eligible bool +} + +func (s *provenanceLocalStore) LocalSessionProvenance(string) (string, bool, error) { + return s.owner, s.eligible, nil +} +func (s *provenanceLocalStore) LocalPromptCreationIdentity(string) (string, string, string, bool, error) { + return s.session, s.inbox, s.project, s.eligible, nil +} + +type provenanceTransport struct { + *fakeCloudTransport + calls []string + registerErr, claimErr error +} + +func (t *provenanceTransport) RegisterSessionAuthority(session, owner string) error { + t.calls = append(t.calls, "register:"+session+":"+owner) + return t.registerErr +} +func (t *provenanceTransport) ClaimPromptPair(session, inbox, syncID, owner, project string) error { + t.calls = append(t.calls, "claim:"+session+":"+inbox+":"+syncID+":"+owner+":"+project) + return t.claimErr +} + +func TestManagerKeyedPromptPreflight(t *testing.T) { + for _, tc := range []struct { + name string + eligible bool + registerErr, claimErr error + wantCalls string + }{ + {"success", true, nil, nil, "[register:s:alpha claim:s:i:p:alpha:beta]"}, + {"missing origin", false, nil, nil, "[]"}, + {"registration fails", true, errors.New("register denied"), nil, "[register:s:alpha]"}, + {"claim fails", true, nil, errors.New("claim denied"), "[register:s:alpha claim:s:i:p:alpha:beta]"}, + } { + t.Run(tc.name, func(t *testing.T) { + ls := &provenanceLocalStore{fakeLocalStore: newFakeLocalStore(), owner: "alpha", session: "s", inbox: "i", project: "beta", eligible: tc.eligible} + ls.mutations = []store.SyncMutation{ + {Seq: 1, Entity: "prompt", EntityKey: "p", Op: "delete", Project: "beta", Payload: `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"beta"}`}, + {Seq: 2, Entity: "obs", EntityKey: "healthy", Op: "upsert", Project: "beta"}, + } + tr := &provenanceTransport{fakeCloudTransport: newFakeTransport(), registerErr: tc.registerErr, claimErr: tc.claimErr} + tr.pushResult = &PushMutationsResult{AcceptedSeqs: []int64{2}} + if tc.name == "success" { + tr.pushResult = &PushMutationsResult{AcceptedSeqs: []int64{1, 2}} + } + err := New(ls, tr, DefaultConfig()).push(context.Background()) + if (err == nil) != (tc.name == "success") { + t.Fatalf("push error: %v", err) + } + if got := fmt.Sprint(tr.calls); got != tc.wantCalls { + t.Fatalf("preflight calls %s, want %s", got, tc.wantCalls) + } + if tc.name != "success" && (len(tr.attempted) != 1 || len(tr.attempted[0]) != 1 || tr.attempted[0][0].EntityKey != "healthy" || fmt.Sprint(ls.ackedSeqs) != "[2]") { + t.Fatalf("unsafe push or ack: attempts=%v ack=%v", tr.attempted, ls.ackedSeqs) + } + }) + } +} + +func TestManagerPromptPreflightNoImplicitAuthority(t *testing.T) { + for _, tc := range []struct { + name, payload string + eligible bool + wantPush bool + }{ + {"idless", `{"sync_id":"p","session_id":"s","project":"beta"}`, false, true}, + {"imported keyed", `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"beta"}`, false, false}, + {"mismatched project", `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"alpha"}`, true, false}, + {"deleted session", `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"beta"}`, false, false}, + } { + t.Run(tc.name, func(t *testing.T) { + ls := &provenanceLocalStore{fakeLocalStore: newFakeLocalStore(), owner: "alpha", session: "s", inbox: "i", project: "beta", eligible: tc.eligible} + ls.mutations = []store.SyncMutation{{Seq: 1, Entity: "prompt", EntityKey: "p", Project: "beta", Op: "delete", Payload: tc.payload}} + tr := &provenanceTransport{fakeCloudTransport: newFakeTransport()} + tr.pushResult = &PushMutationsResult{AcceptedSeqs: []int64{1}} + err := New(ls, tr, DefaultConfig()).push(context.Background()) + if (err == nil) != tc.wantPush || (len(tr.attempted) == 1) != tc.wantPush || (len(ls.ackedSeqs) == 1) != tc.wantPush { + t.Fatalf("error=%v attempts=%v ack=%v", err, tr.attempted, ls.ackedSeqs) + } + if len(tr.calls) != 0 { + t.Fatalf("unexpected authority calls: %v", tr.calls) + } + }) + } +} + +func TestManagerStaleEnqueuedStateDoesNotHidePending(t *testing.T) { + ls := newFakeLocalStore() + ls.staleHighWater = true + ls.mutations = []store.SyncMutation{{Seq: 5, Entity: "obs", EntityKey: "later", Op: "upsert", Project: "alpha"}} + tr := newFakeTransport() + tr.pushResult = &PushMutationsResult{AcceptedSeqs: []int64{5}} + if err := New(ls, tr, DefaultConfig()).push(context.Background()); err != nil { + t.Fatal(err) + } + if fmt.Sprint(ls.ackedSeqs) != "[5]" { + t.Fatalf("pending mutation hidden by stale state: ack=%v", ls.ackedSeqs) + } +} + +func TestManagerBlockedFirstPageDoesNotStarveLaterEntries(t *testing.T) { + ls := &provenanceLocalStore{fakeLocalStore: newFakeLocalStore()} + ls.mutations = []store.SyncMutation{ + {Seq: 1, Entity: "prompt", EntityKey: "unverified", Op: "delete", Project: "beta", Payload: `{"sync_id":"unverified","session_id":"s","source_inbox_id":"i","project":"beta"}`}, + {Seq: 2, Entity: "obs", EntityKey: "healthy", Op: "upsert", Project: "beta"}, + {Seq: 3, Entity: "obs", EntityKey: "later", Op: "upsert", Project: "alpha"}, + } + tr := &provenanceTransport{fakeCloudTransport: newFakeTransport()} + tr.pushResultByProject = map[string]*PushMutationsResult{"beta": {AcceptedSeqs: []int64{2}}, "alpha": {AcceptedSeqs: []int64{3}}} + cfg := DefaultConfig() + cfg.PushBatchSize = 1 + if err := New(ls, tr, cfg).push(context.Background()); err == nil || !strings.Contains(err.Error(), "unverified") { + t.Fatalf("expected visible blocked origin, got %v", err) + } + if fmt.Sprint(ls.ackedSeqs) != "[2 3]" || fmt.Sprint(attemptedProjects(tr.fakeCloudTransport)) != "[beta alpha]" { + t.Fatalf("starved healthy entries: ack=%v attempts=%v", ls.ackedSeqs, attemptedProjects(tr.fakeCloudTransport)) + } +} + // ─── Push ack safety regressions ───────────────────────────────────────────── func TestManagerPushNoPendingDoesNotPushOrAck(t *testing.T) {