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
23 changes: 23 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
14 changes: 14 additions & 0 deletions example_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
110 changes: 91 additions & 19 deletions latest.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,32 +9,96 @@ 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
// recently published value and its arrival time. A module declares *Latest[T]
// 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) {
Expand All @@ -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)

Expand Down Expand Up @@ -91,26 +157,30 @@ 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.
select {
case <-watcher:
default:
}
watcher <- typedValue
watcher <- value
}
}
}
Expand Down Expand Up @@ -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()
}
79 changes: 79 additions & 0 deletions latest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
16 changes: 9 additions & 7 deletions topic.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"reflect"
"slices"
"sync"
"time"
)

// topic carries every value published to one exact Go type from its
Expand All @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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()
}
}