aggregate

package
v0.1.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 5, 2026 License: Apache-2.0 Imports: 12 Imported by: 0

Documentation

Overview

Package aggregate applies metering fact kinds onto a stream aggregate without mutating journal history (requirements 3.2, 3.3, 3.5, 13.6).

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrMixedCurrency is returned when present money observations disagree on
	// normalized currency identity, or when a present money lacks currency.
	ErrMixedCurrency = errors.New("metering/aggregate: mixed or empty currency")
	// ErrOverflow is returned when quantity or money nano aggregation would
	// wrap int64 arithmetic.
	ErrOverflow = errors.New("metering/aggregate: int64 overflow")
)

Stable sentinel errors for Apply (requirements 6.5, 6.7).

View Source
var (
	// ErrAmbiguousCumulative means two non-additive snapshots have the same
	// declared stream position and no revision ordering can select one.
	ErrAmbiguousCumulative = errors.New("metering/aggregate: ambiguous cumulative observation")
	// ErrIdentityConflict is the reducer-facing spelling of the shared replay
	// identity error. It remains one authority, owned by replay.
	ErrIdentityConflict = replay.ErrIdentityConflict
)

Functions

func EffectiveChargeGraphObservations

func EffectiveChargeGraphObservations(snapshot SnapshotV2) []metering.Observation

EffectiveChargeGraphObservations projects the reduced charge snapshot onto canonical observations while preserving the unfiltered effective coverage graph: every non-superseded charge item is retained, including incomplete fields, and every coverage edge is retained verbatim, including edges to an unavailable target. It differs from EffectiveChargeObservations only in not dropping incomplete charges or dangling edges, so a caller that owns closed graph validation can reject or diagnose an invalid graph instead of silently losing the edge. Superseded revisions and charge items removed by a same-item replacement remain excluded because the reducer owns that reduction.

func EffectiveChargeObservations

func EffectiveChargeObservations(snapshot SnapshotV2) []metering.Observation

EffectiveChargeObservations projects the reduced charge snapshot onto canonical observations. Only effective, field-complete charge items are returned, and every coverage edge is retained only when its target is also effective. Unknown or superseded targets therefore remain audit diagnostics rather than participating in effective coverage traversal.

func EffectiveChargeObservationsForCOGS

func EffectiveChargeObservationsForCOGS(snapshot SnapshotV2) []metering.Observation

EffectiveChargeObservationsForCOGS keeps an amount-bearing correction with only a late/unresolved predecessor reference available for COGS attribution. The caller still receives SnapshotV2.PendingSupersedes and must mark the result partial/non-payable. A correction whose known predecessor has an unusable baseline remains excluded, preserving the fail-closed charge taint.

Types

type ReducedCharge

type ReducedCharge struct {
	Scope         Scope
	ObservationID string
	Revision      uint64
	Charge        metering.ReportedCharge
	// Complete is field-local to this charge item. A correction can retain a
	// payable sibling charge while this item remains unresolved.
	Complete bool
	// contains filtered or unexported fields
}

ReducedCharge retains each reported charge item without recomputing or combining monetary amounts. Coverage completeness is exposed separately.

type ReducedMeasure

type ReducedMeasure struct {
	Scope   Scope
	Key     metering.ComponentKey
	Value   metering.Decimal
	Quality string
	// Complete reports that this reduced component has an effective value.
	// Incomplete source revisions remain in SnapshotV2.Observations for audit,
	// but a superseded unavailable revision must not poison its replacement.
	Complete          bool
	LastSequence      uint64
	LastRevision      uint64
	LastObservationID string
}

ReducedMeasure is one exact component value after source-scoped reduction. Value is always present; absent/unknown input measures are represented by Snapshot.Complete and the retained canonical Observations instead.

type Scope

type Scope struct {
	StoreID     string
	TenantID    string
	AccountKey  string
	Origin      string
	Acquisition string
	Perspective metering.EconomicPerspective
	Boundary    metering.Boundary
	Lifecycle   metering.LifecycleScope
	Subject     metering.SubjectRef
	StreamID    string
	ChargeScope string
	// contains filtered or unexported fields
}

Scope is the complete source/subject/charge scope of one reduced measure. A component key is not sufficient identity: local and provider evidence, distinct boundaries/perspectives, accounts, charge events, streams and subjects remain separate.

func ScopeFor

func ScopeFor(observation metering.Observation) Scope

ScopeFor exposes the canonical source scope used by ApplyObservations to consumers that project reduced measures at a later domain boundary. The reducer remains the sole owner of scope construction; callers must not rebuild source identity from component keys alone.

func (Scope) Key

func (s Scope) Key() string

Key returns a delimiter-safe deterministic scope key.

func (Scope) SubjectIdentity

func (s Scope) SubjectIdentity() string

type Snapshot

type Snapshot struct {
	StreamID      string
	Quantities    map[string]int64 // component -> value (present only)
	MoneyNano     int64
	MoneyCurrency string
	MoneyPresent  bool
	Unavailable   []string // fact IDs marked unavailable / unresolved
	Superseded    map[string]struct{}
	LastSequence  int64
}

Snapshot is the idempotent result of replaying an ordered fact stream.

func Apply

func Apply(facts []metering.Fact) (Snapshot, error)

Apply replays facts in Sequence order (stable by FactID on ties). Replaying the same ordered set yields the same Snapshot (restart hydration).

type SnapshotV2

type SnapshotV2 struct {
	Measures          []ReducedMeasure
	Charges           []ReducedCharge
	Observations      []metering.Observation
	PendingCoverage   []metering.ChargeCoverageRef
	PendingSupersedes []metering.ObservationRef
	// UnusablePredecessors identifies resolved correction links whose
	// predecessor did not contain a usable baseline for the corrected field.
	// It is distinct from PendingSupersedes: the revision is known, but the
	// correction cannot be applied safely.
	UnusablePredecessors []metering.ObservationRef
	Unavailable          []string
	Complete             bool
	Payable              bool
	Replayed             int
	LastSequence         uint64
}

SnapshotV2 is the deterministic reduction result for canonical V2 observations. It is source/subject/charge scoped and keeps canonical input observations for audit/replay rather than replacing them with a scalar map.

func ApplyFacts

func ApplyFacts(facts []metering.Fact) (SnapshotV2, error)

ApplyFacts is the explicit one-way V1 bridge. Historical facts are lifted through the SDK's trusted reader and then use this same V2 reduction owner; Fact identity and old hashes are never rewritten.

func ApplyObservations

func ApplyObservations(observations []metering.Observation) (SnapshotV2, error)

ApplyObservations deduplicates exact source-event revisions, validates supersession/coverage graphs, then reduces each source-scoped stream in declared Sequence/Revision order. Graph validation deliberately occurs only after exact replay deduplication so a valid duplicate node is harmless.

func (SnapshotV2) ValueFor

func (s SnapshotV2) ValueFor(observation metering.Observation, key metering.ComponentKey) string

ValueFor returns the reduced value for the same source scope and component key as observation. An invalid or absent value returns an empty string.

func (SnapshotV2) ValueForObservationID

func (s SnapshotV2) ValueForObservationID(observationID, component string) string

ValueForObservationID is a narrow diagnostic lookup for bridge/TCK tests. The observation ID is not itself a reduction key; all identity axes remain retained in the matching reduced measure.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL