Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 13 additions & 1 deletion boatstack/internal/effects/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -208,8 +208,11 @@ func (d Driver) Prepare(ctx context.Context, admission protocol.Admission, trans
}
mutations = append(mutations, bindingMutation)
}
var facetGroups [][2]durable.State
facetGroups = append(facetGroups, [2]durable.State{state, next})
if transition.ID == "workspace.cut" {
parked := parkedSourceState(state, next.Revision, transition.ID, d.clock.Now())
facetGroups = append(facetGroups, [2]durable.State{state, parked})
parkedRaw, encodeErr := durable.EncodeState(parked)
if encodeErr != nil {
return nil, encodeErr
Expand Down Expand Up @@ -271,7 +274,16 @@ func (d Driver) Prepare(ctx context.Context, admission protocol.Admission, trans
mutations[index].Resource = transition.OwnedResources[0]
mutations[index].Owner = transition.Owner
}
prepared := &preparedEffect{mutations: mutations, verifyInvocation: verificationInvocation}
changedFacets, facetErr := changedStateFacets(facetGroups...)
if facetErr != nil {
return nil, facetErr
}
changedFacets, facetErr = validateTransitionStateFacets(transition, changedFacets)
if facetErr != nil {
return nil, facetErr
}
mutations = annotateStateFacetMutations(mutations, changedFacets)
prepared := &preparedEffect{mutations: mutations, verifyInvocation: verificationInvocation, changedStateFacets: changedFacets}
if err := bindPreparedCapabilities(prepared, admission, transition); err != nil {
return nil, err
}
Expand Down
18 changes: 16 additions & 2 deletions boatstack/internal/effects/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"time"

"github.com/operatorstack/boatstack/boatstack/internal/kernel/catalog"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/model"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/ports"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/protocol"
)
Expand Down Expand Up @@ -113,6 +114,12 @@ func readJournal(path string) (journalRecord, error) {
if err := record.Admission.ValidateIdentity(); err != nil || record.Admission.TransitionID != record.TransitionID {
return journalRecord{}, fmt.Errorf("invalid transaction admission in %s: %v", path, err)
}
for _, mutation := range record.Mutations {
facets, err := model.NormalizeStateFacets("journal mutation state facets", mutation.StateFacets)
if err != nil || !slices.Equal(facets, mutation.StateFacets) {
return journalRecord{}, fmt.Errorf("invalid transaction state facets in %s: %v", path, err)
}
}
if record.Receipt != nil {
if err := record.Receipt.Validate(); err != nil || record.Receipt.ID != record.ReceiptID || record.Receipt.AdmissionID != record.Admission.ID || record.Receipt.TransitionID != record.TransitionID {
return journalRecord{}, fmt.Errorf("invalid committed transition fact in %s: %v", path, err)
Expand All @@ -126,7 +133,7 @@ func readJournal(path string) (journalRecord, error) {
receipt.DeliveryID != admission.Goal.DeliveryID || receipt.GoalScope != admission.GoalScope || receipt.GoalStatus != admission.GoalStatus {
return journalRecord{}, fmt.Errorf("committed transition fact in %s does not match its exact admission", path)
}
if err := validateCommittedMutationFacts(record.TransitionClass, record.Mutations, receipt.CommittedEffects); err != nil {
if err := validateCommittedMutationFacts(record.TransitionClass, record.Mutations, receipt.ChangedStateFacets, receipt.CommittedEffects); err != nil {
return journalRecord{}, fmt.Errorf("committed transition fact in %s: %w", path, err)
}
}
Expand All @@ -136,7 +143,14 @@ func readJournal(path string) (journalRecord, error) {
return record, nil
}

func validateCommittedMutationFacts(class catalog.EventClass, mutations []ports.ResourceMutation, facts []protocol.EffectFact) error {
func validateCommittedMutationFacts(class catalog.EventClass, mutations []ports.ResourceMutation, receiptFacets []model.StateFacet, facts []protocol.EffectFact) error {
var mutationFacets []model.StateFacet
for _, mutation := range mutations {
mutationFacets = model.UnionStateFacets(mutationFacets, mutation.StateFacets)
}
if !slices.Equal(mutationFacets, receiptFacets) {
return fmt.Errorf("receipt changed state facets %v do not match staged mutation facets %v", receiptFacets, mutationFacets)
}
resourceFacts := make([]protocol.EffectFact, 0, len(facts))
boundarySettled := false
for _, fact := range facts {
Expand Down
5 changes: 5 additions & 0 deletions boatstack/internal/effects/prepared.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ type preparedEffect struct {
transition catalog.Transition
requiredCapabilities []catalog.Capability
effectiveCapabilities []catalog.Capability
changedStateFacets []model.StateFacet
}

func (p *preparedEffect) ChangedStateFacets() []model.StateFacet {
return append([]model.StateFacet(nil), p.changedStateFacets...)
}

func (p *preparedEffect) Manifest() []ports.ResourceMutation {
Expand Down
4 changes: 3 additions & 1 deletion boatstack/internal/effects/receipts.go
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,7 @@ type processEvent struct {
RequiredCapabilities []catalog.Capability `json:"required_capabilities"`
GrantedCapabilities []catalog.Capability `json:"granted_capabilities"`
CommittedEffects []protocol.EffectFact `json:"committed_effects"`
ChangedStateFacets []model.StateFacet `json:"changed_state_facets"`
Verification protocol.VerificationFact `json:"verification"`
Recovery string `json:"recovery,omitempty"`
Terminal string `json:"terminal"`
Expand Down Expand Up @@ -199,7 +200,7 @@ func (s *ReceiptStore) Project(ctx context.Context, receipt protocol.TransitionR
return err
}
event := processEvent{
SchemaVersion: 4, FlowID: receipt.FlowID, Sequence: receipt.Sequence, Timestamp: s.clock.Now().UTC(), GoalID: receipt.GoalID,
SchemaVersion: 5, FlowID: receipt.FlowID, Sequence: receipt.Sequence, Timestamp: s.clock.Now().UTC(), GoalID: receipt.GoalID,
GoalScope: string(receipt.GoalScope), GoalStatus: string(receipt.GoalStatus),
TransitionID: string(receipt.TransitionID), ProgramID: receipt.Program.ID, ProgramVersion: receipt.Program.Version, ProgramFingerprint: receipt.Program.Fingerprint, PrescriptionID: receipt.PrescriptionID,
PriorStateRevision: receipt.PriorStateRevision, ResultingStateRevision: receipt.ResultingStateRevision,
Expand All @@ -210,6 +211,7 @@ func (s *ReceiptStore) Project(ctx context.Context, receipt protocol.TransitionR
RequiredCapabilities: append([]catalog.Capability(nil), receipt.RequiredCapabilities...),
GrantedCapabilities: append([]catalog.Capability(nil), receipt.GrantedCapabilities...),
CommittedEffects: append([]protocol.EffectFact(nil), receipt.CommittedEffects...),
ChangedStateFacets: append([]model.StateFacet(nil), receipt.ChangedStateFacets...),
Verification: receipt.Verification,
}
for _, source := range receipt.AuthoritySources {
Expand Down
30 changes: 28 additions & 2 deletions boatstack/internal/effects/recovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,12 @@ func (d Driver) prepareRecoveryReplay(ctx context.Context, layout ports.Controll
return nil, err
}
mutations = append(mutations, closure...)
return &preparedEffect{mutations: mutations}, nil
changed, err := recoveryStateFacets(record, transition.ID, admission.Invocation, mutations)
if err != nil {
return nil, err
}
mutations = annotateStateFacetMutations(mutations, changed)
return &preparedEffect{mutations: mutations, changedStateFacets: changed}, nil
}

func (d Driver) prepareWorkspaceCutReconciliation(ctx context.Context, layout ports.ControllerLayout, admission protocol.Admission, record journalRecord, pendingPath string) (ports.PreparedEffect, error) {
Expand Down Expand Up @@ -129,7 +134,28 @@ func (d Driver) prepareWorkspaceCutReconciliation(ctx context.Context, layout po
return nil, err
}
mutations = append(mutations, closure...)
return &preparedEffect{mutations: mutations, verifyInvocation: verificationInvocation}, nil
changed, err := recoveryStateFacets(record, "workspace.reconcile", admission.Invocation, mutations)
if err != nil {
return nil, err
}
mutations = annotateStateFacetMutations(mutations, changed)
return &preparedEffect{mutations: mutations, verifyInvocation: verificationInvocation, changedStateFacets: changed}, nil
}

func recoveryStateFacets(record journalRecord, recovery catalog.TransitionID, invocation model.InvocationContext, mutations []ports.ResourceMutation) ([]model.StateFacet, error) {
staged, err := journalStateFacets(record.Mutations)
if err != nil {
return nil, err
}
allowed := model.UnionStateFacets(catalog.DurableStateWritesForRecovery(record.TransitionID), []model.StateFacet{model.StateFacetControl})
if _, err := validateAllowedStateFacets(recovery, staged, allowed); err != nil {
return nil, err
}
changed, err := mutationStateFacets(invocation, mutations)
if err != nil {
return nil, err
}
return validateAllowedStateFacets(recovery, changed, allowed)
}

func (d Driver) advanceRecoveredState(layout ports.ControllerLayout, admission protocol.Admission, transition catalog.TransitionID, mutations []ports.ResourceMutation) ([]ports.ResourceMutation, error) {
Expand Down
12 changes: 12 additions & 0 deletions boatstack/internal/effects/revision.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

"github.com/operatorstack/boatstack/boatstack/internal/kernel/catalog"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/durable"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/model"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/ports"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/protocol"
)
Expand Down Expand Up @@ -43,6 +44,7 @@ func BindStateRevision(ctx context.Context, prepared ports.PreparedEffect, resol
if state.ProgramFingerprint != "" && state.ProgramFingerprint != admission.ExpectedProgramFingerprint {
return nil, fmt.Errorf("compiled control program changed before revision binding")
}
before := state
state.Revision, err = durable.NextRevision(state.Revision)
if err != nil {
return nil, err
Expand All @@ -62,5 +64,15 @@ func BindStateRevision(ctx context.Context, prepared ports.PreparedEffect, resol
}
mutation.Resource, mutation.Owner = kernelStateResource, kernelStateOwner
effect.mutations = append(effect.mutations, mutation)
changed, err := changedStateFacets([2]durable.State{before, state})
if err != nil {
return nil, err
}
changed, err = validateTransitionStateFacets(transition, changed)
if err != nil {
return nil, err
}
effect.mutations = annotateStateFacetMutations(effect.mutations, changed)
effect.changedStateFacets = model.UnionStateFacets(effect.changedStateFacets, changed)
return effect, nil
}
128 changes: 128 additions & 0 deletions boatstack/internal/effects/state_facet.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
package effects

import (
"fmt"
"path/filepath"

"github.com/operatorstack/boatstack/boatstack/internal/kernel/catalog"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/durable"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/model"
"github.com/operatorstack/boatstack/boatstack/internal/kernel/ports"
)

func validateTransitionStateFacets(transition catalog.Transition, changed []model.StateFacet) ([]model.StateFacet, error) {
policy, err := catalog.DurableStateFacetPolicy(transition)
if err != nil {
return nil, err
}
return validateAllowedStateFacets(transition.ID, changed, policy.Writes)
}

func validateAllowedStateFacets(transition catalog.TransitionID, changed, allowed []model.StateFacet) ([]model.StateFacet, error) {
canonical, err := model.NormalizeStateFacets("changed state facets", changed)
if err != nil {
return nil, err
}
allowedSet := map[model.StateFacet]bool{}
for _, facet := range allowed {
allowedSet[facet] = true
}
for _, facet := range canonical {
if !allowedSet[facet] {
return nil, fmt.Errorf("FACET_OWNERSHIP_VIOLATION: transition %q changed %q outside its kernel-approved durable state facets", transition, facet)
}
}
return canonical, nil
}

func changedStateFacets(groups ...[2]durable.State) ([]model.StateFacet, error) {
var result []model.StateFacet
for _, group := range groups {
changed, err := durable.ChangedFacets(group[0], group[1])
if err != nil {
return nil, err
}
result = model.UnionStateFacets(result, changed)
}
return result, nil
}

// journalStateFacets consumes the semantic state delta that was validated and
// staged with the interrupted transaction. Recovery cannot widen it.
func journalStateFacets(mutations []ports.ResourceMutation) ([]model.StateFacet, error) {
var changed []model.StateFacet
for _, mutation := range mutations {
if filepath.Base(mutation.Path) != "state.json" {
continue
}
stateMutation := false
if mutation.PriorExists {
_, err := durable.DecodeState(mutation.Prior)
stateMutation = err == nil
}
if !stateMutation && !mutation.Delete && mutation.TargetLink == "" {
_, err := durable.DecodeState(mutation.Target)
stateMutation = err == nil
}
if !stateMutation {
continue
}
facets, err := model.NormalizeStateFacets("journal state facets", mutation.StateFacets)
if err != nil || len(facets) == 0 {
return nil, fmt.Errorf("STATE_FACET_UNCLASSIFIED: interrupted durable state mutation %s has no valid staged facets: %v", mutation.Path, err)
}
changed = model.UnionStateFacets(changed, facets)
}
return changed, nil
}

// mutationStateFacets derives the semantic delta that the prepared mutation
// will actually commit. Recovery receipts must describe this delta, not the
// interrupted transition's wider staged envelope.
func mutationStateFacets(invocation model.InvocationContext, mutations []ports.ResourceMutation) ([]model.StateFacet, error) {
var changed []model.StateFacet
for _, mutation := range mutations {
if filepath.Base(mutation.Path) != "state.json" || mutation.Delete || mutation.TargetLink != "" {
continue
}
after, err := durable.DecodeState(mutation.Target)
if err != nil {
continue
}
before := durable.Default(invocation, after.UpdatedAt)
if mutation.PriorExists {
before, err = durable.DecodeState(mutation.Prior)
if err != nil {
return nil, fmt.Errorf("STATE_FACET_UNCLASSIFIED: prepared durable state mutation %s has invalid prior state: %w", mutation.Path, err)
}
}
facets, err := durable.ChangedFacets(before, after)
if err != nil {
return nil, err
}
changed = model.UnionStateFacets(changed, facets)
}
if len(changed) == 0 {
return nil, fmt.Errorf("STATE_FACET_UNCLASSIFIED: prepared effect has no durable state delta")
}
return changed, nil
}

func annotateStateFacetMutations(mutations []ports.ResourceMutation, facets []model.StateFacet) []ports.ResourceMutation {
for index := range mutations {
mutation := &mutations[index]
stateMutation := false
if mutation.PriorExists && filepath.Base(mutation.Path) == "state.json" {
_, err := durable.DecodeState(mutation.Prior)
stateMutation = err == nil
}
if !stateMutation && !mutation.Delete && mutation.TargetLink == "" && filepath.Base(mutation.Path) == "state.json" {
_, err := durable.DecodeState(mutation.Target)
stateMutation = err == nil
}
if stateMutation {
mutation.StateFacets = append([]model.StateFacet(nil), facets...)
}
}
return mutations
}
Loading