From 631db5f8fb70d0601a748c02012b5c4417250d4e Mon Sep 17 00:00:00 2001 From: Daniel Rosales <111561081+dnlrsls@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:53:45 -0500 Subject: [PATCH] fix(autosync): pull despite local prompt provenance blocks --- docs/codebase/prompt-inbox-provenance.md | 2 +- internal/cloud/autosync/manager.go | 56 ++++++++++++--- internal/cloud/autosync/manager_test.go | 90 +++++++++++++++++++++++- 3 files changed, 135 insertions(+), 13 deletions(-) diff --git a/docs/codebase/prompt-inbox-provenance.md b/docs/codebase/prompt-inbox-provenance.md index 72ad5ee07..f56fa56c0 100644 --- a/docs/codebase/prompt-inbox-provenance.md +++ b/docs/codebase/prompt-inbox-provenance.md @@ -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. 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. +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. When all push failures are deterministic local origin or owner denials, autosync still pulls admitted inbound mutations and advances their cursor while retaining a persisted degraded outbound status (`prompt_provenance_blocked`), never acknowledging blocked entries or marking sync healthy. Remote registration/claim errors, local provenance read errors, and mixed push failures involving transport or other unknown failures still skip pull. This inbound progress does not authorize the blocked prompt or its remote delete. 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 3b74c82c2..197aa0f99 100644 --- a/internal/cloud/autosync/manager.go +++ b/internal/cloud/autosync/manager.go @@ -106,6 +106,39 @@ type irreparableSyncMutationQuarantiner interface { QuarantineIrreparableSyncMutations(targetKey, project string, apply bool) (store.SyncMutationQuarantineReport, error) } +// promptPreflightError means this entry was rejected before transport and remains pending. +type promptPreflightError struct{ err error } + +func (e *promptPreflightError) Error() string { return e.err.Error() } +func (e *promptPreflightError) Unwrap() error { return e.err } + +// Only a tree made entirely of known per-entry denials may bypass the push gate. +func safeOutboundBlock(err error) bool { + if err == nil { + return false + } + switch e := err.(type) { + case *nonEnrolledPendingError, *promptPreflightError: + return true + case interface{ Unwrap() []error }: + children := e.Unwrap() + if len(children) == 0 { + return false + } + for _, child := range children { + if !safeOutboundBlock(child) { + return false + } + } + return true + default: + if inner := errors.Unwrap(err); inner != nil { + return safeOutboundBlock(inner) + } + return false + } +} + type nonEnrolledPendingError struct { counts []store.PendingSyncMutationProjectCount } @@ -508,25 +541,28 @@ func (m *Manager) cycle(ctx context.Context) { m.leaseHeld = true m.mu.Unlock() - // Push, then pull. A typed non-enrollment block applies only to outbound - // mutations, so inbound replication can still progress without changing the - // final degraded state that explains the blocked outbound backlog. + // Only exclusively local per-entry preflight blocks can leave inbound + // replication independent of the outbound failure. if err := m.push(ctx); err != nil { - var blocked *nonEnrolledPendingError - if !errors.As(err, &blocked) { + if !safeOutboundBlock(err) { reasonCode := classifyTransportError(err) m.recordFailureWithReason(autosyncFailureMessage(m.cfg.TargetKey, fmt.Sprintf("push: %v", err), err), reasonCode) return } blockedMessage := err.Error() - m.recordBlocked(blockedMessage, constants.ReasonNonEnrolledPendingMutations) + blockedReason := constants.ReasonNonEnrolledPendingMutations + var nonEnrolled *nonEnrolledPendingError + if !errors.As(err, &nonEnrolled) { + blockedReason = "prompt_provenance_blocked" + } + m.recordBlocked(blockedMessage, blockedReason) if err := m.pullPreservingSyncState(ctx); err != nil { reasonCode := classifyTransportError(err) m.recordFailureWithReason(autosyncFailureMessage(m.cfg.TargetKey, fmt.Sprintf("pull: %v", err), err), reasonCode) return } - if err := m.recordBlockedAfterSuccess(blockedMessage, constants.ReasonNonEnrolledPendingMutations); err != nil { + if err := m.recordBlockedAfterSuccess(blockedMessage, blockedReason); err != nil { reasonCode := classifyTransportError(err) m.recordFailureWithReason(autosyncFailureMessage(m.cfg.TargetKey, fmt.Sprintf("persist blocked state after successful pull: %v", err), err), reasonCode) } @@ -772,7 +808,7 @@ func (m *Manager) pushPage(ctx context.Context, pending []store.SyncMutation) [] 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") + return &promptPreflightError{err: fmt.Errorf("unverified keyed prompt mutation identity")} } local, ok := m.store.(localPromptProvenance) if !ok { @@ -787,14 +823,14 @@ func (m *Manager) preflightPrompt(mut store.SyncMutation, syncID, session, inbox return fmt.Errorf("read local prompt origin: %w", err) } if !eligible || originalSession != session || originalInbox != inbox || originalProject != project { - return fmt.Errorf("unverified keyed prompt origin") + return &promptPreflightError{err: 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") + return &promptPreflightError{err: fmt.Errorf("unverified local session owner")} } if err := remote.RegisterSessionAuthority(session, owner); err != nil { return fmt.Errorf("register session authority: %w", err) diff --git a/internal/cloud/autosync/manager_test.go b/internal/cloud/autosync/manager_test.go index 8d2d9e230..afb9a7a20 100644 --- a/internal/cloud/autosync/manager_test.go +++ b/internal/cloud/autosync/manager_test.go @@ -370,13 +370,14 @@ type provenanceLocalStore struct { *fakeLocalStore owner, session, inbox, project string eligible bool + originErr, ownerErr error } func (s *provenanceLocalStore) LocalSessionProvenance(string) (string, bool, error) { - return s.owner, s.eligible, nil + return s.owner, s.eligible, s.ownerErr } func (s *provenanceLocalStore) LocalPromptCreationIdentity(string) (string, string, string, bool, error) { - return s.session, s.inbox, s.project, s.eligible, nil + return s.session, s.inbox, s.project, s.eligible, s.originErr } type provenanceTransport struct { @@ -920,6 +921,91 @@ func TestManagerPushStopsBeforeLaterProjectsWhenCanceled(t *testing.T) { } } +type cursorProvenanceStore struct{ *provenanceLocalStore } + +func (s *cursorProvenanceStore) ApplyPulledMutationPreservingSyncState(target string, mutation store.SyncMutation) error { + if err := s.fakeLocalStore.ApplyPulledMutationPreservingSyncState(target, mutation); err != nil { + return err + } + s.syncState.LastPulledSeq = mutation.Seq + return nil +} + +type cursorProvenanceTransport struct { + *provenanceTransport + since []int64 +} + +func (t *cursorProvenanceTransport) PullMutations(since int64, limit int) (*PullMutationsResponse, error) { + t.since = append(t.since, since) + if since >= 7 { + atomic.AddInt32(&t.pullCalls, 1) + return &PullMutationsResponse{}, nil + } + return t.fakeCloudTransport.PullMutations(since, limit) +} + +func TestManagerCycleKeyedPromptBlockAllowsInboundProgress(t *testing.T) { + for _, tc := range []struct { + name string + eligible bool + registerErr, claimErr error + originErr, ownerErr error + pull bool + }{ + {name: "missing origin", pull: true}, + {name: "missing owner", eligible: true, pull: true, ownerErr: nil}, + {name: "registration transport error", eligible: true, registerErr: errors.New("transport down")}, + {name: "claim transport error", eligible: true, claimErr: errors.New("transport down")}, + {name: "origin read error", eligible: true, originErr: errors.New("database unavailable")}, + {name: "owner read error", eligible: true, ownerErr: errors.New("database unavailable")}, + } { + t.Run(tc.name, func(t *testing.T) { + ls := &provenanceLocalStore{fakeLocalStore: newFakeLocalStore(), owner: "alpha", session: "s", inbox: "i", project: "beta", eligible: tc.eligible, originErr: tc.originErr, ownerErr: tc.ownerErr} + if tc.name == "missing owner" { + ls.owner = "" + } + ls.mutations = []store.SyncMutation{{Seq: 1, Entity: store.SyncEntityPrompt, EntityKey: "p", Op: "delete", Project: "beta", Payload: `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"beta"}`}} + tr := &cursorProvenanceTransport{provenanceTransport: &provenanceTransport{fakeCloudTransport: newFakeTransport(), registerErr: tc.registerErr, claimErr: tc.claimErr}} + tr.pullResult = &PullMutationsResponse{Mutations: []PulledMutation{{Seq: 7, Project: "remote", Entity: store.SyncEntitySession, EntityKey: "remote-s", Op: "upsert", Payload: json.RawMessage(`{"id":"remote-s","project":"remote"}`)}}} + mgr := New(&cursorProvenanceStore{provenanceLocalStore: ls}, tr, DefaultConfig()) + mgr.cycle(context.Background()) + if tc.pull { + mgr.cycle(context.Background()) + } + st := mgr.Status() + if tc.pull { + if atomic.LoadInt32(&tr.pullCalls) != 2 || len(ls.appliedMuts) != 1 || ls.syncState.LastPulledSeq != 7 || fmt.Sprint(tr.since) != "[0 7]" { + t.Fatalf("inbound stalled or replayed: pulls=%d applied=%v cursor=%d since=%v", tr.pullCalls, ls.appliedMuts, ls.syncState.LastPulledSeq, tr.since) + } + } else if tr.pullCalls != 0 || len(ls.appliedMuts) != 0 || st.BackoffUntil == nil || st.ConsecutiveFailures != 1 || ls.failureMessage == "" { + t.Fatalf("uncertain preflight must skip pull and record failure: pulls=%d applied=%v status=%+v persisted=%q", tr.pullCalls, ls.appliedMuts, st, ls.failureMessage) + } + if len(tr.attempted) != 0 || len(ls.ackedSeqs) != 0 { + t.Fatalf("blocked prompt pushed/acked: %v %v", tr.attempted, ls.ackedSeqs) + } + if tc.pull && (st.Phase != PhasePushFailed || st.ReasonCode != "prompt_provenance_blocked" || st.LastSyncAt == nil || ls.blockedReason != st.ReasonCode || ls.blockedAfterSuccessCalls != 2 || ls.healthyCalls != 0) { + t.Fatalf("blocked status not retained: %+v persisted=%q healthy=%d", st, ls.blockedReason, ls.healthyCalls) + } + }) + } +} + +func TestManagerCyclePromptBlockMixedWithTransportFailureSkipsPull(t *testing.T) { + ls := &provenanceLocalStore{fakeLocalStore: newFakeLocalStore()} + ls.mutations = []store.SyncMutation{ + {Seq: 1, Entity: store.SyncEntityPrompt, EntityKey: "p", Op: "delete", Project: "beta", Payload: `{"sync_id":"p","session_id":"s","source_inbox_id":"i","project":"beta"}`}, + {Seq: 2, Entity: "obs", EntityKey: "other", Op: "upsert", Project: "alpha"}, + } + tr := &provenanceTransport{fakeCloudTransport: newFakeTransport()} + tr.pushErrByProject = map[string]error{"alpha": errors.New("transport down")} + mgr := New(ls, tr, DefaultConfig()) + mgr.cycle(context.Background()) + if tr.pullCalls != 0 || mgr.Status().Phase != PhasePushFailed || mgr.Status().BackoffUntil == nil || ls.healthyCalls != 0 || len(ls.ackedSeqs) != 0 { + t.Fatalf("mixed failure must skip pull and stay degraded: pulls=%d status=%+v ack=%v", tr.pullCalls, mgr.Status(), ls.ackedSeqs) + } +} + func TestManagerCyclePartialPushFailureSkipsPullAndHealthyState(t *testing.T) { ls := newFakeLocalStore() ls.mutations = []store.SyncMutation{