reactive

package
v0.1.0 Latest Latest
Warning

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

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

Documentation

Overview

Package reactive carries live page values between the server and the browser over one WebSocket per page. The web package opens and runs the connection, and generated code sends values through it; app code uses Topic and Broadcast.

A template declares a live value with <ssr:var name="lastSeen" type="string" reactive="true"/>. Data sets its first value, and the page's Subscribe hook keeps it current while the page is open. Subscribe usually waits on a Topic or a Broadcast that the code which changes the value publishes to, and returns when its context ends:

func (p *DP) Subscribe(ctx context.Context, r *web.Request, state *ReactiveState) error {
	sub := p.d.LastSeen.Subscribe(r.URLParam("login"))
	defer sub.Close()
	for {
		select {
		case <-ctx.Done():
			return nil
		case v := <-sub.Updates():
			state.SetLastSeen(v)
		}
	}
}

With client-writable="true" the page writes the value too, through ssr:bind or ssr.set, and Validate<Name> checks every write: its web.Error message reaches the page's ssr.onError.

A live value in {{ }} reaches the page as text, never as markup. A viewer has at most 8 live connections at a time. A connection passes the page's access rules and Guards when it opens, and opens again at least every 10 minutes.

Read more in the guide docs/guides/live-values.md and the task docs/tasks/add-live-value.md, which aicoded explain and the MCP tool howto print as guides/live-values and tasks/add-live-value.

Index

Examples

Constants

View Source
const (
	KindText = "text" // shown with textContent
	KindHTML = "html" // markup rendered by the page's escapers
)

Kinds of Binding.

View Source
const (
	CodeUnknownRoute = "unknown_route"
	CodeValidation   = "validation_failed"
	CodeDecode       = "decode_error"
)

Codes of err frames.

View Source
const StatusReauth websocket.StatusCode = 4000

StatusReauth closes a connection that reached its maximum age; the browser reconnects at once, so access is checked again.

Variables

This section is empty.

Functions

func Decode

func Decode[T any](raw json.RawMessage) (T, error)

Decode reads a written value. Native inputs send strings, so a JSON string is parsed into a bool or number T; anything else is JSON for T.

Generated code only.

func ParseValue

func ParseValue[T any](s string) (T, error)

ParseValue parses text a page sent into T: a string, a bool or a number. A number must fit T, and a float must be finite.

Types

type AckMsg

type AckMsg struct {
	T        string `json:"t"`
	RouteKey string `json:"routeKey"`
	Var      string `json:"var"`
}

AckMsg accepts a write.

type Binding

type Binding struct {
	Kind  string `json:"kind"`
	Value string `json:"value"`
}

Binding is the current value of one live site on a page.

func HTMLBinding

func HTMLBinding(markup string) Binding

HTMLBinding returns a binding holding markup the page's escapers produced.

Generated code only.

func TextBinding

func TextBinding(text string) Binding

TextBinding returns a binding the browser shows as plain text.

Generated code only.

type Broadcast

type Broadcast[V any] struct {
	// contains filtered or unexported fields
}

Broadcast is an in-process pub/sub fan-out to every subscriber, for events that concern every connected page: a counter of users online, a maintenance banner.

Each subscription buffers one value. Publish never blocks and the freshest value wins. Use V = struct{} for a plain "something changed" signal.

All methods are safe for concurrent use.

Example

A Broadcast reaches every subscription, such as one for each open home page that counts the visitors online.

package main

import (
	"fmt"

	"aicoded.dev/framework/web/reactive"
)

func main() {
	online := reactive.NewBroadcast[int]()
	a, b := online.Subscribe(), online.Subscribe() // two open home pages
	defer a.Close()

	online.Publish(online.TotalSubs())
	fmt.Println(<-a.Updates(), <-b.Updates())

	b.Close() // the second page closed
	online.Publish(online.TotalSubs())
	fmt.Println(<-a.Updates())
}
Output:
2 2
1

func NewBroadcast

func NewBroadcast[V any]() *Broadcast[V]

NewBroadcast creates an empty Broadcast.

func (*Broadcast[V]) Publish

func (b *Broadcast[V]) Publish(v V)

Publish sends v to every subscription. It never blocks.

func (*Broadcast[V]) Subscribe

func (b *Broadcast[V]) Subscribe() *BroadcastSub[V]

Subscribe registers a new subscription. Close it when done, usually with defer.

func (*Broadcast[V]) TotalSubs

func (b *Broadcast[V]) TotalSubs() int

TotalSubs returns the number of subscriptions.

type BroadcastSub

type BroadcastSub[V any] struct {
	// contains filtered or unexported fields
}

BroadcastSub is a subscription returned by Broadcast.Subscribe.

func (*BroadcastSub[V]) Close

func (s *BroadcastSub[V]) Close()

Close removes the subscription from its Broadcast. Further calls do nothing.

func (*BroadcastSub[V]) Updates

func (s *BroadcastSub[V]) Updates() <-chan V

Updates returns the delivery channel. Read it in a select with ctx.Done(). The channel is never closed.

type CallMsg

type CallMsg struct {
	T        string          `json:"t"`
	ID       uint32          `json:"id"`
	RouteKey string          `json:"routeKey"`
	Name     string          `json:"name"`
	Args     json.RawMessage `json:"args"`
}

CallMsg asks the server to run one of a page's calls.

type Conn

type Conn struct {
	// contains filtered or unexported fields
}

Conn is the live connection of one open page.

func Accept

func Accept(w http.ResponseWriter, r *http.Request) (*Conn, error)

Accept opens a live connection for the request. It clears the read and write deadlines the app's HTTP server set, which would otherwise cut the connection after the request timeout.

Generated code only.

func (*Conn) Ack

func (c *Conn) Ack(ctx context.Context, msg WriteMsg)

Ack accepts a write.

func (*Conn) Close

func (c *Conn) Close(code websocket.StatusCode, reason string)

Close closes the connection with a status and a reason the browser sees.

func (*Conn) Enqueue

func (c *Conn) Enqueue(key string, b Binding)

Enqueue queues b for SendLoop to send as the value of key. A value not sent yet is replaced, so only the latest value of a key goes out. Enqueue never blocks.

func (*Conn) ReadFrames

func (c *Conn) ReadFrames(ctx context.Context, onWrite func(WriteMsg), onCall func(CallMsg)) error

ReadFrames reads the page's frames until reading fails, and returns that error. Write frames go to onWrite and call frames to onCall; a nil handler ignores its frames, and frames of any other type are ignored.

Every frame counts against a budget of writeRate frames a second with bursts of writeBurst; a page over it is closed with StatusPolicyViolation. A frame that is not a JSON text frame closes the connection with StatusUnsupportedData.

func (*Conn) Reject

func (c *Conn) Reject(ctx context.Context, msg WriteMsg, text, code string)

Reject refuses a write with a message and one of the Code constants.

func (*Conn) Send

func (c *Conn) Send(ctx context.Context, v any) error

Send writes v as one JSON frame at once. If the client has not taken the frame within sendTimeout, Send fails and the connection is closed.

func (*Conn) SendLoop

func (c *Conn) SendLoop(ctx context.Context) error

SendLoop sends queued values as patch frames until ctx is done, when it returns nil, or until a send fails. A client that does not take a patch within sendTimeout is too slow and is cut off.

type ErrMsg

type ErrMsg struct {
	T        string `json:"t"`
	RouteKey string `json:"routeKey"`
	Var      string `json:"var"`
	Msg      string `json:"msg"`
	Code     string `json:"code"`
}

ErrMsg refuses a write.

func NewErr

func NewErr(routeKey, varName, text, code string) ErrMsg

NewErr returns an err frame.

type FailMsg

type FailMsg struct {
	T      string `json:"t"`
	ID     uint32 `json:"id"`
	Status int    `json:"status"`
	Msg    string `json:"msg"`
}

FailMsg answers a call that failed.

func NewFail

func NewFail(id uint32, status int, text string) FailMsg

NewFail returns a fail frame.

type InitMsg

type InitMsg struct {
	T        string             `json:"t"`
	Bindings map[string]Binding `json:"bindings"`
}

InitMsg carries every live value of the page when a connection opens.

func NewInit

func NewInit(bindings map[string]Binding) InitMsg

NewInit returns the init frame.

type PatchMsg

type PatchMsg struct {
	T   string `json:"t"`
	Key string `json:"key"`
	Binding
}

PatchMsg carries a changed live value.

type ResultMsg

type ResultMsg struct {
	T     string          `json:"t"`
	ID    uint32          `json:"id"`
	Value json.RawMessage `json:"value"`
}

ResultMsg answers a call.

func NewResult

func NewResult(id uint32, value json.RawMessage) ResultMsg

NewResult returns a result frame; value is already JSON.

type Topic

type Topic[K comparable, V any] struct {
	// contains filtered or unexported fields
}

Topic is a keyed in-process pub/sub fan-out. Use it to wake up Subscribe goroutines when something they care about changes elsewhere in the process. K is the routing key (for example a user id) and V is the payload.

Each subscription buffers one value. Publish never blocks and the freshest value wins: a value still buffered is replaced by the new one, so a slow subscriber sees only the latest.

Use V = struct{} as a dirty bit ("something changed for K, query again"), or a concrete V when the publisher already has the new value.

All methods are safe for concurrent use.

Example

A Topic wakes the Subscribe hooks that watch one key, such as the open pages of one user. A subscription holds only the newest value, and Publish never blocks.

package main

import (
	"fmt"

	"aicoded.dev/framework/web/reactive"
)

func main() {
	lastSeen := reactive.NewTopic[string, string]()

	sub := lastSeen.Subscribe("alice") // in Subscribe of the page /users/alice/info
	defer sub.Close()

	lastSeen.Publish("bob", "10:41") // a page of bob's: nothing for alice's pages
	lastSeen.Publish("alice", "10:42")
	lastSeen.Publish("alice", "10:43") // replaces 10:42, which no one read

	fmt.Println(<-sub.Updates()) // Subscribe passes it to state.SetLastSeen
}
Output:
10:43

func NewTopic

func NewTopic[K comparable, V any]() *Topic[K, V]

NewTopic creates an empty Topic.

func (*Topic[K, V]) Len

func (t *Topic[K, V]) Len() int

Len returns the number of keys with at least one subscription.

func (*Topic[K, V]) Publish

func (t *Topic[K, V]) Publish(key K, v V)

Publish sends v to every subscription for key. It never blocks.

func (*Topic[K, V]) Subscribe

func (t *Topic[K, V]) Subscribe(key K) *TopicSub[V]

Subscribe registers a new subscription for key. Close it when done, usually with defer.

func (*Topic[K, V]) TotalSubs

func (t *Topic[K, V]) TotalSubs() int

TotalSubs returns the number of subscriptions across all keys.

type TopicSub

type TopicSub[V any] struct {
	// contains filtered or unexported fields
}

TopicSub is a subscription returned by Topic.Subscribe.

func (*TopicSub[V]) Close

func (s *TopicSub[V]) Close()

Close removes the subscription from its Topic. Further calls do nothing.

func (*TopicSub[V]) Updates

func (s *TopicSub[V]) Updates() <-chan V

Updates returns the delivery channel. Read it in a select with ctx.Done(). The channel is never closed.

type WriteMsg

type WriteMsg struct {
	T        string          `json:"t"`
	RouteKey string          `json:"routeKey"`
	Var      string          `json:"var"`
	Value    json.RawMessage `json:"value"`
}

WriteMsg is a value the page sends for a client-writable variable.

Directories

Path Synopsis
Package client holds the browser side of live pages as TypeScript.
Package client holds the browser side of live pages as TypeScript.

Jump to

Keyboard shortcuts

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