diff --git a/.github/reviews/codex-accepted-objective-continuity.receipt.json b/.github/reviews/codex-accepted-objective-continuity.receipt.json new file mode 100644 index 0000000..e39be76 --- /dev/null +++ b/.github/reviews/codex-accepted-objective-continuity.receipt.json @@ -0,0 +1,4 @@ +{ + "reviewed_tree": "2746af50573c6b47570bae93280b83748df5c0dc", + "program_fingerprint": "3ca3397ff275d89bdb6d5c934b86b51d3cbdfab0ee628c47fe94d1d4f5767155" +} diff --git a/boatstack/cmd/boatstack-reviewer/instance_store_test.go b/boatstack/cmd/boatstack-reviewer/instance_store_test.go index 3b46fe1..94035d6 100644 --- a/boatstack/cmd/boatstack-reviewer/instance_store_test.go +++ b/boatstack/cmd/boatstack-reviewer/instance_store_test.go @@ -11,16 +11,25 @@ import ( // durable, isolated, idempotently resumable multi-instance kernel Store: // reopening constructs a fresh handle over the same on-disk substrate, so // restart reconstruction exercises real file persistence. +func fileInstanceHarness(t testing.TB) conformance.InstanceStoreHarness { + gitDir := t.TempDir() + store := newFileStore(gitDir, "instance-alpha") + return conformance.InstanceStoreHarness{ + Store: store, + Locker: func() kernel.Locker { return directoryLocker{store: store} }, + Reopen: func(testing.TB) kernel.Store { + return newFileStore(gitDir, "instance-alpha") + }, + } +} + func TestFileStoreInstanceConformance(t *testing.T) { - conformance.InstanceStoreConformance{New: func(t testing.TB) conformance.InstanceStoreHarness { - gitDir := t.TempDir() - store := newFileStore(gitDir, "instance-alpha") - return conformance.InstanceStoreHarness{ - Store: store, - Locker: func() kernel.Locker { return directoryLocker{store: store} }, - Reopen: func(testing.TB) kernel.Store { - return newFileStore(gitDir, "instance-alpha") - }, - } - }}.Run(t) + conformance.InstanceStoreConformance{New: fileInstanceHarness}.Run(t) +} + +// TestFileStoreAcceptedObjectiveContinuityConformance runs the same bounded +// accepted-objective laws through fresh fileStore handles, so restart and BFS +// replay use real on-disk records rather than a memory-only adapter. +func TestFileStoreAcceptedObjectiveContinuityConformance(t *testing.T) { + conformance.AcceptedObjectiveContinuityConformance{New: fileInstanceHarness}.Run(t) } diff --git a/boatstack/kernel/conformance/accepted_objective_continuity.go b/boatstack/kernel/conformance/accepted_objective_continuity.go new file mode 100644 index 0000000..78f6147 --- /dev/null +++ b/boatstack/kernel/conformance/accepted_objective_continuity.go @@ -0,0 +1,1024 @@ +package conformance + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "reflect" + "regexp" + "sync" + "testing" + "time" + + "github.com/operatorstack/boatstack/boatstack/kernel" +) + +const ( + continuityBindTransition = "register.bind" + continuityAdvanceTransition = "register.advance" + continuityFinishTransition = "register.finish" + continuityRecoverTransition = "register.recover" + + continuityOpen = "open" + continuityComplete = "complete" + + continuityProducer = "producer-primary" + continuityVerifier = "verifier-independent" +) + +var continuitySemanticID = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]*$`) + +// AcceptedObjectiveContinuityConformance proves that a bounded hierarchical +// objective remains the exact receipt-backed accepted value across rejection, +// recovery, restart, concurrency, policy revision, and final marking. The +// suite owns the domain fixture; adapters supply only the existing instance +// Store, per-instance locker, and reopen boundary. +type AcceptedObjectiveContinuityConformance struct { + New func(testing.TB) InstanceStoreHarness +} + +// Run executes the accepted-objective continuity laws against a fresh +// harness per law. +func (suite AcceptedObjectiveContinuityConformance) Run(t *testing.T) { + t.Helper() + laws := []struct { + name string + law func(*testing.T, AcceptedObjectiveContinuityConformance) + }{ + {"canonical_rejection_recovery_race_and_completion", acceptedContinuityCanonicalLaw}, + {"concurrent_instances_are_isolated", acceptedContinuityIsolationLaw}, + {"restart_matches_uninterrupted_execution", acceptedContinuityRestartLaw}, + {"bounded_bfs_matches_immutable_reference", acceptedContinuityOracleLaw}, + } + for _, law := range laws { + law := law + t.Run(law.name, func(t *testing.T) { law.law(t, suite) }) + } +} + +type continuityNode struct { + ID string `json:"id"` + ParentID string `json:"parent_id,omitempty"` + Status string `json:"status"` + EvidenceFingerprint string `json:"evidence_fingerprint,omitempty"` +} + +type continuityPolicy struct { + Identity string `json:"identity"` + Version string `json:"version"` +} + +type continuityPayload struct { + InstanceID string `json:"instance_id"` + Nodes []continuityNode `json:"nodes"` + ActiveNodeID string `json:"active_node_id,omitempty"` + ProducerID string `json:"producer_id"` + VerifierID string `json:"verifier_id"` + Policy continuityPolicy `json:"policy"` +} + +func cloneContinuityPayload(value continuityPayload) continuityPayload { + copy := value + copy.Nodes = append([]continuityNode(nil), value.Nodes...) + return copy +} + +func continuityFingerprint(label string) string { + digest := sha256.Sum256([]byte(label)) + return hex.EncodeToString(digest[:]) +} + +func validContinuityFingerprint(value string) bool { + if len(value) != sha256.Size*2 { + return false + } + decoded, err := hex.DecodeString(value) + return err == nil && len(decoded) == sha256.Size && value == hex.EncodeToString(decoded) +} + +// validateContinuityPayload is the fixture's domain contract. It deliberately +// operates on candidate content, not supervisory selection or settlement. +func validateContinuityPayload(payload continuityPayload, marked bool) error { + if !continuitySemanticID.MatchString(payload.InstanceID) { + return fmt.Errorf("instance identity is invalid") + } + if len(payload.Nodes) == 0 || len(payload.Nodes) > 3 { + return fmt.Errorf("node count must be between one and three") + } + if !continuitySemanticID.MatchString(payload.ProducerID) || !continuitySemanticID.MatchString(payload.VerifierID) || payload.ProducerID == payload.VerifierID { + return fmt.Errorf("producer and verifier identities must be distinct") + } + if payload.Policy.Identity == "" || payload.Policy.Version == "" { + return fmt.Errorf("policy identity and version must remain opaque non-empty values") + } + + byID := make(map[string]continuityNode, len(payload.Nodes)) + evidence := map[string]bool{} + for index, node := range payload.Nodes { + if !continuitySemanticID.MatchString(node.ID) || (node.ParentID != "" && !continuitySemanticID.MatchString(node.ParentID)) { + return fmt.Errorf("node identity or parent identity is invalid") + } + if index > 0 && payload.Nodes[index-1].ID >= node.ID { + return fmt.Errorf("nodes must be unique and sorted by identity") + } + if node.Status != continuityOpen && node.Status != continuityComplete { + return fmt.Errorf("node %q has an invalid status", node.ID) + } + if node.Status == continuityComplete && !validContinuityFingerprint(node.EvidenceFingerprint) { + return fmt.Errorf("complete node %q lacks exact evidence", node.ID) + } + if node.Status == continuityComplete && evidence[node.EvidenceFingerprint] { + return fmt.Errorf("complete nodes share an evidence fingerprint") + } + if node.Status == continuityComplete { + evidence[node.EvidenceFingerprint] = true + } + if node.Status == continuityOpen && node.EvidenceFingerprint != "" { + return fmt.Errorf("open node %q carries completion evidence", node.ID) + } + if _, exists := byID[node.ID]; exists { + return fmt.Errorf("node identity %q is duplicated", node.ID) + } + byID[node.ID] = node + } + for _, node := range payload.Nodes { + if node.ParentID != "" { + if _, exists := byID[node.ParentID]; !exists { + return fmt.Errorf("node %q names a missing parent", node.ID) + } + } + } + + // Parent walking detects a cycle before the root-count check, so the + // minimal cyclic mutation is classified by the guard it violates. + for _, node := range payload.Nodes { + seen := map[string]bool{} + current := node + depth := 0 + for current.ParentID != "" { + if seen[current.ID] { + return fmt.Errorf("parent hierarchy is cyclic") + } + seen[current.ID] = true + depth++ + if depth > 2 { + return fmt.Errorf("parent hierarchy exceeds depth two") + } + current = byID[current.ParentID] + } + } + roots := 0 + for _, node := range payload.Nodes { + if node.ParentID == "" { + roots++ + } + } + if roots != 1 { + return fmt.Errorf("payload requires exactly one root") + } + + for _, ancestor := range payload.Nodes { + if ancestor.Status != continuityComplete { + continue + } + for _, candidate := range payload.Nodes { + current := candidate + for current.ParentID != "" { + if current.ParentID == ancestor.ID && candidate.Status != continuityComplete { + return fmt.Errorf("premature completion leaves an open descendant") + } + current = byID[current.ParentID] + } + } + } + + if marked { + if payload.ActiveNodeID != "" { + return fmt.Errorf("marked payload retains an active node") + } + for _, node := range payload.Nodes { + if node.Status != continuityComplete { + return fmt.Errorf("marked payload contains an incomplete node") + } + } + return nil + } + if payload.ActiveNodeID == "" { + return fmt.Errorf("unmarked payload lacks one active open node") + } + active, ok := byID[payload.ActiveNodeID] + if !ok { + return fmt.Errorf("active node is missing") + } + if active.Status != continuityOpen { + return fmt.Errorf("active node is not open") + } + activeLineage := map[string]bool{active.ID: true} + for active.ParentID != "" { + active = byID[active.ParentID] + if active.Status != continuityOpen { + return fmt.Errorf("active node has a closed ancestor") + } + activeLineage[active.ID] = true + } + for _, node := range payload.Nodes { + if node.Status == continuityOpen && !activeLineage[node.ID] { + return fmt.Errorf("open node %q is outside the active ancestry", node.ID) + } + } + return nil +} + +// validateContinuityDelta constrains each immutable successor independently +// from the kernel's lifecycle relation. A successor may add one open leaf, +// complete the active leaf and restore its open parent, change only policy, +// or complete the final root through register.finish. +func validateContinuityDelta(prior, next continuityPayload, transition string) error { + if prior.InstanceID != next.InstanceID || prior.ProducerID != next.ProducerID || prior.VerifierID != next.VerifierID { + return fmt.Errorf("successor changed instance or actor identity") + } + priorByID := make(map[string]continuityNode, len(prior.Nodes)) + for _, node := range prior.Nodes { + priorByID[node.ID] = node + } + nextByID := make(map[string]continuityNode, len(next.Nodes)) + for _, node := range next.Nodes { + nextByID[node.ID] = node + } + for id, old := range priorByID { + updated, exists := nextByID[id] + if !exists || updated.ParentID != old.ParentID { + return fmt.Errorf("successor removed or reparented an existing node") + } + if old.Status == continuityComplete && updated != old { + return fmt.Errorf("successor changed a completed node") + } + } + + if transition == continuityFinishTransition { + if len(next.Nodes) != len(prior.Nodes) || next.ActiveNodeID != "" || next.Policy != prior.Policy { + return fmt.Errorf("finish must retain the hierarchy and clear the active node") + } + for id, old := range priorByID { + updated := nextByID[id] + if old.Status == continuityOpen { + if old.ParentID != "" || updated.Status != continuityComplete || !validContinuityFingerprint(updated.EvidenceFingerprint) { + return fmt.Errorf("finish may complete only the final open root") + } + } + } + return nil + } + if transition != continuityAdvanceTransition { + return nil + } + + if len(next.Nodes) == len(prior.Nodes)+1 { + if next.Policy != prior.Policy { + return fmt.Errorf("adding a leaf cannot substitute revision-bound policy") + } + if next.ActiveNodeID == prior.ActiveNodeID { + return fmt.Errorf("new leaf was not made active") + } + for _, old := range prior.Nodes { + if old.ParentID == prior.ActiveNodeID { + return fmt.Errorf("active node already owns its bounded child") + } + } + added, exists := nextByID[next.ActiveNodeID] + if !exists || added.ParentID != prior.ActiveNodeID || added.Status != continuityOpen { + return fmt.Errorf("advance may add only one open child of the active node") + } + for id, old := range priorByID { + if nextByID[id] != old { + return fmt.Errorf("adding a leaf changed an existing node") + } + } + return nil + } + if len(next.Nodes) != len(prior.Nodes) { + return fmt.Errorf("advance changed the hierarchy by more than one node") + } + + oldActive, oldActiveExists := priorByID[prior.ActiveNodeID] + newActive := nextByID[prior.ActiveNodeID] + if oldActiveExists && oldActive.Status == continuityOpen && newActive.Status == continuityComplete { + if next.Policy != prior.Policy { + return fmt.Errorf("leaf completion cannot substitute revision-bound policy") + } + if oldActive.ParentID == "" || next.ActiveNodeID != oldActive.ParentID { + return fmt.Errorf("completed leaf did not restore its open parent") + } + parent := nextByID[oldActive.ParentID] + if parent.Status != continuityOpen { + return fmt.Errorf("completed leaf restored a closed parent") + } + for id, old := range priorByID { + updated := nextByID[id] + if id == oldActive.ID { + if updated.Status != continuityComplete || !validContinuityFingerprint(updated.EvidenceFingerprint) { + return fmt.Errorf("active leaf completion lacks evidence") + } + continue + } + if updated != old { + return fmt.Errorf("leaf completion changed another node") + } + } + return nil + } + if !reflect.DeepEqual(prior.Nodes, next.Nodes) || prior.ActiveNodeID != next.ActiveNodeID { + return fmt.Errorf("policy-only advance changed hierarchy state") + } + completedDescendant := false + for _, node := range prior.Nodes { + if node.ParentID != "" { + if node.Status != continuityComplete { + return fmt.Errorf("policy-only advance precedes descendant completion") + } + completedDescendant = true + } + } + if prior.ActiveNodeID == "" || !completedDescendant { + return fmt.Errorf("policy-only advance requires a restored active root") + } + if prior.Policy == next.Policy { + return fmt.Errorf("advance changed neither hierarchy nor policy") + } + return nil +} + +type continuityCandidate struct { + Objective kernel.Objective + ContentFingerprint string + Content json.RawMessage + Payload continuityPayload +} + +type continuityCandidates struct { + mu sync.Mutex + byObjective map[string]continuityCandidate +} + +func newContinuityCandidates() *continuityCandidates { + return &continuityCandidates{byObjective: map[string]continuityCandidate{}} +} + +func (c *continuityCandidates) stage(id string, revision uint64, payload continuityPayload) (kernel.Objective, error) { + content, err := json.Marshal(payload) + if err != nil { + return kernel.Objective{}, err + } + digest := sha256.Sum256(content) + contentFingerprint := hex.EncodeToString(digest[:]) + objective, err := kernel.NewObjective(id, revision, struct { + CandidateFingerprint string `json:"candidate_fingerprint"` + }{contentFingerprint}) + if err != nil { + return kernel.Objective{}, err + } + c.mu.Lock() + defer c.mu.Unlock() + if existing, ok := c.byObjective[objective.Fingerprint]; ok { + if existing.ContentFingerprint != contentFingerprint || !reflect.DeepEqual(existing.Content, json.RawMessage(content)) { + return kernel.Objective{}, fmt.Errorf("immutable candidate identity already contains different content") + } + return existing.Objective, nil + } + c.byObjective[objective.Fingerprint] = continuityCandidate{ + Objective: objective, ContentFingerprint: contentFingerprint, + Content: append(json.RawMessage(nil), content...), Payload: cloneContinuityPayload(payload), + } + return objective, nil +} + +func (c *continuityCandidates) resolve(objective kernel.Objective) (continuityCandidate, error) { + if err := objective.Validate(); err != nil { + return continuityCandidate{}, err + } + c.mu.Lock() + candidate, ok := c.byObjective[objective.Fingerprint] + c.mu.Unlock() + if !ok || candidate.Objective.ID != objective.ID || + candidate.Objective.Revision != objective.Revision || + candidate.Objective.Fingerprint != objective.Fingerprint || + !reflect.DeepEqual(candidate.Objective.Reference, objective.Reference) { + return continuityCandidate{}, fmt.Errorf("exact immutable candidate is unavailable") + } + digest := sha256.Sum256(candidate.Content) + if hex.EncodeToString(digest[:]) != candidate.ContentFingerprint { + return continuityCandidate{}, fmt.Errorf("immutable candidate fingerprint mismatch") + } + var reference struct { + CandidateFingerprint string `json:"candidate_fingerprint"` + } + if err := json.Unmarshal(objective.Reference, &reference); err != nil || reference.CandidateFingerprint != candidate.ContentFingerprint { + return continuityCandidate{}, fmt.Errorf("objective does not identify the candidate content") + } + candidate.Content = append(json.RawMessage(nil), candidate.Content...) + candidate.Payload = cloneContinuityPayload(candidate.Payload) + return candidate, nil +} + +func (c *continuityCandidates) resolveBinding(binding kernel.ObjectiveBinding) (continuityCandidate, error) { + c.mu.Lock() + candidate, ok := c.byObjective[binding.ObjectiveFingerprint] + c.mu.Unlock() + if !ok || !binding.Matches(candidate.Objective) { + return continuityCandidate{}, fmt.Errorf("durable binding does not identify an immutable candidate") + } + return c.resolve(candidate.Objective) +} + +type continuityCounters struct { + mu sync.Mutex + provisions map[string]int + attempts map[string]int + commits map[string]int + observations map[string]int + executions map[string]int + verifications map[string]int +} + +func newContinuityCounters() *continuityCounters { + return &continuityCounters{ + provisions: map[string]int{}, attempts: map[string]int{}, commits: map[string]int{}, + observations: map[string]int{}, executions: map[string]int{}, verifications: map[string]int{}, + } +} + +type continuityCounts struct { + Provisions int + Attempts int + Observations int + Executions int + Verifications int + Commits int + Revision uint64 + Receipts int +} + +func (c *continuityCounters) snapshot(instanceID string, record kernel.InstanceRecord) continuityCounts { + c.mu.Lock() + defer c.mu.Unlock() + return continuityCounts{ + Provisions: c.provisions[instanceID], Attempts: c.attempts[instanceID], + Observations: c.observations[instanceID], Executions: c.executions[instanceID], + Verifications: c.verifications[instanceID], Commits: c.commits[instanceID], + Revision: record.State.Revision, Receipts: len(record.Receipts), + } +} + +type continuityCountingStore struct { + inner kernel.Store + counters *continuityCounters +} + +func (s continuityCountingStore) Create(ctx context.Context, instanceID string, initial kernel.ControlState) error { + err := s.inner.Create(ctx, instanceID, initial) + if err == nil { + s.counters.mu.Lock() + s.counters.provisions[instanceID]++ + s.counters.mu.Unlock() + } + return err +} + +func (s continuityCountingStore) Load(ctx context.Context, instanceID string) (kernel.InstanceRecord, error) { + return s.inner.Load(ctx, instanceID) +} + +func (s continuityCountingStore) BeginEffect(ctx context.Context, instanceID string, revision uint64, attempt kernel.ControlState) error { + err := s.inner.BeginEffect(ctx, instanceID, revision, attempt) + if err == nil { + s.counters.mu.Lock() + s.counters.attempts[instanceID]++ + s.counters.mu.Unlock() + } + return err +} + +func (s continuityCountingStore) CommitTransition(ctx context.Context, instanceID string, revision uint64, target kernel.ControlState, receipt kernel.Receipt) error { + err := s.inner.CommitTransition(ctx, instanceID, revision, target, receipt) + if err == nil { + s.counters.mu.Lock() + s.counters.commits[instanceID]++ + s.counters.mu.Unlock() + } + return err +} + +type continuityDomain struct { + candidates *continuityCandidates + counters *continuityCounters +} + +func (d *continuityDomain) Observe(_ context.Context, instanceID string) (kernel.Observation, error) { + d.counters.mu.Lock() + d.counters.observations[instanceID]++ + d.counters.mu.Unlock() + // Staging is intentionally absent. It cannot masquerade as accepted state + // or stale a prescription that names an exact immutable objective. + return continuityObservation(instanceID) +} + +func continuityObservation(instanceID string) (kernel.Observation, error) { + return kernel.NewObservation(struct { + InstanceID string `json:"instance_id"` + Register string `json:"register"` + }{instanceID, "bounded-hierarchical-register"}) +} + +func (d *continuityDomain) Admissible(_ context.Context, evaluation kernel.Evaluation) (bool, string, error) { + switch evaluation.Transition.Operation { + case continuityBindTransition, continuityAdvanceTransition, continuityFinishTransition: + if evaluation.Objective == nil { + return false, "objective mutation requires one exact immutable payload", nil + } + if _, err := d.candidates.resolve(*evaluation.Objective); err != nil { + return false, err.Error(), nil + } + return true, "exact immutable payload is available", nil + case continuityRecoverTransition: + return evaluation.State.Recovery != nil, "an unresolved register attempt requires recovery", nil + default: + return false, "unknown bounded register operation", nil + } +} + +func (d *continuityDomain) Verify(_ context.Context, evaluation kernel.Evaluation, effect kernel.Effect, target kernel.Observation) error { + d.counters.mu.Lock() + d.counters.verifications[evaluation.State.InstanceID]++ + d.counters.mu.Unlock() + if target.Fingerprint != evaluation.Observation.Fingerprint { + return fmt.Errorf("register observation changed during side-effect-free execution") + } + if evaluation.Transition.Operation == continuityRecoverTransition { + want := kernel.EffectFact{Facet: "register.payload", Operation: continuityRecoverTransition, Fingerprint: "recovery-preserved-binding"} + if len(effect.Facts) != 1 || effect.Facts[0] != want { + return fmt.Errorf("recovery effect does not prove preserved binding") + } + return nil + } + if evaluation.Objective == nil { + return fmt.Errorf("candidate verification lacks an objective") + } + candidate, err := d.candidates.resolve(*evaluation.Objective) + if err != nil { + return err + } + if candidate.Payload.InstanceID != evaluation.State.InstanceID { + return fmt.Errorf("candidate belongs to a different control instance") + } + if candidate.Payload.VerifierID != continuityVerifier || candidate.Payload.ProducerID == candidate.Payload.VerifierID { + return fmt.Errorf("candidate lacks an independent verifier identity") + } + marked := evaluation.Transition.ID == continuityFinishTransition + if evaluation.State.ObjectiveBinding != nil && evaluation.Transition.ID == continuityAdvanceTransition { + prior, priorErr := d.candidates.resolveBinding(*evaluation.State.ObjectiveBinding) + if priorErr != nil { + return priorErr + } + // Parent restoration is a transition property. Check it before the + // static active-node rule so this failure remains observable. + if oldActive := nodeByID(prior.Payload, prior.Payload.ActiveNodeID); oldActive != nil && oldActive.ParentID != "" { + if nextActive := nodeByID(candidate.Payload, oldActive.ID); nextActive != nil && oldActive.Status == continuityOpen && nextActive.Status == continuityComplete && candidate.Payload.ActiveNodeID != oldActive.ParentID { + return fmt.Errorf("completed leaf did not restore its open parent") + } + } + } + if err := validateContinuityPayload(candidate.Payload, marked); err != nil { + return err + } + if evaluation.State.ObjectiveBinding != nil { + prior, priorErr := d.candidates.resolveBinding(*evaluation.State.ObjectiveBinding) + if priorErr != nil { + return priorErr + } + if err := validateContinuityDelta(prior.Payload, candidate.Payload, evaluation.Transition.ID); err != nil { + return err + } + } + want := kernel.EffectFact{Facet: "register.payload", Operation: evaluation.Transition.Operation, Fingerprint: candidate.ContentFingerprint} + if len(effect.Facts) != 1 || effect.Facts[0] != want { + return fmt.Errorf("effect does not identify the exact candidate content") + } + return nil +} + +func nodeByID(payload continuityPayload, id string) *continuityNode { + for index := range payload.Nodes { + if payload.Nodes[index].ID == id { + copy := payload.Nodes[index] + return © + } + } + return nil +} + +type continuityOperator struct { + candidates *continuityCandidates + counters *continuityCounters +} + +func (o *continuityOperator) Execute(_ context.Context, operation kernel.Operation) (kernel.Effect, error) { + o.counters.mu.Lock() + o.counters.executions[operation.InstanceID]++ + o.counters.mu.Unlock() + if operation.Transition.Operation == continuityRecoverTransition { + return kernel.Effect{Facts: []kernel.EffectFact{{Facet: "register.payload", Operation: continuityRecoverTransition, Fingerprint: "recovery-preserved-binding"}}}, nil + } + if operation.Objective == nil { + return kernel.Effect{}, fmt.Errorf("register operation lacks an objective") + } + candidate, err := o.candidates.resolve(*operation.Objective) + if err != nil { + return kernel.Effect{}, err + } + return kernel.Effect{Facts: []kernel.EffectFact{{Facet: "register.payload", Operation: operation.Transition.Operation, Fingerprint: candidate.ContentFingerprint}}}, nil +} + +type continuityCapabilities struct{} + +func (continuityCapabilities) RequiredCapabilities(transition kernel.Transition) ([]kernel.Capability, error) { + switch transition.Operation { + case continuityBindTransition: + return []kernel.Capability{"objective.bind", "register.write"}, nil + case continuityAdvanceTransition, continuityFinishTransition: + return []kernel.Capability{"objective.advance", "register.write"}, nil + case continuityRecoverTransition: + return []kernel.Capability{"register.recover"}, nil + default: + return nil, fmt.Errorf("unclassified bounded register operation %q", transition.Operation) + } +} + +func compileContinuityProgram() (kernel.Program, error) { + contract, err := kernel.Fingerprint(struct { + Name string `json:"name"` + MaxNodes int `json:"max_nodes"` + MaxDepth int `json:"max_depth"` + Verifier string `json:"verifier"` + Unmarked string `json:"unmarked"` + Marked string `json:"marked"` + }{"bounded-hierarchical-register", 3, 2, continuityVerifier, "one-active-open-node", "all-complete-no-active-node"}) + if err != nil { + return kernel.Program{}, err + } + mutation := func(id, source, target string, lifecycle kernel.ObjectiveMutation, priority int) kernel.Transition { + capability, _ := lifecycle.Capability() + return kernel.Transition{ + ID: id, SourceModes: []string{source}, TargetMode: target, + ObjectiveScope: kernel.ObjectiveNone, ObjectiveMutation: lifecycle, + RequiredCapabilities: []kernel.Capability{capability, "register.write"}, + OwnedFacets: []string{"register.payload", "supervisor.objective"}, + Operation: id, SelectionRank: 1, Selection: kernel.SelectionExplicitOnly, Priority: priority, + } + } + return kernel.CompileDomainProgram( + "bounded-hierarchical-register", "1.0.0", "kernel-v1", contract, + "supervising", []string{"marked"}, []kernel.Transition{ + mutation(continuityBindTransition, "supervising", "supervising", kernel.BindInitialObjective, 1), + mutation(continuityAdvanceTransition, "supervising", "supervising", kernel.AdvanceObjective, 2), + mutation(continuityFinishTransition, "supervising", "marked", kernel.AdvanceObjective, 3), + { + ID: continuityRecoverTransition, SourceModes: []string{"supervising"}, TargetMode: "supervising", + ObjectiveScope: kernel.ObjectiveOptionalPreserve, ObjectiveMutation: kernel.PreserveObjective, + RequiredCapabilities: []kernel.Capability{"register.recover"}, OwnedFacets: []string{"register.payload"}, + Operation: continuityRecoverTransition, SelectionRank: 1, Selection: kernel.SelectionImplicit, Priority: 10, + Recovers: []string{continuityBindTransition, continuityAdvanceTransition, continuityFinishTransition, continuityRecoverTransition}, + }, + }, + ) +} + +type continuityAccepted struct { + Binding kernel.ObjectiveBinding + CandidateFingerprint string + Payload continuityPayload + Receipt kernel.Receipt +} + +type continuityFixture struct { + harness InstanceStoreHarness + store kernel.Store + program kernel.Program + authority kernel.Authority + clock settlementClock + candidates *continuityCandidates + counters *continuityCounters + domain *continuityDomain + operator *continuityOperator + classifier continuityCapabilities + authorityFingerprint string +} + +func newContinuityFixture(t testing.TB, harness InstanceStoreHarness) *continuityFixture { + t.Helper() + if harness.Store == nil || harness.Locker == nil || harness.Reopen == nil { + t.Fatal("control-law accepted-objective-continuity-harness: store, locker, and reopen ports are required") + } + program, err := compileContinuityProgram() + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-program: %v", err) + } + now := time.Date(2026, 8, 23, 12, 0, 0, 0, time.UTC) + candidates := newContinuityCandidates() + counters := newContinuityCounters() + domain := &continuityDomain{candidates: candidates, counters: counters} + authority := kernel.Authority{Receipts: []kernel.AuthorityReceipt{{ + ID: "continuity-authority", Subject: "bounded-register-fixture", Fingerprint: "continuity-authority", + Capabilities: []kernel.Capability{"objective.bind", "objective.advance", "register.write", "register.recover"}, + IssuedAt: now.Add(-time.Minute), ExpiresAt: now.Add(time.Hour), + }}} + authorityFingerprint, err := kernel.Fingerprint(authority.Receipts) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-authority: %v", err) + } + return &continuityFixture{ + harness: harness, store: continuityCountingStore{inner: harness.Store, counters: counters}, + program: program, clock: settlementClock{now: now}, candidates: candidates, counters: counters, + domain: domain, operator: &continuityOperator{candidates: candidates, counters: counters}, classifier: continuityCapabilities{}, + authority: authority, authorityFingerprint: authorityFingerprint, + } +} + +func (f *continuityFixture) runtime(t testing.TB, store kernel.Store, locker kernel.Locker) kernel.Runtime { + t.Helper() + runtime, err := kernel.NewRuntime(f.program, f.domain, f.operator, f.classifier, store, locker, f.clock) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-runtime: %v", err) + } + return runtime +} + +func (f *continuityFixture) currentRuntime(t testing.TB) kernel.Runtime { + return f.runtime(t, f.store, f.harness.Locker()) +} + +func (f *continuityFixture) reopenRuntime(t testing.TB) kernel.Runtime { + t.Helper() + f.store = continuityCountingStore{inner: f.harness.Reopen(t), counters: f.counters} + return f.runtime(t, f.store, f.harness.Locker()) +} + +func (f *continuityFixture) stage(t testing.TB, objectiveID string, revision uint64, payload continuityPayload) kernel.Objective { + t.Helper() + objective, err := f.candidates.stage(objectiveID, revision, payload) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-stage: %v", err) + } + return objective +} + +func (f *continuityFixture) resolve(t testing.TB, runtime kernel.Runtime, instanceID, transition string, objective *kernel.Objective) (kernel.ResolveRequest, kernel.Prescription) { + t.Helper() + request := kernel.ResolveRequest{InstanceID: instanceID, Objective: objective, Authority: f.authority, Requested: transition} + resolution, err := runtime.Resolve(context.Background(), request) + if err != nil || resolution.Decision.Kind != kernel.Prescribed || resolution.Prescription == nil { + t.Fatalf("control-law accepted-objective-continuity-resolve: transition=%s decision=%#v error=%v", transition, resolution.Decision, err) + } + return request, *resolution.Prescription +} + +func (f *continuityFixture) apply(t testing.TB, runtime kernel.Runtime, instanceID, transition string, objective *kernel.Objective) kernel.Receipt { + t.Helper() + request, prescription := f.resolve(t, runtime, instanceID, transition, objective) + receipt, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: prescription}) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-apply: transition=%s error=%v", transition, err) + } + return receipt +} + +func (f *continuityFixture) accepted(instanceID string) (continuityAccepted, error) { + record, err := f.store.Load(context.Background(), instanceID) + if err != nil { + return continuityAccepted{}, err + } + if err := record.Validate(instanceID); err != nil { + return continuityAccepted{}, err + } + if record.State.ObjectiveBinding == nil { + return continuityAccepted{}, fmt.Errorf("durable objective binding is absent") + } + final, err := f.validateLineage(instanceID, record.Receipts) + if err != nil { + return continuityAccepted{}, err + } + if final == nil || !reflect.DeepEqual(final.ResultObjectiveBinding, record.State.ObjectiveBinding) { + return continuityAccepted{}, fmt.Errorf("complete receipt lineage does not prove the durable binding") + } + candidate, err := f.candidates.resolveBinding(*record.State.ObjectiveBinding) + if err != nil { + return continuityAccepted{}, err + } + if err := validateContinuityPayload(candidate.Payload, record.State.Mode == "marked"); err != nil { + return continuityAccepted{}, fmt.Errorf("accepted payload is invalid: %w", err) + } + return continuityAccepted{ + Binding: *record.State.ObjectiveBinding, CandidateFingerprint: candidate.ContentFingerprint, + Payload: candidate.Payload, Receipt: *final, + }, nil +} + +type continuityReceiptContract struct { + Program kernel.ProgramIdentity + AuthorityFingerprint string + ObservationFingerprint string + CommittedAt time.Time +} + +func continuityPrescriptionID(receipt kernel.Receipt) (string, error) { + bindingFingerprint, err := kernel.Fingerprint(receipt.PriorObjectiveBinding) + if err != nil { + return "", err + } + freshness, err := kernel.NewFreshness( + receipt.InstanceID, receipt.PriorStateRevision, receipt.Program.Fingerprint, + receipt.PriorObservation, bindingFingerprint, receipt.AuthorityFingerprint, + ) + if err != nil { + return "", err + } + prescription := kernel.Prescription{ + SchemaVersion: kernel.PrescriptionSchemaVersion, TransitionID: receipt.TransitionID, + Freshness: freshness, ObjectiveMutation: receipt.ObjectiveMutation, + ExpectedObjectiveBinding: receipt.PriorObjectiveBinding, RequestedObjectiveBinding: receipt.RequestedObjectiveBinding, + RequiredCapabilities: append([]kernel.Capability(nil), receipt.Capabilities...), + } + digest, err := kernel.Fingerprint(prescription) + if err != nil { + return "", err + } + return "prx-" + digest, nil +} + +func (f *continuityFixture) validateLineage(instanceID string, receipts []kernel.Receipt) (*kernel.Receipt, error) { + observation, err := continuityObservation(instanceID) + if err != nil { + return nil, err + } + return validateContinuityLineage(instanceID, receipts, f.candidates, continuityReceiptContract{ + Program: f.program.Identity(), AuthorityFingerprint: f.authorityFingerprint, + ObservationFingerprint: observation.Fingerprint, CommittedAt: f.clock.Now().UTC(), + }) +} + +func validateContinuityLineage(instanceID string, receipts []kernel.Receipt, candidates *continuityCandidates, contract continuityReceiptContract) (*kernel.Receipt, error) { + var chain *kernel.ObjectiveBinding + var final *kernel.Receipt + var priorCandidate *continuityCandidate + seenRevision := map[string]bool{} + var lastResult uint64 + markedSeen := false + for index := range receipts { + receipt := receipts[index] + if err := receipt.Validate(); err != nil { + return nil, fmt.Errorf("committed receipt %q is invalid: %w", receipt.ID, err) + } + if receipt.InstanceID != instanceID || receipt.PriorStateRevision < lastResult { + return nil, fmt.Errorf("committed receipt %q breaks instance or state lineage", receipt.ID) + } + if receipt.Program != contract.Program || receipt.AuthorityFingerprint != contract.AuthorityFingerprint || + receipt.PriorObservation != contract.ObservationFingerprint || receipt.ResultObservation != contract.ObservationFingerprint || + !receipt.CommittedAt.Equal(contract.CommittedAt) { + return nil, fmt.Errorf("committed receipt %q substitutes program, authority, observation, or time facts", receipt.ID) + } + expectedPrescriptionID, err := continuityPrescriptionID(receipt) + if err != nil || receipt.PrescriptionID != expectedPrescriptionID { + return nil, fmt.Errorf("committed receipt %q does not settle its exact prescription", receipt.ID) + } + lastResult = receipt.ResultStateRevision + if markedSeen { + return nil, fmt.Errorf("committed receipt %q follows marked completion", receipt.ID) + } + if receipt.ObjectiveMutation == kernel.PreserveObjective { + if !reflect.DeepEqual(receipt.PriorObjectiveBinding, chain) || !reflect.DeepEqual(receipt.ResultObjectiveBinding, chain) { + return nil, fmt.Errorf("preserve receipt %q changes accepted objective lineage", receipt.ID) + } + wantCapabilities := []kernel.Capability{"register.recover"} + wantEffect := kernel.EffectFact{Facet: "register.payload", Operation: continuityRecoverTransition, Fingerprint: "recovery-preserved-binding"} + if receipt.TransitionID != continuityRecoverTransition || !reflect.DeepEqual(receipt.Capabilities, wantCapabilities) || + len(receipt.Effects) != 1 || receipt.Effects[0] != wantEffect { + return nil, fmt.Errorf("preserve receipt %q does not prove exact recovery facts", receipt.ID) + } + continue + } + if !reflect.DeepEqual(receipt.PriorObjectiveBinding, chain) || receipt.ResultObjectiveBinding == nil { + return nil, fmt.Errorf("committed receipt %q does not extend accepted objective lineage", receipt.ID) + } + key := fmt.Sprintf("%s@%d", receipt.ResultObjectiveBinding.ObjectiveID, receipt.ResultObjectiveBinding.ObjectiveRevision) + if seenRevision[key] { + return nil, fmt.Errorf("objective revision %s has more than one committed receipt", key) + } + seenRevision[key] = true + candidate, err := candidates.resolveBinding(*receipt.ResultObjectiveBinding) + if err != nil { + return nil, fmt.Errorf("receipt %q candidate: %w", receipt.ID, err) + } + marked := receipt.TransitionID == continuityFinishTransition + if candidate.Payload.InstanceID != instanceID { + return nil, fmt.Errorf("receipt %q payload belongs to a different instance", receipt.ID) + } + if err := validateContinuityPayload(candidate.Payload, marked); err != nil { + return nil, fmt.Errorf("receipt %q payload: %w", receipt.ID, err) + } + wantEffect := kernel.EffectFact{Facet: "register.payload", Operation: receipt.TransitionID, Fingerprint: candidate.ContentFingerprint} + if len(receipt.Effects) != 1 || receipt.Effects[0] != wantEffect { + return nil, fmt.Errorf("receipt %q does not prove the exact candidate content", receipt.ID) + } + if priorCandidate != nil { + if err := validateContinuityDelta(priorCandidate.Payload, candidate.Payload, receipt.TransitionID); err != nil { + return nil, fmt.Errorf("receipt %q payload delta: %w", receipt.ID, err) + } + } + switch receipt.TransitionID { + case continuityBindTransition: + if receipt.ObjectiveMutation != kernel.BindInitialObjective || !reflect.DeepEqual(receipt.Capabilities, []kernel.Capability{"objective.bind", "register.write"}) { + return nil, fmt.Errorf("bind receipt %q declares the wrong lifecycle relation", receipt.ID) + } + case continuityAdvanceTransition, continuityFinishTransition: + if receipt.ObjectiveMutation != kernel.AdvanceObjective || !reflect.DeepEqual(receipt.Capabilities, []kernel.Capability{"objective.advance", "register.write"}) { + return nil, fmt.Errorf("advance receipt %q declares the wrong lifecycle relation", receipt.ID) + } + default: + return nil, fmt.Errorf("receipt %q uses an unknown binding transition", receipt.ID) + } + chain = receipt.ResultObjectiveBinding + final = &receipts[index] + candidateCopy := candidate + priorCandidate = &candidateCopy + markedSeen = marked + } + return final, nil +} + +func (f *continuityFixture) record(t testing.TB, instanceID string) kernel.InstanceRecord { + t.Helper() + record, err := f.store.Load(context.Background(), instanceID) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-load: %v", err) + } + return record +} + +func (f *continuityFixture) counts(t testing.TB, instanceID string) continuityCounts { + t.Helper() + return f.counters.snapshot(instanceID, f.record(t, instanceID)) +} + +func baseContinuityPayload(instanceID string) continuityPayload { + return continuityPayload{ + InstanceID: instanceID, + Nodes: []continuityNode{{ID: "root", Status: continuityOpen}}, ActiveNodeID: "root", + ProducerID: continuityProducer, VerifierID: continuityVerifier, + Policy: continuityPolicy{Identity: "policy-alpha", Version: "v1"}, + } +} + +func childContinuityPayload(instanceID string) continuityPayload { + payload := baseContinuityPayload(instanceID) + payload.Nodes = []continuityNode{ + {ID: "leaf", ParentID: "root", Status: continuityOpen}, + {ID: "root", Status: continuityOpen}, + } + payload.ActiveNodeID = "leaf" + return payload +} + +func restoredContinuityPayload(instanceID string) continuityPayload { + payload := childContinuityPayload(instanceID) + payload.Nodes[0].Status = continuityComplete + payload.Nodes[0].EvidenceFingerprint = continuityFingerprint(instanceID + ":leaf-complete") + payload.ActiveNodeID = "root" + return payload +} + +func policyContinuityPayload(instanceID, identity, version string) continuityPayload { + payload := restoredContinuityPayload(instanceID) + payload.Policy = continuityPolicy{Identity: identity, Version: version} + return payload +} + +func markedContinuityPayload(instanceID, identity, version string) continuityPayload { + payload := policyContinuityPayload(instanceID, identity, version) + payload.Nodes[1].Status = continuityComplete + payload.Nodes[1].EvidenceFingerprint = continuityFingerprint(instanceID + ":root-complete") + payload.ActiveNodeID = "" + return payload +} + +func assertContinuityAccepted(t testing.TB, fixture *continuityFixture, instanceID string, objective kernel.Objective, payload continuityPayload) continuityAccepted { + t.Helper() + accepted, err := fixture.accepted(instanceID) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-reader: %v", err) + } + if !accepted.Binding.Matches(objective) || !reflect.DeepEqual(accepted.Payload, payload) { + t.Fatalf("control-law accepted-objective-continuity-reader: accepted=%#v want objective=%#v payload=%#v", accepted, objective, payload) + } + return accepted +} diff --git a/boatstack/kernel/conformance/accepted_objective_continuity_laws.go b/boatstack/kernel/conformance/accepted_objective_continuity_laws.go new file mode 100644 index 0000000..c9c86e7 --- /dev/null +++ b/boatstack/kernel/conformance/accepted_objective_continuity_laws.go @@ -0,0 +1,398 @@ +package conformance + +import ( + "context" + "fmt" + "reflect" + "testing" + + "github.com/operatorstack/boatstack/boatstack/kernel" +) + +const ( + continuityInstanceA = "register-a" + continuityInstanceB = "register-b" + continuityObjectiveA = "objective-a" + continuityObjectiveB = "objective-b" +) + +func freshContinuityFixture(t testing.TB, suite AcceptedObjectiveContinuityConformance) *continuityFixture { + t.Helper() + if suite.New == nil { + t.Fatal("control-law accepted-objective-continuity-harness: fresh fixture factory is required") + } + return newContinuityFixture(t, suite.New(t)) +} + +func attemptContinuityRecovery(runtime kernel.Runtime, fixture *continuityFixture, instanceID, transition string, objective *kernel.Objective) (kernel.ResolveRequest, kernel.Prescription, error) { + request := kernel.ResolveRequest{InstanceID: instanceID, Objective: objective, Authority: fixture.authority, Requested: transition} + resolution, err := runtime.Resolve(context.Background(), request) + if err != nil { + return kernel.ResolveRequest{}, kernel.Prescription{}, err + } + if resolution.Decision.Kind != kernel.Prescribed || resolution.Prescription == nil { + return kernel.ResolveRequest{}, kernel.Prescription{}, fmt.Errorf("transition %q was not prescribed: %#v", transition, resolution.Decision) + } + _, err = runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: *resolution.Prescription}) + if !kernel.IsRecoveryRequired(err) { + return request, *resolution.Prescription, fmt.Errorf("invalid transition returned %v, want recovery-required", err) + } + return request, *resolution.Prescription, nil +} + +func assertContinuityRecordEqual(t testing.TB, label string, want, got kernel.InstanceRecord) { + t.Helper() + if !reflect.DeepEqual(want, got) { + t.Fatalf("control-law accepted-objective-continuity-%s: record changed\nwant=%#v\n got=%#v", label, want, got) + } +} + +func acceptedContinuityCanonicalLaw(t *testing.T, suite AcceptedObjectiveContinuityConformance) { + fixture := freshContinuityFixture(t, suite) + runtime := fixture.currentRuntime(t) + if _, err := runtime.Provision(context.Background(), continuityInstanceA); err != nil { + t.Fatalf("control-law accepted-objective-continuity-provision: %v", err) + } + + a1Payload := baseContinuityPayload(continuityInstanceA) + a1 := fixture.stage(t, continuityObjectiveA, 1, a1Payload) + a1Receipt := fixture.apply(t, runtime, continuityInstanceA, continuityBindTransition, &a1) + if a1Receipt.ObjectiveMutation != kernel.BindInitialObjective { + t.Fatalf("control-law accepted-objective-continuity-bind: receipt declares %q", a1Receipt.ObjectiveMutation) + } + assertContinuityAccepted(t, fixture, continuityInstanceA, a1, a1Payload) + + // Staging and read-only resolution remain observations. Neither is a + // candidate acceptance path. + beforeProjection := fixture.record(t, continuityInstanceA) + stagedOnly := baseContinuityPayload(continuityInstanceA) + stagedOnly.Policy = continuityPolicy{Identity: "policy-staged-only", Version: "v2"} + fixture.stage(t, continuityObjectiveA, 2, stagedOnly) + a2Payload := childContinuityPayload(continuityInstanceA) + a2 := fixture.stage(t, continuityObjectiveA, 2, a2Payload) + fixture.resolve(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a2) + afterProjection := fixture.record(t, continuityInstanceA) + assertContinuityRecordEqual(t, "projection-is-observation", beforeProjection, afterProjection) + assertContinuityAccepted(t, fixture, continuityInstanceA, a1, a1Payload) + + fixture.apply(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a2) + assertContinuityAccepted(t, fixture, continuityInstanceA, a2, a2Payload) + + // A parent completion while its child remains open crosses the effect + // boundary, fails independent verification, and leaves A@2 accepted. + invalidA3Payload := childContinuityPayload(continuityInstanceA) + invalidA3Payload.Nodes[1].Status = continuityComplete + invalidA3Payload.Nodes[1].EvidenceFingerprint = continuityFingerprint("invalid-parent-completion") + invalidA3 := fixture.stage(t, continuityObjectiveA, 3, invalidA3Payload) + beforeRejection := fixture.counts(t, continuityInstanceA) + if _, _, err := attemptContinuityRecovery(runtime, fixture, continuityInstanceA, continuityAdvanceTransition, &invalidA3); err != nil { + t.Fatalf("control-law accepted-objective-continuity-rejection: %v", err) + } + afterRejection := fixture.counts(t, continuityInstanceA) + if afterRejection.Attempts != beforeRejection.Attempts+1 || + afterRejection.Executions != beforeRejection.Executions+1 || + afterRejection.Verifications != beforeRejection.Verifications+1 || + afterRejection.Commits != beforeRejection.Commits || + afterRejection.Receipts != beforeRejection.Receipts || + afterRejection.Revision != beforeRejection.Revision+1 { + t.Fatalf("control-law accepted-objective-continuity-rejection-effects: before=%#v after=%#v", beforeRejection, afterRejection) + } + recoveryRecord := fixture.record(t, continuityInstanceA) + if recoveryRecord.State.Recovery == nil || !recoveryRecord.State.ObjectiveBinding.Matches(a2) { + t.Fatalf("control-law accepted-objective-continuity-recovery-state: %#v", recoveryRecord.State) + } + assertContinuityAccepted(t, fixture, continuityInstanceA, a2, a2Payload) + + recoveryRequest, recoveryPrescription := fixture.resolve(t, runtime, continuityInstanceA, continuityRecoverTransition, nil) + recoveryReceipt, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: recoveryRequest, Prescription: recoveryPrescription}) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-recover: %v", err) + } + if recoveryReceipt.ObjectiveMutation != kernel.PreserveObjective || + !reflect.DeepEqual(recoveryReceipt.PriorObjectiveBinding, recoveryReceipt.ResultObjectiveBinding) || + !recoveryReceipt.ResultObjectiveBinding.Matches(a2) { + t.Fatalf("control-law accepted-objective-continuity-recover-preserves: %#v", recoveryReceipt) + } + assertContinuityAccepted(t, fixture, continuityInstanceA, a2, a2Payload) + beforeRecoveryRetry := fixture.counts(t, continuityInstanceA) + reconciledRecovery, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: recoveryRequest, Prescription: recoveryPrescription}) + if err != nil || reconciledRecovery.ID != recoveryReceipt.ID { + t.Fatalf("control-law accepted-objective-continuity-recover-retry: receipt=%#v error=%v", reconciledRecovery, err) + } + if afterRecoveryRetry := fixture.counts(t, continuityInstanceA); !reflect.DeepEqual(beforeRecoveryRetry, afterRecoveryRetry) { + t.Fatalf("control-law accepted-objective-continuity-recover-retry-effects: before=%#v after=%#v", beforeRecoveryRetry, afterRecoveryRetry) + } + + // Capture a valid A@3 successor before another A@3 settles. It must stay + // stale across a fresh runtime and Store handle. + staleA3Payload := restoredContinuityPayload(continuityInstanceA) + staleA3Payload.Nodes[0].EvidenceFingerprint = continuityFingerprint("stale-a3-evidence") + staleA3 := fixture.stage(t, continuityObjectiveA, 3, staleA3Payload) + staleRequest, stalePrescription := fixture.resolve(t, runtime, continuityInstanceA, continuityAdvanceTransition, &staleA3) + a3Payload := restoredContinuityPayload(continuityInstanceA) + a3 := fixture.stage(t, continuityObjectiveA, 3, a3Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a3) + runtime = fixture.reopenRuntime(t) + assertContinuityAccepted(t, fixture, continuityInstanceA, a3, a3Payload) + beforeStale := fixture.counts(t, continuityInstanceA) + if _, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: staleRequest, Prescription: stalePrescription}); !kernel.IsStale(err) { + t.Fatalf("control-law accepted-objective-continuity-stale-successor: error=%v", err) + } + afterStale := fixture.counts(t, continuityInstanceA) + if afterStale.Executions != beforeStale.Executions || afterStale.Verifications != beforeStale.Verifications || + afterStale.Commits != beforeStale.Commits || afterStale.Receipts != beforeStale.Receipts || afterStale.Revision != beforeStale.Revision { + t.Fatalf("control-law accepted-objective-continuity-stale-effects: before=%#v after=%#v", beforeStale, afterStale) + } + + // A second objective can bind and advance while A remains byte-for-byte + // unchanged in the same substrate. + aBeforeB := fixture.record(t, continuityInstanceA) + if _, err := runtime.Provision(context.Background(), continuityInstanceB); err != nil { + t.Fatalf("control-law accepted-objective-continuity-provision-b: %v", err) + } + b1Payload := baseContinuityPayload(continuityInstanceB) + b1 := fixture.stage(t, continuityObjectiveB, 1, b1Payload) + fixture.apply(t, runtime, continuityInstanceB, continuityBindTransition, &b1) + b2Payload := childContinuityPayload(continuityInstanceB) + b2 := fixture.stage(t, continuityObjectiveB, 2, b2Payload) + fixture.apply(t, runtime, continuityInstanceB, continuityAdvanceTransition, &b2) + assertContinuityAccepted(t, fixture, continuityInstanceB, b2, b2Payload) + assertContinuityRecordEqual(t, "b-progress-isolated", aBeforeB, fixture.record(t, continuityInstanceA)) + + // Both A@4 successors are valid and revision-bound. Per-instance locking + // plus state freshness permits exactly one commit. + a4LeftPayload := policyContinuityPayload(continuityInstanceA, "policy-left", "v2") + a4RightPayload := policyContinuityPayload(continuityInstanceA, "policy-right", "v2") + a4Left := fixture.stage(t, continuityObjectiveA, 4, a4LeftPayload) + a4Right := fixture.stage(t, continuityObjectiveA, 4, a4RightPayload) + leftRequest, leftPrescription := fixture.resolve(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a4Left) + rightRequest, rightPrescription := fixture.resolve(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a4Right) + beforeRace := fixture.counts(t, continuityInstanceA) + type raceResult struct { + left bool + receipt kernel.Receipt + err error + } + start := make(chan struct{}) + results := make(chan raceResult, 2) + for _, item := range []struct { + left bool + request kernel.ResolveRequest + prescription kernel.Prescription + }{{true, leftRequest, leftPrescription}, {false, rightRequest, rightPrescription}} { + item := item + go func() { + <-start + receipt, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: item.request, Prescription: item.prescription}) + results <- raceResult{left: item.left, receipt: receipt, err: err} + }() + } + close(start) + first, second := <-results, <-results + winners := []raceResult{} + losers := []raceResult{} + for _, result := range []raceResult{first, second} { + if result.err == nil { + winners = append(winners, result) + } else { + losers = append(losers, result) + } + } + if len(winners) != 1 || len(losers) != 1 { + t.Fatalf("control-law accepted-objective-continuity-race: winners=%#v losers=%#v", winners, losers) + } + afterRace := fixture.counts(t, continuityInstanceA) + if afterRace.Attempts != beforeRace.Attempts+1 || afterRace.Executions != beforeRace.Executions+1 || + afterRace.Verifications != beforeRace.Verifications+1 || afterRace.Commits != beforeRace.Commits+1 || + afterRace.Receipts != beforeRace.Receipts+1 || afterRace.Revision != beforeRace.Revision+2 { + t.Fatalf("control-law accepted-objective-continuity-race-effects: before=%#v after=%#v loser=%v", beforeRace, afterRace, losers[0].err) + } + loserRequest, loserPrescription := rightRequest, rightPrescription + if losers[0].left { + loserRequest, loserPrescription = leftRequest, leftPrescription + } + if _, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: loserRequest, Prescription: loserPrescription}); !kernel.IsStale(err) { + t.Fatalf("control-law accepted-objective-continuity-race-loser-stale: concurrent_error=%v retry_error=%v", losers[0].err, err) + } + winner := winners[0] + winnerObjective, winnerPayload := a4Right, a4RightPayload + winnerRequest, winnerPrescription := rightRequest, rightPrescription + if winner.left { + winnerObjective, winnerPayload = a4Left, a4LeftPayload + winnerRequest, winnerPrescription = leftRequest, leftPrescription + } + acceptedA4 := assertContinuityAccepted(t, fixture, continuityInstanceA, winnerObjective, winnerPayload) + if acceptedA4.Payload.Policy == a3Payload.Policy || acceptedA4.Binding.ObjectiveRevision != 4 || acceptedA4.Receipt.RequestedObjectiveBinding.ObjectiveFingerprint != winnerObjective.Fingerprint { + t.Fatalf("control-law accepted-objective-continuity-policy-revision: %#v", acceptedA4) + } + + // A lost response retry reconciles the original receipt and has no + // observation, execution, verification, commit, revision, or receipt effect. + beforeRetry := fixture.counts(t, continuityInstanceA) + retried, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: winnerRequest, Prescription: winnerPrescription}) + if err != nil || retried.ID != winner.receipt.ID { + t.Fatalf("control-law accepted-objective-continuity-retry: receipt=%#v error=%v", retried, err) + } + afterRetry := fixture.counts(t, continuityInstanceA) + if !reflect.DeepEqual(beforeRetry, afterRetry) { + t.Fatalf("control-law accepted-objective-continuity-retry-effects: before=%#v after=%#v", beforeRetry, afterRetry) + } + + // The same A prescription cannot settle B, even though both are valid + // instances in the same store. + aBeforeCross := fixture.record(t, continuityInstanceA) + bBeforeCross := fixture.record(t, continuityInstanceB) + crossRequest := winnerRequest + crossRequest.InstanceID = continuityInstanceB + if _, err := runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: crossRequest, Prescription: winnerPrescription}); !kernel.IsStale(err) { + t.Fatalf("control-law accepted-objective-continuity-cross-instance: error=%v", err) + } + assertContinuityRecordEqual(t, "cross-instance-a", aBeforeCross, fixture.record(t, continuityInstanceA)) + assertContinuityRecordEqual(t, "cross-instance-b", bBeforeCross, fixture.record(t, continuityInstanceB)) + + a5Payload := markedContinuityPayload(continuityInstanceA, winnerPayload.Policy.Identity, winnerPayload.Policy.Version) + a5 := fixture.stage(t, continuityObjectiveA, 5, a5Payload) + a5Receipt := fixture.apply(t, runtime, continuityInstanceA, continuityFinishTransition, &a5) + acceptedA5 := assertContinuityAccepted(t, fixture, continuityInstanceA, a5, a5Payload) + finalRecord := fixture.record(t, continuityInstanceA) + if finalRecord.State.Mode != "marked" || finalRecord.State.Recovery != nil || a5Receipt.ObjectiveMutation != kernel.AdvanceObjective || acceptedA5.Binding.ObjectiveRevision != 5 { + t.Fatalf("control-law accepted-objective-continuity-marked: record=%#v receipt=%#v accepted=%#v", finalRecord, a5Receipt, acceptedA5) + } + for _, node := range acceptedA5.Payload.Nodes { + if node.Status != continuityComplete || !validContinuityFingerprint(node.EvidenceFingerprint) { + t.Fatalf("control-law accepted-objective-continuity-marked-evidence: node=%#v", node) + } + } + if acceptedA5.Payload.ActiveNodeID != "" { + t.Fatalf("control-law accepted-objective-continuity-marked-active: %#v", acceptedA5.Payload) + } + + counts := fixture.counts(t, continuityInstanceA) + if counts.Provisions != 1 || counts.Attempts != 7 || counts.Executions != 7 || counts.Verifications != 7 || + counts.Commits != 6 || counts.Receipts != 6 || counts.Revision != 14 || counts.Observations <= counts.Executions { + t.Fatalf("control-law accepted-objective-continuity-accounting: %#v", counts) + } + if _, err := fixture.validateLineage(continuityInstanceA, finalRecord.Receipts); err != nil { + t.Fatalf("control-law accepted-objective-continuity-lineage: %v", err) + } +} + +func acceptedContinuityIsolationLaw(t *testing.T, suite AcceptedObjectiveContinuityConformance) { + fixture := freshContinuityFixture(t, suite) + runtime := fixture.currentRuntime(t) + for _, instanceID := range []string{continuityInstanceA, continuityInstanceB} { + if _, err := runtime.Provision(context.Background(), instanceID); err != nil { + t.Fatalf("control-law accepted-objective-continuity-isolation-provision: %v", err) + } + } + a1Payload, b1Payload := baseContinuityPayload(continuityInstanceA), baseContinuityPayload(continuityInstanceB) + a1 := fixture.stage(t, continuityObjectiveA, 1, a1Payload) + b1 := fixture.stage(t, continuityObjectiveB, 1, b1Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityBindTransition, &a1) + fixture.apply(t, runtime, continuityInstanceB, continuityBindTransition, &b1) + a2Payload, b2Payload := childContinuityPayload(continuityInstanceA), childContinuityPayload(continuityInstanceB) + a2 := fixture.stage(t, continuityObjectiveA, 2, a2Payload) + b2 := fixture.stage(t, continuityObjectiveB, 2, b2Payload) + aRequest, aPrescription := fixture.resolve(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a2) + bRequest, bPrescription := fixture.resolve(t, runtime, continuityInstanceB, continuityAdvanceTransition, &b2) + + start := make(chan struct{}) + errs := make(chan error, 2) + for _, apply := range []kernel.ApplyRequest{ + {ResolveRequest: aRequest, Prescription: aPrescription}, + {ResolveRequest: bRequest, Prescription: bPrescription}, + } { + apply := apply + go func() { + <-start + _, err := runtime.Apply(context.Background(), apply) + errs <- err + }() + } + close(start) + if err := <-errs; err != nil { + t.Fatalf("control-law accepted-objective-continuity-isolation-a: %v", err) + } + if err := <-errs; err != nil { + t.Fatalf("control-law accepted-objective-continuity-isolation-b: %v", err) + } + acceptedA := assertContinuityAccepted(t, fixture, continuityInstanceA, a2, a2Payload) + acceptedB := assertContinuityAccepted(t, fixture, continuityInstanceB, b2, b2Payload) + if acceptedA.Binding == acceptedB.Binding || acceptedA.Payload.InstanceID == acceptedB.Payload.InstanceID { + t.Fatalf("control-law accepted-objective-continuity-isolation-alias: a=%#v b=%#v", acceptedA, acceptedB) + } + for instanceID, record := range map[string]kernel.InstanceRecord{ + continuityInstanceA: fixture.record(t, continuityInstanceA), + continuityInstanceB: fixture.record(t, continuityInstanceB), + } { + if len(record.Receipts) != 2 || record.State.Revision != 5 { + t.Fatalf("control-law accepted-objective-continuity-isolation-history: instance=%s record=%#v", instanceID, record) + } + for _, receipt := range record.Receipts { + if receipt.InstanceID != instanceID { + t.Fatalf("control-law accepted-objective-continuity-isolation-routing: instance=%s receipt=%#v", instanceID, receipt) + } + } + } +} + +type continuityTrajectoryResult struct { + Record kernel.InstanceRecord + Accepted continuityAccepted + Counts continuityCounts +} + +func executeContinuityTrajectory(t testing.TB, suite AcceptedObjectiveContinuityConformance, reopen bool) continuityTrajectoryResult { + t.Helper() + fixture := freshContinuityFixture(t, suite) + runtime := fixture.currentRuntime(t) + instanceID, objectiveID := "register-equivalence", "objective-equivalence" + if _, err := runtime.Provision(context.Background(), instanceID); err != nil { + t.Fatal(err) + } + reopenRuntime := func() { + if reopen { + runtime = fixture.reopenRuntime(t) + } + } + a1Payload := baseContinuityPayload(instanceID) + a1 := fixture.stage(t, objectiveID, 1, a1Payload) + fixture.apply(t, runtime, instanceID, continuityBindTransition, &a1) + reopenRuntime() + a2Payload := childContinuityPayload(instanceID) + a2 := fixture.stage(t, objectiveID, 2, a2Payload) + fixture.apply(t, runtime, instanceID, continuityAdvanceTransition, &a2) + reopenRuntime() + invalid := childContinuityPayload(instanceID) + invalid.Nodes[1].Status = continuityComplete + invalid.Nodes[1].EvidenceFingerprint = continuityFingerprint("equivalence-invalid-parent") + invalidObjective := fixture.stage(t, objectiveID, 3, invalid) + if _, _, err := attemptContinuityRecovery(runtime, fixture, instanceID, continuityAdvanceTransition, &invalidObjective); err != nil { + t.Fatal(err) + } + reopenRuntime() + fixture.apply(t, runtime, instanceID, continuityRecoverTransition, nil) + reopenRuntime() + a3Payload := restoredContinuityPayload(instanceID) + a3 := fixture.stage(t, objectiveID, 3, a3Payload) + fixture.apply(t, runtime, instanceID, continuityAdvanceTransition, &a3) + reopenRuntime() + a4Payload := policyContinuityPayload(instanceID, "policy-equivalence", "v2") + a4 := fixture.stage(t, objectiveID, 4, a4Payload) + fixture.apply(t, runtime, instanceID, continuityAdvanceTransition, &a4) + reopenRuntime() + a5Payload := markedContinuityPayload(instanceID, a4Payload.Policy.Identity, a4Payload.Policy.Version) + a5 := fixture.stage(t, objectiveID, 5, a5Payload) + fixture.apply(t, runtime, instanceID, continuityFinishTransition, &a5) + reopenRuntime() + accepted := assertContinuityAccepted(t, fixture, instanceID, a5, a5Payload) + return continuityTrajectoryResult{Record: fixture.record(t, instanceID), Accepted: accepted, Counts: fixture.counts(t, instanceID)} +} + +func acceptedContinuityRestartLaw(t *testing.T, suite AcceptedObjectiveContinuityConformance) { + uninterrupted := executeContinuityTrajectory(t, suite, false) + restarted := executeContinuityTrajectory(t, suite, true) + if !reflect.DeepEqual(uninterrupted, restarted) { + t.Fatalf("control-law accepted-objective-continuity-restart-equivalence:\nuninterrupted=%#v\nrestarted=%#v", uninterrupted, restarted) + } +} diff --git a/boatstack/kernel/conformance/accepted_objective_continuity_oracle.go b/boatstack/kernel/conformance/accepted_objective_continuity_oracle.go new file mode 100644 index 0000000..c0e5215 --- /dev/null +++ b/boatstack/kernel/conformance/accepted_objective_continuity_oracle.go @@ -0,0 +1,1010 @@ +package conformance + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "reflect" + "sort" + "strings" + "testing" + + "github.com/operatorstack/boatstack/boatstack/kernel" +) + +const continuityOracleMaxSteps = 6 + +type continuityOracleOperation struct { + Kind string + InstanceID string +} + +const ( + oracleProvision = "provision" + oracleBind = "bind-root" + oracleOpenChild = "open-child" + oracleCompleteChild = "complete-child" + oracleChangePolicy = "change-policy" + oracleFinish = "finish-root" + oracleRecover = "recover" + oracleCyclicParents = "cyclic-parents" + oracleMissingActive = "missing-active-node" + oracleClosedActive = "closed-active-node" + oracleParentRestoration = "failed-parent-restoration" + oraclePrematureComplete = "premature-completion" + oracleSelfVerification = "self-verification" + oracleLatestStagedRead = "latest-staged-read" + oracleAliasedInstances = "aliased-instances" + oracleStaleOverwrite = "stale-overwrite" + oraclePolicySubstitution = "same-revision-policy-substitution" +) + +var continuityOracleGuards = []string{ + oracleCyclicParents, + oracleMissingActive, + oracleClosedActive, + oracleParentRestoration, + oraclePrematureComplete, + oracleSelfVerification, + oracleLatestStagedRead, + oracleAliasedInstances, + oracleStaleOverwrite, + oraclePolicySubstitution, +} + +func (o continuityOracleOperation) String() string { + if o.InstanceID == "" { + return o.Kind + } + return o.InstanceID + ":" + o.Kind +} + +func (o continuityOracleOperation) semanticStep() bool { + // Staging and accepted reads are observations. Provisioning, recovery, + // and every requested transition (including a refusal) are semantic steps. + return o.Kind != oracleLatestStagedRead +} + +type continuityOracleReceipt struct { + Transition string + ObjectiveMutation kernel.ObjectiveMutation + ObjectiveRevision uint64 + ResultRevision uint64 +} + +type continuityReferenceInstance struct { + Exists bool + ObjectiveID string + ObjectiveRevision uint64 + Payload continuityPayload + StateRevision uint64 + Mode string + Recovery bool + Receipts []continuityOracleReceipt + HasStale bool +} + +type continuityReferenceWorld struct { + Instances map[string]continuityReferenceInstance +} + +func newContinuityReferenceWorld() continuityReferenceWorld { + return continuityReferenceWorld{Instances: map[string]continuityReferenceInstance{}} +} + +func cloneContinuityReferenceWorld(world continuityReferenceWorld) continuityReferenceWorld { + copy := newContinuityReferenceWorld() + for id, instance := range world.Instances { + instance.Payload = cloneContinuityPayload(instance.Payload) + instance.Receipts = append([]continuityOracleReceipt(nil), instance.Receipts...) + copy.Instances[id] = instance + } + return copy +} + +func continuityOracleObjectiveID(instanceID string) string { + if instanceID == continuityInstanceA { + return continuityObjectiveA + } + return continuityObjectiveB +} + +// referenceContinuityTransition is immutable and independent of Boatstack's +// selector and settlement implementation. It models only the accepted +// register contract, state-revision effects, and receipt semantics. +func referenceContinuityTransition(world continuityReferenceWorld, operation continuityOracleOperation) (continuityReferenceWorld, bool, error) { + next := cloneContinuityReferenceWorld(world) + instance := next.Instances[operation.InstanceID] + rejected := false + appendReceipt := func(transition string, mutation kernel.ObjectiveMutation) { + instance.StateRevision += 2 + instance.Receipts = append(instance.Receipts, continuityOracleReceipt{ + Transition: transition, ObjectiveMutation: mutation, + ObjectiveRevision: instance.ObjectiveRevision, ResultRevision: instance.StateRevision, + }) + } + + switch operation.Kind { + case oracleProvision: + if instance.Exists { + return world, false, fmt.Errorf("reference provision duplicates %q", operation.InstanceID) + } + instance = continuityReferenceInstance{Exists: true, StateRevision: 1, Mode: "supervising"} + case oracleBind: + if !instance.Exists || instance.ObjectiveRevision != 0 || instance.Recovery { + return world, false, fmt.Errorf("reference bind is inadmissible") + } + instance.ObjectiveID = continuityOracleObjectiveID(operation.InstanceID) + instance.ObjectiveRevision = 1 + instance.Payload = referenceBasePayload(operation.InstanceID) + instance.HasStale = true + appendReceipt(continuityBindTransition, kernel.BindInitialObjective) + case oracleOpenChild: + if !referenceCanOpenChild(instance) { + return world, false, fmt.Errorf("reference open-child is inadmissible") + } + instance.ObjectiveRevision++ + instance.Payload = referenceAddedChildPayload(instance.Payload) + instance.HasStale = true + appendReceipt(continuityAdvanceTransition, kernel.AdvanceObjective) + case oracleCompleteChild: + if !referenceCanCompleteChild(instance) { + return world, false, fmt.Errorf("reference complete-child is inadmissible") + } + instance.ObjectiveRevision++ + instance.Payload = referenceCompletedActivePayload(instance.Payload) + instance.HasStale = true + appendReceipt(continuityAdvanceTransition, kernel.AdvanceObjective) + case oracleChangePolicy: + if !referenceCanChangePolicy(instance) { + return world, false, fmt.Errorf("reference policy change is inadmissible") + } + instance.ObjectiveRevision++ + instance.Payload = cloneContinuityPayload(instance.Payload) + instance.Payload.Policy = continuityPolicy{Identity: "policy-oracle", Version: "v2"} + instance.HasStale = true + appendReceipt(continuityAdvanceTransition, kernel.AdvanceObjective) + case oracleFinish: + if !referenceCanFinish(instance) { + return world, false, fmt.Errorf("reference finish is inadmissible") + } + instance.ObjectiveRevision++ + instance.Payload = cloneContinuityPayload(instance.Payload) + for index := range instance.Payload.Nodes { + if instance.Payload.Nodes[index].ParentID == "" { + instance.Payload.Nodes[index].Status = continuityComplete + instance.Payload.Nodes[index].EvidenceFingerprint = continuityFingerprint(operation.InstanceID + ":root-complete") + } + } + instance.Payload.ActiveNodeID = "" + instance.Mode = "marked" + instance.HasStale = true + appendReceipt(continuityFinishTransition, kernel.AdvanceObjective) + case oracleRecover: + if !instance.Exists || !instance.Recovery { + return world, false, fmt.Errorf("reference recovery is inadmissible") + } + instance.Recovery = false + appendReceipt(continuityRecoverTransition, kernel.PreserveObjective) + case oracleCyclicParents, oracleMissingActive, oracleClosedActive, + oracleParentRestoration, oraclePrematureComplete, oracleSelfVerification: + if !instance.Exists || instance.Recovery || instance.Mode == "marked" { + return world, false, fmt.Errorf("reference invalid mutation is outside its minimal guard state") + } + instance.StateRevision++ + instance.Recovery = true + rejected = true + case oracleLatestStagedRead, oraclePolicySubstitution, oracleStaleOverwrite: + if !instance.Exists || instance.ObjectiveRevision == 0 || instance.Recovery || instance.Mode == "marked" { + return world, false, fmt.Errorf("reference observation/refusal is inadmissible") + } + if operation.Kind == oracleStaleOverwrite && !instance.HasStale { + return world, false, fmt.Errorf("reference stale overwrite lacks a captured successor") + } + rejected = true + case oracleAliasedInstances: + a := next.Instances[continuityInstanceA] + b := next.Instances[continuityInstanceB] + if !a.Exists || !b.Exists || a.ObjectiveRevision == 0 || b.ObjectiveRevision == 0 || a.Recovery || b.Recovery || a.Mode == "marked" { + return world, false, fmt.Errorf("reference alias attempt is inadmissible") + } + rejected = true + default: + return world, false, fmt.Errorf("unknown reference operation %q", operation.Kind) + } + next.Instances[operation.InstanceID] = instance + return next, rejected, nil +} + +func referenceBasePayload(instanceID string) continuityPayload { + return continuityPayload{ + InstanceID: instanceID, + Nodes: []continuityNode{{ID: "root", Status: continuityOpen}}, + ActiveNodeID: "root", ProducerID: continuityProducer, VerifierID: continuityVerifier, + Policy: continuityPolicy{Identity: "policy-alpha", Version: "v1"}, + } +} + +func referenceAddedChildPayload(current continuityPayload) continuityPayload { + next := cloneContinuityPayload(current) + childID := "leaf" + if current.ActiveNodeID == "leaf" { + childID = "twig" + } + next.Nodes = append(next.Nodes, continuityNode{ID: childID, ParentID: current.ActiveNodeID, Status: continuityOpen}) + sort.Slice(next.Nodes, func(i, j int) bool { return next.Nodes[i].ID < next.Nodes[j].ID }) + next.ActiveNodeID = childID + return next +} + +func referenceCompletedActivePayload(current continuityPayload) continuityPayload { + next := cloneContinuityPayload(current) + for index := range next.Nodes { + if next.Nodes[index].ID == current.ActiveNodeID { + next.Nodes[index].Status = continuityComplete + next.Nodes[index].EvidenceFingerprint = continuityFingerprint(current.InstanceID + ":" + current.ActiveNodeID + "-complete") + next.ActiveNodeID = next.Nodes[index].ParentID + break + } + } + return next +} + +func referenceCanOpenChild(instance continuityReferenceInstance) bool { + if !instance.Exists || instance.Recovery || instance.Mode != "supervising" || + instance.ObjectiveRevision == 0 || len(instance.Payload.Nodes) >= 3 || instance.Payload.ActiveNodeID == "" { + return false + } + for _, node := range instance.Payload.Nodes { + if node.ParentID == instance.Payload.ActiveNodeID { + return false + } + } + return true +} + +func referenceCanCompleteChild(instance continuityReferenceInstance) bool { + if !instance.Exists || instance.Recovery || instance.Mode != "supervising" || instance.Payload.ActiveNodeID == "" || instance.Payload.ActiveNodeID == "root" { + return false + } + for _, node := range instance.Payload.Nodes { + if node.ID == instance.Payload.ActiveNodeID && node.Status == continuityOpen { + for _, descendant := range instance.Payload.Nodes { + if descendant.ParentID == node.ID && descendant.Status != continuityComplete { + return false + } + } + return true + } + } + return false +} + +func referenceCanChangePolicy(instance continuityReferenceInstance) bool { + if !instance.Exists || instance.Recovery || instance.Mode != "supervising" || len(instance.Payload.Nodes) < 2 || instance.Payload.ActiveNodeID != "root" || instance.Payload.Policy.Version != "v1" { + return false + } + for _, node := range instance.Payload.Nodes { + if node.ParentID != "" && node.Status != continuityComplete { + return false + } + } + return true +} + +func referenceCanFinish(instance continuityReferenceInstance) bool { + if !instance.Exists || instance.Recovery || instance.Mode != "supervising" || len(instance.Payload.Nodes) < 2 || instance.Payload.ActiveNodeID != "root" || instance.Payload.Policy.Version != "v2" { + return false + } + for _, node := range instance.Payload.Nodes { + if node.ParentID != "" && node.Status != continuityComplete { + return false + } + } + return true +} + +func enumerateContinuityOperations(world continuityReferenceWorld) []continuityOracleOperation { + var operations []continuityOracleOperation + for _, instanceID := range []string{continuityInstanceA, continuityInstanceB} { + instance := world.Instances[instanceID] + if !instance.Exists { + operations = append(operations, continuityOracleOperation{oracleProvision, instanceID}) + continue + } + if instance.Recovery { + operations = append(operations, continuityOracleOperation{oracleRecover, instanceID}) + continue + } + if instance.Mode == "marked" { + continue + } + if instance.ObjectiveRevision == 0 { + operations = append(operations, continuityOracleOperation{oracleBind, instanceID}) + if instanceID == continuityInstanceA { + operations = append(operations, continuityOracleOperation{oracleCyclicParents, instanceID}) + } + continue + } + if referenceCanOpenChild(instance) { + operations = append(operations, continuityOracleOperation{oracleOpenChild, instanceID}) + } + if referenceCanCompleteChild(instance) { + operations = append(operations, continuityOracleOperation{oracleCompleteChild, instanceID}) + } + if referenceCanChangePolicy(instance) { + operations = append(operations, continuityOracleOperation{oracleChangePolicy, instanceID}) + } + if referenceCanFinish(instance) { + operations = append(operations, continuityOracleOperation{oracleFinish, instanceID}) + } + if instanceID != continuityInstanceA { + continue + } + operations = append(operations, + continuityOracleOperation{oracleLatestStagedRead, instanceID}, + continuityOracleOperation{oraclePolicySubstitution, instanceID}, + continuityOracleOperation{oracleSelfVerification, instanceID}, + ) + if instance.HasStale { + operations = append(operations, continuityOracleOperation{oracleStaleOverwrite, instanceID}) + } + if referenceCanOpenChild(instance) { + operations = append(operations, + continuityOracleOperation{oracleMissingActive, instanceID}, + continuityOracleOperation{oracleClosedActive, instanceID}, + ) + } + if referenceCanCompleteChild(instance) { + operations = append(operations, + continuityOracleOperation{oracleParentRestoration, instanceID}, + continuityOracleOperation{oraclePrematureComplete, instanceID}, + ) + } + } + a, b := world.Instances[continuityInstanceA], world.Instances[continuityInstanceB] + if a.Exists && b.Exists && a.ObjectiveRevision > 0 && b.ObjectiveRevision > 0 && !a.Recovery && !b.Recovery && a.Mode != "marked" { + operations = append(operations, continuityOracleOperation{oracleAliasedInstances, continuityInstanceA}) + } + sort.Slice(operations, func(i, j int) bool { return operations[i].String() < operations[j].String() }) + return operations +} + +type continuityOracleRunner struct { + fixture *continuityFixture + runtime kernel.Runtime + stale map[string]kernel.ApplyRequest +} + +func newContinuityOracleRunner(t testing.TB, suite AcceptedObjectiveContinuityConformance) *continuityOracleRunner { + fixture := freshContinuityFixture(t, suite) + return &continuityOracleRunner{fixture: fixture, runtime: fixture.currentRuntime(t), stale: map[string]kernel.ApplyRequest{}} +} + +func (r *continuityOracleRunner) resolve(instanceID, transition string, objective *kernel.Objective) (kernel.ResolveRequest, kernel.Prescription, error) { + request := kernel.ResolveRequest{InstanceID: instanceID, Objective: objective, Authority: r.fixture.authority, Requested: transition} + resolution, err := r.runtime.Resolve(context.Background(), request) + if err != nil { + return kernel.ResolveRequest{}, kernel.Prescription{}, err + } + if resolution.Decision.Kind != kernel.Prescribed || resolution.Prescription == nil { + return kernel.ResolveRequest{}, kernel.Prescription{}, fmt.Errorf("transition %q was not prescribed: %#v", transition, resolution.Decision) + } + return request, *resolution.Prescription, nil +} + +func (r *continuityOracleRunner) current(instanceID string) (kernel.InstanceRecord, *continuityAccepted, error) { + record, err := r.fixture.store.Load(context.Background(), instanceID) + if err != nil { + return kernel.InstanceRecord{}, nil, err + } + if record.State.ObjectiveBinding == nil { + return record, nil, nil + } + accepted, err := r.fixture.accepted(instanceID) + if err != nil { + return kernel.InstanceRecord{}, nil, err + } + return record, &accepted, nil +} + +func (r *continuityOracleRunner) stage(instanceID string, revision uint64, payload continuityPayload) (kernel.Objective, error) { + return r.fixture.candidates.stage(continuityOracleObjectiveID(instanceID), revision, payload) +} + +func (r *continuityOracleRunner) step(operation continuityOracleOperation) (bool, error) { + if operation.Kind == oracleProvision { + _, err := r.runtime.Provision(context.Background(), operation.InstanceID) + return false, err + } + if operation.Kind == oracleRecover { + request, prescription, err := r.resolve(operation.InstanceID, continuityRecoverTransition, nil) + if err != nil { + return false, err + } + _, err = r.runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: prescription}) + return false, err + } + _, accepted, err := r.current(operation.InstanceID) + if err != nil { + return false, err + } + + switch operation.Kind { + case oracleLatestStagedRead: + before := *accepted + payload, _, err := actualSuccessorPayload(operation.InstanceID, accepted.Payload, oracleNextValidKind(accepted.Payload)) + if err != nil { + return false, err + } + payload.Policy.Identity += "-staged" + if _, err := r.stage(operation.InstanceID, accepted.Binding.ObjectiveRevision+1, payload); err != nil { + return false, err + } + after, err := r.fixture.accepted(operation.InstanceID) + if err != nil { + return false, err + } + if !reflect.DeepEqual(before, after) { + return false, fmt.Errorf("latest staged candidate masqueraded as accepted state") + } + return true, nil + case oraclePolicySubstitution: + payload := cloneContinuityPayload(accepted.Payload) + payload.Policy = continuityPolicy{Identity: "policy-substituted", Version: payload.Policy.Version} + objective, err := r.stage(operation.InstanceID, accepted.Binding.ObjectiveRevision, payload) + if err != nil { + return false, err + } + request := kernel.ResolveRequest{InstanceID: operation.InstanceID, Objective: &objective, Authority: r.fixture.authority, Requested: continuityAdvanceTransition} + resolution, err := r.runtime.Resolve(context.Background(), request) + if err != nil { + return false, err + } + if resolution.Decision.Kind == kernel.Prescribed || resolution.Prescription != nil { + return false, fmt.Errorf("same-revision policy substitution was prescribed") + } + return true, nil + case oracleStaleOverwrite: + stale, ok := r.stale[operation.InstanceID] + if !ok { + return false, fmt.Errorf("no stale successor was captured") + } + _, err := r.runtime.Apply(context.Background(), stale) + if !kernel.IsStale(err) { + return false, fmt.Errorf("stale successor returned %v", err) + } + return true, nil + case oracleAliasedInstances: + if accepted == nil { + return false, fmt.Errorf("alias source has no accepted objective") + } + kind := oracleNextValidKind(accepted.Payload) + payload, transition, err := actualSuccessorPayload(continuityInstanceA, accepted.Payload, kind) + if err != nil { + return false, err + } + objective, err := r.stage(continuityInstanceA, accepted.Binding.ObjectiveRevision+1, payload) + if err != nil { + return false, err + } + request, prescription, err := r.resolve(continuityInstanceA, transition, &objective) + if err != nil { + return false, err + } + request.InstanceID = continuityInstanceB + _, err = r.runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: prescription}) + if !kernel.IsStale(err) { + return false, fmt.Errorf("cross-instance prescription returned %v", err) + } + return true, nil + } + + if isOracleInvalidMutation(operation.Kind) { + payload, transition, err := actualInvalidPayload(operation.Kind, operation.InstanceID, accepted) + if err != nil { + return false, err + } + revision := uint64(1) + if accepted != nil { + revision = accepted.Binding.ObjectiveRevision + 1 + } + objective, err := r.stage(operation.InstanceID, revision, payload) + if err != nil { + return false, err + } + request, prescription, err := r.resolve(operation.InstanceID, transition, &objective) + if err != nil { + return false, err + } + _, err = r.runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: prescription}) + if !kernel.IsRecoveryRequired(err) { + return false, fmt.Errorf("guard mutation %q returned %v", operation.Kind, err) + } + return true, nil + } + + payload, transition, err := actualSuccessorPayload(operation.InstanceID, continuityPayload{}, operation.Kind) + revision := uint64(1) + if accepted != nil { + payload, transition, err = actualSuccessorPayload(operation.InstanceID, accepted.Payload, operation.Kind) + revision = accepted.Binding.ObjectiveRevision + 1 + } + if err != nil { + return false, err + } + objective, err := r.stage(operation.InstanceID, revision, payload) + if err != nil { + return false, err + } + // Resolve one alternative before the winner. After the winner commits, + // this exact prescription becomes the minimal stale-overwrite probe. + alternate := alternateContinuityPayload(operation.InstanceID, payload, operation.Kind) + alternateObjective, err := r.stage(operation.InstanceID, revision, alternate) + if err != nil { + return false, err + } + staleRequest, stalePrescription, err := r.resolve(operation.InstanceID, transition, &alternateObjective) + if err != nil { + return false, err + } + request, prescription, err := r.resolve(operation.InstanceID, transition, &objective) + if err != nil { + return false, err + } + _, err = r.runtime.Apply(context.Background(), kernel.ApplyRequest{ResolveRequest: request, Prescription: prescription}) + if err != nil { + return false, err + } + r.stale[operation.InstanceID] = kernel.ApplyRequest{ResolveRequest: staleRequest, Prescription: stalePrescription} + return false, nil +} + +func actualSuccessorPayload(instanceID string, current continuityPayload, kind string) (continuityPayload, string, error) { + switch kind { + case oracleBind: + return baseContinuityPayload(instanceID), continuityBindTransition, nil + case oracleOpenChild: + payload := cloneContinuityPayload(current) + childID := "leaf" + if current.ActiveNodeID == "leaf" { + childID = "twig" + } + payload.Nodes = append(payload.Nodes, continuityNode{ID: childID, ParentID: current.ActiveNodeID, Status: continuityOpen}) + sort.Slice(payload.Nodes, func(i, j int) bool { return payload.Nodes[i].ID < payload.Nodes[j].ID }) + payload.ActiveNodeID = childID + return payload, continuityAdvanceTransition, nil + case oracleCompleteChild: + payload := cloneContinuityPayload(current) + for index := range payload.Nodes { + if payload.Nodes[index].ID == current.ActiveNodeID { + payload.Nodes[index].Status = continuityComplete + payload.Nodes[index].EvidenceFingerprint = continuityFingerprint(instanceID + ":" + current.ActiveNodeID + "-complete") + payload.ActiveNodeID = payload.Nodes[index].ParentID + break + } + } + return payload, continuityAdvanceTransition, nil + case oracleChangePolicy: + payload := cloneContinuityPayload(current) + payload.Policy = continuityPolicy{Identity: "policy-oracle", Version: "v2"} + return payload, continuityAdvanceTransition, nil + case oracleFinish: + payload := cloneContinuityPayload(current) + for index := range payload.Nodes { + if payload.Nodes[index].ParentID == "" { + payload.Nodes[index].Status = continuityComplete + payload.Nodes[index].EvidenceFingerprint = continuityFingerprint(instanceID + ":root-complete") + } + } + payload.ActiveNodeID = "" + return payload, continuityFinishTransition, nil + default: + return continuityPayload{}, "", fmt.Errorf("operation %q has no valid successor", kind) + } +} + +func alternateContinuityPayload(instanceID string, payload continuityPayload, kind string) continuityPayload { + alternate := cloneContinuityPayload(payload) + switch kind { + case oracleBind: + alternate.Policy.Identity = "policy-bind-alternate" + case oracleOpenChild: + for index := range alternate.Nodes { + if alternate.Nodes[index].ID == alternate.ActiveNodeID { + alternate.Nodes[index].ID += "-alt" + alternate.ActiveNodeID = alternate.Nodes[index].ID + break + } + } + sort.Slice(alternate.Nodes, func(i, j int) bool { return alternate.Nodes[i].ID < alternate.Nodes[j].ID }) + case oracleCompleteChild: + for index := range alternate.Nodes { + if alternate.Nodes[index].Status == continuityComplete && alternate.Nodes[index].EvidenceFingerprint != "" { + alternate.Nodes[index].EvidenceFingerprint = continuityFingerprint(instanceID + ":" + alternate.Nodes[index].ID + "-alternate") + break + } + } + case oracleChangePolicy: + alternate.Policy.Identity = "policy-oracle-alternate" + case oracleFinish: + for index := range alternate.Nodes { + if alternate.Nodes[index].ParentID == "" { + alternate.Nodes[index].EvidenceFingerprint = continuityFingerprint(instanceID + ":root-alternate") + } + } + } + return alternate +} + +func oracleNextValidKind(payload continuityPayload) string { + switch { + case payload.ActiveNodeID != "root": + return oracleCompleteChild + case len(payload.Nodes) == 1: + return oracleOpenChild + case payload.Policy.Version == "v1": + return oracleChangePolicy + default: + return oracleFinish + } +} + +func isOracleInvalidMutation(kind string) bool { + switch kind { + case oracleCyclicParents, oracleMissingActive, oracleClosedActive, + oracleParentRestoration, oraclePrematureComplete, oracleSelfVerification: + return true + } + return false +} + +func actualInvalidPayload(kind, instanceID string, accepted *continuityAccepted) (continuityPayload, string, error) { + if kind == oracleCyclicParents { + return continuityPayload{ + InstanceID: instanceID, + Nodes: []continuityNode{ + {ID: "leaf", ParentID: "root", Status: continuityOpen}, + {ID: "root", ParentID: "leaf", Status: continuityOpen}, + }, + ActiveNodeID: "leaf", ProducerID: continuityProducer, VerifierID: continuityVerifier, + Policy: continuityPolicy{Identity: "policy-alpha", Version: "v1"}, + }, continuityBindTransition, nil + } + if accepted == nil { + return continuityPayload{}, "", fmt.Errorf("mutation %q requires an accepted payload", kind) + } + payload := cloneContinuityPayload(accepted.Payload) + payload.Policy.Identity += "-guard" + switch kind { + case oracleMissingActive: + payload.ActiveNodeID = "" + case oracleClosedActive: + for index := range payload.Nodes { + if payload.Nodes[index].ID == payload.ActiveNodeID { + payload.Nodes[index].Status = continuityComplete + payload.Nodes[index].EvidenceFingerprint = continuityFingerprint("closed-active") + } + } + case oracleParentRestoration: + for index := range payload.Nodes { + if payload.Nodes[index].ID == payload.ActiveNodeID { + payload.Nodes[index].Status = continuityComplete + payload.Nodes[index].EvidenceFingerprint = continuityFingerprint("parent-not-restored") + } + } + payload.ActiveNodeID = "" + case oraclePrematureComplete: + for index := range payload.Nodes { + if payload.Nodes[index].ParentID == "" { + payload.Nodes[index].Status = continuityComplete + payload.Nodes[index].EvidenceFingerprint = continuityFingerprint("premature-root") + } + } + case oracleSelfVerification: + payload.VerifierID = payload.ProducerID + default: + return continuityPayload{}, "", fmt.Errorf("unknown invalid mutation %q", kind) + } + return payload, continuityAdvanceTransition, nil +} + +type continuityActualReceipt struct { + ID string + Transition string + ObjectiveMutation kernel.ObjectiveMutation + ObjectiveRevision uint64 + ResultRevision uint64 +} + +type continuityActualInstance struct { + Exists bool + ObjectiveID string + ObjectiveRevision uint64 + ObjectiveFingerprint string + Payload continuityPayload + CandidateFingerprint string + StateRevision uint64 + Mode string + RecoveryPrescriptionID string + RecoveryTransitionID string + Receipts []continuityActualReceipt +} + +type continuityActualWorld struct { + Instances []struct { + ID string + Instance continuityActualInstance + } +} + +func snapshotContinuityActual(runner *continuityOracleRunner) (continuityActualWorld, error) { + world := continuityActualWorld{} + for _, instanceID := range []string{continuityInstanceA, continuityInstanceB} { + record, err := runner.fixture.store.Load(context.Background(), instanceID) + if kernel.IsInstanceNotFound(err) { + world.Instances = append(world.Instances, struct { + ID string + Instance continuityActualInstance + }{ID: instanceID}) + continue + } + if err != nil { + return continuityActualWorld{}, err + } + if err := record.Validate(instanceID); err != nil { + return continuityActualWorld{}, err + } + actual := continuityActualInstance{Exists: true, StateRevision: record.State.Revision, Mode: record.State.Mode} + if record.State.Recovery != nil { + actual.RecoveryPrescriptionID = record.State.Recovery.PrescriptionID + actual.RecoveryTransitionID = record.State.Recovery.TransitionID + } + if record.State.ObjectiveBinding != nil { + accepted, err := runner.fixture.accepted(instanceID) + if err != nil { + return continuityActualWorld{}, err + } + actual.ObjectiveID = accepted.Binding.ObjectiveID + actual.ObjectiveRevision = accepted.Binding.ObjectiveRevision + actual.ObjectiveFingerprint = accepted.Binding.ObjectiveFingerprint + actual.Payload = accepted.Payload + actual.CandidateFingerprint = accepted.CandidateFingerprint + } else { + final, err := runner.fixture.validateLineage(instanceID, record.Receipts) + if err != nil { + return continuityActualWorld{}, err + } + if final != nil { + return continuityActualWorld{}, fmt.Errorf("unbound instance %q has a committed objective lineage", instanceID) + } + } + for _, receipt := range record.Receipts { + revision := uint64(0) + if receipt.ResultObjectiveBinding != nil { + revision = receipt.ResultObjectiveBinding.ObjectiveRevision + } + actual.Receipts = append(actual.Receipts, continuityActualReceipt{ + ID: receipt.ID, Transition: receipt.TransitionID, ObjectiveMutation: receipt.ObjectiveMutation, + ObjectiveRevision: revision, ResultRevision: receipt.ResultStateRevision, + }) + } + world.Instances = append(world.Instances, struct { + ID string + Instance continuityActualInstance + }{ID: instanceID, Instance: actual}) + } + return world, nil +} + +func compareContinuityWorld(reference continuityReferenceWorld, actual continuityActualWorld) error { + if len(actual.Instances) != 2 { + return fmt.Errorf("actual world has %d instance slots", len(actual.Instances)) + } + for _, item := range actual.Instances { + want := reference.Instances[item.ID] + got := item.Instance + if want.Exists != got.Exists { + return fmt.Errorf("instance %s existence=%v want %v", item.ID, got.Exists, want.Exists) + } + if !want.Exists { + continue + } + if got.StateRevision != want.StateRevision || got.Mode != want.Mode || (got.RecoveryPrescriptionID != "") != want.Recovery { + return fmt.Errorf("instance %s control state=(revision=%d mode=%s recovery=%v) want (%d %s %v)", item.ID, got.StateRevision, got.Mode, got.RecoveryPrescriptionID != "", want.StateRevision, want.Mode, want.Recovery) + } + if got.ObjectiveID != want.ObjectiveID || got.ObjectiveRevision != want.ObjectiveRevision || !reflect.DeepEqual(got.Payload, want.Payload) { + return fmt.Errorf("instance %s accepted objective=(%s@%d %#v) want (%s@%d %#v)", item.ID, got.ObjectiveID, got.ObjectiveRevision, got.Payload, want.ObjectiveID, want.ObjectiveRevision, want.Payload) + } + if len(got.Receipts) != len(want.Receipts) { + return fmt.Errorf("instance %s receipt count=%d want %d", item.ID, len(got.Receipts), len(want.Receipts)) + } + for index := range got.Receipts { + actualReceipt, referenceReceipt := got.Receipts[index], want.Receipts[index] + if actualReceipt.Transition != referenceReceipt.Transition || + actualReceipt.ObjectiveMutation != referenceReceipt.ObjectiveMutation || + actualReceipt.ResultRevision != referenceReceipt.ResultRevision { + return fmt.Errorf("instance %s receipt %d=%#v want %#v", item.ID, index, actualReceipt, referenceReceipt) + } + if actualReceipt.ObjectiveMutation != kernel.PreserveObjective && actualReceipt.ObjectiveRevision != referenceReceipt.ObjectiveRevision { + return fmt.Errorf("instance %s receipt %d objective revision=%d want %d", item.ID, index, actualReceipt.ObjectiveRevision, referenceReceipt.ObjectiveRevision) + } + } + } + return nil +} + +func continuityWorldHash(actual continuityActualWorld) (string, error) { + // The hash retains sorted per-instance control state, accepted payload and + // content identity, recovery identity, and every committed receipt ID. + encoded, err := json.Marshal(actual) + if err != nil { + return "", err + } + digest := sha256.Sum256(encoded) + return hex.EncodeToString(digest[:]), nil +} + +type continuityOracleReplay struct { + Actual continuityActualWorld + Rejected []bool + Coverage map[string]bool +} + +func replayContinuityTrace(t testing.TB, suite AcceptedObjectiveContinuityConformance, trace []continuityOracleOperation) (continuityOracleReplay, error) { + runner := newContinuityOracleRunner(t, suite) + replay := continuityOracleReplay{Coverage: map[string]bool{}} + for _, operation := range trace { + rejected, err := runner.step(operation) + if err != nil { + return continuityOracleReplay{}, fmt.Errorf("%s: %w", operation, err) + } + replay.Rejected = append(replay.Rejected, rejected) + if rejected { + for _, guard := range continuityOracleGuards { + if operation.Kind == guard { + replay.Coverage[guard] = true + } + } + } + } + actual, err := snapshotContinuityActual(runner) + if err != nil { + return continuityOracleReplay{}, err + } + replay.Actual = actual + return replay, nil +} + +type continuityOracleNode struct { + Trace []continuityOracleOperation + Reference continuityReferenceWorld + ExpectedRejects []bool + Steps int +} + +type continuityOracleReport struct { + Replayed int + States int + MaxNodes int + MaxDepth int + Coverage []string +} + +func continuityPayloadDepth(payload continuityPayload) int { + byID := map[string]continuityNode{} + for _, node := range payload.Nodes { + byID[node.ID] = node + } + maximum := 0 + for _, node := range payload.Nodes { + depth := 0 + for node.ParentID != "" { + depth++ + node = byID[node.ParentID] + } + if depth > maximum { + maximum = depth + } + } + return maximum +} + +func formatContinuityTrace(trace []continuityOracleOperation) string { + parts := make([]string, len(trace)) + for index, operation := range trace { + parts[index] = operation.String() + } + if len(parts) == 0 { + return "" + } + return strings.Join(parts, " -> ") +} + +func runContinuityOracle(t testing.TB, suite AcceptedObjectiveContinuityConformance) (continuityOracleReport, error) { + t.Helper() + queue := []continuityOracleNode{{Reference: newContinuityReferenceWorld()}} + visited := map[string]bool{} + coverage := map[string]bool{} + replayed := 0 + maxNodes, maxDepth := 0, 0 + for len(queue) > 0 { + node := queue[0] + queue = queue[1:] + replay, err := replayContinuityTrace(t, suite, node.Trace) + replayed++ + if err != nil { + return continuityOracleReport{}, fmt.Errorf("shortest divergence after %s: %w", formatContinuityTrace(node.Trace), err) + } + if !reflect.DeepEqual(replay.Rejected, node.ExpectedRejects) { + return continuityOracleReport{}, fmt.Errorf("shortest divergence after %s: rejected outcomes=%v want %v", formatContinuityTrace(node.Trace), replay.Rejected, node.ExpectedRejects) + } + if err := compareContinuityWorld(node.Reference, replay.Actual); err != nil { + return continuityOracleReport{}, fmt.Errorf("shortest divergence after %s: %w", formatContinuityTrace(node.Trace), err) + } + for guard := range replay.Coverage { + coverage[guard] = true + } + for _, item := range replay.Actual.Instances { + if nodes := len(item.Instance.Payload.Nodes); nodes > maxNodes { + maxNodes = nodes + } + if depth := continuityPayloadDepth(item.Instance.Payload); depth > maxDepth { + maxDepth = depth + } + } + hash, err := continuityWorldHash(replay.Actual) + if err != nil { + return continuityOracleReport{}, err + } + if visited[hash] { + continue + } + visited[hash] = true + for _, operation := range enumerateContinuityOperations(node.Reference) { + steps := node.Steps + if operation.semanticStep() { + steps++ + } + if steps > continuityOracleMaxSteps { + continue + } + next, rejected, err := referenceContinuityTransition(node.Reference, operation) + if err != nil { + return continuityOracleReport{}, fmt.Errorf("reference enumeration admitted %s after %s: %w", operation, formatContinuityTrace(node.Trace), err) + } + queue = append(queue, continuityOracleNode{ + Trace: append(append([]continuityOracleOperation(nil), node.Trace...), operation), + Reference: next, ExpectedRejects: append(append([]bool(nil), node.ExpectedRejects...), rejected), Steps: steps, + }) + } + } + missing := []string{} + for _, guard := range continuityOracleGuards { + if !coverage[guard] { + missing = append(missing, guard) + } + } + if len(missing) != 0 { + return continuityOracleReport{}, fmt.Errorf("oracle did not exercise minimal guard violations: %v", missing) + } + covered := make([]string, 0, len(coverage)) + for guard := range coverage { + covered = append(covered, guard) + } + sort.Strings(covered) + if maxNodes != 3 || maxDepth != 2 { + return continuityOracleReport{}, fmt.Errorf("oracle bounded coverage reached nodes=%d depth=%d, want nodes=3 depth=2", maxNodes, maxDepth) + } + return continuityOracleReport{Replayed: replayed, States: len(visited), MaxNodes: maxNodes, MaxDepth: maxDepth, Coverage: covered}, nil +} + +func acceptedContinuityOracleLaw(t *testing.T, suite AcceptedObjectiveContinuityConformance) { + report, err := runContinuityOracle(t, suite) + if err != nil { + t.Fatalf("control-law accepted-objective-continuity-oracle: %v", err) + } + t.Logf("control-law accepted-objective-continuity-oracle: replayed=%d canonical_states=%d max_steps=%d max_nodes=%d max_depth=%d rejected_guards=%s", report.Replayed, report.States, continuityOracleMaxSteps, report.MaxNodes, report.MaxDepth, strings.Join(report.Coverage, ",")) +} diff --git a/boatstack/kernel/conformance/accepted_objective_continuity_test.go b/boatstack/kernel/conformance/accepted_objective_continuity_test.go new file mode 100644 index 0000000..d728a97 --- /dev/null +++ b/boatstack/kernel/conformance/accepted_objective_continuity_test.go @@ -0,0 +1,373 @@ +package conformance + +import ( + "context" + "encoding/json" + "strings" + "testing" + "time" + + "github.com/operatorstack/boatstack/boatstack/kernel" +) + +func TestAcceptedObjectiveContinuityConformanceMemoryStore(t *testing.T) { + AcceptedObjectiveContinuityConformance{New: memoryInstanceHarness}.Run(t) +} + +func TestAcceptedObjectiveContinuityConformanceRegisterStore(t *testing.T) { + AcceptedObjectiveContinuityConformance{New: func(testing.TB) InstanceStoreHarness { + store := ®isterStore{instances: map[string]*registerInstance{}} + return InstanceStoreHarness{ + Store: store, Locker: func() kernel.Locker { return NewKeyedMemoryLocker() }, + Reopen: func(testing.TB) kernel.Store { return store }, + } + }}.Run(t) +} + +func TestAcceptedObjectiveContinuityPayloadMutationGuards(t *testing.T) { + cyclic := continuityPayload{ + InstanceID: continuityInstanceA, + Nodes: []continuityNode{ + {ID: "leaf", ParentID: "root", Status: continuityOpen}, + {ID: "root", ParentID: "leaf", Status: continuityOpen}, + }, + ActiveNodeID: "leaf", ProducerID: continuityProducer, VerifierID: continuityVerifier, + Policy: continuityPolicy{Identity: "policy-alpha", Version: "v1"}, + } + missingActive := baseContinuityPayload(continuityInstanceA) + missingActive.ActiveNodeID = "" + closedActive := baseContinuityPayload(continuityInstanceA) + closedActive.Nodes[0].Status = continuityComplete + closedActive.Nodes[0].EvidenceFingerprint = continuityFingerprint("closed-active") + premature := childContinuityPayload(continuityInstanceA) + premature.Nodes[1].Status = continuityComplete + premature.Nodes[1].EvidenceFingerprint = continuityFingerprint("premature-parent") + selfVerified := baseContinuityPayload(continuityInstanceA) + selfVerified.VerifierID = selfVerified.ProducerID + wrongActive := childContinuityPayload(continuityInstanceA) + wrongActive.ActiveNodeID = "root" + + for _, test := range []struct { + name string + payload continuityPayload + want string + }{ + {"cyclic parents", cyclic, "cyclic"}, + {"missing active node", missingActive, "lacks one active"}, + {"closed active node", closedActive, "not open"}, + {"premature completion", premature, "premature completion"}, + {"self verification", selfVerified, "distinct"}, + {"open node outside active ancestry", wrongActive, "outside the active ancestry"}, + } { + t.Run(test.name, func(t *testing.T) { + err := validateContinuityPayload(test.payload, false) + if err == nil || !strings.Contains(err.Error(), test.want) { + t.Fatalf("control-law accepted-objective-continuity-mutation-%s: error=%v want substring %q", test.name, err, test.want) + } + }) + } + + prior := childContinuityPayload(continuityInstanceA) + next := cloneContinuityPayload(prior) + next.Nodes[0].Status = continuityComplete + next.Nodes[0].EvidenceFingerprint = continuityFingerprint("leaf-complete-without-restoration") + next.ActiveNodeID = "" + if err := validateContinuityDelta(prior, next, continuityAdvanceTransition); err == nil || !strings.Contains(err.Error(), "restore") { + t.Fatalf("control-law accepted-objective-continuity-mutation-parent-restoration: %v", err) + } +} + +func newContinuityMemoryTestFixture(t testing.TB) (*continuityFixture, kernel.Runtime) { + t.Helper() + store := NewMemoryStateStore() + fixture := newContinuityFixture(t, InstanceStoreHarness{ + Store: store, Locker: func() kernel.Locker { return NewKeyedMemoryLocker() }, + Reopen: func(testing.TB) kernel.Store { return store }, + }) + return fixture, fixture.currentRuntime(t) +} + +func TestAcceptedObjectiveContinuityRejectsOpenNodeOutsideActiveAncestry(t *testing.T) { + fixture, runtime := newContinuityMemoryTestFixture(t) + if _, err := runtime.Provision(context.Background(), continuityInstanceA); err != nil { + t.Fatal(err) + } + payload := childContinuityPayload(continuityInstanceA) + payload.ActiveNodeID = "root" + objective := fixture.stage(t, continuityObjectiveA, 1, payload) + if _, _, err := attemptContinuityRecovery(runtime, fixture, continuityInstanceA, continuityBindTransition, &objective); err != nil { + t.Fatalf("control-law accepted-objective-continuity-active-ancestry: %v", err) + } + record := fixture.record(t, continuityInstanceA) + if record.State.Recovery == nil || record.State.ObjectiveBinding != nil || len(record.Receipts) != 0 { + t.Fatalf("control-law accepted-objective-continuity-active-ancestry: invalid bind became accepted: %#v", record) + } + if accepted, err := fixture.accepted(continuityInstanceA); err == nil { + t.Fatalf("control-law accepted-objective-continuity-active-ancestry: accepted=%#v", accepted) + } +} + +func actualContinuityWorldFromReference(reference continuityReferenceWorld) continuityActualWorld { + actual := continuityActualWorld{} + for _, instanceID := range []string{continuityInstanceA, continuityInstanceB} { + want := reference.Instances[instanceID] + instance := continuityActualInstance{ + Exists: want.Exists, ObjectiveID: want.ObjectiveID, ObjectiveRevision: want.ObjectiveRevision, + Payload: cloneContinuityPayload(want.Payload), StateRevision: want.StateRevision, Mode: want.Mode, + } + if want.Recovery { + instance.RecoveryPrescriptionID = "prx-reference-recovery" + } + for index, receipt := range want.Receipts { + instance.Receipts = append(instance.Receipts, continuityActualReceipt{ + ID: "receipt-reference-" + string(rune('a'+index)), Transition: receipt.Transition, + ObjectiveMutation: receipt.ObjectiveMutation, ObjectiveRevision: receipt.ObjectiveRevision, + ResultRevision: receipt.ResultRevision, + }) + } + actual.Instances = append(actual.Instances, struct { + ID string + Instance continuityActualInstance + }{ID: instanceID, Instance: instance}) + } + return actual +} + +func cloneContinuityActualWorld(t testing.TB, world continuityActualWorld) continuityActualWorld { + t.Helper() + encoded, err := json.Marshal(world) + if err != nil { + t.Fatal(err) + } + var copy continuityActualWorld + if err := json.Unmarshal(encoded, ©); err != nil { + t.Fatal(err) + } + return copy +} + +func TestAcceptedObjectiveContinuityOracleRejectsDishonestAcceptedState(t *testing.T) { + reference := newContinuityReferenceWorld() + apply := func(operation continuityOracleOperation) { + var err error + reference, _, err = referenceContinuityTransition(reference, operation) + if err != nil { + t.Fatal(err) + } + } + for _, operation := range []continuityOracleOperation{ + {oracleProvision, continuityInstanceA}, {oracleBind, continuityInstanceA}, + {oracleOpenChild, continuityInstanceA}, {oracleCompleteChild, continuityInstanceA}, + {oracleChangePolicy, continuityInstanceA}, + {oracleProvision, continuityInstanceB}, {oracleBind, continuityInstanceB}, + } { + apply(operation) + } + trusted := actualContinuityWorldFromReference(reference) + if err := compareContinuityWorld(reference, trusted); err != nil { + t.Fatalf("control-law accepted-objective-continuity-oracle-baseline: %v", err) + } + + for _, mutation := range []struct { + name string + mutate func(*continuityActualWorld) + }{ + {"latest staged read", func(world *continuityActualWorld) { + world.Instances[0].Instance.ObjectiveRevision++ + world.Instances[0].Instance.Payload = markedContinuityPayload(continuityInstanceA, "policy-oracle", "v2") + }}, + {"aliased instances", func(world *continuityActualWorld) { + world.Instances[1].Instance = world.Instances[0].Instance + }}, + {"stale overwrite", func(world *continuityActualWorld) { + world.Instances[0].Instance.ObjectiveRevision = 2 + world.Instances[0].Instance.Payload = childContinuityPayload(continuityInstanceA) + }}, + {"same revision policy substitution", func(world *continuityActualWorld) { + world.Instances[0].Instance.Payload.Policy.Identity = "policy-substituted" + }}, + } { + t.Run(mutation.name, func(t *testing.T) { + dishonest := cloneContinuityActualWorld(t, trusted) + mutation.mutate(&dishonest) + if err := compareContinuityWorld(reference, dishonest); err == nil { + t.Fatalf("control-law accepted-objective-continuity-oracle-mutation-%s: dishonest state was accepted", mutation.name) + } + }) + } +} + +type continuityLoadMutationStore struct { + kernel.Store + mutate func(string, kernel.InstanceRecord) kernel.InstanceRecord +} + +func (s continuityLoadMutationStore) Load(ctx context.Context, instanceID string) (kernel.InstanceRecord, error) { + record, err := s.Store.Load(ctx, instanceID) + if err != nil { + return kernel.InstanceRecord{}, err + } + return s.mutate(instanceID, record), nil +} + +func newRecoveredContinuityTestFixture(t testing.TB) *continuityFixture { + t.Helper() + fixture, runtime := newContinuityMemoryTestFixture(t) + if _, err := runtime.Provision(context.Background(), continuityInstanceA); err != nil { + t.Fatal(err) + } + a1Payload := baseContinuityPayload(continuityInstanceA) + a1 := fixture.stage(t, continuityObjectiveA, 1, a1Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityBindTransition, &a1) + invalid := childContinuityPayload(continuityInstanceA) + invalid.ActiveNodeID = "root" + invalidObjective := fixture.stage(t, continuityObjectiveA, 2, invalid) + if _, _, err := attemptContinuityRecovery(runtime, fixture, continuityInstanceA, continuityAdvanceTransition, &invalidObjective); err != nil { + t.Fatal(err) + } + fixture.apply(t, runtime, continuityInstanceA, continuityRecoverTransition, nil) + return fixture +} + +func TestAcceptedObjectiveContinuityReaderRejectsSubstitutedReceiptFacts(t *testing.T) { + t.Run("recovery effect", func(t *testing.T) { + fixture := newRecoveredContinuityTestFixture(t) + fixture.store = continuityLoadMutationStore{Store: fixture.store, mutate: func(_ string, record kernel.InstanceRecord) kernel.InstanceRecord { + for index := range record.Receipts { + if record.Receipts[index].TransitionID == continuityRecoverTransition { + record.Receipts[index].Effects = append([]kernel.EffectFact(nil), record.Receipts[index].Effects...) + record.Receipts[index].Effects[0].Fingerprint = continuityFingerprint("substituted-recovery-effect") + record.Receipts[index] = remintReceipt(record.Receipts[index]) + } + } + return record + }} + if accepted, err := fixture.accepted(continuityInstanceA); err == nil { + t.Fatalf("control-law accepted-objective-continuity-receipt-recovery: accepted=%#v", accepted) + } + }) + + t.Run("recovery prescription", func(t *testing.T) { + fixture := newRecoveredContinuityTestFixture(t) + fixture.store = continuityLoadMutationStore{Store: fixture.store, mutate: func(_ string, record kernel.InstanceRecord) kernel.InstanceRecord { + for index := range record.Receipts { + if record.Receipts[index].TransitionID == continuityRecoverTransition { + record.Receipts[index].PrescriptionID = "prx-substituted-recovery" + record.Receipts[index] = remintReceipt(record.Receipts[index]) + } + } + return record + }} + if accepted, err := fixture.accepted(continuityInstanceA); err == nil { + t.Fatalf("control-law accepted-objective-continuity-receipt-recovery-prescription: accepted=%#v", accepted) + } + }) + + mutations := []struct { + name string + mutate func(*kernel.Receipt) + }{ + {"program", func(receipt *kernel.Receipt) { + receipt.Program = kernel.ProgramIdentity{ID: "substituted-program", Version: "9.0.0", Fingerprint: continuityFingerprint("substituted-program")} + }}, + {"authority", func(receipt *kernel.Receipt) { receipt.AuthorityFingerprint = "substituted-authority" }}, + {"capabilities", func(receipt *kernel.Receipt) { receipt.Capabilities = []kernel.Capability{"objective.advance"} }}, + {"observation", func(receipt *kernel.Receipt) { + receipt.ResultObservation = continuityFingerprint("substituted-observation") + }}, + {"commit time", func(receipt *kernel.Receipt) { receipt.CommittedAt = receipt.CommittedAt.Add(time.Second) }}, + {"effect", func(receipt *kernel.Receipt) { + receipt.Effects = append([]kernel.EffectFact(nil), receipt.Effects...) + receipt.Effects[0].Fingerprint = continuityFingerprint("substituted-advance-effect") + }}, + } + for _, mutation := range mutations { + mutation := mutation + t.Run("advance "+mutation.name, func(t *testing.T) { + fixture, runtime := newContinuityMemoryTestFixture(t) + if _, err := runtime.Provision(context.Background(), continuityInstanceA); err != nil { + t.Fatal(err) + } + a1Payload := baseContinuityPayload(continuityInstanceA) + a1 := fixture.stage(t, continuityObjectiveA, 1, a1Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityBindTransition, &a1) + a2Payload := childContinuityPayload(continuityInstanceA) + a2 := fixture.stage(t, continuityObjectiveA, 2, a2Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a2) + fixture.store = continuityLoadMutationStore{Store: fixture.store, mutate: func(_ string, record kernel.InstanceRecord) kernel.InstanceRecord { + last := len(record.Receipts) - 1 + mutation.mutate(&record.Receipts[last]) + prescriptionID, err := continuityPrescriptionID(record.Receipts[last]) + if err != nil { + t.Fatal(err) + } + record.Receipts[last].PrescriptionID = prescriptionID + record.Receipts[last] = remintReceipt(record.Receipts[last]) + return record + }} + if accepted, err := fixture.accepted(continuityInstanceA); err == nil { + t.Fatalf("control-law accepted-objective-continuity-receipt-%s: accepted=%#v", mutation.name, accepted) + } + }) + } +} + +func TestAcceptedObjectiveContinuityReaderRejectsDishonestPayloadLineage(t *testing.T) { + fixture, runtime := newContinuityMemoryTestFixture(t) + if _, err := runtime.Provision(context.Background(), continuityInstanceA); err != nil { + t.Fatal(err) + } + a1Payload := baseContinuityPayload(continuityInstanceA) + a1 := fixture.stage(t, continuityObjectiveA, 1, a1Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityBindTransition, &a1) + a2Payload := childContinuityPayload(continuityInstanceA) + a2 := fixture.stage(t, continuityObjectiveA, 2, a2Payload) + fixture.apply(t, runtime, continuityInstanceA, continuityAdvanceTransition, &a2) + + // The forged A@3 payload is statically valid, and its receipt has a valid + // content identity and objective lifecycle. It is still dishonest because + // it removes the durable A@2 child instead of extending the payload lineage. + forgedPayload := baseContinuityPayload(continuityInstanceA) + forgedPayload.Policy = continuityPolicy{Identity: "policy-forged", Version: "v2"} + forgedObjective := fixture.stage(t, continuityObjectiveA, 3, forgedPayload) + forgedCandidate, err := fixture.candidates.resolve(forgedObjective) + if err != nil { + t.Fatal(err) + } + forgedBinding, err := kernel.BindObjective(forgedObjective) + if err != nil { + t.Fatal(err) + } + trusted := fixture.record(t, continuityInstanceA) + priorBinding := *trusted.State.ObjectiveBinding + forgedReceipt := trusted.Receipts[len(trusted.Receipts)-1] + forgedReceipt.ID = "" + forgedReceipt.PriorStateRevision = trusted.State.Revision + forgedReceipt.AttemptStateRevision = trusted.State.Revision + 1 + forgedReceipt.ResultStateRevision = trusted.State.Revision + 2 + forgedReceipt.PriorObjectiveBinding = &priorBinding + forgedReceipt.RequestedObjectiveBinding = &forgedBinding + forgedReceipt.ResultObjectiveBinding = &forgedBinding + forgedReceipt.Effects = []kernel.EffectFact{{ + Facet: "register.payload", Operation: continuityAdvanceTransition, Fingerprint: forgedCandidate.ContentFingerprint, + }} + forgedReceipt.PrescriptionID, err = continuityPrescriptionID(forgedReceipt) + if err != nil { + t.Fatal(err) + } + forgedReceipt = remintReceipt(forgedReceipt) + forgedState := trusted.State + forgedState.Revision = forgedReceipt.ResultStateRevision + forgedState.ObjectiveBinding = &forgedBinding + + fixture.store = continuityLoadMutationStore{Store: fixture.store, mutate: func(instanceID string, record kernel.InstanceRecord) kernel.InstanceRecord { + if instanceID == continuityInstanceA { + record.State = forgedState + record.Receipts = append(record.Receipts, forgedReceipt) + } + return record + }} + if accepted, err := fixture.accepted(continuityInstanceA); err == nil || !strings.Contains(err.Error(), "payload delta") { + t.Fatalf("control-law accepted-objective-continuity-reader-lineage: accepted=%#v error=%v", accepted, err) + } +} diff --git a/release-notes/2026-08-23-accepted-objective-continuity-conformance.md b/release-notes/2026-08-23-accepted-objective-continuity-conformance.md new file mode 100644 index 0000000..dbc82ce --- /dev/null +++ b/release-notes/2026-08-23-accepted-objective-continuity-conformance.md @@ -0,0 +1,18 @@ +### Accepted objectives now have a reusable continuity conformance suite + +Store adapters can run one domain-neutral suite that follows an immutable, +bounded hierarchical register from its first objective binding through +advancement, rejection, recovery, restart, a concurrent successor race, +policy revision, and marked completion. The suite reconstructs acceptance +only from the durable binding and complete committed receipt lineage. Staged +candidates and read-only resolution cannot become accepted state. + +The laws also prove per-instance isolation, side-effect-free committed retry, +cross-instance refusal, and restarted-versus-uninterrupted equivalence. A +deterministic six-step breadth-first oracle replays every explored trace from +a fresh harness and reports the shortest divergence. Its guard mutations +cover cyclic parents, invalid active nodes, missing parent restoration, +premature completion, self-verification, staged-value substitution, instance +aliasing, stale overwrites, and same-revision policy substitution. The memory, +revisioned-register, and reviewer file stores all run the same suite without a +new kernel API.