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
7 changes: 6 additions & 1 deletion .github/tests/test_detached_supervision.py
Original file line number Diff line number Diff line change
Expand Up @@ -310,7 +310,12 @@ def test_authority_free_frontier_does_not_block_authorized_plan_creation(self) -
self.assertEqual(applied_process.stderr, "")
self.assertEqual(applied["receipt"]["transition_id"], "plan.create")
self.assertEqual(applied["receipt"]["flow_id"], "flow-codex-driver-authority-triggers")
self.assertEqual(applied["receipt"]["outcome"], "succeeded")
self.assertEqual(applied["receipt"]["kind"], "transition-committed")
self.assertTrue(applied["receipt"]["program"]["id"])
self.assertTrue(applied["receipt"]["program"]["version"])
self.assertTrue(applied["receipt"]["program"]["fingerprint"])
self.assertTrue(applied["receipt"]["committed_effects"])
self.assertEqual(applied["receipt"]["verification"]["result"], "satisfied")
self.assertTrue(applied["receipt"]["target_fingerprint"])
self.assertEqual(applied["receipt"]["recovery"], "recovery.resume")
self.assertEqual(applied["snapshot"]["plan"]["value"], "draft")
Expand Down
11 changes: 9 additions & 2 deletions .github/tests/test_repository_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -573,7 +573,9 @@ def test_offline_installer_initializes_updates_and_guards_through_kernel(self) -
self.assertTrue(event["authority_fingerprint"])
self.assertTrue(event["required_capabilities"])
self.assertTrue(event["granted_capabilities"])
self.assertTrue(event["exercised_capabilities"])
self.assertNotIn("exercised_capabilities", event)
self.assertTrue(event["committed_effects"])
self.assertEqual(event["verification"]["result"], "satisfied")

def test_program_changing_update_is_explicit_atomic_and_dormant_safe(self) -> None:
# control-law: accepted-program-delta-atomically-pins-runtime-and-program
Expand Down Expand Up @@ -692,9 +694,14 @@ def test_program_changing_update_is_explicit_atomic_and_dormant_safe(self) -> No
if receipt["transition_id"] == "installation.reconcile-update"
)
self.assertTrue(update["program_change_accepted"])
self.assertEqual(update["kind"], "transition-committed")
self.assertRegex(update["prior_program_fingerprint"], r"^[0-9a-f]{64}$")
self.assertRegex(update["program_fingerprint"], r"^[0-9a-f]{64}$")
self.assertTrue(update["program"]["id"])
self.assertTrue(update["program"]["version"])
self.assertRegex(update["program"]["fingerprint"], r"^[0-9a-f]{64}$")
self.assertRegex(update["program_delta_fingerprint"], r"^[0-9a-f]{64}$")
self.assertTrue(update["committed_effects"])
self.assertEqual(update["verification"]["result"], "satisfied")
self.assertEqual(
update["runtime_fingerprint"], hashlib.sha256(self.helper.read_bytes()).hexdigest()
)
Expand Down
15 changes: 14 additions & 1 deletion boatstack/internal/effects/cas_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ func TestConcurrentApplyConsumesOneRevisionExactlyOnce(t *testing.T) {
t.Fatalf("concurrent results: success=%d stale=%d one=%v two=%v three=%v", successes, stale, one.err, two.err, three.err)
}
if committed.Receipt == nil || committed.Receipt.PriorStateRevision != 1 || committed.Receipt.ResultingStateRevision != 2 ||
committed.Receipt.ProgramFingerprint != program.Fingerprint() || committed.Receipt.PrescriptionID != request.Prescription.ID {
committed.Receipt.Program.Fingerprint != program.Fingerprint() || committed.Receipt.PrescriptionID != request.Prescription.ID {
t.Fatalf("commit receipt does not prove the consumed revision/program pair: %#v", committed.Receipt)
}

Expand All @@ -136,6 +136,19 @@ func TestConcurrentApplyConsumesOneRevisionExactlyOnce(t *testing.T) {
if err != nil || bytes.Count(receiptRaw, []byte("\n")) != 1 {
t.Fatalf("receipt stream contains more than one commit: %v %q", err, receiptRaw)
}
committedJournals, err := filepath.Glob(filepath.Join(layout.JournalRoot, "*.committed"))
if err != nil || len(committedJournals) != 1 {
t.Fatalf("canonical committed journal count=%d err=%v", len(committedJournals), err)
}
committedRaw, err := os.ReadFile(committedJournals[0])
if err != nil || !bytes.Contains(committedRaw, []byte(committed.Receipt.ID)) || !bytes.Contains(committedRaw, []byte("committed_effects")) {
t.Fatalf("committed journal lacks its complete transition fact: %v %q", err, committedRaw)
}
// Simulate a crash after canonical commit but before the passive receipt
// projection reaches its consumer. Replay must recover from the journal fact.
if err := os.Remove(layout.ReceiptPath); err != nil {
t.Fatal(err)
}

replayRequest := request
replayRequest.IdempotencyKey = committed.Receipt.IdempotencyKey
Expand Down
16 changes: 9 additions & 7 deletions boatstack/internal/effects/integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@ type fixedClock struct{ value time.Time }

const testProgramFingerprint = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"

var testProgramIdentity = protocol.ProgramIdentity{ID: "standard", Version: "test", Fingerprint: testProgramFingerprint}

func testGoalContracts() catalog.GoalContracts {
manifest, err := standard.Definition().RuntimeManifest(context.Background())
if err != nil {
Expand Down Expand Up @@ -158,7 +160,7 @@ func TestConcreteBoundaryAppliesAndReceiptsOneTransition(t *testing.T) {
if err != nil {
t.Fatal(err)
}
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramFingerprint, observer, clock, locker, journal, driver, receipts)
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramIdentity, observer, clock, locker, journal, driver, receipts)
if err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -320,7 +322,7 @@ func TestProgramDriftRequiresAtomicInstallationReconciliation(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if initialized.Receipt == nil || initialized.Receipt.ProgramFingerprint != oldProgram.Fingerprint() {
if initialized.Receipt == nil || initialized.Receipt.Program.Fingerprint != oldProgram.Fingerprint() {
t.Fatalf("initial receipt did not freeze old program: %#v", initialized.Receipt)
}
if initialized.Snapshot == nil || initialized.Snapshot.Goal.Status != model.FactAbsent || initialized.Receipt.GoalStatus != model.FactAbsent || initialized.Receipt.GoalID != "" {
Expand Down Expand Up @@ -390,7 +392,7 @@ func TestProgramDriftRequiresAtomicInstallationReconciliation(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if reconciled.Receipt == nil || reconciled.Receipt.ProgramFingerprint != newProgram.Fingerprint() ||
if reconciled.Receipt == nil || reconciled.Receipt.Program.Fingerprint != newProgram.Fingerprint() ||
reconciled.Receipt.PriorProgramFingerprint != oldProgram.Fingerprint() || reconciled.Receipt.ProgramDeltaFingerprint == "" ||
!reconciled.Receipt.ProgramChangeAccepted || reconciled.Receipt.RuntimeFingerprint != digestBytes(runtimeRaw) ||
reconciled.Receipt.RuntimeSourceRevision != "program-new" || reconciled.Snapshot == nil ||
Expand Down Expand Up @@ -506,7 +508,7 @@ func TestReferenceExtensionUsesKernelAdmissionVerificationAndReceiptPath(t *test
{Name: "source_revision", Value: "extension-fixture"}, {Name: "runtime_version", Value: runtimeVersion}, {Name: "runtime_sha256", Value: digestBytes(runtimeRaw)},
{Name: "config_path", Value: configPath}, {Name: "config_sha256", Value: configFingerprint(t, configRaw)},
})
if initialized.Receipt == nil || initialized.Receipt.AuthorityFingerprint == "" || len(initialized.Receipt.AuthoritySources) != 1 || len(initialized.Receipt.RequiredCapabilities) == 0 || len(initialized.Receipt.GrantedCapabilities) == 0 || len(initialized.Receipt.ExercisedCapabilities) == 0 {
if initialized.Receipt == nil || initialized.Receipt.AuthorityFingerprint == "" || len(initialized.Receipt.AuthoritySources) != 1 || len(initialized.Receipt.RequiredCapabilities) == 0 || len(initialized.Receipt.GrantedCapabilities) == 0 || len(initialized.Receipt.ExercisedCapabilities) != 0 || len(initialized.Receipt.CommittedEffects) == 0 || initialized.Receipt.Verification.Result != protocol.VerificationSatisfied {
t.Fatalf("receipt lost capability or authority provenance: %#v", initialized.Receipt)
}
apply("goal.configure", authority(catalog.AuthorityHuman), protocol.Parameters{{Name: "goal_kind", Value: string(goal.Kind)}, {Name: "delivery_id", Value: goal.DeliveryID}})
Expand Down Expand Up @@ -549,7 +551,7 @@ func TestReferenceExtensionUsesKernelAdmissionVerificationAndReceiptPath(t *test
t.Fatalf("unmet extension obligation decision = %#v", next.Decision)
}
completed := apply(releasenote.Transition, authority(catalog.AuthorityRepository), nil)
if completed.Receipt == nil || completed.Receipt.TransitionID != releasenote.Transition || completed.Receipt.ProgramFingerprint != program.Fingerprint() ||
if completed.Receipt == nil || completed.Receipt.TransitionID != releasenote.Transition || completed.Receipt.Program.Fingerprint != program.Fingerprint() ||
completed.Snapshot == nil || completed.Snapshot.ExtensionFacts[releasenote.FactID].Value != "verified" {
t.Fatalf("extension did not traverse verified receipt path: %#v", completed)
}
Expand Down Expand Up @@ -588,7 +590,7 @@ func TestConcreteWorkflowPreservesConfigurationProofAndGoalTerminals(t *testing.
journal, _ := effects.NewJournal(resolver, clock)
receipts, _ := effects.NewReceiptStore(resolver, clock)
driver, _ := effects.NewDriver(resolver, clock, effects.NewNativeBoundary())
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramFingerprint, observer, clock, locker, journal, driver, receipts)
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramIdentity, observer, clock, locker, journal, driver, receipts)
if err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -734,7 +736,7 @@ func TestWorkspaceCutTransfersAuthorityToExactDestinationWorktree(t *testing.T)
journal, _ := effects.NewJournal(resolver, clock)
receipts, _ := effects.NewReceiptStore(resolver, clock)
driver, _ := effects.NewDriver(resolver, clock, effects.NewNativeBoundary())
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramFingerprint, observer, clock, locker, journal, driver, receipts)
kernel, err := engine.New(testprogram.StandardRegistry(), testGoalContracts(), testProgramIdentity, observer, clock, locker, journal, driver, receipts)
if err != nil {
t.Fatal(err)
}
Expand Down
142 changes: 130 additions & 12 deletions boatstack/internal/effects/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,12 @@ package effects
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"slices"
"strings"
"time"

Expand All @@ -30,17 +32,18 @@ func NewJournal(resolver ports.InvocationResolver, clock ports.Clock) (*Journal,
}

type journalRecord struct {
SchemaVersion int `json:"schema_version"`
Admission protocol.Admission `json:"admission"`
TransitionID catalog.TransitionID `json:"transition_id"`
TransitionClass catalog.EventClass `json:"transition_class"`
ReconcilesProgram bool `json:"reconciles_program,omitempty"`
Status string `json:"status"`
Mutations []ports.ResourceMutation `json:"mutations,omitempty"`
Reason string `json:"reason,omitempty"`
ReceiptID string `json:"receipt_id,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
SchemaVersion int `json:"schema_version"`
Admission protocol.Admission `json:"admission"`
TransitionID catalog.TransitionID `json:"transition_id"`
TransitionClass catalog.EventClass `json:"transition_class"`
ReconcilesProgram bool `json:"reconciles_program,omitempty"`
Status string `json:"status"`
Mutations []ports.ResourceMutation `json:"mutations,omitempty"`
Reason string `json:"reason,omitempty"`
ReceiptID string `json:"receipt_id,omitempty"`
Receipt *protocol.TransitionReceipt `json:"receipt,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}

func journalName(id, suffix string) (string, error) {
Expand Down Expand Up @@ -110,9 +113,72 @@ 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)
}
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)
}
receipt := record.Receipt
admission := record.Admission
if receipt.PrescriptionID != admission.PrescriptionID || receipt.TransitionVersion != admission.TransitionVersion || receipt.Program.Fingerprint != admission.ExpectedProgramFingerprint ||
receipt.PriorStateRevision != admission.ExpectedStateRevision || receipt.SourceFingerprint != admission.ExpectedSnapshotFingerprint ||
receipt.AuthorityFingerprint != admission.AuthorityFingerprint || !slices.Equal(receipt.RequiredCapabilities, admission.RequiredCapabilities) ||
!slices.Equal(receipt.GrantedCapabilities, admission.GrantedCapabilities) || receipt.GoalID != admission.Goal.ID || receipt.GoalKind != admission.Goal.Kind ||
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 {
return journalRecord{}, fmt.Errorf("committed transition fact in %s: %w", path, err)
}
}
if strings.HasSuffix(path, ".committed") && (record.Status != "committed" || record.Receipt == nil) {
return journalRecord{}, fmt.Errorf("committed transaction journal %s lacks its canonical transition fact", path)
}
return record, nil
}

func validateCommittedMutationFacts(class catalog.EventClass, mutations []ports.ResourceMutation, facts []protocol.EffectFact) error {
resourceFacts := make([]protocol.EffectFact, 0, len(facts))
boundarySettled := false
for _, fact := range facts {
if fact.Kind == protocol.EffectResourceMutation {
resourceFacts = append(resourceFacts, fact)
} else if fact.Kind == protocol.EffectBoundarySettled {
boundarySettled = true
}
}
if class == catalog.EventOwnedExternal && !boundarySettled {
return fmt.Errorf("owned external transaction lacks a settled boundary fact")
}
if len(resourceFacts) != len(mutations) {
return fmt.Errorf("resource fact count %d does not match staged mutation count %d", len(resourceFacts), len(mutations))
}
matched := make([]bool, len(resourceFacts))
for _, mutation := range mutations {
operation := "update"
switch {
case mutation.Delete:
operation = "delete"
case mutation.TargetLink != "":
operation = "symlink"
case !mutation.PriorExists:
operation = "create"
}
prior := mutationStateFingerprint(mutation.PriorExists, mutation.Prior, mutation.PriorLink, mutation.Mode)
result := mutationStateFingerprint(!mutation.Delete, mutation.Target, mutation.TargetLink, mutation.Mode)
found := false
for index, fact := range resourceFacts {
if !matched[index] && fact.Target == mutation.Path && fact.Operation == operation && fact.PriorFingerprint == prior && fact.ResultingFingerprint == result {
matched[index], found = true, true
break
}
}
if !found {
return fmt.Errorf("staged mutation %s has no exact committed effect fact", mutation.Path)
}
}
return nil
}

func (j *Journal) update(ctx context.Context, admissionID string, update func(*journalRecord)) error {
name, err := journalName(admissionID, ".pending")
if err != nil {
Expand Down Expand Up @@ -192,6 +258,10 @@ func (j *Journal) finalize(ctx context.Context, admissionID, suffix, status, rea
return err
}
record.Status, record.Reason, record.ReceiptID, record.UpdatedAt = status, reason, receiptID, j.clock.Now().UTC()
if status != "committed" {
record.ReceiptID = ""
record.Receipt = nil
}
raw, err := encodeJSON(record)
if err != nil {
return err
Expand All @@ -208,7 +278,54 @@ func (j *Journal) finalize(ctx context.Context, admissionID, suffix, status, rea
}

func (j *Journal) Commit(ctx context.Context, receipt protocol.TransitionReceipt) error {
return j.finalize(ctx, receipt.AdmissionID, ".committed", "committed", "", receipt.ID)
if err := receipt.Validate(); err != nil {
return err
}
name, err := journalName(receipt.AdmissionID, ".pending")
if err != nil {
return err
}
path, ok := j.activePath(name)
if !ok {
return fmt.Errorf("transaction journal path is not bound for %s", receipt.AdmissionID)
}
record, err := readJournal(path)
if err != nil {
return err
}
if record.Admission.ID != receipt.AdmissionID || record.TransitionID != receipt.TransitionID {
return fmt.Errorf("transition fact does not match its transaction journal")
}
record.Status, record.Reason, record.ReceiptID, record.UpdatedAt = "committed", "", receipt.ID, j.clock.Now().UTC()
receiptCopy := receipt
record.Receipt = &receiptCopy
raw, err := encodeJSON(record)
if err != nil {
return err
}
// Persist the complete fact into the pending record, then atomically rename
// that same record. A crash exposes either recovery-required pending work or
// one canonical committed fact, never a separate success that outruns it.
if err := atomicWrite(path, raw, 0o600); err != nil {
return err
}
finalPath := strings.TrimSuffix(path, ".pending") + ".committed"
if _, statErr := os.Stat(finalPath); statErr == nil {
return fmt.Errorf("committed transaction journal already exists for %s", receipt.AdmissionID)
} else if !os.IsNotExist(statErr) {
return statErr
}
if err := replaceFile(path, finalPath); err != nil {
return err
}
if err := syncDirectory(filepath.Dir(path)); err != nil {
if rollbackErr := replaceFile(finalPath, path); rollbackErr != nil {
return errors.Join(err, fmt.Errorf("restore pending journal after directory sync failure: %w", rollbackErr))
}
return err
}
j.unbind(name)
return nil
}

func (j *Journal) Abort(ctx context.Context, admissionID, reason string) error {
Expand All @@ -218,6 +335,7 @@ func (j *Journal) Abort(ctx context.Context, admissionID, reason string) error {
func (j *Journal) RequireRecovery(ctx context.Context, admissionID, reason string) error {
return j.update(ctx, admissionID, func(record *journalRecord) {
record.Status, record.Reason = "recovery-required", reason
record.ReceiptID, record.Receipt = "", nil
})
}

Expand Down
Loading