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
27 changes: 27 additions & 0 deletions cmd/engram/autosync_e2e_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
6 changes: 3 additions & 3 deletions docs/codebase/prompt-inbox-provenance.md
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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

Expand Down 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. 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.
148 changes: 131 additions & 17 deletions internal/cloud/autosync/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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})
Expand All @@ -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 ────────────────────────────────────────────────────────────────────
Expand Down
Loading
Loading