diff --git a/README.md b/README.md index 6b74190..babdd16 100644 --- a/README.md +++ b/README.md @@ -150,6 +150,29 @@ newest value and its arrival time; `Watch` gives dynamic consumers (an HTTP handler, an SSE stream) latest-wins updates that may skip intermediate values but can never stall the publishers. +Code outside the runtime can project a channel through the same constructor the +runtime uses: + +```go +updates := make(chan printerState, 1) +latest := backplane.NewLatest((<-chan printerState)(updates)) + +updates <- printing +``` + +The caller owns `updates`: `NewLatest` reads it but never closes it. Closing the +channel projects any buffered values, retains the final one, and then closes all +watchers. The projection runs asynchronously, so a completed send or close is +not a read-after-write barrier for `Load`; wait for `Watch` when synchronisation +matters. Passing a nil channel panics. `Latest` itself deliberately has no +`Publish` or `Close` method, so passing it to a module does not also hand that +module control of the source. + +Runtime-owned projections use a private one-value, latest-wins input so they do +not add backpressure to the topic. Topic completion waits for that input to be +drained, preserving the guarantee that `Load` sees the final accepted value +after `Run` returns. + ## What backplane is not - **Not a message broker** — no persistence, cross-process transport, or QoS. diff --git a/doc.go b/doc.go index 16bab53..acc9483 100644 --- a/doc.go +++ b/doc.go @@ -76,6 +76,8 @@ // [Latest] is the deliberate boundary between streams and state: it retains // the newest value and its arrival time, and its watchers get latest-wins // delivery that may skip intermediate values but never backpressures the -// topic. Use a subscriber when every value matters, a Latest when only the -// current value does. +// source. Use a subscriber when every value matters, a Latest when only the +// current value does. The runtime projects topics through [NewLatest]; callers +// outside the runtime can use the same constructor with a channel they own, +// without gaining write or close methods on Latest itself. package backplane diff --git a/example_test.go b/example_test.go index 46f4223..80dd1c0 100644 --- a/example_test.go +++ b/example_test.go @@ -118,3 +118,17 @@ func ExampleLatest() { // Output: // printer-1 is idle } + +func ExampleNewLatest() { + updates := make(chan string, 1) + latest := backplane.NewLatest((<-chan string)(updates)) + + updates <- "ready" + close(updates) + + for state := range latest.Watch(context.Background()) { + fmt.Println(state) + } + // Output: + // ready +} diff --git a/latest.go b/latest.go index 3d19dc3..1146868 100644 --- a/latest.go +++ b/latest.go @@ -9,12 +9,49 @@ import ( var latestProjectionType = reflect.TypeFor[latestProjection]() -// latestProjection is the untyped view of a *Latest[T] used by the runtime to -// detect the parameter, feed it values, and close it when its topic completes. +// latestProjection is the untyped view used to recognise *Latest[T] during +// signature inspection. type latestProjection interface { messageType() reflect.Type - publish(value reflect.Value, receivedAt time.Time) - close() +} + +// latestFactory bridges signature reflection back into a method where T is +// known, allowing the runtime to use NewLatest rather than constructing or +// mutating a projection through a second path. +type latestFactory interface { + newLatest() (latestProjection, latestInput) +} + +// latestInput is the runtime-owned write side of the channel consumed by a +// Latest. It is deliberately private: modules only receive the projection. +type latestInput interface { + offer(reflect.Value) + closeAndWait() +} + +type typedLatestInput[T any] struct { + updates chan T + done <-chan struct{} +} + +func (input *typedLatestInput[T]) offer(value reflect.Value) { + update := value.Interface().(T) + select { + case input.updates <- update: + default: + // Runtime inputs hold one pending value. Replace it instead of + // backpressuring the topic; this projection is latest-wins by design. + select { + case <-input.updates: + default: + } + input.updates <- update + } +} + +func (input *typedLatestInput[T]) closeAndWait() { + close(input.updates) + <-input.done } // Latest is a lossy, latest-wins projection of a topic: it retains the most @@ -22,19 +59,46 @@ type latestProjection interface { // instead of <-chan T when it needs current state, or change notifications // that must never backpressure the topic. // -// The runtime creates one Latest per topic and hands the same instance to -// every module that declares it. The zero value holds no value and stays -// empty forever; only the runtime populates a Latest. +// The runtime creates one Latest per topic with [NewLatest] and hands the same +// instance to every module that declares it. External callers can use NewLatest +// to project a channel they own. Latest itself exposes no operation that can +// publish a value or complete the projection. The zero value holds no value and +// stays empty forever. type Latest[T any] struct { mu sync.Mutex value T receivedAt time.Time hasValue bool closed bool - done chan struct{} // created on first Watch, closed with the topic + done chan struct{} // created by NewLatest or first Watch, closed with the source watchers map[chan T]struct{} } +// NewLatest projects updates from a channel into a Latest. It starts a +// goroutine that consumes updates until the channel closes. Each value is +// timestamped when that goroutine receives it; closing updates closes every +// watcher after all buffered values have been projected. NewLatest never closes +// the supplied channel. +// +// Projection is asynchronous: a send or close returning does not by itself +// guarantee that Load has observed the value. Use Watch when synchronisation is +// required. NewLatest panics if updates is nil. +func NewLatest[T any](updates <-chan T) *Latest[T] { + if updates == nil { + panic("backplane: NewLatest called with nil update channel") + } + + latest := &Latest[T]{done: make(chan struct{})} + + go func() { + for value := range updates { + latest.store(value, time.Now()) + } + latest.close() + }() + return latest +} + // Load returns the most recent value, its arrival time, and whether any value // has been observed yet. func (l *Latest[T]) Load() (value T, receivedAt time.Time, ok bool) { @@ -46,9 +110,11 @@ func (l *Latest[T]) Load() (value T, receivedAt time.Time, ok bool) { // Watch returns a channel that converges on the most recent value: the // current value (if any) is delivered immediately, and each newer value // overwrites any undelivered one, so a slow watcher misses intermediate -// values rather than backpressuring the topic. The channel closes when the -// topic completes or ctx is cancelled, whichever comes first. Watching after -// the topic has completed yields the final value, then a closed channel. +// values rather than backpressuring the source. The channel closes when the +// source completes or ctx is cancelled, whichever comes first. For a runtime +// projection the source completes with its topic; for a projection made with +// NewLatest it completes when the updates channel closes. Watching after the +// source has completed yields the final value, then a closed channel. func (l *Latest[T]) Watch(ctx context.Context) <-chan T { watcher := make(chan T, 1) @@ -91,18 +157,22 @@ func (*Latest[T]) messageType() reflect.Type { return reflect.TypeFor[T]() } -func (l *Latest[T]) publish(value reflect.Value, receivedAt time.Time) { - typedValue := value.Interface().(T) +func (*Latest[T]) newLatest() (latestProjection, latestInput) { + updates := make(chan T, 1) + latest := NewLatest((<-chan T)(updates)) + return latest, &typedLatestInput[T]{updates: updates, done: latest.done} +} +func (l *Latest[T]) store(value T, receivedAt time.Time) { l.mu.Lock() defer l.mu.Unlock() - l.value = typedValue + l.value = value l.receivedAt = receivedAt l.hasValue = true for watcher := range l.watchers { select { - case watcher <- typedValue: + case watcher <- value: default: // The watcher has an undelivered value: replace it. Nothing else // sends on watcher, so after the drain the send cannot block. @@ -110,7 +180,7 @@ func (l *Latest[T]) publish(value reflect.Value, receivedAt time.Time) { case <-watcher: default: } - watcher <- typedValue + watcher <- value } } } @@ -138,9 +208,11 @@ func latestMessageType(parameterType reflect.Type) (reflect.Type, bool) { if parameterType.Kind() != reflect.Pointer || !parameterType.Implements(latestProjectionType) { return nil, false } - return newLatestProjection(parameterType).messageType(), true + projection := reflect.New(parameterType.Elem()).Interface().(latestProjection) + return projection.messageType(), true } -func newLatestProjection(parameterType reflect.Type) latestProjection { - return reflect.New(parameterType.Elem()).Interface().(latestProjection) +func newLatestProjection(parameterType reflect.Type) (latestProjection, latestInput) { + factory := reflect.New(parameterType.Elem()).Interface().(latestFactory) + return factory.newLatest() } diff --git a/latest_test.go b/latest_test.go index 1daa7c3..a6c41a2 100644 --- a/latest_test.go +++ b/latest_test.go @@ -26,6 +26,85 @@ func TestLatestZeroValueHoldsNothing(t *testing.T) { } } +func TestNewLatestRejectsNilUpdateChannel(t *testing.T) { + defer func() { + if recover() == nil { + t.Fatal("NewLatest accepted a nil update channel") + } + }() + + backplane.NewLatest[int](nil) +} + +func TestNewLatestProjectsChannelUpdates(t *testing.T) { + updates := make(chan int) + latest := backplane.NewLatest(updates) + defer close(updates) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + watcher := latest.Watch(ctx) + + sentAt := time.Now() + select { + case updates <- 42: + case <-ctx.Done(): + t.Fatal("NewLatest did not receive an update") + } + + select { + case got := <-watcher: + if got != 42 { + t.Fatalf("Watch() delivered %d, want 42", got) + } + case <-ctx.Done(): + t.Fatal("Watch() did not deliver the channel update") + } + + got, receivedAt, ok := latest.Load() + if !ok || got != 42 { + t.Fatalf("Load() = %d, %v; want 42, true", got, ok) + } + if receivedAt.Before(sentAt) || receivedAt.After(time.Now()) { + t.Fatalf("Load() arrival time = %v; want a time after the update was sent", receivedAt) + } +} + +func TestNewLatestClosesAfterUpdateChannelCloses(t *testing.T) { + updates := make(chan string) + latest := backplane.NewLatest(updates) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + watcher := latest.Watch(ctx) + + select { + case updates <- "final": + case <-ctx.Done(): + t.Fatal("NewLatest did not receive the final update") + } + select { + case <-watcher: + case <-ctx.Done(): + t.Fatal("Watch() did not deliver the final update") + } + close(updates) + + select { + case _, ok := <-watcher: + if ok { + t.Fatal("Watch() remained open after the update channel closed") + } + case <-ctx.Done(): + t.Fatal("Watch() did not close with the update channel") + } + + got, _, ok := latest.Load() + if !ok || got != "final" { + t.Fatalf("Load() = %q, %v; want final, true", got, ok) + } +} + func TestLatestRetainsTheFinalTopicValue(t *testing.T) { type printerState string diff --git a/topic.go b/topic.go index 1f33937..6d441cc 100644 --- a/topic.go +++ b/topic.go @@ -5,7 +5,6 @@ import ( "reflect" "slices" "sync" - "time" ) // topic carries every value published to one exact Go type from its @@ -21,6 +20,7 @@ type topic struct { inputs []*publisherEndpoint subscribers []subscription latest latestProjection + latestInput latestInput mu sync.Mutex remaining int // publisher endpoints not yet closed or returned @@ -71,7 +71,7 @@ func (t *topic) addSubscriber(parameterType reflect.Type, moduleDone chan struct func (t *topic) addLatest(parameterType reflect.Type) reflect.Value { if t.latest == nil { - t.latest = newLatestProjection(parameterType) + t.latest, t.latestInput = newLatestProjection(parameterType) } return reflect.ValueOf(t.latest) } @@ -118,8 +118,8 @@ func (t *topic) pump(ctx context.Context) { inputs = slices.Delete(inputs, chosen-2, chosen-1) case cancelled: // cancellation interrupts delivery: drain and drop default: - if t.latest != nil { - t.latest.publish(value, time.Now()) + if t.latestInput != nil { + t.latestInput.offer(value) } cancelled = !t.deliver(value, contextDone) } @@ -151,10 +151,12 @@ func (t *topic) deliver(value reflect.Value, contextDone reflect.SelectCase) boo } func (t *topic) finish() { + // Finish the projection before closing subscribers so every observer agrees + // on the topic's final state when subscription completion becomes visible. + if t.latestInput != nil { + t.latestInput.closeAndWait() + } for _, sub := range t.subscribers { sub.channel.Close() } - if t.latest != nil { - t.latest.close() - } }