@@ -3,10 +3,12 @@ package effects
33import (
44 "context"
55 "encoding/json"
6+ "errors"
67 "fmt"
78 "io"
89 "os"
910 "path/filepath"
11+ "slices"
1012 "strings"
1113 "time"
1214
@@ -30,17 +32,18 @@ func NewJournal(resolver ports.InvocationResolver, clock ports.Clock) (*Journal,
3032}
3133
3234type journalRecord struct {
33- SchemaVersion int `json:"schema_version"`
34- Admission protocol.Admission `json:"admission"`
35- TransitionID catalog.TransitionID `json:"transition_id"`
36- TransitionClass catalog.EventClass `json:"transition_class"`
37- ReconcilesProgram bool `json:"reconciles_program,omitempty"`
38- Status string `json:"status"`
39- Mutations []ports.ResourceMutation `json:"mutations,omitempty"`
40- Reason string `json:"reason,omitempty"`
41- ReceiptID string `json:"receipt_id,omitempty"`
42- CreatedAt time.Time `json:"created_at"`
43- UpdatedAt time.Time `json:"updated_at"`
35+ SchemaVersion int `json:"schema_version"`
36+ Admission protocol.Admission `json:"admission"`
37+ TransitionID catalog.TransitionID `json:"transition_id"`
38+ TransitionClass catalog.EventClass `json:"transition_class"`
39+ ReconcilesProgram bool `json:"reconciles_program,omitempty"`
40+ Status string `json:"status"`
41+ Mutations []ports.ResourceMutation `json:"mutations,omitempty"`
42+ Reason string `json:"reason,omitempty"`
43+ ReceiptID string `json:"receipt_id,omitempty"`
44+ Receipt * protocol.TransitionReceipt `json:"receipt,omitempty"`
45+ CreatedAt time.Time `json:"created_at"`
46+ UpdatedAt time.Time `json:"updated_at"`
4447}
4548
4649func journalName (id , suffix string ) (string , error ) {
@@ -110,9 +113,72 @@ func readJournal(path string) (journalRecord, error) {
110113 if err := record .Admission .ValidateIdentity (); err != nil || record .Admission .TransitionID != record .TransitionID {
111114 return journalRecord {}, fmt .Errorf ("invalid transaction admission in %s: %v" , path , err )
112115 }
116+ if record .Receipt != nil {
117+ if err := record .Receipt .Validate (); err != nil || record .Receipt .ID != record .ReceiptID || record .Receipt .AdmissionID != record .Admission .ID || record .Receipt .TransitionID != record .TransitionID {
118+ return journalRecord {}, fmt .Errorf ("invalid committed transition fact in %s: %v" , path , err )
119+ }
120+ receipt := record .Receipt
121+ admission := record .Admission
122+ if receipt .PrescriptionID != admission .PrescriptionID || receipt .TransitionVersion != admission .TransitionVersion || receipt .Program .Fingerprint != admission .ExpectedProgramFingerprint ||
123+ receipt .PriorStateRevision != admission .ExpectedStateRevision || receipt .SourceFingerprint != admission .ExpectedSnapshotFingerprint ||
124+ receipt .AuthorityFingerprint != admission .AuthorityFingerprint || ! slices .Equal (receipt .RequiredCapabilities , admission .RequiredCapabilities ) ||
125+ ! slices .Equal (receipt .GrantedCapabilities , admission .GrantedCapabilities ) || receipt .GoalID != admission .Goal .ID || receipt .GoalKind != admission .Goal .Kind ||
126+ receipt .DeliveryID != admission .Goal .DeliveryID || receipt .GoalScope != admission .GoalScope || receipt .GoalStatus != admission .GoalStatus {
127+ return journalRecord {}, fmt .Errorf ("committed transition fact in %s does not match its exact admission" , path )
128+ }
129+ if err := validateCommittedMutationFacts (record .TransitionClass , record .Mutations , receipt .CommittedEffects ); err != nil {
130+ return journalRecord {}, fmt .Errorf ("committed transition fact in %s: %w" , path , err )
131+ }
132+ }
133+ if strings .HasSuffix (path , ".committed" ) && (record .Status != "committed" || record .Receipt == nil ) {
134+ return journalRecord {}, fmt .Errorf ("committed transaction journal %s lacks its canonical transition fact" , path )
135+ }
113136 return record , nil
114137}
115138
139+ func validateCommittedMutationFacts (class catalog.EventClass , mutations []ports.ResourceMutation , facts []protocol.EffectFact ) error {
140+ resourceFacts := make ([]protocol.EffectFact , 0 , len (facts ))
141+ boundarySettled := false
142+ for _ , fact := range facts {
143+ if fact .Kind == protocol .EffectResourceMutation {
144+ resourceFacts = append (resourceFacts , fact )
145+ } else if fact .Kind == protocol .EffectBoundarySettled {
146+ boundarySettled = true
147+ }
148+ }
149+ if class == catalog .EventOwnedExternal && ! boundarySettled {
150+ return fmt .Errorf ("owned external transaction lacks a settled boundary fact" )
151+ }
152+ if len (resourceFacts ) != len (mutations ) {
153+ return fmt .Errorf ("resource fact count %d does not match staged mutation count %d" , len (resourceFacts ), len (mutations ))
154+ }
155+ matched := make ([]bool , len (resourceFacts ))
156+ for _ , mutation := range mutations {
157+ operation := "update"
158+ switch {
159+ case mutation .Delete :
160+ operation = "delete"
161+ case mutation .TargetLink != "" :
162+ operation = "symlink"
163+ case ! mutation .PriorExists :
164+ operation = "create"
165+ }
166+ prior := mutationStateFingerprint (mutation .PriorExists , mutation .Prior , mutation .PriorLink , mutation .Mode )
167+ result := mutationStateFingerprint (! mutation .Delete , mutation .Target , mutation .TargetLink , mutation .Mode )
168+ found := false
169+ for index , fact := range resourceFacts {
170+ if ! matched [index ] && fact .Target == mutation .Path && fact .Operation == operation && fact .PriorFingerprint == prior && fact .ResultingFingerprint == result {
171+ matched [index ], found = true , true
172+ break
173+ }
174+ }
175+ if ! found {
176+ return fmt .Errorf ("staged mutation %s has no exact committed effect fact" , mutation .Path )
177+ }
178+ }
179+ return nil
180+ }
181+
116182func (j * Journal ) update (ctx context.Context , admissionID string , update func (* journalRecord )) error {
117183 name , err := journalName (admissionID , ".pending" )
118184 if err != nil {
@@ -192,6 +258,10 @@ func (j *Journal) finalize(ctx context.Context, admissionID, suffix, status, rea
192258 return err
193259 }
194260 record .Status , record .Reason , record .ReceiptID , record .UpdatedAt = status , reason , receiptID , j .clock .Now ().UTC ()
261+ if status != "committed" {
262+ record .ReceiptID = ""
263+ record .Receipt = nil
264+ }
195265 raw , err := encodeJSON (record )
196266 if err != nil {
197267 return err
@@ -208,7 +278,54 @@ func (j *Journal) finalize(ctx context.Context, admissionID, suffix, status, rea
208278}
209279
210280func (j * Journal ) Commit (ctx context.Context , receipt protocol.TransitionReceipt ) error {
211- return j .finalize (ctx , receipt .AdmissionID , ".committed" , "committed" , "" , receipt .ID )
281+ if err := receipt .Validate (); err != nil {
282+ return err
283+ }
284+ name , err := journalName (receipt .AdmissionID , ".pending" )
285+ if err != nil {
286+ return err
287+ }
288+ path , ok := j .activePath (name )
289+ if ! ok {
290+ return fmt .Errorf ("transaction journal path is not bound for %s" , receipt .AdmissionID )
291+ }
292+ record , err := readJournal (path )
293+ if err != nil {
294+ return err
295+ }
296+ if record .Admission .ID != receipt .AdmissionID || record .TransitionID != receipt .TransitionID {
297+ return fmt .Errorf ("transition fact does not match its transaction journal" )
298+ }
299+ record .Status , record .Reason , record .ReceiptID , record .UpdatedAt = "committed" , "" , receipt .ID , j .clock .Now ().UTC ()
300+ receiptCopy := receipt
301+ record .Receipt = & receiptCopy
302+ raw , err := encodeJSON (record )
303+ if err != nil {
304+ return err
305+ }
306+ // Persist the complete fact into the pending record, then atomically rename
307+ // that same record. A crash exposes either recovery-required pending work or
308+ // one canonical committed fact, never a separate success that outruns it.
309+ if err := atomicWrite (path , raw , 0o600 ); err != nil {
310+ return err
311+ }
312+ finalPath := strings .TrimSuffix (path , ".pending" ) + ".committed"
313+ if _ , statErr := os .Stat (finalPath ); statErr == nil {
314+ return fmt .Errorf ("committed transaction journal already exists for %s" , receipt .AdmissionID )
315+ } else if ! os .IsNotExist (statErr ) {
316+ return statErr
317+ }
318+ if err := replaceFile (path , finalPath ); err != nil {
319+ return err
320+ }
321+ if err := syncDirectory (filepath .Dir (path )); err != nil {
322+ if rollbackErr := replaceFile (finalPath , path ); rollbackErr != nil {
323+ return errors .Join (err , fmt .Errorf ("restore pending journal after directory sync failure: %w" , rollbackErr ))
324+ }
325+ return err
326+ }
327+ j .unbind (name )
328+ return nil
212329}
213330
214331func (j * Journal ) Abort (ctx context.Context , admissionID , reason string ) error {
@@ -218,6 +335,7 @@ func (j *Journal) Abort(ctx context.Context, admissionID, reason string) error {
218335func (j * Journal ) RequireRecovery (ctx context.Context , admissionID , reason string ) error {
219336 return j .update (ctx , admissionID , func (record * journalRecord ) {
220337 record .Status , record .Reason = "recovery-required" , reason
338+ record .ReceiptID , record .Receipt = "" , nil
221339 })
222340}
223341
0 commit comments