hostclient
packageAPI reference for the hostclient
package.
Imports
(27)context
STD
strings
STD
testing
STD
time
INT
github.com/rfwlab/rfw/v2/js
STD
regexp
STD
fmt
INT
github.com/rfwlab/rfw/v2/dom
STD
encoding/json
STD
errors
STD
log
STD
sync
STD
sync/atomic
PKG
github.com/mirkobrombin/go-foundation/v2/core/caching
PKG
github.com/mirkobrombin/go-foundation/v2/core/resiliency
INT
github.com/rfwlab/rfw/v2/core
INT
github.com/rfwlab/rfw/v2/state
STD
crypto/sha256
STD
encoding/hex
STD
bufio
STD
encoding/binary
STD
io
STD
net
STD
net/url
STD
strconv
STD
syscall/js
STD
encoding/base64
installFakeCapacitorSSC
func installFakeCapacitorSSC(t *testing.T) js.Value
{
t.Helper()
originalCapacitor := js.Get("Capacitor")
originalTransport := js.Get("RFW_SSC_TRANSPORT")
js.Call("eval", fakeCapacitorSSCSource)
js.Set("RFW_SSC_TRANSPORT", sscTransportCapacitor)
t.Cleanup(func() {
js.Set("Capacitor", originalCapacitor)
js.Set("RFW_SSC_TRANSPORT", originalTransport)
js.Global().Delete("__fakeRFWSSC")
})
return js.Get("__fakeRFWSSC")
}
TestSSCTransportDefaultsToBrowser
Parameters
func TestSSCTransportDefaultsToBrowser(t *testing.T)
{
original := js.Get("RFW_SSC_TRANSPORT")
js.Global().Delete("RFW_SSC_TRANSPORT")
t.Cleanup(func() { js.Set("RFW_SSC_TRANSPORT", original) })
if got := sscTransport(); got != sscTransportBrowser {
t.Fatalf("transport = %q, want browser", got)
}
}
TestCapacitorTransportFailsClosedWhenPluginIsMissing
Parameters
func TestCapacitorTransportFailsClosedWhenPluginIsMissing(t *testing.T)
{
originalCapacitor := js.Get("Capacitor")
originalTransport := js.Get("RFW_SSC_TRANSPORT")
js.Set("Capacitor", js.NewDict().Value)
js.Set("RFW_SSC_TRANSPORT", sscTransportCapacitor)
t.Cleanup(func() {
js.Set("Capacitor", originalCapacitor)
js.Set("RFW_SSC_TRANSPORT", originalTransport)
})
_, err := dial(context.Background(), "wss://api.example.com/ws")
if err == nil || !strings.Contains(err.Error(), "plugins are unavailable") {
t.Fatalf("dial error = %v, want missing-plugin failure", err)
}
}
TestCapacitorTransportConnectsWritesReadsAndCloses
Parameters
func TestCapacitorTransportConnectsWritesReadsAndCloses(t *testing.T)
{
plugin := installFakeCapacitorSSC(t)
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
conn, err := dial(ctx, "wss://api.example.com/ws")
if err != nil {
t.Fatalf("dial: %v", err)
}
if got := plugin.Get("connects").Get("length").Int(); got != 1 {
t.Fatalf("connect calls = %d, want 1", got)
}
connect := plugin.Get("connects").Index(0)
if got := connect.Get("url").String(); got != "wss://api.example.com/ws" {
t.Fatalf("connect url = %q", got)
}
if connect.Get("id").String() == "" {
t.Fatal("connect id is empty")
}
if err := conn.writeJSON(wireMessage{Component: "ticker", Sequence: 1}); err != nil {
t.Fatalf("write: %v", err)
}
if got := plugin.Get("sends").Get("length").Int(); got != 1 {
t.Fatalf("send calls = %d, want 1", got)
}
if frame := plugin.Get("sends").Index(0).Get("data").String(); !strings.Contains(frame, `"component":"ticker"`) {
t.Fatalf("sent frame = %q", frame)
}
event := js.NewDict()
event.Set("type", "message")
event.Set("encoding", "text")
event.Set("data", `{"component":"ticker","sequence":1}`)
plugin.Get("callback").Invoke(event.Value)
frame, err := conn.read(ctx)
if err != nil {
t.Fatalf("read: %v", err)
}
if got := string(frame); got != `{"component":"ticker","sequence":1}` {
t.Fatalf("frame = %q", got)
}
if err := conn.close(); err != nil {
t.Fatalf("close: %v", err)
}
if got := plugin.Get("closes").Get("length").Int(); got != 1 {
t.Fatalf("close calls = %d, want 1", got)
}
}
TestCapacitorTransportRejectsMalformedBinaryFrame
Parameters
func TestCapacitorTransportRejectsMalformedBinaryFrame(t *testing.T)
{
plugin := installFakeCapacitorSSC(t)
conn, err := dial(context.Background(), "wss://api.example.com/ws")
if err != nil {
t.Fatalf("dial: %v", err)
}
t.Cleanup(func() { _ = conn.close() })
event := js.NewDict()
event.Set("type", "message")
event.Set("encoding", "base64")
event.Set("data", "not base64")
plugin.Get("callback").Invoke(event.Value)
if _, err := conn.read(context.Background()); err == nil || !strings.Contains(err.Error(), "invalid base64") {
t.Fatalf("read error = %v, want invalid base64", err)
}
}
TestCapacitorTransportKeepsCallbackAliveThroughAsyncClose
Parameters
func TestCapacitorTransportKeepsCallbackAliveThroughAsyncClose(t *testing.T)
{
plugin := installFakeCapacitorSSC(t)
plugin.Set("closeDelayMs", 10)
conn, err := dial(context.Background(), "wss://api.example.com/ws")
if err != nil {
t.Fatalf("dial: %v", err)
}
if err := conn.close(); err != nil {
t.Fatalf("close: %v", err)
}
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
if _, err := conn.read(ctx); err == nil || !strings.Contains(err.Error(), "code 1000") {
t.Fatalf("read error = %v, want native close", err)
}
if got := plugin.Get("terminalCallbacks").Int(); got != 1 {
t.Fatalf("terminal callbacks = %d, want 1", got)
}
}
TestCapacitorTransportDoesNotFallBackForUnknownMode
Parameters
func TestCapacitorTransportDoesNotFallBackForUnknownMode(t *testing.T)
{
original := js.Get("RFW_SSC_TRANSPORT")
js.Set("RFW_SSC_TRANSPORT", "native")
t.Cleanup(func() { js.Set("RFW_SSC_TRANSPORT", original) })
_, err := dial(context.Background(), "wss://api.example.com/ws")
if err == nil || !strings.Contains(err.Error(), "unsupported SSC transport") {
t.Fatalf("dial error = %v, want unsupported transport", err)
}
}
fakeElement
type fakeElement struct
Methods
Parameters
Returns
func (*fakeElement) Attr(name string) string
{
if name == hostExpectedAttr {
return e.expected
}
if e.attrStore != nil {
return e.attrStore[name]
}
return ""
}
Parameters
func (*fakeElement) SetAttr(name, value string)
{
if name == hostExpectedAttr {
e.expected = value
return
}
if e.attrStore == nil {
e.attrStore = make(map[string]string)
}
e.attrStore[name] = value
}
Fields
| Name | Type | Description |
|---|---|---|
| text | string | |
| expected | string | |
| exists | bool | |
| attrStore | map[string]string |
fakeRoot
type fakeRoot struct
Methods
Parameters
Returns
func (*fakeRoot) HostVar(name string) hostVarElement
{
if el, ok := r.elems[name]; ok {
return el
}
return &fakeElement{}
}
Parameters
func (*fakeRoot) SetHTML(html string)
{
r.html = html
r.elems = make(map[string]*fakeElement)
re := regexp.MustCompile(`<span[^>]*data-host-var="([^"]+)"[^>]*data-host-expected="([^"]*)"[^>]*>([^<]*)</span>`)
matches := re.FindAllStringSubmatch(html, -1)
for _, m := range matches {
name := m[1]
expected := m[2]
text := m[3]
r.elems[name] = &fakeElement{exists: true, expected: expected, text: text}
}
}
Fields
| Name | Type | Description |
|---|---|---|
| elems | map[string]*fakeElement | |
| html | string |
newFakeRoot
Returns
func newFakeRoot() *fakeRoot
{
return &fakeRoot{elems: make(map[string]*fakeElement)}
}
TestHandleHostPayloadMismatchTriggersResync
Parameters
func TestHandleHostPayloadMismatchTriggersResync(t *testing.T)
{
root := newFakeRoot()
root.elems["greeting"] = &fakeElement{
exists: true,
expected: encodeExpectation("server"),
text: "tampered",
}
payload := map[string]any{"greeting": "fresh"}
mismatches := handleHostPayload(root, payload, nil)
if len(mismatches) != 1 {
t.Fatalf("expected 1 mismatch, got %d", len(mismatches))
}
if root.elems["greeting"].text != "tampered" {
t.Fatalf("text was updated despite mismatch")
}
resync := buildResyncPayload(mismatches)
body, ok := resync["resync"].(map[string]any)
if !ok {
t.Fatalf("resync payload missing body")
}
if body["reason"] != "host-var-mismatch" {
t.Fatalf("unexpected reason %v", body["reason"])
}
vars, ok := body["vars"].([]map[string]string)
if ok {
if vars[0]["var"] != "greeting" {
t.Fatalf("unexpected var name %s", vars[0]["var"])
}
if vars[0]["expected"] == vars[0]["actualHash"] {
t.Fatalf("expected hashes to differ on mismatch")
}
}
}
TestLegacyExpectationRequiresResync
Parameters
func TestLegacyExpectationRequiresResync(t *testing.T)
{
root := newFakeRoot()
root.elems["greeting"] = &fakeElement{
exists: true,
expected: "sha1:2b42fba6b3f0c7b0d352c30b63f055c1b2f507a2",
text: "hello",
}
if mismatches := handleHostPayload(root, map[string]any{"greeting": "updated"}, nil); len(mismatches) != 1 {
t.Fatalf("legacy expectation was trusted without verification: %+v", mismatches)
}
}
TestInitSnapshotRecoveryAndUpdate
Parameters
func TestInitSnapshotRecoveryAndUpdate(t *testing.T)
{
root := newFakeRoot()
root.elems["count"] = &fakeElement{
exists: true,
expected: encodeExpectation("1"),
text: "0",
}
if mismatches := handleHostPayload(root, map[string]any{"count": "2"}, nil); len(mismatches) == 0 {
t.Fatalf("expected mismatch when expectation diverges")
}
snapHTML := `<span data-host-var="count" data-host-expected="` + encodeExpectation("1") + `">1</span>`
applyInitSnapshot(root, &initSnapshotPayload{HTML: snapHTML})
if mismatches := handleHostPayload(root, map[string]any{"count": "3"}, nil); len(mismatches) != 0 {
t.Fatalf("expected clean hydration after snapshot")
}
elem := root.HostVar("count").(*fakeElement)
if elem.text != "3" {
t.Fatalf("expected text to update to 3, got %s", elem.text)
}
if elem.expected != encodeExpectation("3") {
t.Fatalf("expected hash to reflect new value")
}
}
domComponentRoot
type domComponentRoot struct
Methods
Parameters
Returns
func (domComponentRoot) HostVar(name string) hostVarElement
{
selector := fmt.Sprintf(`[%s="%s"]`, hostVarAttr, name)
return domHostVarElement{r.Query(selector)}
}
Parameters
func (domComponentRoot) SetHTML(html string)
{
r.Element.SetHTML(html)
}
newComponentRoot
Parameters
Returns
func newComponentRoot(el dom.Element) componentRoot
{
return domComponentRoot{el}
}
domHostVarElement
type domHostVarElement struct
Methods
Parameters
func (domHostVarElement) SetText(value string)
{ e.Element.SetText(value) }
Parameters
Returns
func (domHostVarElement) Attr(name string) string
{ return e.Element.Attr(name) }
Parameters
func (domHostVarElement) SetAttr(name, value string)
{ e.Element.SetAttr(name, value) }
componentBinding
type componentBinding struct
Fields
| Name | Type | Description |
|---|---|---|
| id | string | |
| vars | []string | |
| gate | *deliveryGate |
deliveryGate
deliveryGate is the revocation switch a registration hands to every frame
delivered under it. A release closes it inside the same section that drops
the registration and then waits out the writes it may already have allowed,
so a frame the read loop snapshotted before a cleanup cannot update a host
signal once that cleanup returned. Reading it costs one atomic load and takes
no lock, so a signal setter that re-enters registration or release cannot
deadlock against a delivery holding it.
type deliveryGate struct
Methods
open reports whether the registration this gate belongs to is still the live one. The zero binding carries no gate and owns nothing.
Returns
func (*deliveryGate) open() bool
{ return g != nil && !g.closed.Load() }
func (*deliveryGate) close()
{
if g != nil {
g.closed.Store(true)
}
}
Fields
| Name | Type | Description |
|---|---|---|
| closed | atomic.Bool |
gatedHostSetter
gatedHostSetter is a host signal that can fold the caller’s ownership check
into its own store, so the two are one step against a concurrent release.
type gatedHostSetter interface
Methods
Parameters
Returns
func SetFromHostGated(...)
hostSetter
hostSetter is the plain host signal contract, without the gate.
type hostSetter interface
Methods
hostWriteBarrier
hostWriteBarrier waits out a gated write that was already allowed.
type hostWriteBarrier interface
Methods
func HostWriteBarrier(...)
message
type message struct
Fields
| Name | Type | Description |
|---|---|---|
| name | string | |
| action | string | |
| id | string | |
| payload | any | |
| sequence | uint64 |
wireMessage
type wireMessage struct
Fields
| Name | Type | Description |
|---|---|---|
| Component | string | json:"component,omitempty" |
| Action | string | json:"action,omitempty" |
| Control | string | json:"control,omitempty" |
| ID | string | json:"id,omitempty" |
| Payload | any | json:"payload,omitempty" |
| Sequence | uint64 | json:"sequence" |
| Ack | uint64 | json:"ack,omitempty" |
| ResumeToken | string | json:"resumeToken,omitempty" |
messageWriter
type messageWriter func(context.Context, *hostConn, wireMessage) error
actionReply
type actionReply struct
Fields
| Name | Type | Description |
|---|---|---|
| payload | any | |
| err | *ActionError | |
| resetErr | error |
decodeInitSnapshotPayload
Parameters
Returns
func decodeInitSnapshotPayload(raw any) *initSnapshotPayload
{
if raw == nil {
return nil
}
m, ok := raw.(map[string]any)
if !ok {
return nil
}
html, _ := m["html"].(string)
if html == "" {
return nil
}
var vars []string
if list, ok := m["vars"].([]any); ok {
vars = make([]string, 0, len(list))
for _, item := range list {
if s, ok := item.(string); ok {
vars = append(vars, s)
}
}
} else if list, ok := m["vars"].([]string); ok {
vars = append(vars, list...)
}
return &initSnapshotPayload{HTML: html, Vars: vars}
}
ActionError
ActionError is a machine-readable error returned by a typed host action.
type ActionError struct
Methods
Returns
func (*ActionError) Error() string
{
if e == nil {
return ""
}
return e.Code + ": " + e.Message
}
Fields
| Name | Type | Description |
|---|---|---|
| Code | string | json:"code" |
| Message | string | json:"message" |
| Fields | map[string]string | json:"fields,omitempty" |
init
func init()
{
cb = fnres.NewCircuitBreaker(5, 30*time.Second)
cb.OnStateChange(func(from, to fnres.State) {
if debug {
log.Printf("hostclient: circuit %v -> %v", from, to)
}
})
hydrateCB = fnres.NewCircuitBreaker(3, 15*time.Second)
sendCache = fncaching.NewInMemory[string](
fncaching.WithMaxEntries[string](256),
fncaching.WithTTL[string](5*time.Second),
)
}
connect
func connect()
{
once.Do(func() {
go func() {
for {
js.Guard("host connection loop", connectionLoop)
time.Sleep(time.Second)
}
}()
})
}
hostWSURL
hostWSURL builds the WebSocket URL the client uses to reach its host.
The endpoint is resolved in order of precedence: a full URL in
window.RFW_HOST_URL (ws, wss, http, https, or a bare host[:port] with an
optional path), the legacy host[:port] in window.RFW_HOST, or the page
origin. The path defaults to /ws when the endpoint carries none.
Returns
func hostWSURL() string
{
if u := js.Get("RFW_HOST_URL"); u.Truthy() {
if s := normalizeWSURL(u.String()); s != "" {
return s
}
}
host := js.Location().Get("host").String()
if h := js.Get("RFW_HOST"); h.Truthy() {
host = h.String()
}
return normalizeWSURL(host)
}
normalizeWSURL
normalizeWSURL turns a configured endpoint into a WebSocket URL: http and
https map to ws and wss, a bare host takes the page scheme, and /ws is
appended when the endpoint carries no path.
Parameters
Returns
func normalizeWSURL(raw string) string
{
raw = strings.TrimSpace(raw)
if raw == "" {
return ""
}
switch {
case strings.HasPrefix(raw, "ws://"), strings.HasPrefix(raw, "wss://"):
case strings.HasPrefix(raw, "http://"):
raw = "ws://" + strings.TrimPrefix(raw, "http://")
case strings.HasPrefix(raw, "https://"):
raw = "wss://" + strings.TrimPrefix(raw, "https://")
default:
scheme := "wss"
if js.Location().Get("protocol").String() == "http:" {
scheme = "ws"
}
raw = scheme + "://" + raw
}
// An endpoint that carries no path gets the default one. A bare trailing
// slash is no path either: "https://host/" would otherwise dial the root,
// which the host does not serve.
rest := raw[strings.Index(raw, "://")+3:]
if slash := strings.Index(rest, "/"); slash == -1 || strings.Trim(rest[slash:], "/") == "" {
raw = strings.TrimRight(raw, "/") + "/ws"
}
return raw
}
connectionLoop
func connectionLoop()
{
for {
url := hostWSURL()
connectionState.Set(ConnectionConnecting)
err := fnres.Retry(context.Background(), func() error {
return cb.Execute(func() error {
generation := connectionGeneration.Load()
if debug {
log.Printf("hostclient: dialing %s", url)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
c, derr := dial(ctx, url)
if derr != nil {
return derr
}
c.generation = generation
lifecycleMu.RLock()
if generation != connectionGeneration.Load() {
lifecycleMu.RUnlock()
_ = c.close()
return nil
}
sendMu.Lock()
mu.Lock()
conn = c
pend := pending
pending = nil
mu.Unlock()
if debug {
log.Printf("hostclient: connected")
}
connectionState.Set(ConnectionConnected)
mu.RLock()
names := make([]string, 0, len(bindings)+len(handlers))
for name := range bindings {
names = append(names, name)
}
for name := range handlers {
if _, bound := bindings[name]; !bound {
names = append(names, name)
}
}
mu.RUnlock()
deliveryMu.Lock()
unacknowledged := make([]message, 0, len(outbox))
for sequence := uint64(1); sequence <= nextOutbound; sequence++ {
if msg, ok := outbox[sequence]; ok {
unacknowledged = append(unacknowledged, msg)
}
}
deliveryMu.Unlock()
initialized := make(map[string]struct{})
for _, msg := range unacknowledged {
sendMessageUnlocked(c, msg)
if name, ok := initMessageName(msg); ok {
initialized[name] = struct{}{}
}
}
for _, msg := range pend {
sendMessageUnlocked(c, msg)
if name, ok := initMessageName(msg); ok {
initialized[name] = struct{}{}
}
}
for _, name := range names {
if _, sent := initialized[name]; sent {
continue
}
sendMessageUnlocked(c, message{name: name, payload: map[string]any{"init": true}})
}
sendMu.Unlock()
lifecycleMu.RUnlock()
ctx2, cancel2 := context.WithCancel(context.Background())
defer cancel2()
errCh := make(chan error, 2)
go func() { errCh <- guardedLoop("host read loop", func() error { return readLoop(ctx2, c) }) }()
go func() { errCh <- guardedLoop("host heartbeat loop", func() error { return heartbeatLoop(ctx2, c) }) }()
loopErr := <-errCh
cancel2()
closeErr := c.close()
mu.Lock()
if conn == c {
conn = nil
}
mu.Unlock()
connectionState.Set(ConnectionDisconnected)
if generation != connectionGeneration.Load() {
return nil
}
if loopErr != nil {
return loopErr
}
return closeErr
})
},
fnres.WithAttempts(5),
fnres.WithDelay(time.Second, 30*time.Second),
fnres.WithFactor(2),
fnres.WithJitter(0.1),
fnres.WithRetryIf(func(err error) bool { return err != nil }),
)
if err != nil && debug {
log.Printf("hostclient: connection attempt failed: %v", err)
}
connectionState.Set(ConnectionDisconnected)
// Back off before reconnecting to avoid tight loops on persistent failures.
time.Sleep(time.Second)
}
}
guardedLoop
Parameters
Returns
func guardedLoop(context string, fn func() error) error
{
var err error
if !js.Guard(context, func() { err = fn() }) {
return fmt.Errorf("%s panicked", context)
}
return err
}
heartbeatLoop
Parameters
Returns
func heartbeatLoop(ctx context.Context, c *hostConn) error
{
return c.heartbeat(ctx, heartbeatInterval, heartbeatTimeout)
}
readLoop
Parameters
Returns
func readLoop(ctx context.Context, c *hostConn) error
{
for {
var msg struct {
Component string `json:"component"`
Action string `json:"action"`
Control string `json:"control"`
ID string `json:"id"`
Payload any `json:"payload"`
Error *ActionError `json:"error"`
Session string `json:"session"`
Sequence uint64 `json:"sequence"`
Ack uint64 `json:"ack"`
ResumeToken string `json:"resumeToken"`
}
frame, err := c.read(ctx)
if err != nil {
return err
}
if c.generation != connectionGeneration.Load() {
return ErrSessionReset
}
if err := json.Unmarshal(frame, &msg); err != nil {
return err
}
if debug {
log.Printf("hostclient: recv %s %v", msg.Component, msg.Payload)
}
prepareInboundDelivery(msg.Session, msg.Control)
deliveryMu.Lock()
for sequence := range outbox {
if sequence <= msg.Ack {
delete(outbox, sequence)
}
}
if msg.Sequence != 0 {
if msg.Sequence <= lastInbound {
deliveryMu.Unlock()
continue
}
if lastInbound != 0 && msg.Sequence != lastInbound+1 {
deliveryMu.Unlock()
connectionState.Set(ConnectionDesynced)
return errors.New("hostclient: server message sequence gap")
}
lastInbound = msg.Sequence
}
if msg.ResumeToken != "" {
resumeToken = msg.ResumeToken
}
deliveryMu.Unlock()
if msg.ID != "" {
callMu.Lock()
replyChannel := pendingCalls[msg.ID]
if replyChannel != nil {
delete(pendingCalls, msg.ID)
}
callMu.Unlock()
if replyChannel != nil {
replyChannel <- actionReply{payload: msg.Payload, err: msg.Error}
continue
}
}
if msg.Control != "" {
continue
}
payload, _ := msg.Payload.(map[string]any)
if payload == nil {
payload = make(map[string]any)
}
mu.RLock()
h, hasHandler := handlers[msg.Component]
_, hasBinding := bindings[msg.Component]
token := bindingTokens[msg.Component]
barrier := afterBindingSnapshot
mu.RUnlock()
if hasHandler {
if msg.Session != "" {
payload["_session"] = msg.Session
}
js.Guard("host handler: "+msg.Component, func() { h(payload) })
continue
}
if hasBinding {
if barrier != nil {
barrier(msg.Component)
}
js.Guard("host binding: "+msg.Component, func() {
applyHostBinding(msg.Component, payload, token)
})
}
}
}
hostSignalUpdate
hostSignalUpdate is one host variable a delivered frame carries, held until
the DOM work is done and its signal can be set outside bindingMu.
type hostSignalUpdate struct
Fields
| Name | Type | Description |
|---|---|---|
| name | string | |
| value | any |
applyHostBinding
applyHostBinding delivers one host frame to the binding that owns it: the DOM
of the root under bindingMu, then the signals and any resync with the lock
released, since both run code that reacquires it.
Parameters
func applyHostBinding(component string, payload map[string]any, token uint64)
{
updates, mismatches := deliverHostFrame(component, payload, token)
applyHostSignals(component, token, updates)
if len(mismatches) == 0 {
return
}
for _, mismatch := range mismatches {
log.Printf("hostclient: hydration mismatch component=%s var=%s expected=%s actualHash=%s actual=%q", component, mismatch.VarName, mismatch.Expected, mismatch.ActualHash, mismatch.Actual)
}
resyncErr := hydrateCB.Execute(func() error {
Send(component, buildResyncPayload(mismatches))
return nil
})
if resyncErr != nil {
log.Printf("hostclient: hydration circuit open, skipping resync for %s", component)
}
}
deliverHostFrame
deliverHostFrame updates the DOM of the root the binding owns and reports the
signal updates the frame carries. It runs under bindingMu and re-reads the
registration token there: a frame the read loop snapshotted before a cleanup
blocks until the release completes and then finds the token gone, so it
reaches neither the released root nor the replacement registered under the
same name. mu is taken only for the token read and the snapshot bookkeeping.
Parameters
Returns
func deliverHostFrame(component string, payload map[string]any, token uint64) ([]hostSignalUpdate, []hydrationMismatch)
{
bindingMu.Lock()
defer bindingMu.Unlock()
binding, live := liveBinding(component, token)
if !live {
return nil, nil
}
rootEl := hostComponentRoot(binding.id)
if !rootEl.Truthy() {
return nil, nil
}
root := newComponentRoot(rootEl)
if snap := decodeInitSnapshotPayload(payload["initSnapshot"]); snap != nil {
applyInitSnapshot(root, snap)
if len(snap.Vars) > 0 {
binding.vars = append([]string(nil), snap.Vars...)
mu.Lock()
// Only the registration this frame was validated against may be
// updated: a snapshot must not reinstate a released binding nor
// overwrite a newer one.
if bindingTokens[component] == token {
bindings[component] = binding
}
mu.Unlock()
}
return nil, nil
}
var updates []hostSignalUpdate
mismatches := handleHostPayload(root, payload, func(name string, raw any) {
updates = append(updates, hostSignalUpdate{name: name, value: raw})
})
return updates, mismatches
}
applyHostSignals
applyHostSignals pushes a delivered frame into the component’s host signals.
A setter runs application code, which may register or release a host binding
in turn, so it must not run under bindingMu. The registration’s gate carries
the ownership check into the write instead: the signal evaluates it under the
same lock it stores the value with, so a release either refused the write or
waited it out before returning, and neither side holds a binding lock while
application code runs.
Parameters
func applyHostSignals(component string, token uint64, updates []hostSignalUpdate)
{
if len(updates) == 0 {
return
}
binding, live := liveBinding(component, token)
if !live {
return
}
signals := dom.SnapshotComponentSignals(binding.id)
if len(signals) == 0 {
return
}
hook := hostSignalWriteHook()
for _, update := range updates {
signal, ok := signals[update.name]
if !ok {
continue
}
if hook != nil {
hook(component, update.name)
}
if !applyHostSignal(signal, binding.gate, update.value) {
// The registration was released mid-frame: the rest of the frame
// was addressed to it too.
return
}
}
}
applyHostSignal
applyHostSignal writes one update and reports whether the registration was
still live for it. A signal that predates the gated setter cannot fold the
check into its store, so it is checked before the write instead, which leaves
the window the gate exists to close.
Parameters
Returns
func applyHostSignal(signal any, gate *deliveryGate, value any) bool
{
if gated, ok := signal.(gatedHostSetter); ok {
if gated.SetFromHostGated(value, gate.open) {
return true
}
// A refused write is either a revoked registration or a payload the
// signal cannot represent; only the first one ends the frame.
return gate.open()
}
setter, ok := signal.(hostSetter)
if !ok {
return true
}
if !gate.open() {
return false
}
setter.SetFromHost(value)
return true
}
fenceHostSignalWrites
fenceHostSignalWrites waits out the host signal writes a now closed gate had
already allowed. Each barrier takes only that signal’s own value lock, which
never covers application code, so a release re-entered from a setter never
waits on itself. It covers the signals the component still has registered,
which is why core unmounts a component by releasing its host bindings before
it drops its signals.
Parameters
func fenceHostSignalWrites(id string)
{
if id == "" {
return
}
for _, signal := range dom.SnapshotComponentSignals(id) {
if barrier, ok := signal.(hostWriteBarrier); ok {
barrier.HostWriteBarrier()
}
}
}
hostSignalWriteHook
Returns
func hostSignalWriteHook() func(component, name string)
{
mu.RLock()
defer mu.RUnlock()
return beforeHostSignalWrite
}
liveBinding
liveBinding returns the binding still registered under token. A token that no
longer matches means the registration was released or replaced, so the caller
owns nothing to update.
Parameters
Returns
func liveBinding(component string, token uint64) (componentBinding, bool)
{
mu.RLock()
defer mu.RUnlock()
if bindingTokens[component] != token {
return componentBinding{}, false
}
binding, bound := bindings[component]
return binding, bound
}
hostComponentRoot
hostComponentRoot resolves the exact root a binding owns. dom.ComponentRoot
falls back to #app when the id matches nothing, which for host delivery would
mean a frame addressed to an unmounted component writing into the application
shell or into whatever mounted after it. A missing root is a frame to ignore.
RegisterComponent is public and takes any id, so the id never reaches a
selector unescaped: a quote or a bracket in it would otherwise throw a
DOMException out of querySelector, or select a root the caller never owned.
Parameters
Returns
func hostComponentRoot(id string) dom.Element
{
if id == "" {
return dom.Element{Value: js.Null()}
}
if selector, ok := componentIDSelector(id); ok {
return dom.Doc().Query(selector)
}
return scanComponentRoot(id)
}
componentIDSelector
componentIDSelector builds the attribute selector for id and leaves the
escaping rules to the browser. CSS.escape returns an identifier, which is
what the value position of an attribute selector accepts.
Parameters
Returns
func componentIDSelector(id string) (string, bool)
{
css := js.Get("CSS")
if !css.Truthy() || css.Get("escape").Type() != js.TypeFunction {
return "", false
}
return "[data-component-id=" + css.Call("escape", id).String() + "]", true
}
scanComponentRoot
scanComponentRoot is the fallback for a runtime without CSS.escape: the
candidates are selected on the bare attribute and compared on its value, so
no part of the id is ever parsed as a selector.
Parameters
Returns
func scanComponentRoot(id string) dom.Element
{
roots := dom.Doc().QueryAll("[data-component-id]")
if !roots.Truthy() {
return dom.Element{Value: js.Null()}
}
for index := 0; index < roots.Length(); index++ {
candidate := roots.Index(index)
if candidate.Attr("data-component-id") == id {
return candidate
}
}
return dom.Element{Value: js.Null()}
}
prepareInboundDelivery
Parameters
func prepareInboundDelivery(remoteSession, control string)
{
sessionMu.Lock()
previousSession := sessionID
if remoteSession != "" {
sessionID = remoteSession
}
sessionMu.Unlock()
if control != "resume_rejected" && (remoteSession == "" || previousSession == "" || remoteSession == previousSession) {
return
}
deliveryMu.Lock()
lastInbound = 0
resumeToken = ""
deliveryMu.Unlock()
}
RegisterComponent
RegisterComponent binds a client component to a host component name. The
binding lives until another registration replaces it. Use
RegisterComponentOwned when the caller has to release the binding again, for
instance because its component can unmount.
Parameters
func RegisterComponent(id, name string, vars []string)
{
registerComponent(id, name, vars)
}
RegisterComponentOwned
RegisterComponentOwned binds like RegisterComponent and returns an idempotent
cleanup that owns the binding’s lifecycle. Releasing removes the binding from
reconnect hydration and tells the active host session to stop broadcasts for
the component; once the cleanup returns, no frame still in flight and no
later frame can update the root it owned or write into its host signals. A
stale cleanup closure never removes a newer binding registered under the same
name. The cleanup is safe to call from a host signal setter, which is where a
component that unmounts on an update ends up calling it.
Parameters
Returns
func RegisterComponentOwned(id, name string, vars []string) func()
{
token := registerComponent(id, name, vars)
var once sync.Once
return func() {
once.Do(func() { releaseComponent(name, token) })
}
}
registerComponent
registerComponent installs the binding and returns the token identifying this
registration. Delivery validates against it, so registering also invalidates
the cleanup of the binding it replaced.
Parameters
Returns
func registerComponent(id, name string, vars []string) uint64
{
token := bindingSequence.Add(1)
bindingMu.Lock()
mu.Lock()
previous := bindings[name]
bindings[name] = componentBinding{id: id, vars: vars, gate: &deliveryGate{}}
bindingTokens[name] = token
current := conn
recordPendingControl(name, map[string]any{"init": true}, current)
mu.Unlock()
// The registration this one replaces stops being deliverable here. Its
// frames already fail the token check; closing its gate stops one that was
// snapshotted before the replacement from writing a signal on its way out.
previous.gate.close()
bindingMu.Unlock()
connect()
if current != nil {
sendMessage(current, message{name: name, payload: map[string]any{"init": true}})
}
return token
}
releaseComponent
Parameters
func releaseComponent(name string, token uint64)
{
bindingMu.Lock()
mu.Lock()
if bindingTokens[name] != token {
mu.Unlock()
bindingMu.Unlock()
return
}
binding := bindings[name]
delete(bindings, name)
delete(bindingTokens, name)
current := conn
recordPendingControl(name, map[string]any{"unsubscribe": true}, current)
mu.Unlock()
// Closing the gate in the section that dropped the registration means no
// delivery can read it open again afterwards.
binding.gate.close()
// Every delivery for this binding is now either finished or bound to fail
// its token check, so the unsubscribe can go out without the lock.
bindingMu.Unlock()
// The root is safe by now, the signals are not: a frame that read the gate
// open before it closed may be committing a write. Waiting that write out,
// with no binding lock held, is what makes the cleanup final.
fenceHostSignalWrites(binding.id)
if current != nil {
sendMessage(current, message{name: name, payload: map[string]any{"unsubscribe": true}})
}
}
recordPendingControl
recordPendingControl keeps the reconnect queue holding one registration
control per host component, the latest desired state: a registration
supersedes a queued unsubscribe and a release supersedes a queued init, so
route churn while offline cannot retain one message per cycle. The control
the caller just decided on is queued only when there is no connection to send
it on; a superseded one is dropped either way, since sending the new state
makes replaying the old one wrong.
The new state takes the queue position of the control it supersedes rather
than the tail: everything else in the queue, the controls of other names and
this name’s own messages, keeps the order it was queued in, so a reconnect
replays the component’s init before the command that assumed it.
It must be called with mu held.
Parameters
func recordPendingControl(name string, payload map[string]any, current *hostConn)
{
replaced := false
filtered := pending[:0]
for _, queued := range pending {
if queued.name != name || !isRegistrationControl(queued.payload) {
filtered = append(filtered, queued)
continue
}
if current != nil || replaced {
continue
}
queued.payload = payload
filtered = append(filtered, queued)
replaced = true
}
pending = filtered
if current == nil && !replaced {
pending = append(pending, message{name: name, payload: payload})
}
}
isRegistrationControl
Parameters
Returns
func isRegistrationControl(payload any) bool
{
values, ok := payload.(map[string]any)
if !ok {
return false
}
return values["init"] == true || values["unsubscribe"] == true
}
EnableSendDedup
EnableSendDedup turns on payload-based deduplication for the named channel:
identical payloads sent within a 5 second window are dropped. Dedup is off
by default because repeated identical messages are usually intentional user
actions (e.g. clicking +1 twice); opt in only for channels where duplicate
suppression is the desired semantic.
Parameters
func EnableSendDedup(name string)
{
mu.Lock()
dedup[name] = struct{}{}
mu.Unlock()
}
dedupEnabled
Parameters
Returns
func dedupEnabled(name string) bool
{
mu.RLock()
_, ok := dedup[name]
mu.RUnlock()
return ok
}
Send
Send queues or transmits a host component message.
Parameters
func Send(name string, payload any)
{
lifecycleMu.RLock()
defer lifecycleMu.RUnlock()
connect()
if dedupEnabled(name) {
key := fmt.Sprintf("%d|%s|%v", connectionGeneration.Load(), name, payload)
if _, ok, _ := sendCache.Get(context.Background(), key); ok {
return
}
if err := sendCache.Set(context.Background(), key, "sent", 5*time.Second); err != nil {
log.Printf("hostclient: dedup cache set failed: %v", err)
}
}
mu.RLock()
c := conn
mu.RUnlock()
if c == nil {
mu.Lock()
pending = append(pending, message{name: name, payload: payload})
mu.Unlock()
return
}
if debug {
log.Printf("hostclient: send %s %v", name, payload)
}
sendMessage(c, message{name: name, payload: payload})
}
RegisterHandler
RegisterHandler registers a handler for host messages and returns an
idempotent unsubscribe function. Unsubscribing removes the handler from
reconnect hydration and tells the active host session to stop broadcasts for
the component. A stale unsubscribe closure never removes a newer handler
registered under the same name.
Parameters
Returns
func RegisterHandler(name string, h func(map[string]any)) func()
{
token := handlerSequence.Add(1)
mu.Lock()
handlers[name] = h
handlerTokens[name] = token
current := conn
recordPendingControl(name, map[string]any{"init": true}, current)
mu.Unlock()
connect()
if current != nil {
sendMessage(current, message{name: name, payload: map[string]any{"init": true}})
}
var once sync.Once
return func() {
once.Do(func() {
mu.Lock()
if handlerTokens[name] != token {
mu.Unlock()
return
}
delete(handlers, name)
delete(handlerTokens, name)
current := conn
recordPendingControl(name, map[string]any{"unsubscribe": true}, current)
mu.Unlock()
if current != nil {
sendMessage(current, message{name: name, payload: map[string]any{"unsubscribe": true}})
}
})
}
}
SessionID
SessionID returns the current SSC session ID.
Returns
func SessionID() string
{
sessionMu.RLock()
defer sessionMu.RUnlock()
return sessionID
}
ResetSession
ResetSession closes the active SSC transport and discards all delivery
state owned by its authenticated session. Live registrations are preserved;
callers must unmount user-owned scopes before invoking it.
func ResetSession()
{
lifecycleMu.Lock()
defer lifecycleMu.Unlock()
connectionGeneration.Add(1)
connectionState.Set(ConnectionDisconnected)
mu.Lock()
current := conn
conn = nil
pending = nil
mu.Unlock()
sessionMu.Lock()
sessionID = ""
sessionMu.Unlock()
deliveryMu.Lock()
resumeToken = ""
nextOutbound = 0
lastInbound = 0
outbox = map[uint64]message{}
deliveryMu.Unlock()
callMu.Lock()
calls := pendingCalls
pendingCalls = map[string]chan actionReply{}
callMu.Unlock()
for _, reply := range calls {
reply <- actionReply{resetErr: ErrSessionReset}
}
if current != nil {
_ = current.close()
}
}
sendMessage
func sendMessage(c *hostConn, msg message)
{
sendMessageWithWriter(c, msg, writeMessage)
}
Uses
sendMessageWithWriter
Parameters
func sendMessageWithWriter(c *hostConn, msg message, writer messageWriter)
{
sendMu.Lock()
defer sendMu.Unlock()
sendMessageUnlockedWithWriter(c, msg, writer)
}
sendMessageUnlocked
func sendMessageUnlocked(c *hostConn, msg message)
{
sendMessageUnlockedWithWriter(c, msg, writeMessage)
}
Uses
sendMessageUnlockedWithWriter
Parameters
func sendMessageUnlockedWithWriter(c *hostConn, msg message, writer messageWriter)
{
deliveryMu.Lock()
if msg.sequence == 0 {
nextOutbound++
msg.sequence = nextOutbound
outbox[msg.sequence] = msg
}
token := resumeToken
ack := lastInbound
deliveryMu.Unlock()
outbound := wireMessage{
Component: msg.name,
Action: msg.action,
ID: msg.id,
Payload: msg.payload,
Sequence: msg.sequence,
Ack: ack,
ResumeToken: token,
}
ctx := context.Background()
_ = writer(ctx, c, outbound)
}
writeMessage
Parameters
Returns
func writeMessage(_ context.Context, c *hostConn, message wireMessage) error
{
return c.writeJSON(message)
}
initMessageName
Parameters
Returns
func initMessageName(msg message) (string, bool)
{
if msg.name == "" || msg.action != "" {
return "", false
}
payload, ok := msg.payload.(map[string]any)
if !ok || payload["init"] != true {
return "", false
}
return msg.name, true
}
Uses
Call
Call invokes a typed SSC action and waits for its correlated response.
Parameters
Returns
func Call[Request, Response any](ctx context.Context, action string, request Request) (Response, error)
{
var zero Response
if ctx == nil {
ctx = context.Background()
}
if action == "" {
return zero, errors.New("hostclient: empty action name")
}
lifecycleMu.RLock()
connect()
id := fmt.Sprintf("call-%d", callSequence.Add(1))
replyChannel := make(chan actionReply, 1)
callMu.Lock()
pendingCalls[id] = replyChannel
callMu.Unlock()
msg := message{action: action, id: id, payload: request}
mu.RLock()
current := conn
mu.RUnlock()
if current == nil {
mu.Lock()
pending = append(pending, msg)
mu.Unlock()
} else {
sendMessage(current, msg)
}
lifecycleMu.RUnlock()
select {
case reply := <-replyChannel:
if reply.resetErr != nil {
return zero, reply.resetErr
}
if reply.err != nil {
return zero, reply.err
}
data, err := json.Marshal(reply.payload)
if err != nil {
return zero, fmt.Errorf("hostclient: encode action response: %w", err)
}
if err := json.Unmarshal(data, &zero); err != nil {
return zero, fmt.Errorf("hostclient: decode action response: %w", err)
}
return zero, nil
case <-ctx.Done():
callMu.Lock()
delete(pendingCalls, id)
callMu.Unlock()
return zero, ctx.Err()
}
}
FormResponse
FormResponse is the typed result returned by host.RegisterForm.
type FormResponse struct
Fields
| Name | Type | Description |
|---|---|---|
| Data | Response | json:"data,omitempty" |
| Fields | map[string]string | json:"fields,omitempty" |
| Valid | bool | json:"valid" |
SubmitForm
SubmitForm invokes a typed SSC form action.
Parameters
Returns
func SubmitForm[Values, Response any](ctx context.Context, action string, values Values) (FormResponse[Response], error)
{
return Call[Values, FormResponse[Response]](ctx, action, values)
}
EnableDebug
EnableDebug enables host client debug logging.
func EnableDebug()
{ debug = true }
uniqueRuntimeTestName
Parameters
Returns
func uniqueRuntimeTestName(prefix string) string
{
return fmt.Sprintf("%s-%d", prefix, runtimeTestSequence.Add(1))
}
TestGuardedLoopConvertsPanicAndNextLoopRuns
Parameters
func TestGuardedLoopConvertsPanicAndNextLoopRuns(t *testing.T)
{
previous := js.OnRuntimePanic
defer func() { js.OnRuntimePanic = previous }()
recovered := 0
js.OnRuntimePanic = func(any, string, []byte) { recovered++ }
if err := guardedLoop("read", func() error { panic("bad push") }); err == nil {
t.Fatal("panicking loop returned nil")
}
if err := guardedLoop("read", func() error { return nil }); err != nil {
t.Fatalf("next loop returned %v", err)
}
if recovered != 1 {
t.Fatalf("recovered loop panics = %d, want 1", recovered)
}
}
TestInboundMessageLimitSupportsHydrationSnapshots
Parameters
func TestInboundMessageLimitSupportsHydrationSnapshots(t *testing.T)
{
if maxInboundMessageBytes < 1<<20 {
t.Fatalf("inbound message limit = %d, want at least 1 MiB", maxInboundMessageBytes)
}
if maxInboundMessageBytes > 16<<20 {
t.Fatalf("inbound message limit = %d, want a bounded ceiling", maxInboundMessageBytes)
}
}
pendingCount
Returns
func pendingCount() int
{
mu.RLock()
defer mu.RUnlock()
return len(pending)
}
TestRegisterHandlerUnsubscribeQueuesWireUnsubscribe
Parameters
func TestRegisterHandlerUnsubscribeQueuesWireUnsubscribe(t *testing.T)
{
name := uniqueRuntimeTestName("scoped-handler")
before := pendingCount()
unsubscribe := RegisterHandler(name, func(map[string]any) {})
if got := pendingCount() - before; got != 1 {
t.Fatalf("queued subscribe messages = %d, want 1", got)
}
unsubscribe()
mu.RLock()
_, stillRegistered := handlers[name]
queued := append([]message(nil), pending...)
mu.RUnlock()
if stillRegistered {
t.Fatal("handler remained registered after unsubscribe")
}
count := 0
for _, item := range queued {
if item.name == name {
values, _ := item.payload.(map[string]any)
if values["unsubscribe"] == true {
count++
}
if values["init"] == true {
t.Fatal("stale init remained queued after unsubscribe")
}
}
}
if count != 1 {
t.Fatalf("queued unsubscribe messages = %d, want 1", count)
}
}
TestStaleUnsubscribeDoesNotRemoveReplacementHandler
Parameters
func TestStaleUnsubscribeDoesNotRemoveReplacementHandler(t *testing.T)
{
name := uniqueRuntimeTestName("replacement-handler")
first := RegisterHandler(name, func(map[string]any) {})
second := RegisterHandler(name, func(map[string]any) {})
first()
mu.RLock()
_, registered := handlers[name]
mu.RUnlock()
if !registered {
t.Fatal("stale unsubscribe removed replacement handler")
}
second()
}
TestSendRepeatedMessagesNotDeduped
Repeated identical messages must go through by default: two identical user
actions within the dedup window (e.g. clicking +1 twice) are intentional.
Parameters
func TestSendRepeatedMessagesNotDeduped(t *testing.T)
{
before := pendingCount()
Send("CounterHost", map[string]any{"cmd": "increment"})
Send("CounterHost", map[string]any{"cmd": "increment"})
if got := pendingCount() - before; got != 2 {
t.Fatalf("expected 2 queued messages, got %d", got)
}
}
TestSendDedupOptIn
Dedup is opt-in per channel: after EnableSendDedup identical payloads within
the TTL window are dropped.
Parameters
func TestSendDedupOptIn(t *testing.T)
{
name := uniqueRuntimeTestName("DedupHost")
EnableSendDedup(name)
before := pendingCount()
Send(name, map[string]any{"cmd": "refresh"})
Send(name, map[string]any{"cmd": "refresh"})
if got := pendingCount() - before; got != 1 {
t.Fatalf("expected 1 queued message after dedup, got %d", got)
}
}
TestSendMessageSerializesSequenceAndWrite
Parameters
func TestSendMessageSerializesSequenceAndWrite(t *testing.T)
{
deliveryMu.Lock()
savedNext := nextOutbound
savedOutbox := outbox
nextOutbound = 0
outbox = map[uint64]message{}
deliveryMu.Unlock()
defer func() {
deliveryMu.Lock()
nextOutbound = savedNext
outbox = savedOutbox
deliveryMu.Unlock()
}()
firstEntered := make(chan struct{}, 1)
firstRelease := make(chan struct{})
secondEntered := make(chan struct{}, 1)
firstDone := make(chan struct{})
secondDone := make(chan struct{})
firstWriter := func(_ context.Context, _ *hostConn, message wireMessage) error {
firstEntered <- struct{}{}
<-firstRelease
if message.Sequence != 1 {
t.Errorf("first sequence = %d, want 1", message.Sequence)
}
return nil
}
secondWriter := func(_ context.Context, _ *hostConn, message wireMessage) error {
secondEntered <- struct{}{}
if message.Sequence != 2 {
t.Errorf("second sequence = %d, want 2", message.Sequence)
}
return nil
}
go func() {
sendMessageWithWriter(nil, message{name: "first"}, firstWriter)
close(firstDone)
}()
<-firstEntered
go func() {
sendMessageWithWriter(nil, message{name: "second"}, secondWriter)
close(secondDone)
}()
select {
case <-secondEntered:
close(firstRelease)
<-firstDone
<-secondDone
t.Fatal("second message reached the writer before the first completed")
case <-time.After(50 * time.Millisecond):
}
close(firstRelease)
<-firstDone
<-secondDone
}
TestPrepareInboundDeliveryResetsNewSessionState
Parameters
func TestPrepareInboundDeliveryResetsNewSessionState(t *testing.T)
{
sessionMu.Lock()
savedSession := sessionID
sessionID = "old-session"
sessionMu.Unlock()
deliveryMu.Lock()
savedInbound := lastInbound
savedToken := resumeToken
lastInbound = 9
resumeToken = "old-token"
deliveryMu.Unlock()
defer func() {
sessionMu.Lock()
sessionID = savedSession
sessionMu.Unlock()
deliveryMu.Lock()
lastInbound = savedInbound
resumeToken = savedToken
deliveryMu.Unlock()
}()
prepareInboundDelivery("new-session", "")
sessionMu.RLock()
currentSession := sessionID
sessionMu.RUnlock()
deliveryMu.Lock()
currentInbound := lastInbound
currentToken := resumeToken
deliveryMu.Unlock()
if currentSession != "new-session" || currentInbound != 0 || currentToken != "" {
t.Fatalf("delivery state was not reset: session=%q inbound=%d token=%q", currentSession, currentInbound, currentToken)
}
}
TestPrepareInboundDeliveryResetsRejectedResume
Parameters
func TestPrepareInboundDeliveryResetsRejectedResume(t *testing.T)
{
sessionMu.Lock()
savedSession := sessionID
sessionID = "current-session"
sessionMu.Unlock()
deliveryMu.Lock()
savedInbound := lastInbound
savedToken := resumeToken
lastInbound = 9
resumeToken = "old-token"
deliveryMu.Unlock()
defer func() {
sessionMu.Lock()
sessionID = savedSession
sessionMu.Unlock()
deliveryMu.Lock()
lastInbound = savedInbound
resumeToken = savedToken
deliveryMu.Unlock()
}()
prepareInboundDelivery("current-session", "resume_rejected")
deliveryMu.Lock()
currentInbound := lastInbound
currentToken := resumeToken
deliveryMu.Unlock()
if currentInbound != 0 || currentToken != "" {
t.Fatalf("rejected resume state was not reset: inbound=%d token=%q", currentInbound, currentToken)
}
}
TestResetSessionClearsIdentityDeliveryAndPendingCalls
Parameters
func TestResetSessionClearsIdentityDeliveryAndPendingCalls(t *testing.T)
{
lifecycleMu.Lock()
savedGeneration := connectionGeneration.Load()
mu.Lock()
savedConn := conn
savedPending := pending
savedBinding := bindings["reset-live"]
_, hadBinding := bindings["reset-live"]
savedHandler := handlers["reset-live"]
_, hadHandler := handlers["reset-live"]
current := newHostConn()
closed := 0
current.shutdown = func() error { closed++; current.fail(errConnectionClosed); return nil }
current.generation = savedGeneration
conn = current
pending = []message{{name: "old-user", payload: map[string]any{"cmd": "stale"}}}
bindings["reset-live"] = componentBinding{id: "component-live"}
handlers["reset-live"] = func(map[string]any) {}
mu.Unlock()
sessionMu.Lock()
savedSession := sessionID
sessionID = "ssc-old"
sessionMu.Unlock()
deliveryMu.Lock()
savedToken := resumeToken
savedNext := nextOutbound
savedInbound := lastInbound
savedOutbox := outbox
resumeToken = "resume-old"
nextOutbound = 4
lastInbound = 7
outbox = map[uint64]message{4: {name: "old-user", sequence: 4}}
deliveryMu.Unlock()
callMu.Lock()
savedCalls := pendingCalls
reply := make(chan actionReply, 1)
pendingCalls = map[string]chan actionReply{"call-old": reply}
callMu.Unlock()
lifecycleMu.Unlock()
t.Cleanup(func() {
lifecycleMu.Lock()
connectionGeneration.Store(savedGeneration)
mu.Lock()
conn = savedConn
pending = savedPending
if hadBinding {
bindings["reset-live"] = savedBinding
} else {
delete(bindings, "reset-live")
}
if hadHandler {
handlers["reset-live"] = savedHandler
} else {
delete(handlers, "reset-live")
}
mu.Unlock()
sessionMu.Lock()
sessionID = savedSession
sessionMu.Unlock()
deliveryMu.Lock()
resumeToken = savedToken
nextOutbound = savedNext
lastInbound = savedInbound
outbox = savedOutbox
deliveryMu.Unlock()
callMu.Lock()
pendingCalls = savedCalls
callMu.Unlock()
lifecycleMu.Unlock()
})
ResetSession()
if closed != 1 {
t.Fatalf("transport closes = %d, want 1", closed)
}
if got := connectionGeneration.Load(); got != savedGeneration+1 {
t.Fatalf("generation = %d, want %d", got, savedGeneration+1)
}
mu.RLock()
currentConn := conn
queued := len(pending)
_, bindingPreserved := bindings["reset-live"]
_, handlerPreserved := handlers["reset-live"]
mu.RUnlock()
if currentConn != nil || queued != 0 {
t.Fatalf("connection state after reset: conn=%v pending=%d", currentConn, queued)
}
if !bindingPreserved || !handlerPreserved {
t.Fatal("reset removed a live binding or handler")
}
if got := SessionID(); got != "" {
t.Fatalf("session ID after reset = %q", got)
}
deliveryMu.Lock()
gotToken, gotNext, gotInbound, gotOutbox := resumeToken, nextOutbound, lastInbound, len(outbox)
deliveryMu.Unlock()
if gotToken != "" || gotNext != 0 || gotInbound != 0 || gotOutbox != 0 {
t.Fatalf("delivery after reset: token=%q next=%d inbound=%d outbox=%d", gotToken, gotNext, gotInbound, gotOutbox)
}
select {
case got := <-reply:
if !errors.Is(got.resetErr, ErrSessionReset) {
t.Fatalf("pending call error = %v, want ErrSessionReset", got.resetErr)
}
default:
t.Fatal("reset did not fail the pending call")
}
}
TestOldGenerationFrameIsRejectedAfterSessionReset
Parameters
func TestOldGenerationFrameIsRejectedAfterSessionReset(t *testing.T)
{
lifecycleMu.Lock()
savedGeneration := connectionGeneration.Load()
connectionGeneration.Store(savedGeneration + 1)
mu.Lock()
savedHandler := handlers["old-generation"]
_, hadHandler := handlers["old-generation"]
deliveries := 0
handlers["old-generation"] = func(map[string]any) { deliveries++ }
mu.Unlock()
lifecycleMu.Unlock()
t.Cleanup(func() {
lifecycleMu.Lock()
connectionGeneration.Store(savedGeneration)
mu.Lock()
if hadHandler {
handlers["old-generation"] = savedHandler
} else {
delete(handlers, "old-generation")
}
mu.Unlock()
lifecycleMu.Unlock()
})
stale := newHostConn()
stale.generation = savedGeneration
stale.deliver([]byte(`{"component":"old-generation","payload":{"value":"stale"},"sequence":1}`))
err := readLoop(context.Background(), stale)
if !errors.Is(err, ErrSessionReset) {
t.Fatalf("stale frame error = %v, want ErrSessionReset", err)
}
if deliveries != 0 {
t.Fatalf("stale frame deliveries = %d, want 0", deliveries)
}
}
bindingSnapshot
Parameters
Returns
func bindingSnapshot(name string) (componentBinding, bool)
{
mu.RLock()
defer mu.RUnlock()
binding, bound := bindings[name]
return binding, bound
}
pendingFor
Parameters
Returns
func pendingFor(name string) (inits, unsubscribes int)
{
mu.RLock()
defer mu.RUnlock()
for _, queued := range pending {
if queued.name != name {
continue
}
values, _ := queued.payload.(map[string]any)
if values["init"] == true {
inits++
}
if values["unsubscribe"] == true {
unsubscribes++
}
}
return inits, unsubscribes
}
pendingLabels
pendingLabels renders the queued messages for the given names in queue order,
so a test can assert where a control sits and not only how many there are.
Parameters
Returns
func pendingLabels(names ...string) []string
{
wanted := make(map[string]struct{}, len(names))
for _, name := range names {
wanted[name] = struct{}{}
}
mu.RLock()
defer mu.RUnlock()
labels := make([]string, 0, len(pending))
for _, queued := range pending {
if _, ok := wanted[queued.name]; !ok {
continue
}
labels = append(labels, queued.name+"/"+pendingKind(queued.payload))
}
return labels
}
pendingKind
Parameters
Returns
func pendingKind(payload any) string
{
values, ok := payload.(map[string]any)
if !ok {
return "message"
}
switch {
case values["init"] == true:
return "init"
case values["unsubscribe"] == true:
return "unsubscribe"
}
if cmd, ok := values["cmd"].(string); ok {
return "cmd:" + cmd
}
return "message"
}
forgetPending
forgetPending drops the messages a test queued, so package state does not
leak into the next one.
Parameters
func forgetPending(t *testing.T, name string)
{
t.Cleanup(func() {
mu.Lock()
filtered := pending[:0]
for _, queued := range pending {
if queued.name == name {
continue
}
filtered = append(filtered, queued)
}
pending = filtered
mu.Unlock()
})
}
setConnection
setConnection makes the package treat c as its host connection, the way the
connection loop does once a handshake completes. A nil c is the offline
state, where registration controls go to the reconnect queue.
func setConnection(t *testing.T, c *hostConn)
{
t.Helper()
mu.Lock()
conn = c
mu.Unlock()
t.Cleanup(func() {
mu.Lock()
conn = nil
mu.Unlock()
})
}
setSnapshotBarrier
setSnapshotBarrier installs the delivery barrier under the lock the read loop
reads it with.
Parameters
func setSnapshotBarrier(fn func(component string))
{
mu.Lock()
afterBindingSnapshot = fn
mu.Unlock()
}
TestRegisterComponentBindsWithoutACleanup
A registration made through the plain registrar owns no cleanup: it lives
until another registration replaces it, which is the behavior code calling
RegisterComponent as a statement has always had.
Parameters
func TestRegisterComponentBindsWithoutACleanup(t *testing.T)
{
name := "plain-binding"
forgetPending(t, name)
RegisterComponent("plain-root", name, []string{"value"})
binding, bound := bindingSnapshot(name)
if !bound || binding.id != "plain-root" {
t.Fatalf("registration bound %#v bound=%v", binding, bound)
}
release := RegisterComponentOwned("owned-root", name, nil)
if binding, _ := bindingSnapshot(name); binding.id != "owned-root" {
t.Fatalf("the owned registration did not replace the plain one: %#v", binding)
}
release()
if _, bound := bindingSnapshot(name); bound {
t.Fatal("binding remained after the owned cleanup")
}
}
TestRegisterComponentCleanupReleasesTheBinding
Parameters
func TestRegisterComponentCleanupReleasesTheBinding(t *testing.T)
{
name := "scoped-binding"
setConnection(t, nil)
forgetPending(t, name)
release := RegisterComponentOwned("scoped-root", name, []string{"value"})
if _, bound := bindingSnapshot(name); !bound {
t.Fatal("registration did not bind the component")
}
if inits, _ := pendingFor(name); inits != 1 {
t.Fatalf("queued init messages = %d, want 1", inits)
}
release()
release()
if _, bound := bindingSnapshot(name); bound {
t.Fatal("binding remained after cleanup")
}
inits, unsubscribes := pendingFor(name)
if inits != 0 {
t.Fatalf("stale init messages = %d, want 0", inits)
}
if unsubscribes != 1 {
t.Fatalf("queued unsubscribe messages = %d, want 1", unsubscribes)
}
}
TestStaleComponentCleanupKeepsTheReplacementBinding
Parameters
func TestStaleComponentCleanupKeepsTheReplacementBinding(t *testing.T)
{
name := "replacement-binding"
forgetPending(t, name)
first := RegisterComponentOwned("first-root", name, nil)
second := RegisterComponentOwned("second-root", name, nil)
first()
binding, bound := bindingSnapshot(name)
if !bound || binding.id != "second-root" {
t.Fatalf("stale cleanup removed the replacement binding: %#v bound=%v", binding, bound)
}
second()
if _, bound := bindingSnapshot(name); bound {
t.Fatal("the replacement binding survived its own cleanup")
}
}
TestActiveConnectionCleanupSendsOneUnsubscribe
With a live connection the lifecycle talks to the host directly: one init on
registration and one unsubscribe on cleanup, however often the cleanup runs,
and nothing left queued for a reconnect.
Parameters
func TestActiveConnectionCleanupSendsOneUnsubscribe(t *testing.T)
{
c, socket := dialFake(t)
setConnection(t, c)
name := "wired-ticker"
forgetPending(t, name)
release := RegisterComponentOwned("wired-root", name, nil)
release()
release()
inits, unsubscribes := 0, 0
for _, frame := range socket.sent() {
if !strings.Contains(frame, `"component":"`+name+`"`) {
continue
}
if strings.Contains(frame, `"init":true`) {
inits++
}
if strings.Contains(frame, `"unsubscribe":true`) {
unsubscribes++
}
}
if inits != 1 || unsubscribes != 1 {
t.Fatalf("frames sent: init=%d unsubscribe=%d, want 1 each", inits, unsubscribes)
}
if _, queued := pendingFor(name); queued != 0 {
t.Fatalf("a live connection queued %d unsubscribes", queued)
}
}
TestComponentRegistrationCyclesRetainNothing
Route entry and exit repeated: every mount owns one binding and one init, and
leaves nothing behind when it unmounts. Offline, the queue keeps only the
latest desired state for the name, so churn does not retain protocol work
proportional to the number of cycles.
Parameters
func TestComponentRegistrationCyclesRetainNothing(t *testing.T)
{
name := "cycled-binding"
setConnection(t, nil)
forgetPending(t, name)
for cycle := 0; cycle < 100; cycle++ {
release := RegisterComponentOwned("cycled-root", name, nil)
if _, bound := bindingSnapshot(name); !bound {
t.Fatalf("cycle %d: registration did not bind the component", cycle)
}
release()
if _, bound := bindingSnapshot(name); bound {
t.Fatalf("cycle %d: binding retained after cleanup", cycle)
}
}
mu.RLock()
_, tracked := bindingTokens[name]
mu.RUnlock()
if tracked {
t.Fatal("registration bookkeeping retained the released name")
}
inits, unsubscribes := pendingFor(name)
if inits != 0 || unsubscribes != 1 {
t.Fatalf("queued messages: init=%d unsubscribe=%d, want 0 and 1", inits, unsubscribes)
}
}
TestPendingRegistrationControlKeepsOnlyTheLatestState
The pending queue is a per-name desired state, not a log: a control is
superseded where it stands, so the messages that are not registration
controls, and the controls of other names, keep their exact place.
Parameters
func TestPendingRegistrationControlKeepsOnlyTheLatestState(t *testing.T)
{
name := "queued-binding"
other := "other-binding"
setConnection(t, nil)
forgetPending(t, name)
forgetPending(t, other)
assertQueue := func(step string, want ...string) {
t.Helper()
got := strings.Join(pendingLabels(name, other), " ")
if expected := strings.Join(want, " "); got != expected {
t.Fatalf("%s: queue = %q, want %q", step, got, expected)
}
}
release := RegisterComponentOwned("queued-root", name, nil)
Send(name, map[string]any{"cmd": "refresh"})
otherRelease := RegisterComponentOwned("other-root", other, nil)
assertQueue("after registration",
name+"/init", name+"/cmd:refresh", other+"/init")
release()
assertQueue("after release",
name+"/unsubscribe", name+"/cmd:refresh", other+"/init")
second := RegisterComponentOwned("queued-root", name, nil)
assertQueue("after re-registration",
name+"/init", name+"/cmd:refresh", other+"/init")
otherRelease()
assertQueue("after the other release",
name+"/init", name+"/cmd:refresh", other+"/unsubscribe")
second()
assertQueue("after the second release",
name+"/unsubscribe", name+"/cmd:refresh", other+"/unsubscribe")
}
TestConnectedRegistrationControlDropsTheQueuedControl
A control decided with the connection up goes out on the wire, so the queue
keeps neither it nor the state it superseded, and everything queued around it
stays where it was.
Parameters
func TestConnectedRegistrationControlDropsTheQueuedControl(t *testing.T)
{
name := "wired-queue-binding"
other := "wired-queue-other"
setConnection(t, nil)
forgetPending(t, name)
forgetPending(t, other)
release := RegisterComponentOwned("wired-queue-root", name, nil)
Send(name, map[string]any{"cmd": "refresh"})
otherRelease := RegisterComponentOwned("wired-queue-other-root", other, nil)
t.Cleanup(otherRelease)
c, _ := dialFake(t)
setConnection(t, c)
release()
got := strings.Join(pendingLabels(name, other), " ")
want := strings.Join([]string{name + "/cmd:refresh", other + "/init"}, " ")
if got != want {
t.Fatalf("queue = %q, want %q", got, want)
}
}
TestCoreComponentLifecycleOwnsItsBinding
This package installs the core registration hook from its own init, so a
mounted component binds through the real boundary and unmounting releases it,
with no test double in between.
Parameters
func TestCoreComponentLifecycleOwnsItsBinding(t *testing.T)
{
setConnection(t, nil)
forgetPending(t, "BoundaryHost")
component := core.NewHTMLComponent("BoundaryComponent", []byte(`<root></root>`), nil)
component.AddHostComponent("BoundaryHost")
component.Init(nil)
for cycle := 0; cycle < 100; cycle++ {
component.Mount()
component.Render()
binding, bound := bindingSnapshot("BoundaryHost")
if !bound || binding.id != component.ID {
t.Fatalf("cycle %d: mount bound %#v, want the component id %q", cycle, binding, component.ID)
}
component.Unmount()
if _, bound := bindingSnapshot("BoundaryHost"); bound {
t.Fatalf("cycle %d: unmount left the binding in place", cycle)
}
}
mu.RLock()
_, tracked := bindingTokens["BoundaryHost"]
mu.RUnlock()
if tracked {
t.Fatal("route churn retained the registration token")
}
inits, unsubscribes := pendingFor("BoundaryHost")
if inits != 0 || unsubscribes != 1 {
t.Fatalf("queued messages: init=%d unsubscribe=%d, want 0 and 1", inits, unsubscribes)
}
}
hostVarRoot
hostVarRoot builds a component root carrying one host variable element.
Parameters
Returns
func hostVarRoot(t *testing.T, id string) dom.Element
{
t.Helper()
root := dom.CreateElement("div")
root.SetAttr("data-component-id", id)
variable := dom.CreateElement("span")
variable.SetAttr("data-host-var", "value")
root.AppendChild(variable)
dom.Doc().Body().AppendChild(root)
t.Cleanup(func() { root.Call("remove") })
return root
}
appSentinel
appSentinel puts a host variable under #app, which is where
dom.ComponentRoot falls back when a component id resolves to nothing. Host
delivery must never take that fallback, so this element stays untouched.
Parameters
Returns
func appSentinel(t *testing.T) dom.Element
{
t.Helper()
app := dom.ByID("app")
if app.IsNull() || app.IsUndefined() {
app = dom.CreateElement("div")
app.SetAttr("id", "app")
dom.Doc().Body().AppendChild(app)
t.Cleanup(func() { app.Call("remove") })
}
sentinel := dom.CreateElement("span")
sentinel.SetAttr("data-host-var", "value")
sentinel.SetText("untouched")
app.AppendChild(sentinel)
t.Cleanup(func() { sentinel.Call("remove") })
return sentinel
}
hostVarText
Parameters
Returns
func hostVarText(root dom.Element) string
{
return root.Query(`[data-host-var="value"]`).Text()
}
waitForHostVar
Parameters
func waitForHostVar(t *testing.T, root dom.Element, want string)
{
t.Helper()
for i := 0; i < 200; i++ {
if hostVarText(root) == want {
return
}
time.Sleep(5 * time.Millisecond)
}
t.Fatalf("host variable = %q, want %q", hostVarText(root), want)
}
TestReleasedBindingIgnoresLateDelivery
A component that unmounts stops receiving host pushes: the frame that arrives
after its cleanup reaches neither the root it owned, which route exit removed
from the document, nor the #app shell the generic root lookup would fall back
to, while a component registered afterwards keeps receiving its own.
Parameters
func TestReleasedBindingIgnoresLateDelivery(t *testing.T)
{
conn, socket := dialFake(t)
forgetPending(t, "released-ticker")
forgetPending(t, "live-ticker")
sentinel := appSentinel(t)
releasedRoot := hostVarRoot(t, "released-root")
liveRoot := hostVarRoot(t, "live-root")
release := RegisterComponentOwned("released-root", "released-ticker", []string{"value"})
t.Cleanup(release)
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"released-ticker","payload":{"value":"1"},"sequence":1}`)
waitForHostVar(t, releasedRoot, "1")
// Route exit: the cleanup runs and the root it owned leaves the document.
release()
releasedRoot.Call("remove")
socket.deliver(`{"component":"released-ticker","payload":{"value":"2"},"sequence":2}`)
// Frames are applied in order, so a later frame reaching a live binding
// proves the released one was already processed.
releaseLive := RegisterComponentOwned("live-root", "live-ticker", []string{"value"})
t.Cleanup(releaseLive)
socket.deliver(`{"component":"live-ticker","payload":{"value":"3"},"sequence":3}`)
waitForHostVar(t, liveRoot, "3")
if got := hostVarText(releasedRoot); got != "1" {
t.Fatalf("the released root was updated to %q", got)
}
if got := sentinel.Text(); got != "untouched" {
t.Fatalf("a released frame fell back to #app and wrote %q", got)
}
}
withoutCSSEscape
withoutCSSEscape hides the CSS object for the duration of a test, which is
the runtime where the root has to be found by scanning the attribute.
Parameters
func withoutCSSEscape(t *testing.T)
{
t.Helper()
original := js.Get("CSS")
js.Set("CSS", js.Null())
t.Cleanup(func() { js.Set("CSS", original) })
}
TestDeliveryResolvesRootsWithCSSMetacharacters
RegisterComponent is public and takes any id, so an application can bind a
root whose id is not a CSS identifier. Delivery has to resolve exactly that
root: an id spliced into a selector raw would throw a DOMException out of
querySelector, or select a root the binding does not own.
Parameters
func TestDeliveryResolvesRootsWithCSSMetacharacters(t *testing.T)
{
cases := []struct {
label string
component string
id string
escape bool
}{
{label: "css escape", component: "meta-escaped", id: `meta"root'\]:.#[ 1`, escape: true},
{label: "attribute scan", component: "meta-scanned", id: `meta"root'\]:.#[ 2`, escape: false},
}
for _, tc := range cases {
t.Run(tc.label, func(t *testing.T) {
if !tc.escape {
withoutCSSEscape(t)
}
conn, socket := dialFake(t)
forgetPending(t, tc.component)
sentinel := appSentinel(t)
decoy := hostVarRoot(t, "meta-decoy")
root := hostVarRoot(t, tc.id)
release := RegisterComponentOwned(tc.id, tc.component, []string{"value"})
t.Cleanup(release)
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"` + tc.component + `","payload":{"value":"1"},"sequence":1}`)
waitForHostVar(t, root, "1")
if got := hostVarText(decoy); got != "" {
t.Fatalf("the frame also wrote %q into a root it does not own", got)
}
if got := sentinel.Text(); got != "untouched" {
t.Fatalf("the frame fell back to #app and wrote %q", got)
}
})
}
}
TestCleanupDuringAnInFlightDeliveryWins
The window a cleanup has to close is the one between the read loop
snapshotting a binding and the delivery touching the DOM. The barrier runs
inside exactly that window, on a frame the loop has already accepted, so the
overlap is reproduced rather than raced for.
Parameters
func TestCleanupDuringAnInFlightDeliveryWins(t *testing.T)
{
conn, socket := dialFake(t)
forgetPending(t, "overlap-ticker")
forgetPending(t, "overlap-live")
sentinel := appSentinel(t)
overlapRoot := hostVarRoot(t, "overlap-root")
liveRoot := hostVarRoot(t, "overlap-live-root")
release := RegisterComponentOwned("overlap-root", "overlap-ticker", []string{"value"})
t.Cleanup(release)
releaseLive := RegisterComponentOwned("overlap-live-root", "overlap-live", []string{"value"})
t.Cleanup(releaseLive)
barriers := make(chan struct{}, 1)
setSnapshotBarrier(func(component string) {
if component != "overlap-ticker" {
return
}
setSnapshotBarrier(nil)
release()
overlapRoot.Call("remove")
barriers <- struct{}{}
})
t.Cleanup(func() { setSnapshotBarrier(nil) })
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"overlap-ticker","payload":{"value":"9"},"sequence":1}`)
socket.deliver(`{"component":"overlap-live","payload":{"value":"10"},"sequence":2}`)
waitForHostVar(t, liveRoot, "10")
select {
case <-barriers:
default:
t.Fatal("the delivery never reached the snapshot barrier")
}
if got := hostVarText(overlapRoot); got != "" {
t.Fatalf("a delivery snapshotted before the cleanup still wrote %q", got)
}
if got := sentinel.Text(); got != "untouched" {
t.Fatalf("the overlapping delivery fell back to #app and wrote %q", got)
}
}
ConnectionState
ConnectionState describes the SSC transport state.
type ConnectionState string
ConnectionStateSignal
ConnectionStateSignal returns the reactive SSC connection state.
Returns
func ConnectionStateSignal() *state.Signal[ConnectionState]
{
return connectionState
}
setHostSignalWriteHook
setHostSignalWriteHook installs the delivery hook under the lock the read
loop reads it with.
Parameters
func setHostSignalWriteHook(fn func(component, name string))
{
mu.Lock()
beforeHostSignalWrite = fn
mu.Unlock()
}
hostSignalRoot
hostSignalRoot builds the root a binding owns and registers one host signal
per name against it, which is the state a rendered component leaves behind.
Parameters
Returns
func hostSignalRoot(t *testing.T, id string, names ...string) map[string]*state.Signal[string]
{
t.Helper()
root := dom.CreateElement("div")
root.SetAttr("data-component-id", id)
dom.Doc().Body().AppendChild(root)
signals := make(map[string]*state.Signal[string], len(names))
for _, name := range names {
signal := state.NewSignal("initial")
dom.RegisterSignal(id, name, signal)
signals[name] = signal
}
t.Cleanup(func() {
dom.RemoveComponentSignals(id)
root.Call("remove")
})
return signals
}
waitForSignal
Parameters
func waitForSignal(t *testing.T, signal *state.Signal[string], want string)
{
t.Helper()
for i := 0; i < 200; i++ {
if signal.Get() == want {
return
}
time.Sleep(5 * time.Millisecond)
}
t.Fatalf("host signal = %q, want %q", signal.Get(), want)
}
untouched
untouched reports the signals of a component that still hold their initial
value, so a test can count how much of a frame landed.
Parameters
Returns
func untouched(signals map[string]*state.Signal[string]) int
{
count := 0
for _, signal := range signals {
if signal.Get() == "initial" {
count++
}
}
return count
}
TestCleanupBetweenTheSignalCheckAndTheWriteWins
The window a cleanup has to close for host signals is the one between the
ownership check a delivery makes and the write it was about to perform. The
hook runs inside exactly that window, on a frame already accepted, so the
overlap is reproduced rather than raced for: the cleanup returns first, and
the frame it overlapped leaves the released component’s signals alone.
Parameters
func TestCleanupBetweenTheSignalCheckAndTheWriteWins(t *testing.T)
{
conn, socket := dialFake(t)
forgetPending(t, "signal-ticker")
forgetPending(t, "signal-live")
released := hostSignalRoot(t, "signal-root", "value")
live := hostSignalRoot(t, "signal-live-root", "value")
release := RegisterComponentOwned("signal-root", "signal-ticker", []string{"value"})
t.Cleanup(release)
releaseLive := RegisterComponentOwned("signal-live-root", "signal-live", []string{"value"})
t.Cleanup(releaseLive)
cleanupReturned := make(chan struct{})
setHostSignalWriteHook(func(component, _ string) {
if component != "signal-ticker" {
return
}
setHostSignalWriteHook(nil)
// Route exit runs on its own goroutine, as it does in a browser, and
// the delivery waits for the cleanup to return before its write.
go func() {
release()
close(cleanupReturned)
}()
select {
case <-cleanupReturned:
case <-time.After(2 * time.Second):
t.Error("the cleanup did not return while a delivery was mid-frame")
}
})
t.Cleanup(func() { setHostSignalWriteHook(nil) })
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"signal-ticker","payload":{"value":"stale"},"sequence":1}`)
// Frames are applied in order, so a later frame reaching a live binding
// proves the overlapped one was already processed.
socket.deliver(`{"component":"signal-live","payload":{"value":"fresh"},"sequence":2}`)
waitForSignal(t, live["value"], "fresh")
select {
case <-cleanupReturned:
default:
t.Fatal("the delivery never reached the signal write hook")
}
if got := released["value"].Get(); got != "initial" {
t.Fatalf("a frame that overlapped the cleanup wrote %q into a released host signal", got)
}
}
TestSetterReleasingItsOwnBindingEndsTheFrame
A host signal setter that unmounts its own component re-enters the release
from inside the frame being applied. The release must not wait on the
delivery it is running under, and the rest of that frame belongs to the
registration the setter just dropped.
Parameters
func TestSetterReleasingItsOwnBindingEndsTheFrame(t *testing.T)
{
conn, socket := dialFake(t)
forgetPending(t, "reentrant-ticker")
forgetPending(t, "reentrant-live")
signals := hostSignalRoot(t, "reentrant-root", "first", "second")
live := hostSignalRoot(t, "reentrant-live-root", "value")
release := RegisterComponentOwned("reentrant-root", "reentrant-ticker", []string{"first", "second"})
t.Cleanup(release)
releaseLive := RegisterComponentOwned("reentrant-live-root", "reentrant-live", []string{"value"})
t.Cleanup(releaseLive)
var unmount sync.Once
for _, signal := range signals {
signal.OnChange(func(string) { unmount.Do(release) })
}
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"reentrant-ticker","payload":{"first":"one","second":"two"},"sequence":1}`)
// The witness frame is what proves the read loop came back from the
// re-entrant release instead of deadlocking inside it.
socket.deliver(`{"component":"reentrant-live","payload":{"value":"fresh"},"sequence":2}`)
waitForSignal(t, live["value"], "fresh")
if remaining := untouched(signals); remaining != 1 {
t.Fatalf("host signals left untouched = %d, want 1: the release from a setter did not end the frame", remaining)
}
if _, bound := bindingSnapshot("reentrant-ticker"); bound {
t.Fatal("the binding survived the release its own setter ran")
}
}
TestSetterRebindingItsComponentEndsTheFrame
A setter that mounts the component again re-enters registration from inside
the frame. The replacement is what stays bound, and the frame the superseded
registration was carrying ends where it was superseded.
Parameters
func TestSetterRebindingItsComponentEndsTheFrame(t *testing.T)
{
conn, socket := dialFake(t)
forgetPending(t, "rebound-ticker")
forgetPending(t, "rebound-live")
signals := hostSignalRoot(t, "rebound-root", "first", "second")
live := hostSignalRoot(t, "rebound-live-root", "value")
release := RegisterComponentOwned("rebound-root", "rebound-ticker", []string{"first", "second"})
t.Cleanup(release)
releaseLive := RegisterComponentOwned("rebound-live-root", "rebound-live", []string{"value"})
t.Cleanup(releaseLive)
var remount sync.Once
rebound := make(chan func(), 1)
for _, signal := range signals {
signal.OnChange(func(string) {
remount.Do(func() {
rebound <- RegisterComponentOwned("remounted-root", "rebound-ticker", []string{"first"})
})
})
}
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"rebound-ticker","payload":{"first":"one","second":"two"},"sequence":1}`)
socket.deliver(`{"component":"rebound-live","payload":{"value":"fresh"},"sequence":2}`)
waitForSignal(t, live["value"], "fresh")
select {
case releaseRemounted := <-rebound:
t.Cleanup(releaseRemounted)
default:
t.Fatal("the setter never re-registered the component")
}
if remaining := untouched(signals); remaining != 1 {
t.Fatalf("host signals left untouched = %d, want 1: the registration from a setter did not end the frame", remaining)
}
binding, bound := bindingSnapshot("rebound-ticker")
if !bound || binding.id != "remounted-root" {
t.Fatalf("the replacement registration is not the live one: %#v bound=%v", binding, bound)
}
}
hostVarElement
type hostVarElement interface
componentRoot
type componentRoot interface
Methods
hydrationMismatch
type hydrationMismatch struct
Fields
| Name | Type | Description |
|---|---|---|
| VarName | string | |
| Expected | string | |
| Actual | string | |
| ActualHash | string | |
| ExpectedAlg | string |
initSnapshotPayload
type initSnapshotPayload struct
Fields
| Name | Type | Description |
|---|---|---|
| HTML | string | |
| Vars | []string |
encodeExpectation
Parameters
Returns
func encodeExpectation(value string) string
{
sum := sha256.Sum256([]byte(value))
return fmt.Sprintf("%s:%s", expectationHashAlg, hex.EncodeToString(sum[:]))
}
expectationMatches
Parameters
Returns
func expectationMatches(expectedAttr, actual string) (bool, string, string)
{
actualHash := encodeExpectation(actual)
if expectedAttr == "" {
return true, expectationHashAlg, actualHash
}
if strings.HasPrefix(expectedAttr, expectationHashAlg+":") {
return expectedAttr == actualHash, expectationHashAlg, actualHash
}
return expectedAttr == actual, "raw", actualHash
}
updateHostVar
Parameters
Returns
func updateHostVar(root componentRoot, name, value string) *hydrationMismatch
{
node := root.HostVar(name)
if !node.Exists() {
return nil
}
expectedAttr := node.Attr(hostExpectedAttr)
actualText := node.Text()
matches, alg, actualHash := expectationMatches(expectedAttr, actualText)
if !matches {
return &hydrationMismatch{
VarName: name,
Expected: expectedAttr,
Actual: actualText,
ActualHash: actualHash,
ExpectedAlg: alg,
}
}
node.SetText(value)
node.SetAttr(hostExpectedAttr, encodeExpectation(value))
return nil
}
handleHostPayload
Parameters
Returns
func handleHostPayload(root componentRoot, payload map[string]any, updateSignal func(name string, raw any)) []hydrationMismatch
{
mismatches := make([]hydrationMismatch, 0)
for key, raw := range payload {
if key == "initSnapshot" || strings.HasPrefix(key, "_") {
continue
}
mismatch := updateHostVar(root, key, fmt.Sprintf("%v", raw))
if mismatch != nil {
mismatches = append(mismatches, *mismatch)
}
if updateSignal != nil {
updateSignal(key, raw)
}
}
return mismatches
}
applyInitSnapshot
Parameters
func applyInitSnapshot(root componentRoot, payload *initSnapshotPayload)
{
if payload == nil {
return
}
root.SetHTML(payload.HTML)
}
buildResyncPayload
Parameters
Returns
func buildResyncPayload(mismatches []hydrationMismatch) map[string]any
{
entries := make([]map[string]string, 0, len(mismatches))
for _, m := range mismatches {
entries = append(entries, map[string]string{
"var": m.VarName,
"expected": m.Expected,
"expectedAlg": m.ExpectedAlg,
"actual": m.Actual,
"actualHash": m.ActualHash,
})
}
return map[string]any{
"resync": map[string]any{
"reason": "host-var-mismatch",
"vars": entries,
},
}
}
preferredHostTransport
Returns
func preferredHostTransport() string
{
value := strings.ToLower(strings.TrimSpace(stdjs.Global().Get("RFW_TRANSPORT").String()))
switch value {
case hostTransportStreamBus, "webtransport", "warp-streambus":
return hostTransportStreamBus
case hostTransportAuto:
return hostTransportAuto
default:
return hostTransportWebSocket
}
}
hostStreamBusURL
Returns
func hostStreamBusURL() string
{
if configured := stdjs.Global().Get("RFW_STREAMBUS_URL"); configured.Truthy() {
return normalizeStreamBusURL(configured.String(), false)
}
return normalizeStreamBusURL(hostWSURL(), true)
}
normalizeStreamBusURL
Parameters
Returns
func normalizeStreamBusURL(raw string, incrementHTTPPort bool) string
{
raw = strings.TrimSpace(raw)
if raw == "" {
return ""
}
switch {
case strings.HasPrefix(raw, "wss://"):
raw = "https://" + strings.TrimPrefix(raw, "wss://")
incrementHTTPPort = false
case strings.HasPrefix(raw, "ws://"):
raw = "https://" + strings.TrimPrefix(raw, "ws://")
case strings.HasPrefix(raw, "http://"):
raw = "https://" + strings.TrimPrefix(raw, "http://")
case strings.HasPrefix(raw, "https://"):
incrementHTTPPort = false
default:
raw = "https://" + raw
}
parsed, err := url.Parse(raw)
if err != nil || parsed.Hostname() == "" {
return ""
}
if incrementHTTPPort {
if port, err := strconv.Atoi(parsed.Port()); err == nil && port > 0 {
parsed.Host = net.JoinHostPort(parsed.Hostname(), strconv.Itoa(port+1))
}
}
if parsed.Path == "" || parsed.Path == "/" || parsed.Path == "/ws" {
parsed.Path = streamBusPath
}
parsed.RawQuery = ""
parsed.Fragment = ""
return parsed.String()
}
dialStreamBus
Parameters
Returns
func dialStreamBus(ctx context.Context, endpoint string) (*hostConn, error)
{
constructor := stdjs.Global().Get("WebTransport")
if constructor.Type() != stdjs.TypeFunction {
return nil, errors.New("hostclient: WebTransport API is unavailable")
}
config, _ := loadStreamBusConfig(ctx)
if config.port != "" {
endpoint = replaceURLPort(endpoint, config.port)
}
var transport stdjs.Value
if config.options.Type() == stdjs.TypeObject {
transport = constructor.New(endpoint, config.options)
} else {
transport = constructor.New(endpoint)
}
if _, err := awaitPromise(ctx, transport.Get("ready")); err != nil {
transport.Call("close")
return nil, err
}
stream, err := awaitPromise(ctx, transport.Call("createBidirectionalStream"))
if err != nil {
transport.Call("close")
return nil, err
}
reader := &jsReader{reader: stream.Get("readable").Call("getReader")}
writer := &jsWriter{writer: stream.Get("writable").Call("getWriter")}
buffered := bufio.NewReader(reader)
c := newHostConn()
c.send = func(data string) error {
return writeStreamBusFrame(context.Background(), writer, []byte(data))
}
c.shutdown = func() error {
reader.reader.Call("cancel")
writer.writer.Call("close")
transport.Call("close")
return nil
}
c.openOnce.Do(func() { close(c.open) })
go func() {
for {
payload, readErr := readStreamBusFrame(buffered, int(maxInboundMessageBytes))
if readErr != nil {
c.fail(readErr)
return
}
c.deliver(payload)
}
}()
return c, nil
}
streamBusClientConfig
type streamBusClientConfig struct
Fields
| Name | Type | Description |
|---|---|---|
| options | stdjs.Value | |
| port | string |
loadStreamBusConfig
Parameters
Returns
func loadStreamBusConfig(ctx context.Context) (streamBusClientConfig, error)
{
fetch := stdjs.Global().Get("fetch")
if fetch.Type() != stdjs.TypeFunction {
return streamBusClientConfig{options: stdjs.Undefined()}, nil
}
response, err := awaitPromise(ctx, fetch.Invoke("/__rfw/streambus-config"))
if err != nil || !response.Get("ok").Bool() {
return streamBusClientConfig{options: stdjs.Undefined()}, err
}
config, err := awaitPromise(ctx, response.Call("json"))
if err != nil {
return streamBusClientConfig{options: stdjs.Undefined()}, err
}
hash, err := hex.DecodeString(config.Get("certificateHash").String())
if err != nil || len(hash) != 32 {
return streamBusClientConfig{options: stdjs.Undefined()}, errors.New("hostclient: invalid StreamBus certificate hash")
}
value := stdjs.Global().Get("Uint8Array").New(len(hash))
stdjs.CopyBytesToJS(value, hash)
descriptor := stdjs.Global().Get("Object").New()
descriptor.Set("algorithm", "sha-256")
descriptor.Set("value", value.Get("buffer"))
list := stdjs.Global().Get("Array").New()
list.Call("push", descriptor)
options := stdjs.Global().Get("Object").New()
options.Set("serverCertificateHashes", list)
return streamBusClientConfig{options: options, port: config.Get("port").String()}, nil
}
replaceURLPort
Parameters
Returns
func replaceURLPort(raw, port string) string
{
parsed, err := url.Parse(raw)
if err != nil || parsed.Hostname() == "" || port == "" {
return raw
}
parsed.Host = net.JoinHostPort(parsed.Hostname(), port)
return parsed.String()
}
writeStreamBusFrame
Parameters
Returns
func writeStreamBusFrame(ctx context.Context, writer *jsWriter, payload []byte) error
{
var prefix [binary.MaxVarintLen64]byte
n := binary.PutUvarint(prefix[:], uint64(len(payload)))
if err := writer.write(ctx, prefix[:n]); err != nil {
return err
}
return writer.write(ctx, payload)
}
readStreamBusFrame
Parameters
Returns
func readStreamBusFrame(reader *bufio.Reader, maximum int) ([]byte, error)
{
size, err := binary.ReadUvarint(reader)
if err != nil {
return nil, err
}
if maximum > 0 && size > uint64(maximum) {
return nil, fmt.Errorf("hostclient: StreamBus frame of %d bytes exceeds %d", size, maximum)
}
if size > uint64(^uint(0)>>1) {
return nil, fmt.Errorf("hostclient: StreamBus frame of %d bytes exceeds platform limit", size)
}
payload := make([]byte, size)
if _, err := io.ReadFull(reader, payload); err != nil {
return nil, err
}
return payload, nil
}
jsReader
type jsReader struct
Methods
Parameters
Returns
func (*jsReader) Read(payload []byte) (int, error)
{
if len(r.buffer) == 0 {
result, err := awaitPromise(context.Background(), r.reader.Call("read"))
if err != nil {
return 0, err
}
if result.Get("done").Bool() {
return 0, io.EOF
}
value := result.Get("value")
r.buffer = make([]byte, value.Get("byteLength").Int())
stdjs.CopyBytesToGo(r.buffer, value)
}
n := copy(payload, r.buffer)
r.buffer = r.buffer[n:]
return n, nil
}
Fields
| Name | Type | Description |
|---|---|---|
| reader | stdjs.Value | |
| buffer | []byte |
jsWriter
type jsWriter struct
Methods
Parameters
Returns
func (*jsWriter) write(ctx context.Context, payload []byte) error
{
array := stdjs.Global().Get("Uint8Array").New(len(payload))
stdjs.CopyBytesToJS(array, payload)
_, err := awaitPromise(ctx, w.writer.Call("write", array))
return err
}
Fields
| Name | Type | Description |
|---|---|---|
| writer | stdjs.Value |
promiseResult
type promiseResult struct
Fields
| Name | Type | Description |
|---|---|---|
| value | stdjs.Value | |
| err | error |
awaitPromise
Parameters
Returns
func awaitPromise(ctx context.Context, promise stdjs.Value) (stdjs.Value, error)
{
result := make(chan promiseResult, 1)
resolve := stdjs.FuncOf(func(_ stdjs.Value, args []stdjs.Value) any {
value := stdjs.Undefined()
if len(args) > 0 {
value = args[0]
}
result <- promiseResult{value: value}
return nil
})
reject := stdjs.FuncOf(func(_ stdjs.Value, args []stdjs.Value) any {
message := "promise rejected"
if len(args) > 0 {
message = args[0].String()
}
result <- promiseResult{err: errors.New(message)}
return nil
})
promise.Call("then", resolve).Call("catch", reject)
select {
case settled := <-result:
resolve.Release()
reject.Release()
return settled.value, settled.err
case <-ctx.Done():
go func() {
<-result
resolve.Release()
reject.Release()
}()
return stdjs.Undefined(), ctx.Err()
}
}
TestStreamBusURLUsesAdvertisedHTTP3Port
Parameters
func TestStreamBusURLUsesAdvertisedHTTP3Port(t *testing.T)
{
got := replaceURLPort("https://localhost:8081/streambus", "8083")
if got != "https://localhost:8083/streambus" {
t.Fatalf("URL = %q", got)
}
}
TestNormalizeStreamBusURLFromDevelopmentWebSocket
Parameters
func TestNormalizeStreamBusURLFromDevelopmentWebSocket(t *testing.T)
{
got := normalizeStreamBusURL("ws://localhost:8080/ws", true)
if got != "https://localhost:8081/streambus" {
t.Fatalf("URL = %q", got)
}
}
TestPreferredHostTransportReadsGeneratedConfig
Parameters
func TestPreferredHostTransportReadsGeneratedConfig(t *testing.T)
{
global := stdjs.Global()
previous := global.Get("RFW_TRANSPORT")
global.Set("RFW_TRANSPORT", "streambus")
t.Cleanup(func() {
if previous.Type() == stdjs.TypeUndefined {
global.Delete("RFW_TRANSPORT")
} else {
global.Set("RFW_TRANSPORT", previous)
}
})
if got := preferredHostTransport(); got != hostTransportStreamBus {
t.Fatalf("transport = %q", got)
}
}
hostConn
hostConn adapts the browser WebSocket to the blocking read and write calls
the connection loops expect.
type hostConn struct
Methods
Parameters
func (*hostConn) deliver(payload []byte)
{
if int64(len(payload)) > maxInboundMessageBytes {
c.fail(fmt.Errorf(
"hostclient: inbound message of %d bytes exceeds the %d byte limit",
len(payload), maxInboundMessageBytes,
))
return
}
select {
case c.messages <- payload:
c.received.Add(1)
default:
c.fail(errors.New("hostclient: inbound frame queue overflow"))
}
}
Parameters
func (*hostConn) fail(err error)
{
c.closeOnce.Do(func() {
if err == nil {
err = errConnectionClosed
}
c.closeErr.Store(&err)
close(c.done)
})
}
err reports why the connection ended, once done is closed.
Returns
func (*hostConn) err() error
{
if stored := c.closeErr.Load(); stored != nil {
return *stored
}
return errConnectionClosed
}
read returns the next frame. Buffered frames always win over a close: a select with both cases ready picks at random, so the queue is drained without blocking first and only an empty queue lets the close through. A connection that ends mid-burst still delivers every frame the host sent, which is what ordered delivery depends on.
Parameters
Returns
func (*hostConn) read(ctx context.Context) ([]byte, error)
{
select {
case payload := <-c.messages:
return payload, nil
default:
}
select {
case payload := <-c.messages:
return payload, nil
case <-c.done:
// The close may have raced a frame that landed between the two
// selects. done stays closed, so returning the frame now still lets
// the next read report the failure.
select {
case payload := <-c.messages:
return payload, nil
default:
}
return nil, c.err()
case <-ctx.Done():
return nil, ctx.Err()
}
}
Parameters
Returns
func (*hostConn) writeJSON(payload any) error
{
data, err := json.Marshal(payload)
if err != nil {
return err
}
if c.send == nil {
return errConnectionClosed
}
return c.send(string(data))
}
Returns
func (*hostConn) close() error
{
var err error
if c.shutdown != nil {
err = c.shutdown()
} else if c.socket != nil {
err = c.socket.Close(normalClosure, "connection closed")
}
if c.release != nil {
c.release()
c.release = nil
}
return err
}
heartbeat probes liveness with a control message. The browser WebSocket API gives script no access to protocol ping and pong frames, so a half-open connection can only be detected by asking the host to answer over the same JSON channel.
Parameters
Returns
func (*hostConn) heartbeat(ctx context.Context, interval, timeout time.Duration) error
{
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
case <-ctx.Done():
return ctx.Err()
}
// Snapshot before probing. deliver counts a frame when the browser
// hands it over, not when readLoop drains it, so every frame already
// queued is behind this number: a half-open socket cannot hide behind
// a backlog. Any arrival counts, not just the pong, which is what lets
// a host that predates the control message answer with its generic ack.
delivered := c.received.Load()
if err := sendControl(c, "ping"); err != nil {
return err
}
timer := time.NewTimer(timeout)
select {
case <-timer.C:
case <-ctx.Done():
timer.Stop()
return ctx.Err()
}
if c.received.Load() == delivered {
return errors.New("hostclient: heartbeat timed out")
}
}
}
Fields
| Name | Type | Description |
|---|---|---|
| socket | *js.Socket | |
| send | func(string) error | |
| shutdown | func() error | |
| release | func() | |
| generation | uint64 | |
| messages | chan []byte | |
| open | chan struct{} | |
| done | chan struct{} | |
| openOnce | sync.Once | |
| closeOnce | sync.Once | |
| closeErr | atomic.Pointer[error] | |
| received | atomic.Uint64 |
newHostConn
Returns
func newHostConn() *hostConn
{
return &hostConn{
messages: make(chan []byte, inboundFrames),
open: make(chan struct{}),
done: make(chan struct{}),
}
}
dialBrowser
dialBrowser opens a connection and waits for the browser handshake to
complete. It remains the default SSC transport for existing applications.
Parameters
Returns
func dialBrowser(ctx context.Context, url string) (*hostConn, error)
{
c := newHostConn()
socket, err := js.OpenSocket(url, js.SocketHandlers{
MaxMessageBytes: int(maxInboundMessageBytes),
Open: func() { c.openOnce.Do(func() { close(c.open) }) },
Text: func(text string) { c.deliver([]byte(text)) },
Bytes: c.deliver,
Error: c.fail,
Close: func(code int, reason string, _ bool) {
c.fail(fmt.Errorf("hostclient: connection closed with code %d: %s", code, reason))
},
})
if err != nil {
return nil, err
}
c.socket = socket
c.send = socket.SendText
c.shutdown = func() error { return socket.Close(normalClosure, "connection closed") }
select {
case <-c.open:
return c, nil
case <-c.done:
_ = c.close()
return nil, c.err()
case <-ctx.Done():
_ = c.close()
return nil, ctx.Err()
}
}
sendControl
sendControl writes an unsequenced control frame. It carries no sequence so it
never enters the outbox and is never replayed after a reconnect.
Parameters
Returns
func sendControl(c *hostConn, control string) error
{
deliveryMu.Lock()
ack := lastInbound
deliveryMu.Unlock()
sendMu.Lock()
defer sendMu.Unlock()
return c.writeJSON(wireMessage{Control: control, Ack: ack})
}
fakeSocket
type fakeSocket struct
Methods
open moves the socket to OPEN and fires onopen, the way a browser does once the handshake completes.
func (fakeSocket) open()
{
f.value.Set("readyState", js.ValueOf(1))
f.value.Call("onopen", js.NewDict().Value)
}
deliver fires onmessage with a text frame.
Parameters
func (fakeSocket) deliver(frame string)
{
event := js.NewDict()
event.Set("data", frame)
f.value.Call("onmessage", event.Value)
}
serverClose fires onclose the way a host closing the connection does.
Parameters
func (fakeSocket) serverClose(code int, reason string)
{
f.value.Set("readyState", js.ValueOf(3))
event := js.NewDict()
event.Set("code", code)
event.Set("reason", reason)
event.Set("wasClean", true)
f.value.Call("onclose", event.Value)
}
sent returns every frame the client wrote to this socket.
Returns
func (fakeSocket) sent() []string
{
raw := f.value.Get("sent")
out := make([]string, raw.Get("length").Int())
for i := range out {
out[i] = raw.Index(i).String()
}
return out
}
handlersBound reports whether the socket still has Go callbacks attached.
Returns
func (fakeSocket) handlersBound() bool
{
for _, event := range []string{"onopen", "onmessage", "onerror", "onclose"} {
if f.value.Get(event).Truthy() {
return true
}
}
return false
}
installFakeSockets
installFakeSockets swaps window.WebSocket for the scripted constructor and
restores the real one on cleanup. Cleanup rather than a deferred call in the
test body: a failed assertion must not leave the fake installed, or the
tests in js/websocket_test.go that dial for real start failing instead.
Parameters
func installFakeSockets(t *testing.T)
{
t.Helper()
js.Call("eval", fakeSocketSource)
original := js.Get("WebSocket")
js.Set("WebSocket", js.Get("__FakeWebSocket"))
js.Set("__fakeSockets", js.NewArray().Value)
t.Cleanup(func() {
js.Set("WebSocket", original)
js.Global().Delete("__fakeSockets")
js.Global().Delete("__FakeWebSocket")
})
}
lastSocket
lastSocket returns the most recently constructed fake.
Parameters
Returns
func lastSocket(t *testing.T) fakeSocket
{
t.Helper()
sockets := js.Get("__fakeSockets")
length := sockets.Get("length").Int()
if length == 0 {
t.Fatal("no socket was constructed")
}
return fakeSocket{t: t, value: sockets.Index(length - 1)}
}
socketCount
Returns
func socketCount() int
{ return js.Get("__fakeSockets").Get("length").Int() }
dialFake
dialFake opens a connection through the fake and completes the handshake.
Parameters
Returns
func dialFake(t *testing.T) (*hostConn, fakeSocket)
{
t.Helper()
installFakeSockets(t)
resetDeliveryState(t)
type dialed struct {
conn *hostConn
err error
}
done := make(chan dialed, 1)
go func() {
c, err := dial(context.Background(), "ws://host.invalid/ws")
done <- dialed{conn: c, err: err}
}()
socket := waitForSocket(t)
socket.open()
select {
case result := <-done:
if result.err != nil {
t.Fatalf("dial: %v", result.err)
}
t.Cleanup(func() { _ = result.conn.close() })
return result.conn, socket
case <-time.After(2 * time.Second):
t.Fatal("dial did not complete after the handshake")
return nil, fakeSocket{}
}
}
waitForSocket
waitForSocket yields to the JavaScript event loop until the constructor has
run. The dial goroutine reaches OpenSocket only once it is scheduled.
Parameters
Returns
func waitForSocket(t *testing.T) fakeSocket
{
t.Helper()
for i := 0; i < 200; i++ {
if socketCount() > 0 {
return lastSocket(t)
}
time.Sleep(5 * time.Millisecond)
}
t.Fatal("the client never constructed a socket")
return fakeSocket{}
}
resetDeliveryState
resetDeliveryState clears the package level delivery bookkeeping so tests do
not inherit sequence numbers or pending calls from each other.
Parameters
func resetDeliveryState(t *testing.T)
{
t.Helper()
deliveryMu.Lock()
lastInbound = 0
nextOutbound = 0
resumeToken = ""
outbox = map[uint64]message{}
deliveryMu.Unlock()
callMu.Lock()
pendingCalls = map[string]chan actionReply{}
callMu.Unlock()
t.Cleanup(func() {
deliveryMu.Lock()
lastInbound = 0
nextOutbound = 0
resumeToken = ""
outbox = map[uint64]message{}
deliveryMu.Unlock()
callMu.Lock()
pendingCalls = map[string]chan actionReply{}
callMu.Unlock()
})
}
readOnce
readOnce runs readLoop until it returns, so a test can assert on the error a
single frame produces.
Parameters
Returns
func readOnce(c *hostConn, timeout time.Duration) error
{
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
return readLoop(ctx, c)
}
TestBrowserSocketCompletesTheHandshakeAndWrites
Parameters
func TestBrowserSocketCompletesTheHandshakeAndWrites(t *testing.T)
{
conn, socket := dialFake(t)
if got := socket.value.Get("binaryType").String(); got != "arraybuffer" {
t.Fatalf("binaryType = %q, want arraybuffer", got)
}
if got := socket.value.Get("url").String(); got != "ws://host.invalid/ws" {
t.Fatalf("url = %q", got)
}
if err := conn.writeJSON(wireMessage{Component: "counter", Sequence: 1}); err != nil {
t.Fatalf("write: %v", err)
}
frames := socket.sent()
if len(frames) != 1 || !strings.Contains(frames[0], `"component":"counter"`) {
t.Fatalf("frames written = %v", frames)
}
}
TestBrowserSocketDeliversAnActionReply
Parameters
func TestBrowserSocketDeliversAnActionReply(t *testing.T)
{
conn, socket := dialFake(t)
reply := make(chan actionReply, 1)
callMu.Lock()
pendingCalls["call-1"] = reply
callMu.Unlock()
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"id":"call-1","payload":{"total":7},"sequence":1}`)
select {
case got := <-reply:
payload, _ := got.payload.(map[string]any)
if payload["total"] != float64(7) {
t.Fatalf("reply payload = %#v", got.payload)
}
case <-time.After(2 * time.Second):
t.Fatal("the action reply never reached its caller")
}
}
TestBrowserSocketDeliversASubscriptionPayload
Parameters
func TestBrowserSocketDeliversASubscriptionPayload(t *testing.T)
{
conn, socket := dialFake(t)
received := make(chan map[string]any, 1)
mu.Lock()
handlers["ticker"] = func(payload map[string]any) { received <- payload }
mu.Unlock()
t.Cleanup(func() {
mu.Lock()
delete(handlers, "ticker")
mu.Unlock()
})
go func() { _ = readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"ticker","payload":{"price":1.5},"sequence":1}`)
select {
case payload := <-received:
if payload["price"] != 1.5 {
t.Fatalf("handler payload = %#v", payload)
}
case <-time.After(2 * time.Second):
t.Fatal("the subscription payload never reached its handler")
}
}
TestBrowserSocketRejectsAMalformedFrame
Parameters
func TestBrowserSocketRejectsAMalformedFrame(t *testing.T)
{
conn, socket := dialFake(t)
errCh := make(chan error, 1)
go func() { errCh <- readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":`)
select {
case err := <-errCh:
if err == nil {
t.Fatal("a malformed frame was accepted")
}
if _, ok := err.(*json.SyntaxError); !ok {
t.Fatalf("malformed frame error = %v (%T), want a JSON syntax error", err, err)
}
case <-time.After(2 * time.Second):
t.Fatal("a malformed frame froze the read loop")
}
}
TestBrowserSocketReportsAServerClosure
Parameters
func TestBrowserSocketReportsAServerClosure(t *testing.T)
{
conn, socket := dialFake(t)
errCh := make(chan error, 1)
go func() { errCh <- readOnce(conn, 2*time.Second) }()
socket.serverClose(1001, "going away")
select {
case err := <-errCh:
if err == nil {
t.Fatal("a server closure was not reported")
}
if !strings.Contains(err.Error(), "1001") {
t.Fatalf("closure error = %v, want the close code", err)
}
case <-time.After(2 * time.Second):
t.Fatal("a server closure left the read loop running")
}
}
TestBrowserSocketRejectsASequenceGap
A frame that skips a sequence desyncs the session instead of applying out of
order, which is what forces the reconnect and replay path.
Parameters
func TestBrowserSocketRejectsASequenceGap(t *testing.T)
{
conn, socket := dialFake(t)
errCh := make(chan error, 1)
go func() { errCh <- readOnce(conn, 2*time.Second) }()
socket.deliver(`{"component":"ticker","sequence":1}`)
socket.deliver(`{"component":"ticker","sequence":3}`)
select {
case err := <-errCh:
if err == nil || !strings.Contains(err.Error(), "sequence gap") {
t.Fatalf("sequence gap error = %v", err)
}
if got := connectionState.Get(); got != ConnectionDesynced {
t.Fatalf("connection state = %v, want desynced", got)
}
case <-time.After(2 * time.Second):
t.Fatal("a sequence gap did not end the read loop")
}
}
TestBrowserSocketConsumesARejectedResume
The host rejects a resume it cannot honour with a control frame. The client
must consume it without treating it as component traffic.
Parameters
func TestBrowserSocketConsumesARejectedResume(t *testing.T)
{
conn, socket := dialFake(t)
mu.Lock()
handlers[""] = func(map[string]any) { t.Error("a control frame reached a component handler") }
mu.Unlock()
t.Cleanup(func() {
mu.Lock()
delete(handlers, "")
mu.Unlock()
})
errCh := make(chan error, 1)
go func() { errCh <- readOnce(conn, 400*time.Millisecond) }()
socket.deliver(`{"control":"resume_rejected","error":{"code":"resume_rejected","message":"session could not be resumed"}}`)
if err := <-errCh; err != context.DeadlineExceeded {
t.Fatalf("read loop ended with %v, want it still running until the deadline", err)
}
}
TestBrowserSocketReconnectsOnANewSocket
Parameters
func TestBrowserSocketReconnectsOnANewSocket(t *testing.T)
{
conn, socket := dialFake(t)
socket.serverClose(1006, "abnormal")
if _, err := conn.read(context.Background()); err == nil {
t.Fatal("the closed connection kept reading")
}
before := socketCount()
type dialed struct {
conn *hostConn
err error
}
done := make(chan dialed, 1)
go func() {
c, err := dial(context.Background(), "ws://host.invalid/ws")
done <- dialed{conn: c, err: err}
}()
for i := 0; i < 200 && socketCount() == before; i++ {
time.Sleep(5 * time.Millisecond)
}
if socketCount() <= before {
t.Fatal("the reconnect reused the closed socket")
}
lastSocket(t).open()
select {
case result := <-done:
if result.err != nil {
t.Fatalf("reconnect dial: %v", result.err)
}
defer func() { _ = result.conn.close() }()
case <-time.After(2 * time.Second):
t.Fatal("the reconnect never completed its handshake")
}
}
TestBrowserSocketReleasesEverythingOnClose
Closing a connection must leave nothing attached: no Go callbacks on the
socket and no js.Func values still registered with the runtime.
Parameters
func TestBrowserSocketReleasesEverythingOnClose(t *testing.T)
{
conn, socket := dialFake(t)
if !socket.handlersBound() {
t.Fatal("the socket had no callbacks bound while open")
}
if err := conn.close(); err != nil {
t.Fatalf("close: %v", err)
}
if socket.handlersBound() {
t.Fatal("callbacks stayed attached after close")
}
if state := conn.socket.ReadyState(); state != js.SocketClosed {
t.Fatalf("socket readyState after close = %d, want closed", state)
}
// A second close is what a reconnect after an error does; it must not
// panic or re-release the Go functions.
if err := conn.close(); err != nil {
t.Fatalf("second close: %v", err)
}
}
dial
dial selects the transport stamped into rfw_config.js. Browser WebSocket is
the compatibility default. A configured native transport never falls back:
silently doing so would restore the cookie boundary the native transport was
selected to avoid.
Parameters
Returns
func dial(ctx context.Context, url string) (*hostConn, error)
{
switch sscTransport() {
case sscTransportBrowser:
mode := preferredHostTransport()
if mode == hostTransportStreamBus || mode == hostTransportAuto {
if connection, err := dialStreamBus(ctx, hostStreamBusURL()); err == nil {
return connection, nil
} else if debug {
fmt.Printf("hostclient: StreamBus unavailable, falling back to WebSocket: %v\n", err)
}
}
return dialBrowser(ctx, url)
case sscTransportCapacitor:
return dialCapacitor(ctx, url)
default:
return nil, fmt.Errorf("hostclient: unsupported SSC transport %q", sscTransport())
}
}
sscTransport
Returns
func sscTransport() string
{
configured := js.Get("RFW_SSC_TRANSPORT")
if !configured.Truthy() {
return sscTransportBrowser
}
return strings.TrimSpace(configured.String())
}
dialCapacitor
dialCapacitor connects through the RFW Capacitor plugin. Authentication
cookies and the Origin header are owned by the native URLSession; Go/WASM
receives protocol frames and lifecycle events, never cookie values.
Parameters
Returns
func dialCapacitor(ctx context.Context, url string) (*hostConn, error)
{
plugin, err := capacitorSSCPlugin()
if err != nil {
return nil, err
}
c := newHostConn()
id := fmt.Sprintf("rfw-ssc-%d", nativeConnectionSequence.Add(1))
var callback js.Func
var releaseOnce sync.Once
releaseCallback := func() { releaseOnce.Do(callback.Release) }
releaseAfterCallback := func() { time.AfterFunc(time.Millisecond, releaseCallback) }
callback = js.SafeFuncOf(func(_ js.Value, args []js.Value) any {
if len(args) > 1 && args[1].Truthy() {
c.fail(errors.New("hostclient: native SSC callback failed"))
releaseAfterCallback()
return nil
}
if len(args) == 0 || !args[0].Truthy() {
c.fail(errors.New("hostclient: native SSC callback returned no event"))
releaseAfterCallback()
return nil
}
event := args[0]
switch event.Get("type").String() {
case "open":
c.openOnce.Do(func() { close(c.open) })
case "message":
payload, decodeErr := decodeNativeFrame(event)
if decodeErr != nil {
c.fail(decodeErr)
return nil
}
c.deliver(payload)
case "error":
c.fail(nativeEventError(event, "native SSC transport failed"))
releaseAfterCallback()
case "close":
c.fail(fmt.Errorf(
"hostclient: native SSC connection closed with code %d: %s",
event.Get("code").Int(), event.Get("reason").String(),
))
releaseAfterCallback()
default:
c.fail(errors.New("hostclient: native SSC transport returned an unknown event"))
releaseAfterCallback()
}
return nil
})
// Native close crosses the Capacitor bridge asynchronously. Keep the Go
// callback alive long enough for the final close/error event; a terminal
// callback releases it sooner through the same once guard.
c.release = func() { time.AfterFunc(time.Second, releaseCallback) }
c.send = func(data string) error {
options := js.NewDict()
options.Set("id", id)
options.Set("data", data)
return callNativePlugin(plugin, "send", options.Value)
}
c.shutdown = func() error {
options := js.NewDict()
options.Set("id", id)
return callNativePlugin(plugin, "close", options.Value)
}
options := js.NewDict()
options.Set("id", id)
options.Set("url", url)
if err := callNativePlugin(plugin, "connect", options.Value, callback); err != nil {
_ = c.close()
return nil, err
}
select {
case <-c.open:
return c, nil
case <-c.done:
_ = c.close()
return nil, c.err()
case <-ctx.Done():
_ = c.close()
return nil, ctx.Err()
}
}
capacitorSSCPlugin
Returns
func capacitorSSCPlugin() (js.Value, error)
{
capacitor := js.Get("Capacitor")
if !capacitor.Truthy() {
return js.Undefined(), errors.New("hostclient: Capacitor is unavailable for native SSC transport")
}
plugins := capacitor.Get("Plugins")
if !plugins.Truthy() {
return js.Undefined(), errors.New("hostclient: Capacitor plugins are unavailable for native SSC transport")
}
plugin := plugins.Get(capacitorPluginName)
if !plugin.Truthy() {
return js.Undefined(), errors.New("hostclient: RFWSSC Capacitor plugin is not installed")
}
return plugin, nil
}
decodeNativeFrame
Parameters
Returns
func decodeNativeFrame(event js.Value) ([]byte, error)
{
data := event.Get("data").String()
switch event.Get("encoding").String() {
case "text":
return []byte(data), nil
case "base64":
payload, err := base64.StdEncoding.DecodeString(data)
if err != nil {
return nil, errors.New("hostclient: native SSC returned invalid base64")
}
return payload, nil
default:
return nil, errors.New("hostclient: native SSC returned an unknown frame encoding")
}
}
nativeEventError
Parameters
Returns
func nativeEventError(event js.Value, fallback string) error
{
message := strings.TrimSpace(event.Get("message").String())
if message == "" {
message = fallback
}
return errors.New("hostclient: " + message)
}
callNativePlugin
Parameters
Returns
func callNativePlugin(plugin js.Value, method string, args ...any) (err error)
{
defer func() {
if recovered := recover(); recovered != nil {
err = fmt.Errorf("hostclient: native SSC %s call failed", method)
}
}()
plugin.Call(method, args...)
return nil
}
TestReadDrainsBufferedFramesBeforeReportingTheClose
Parameters
func TestReadDrainsBufferedFramesBeforeReportingTheClose(t *testing.T)
{
c := newHostConn()
// Several frames, not one: a single frame is drained by the first
// non-blocking select and would pass even if the close raced the queue.
frames := []string{`{"sequence":1}`, `{"sequence":2}`, `{"sequence":3}`}
for _, frame := range frames {
c.deliver([]byte(frame))
}
c.fail(errors.New("host went away"))
for i, want := range frames {
frame, err := c.read(context.Background())
if err != nil {
t.Fatalf("read frame %d: %v", i, err)
}
if string(frame) != want {
t.Fatalf("frame %d = %s, want %s", i, frame, want)
}
}
if _, err := c.read(context.Background()); err == nil {
t.Fatal("read after the queue drained returned no error")
}
// The failure must stay readable: a reconnect asks again.
if _, err := c.read(context.Background()); err == nil {
t.Fatal("second read after close returned no error")
}
}
TestDeliveryCountsArrivalNotConsumption
Liveness counts frames as the browser delivers them, not as readLoop drains
them, so a slow handler cannot look like a dead connection.
Parameters
func TestDeliveryCountsArrivalNotConsumption(t *testing.T)
{
c := newHostConn()
c.deliver([]byte("{}"))
c.deliver([]byte("{}"))
if got := c.received.Load(); got != 2 {
t.Fatalf("received frames = %d, want 2", got)
}
if _, err := c.read(context.Background()); err != nil {
t.Fatalf("read frame: %v", err)
}
if got := c.received.Load(); got != 2 {
t.Fatalf("received frames after a read = %d, want 2", got)
}
}
TestInboundQueueOverflowClosesTheConnection
Parameters
func TestInboundQueueOverflowClosesTheConnection(t *testing.T)
{
c := newHostConn()
for i := 0; i < inboundFrames; i++ {
c.deliver([]byte("{}"))
}
c.deliver([]byte("{}"))
select {
case <-c.done:
if c.err() == nil {
t.Fatal("overflow reported a nil error")
}
default:
t.Fatal("overflow left the connection open")
}
}
TestHeartbeatFailsWithoutInboundTraffic
Parameters
func TestHeartbeatFailsWithoutInboundTraffic(t *testing.T)
{
c := newHostConn()
c.send = func(string) error { return nil }
if err := c.heartbeat(context.Background(), 5*time.Millisecond, 20*time.Millisecond); err == nil {
t.Fatal("heartbeat accepted a silent connection")
}
}
TestHeartbeatSurvivesAnsweredProbes
Parameters
func TestHeartbeatSurvivesAnsweredProbes(t *testing.T)
{
c := newHostConn()
c.send = func(string) error {
c.received.Add(1)
return nil
}
ctx, cancel := context.WithTimeout(context.Background(), 80*time.Millisecond)
defer cancel()
if err := c.heartbeat(ctx, 5*time.Millisecond, 10*time.Millisecond); !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("heartbeat error = %v, want a context deadline", err)
}
}
TestControlFrameIsUnsequencedAndNotRetained
Parameters
func TestControlFrameIsUnsequencedAndNotRetained(t *testing.T)
{
deliveryMu.Lock()
savedNext := nextOutbound
savedOutbox := outbox
savedInbound := lastInbound
nextOutbound = 7
outbox = map[uint64]message{}
lastInbound = 4
deliveryMu.Unlock()
defer func() {
deliveryMu.Lock()
nextOutbound = savedNext
outbox = savedOutbox
lastInbound = savedInbound
deliveryMu.Unlock()
}()
frames := make(chan string, 1)
c := newHostConn()
c.send = func(data string) error {
frames <- data
return nil
}
if err := sendControl(c, "ping"); err != nil {
t.Fatalf("send control frame: %v", err)
}
var frame wireMessage
if err := json.Unmarshal([]byte(<-frames), &frame); err != nil {
t.Fatalf("decode control frame: %v", err)
}
if frame.Control != "ping" {
t.Fatalf("control = %q, want ping", frame.Control)
}
if frame.Sequence != 0 {
t.Fatalf("control sequence = %d, want 0", frame.Sequence)
}
if frame.Ack != 4 {
t.Fatalf("control ack = %d, want 4", frame.Ack)
}
deliveryMu.Lock()
retained := len(outbox)
next := nextOutbound
deliveryMu.Unlock()
if retained != 0 {
t.Fatalf("control frame entered the outbox: %d retained", retained)
}
if next != 7 {
t.Fatalf("control frame consumed a sequence: nextOutbound = %d, want 7", next)
}
}
TestNormalizeWSURLDerivesTheScheme
The endpoint is derived from window.location unless the app configured one.
A page served over https must not downgrade its socket to ws.
Parameters
func TestNormalizeWSURLDerivesTheScheme(t *testing.T)
{
cases := map[string]string{
"ws://host.invalid/ws": "ws://host.invalid/ws",
"wss://host.invalid/ws": "wss://host.invalid/ws",
"http://host.invalid/ws": "ws://host.invalid/ws",
"https://host.invalid/ws": "wss://host.invalid/ws",
"http://host.invalid": "ws://host.invalid/ws",
"https://host.invalid": "wss://host.invalid/ws",
"https://host.invalid/": "wss://host.invalid/ws",
"https://host.invalid//": "wss://host.invalid/ws",
"wss://host.invalid/live": "wss://host.invalid/live",
" ws://host.invalid/ws ": "ws://host.invalid/ws",
"": "",
}
for raw, want := range cases {
if got := normalizeWSURL(raw); got != want {
t.Errorf("normalizeWSURL(%q) = %q, want %q", raw, got, want)
}
}
}
TestNormalizeWSURLFollowsThePageProtocol
A bare host takes the scheme of the page. wasmbrowsertest serves over plain
HTTP, so the socket must come out as ws and carry the default path.
Parameters
func TestNormalizeWSURLFollowsThePageProtocol(t *testing.T)
{
got := normalizeWSURL("host.invalid:8080")
protocol := js.Location().Get("protocol").String()
want := "wss://host.invalid:8080/ws"
if protocol == "http:" {
want = "ws://host.invalid:8080/ws"
}
if got != want {
t.Fatalf("normalizeWSURL on a %s page = %q, want %q", protocol, got, want)
}
}
init
Linking this package wires SSC component registration into core. Apps that
never use SSC never import hostclient, so the websocket and net/http stacks
stay out of their bundle.
func init()
{ core.SetHostRegistrar(RegisterComponentOwned) }