Add the conformance suite, echo fixture and integration guide
Some checks failed
ci / build (push) Has been cancelled

Completes FLUID-WP-0007. The seven minimal-conformance requirements and
the mechanically checkable architectural invariants are asserted as
tests rather than claimed in a README, because a conformance claim
nobody re-checks is one that quietly stops being true. Only the
checkable subset of the invariants is asserted; pretending a test can
settle the rest would be worse than leaving them to review.

TestFirstVerticalSlice runs all eleven steps of Blueprint 50 with no
human steps: two revisions, explicit routing, telemetry, a cohort
dimension, detected pressure, a hypothesis, a candidate, a 90/10
experiment, fitness comparison, promotion, and a complete audit trail.
Requests per completed task fall from 5.65 to 1.00 against a 1.20
target. A companion test runs the loop twice and requires the same
verdict, since a loop whose conclusion depended on run order would be
measuring the harness rather than the interface.

The failure-containment matrix covers Blueprint 34 directly: the data
plane keeps serving with the evidence store closed, with telemetry
wedged against a sink that never returns, after a failed build, after an
experiment rollback, and with the adaptive concurrency limit saturated.

Fixes a real bug the suite exposed. Drain closed the emitter outright,
so every request after the first flush emitted into a dead emitter and
was silently lost -- the kind of fault that makes a later measurement
quietly wrong rather than loudly broken. Emitter.Flush now waits for
delivery without stopping it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014KmVxhJ35tCo7rE7UnLwWu

Assistant: claude-code
Assistant-Model: opus
Assistant-Process: 1116572@bnt-lap001
Assistant-Session: 8ba9bb93-a72a-4883-b189-2499cce5c400
This commit is contained in:
tegwick 2026-09-04 08:21:49 +02:00
parent 61d8d8cabe
commit 55363905bc
16 changed files with 1885 additions and 23 deletions

View file

@ -29,6 +29,7 @@ type Emitter struct {
dropped atomic.Int64
written atomic.Int64
failed atomic.Int64
inflight atomic.Int64
stopOnce sync.Once
done chan struct{}
wg sync.WaitGroup
@ -76,6 +77,7 @@ func NewEmitter(sink Sink, opts EmitterOptions) *Emitter {
func (e *Emitter) Emit(ev contract.FluidTelemetry) {
select {
case e.ch <- ev:
e.inflight.Add(1)
default:
e.dropped.Add(1)
}
@ -89,25 +91,13 @@ func (e *Emitter) run(timeout time.Duration) {
if !ok {
return
}
ctx, cancel := context.WithTimeout(context.Background(), timeout)
if err := e.sink.Write(ctx, ev); err != nil {
e.failed.Add(1)
} else {
e.written.Add(1)
}
cancel()
e.deliver(ev, timeout)
case <-e.done:
// Drain what is already buffered, then stop.
for {
select {
case ev := <-e.ch:
ctx, cancel := context.WithTimeout(context.Background(), timeout)
if err := e.sink.Write(ctx, ev); err != nil {
e.failed.Add(1)
} else {
e.written.Add(1)
}
cancel()
e.deliver(ev, timeout)
default:
return
}
@ -116,6 +106,39 @@ func (e *Emitter) run(timeout time.Duration) {
}
}
// deliver writes one event and settles its in-flight accounting.
func (e *Emitter) deliver(ev contract.FluidTelemetry, timeout time.Duration) {
ctx, cancel := context.WithTimeout(context.Background(), timeout)
if err := e.sink.Write(ctx, ev); err != nil {
e.failed.Add(1)
} else {
e.written.Add(1)
}
cancel()
e.inflight.Add(-1)
}
// Flush waits for queued telemetry to reach the sink without stopping delivery.
//
// It exists for tests and for operational tooling that needs to read back what
// it just emitted. Close would also flush, but closing an emitter that is still
// serving traffic silently drops everything emitted afterwards -- which is the
// kind of bug that makes a later measurement quietly wrong rather than loudly
// broken.
//
// It returns false if the deadline passes with work still outstanding, so a
// caller can tell a slow sink from an empty one.
func (e *Emitter) Flush(timeout time.Duration) bool {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if e.inflight.Load() == 0 && len(e.ch) == 0 {
return true
}
time.Sleep(time.Millisecond)
}
return e.inflight.Load() == 0 && len(e.ch) == 0
}
// Close stops delivery after draining the buffer.
func (e *Emitter) Close() {
e.stopOnce.Do(func() { close(e.done) })