diff --git a/README.md b/README.md index e6ad716..63840d8 100644 --- a/README.md +++ b/README.md @@ -93,6 +93,30 @@ rules: `Compile` is a one-shot convenience. To compile many sources with the same configuration, build a `Compiler` once with `mapper.NewCompiler(opts...)` and reuse it. +### Batch aggregation + +When processing a batch of events, call `Evaluate` per event and accumulate the results. Before writing to OpenFGA, pass the combined tuple slice to `Compact` to deduplicate and detect conflicts: + +```go +var tuples []mapper.Tuple +for _, event := range events { + result, err := m.Evaluate(ctx, event) + if err != nil { + panic(err) + } + tuples = append(tuples, result.Tuples...) +} + +compacted, err := mapper.Compact(tuples) +if err != nil { + // err is a *mapper.ConflictError — two incompatible desired states on the same relationship. + panic(err) +} +// compacted is ready to write to OpenFGA. +``` + +`Compact` applies the same dedup and conflict-detection semantics `Evaluate` runs per-record: exact-identity duplicates collapse to the first occurrence, repeated deletes on the same `(user, relation, object)` collapse, and incompatible desired states (write + delete, or two writes with differing condition/context on the same relationship) are reported as a `*ConflictError`. + ## Documentation - [Language specification](./docs/language-spec.md) — the canonical, user-facing mapping language spec. diff --git a/types.go b/types.go index fb895c0..5e0f52c 100644 --- a/types.go +++ b/types.go @@ -123,46 +123,43 @@ func uroKey(t language.Tuple) string { return k.Key() } -// postProcess deduplicates result tuples in place and detects conflicts in a single pass. -// A relationship is identified by its (user, relation, object) triple, matching how -// OpenFGA stores tuples and validates a Write batch. Two desired states on the same URO -// that cannot both be satisfied in one batch are conflicts: -// - a write and a delete on the same URO (ConflictWriteDelete); -// - two writes on the same URO whose condition or context differ (ConflictCompetingWrites). +// collapseTuples deduplicates and conflict-checks tuples in a single pass, +// reusing the tuples slice's backing array (write-pointer pattern). It +// implements the core loop logic shared by postProcess and Collapse. // -// Dedup: exact-identity duplicates (including condition and context) collapse to the first, -// as do repeated deletes on the same URO; order is preserved and the backing array reused. -// When tracing is enabled, all metadata is stored on Trace.PostProcess and all conflicts are -// collected. When tracing is disabled, returns immediately on the first conflict (fast path). -// Returns the first conflict as a *ConflictError, or nil. -func (r *Result) postProcess() error { - tracing := r.Trace != nil - +// When collectAll is false (fast path), the function returns immediately on +// the first ConflictError, returning the partially-collapsed prefix and the +// error. It always returns immediately on a ValidationError regardless of +// collectAll. When collectAll is true, all removed duplicates and all +// conflicts are accumulated; a ValidationError still causes an early return. +// +// Nil-interface invariant: firstErr is always a true nil interface or a +// non-nil concrete error value — a (*ConflictError)(nil) is never returned +// inside the interface. +func collapseTuples(tuples []language.Tuple, collectAll bool) (collapsed []language.Tuple, removed []language.Tuple, allConflicts []Conflict, firstErr error) { type uroState struct { hasWrite bool writeTuple language.Tuple hasDelete bool } - seenFull := make(map[string]struct{}, len(r.Tuples)) - states := make(map[string]*uroState, len(r.Tuples)) + seenFull := make(map[string]struct{}, len(tuples)) + states := make(map[string]*uroState, len(tuples)) - var removed []language.Tuple - var conflicts []Conflict var firstConflictErr *ConflictError recordConflict := func(ce *ConflictError) { if firstConflictErr == nil { firstConflictErr = ce } - conflicts = append(conflicts, ce.Conflict) + allConflicts = append(allConflicts, ce.Conflict) } w := 0 - for _, t := range r.Tuples { + for _, t := range tuples { if t.Action != language.ActionWrite && t.Action != language.ActionDelete { - r.Tuples = r.Tuples[:w] - return &language.ValidationError{ + tuples = tuples[:w] + return tuples, removed, allConflicts, &language.ValidationError{ Field: fmt.Sprintf("(%s, %s, %s)", t.User, t.Relation, t.Object), Message: fmt.Sprintf("unknown tuple action %q", t.Action), } @@ -170,7 +167,7 @@ func (r *Result) postProcess() error { fullKey := t.Key() if _, dup := seenFull[fullKey]; dup { - if tracing { + if collectAll { removed = append(removed, t) } continue @@ -191,7 +188,7 @@ func (r *Result) postProcess() error { case st.hasDelete: ce = &ConflictError{Conflict: conflict, Condition: t.Condition, Kind: ConflictWriteDelete} case st.hasWrite: - // Distinct from the stored write (an identical one is caught by seenFull), + // Distinct from the stored write (identical ones caught by seenFull), // so this is a second, incompatible desired state for the same relationship. ce = &ConflictError{Conflict: conflict, Kind: ConflictCompetingWrites} } @@ -202,7 +199,7 @@ func (r *Result) postProcess() error { case st.hasDelete: // A delete targets a relationship by URO regardless of condition, so a // second delete on the same URO is a duplicate. - if tracing { + if collectAll { removed = append(removed, t) } continue @@ -210,9 +207,9 @@ func (r *Result) postProcess() error { } if ce != nil { - if !tracing { - r.Tuples = r.Tuples[:w] - return ce + if !collectAll { + tuples = tuples[:w] + return tuples, nil, nil, ce } recordConflict(ce) } @@ -228,10 +225,41 @@ func (r *Result) postProcess() error { } seenFull[fullKey] = struct{}{} - r.Tuples[w] = t + tuples[w] = t w++ } - r.Tuples = r.Tuples[:w] + + collapsed = tuples[:w] + if firstConflictErr != nil { + firstErr = firstConflictErr + } + return collapsed, removed, allConflicts, firstErr +} + +// postProcess deduplicates result tuples in place and detects conflicts in a single pass. +// A relationship is identified by its (user, relation, object) triple, matching how +// OpenFGA stores tuples and validates a Write batch. Two desired states on the same URO +// that cannot both be satisfied in one batch are conflicts: +// - a write and a delete on the same URO (ConflictWriteDelete); +// - two writes on the same URO whose condition or context differ (ConflictCompetingWrites). +// +// Dedup: exact-identity duplicates (including condition and context) collapse to the first, +// as do repeated deletes on the same URO; order is preserved and the backing array reused. +// When tracing is enabled, all metadata is stored on Trace.PostProcess and all conflicts are +// collected. When tracing is disabled, returns immediately on the first conflict (fast path). +// Returns the first conflict as a *ConflictError, or nil. +func (r *Result) postProcess() error { + tracing := r.Trace != nil + + collapsed, removed, conflicts, firstErr := collapseTuples(r.Tuples, tracing) + r.Tuples = collapsed + + // On a ValidationError (any mode) or a ConflictError in non-tracing mode, + // skip TupleFilterOperations dedup — preserve the original early-exit behaviour. + _, isVal := firstErr.(*language.ValidationError) + if isVal || (firstErr != nil && !tracing) { + return firstErr + } // Dedup each TupleFilterOperation's desired-state tuples independently. // Desired state is a set per rule; duplicates from iterators are wasteful. @@ -246,10 +274,19 @@ func (r *Result) postProcess() error { } } - if firstConflictErr != nil { - return firstConflictErr - } - return nil + return firstErr +} + +// Compact deduplicates tuples by full identity and detects write/delete and +// competing-write conflicts on the same (user, relation, object) relationship. +// It applies the same post-processing semantics Evaluate runs per-record, but +// over an arbitrary tuple set — for callers reconciling tuples produced across +// multiple Evaluate calls (e.g. a batch of input records). Order is preserved, +// first occurrence wins. Returns the compacted tuples and the first conflict as +// a *ConflictError, or nil. +func Compact(tuples []language.Tuple) ([]language.Tuple, error) { + collapsed, _, _, firstErr := collapseTuples(tuples, false) + return collapsed, firstErr } // dedupTuples removes duplicate tuples by their full identity (including condition and context), diff --git a/types_test.go b/types_test.go index a483cc9..2eb75a8 100644 --- a/types_test.go +++ b/types_test.go @@ -196,6 +196,158 @@ func TestPostProcess_UROConflict(t *testing.T) { }) } +func TestCompact(t *testing.T) { + t.Parallel() + + t.Run("nil input — returns nil, no error", func(t *testing.T) { + collapsed, err := Compact(nil) + require.NoError(t, err) + assert.Nil(t, collapsed) + }) + + t.Run("empty input — returns empty, no error", func(t *testing.T) { + collapsed, err := Compact([]language.Tuple{}) + require.NoError(t, err) + assert.Empty(t, collapsed) + }) + + t.Run("single tuple — returned unchanged", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite}, + } + collapsed, err := Compact(input) + require.NoError(t, err) + assert.Equal(t, input, collapsed) + }) + + t.Run("distinct tuples — all preserved in order", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite}, + {User: "user:bob", Relation: "viewer", Object: "doc:1", Action: language.ActionDelete}, + } + collapsed, err := Compact(input) + require.NoError(t, err) + assert.Equal(t, input, collapsed) + }) + + t.Run("exact-identity duplicates — first occurrence wins, order preserved", func(t *testing.T) { + t1 := language.Tuple{User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite} + t2 := language.Tuple{User: "user:bob", Relation: "viewer", Object: "org:1", Action: language.ActionWrite} + input := []language.Tuple{t1, t2, t1} // t1 appears twice, t2 in between + collapsed, err := Compact(input) + require.NoError(t, err) + assert.Equal(t, []language.Tuple{t1, t2}, collapsed) + }) + + t.Run("order preservation with mixed duplicates", func(t *testing.T) { + t1 := language.Tuple{User: "user:a", Relation: "r", Object: "o:1", Action: language.ActionWrite} + t2 := language.Tuple{User: "user:b", Relation: "r", Object: "o:1", Action: language.ActionWrite} + t3 := language.Tuple{User: "user:c", Relation: "r", Object: "o:1", Action: language.ActionWrite} + input := []language.Tuple{t1, t2, t1, t3, t2} + collapsed, err := Compact(input) + require.NoError(t, err) + assert.Equal(t, []language.Tuple{t1, t2, t3}, collapsed) + }) + + t.Run("same URO differing condition — two writes are a ConflictCompetingWrites", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite, Condition: "cond_a"}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite, Condition: "cond_b"}, + } + _, err := Compact(input) + require.Error(t, err) + var ce *ConflictError + require.ErrorAs(t, err, &ce) + assert.Equal(t, ConflictCompetingWrites, ce.Kind) + assert.Equal(t, "user:alice", ce.User) + assert.Equal(t, "member", ce.Relation) + assert.Equal(t, "org:1", ce.Object) + }) + + t.Run("same URO differing context — two writes are a ConflictCompetingWrites", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite, Condition: "c", Context: map[string]any{"k": "v1"}}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite, Condition: "c", Context: map[string]any{"k": "v2"}}, + } + _, err := Compact(input) + require.Error(t, err) + var ce *ConflictError + require.ErrorAs(t, err, &ce) + assert.Equal(t, ConflictCompetingWrites, ce.Kind) + }) + + t.Run("write then delete on same URO — ConflictWriteDelete, condition from write", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite, Condition: "cond_a"}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete}, + } + _, err := Compact(input) + require.Error(t, err) + var ce *ConflictError + require.ErrorAs(t, err, &ce) + assert.Equal(t, ConflictWriteDelete, ce.Kind) + assert.Equal(t, "cond_a", ce.Condition) + }) + + t.Run("delete then write on same URO — ConflictWriteDelete", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite}, + } + _, err := Compact(input) + require.Error(t, err) + var ce *ConflictError + require.ErrorAs(t, err, &ce) + assert.Equal(t, ConflictWriteDelete, ce.Kind) + }) + + t.Run("repeated identical deletes — collapse to one, no error", func(t *testing.T) { + del := language.Tuple{User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete} + input := []language.Tuple{del, del, del} + collapsed, err := Compact(input) + require.NoError(t, err) + assert.Equal(t, []language.Tuple{del}, collapsed) + }) + + t.Run("repeated deletes on same URO with differing conditions — collapse to first, no error", func(t *testing.T) { + // URO-level dedup: delete targets (user,relation,object) regardless of condition. + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete, Condition: "cond_a"}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete, Condition: "cond_b"}, + } + collapsed, err := Compact(input) + require.NoError(t, err) + require.Len(t, collapsed, 1) + assert.Equal(t, "cond_a", collapsed[0].Condition) + }) + + t.Run("unknown action — returns ValidationError", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: "bogus"}, + } + _, err := Compact(input) + require.Error(t, err) + var ve *language.ValidationError + require.ErrorAs(t, err, &ve) + assert.Contains(t, ve.Message, `unknown tuple action`) + assert.Contains(t, ve.Message, `"bogus"`) + }) + + t.Run("fast path — only first conflict returned", func(t *testing.T) { + input := []language.Tuple{ + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionWrite}, + {User: "user:alice", Relation: "member", Object: "org:1", Action: language.ActionDelete}, + {User: "user:bob", Relation: "viewer", Object: "org:2", Action: language.ActionWrite}, + {User: "user:bob", Relation: "viewer", Object: "org:2", Action: language.ActionDelete}, + } + _, err := Compact(input) + require.Error(t, err) + var ce *ConflictError + require.ErrorAs(t, err, &ce) + assert.Equal(t, "user:alice", ce.User) // first conflict, not the second + }) +} + func TestErrorUnwrap(t *testing.T) { t.Parallel() tests := []struct {