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
2 changes: 1 addition & 1 deletion docs/codebase/prompt-inbox-provenance.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
56 changes: 46 additions & 10 deletions internal/cloud/autosync/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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 {
Expand All @@ -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)
Expand Down
90 changes: 88 additions & 2 deletions internal/cloud/autosync/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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{
Expand Down
Loading