diff --git a/go.mod b/go.mod index c33a0cd..87015ca 100644 --- a/go.mod +++ b/go.mod @@ -1,3 +1,21 @@ module github.com/tegwick/fluid-core go 1.22.2 + +require ( + github.com/dustin/go-humanize v1.0.1 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/ncruces/go-strftime v0.1.9 // indirect + github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect + golang.org/x/sys v0.19.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect + modernc.org/gc/v3 v3.0.0-20240107210532-573471604cb6 // indirect + modernc.org/libc v1.49.3 // indirect + modernc.org/mathutil v1.6.0 // indirect + modernc.org/memory v1.8.0 // indirect + modernc.org/sqlite v1.29.10 // indirect + modernc.org/strutil v1.2.0 // indirect + modernc.org/token v1.1.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..7b76d67 --- /dev/null +++ b/go.sum @@ -0,0 +1,32 @@ +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/ncruces/go-strftime v0.1.9 h1:bY0MQC28UADQmHmaF5dgpLmImcShSi2kHU9XLdhx/f4= +github.com/ncruces/go-strftime v0.1.9/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.19.0 h1:q5f1RH2jigJ1MoAWp2KTp3gm5zAGFUTarQZ5U386+4o= +golang.org/x/sys v0.19.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +modernc.org/gc/v3 v3.0.0-20240107210532-573471604cb6 h1:5D53IMaUuA5InSeMu9eJtlQXS2NxAhyWQvkKEgXZhHI= +modernc.org/gc/v3 v3.0.0-20240107210532-573471604cb6/go.mod h1:Qz0X07sNOR1jWYCrJMEnbW/X55x206Q7Vt4mz6/wHp4= +modernc.org/libc v1.49.3 h1:j2MRCRdwJI2ls/sGbeSk0t2bypOG/uvPZUsGQFDulqg= +modernc.org/libc v1.49.3/go.mod h1:yMZuGkn7pXbKfoT/M35gFJOAEdSKdxL0q64sF7KqCDo= +modernc.org/mathutil v1.6.0 h1:fRe9+AmYlaej+64JsEEhoWuAYBkOtQiMEU7n/XgfYi4= +modernc.org/mathutil v1.6.0/go.mod h1:Ui5Q9q1TR2gFm0AQRqQUaBWFLAhQpCwNcuhBOSedWPo= +modernc.org/memory v1.8.0 h1:IqGTL6eFMaDZZhEWwcREgeMXYwmW83LYW8cROZYkg+E= +modernc.org/memory v1.8.0/go.mod h1:XPZ936zp5OMKGWPqbD3JShgd/ZoQ7899TUuQqxY+peU= +modernc.org/sqlite v1.29.10 h1:3u93dz83myFnMilBGCOLbr+HjklS6+5rJLx4q86RDAg= +modernc.org/sqlite v1.29.10/go.mod h1:ItX2a1OVGgNsFh6Dv60JQvGfJfTPHPVpV6DF59akYOA= +modernc.org/strutil v1.2.0 h1:agBi9dp1I+eOnxXeiZawM8F4LawKv4NzGWSaLfyeNZA= +modernc.org/strutil v1.2.0/go.mod h1:/mdcBmfOibveCTBxUl5B5l6W+TTH1FXPLHZE6bTosX0= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= diff --git a/internal/runtime/registry.go b/internal/runtime/registry.go new file mode 100644 index 0000000..0a3d6db --- /dev/null +++ b/internal/runtime/registry.go @@ -0,0 +1,170 @@ +// Package runtime implements the FLUID deterministic data plane: the gateway, +// revision resolution, routing, backend connection and telemetry emission. +// +// Nothing here may depend on the evolution control plane to serve a request. +// ArchitectureBlueprint.md section 2 makes this an invariant: the interface +// runtime must continue to function when the Daimon, the model provider, the +// hypothesis store, the experiment controller and the AI budget are all +// unavailable. +package runtime + +import ( + "errors" + "fmt" + "sync" + + "github.com/tegwick/fluid-core/internal/contract" +) + +// Registry is the gateway's cached view of published control-plane state. +// +// It is deliberately a snapshot rather than a client. When the control plane +// dies the registry keeps answering from what it last held, which is what +// ArchitectureBlueprint.md section 34.6 requires: the runtime continues using +// cached published configuration, and no new promotions occur until +// control-plane consistency is restored. +type Registry struct { + mu sync.RWMutex + iface contract.InterfaceID + revisions map[contract.RevisionID]contract.Revision + policy contract.RoutingPolicy + hasPolicy bool +} + +// NewRegistry returns an empty registry for one interface. +func NewRegistry(iface contract.InterfaceID) *Registry { + return &Registry{ + iface: iface, + revisions: make(map[contract.RevisionID]contract.Revision), + } +} + +var ( + // ErrUnknownRevision is returned for a revision the registry has never seen. + ErrUnknownRevision = errors.New("unknown revision") + // ErrRevisionNotRoutable is returned for a revision that exists but must not + // receive traffic. + ErrRevisionNotRoutable = errors.New("revision not routable") + // ErrNoPolicy is returned before any routing policy has been loaded. + ErrNoPolicy = errors.New("no routing policy loaded") + // ErrWrongInterface guards against loading another interface's artifacts. + ErrWrongInterface = errors.New("artifact belongs to a different interface") + // ErrStalePolicy is returned when an older policy generation is offered. + ErrStalePolicy = errors.New("routing policy generation is not newer") +) + +// PutRevision publishes a revision descriptor into the registry. +// +// The descriptor is expected to have been signature-verified already; this +// method enforces only the structural conditions the router depends on. +func (r *Registry) PutRevision(d contract.Revision) error { + if d.Interface != r.iface { + return fmt.Errorf("%w: descriptor is for %q, registry serves %q", + ErrWrongInterface, d.Interface, r.iface) + } + if d.Runtime.Upstream == "" { + return fmt.Errorf("revision %s: descriptor has no runtime upstream", d.ID) + } + if !d.State.Valid() { + return fmt.Errorf("revision %s: unknown state %q", d.ID, d.State) + } + + r.mu.Lock() + defer r.mu.Unlock() + r.revisions[d.ID] = d + return nil +} + +// Revision returns a published descriptor. +func (r *Registry) Revision(id contract.RevisionID) (contract.Revision, error) { + r.mu.RLock() + defer r.mu.RUnlock() + d, ok := r.revisions[id] + if !ok { + return contract.Revision{}, fmt.Errorf("%w: %s", ErrUnknownRevision, id) + } + return d, nil +} + +// PutPolicy installs a routing policy. +// +// Generations are monotonic: an older policy is refused rather than applied. +// Without this a delayed delivery could silently roll traffic back to a +// superseded allocation, which would corrupt an in-flight experiment's +// measurement window. +func (r *Registry) PutPolicy(p contract.RoutingPolicy) error { + if p.Interface != r.iface { + return fmt.Errorf("%w: policy is for %q, registry serves %q", + ErrWrongInterface, p.Interface, r.iface) + } + + r.mu.Lock() + defer r.mu.Unlock() + + if r.hasPolicy && p.Generation <= r.policy.Generation { + return fmt.Errorf("%w: offered %d, holding %d", + ErrStalePolicy, p.Generation, r.policy.Generation) + } + r.policy = p + r.hasPolicy = true + return nil +} + +// Policy returns the current routing policy. +func (r *Registry) Policy() (contract.RoutingPolicy, error) { + r.mu.RLock() + defer r.mu.RUnlock() + if !r.hasPolicy { + return contract.RoutingPolicy{}, ErrNoPolicy + } + return r.policy, nil +} + +// routableStates lists the descriptor states the router may send traffic to. +// +// ArchitectureBlueprint.md section 5.3 requires rejecting unpublished, failed +// and retired revisions. "created" is unpublished; "verified" has passed tests +// but has not been exposed; "retired" is finished. +var routableStates = map[contract.RevisionState]bool{ + contract.RevisionStateExperiment: true, + contract.RevisionStateCandidate: true, + contract.RevisionStateStable: true, + contract.RevisionStateDeprecated: true, +} + +// CheckRoutable reports whether a revision may currently receive traffic from +// the given cohort. +func (r *Registry) CheckRoutable(id contract.RevisionID, cohort contract.CohortID) error { + d, err := r.Revision(id) + if err != nil { + return err + } + + if !routableStates[d.State] { + return fmt.Errorf("%w: %s is %s", ErrRevisionNotRoutable, id, d.State) + } + if d.Policy.SecurityCheck != contract.RevisionPolicySecurityCheckPassed { + return fmt.Errorf("%w: %s has security_check=%s", + ErrRevisionNotRoutable, id, d.Policy.SecurityCheck) + } + if d.Policy.PolicyCheck != nil && *d.Policy.PolicyCheck == contract.RevisionPolicyPolicyCheckFailed { + return fmt.Errorf("%w: %s failed its policy check", ErrRevisionNotRoutable, id) + } + + if d.Routing != nil && len(d.Routing.EligibleCohorts) > 0 { + if !containsCohort(d.Routing.EligibleCohorts, cohort) { + return fmt.Errorf("%w: cohort %q is not eligible for %s", + ErrRevisionNotRoutable, cohort, id) + } + } + return nil +} + +func containsCohort(list []contract.CohortID, want contract.CohortID) bool { + for _, c := range list { + if c == want { + return true + } + } + return false +} diff --git a/internal/runtime/resolver.go b/internal/runtime/resolver.go new file mode 100644 index 0000000..03f0aed --- /dev/null +++ b/internal/runtime/resolver.go @@ -0,0 +1,212 @@ +package runtime + +import ( + "fmt" + "hash/fnv" + "sort" + + "github.com/tegwick/fluid-core/internal/contract" +) + +// Request is the subset of an inbound request that revision resolution may +// consider. Nothing else is allowed to influence the decision: resolution must +// be a pure function of these fields and the loaded policy, or it stops being +// auditable (ArchitectureBlueprint.md section 5.2). +type Request struct { + // ExplicitRevision is a revision the consumer asked for by name. + ExplicitRevision contract.RevisionID + // BoundRevision comes from a client contract binding. + BoundRevision contract.RevisionID + // Cohort is the consumer's cohort assignment. + Cohort contract.CohortID + // Tenant identifies the calling tenant, where the interface is multi-tenant. + Tenant string + // ConsumerRef is the pseudonymous, stable consumer identity used to keep a + // long-lived consumer on one side of an experiment. + ConsumerRef string + // CorrelationID ties this request to its telemetry. + CorrelationID string +} + +// Resolution is the outcome of revision resolution, including why. +// +// The reason is not decoration. Blueprint section 5.2 requires resolution to be +// auditable, and "which revision served this request" is unanswerable later +// without recording how it was chosen. +type Resolution struct { + Revision contract.RevisionID + Reason contract.FluidTelemetryResolutionReason + Experiment *contract.ExperimentID + // PolicyGeneration records which policy produced this decision. + PolicyGeneration int64 +} + +// Resolver implements the deterministic precedence chain. +type Resolver struct { + registry *Registry + // allowExplicit controls whether consumers may pin a revision by name. Some + // interfaces want this for migration testing; others must not expose it. + allowExplicit bool +} + +// NewResolver returns a resolver over reg. +func NewResolver(reg *Registry, allowExplicit bool) *Resolver { + return &Resolver{registry: reg, allowExplicit: allowExplicit} +} + +// Resolve selects the revision that will serve req. +// +// The order is fixed by ArchitectureBlueprint.md section 5.2: +// +// explicit revision -> bound client contract -> experiment assignment -> stable default +// +// Each step is skipped rather than failed when the candidate is not routable, +// so a retired pin or an ineligible cohort degrades to the default instead of +// erroring the request. +func (r *Resolver) Resolve(req Request) (Resolution, error) { + policy, err := r.registry.Policy() + if err != nil { + return Resolution{}, err + } + + if r.allowExplicit && req.ExplicitRevision != "" { + if err := r.registry.CheckRoutable(req.ExplicitRevision, req.Cohort); err != nil { + // An explicit request for something unroutable is a consumer error + // worth surfacing, not something to silently reinterpret. + return Resolution{}, fmt.Errorf("explicit revision %s: %w", req.ExplicitRevision, err) + } + return Resolution{ + Revision: req.ExplicitRevision, + Reason: contract.FluidTelemetryResolutionReasonExplicitRevision, + PolicyGeneration: policy.Generation, + }, nil + } + + if req.BoundRevision != "" { + if err := r.registry.CheckRoutable(req.BoundRevision, req.Cohort); err == nil { + return Resolution{ + Revision: req.BoundRevision, + Reason: contract.FluidTelemetryResolutionReasonBoundContract, + PolicyGeneration: policy.Generation, + }, nil + } + } + + if rule, ok := matchRule(policy.Rules, req); ok { + chosen, ok := allocate(rule, req, policy.DefaultRevision) + if ok { + if err := r.registry.CheckRoutable(chosen, req.Cohort); err == nil { + reason := contract.FluidTelemetryResolutionReasonCohortRule + if rule.Experiment != nil { + reason = contract.FluidTelemetryResolutionReasonExperimentAssignment + } + return Resolution{ + Revision: chosen, + Reason: reason, + Experiment: rule.Experiment, + PolicyGeneration: policy.Generation, + }, nil + } + } + } + + if err := r.registry.CheckRoutable(policy.DefaultRevision, req.Cohort); err != nil { + return Resolution{}, fmt.Errorf("default revision %s: %w", policy.DefaultRevision, err) + } + return Resolution{ + Revision: policy.DefaultRevision, + Reason: contract.FluidTelemetryResolutionReasonStableDefault, + PolicyGeneration: policy.Generation, + }, nil +} + +// matchRule returns the first rule matching the request. Rules are evaluated in +// document order and the first match wins, so policy authors control precedence +// by ordering rather than by scoring. +func matchRule(rules []contract.RoutingPolicyRulesItem, req Request) (contract.RoutingPolicyRulesItem, bool) { + for _, rule := range rules { + if rule.Cohort != nil && *rule.Cohort != req.Cohort { + continue + } + if rule.Tenant != "" && rule.Tenant != req.Tenant { + continue + } + return rule, true + } + return contract.RoutingPolicyRulesItem{}, false +} + +// allocate picks a revision from a rule's traffic shares. +// +// Assignment is a deterministic function of the sticky key, so a given consumer +// lands on the same side of an experiment for its whole duration. Random +// per-request assignment would make within-consumer comparisons meaningless and +// would let a client observe both revisions at once. +func allocate(rule contract.RoutingPolicyRulesItem, req Request, fallback contract.RevisionID) (contract.RevisionID, bool) { + if len(rule.Allocation) == 0 { + return "", false + } + + // Sorting makes the traversal order independent of Go's map iteration, which + // is what turns a hash bucket into a stable assignment. + ids := make([]string, 0, len(rule.Allocation)) + var total float64 + for id, share := range rule.Allocation { + ids = append(ids, id) + total += float64(share) + } + sort.Strings(ids) + + if total <= 0 { + return "", false + } + + key := stickyKey(rule, req) + position := bucket(key) * total + + var cumulative float64 + for _, id := range ids { + cumulative += float64(rule.Allocation[id]) + if position < cumulative { + return contract.RevisionID(id), true + } + } + // Floating-point drift at the top of the range. + return contract.RevisionID(ids[len(ids)-1]), true +} + +// stickyKey chooses what keeps a consumer on one side of an experiment. +func stickyKey(rule contract.RoutingPolicyRulesItem, req Request) string { + mode := contract.RoutingPolicyRulesItemStickyByConsumerID + if rule.StickyBy != nil { + mode = *rule.StickyBy + } + + var subject string + switch mode { + case contract.RoutingPolicyRulesItemStickyByTenant: + subject = req.Tenant + case contract.RoutingPolicyRulesItemStickyByCorrelationID: + subject = req.CorrelationID + case contract.RoutingPolicyRulesItemStickyByNone: + subject = req.CorrelationID + default: + subject = req.ConsumerRef + } + + // Namespacing by experiment stops one consumer from landing in the same + // arm of every concurrent experiment, which would confound their results. + if rule.Experiment != nil { + return string(*rule.Experiment) + "\x00" + subject + } + return subject +} + +// bucket maps a key into [0, 1). +func bucket(key string) float64 { + h := fnv.New64a() + _, _ = h.Write([]byte(key)) + // 53 bits keeps the result exactly representable as a float64. + const mask = 1<<53 - 1 + return float64(h.Sum64()&mask) / float64(mask+1) +} diff --git a/internal/runtime/resolver_test.go b/internal/runtime/resolver_test.go new file mode 100644 index 0000000..20c7ca2 --- /dev/null +++ b/internal/runtime/resolver_test.go @@ -0,0 +1,275 @@ +package runtime + +import ( + "errors" + "testing" + + "github.com/tegwick/fluid-core/internal/contract" +) + +const testInterface contract.InterfaceID = "hall-publishing" + +// descriptor builds a minimal routable revision descriptor. +func descriptor(id contract.RevisionID, state contract.RevisionState, cohorts ...contract.CohortID) contract.Revision { + d := contract.Revision{ + SchemaVersion: "0.1", + ID: id, + Interface: testInterface, + State: state, + Contract: contract.RevisionContract{ + Type: contract.RevisionContractTypeOpenapi, + Digest: contract.Digest("sha256:" + zeros(64)), + }, + Runtime: contract.RevisionRuntime{Upstream: "http://adapter:8080"}, + Intent: contract.RevisionIntent{Version: "IEI-1"}, + Policy: contract.RevisionPolicy{ + Compatibility: contract.RevisionPolicyCompatibilityAdditive, + SecurityCheck: contract.RevisionPolicySecurityCheckPassed, + }, + } + if len(cohorts) > 0 { + d.Routing = &contract.RevisionRouting{EligibleCohorts: cohorts} + } + return d +} + +func zeros(n int) string { + b := make([]byte, n) + for i := range b { + b[i] = '0' + } + return string(b) +} + +func policyWith(defaultRev contract.RevisionID, gen int64, rules ...contract.RoutingPolicyRulesItem) contract.RoutingPolicy { + return contract.RoutingPolicy{ + SchemaVersion: "0.1", + Interface: testInterface, + Generation: gen, + DefaultRevision: defaultRev, + Rules: rules, + } +} + +func newFixture(t *testing.T, revs ...contract.Revision) *Registry { + t.Helper() + reg := NewRegistry(testInterface) + for _, r := range revs { + if err := reg.PutRevision(r); err != nil { + t.Fatalf("PutRevision(%s): %v", r.ID, err) + } + } + return reg +} + +func TestResolvePrecedence(t *testing.T) { + reg := newFixture(t, + descriptor("R-1", contract.RevisionStateStable), + descriptor("R-2", contract.RevisionStateExperiment), + descriptor("R-3", contract.RevisionStateCandidate), + ) + if err := reg.PutPolicy(policyWith("R-1", 1)); err != nil { + t.Fatal(err) + } + + res := NewResolver(reg, true) + + t.Run("explicit wins", func(t *testing.T) { + got, err := res.Resolve(Request{ExplicitRevision: "R-3", BoundRevision: "R-2"}) + if err != nil { + t.Fatal(err) + } + if got.Revision != "R-3" { + t.Errorf("got %s, want R-3", got.Revision) + } + if got.Reason != contract.FluidTelemetryResolutionReasonExplicitRevision { + t.Errorf("reason = %s", got.Reason) + } + }) + + t.Run("bound contract beats default", func(t *testing.T) { + got, err := res.Resolve(Request{BoundRevision: "R-2"}) + if err != nil { + t.Fatal(err) + } + if got.Revision != "R-2" || got.Reason != contract.FluidTelemetryResolutionReasonBoundContract { + t.Errorf("got %s via %s, want R-2 via bound_contract", got.Revision, got.Reason) + } + }) + + t.Run("falls through to stable default", func(t *testing.T) { + got, err := res.Resolve(Request{}) + if err != nil { + t.Fatal(err) + } + if got.Revision != "R-1" || got.Reason != contract.FluidTelemetryResolutionReasonStableDefault { + t.Errorf("got %s via %s, want R-1 via stable_default", got.Revision, got.Reason) + } + }) + + t.Run("explicit disabled is ignored", func(t *testing.T) { + strict := NewResolver(reg, false) + got, err := strict.Resolve(Request{ExplicitRevision: "R-3"}) + if err != nil { + t.Fatal(err) + } + if got.Revision != "R-1" { + t.Errorf("got %s, want the default R-1 when pinning is disabled", got.Revision) + } + }) +} + +func TestUnroutableRevisionsAreRefused(t *testing.T) { + reg := newFixture(t, + descriptor("R-1", contract.RevisionStateStable), + descriptor("R-created", contract.RevisionStateCreated), + descriptor("R-retired", contract.RevisionStateRetired), + ) + failed := descriptor("R-insecure", contract.RevisionStateStable) + failed.Policy.SecurityCheck = contract.RevisionPolicySecurityCheckFailed + if err := reg.PutRevision(failed); err != nil { + t.Fatal(err) + } + if err := reg.PutPolicy(policyWith("R-1", 1)); err != nil { + t.Fatal(err) + } + res := NewResolver(reg, true) + + for _, id := range []contract.RevisionID{"R-created", "R-retired", "R-insecure"} { + if _, err := res.Resolve(Request{ExplicitRevision: id}); err == nil { + t.Errorf("resolving %s should have been refused", id) + } else if !errors.Is(err, ErrRevisionNotRoutable) { + t.Errorf("resolving %s: got %v, want ErrRevisionNotRoutable", id, err) + } + } +} + +func TestBoundRevisionDegradesToDefault(t *testing.T) { + // A consumer pinned to a revision that has since retired should keep being + // served rather than start failing. + reg := newFixture(t, + descriptor("R-1", contract.RevisionStateStable), + descriptor("R-old", contract.RevisionStateRetired), + ) + if err := reg.PutPolicy(policyWith("R-1", 1)); err != nil { + t.Fatal(err) + } + + got, err := NewResolver(reg, false).Resolve(Request{BoundRevision: "R-old"}) + if err != nil { + t.Fatalf("a retired binding should degrade, not fail: %v", err) + } + if got.Revision != "R-1" { + t.Errorf("got %s, want R-1", got.Revision) + } +} + +func TestCohortEligibility(t *testing.T) { + reg := newFixture(t, + descriptor("R-1", contract.RevisionStateStable), + descriptor("R-2", contract.RevisionStateExperiment, "coding-agents"), + ) + if err := reg.PutPolicy(policyWith("R-1", 1)); err != nil { + t.Fatal(err) + } + + if err := reg.CheckRoutable("R-2", "coding-agents"); err != nil { + t.Errorf("eligible cohort refused: %v", err) + } + if err := reg.CheckRoutable("R-2", "partner-integrations"); err == nil { + t.Error("ineligible cohort accepted") + } +} + +func TestAllocationIsStickyAndProportional(t *testing.T) { + reg := newFixture(t, + descriptor("R-1", contract.RevisionStateStable), + descriptor("R-2", contract.RevisionStateExperiment), + ) + exp := contract.ExperimentID("E-1") + rule := contract.RoutingPolicyRulesItem{ + Experiment: &exp, + Allocation: map[string]contract.UnitInterval{"R-1": 0.9, "R-2": 0.1}, + } + if err := reg.PutPolicy(policyWith("R-1", 1, rule)); err != nil { + t.Fatal(err) + } + res := NewResolver(reg, false) + + // Stickiness: the same consumer must resolve identically every time. + first, err := res.Resolve(Request{ConsumerRef: "consumer-42"}) + if err != nil { + t.Fatal(err) + } + for i := 0; i < 50; i++ { + again, err := res.Resolve(Request{ConsumerRef: "consumer-42"}) + if err != nil { + t.Fatal(err) + } + if again.Revision != first.Revision { + t.Fatalf("assignment drifted: %s then %s", first.Revision, again.Revision) + } + } + + // Proportionality: roughly a tenth of consumers should see the candidate. + const n = 4000 + candidate := 0 + for i := 0; i < n; i++ { + got, err := res.Resolve(Request{ConsumerRef: consumerName(i)}) + if err != nil { + t.Fatal(err) + } + if got.Revision == "R-2" { + candidate++ + } + if got.Experiment == nil || *got.Experiment != exp { + t.Fatalf("experiment not recorded on resolution for consumer %d", i) + } + } + share := float64(candidate) / n + if share < 0.07 || share > 0.13 { + t.Errorf("candidate share %.3f, want approximately 0.10", share) + } +} + +func consumerName(i int) string { + digits := "0123456789" + out := []byte("consumer-") + if i == 0 { + return string(append(out, '0')) + } + var rev []byte + for i > 0 { + rev = append(rev, digits[i%10]) + i /= 10 + } + for j := len(rev) - 1; j >= 0; j-- { + out = append(out, rev[j]) + } + return string(out) +} + +func TestStalePolicyRefused(t *testing.T) { + reg := newFixture(t, descriptor("R-1", contract.RevisionStateStable)) + if err := reg.PutPolicy(policyWith("R-1", 5)); err != nil { + t.Fatal(err) + } + if err := reg.PutPolicy(policyWith("R-1", 4)); !errors.Is(err, ErrStalePolicy) { + t.Errorf("older generation accepted: %v", err) + } + if err := reg.PutPolicy(policyWith("R-1", 5)); !errors.Is(err, ErrStalePolicy) { + t.Errorf("equal generation accepted: %v", err) + } + if err := reg.PutPolicy(policyWith("R-1", 6)); err != nil { + t.Errorf("newer generation refused: %v", err) + } +} + +func TestWrongInterfaceRefused(t *testing.T) { + reg := NewRegistry(testInterface) + other := descriptor("R-1", contract.RevisionStateStable) + other.Interface = "some-other-api" + if err := reg.PutRevision(other); !errors.Is(err, ErrWrongInterface) { + t.Errorf("foreign descriptor accepted: %v", err) + } +}