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
24 changes: 24 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
105 changes: 71 additions & 34 deletions types.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,54 +123,51 @@ 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),
}
}

fullKey := t.Key()
if _, dup := seenFull[fullKey]; dup {
if tracing {
if collectAll {
removed = append(removed, t)
}
continue
Expand All @@ -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}
}
Expand All @@ -202,17 +199,17 @@ 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
}
}

if ce != nil {
if !tracing {
r.Tuples = r.Tuples[:w]
return ce
if !collectAll {
tuples = tuples[:w]
return tuples, nil, nil, ce
}
recordConflict(ce)
}
Expand All @@ -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.
Expand All @@ -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),
Expand Down
152 changes: 152 additions & 0 deletions types_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading