hostclient API

hostclient

package

API reference for the hostclient package.

F
function

installFakeCapacitorSSC

Parameters

Returns

hostclient/capacitor_transport_wasm_test.go:44-56
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")
}
F
function

TestSSCTransportDefaultsToBrowser

Parameters

hostclient/capacitor_transport_wasm_test.go:58-66
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)
	}
}
F
function

TestCapacitorTransportFailsClosedWhenPluginIsMissing

Parameters

hostclient/capacitor_transport_wasm_test.go:68-82
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)
	}
}
F
function

TestCapacitorTransportConnectsWritesReadsAndCloses

Parameters

hostclient/capacitor_transport_wasm_test.go:84-133
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)
	}
}
F
function

TestCapacitorTransportRejectsMalformedBinaryFrame

Parameters

hostclient/capacitor_transport_wasm_test.go:135-152
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)
	}
}
F
function

TestCapacitorTransportKeepsCallbackAliveThroughAsyncClose

Parameters

hostclient/capacitor_transport_wasm_test.go:154-173
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)
	}
}
F
function

TestCapacitorTransportDoesNotFallBackForUnknownMode

Parameters

hostclient/capacitor_transport_wasm_test.go:175-184
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)
	}
}
S
struct

fakeElement

hostclient/hostclient_test.go:8-13
type fakeElement struct

Methods

Exists
Method

Returns

bool
func (*fakeElement) Exists() bool
{ return e.exists }
Text
Method

Returns

string
func (*fakeElement) Text() string
{ return e.text }
SetText
Method

Parameters

v string
func (*fakeElement) SetText(v string)
{ e.text = v }
Attr
Method

Parameters

name string

Returns

string
func (*fakeElement) Attr(name string) string
{
	if name == hostExpectedAttr {
		return e.expected
	}
	if e.attrStore != nil {
		return e.attrStore[name]
	}
	return ""
}
SetAttr
Method

Parameters

name string
value string
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
S
struct

fakeRoot

hostclient/hostclient_test.go:42-45
type fakeRoot struct

Methods

HostVar
Method

Parameters

name string

Returns

func (*fakeRoot) HostVar(name string) hostVarElement
{
	if el, ok := r.elems[name]; ok {
		return el
	}
	return &fakeElement{}
}
SetHTML
Method

Parameters

html string
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
F
function

newFakeRoot

Returns

hostclient/hostclient_test.go:47-49
func newFakeRoot() *fakeRoot

{
	return &fakeRoot{elems: make(map[string]*fakeElement)}
}
F
function

TestHandleHostPayloadMismatchTriggersResync

Parameters

hostclient/hostclient_test.go:71-104
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")
		}
	}
}
F
function

TestLegacyExpectationRequiresResync

Parameters

hostclient/hostclient_test.go:106-117
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)
	}
}
F
function

TestInitSnapshotRecoveryAndUpdate

Parameters

hostclient/hostclient_test.go:119-145
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")
	}
}
S
struct

domComponentRoot

hostclient/hydration_dom.go:11-11
type domComponentRoot struct

Methods

HostVar
Method

Parameters

name string

Returns

func (domComponentRoot) HostVar(name string) hostVarElement
{
	selector := fmt.Sprintf(`[%s="%s"]`, hostVarAttr, name)
	return domHostVarElement{r.Query(selector)}
}
SetHTML
Method

Parameters

html string
func (domComponentRoot) SetHTML(html string)
{
	r.Element.SetHTML(html)
}
F
function

newComponentRoot

Parameters

Returns

hostclient/hydration_dom.go:13-15
func newComponentRoot(el dom.Element) componentRoot

{
	return domComponentRoot{el}
}
S
struct

domHostVarElement

hostclient/hydration_dom.go:26-26
type domHostVarElement struct

Methods

Exists
Method

Returns

bool
func (domHostVarElement) Exists() bool
{ return e.Truthy() }
Text
Method

Returns

string
func (domHostVarElement) Text() string
{ return e.Element.Text() }
SetText
Method

Parameters

value string
func (domHostVarElement) SetText(value string)
{ e.Element.SetText(value) }
Attr
Method

Parameters

name string

Returns

string
func (domHostVarElement) Attr(name string) string
{ return e.Element.Attr(name) }
SetAttr
Method

Parameters

name string
value string
func (domHostVarElement) SetAttr(name, value string)
{ e.Element.SetAttr(name, value) }
S
struct

componentBinding

hostclient/runtime_foundation.go:23-29
type componentBinding struct

Fields

Name Type Description
id string
vars []string
gate *deliveryGate
S
struct

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.

hostclient/runtime_foundation.go:38-38
type deliveryGate struct

Methods

open
Method

open reports whether the registration this gate belongs to is still the live one. The zero binding carries no gate and owns nothing.

Returns

bool
func (*deliveryGate) open() bool
{ return g != nil && !g.closed.Load() }
close
Method
func (*deliveryGate) close()
{
	if g != nil {
		g.closed.Store(true)
	}
}

Fields

Name Type Description
closed atomic.Bool
I
interface

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.

hostclient/runtime_foundation.go:52-54
type gatedHostSetter interface

Methods

Parameters

raw any
allow func() bool

Returns

bool
func SetFromHostGated(...)
I
interface

hostSetter

hostSetter is the plain host signal contract, without the gate.

hostclient/runtime_foundation.go:57-57
type hostSetter interface

Methods

SetFromHost
Method

Parameters

any
func SetFromHost(...)
I
interface

hostWriteBarrier

hostWriteBarrier waits out a gated write that was already allowed.

hostclient/runtime_foundation.go:60-60
type hostWriteBarrier interface

Methods

func HostWriteBarrier(...)
S
struct

message

hostclient/runtime_foundation.go:120-126
type message struct

Fields

Name Type Description
name string
action string
id string
payload any
sequence uint64
S
struct

wireMessage

hostclient/runtime_foundation.go:128-137
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"
T
type

messageWriter

hostclient/runtime_foundation.go:139-139
type messageWriter func(context.Context, *hostConn, wireMessage) error
S
struct

actionReply

hostclient/runtime_foundation.go:141-145
type actionReply struct

Fields

Name Type Description
payload any
err *ActionError
resetErr error
F
function

decodeInitSnapshotPayload

Parameters

raw
any
hostclient/runtime_foundation.go:151-175
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}
}
S
struct

ActionError

ActionError is a machine-readable error returned by a typed host action.

hostclient/runtime_foundation.go:178-182
type ActionError struct

Methods

Error
Method

Returns

string
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"
F
function

init

hostclient/runtime_foundation.go:191-203
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),
	)
}
F
function

connect

hostclient/runtime_foundation.go:205-214
func connect()

{
	once.Do(func() {
		go func() {
			for {
				js.Guard("host connection loop", connectionLoop)
				time.Sleep(time.Second)
			}
		}()
	})
}
F
function

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

string
hostclient/runtime_foundation.go:221-232
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)
}
F
function

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

raw
string

Returns

string
hostclient/runtime_foundation.go:237-263
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
}
F
function

connectionLoop

hostclient/runtime_foundation.go:265-382
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)
	}
}
F
function

guardedLoop

Parameters

context
string
fn
func() error

Returns

error
hostclient/runtime_foundation.go:384-390
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
}
F
function

heartbeatLoop

Parameters

Returns

error
hostclient/runtime_foundation.go:392-394
func heartbeatLoop(ctx context.Context, c *hostConn) error

{
	return c.heartbeat(ctx, heartbeatInterval, heartbeatTimeout)
}
F
function

readLoop

Parameters

Returns

error
hostclient/runtime_foundation.go:396-487
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)
			})
		}
	}
}
S
struct

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.

hostclient/runtime_foundation.go:491-494
type hostSignalUpdate struct

Fields

Name Type Description
name string
value any
F
function

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

component
string
payload
map[string]any
token
uint64
hostclient/runtime_foundation.go:499-515
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)
	}
}
F
function

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

component
string
payload
map[string]any
token
uint64
hostclient/runtime_foundation.go:523-557
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
}
F
function

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

component
string
token
uint64
updates
hostclient/runtime_foundation.go:566-593
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
		}
	}
}
F
function

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

signal
any
gate
value
any

Returns

bool
hostclient/runtime_foundation.go:599-617
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
}
F
function

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

id
string
hostclient/runtime_foundation.go:625-634
func fenceHostSignalWrites(id string)

{
	if id == "" {
		return
	}
	for _, signal := range dom.SnapshotComponentSignals(id) {
		if barrier, ok := signal.(hostWriteBarrier); ok {
			barrier.HostWriteBarrier()
		}
	}
}
F
function

hostSignalWriteHook

Returns

func(component,
name string)
hostclient/runtime_foundation.go:636-640
func hostSignalWriteHook() func(component, name string)

{
	mu.RLock()
	defer mu.RUnlock()
	return beforeHostSignalWrite
}
F
function

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

component
string
token
uint64

Returns

hostclient/runtime_foundation.go:645-653
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
}
F
function

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

id
string

Returns

hostclient/runtime_foundation.go:662-670
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)
}
F
function

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

id
string

Returns

string
bool
hostclient/runtime_foundation.go:675-681
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
}
F
function

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

id
string

Returns

hostclient/runtime_foundation.go:686-698
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()}
}
F
function

prepareInboundDelivery

Parameters

remoteSession
string
control
string
hostclient/runtime_foundation.go:700-714
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()
}
F
function

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

id
string
name
string
vars
[]string
hostclient/runtime_foundation.go:720-722
func RegisterComponent(id, name string, vars []string)

{
	registerComponent(id, name, vars)
}
F
function

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

id
string
name
string
vars
[]string

Returns

func()
hostclient/runtime_foundation.go:732-738
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) })
	}
}
F
function

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

id
string
name
string
vars
[]string

Returns

uint64
hostclient/runtime_foundation.go:743-763
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
}
F
function

releaseComponent

Parameters

name
string
token
uint64
hostclient/runtime_foundation.go:765-794
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}})
	}
}
F
function

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

name
string
payload
map[string]any
current
hostclient/runtime_foundation.go:809-828
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})
	}
}
F
function

isRegistrationControl

Parameters

payload
any

Returns

bool
hostclient/runtime_foundation.go:830-836
func isRegistrationControl(payload any) bool

{
	values, ok := payload.(map[string]any)
	if !ok {
		return false
	}
	return values["init"] == true || values["unsubscribe"] == true
}
F
function

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

name
string
hostclient/runtime_foundation.go:843-847
func EnableSendDedup(name string)

{
	mu.Lock()
	dedup[name] = struct{}{}
	mu.Unlock()
}
F
function

dedupEnabled

Parameters

name
string

Returns

bool
hostclient/runtime_foundation.go:849-854
func dedupEnabled(name string) bool

{
	mu.RLock()
	_, ok := dedup[name]
	mu.RUnlock()
	return ok
}
F
function

Send

Send queues or transmits a host component message.

Parameters

name
string
payload
any
hostclient/runtime_foundation.go:857-884
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})
}
F
function

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

name
string
h
func(map[string]any)

Returns

func()
hostclient/runtime_foundation.go:891-922
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}})
			}
		})
	}
}
F
function

SessionID

SessionID returns the current SSC session ID.

Returns

string
hostclient/runtime_foundation.go:925-929
func SessionID() string

{
	sessionMu.RLock()
	defer sessionMu.RUnlock()
	return sessionID
}
F
function

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.

hostclient/runtime_foundation.go:934-964
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()
	}
}
F
function

sendMessage

Parameters

msg
hostclient/runtime_foundation.go:966-968
func sendMessage(c *hostConn, msg message)

{
	sendMessageWithWriter(c, msg, writeMessage)
}
F
function

sendMessageWithWriter

Parameters

hostclient/runtime_foundation.go:970-974
func sendMessageWithWriter(c *hostConn, msg message, writer messageWriter)

{
	sendMu.Lock()
	defer sendMu.Unlock()
	sendMessageUnlockedWithWriter(c, msg, writer)
}
F
function

sendMessageUnlocked

Parameters

msg
hostclient/runtime_foundation.go:976-978
func sendMessageUnlocked(c *hostConn, msg message)

{
	sendMessageUnlockedWithWriter(c, msg, writeMessage)
}
F
function

sendMessageUnlockedWithWriter

Parameters

hostclient/runtime_foundation.go:980-1001
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)
}
F
function

writeMessage

Parameters

Returns

error
hostclient/runtime_foundation.go:1003-1005
func writeMessage(_ context.Context, c *hostConn, message wireMessage) error

{
	return c.writeJSON(message)
}
F
function

initMessageName

Parameters

msg

Returns

string
bool
hostclient/runtime_foundation.go:1007-1016
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
}
F
function

Call

Call invokes a typed SSC action and waits for its correlated response.

Parameters

action
string
request

Returns

error
hostclient/runtime_foundation.go:1019-1070
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()
	}
}
S
struct

FormResponse

FormResponse is the typed result returned by host.RegisterForm.

hostclient/runtime_foundation.go:1073-1077
type FormResponse struct

Fields

Name Type Description
Data Response json:"data,omitempty"
Fields map[string]string json:"fields,omitempty"
Valid bool json:"valid"
F
function

SubmitForm

SubmitForm invokes a typed SSC form action.

Parameters

action
string
values
Values

Returns

FormResponse[Response]
error
hostclient/runtime_foundation.go:1080-1082
func SubmitForm[Values, Response any](ctx context.Context, action string, values Values) (FormResponse[Response], error)

{
	return Call[Values, FormResponse[Response]](ctx, action, values)
}
F
function

EnableDebug

EnableDebug enables host client debug logging.

hostclient/runtime_foundation.go:1085-1085
func EnableDebug()

{ debug = true }
F
function

uniqueRuntimeTestName

Parameters

prefix
string

Returns

string
hostclient/runtime_wasm_test.go:18-20
func uniqueRuntimeTestName(prefix string) string

{
	return fmt.Sprintf("%s-%d", prefix, runtimeTestSequence.Add(1))
}
F
function

TestGuardedLoopConvertsPanicAndNextLoopRuns

Parameters

hostclient/runtime_wasm_test.go:22-37
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)
	}
}
F
function

TestInboundMessageLimitSupportsHydrationSnapshots

Parameters

hostclient/runtime_wasm_test.go:39-46
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)
	}
}
F
function

pendingCount

Returns

int
hostclient/runtime_wasm_test.go:48-52
func pendingCount() int

{
	mu.RLock()
	defer mu.RUnlock()
	return len(pending)
}
F
function

TestRegisterHandlerUnsubscribeQueuesWireUnsubscribe

Parameters

hostclient/runtime_wasm_test.go:54-85
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)
	}
}
F
function

TestStaleUnsubscribeDoesNotRemoveReplacementHandler

Parameters

hostclient/runtime_wasm_test.go:87-99
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()
}
F
function

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

hostclient/runtime_wasm_test.go:103-110
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)
	}
}
F
function

TestSendDedupOptIn

Dedup is opt-in per channel: after EnableSendDedup identical payloads within
the TTL window are dropped.

Parameters

hostclient/runtime_wasm_test.go:114-123
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)
	}
}
F
function

TestSendMessageSerializesSequenceAndWrite

Parameters

hostclient/runtime_wasm_test.go:125-181
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
}
F
function

TestPrepareInboundDeliveryResetsNewSessionState

Parameters

hostclient/runtime_wasm_test.go:183-216
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)
	}
}
F
function

TestPrepareInboundDeliveryResetsRejectedResume

Parameters

hostclient/runtime_wasm_test.go:218-248
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)
	}
}
F
function

TestResetSessionClearsIdentityDeliveryAndPendingCalls

Parameters

hostclient/runtime_wasm_test.go:250-362
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")
	}
}
F
function

TestOldGenerationFrameIsRejectedAfterSessionReset

Parameters

hostclient/runtime_wasm_test.go:364-398
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)
	}
}
F
function

bindingSnapshot

Parameters

name
string

Returns

hostclient/component_binding_wasm_test.go:22-27
func bindingSnapshot(name string) (componentBinding, bool)

{
	mu.RLock()
	defer mu.RUnlock()
	binding, bound := bindings[name]
	return binding, bound
}
F
function

pendingFor

Parameters

name
string

Returns

inits
int
unsubscribes
int
hostclient/component_binding_wasm_test.go:29-45
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
}
F
function

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

names
...string

Returns

[]string
hostclient/component_binding_wasm_test.go:49-64
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
}
F
function

pendingKind

Parameters

payload
any

Returns

string
hostclient/component_binding_wasm_test.go:66-81
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"
}
F
function

forgetPending

forgetPending drops the messages a test queued, so package state does not
leak into the next one.

Parameters

name
string
hostclient/component_binding_wasm_test.go:85-98
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()
	})
}
F
function

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.

Parameters

hostclient/component_binding_wasm_test.go:103-113
func setConnection(t *testing.T, c *hostConn)

{
	t.Helper()
	mu.Lock()
	conn = c
	mu.Unlock()
	t.Cleanup(func() {
		mu.Lock()
		conn = nil
		mu.Unlock()
	})
}
F
function

setSnapshotBarrier

setSnapshotBarrier installs the delivery barrier under the lock the read loop
reads it with.

Parameters

fn
func(component string)
hostclient/component_binding_wasm_test.go:117-121
func setSnapshotBarrier(fn func(component string))

{
	mu.Lock()
	afterBindingSnapshot = fn
	mu.Unlock()
}
F
function

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

hostclient/component_binding_wasm_test.go:126-144
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")
	}
}
F
function

TestRegisterComponentCleanupReleasesTheBinding

Parameters

hostclient/component_binding_wasm_test.go:146-172
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)
	}
}
F
function

TestStaleComponentCleanupKeepsTheReplacementBinding

Parameters

hostclient/component_binding_wasm_test.go:174-190
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")
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:195-223
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)
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:229-255
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)
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:260-296
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")
}
F
function

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

hostclient/component_binding_wasm_test.go:301-322
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)
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:327-358
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)
	}
}
F
function

hostVarRoot

hostVarRoot builds a component root carrying one host variable element.

Parameters

id
string

Returns

hostclient/component_binding_wasm_test.go:361-371
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
}
F
function

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

hostclient/component_binding_wasm_test.go:376-391
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
}
F
function

hostVarText

Parameters

Returns

string
hostclient/component_binding_wasm_test.go:393-395
func hostVarText(root dom.Element) string

{
	return root.Query(`[data-host-var="value"]`).Text()
}
F
function

waitForHostVar

Parameters

want
string
hostclient/component_binding_wasm_test.go:397-406
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)
}
F
function

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

hostclient/component_binding_wasm_test.go:412-445
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)
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:449-454
func withoutCSSEscape(t *testing.T)

{
	t.Helper()
	original := js.Get("CSS")
	js.Set("CSS", js.Null())
	t.Cleanup(func() { js.Set("CSS", original) })
}
F
function

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

hostclient/component_binding_wasm_test.go:460-496
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)
			}
		})
	}
}
F
function

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

hostclient/component_binding_wasm_test.go:502-543
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)
	}
}
T
type

ConnectionState

ConnectionState describes the SSC transport state.

hostclient/connection_state.go:7-7
type ConnectionState string
F
function

ConnectionStateSignal

ConnectionStateSignal returns the reactive SSC connection state.

hostclient/connection_state.go:23-25
func ConnectionStateSignal() *state.Signal[ConnectionState]

{
	return connectionState
}
F
function

setHostSignalWriteHook

setHostSignalWriteHook installs the delivery hook under the lock the read
loop reads it with.

Parameters

fn
func(component, name string)
hostclient/host_signal_delivery_wasm_test.go:16-20
func setHostSignalWriteHook(fn func(component, name string))

{
	mu.Lock()
	beforeHostSignalWrite = fn
	mu.Unlock()
}
F
function

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

id
string
names
...string

Returns

map[string]*state.Signal[string]
hostclient/host_signal_delivery_wasm_test.go:24-40
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
}
F
function

waitForSignal

Parameters

want
string
hostclient/host_signal_delivery_wasm_test.go:42-51
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)
}
F
function

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

signals
map[string]*state.Signal[string]

Returns

int
hostclient/host_signal_delivery_wasm_test.go:55-63
func untouched(signals map[string]*state.Signal[string]) int

{
	count := 0
	for _, signal := range signals {
		if signal.Get() == "initial" {
			count++
		}
	}
	return count
}
F
function

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

hostclient/host_signal_delivery_wasm_test.go:70-117
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)
	}
}
F
function

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

hostclient/host_signal_delivery_wasm_test.go:123-153
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")
	}
}
F
function

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

hostclient/host_signal_delivery_wasm_test.go:158-198
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)
	}
}
I
interface

hostVarElement

hostclient/hydration_shared.go:16-22
type hostVarElement interface

Methods

Exists
Method

Returns

bool
func Exists(...)
Text
Method

Returns

string
func Text(...)
SetText
Method

Parameters

string
func SetText(...)
Attr
Method

Parameters

string

Returns

string
func Attr(...)
SetAttr
Method

Parameters

string
string
func SetAttr(...)
I
interface

componentRoot

hostclient/hydration_shared.go:24-27
type componentRoot interface

Methods

HostVar
Method

Parameters

string

Returns

func HostVar(...)
SetHTML
Method

Parameters

string
func SetHTML(...)
S
struct

hydrationMismatch

hostclient/hydration_shared.go:29-35
type hydrationMismatch struct

Fields

Name Type Description
VarName string
Expected string
Actual string
ActualHash string
ExpectedAlg string
S
struct

initSnapshotPayload

hostclient/hydration_shared.go:37-40
type initSnapshotPayload struct

Fields

Name Type Description
HTML string
Vars []string
F
function

encodeExpectation

Parameters

value
string

Returns

string
hostclient/hydration_shared.go:42-45
func encodeExpectation(value string) string

{
	sum := sha256.Sum256([]byte(value))
	return fmt.Sprintf("%s:%s", expectationHashAlg, hex.EncodeToString(sum[:]))
}
F
function

expectationMatches

Parameters

expectedAttr
string
actual
string

Returns

bool
string
string
hostclient/hydration_shared.go:47-56
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
}
F
function

updateHostVar

Parameters

name
string
value
string
hostclient/hydration_shared.go:58-78
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
}
F
function

handleHostPayload

Parameters

payload
map[string]any
updateSignal
func(name string, raw any)

Returns

hostclient/hydration_shared.go:80-95
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
}
F
function

applyInitSnapshot

Parameters

hostclient/hydration_shared.go:97-102
func applyInitSnapshot(root componentRoot, payload *initSnapshotPayload)

{
	if payload == nil {
		return
	}
	root.SetHTML(payload.HTML)
}
F
function

buildResyncPayload

Parameters

mismatches

Returns

map[string]any
hostclient/hydration_shared.go:104-121
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,
		},
	}
}
F
function

preferredHostTransport

Returns

string
hostclient/streambus_transport_wasm.go:27-37
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
	}
}
F
function

hostStreamBusURL

Returns

string
hostclient/streambus_transport_wasm.go:39-44
func hostStreamBusURL() string

{
	if configured := stdjs.Global().Get("RFW_STREAMBUS_URL"); configured.Truthy() {
		return normalizeStreamBusURL(configured.String(), false)
	}
	return normalizeStreamBusURL(hostWSURL(), true)
}
F
function

normalizeStreamBusURL

Parameters

raw
string
incrementHTTPPort
bool

Returns

string
hostclient/streambus_transport_wasm.go:46-79
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()
}
F
function

dialStreamBus

Parameters

endpoint
string

Returns

error
hostclient/streambus_transport_wasm.go:81-130
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
}
S
struct

streamBusClientConfig

hostclient/streambus_transport_wasm.go:132-135
type streamBusClientConfig struct

Fields

Name Type Description
options stdjs.Value
port string
F
function

loadStreamBusConfig

Parameters

Returns

hostclient/streambus_transport_wasm.go:137-164
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
}
F
function

replaceURLPort

Parameters

raw
string
port
string

Returns

string
hostclient/streambus_transport_wasm.go:166-173
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()
}
F
function

writeStreamBusFrame

Parameters

writer
payload
[]byte

Returns

error
hostclient/streambus_transport_wasm.go:175-182
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)
}
F
function

readStreamBusFrame

Parameters

reader
maximum
int

Returns

[]byte
error
hostclient/streambus_transport_wasm.go:184-200
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
}
S
struct

jsReader

hostclient/streambus_transport_wasm.go:202-205
type jsReader struct

Methods

Read
Method

Parameters

payload []byte

Returns

int
error
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
S
struct

jsWriter

hostclient/streambus_transport_wasm.go:225-225
type jsWriter struct

Methods

write
Method

Parameters

payload []byte

Returns

error
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
S
struct

promiseResult

hostclient/streambus_transport_wasm.go:234-237
type promiseResult struct

Fields

Name Type Description
value stdjs.Value
err error
F
function

awaitPromise

Parameters

Returns

error
hostclient/streambus_transport_wasm.go:239-271
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()
	}
}
F
function

TestStreamBusURLUsesAdvertisedHTTP3Port

Parameters

hostclient/streambus_transport_wasm_test.go:10-15
func TestStreamBusURLUsesAdvertisedHTTP3Port(t *testing.T)

{
	got := replaceURLPort("https://localhost:8081/streambus", "8083")
	if got != "https://localhost:8083/streambus" {
		t.Fatalf("URL = %q", got)
	}
}
F
function

TestNormalizeStreamBusURLFromDevelopmentWebSocket

Parameters

hostclient/streambus_transport_wasm_test.go:17-22
func TestNormalizeStreamBusURLFromDevelopmentWebSocket(t *testing.T)

{
	got := normalizeStreamBusURL("ws://localhost:8080/ws", true)
	if got != "https://localhost:8081/streambus" {
		t.Fatalf("URL = %q", got)
	}
}
F
function

TestPreferredHostTransportReadsGeneratedConfig

Parameters

hostclient/streambus_transport_wasm_test.go:24-38
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)
	}
}
S
struct

hostConn

hostConn adapts the browser WebSocket to the blocking read and write calls
the connection loops expect.

hostclient/transport.go:35-56
type hostConn struct

Methods

deliver
Method

Parameters

payload []byte
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"))
	}
}
fail
Method

Parameters

err error
func (*hostConn) fail(err error)
{
	c.closeOnce.Do(func() {
		if err == nil {
			err = errConnectionClosed
		}
		c.closeErr.Store(&err)
		close(c.done)
	})
}
err
Method

err reports why the connection ended, once done is closed.

Returns

error
func (*hostConn) err() error
{
	if stored := c.closeErr.Load(); stored != nil {
		return *stored
	}
	return errConnectionClosed
}
read
Method

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

[]byte
error
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()
	}
}
writeJSON
Method

Parameters

payload any

Returns

error
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))
}
close
Method

Returns

error
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
Method

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

error
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
F
function

newHostConn

Returns

hostclient/transport.go:58-64
func newHostConn() *hostConn

{
	return &hostConn{
		messages: make(chan []byte, inboundFrames),
		open:     make(chan struct{}),
		done:     make(chan struct{}),
	}
}
F
function

dialBrowser

dialBrowser opens a connection and waits for the browser handshake to
complete. It remains the default SSC transport for existing applications.

Parameters

url
string

Returns

error
hostclient/transport.go:68-97
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()
	}
}
F
function

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

control
string

Returns

error
hostclient/transport.go:224-231
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})
}
S
struct

fakeSocket

hostclient/browser_socket_wasm_test.go:43-46
type fakeSocket struct

Methods

open
Method

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
Method

deliver fires onmessage with a text frame.

Parameters

frame string
func (fakeSocket) deliver(frame string)
{
	event := js.NewDict()
	event.Set("data", frame)
	f.value.Call("onmessage", event.Value)
}
serverClose
Method

serverClose fires onclose the way a host closing the connection does.

Parameters

code int
reason string
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
Method

sent returns every frame the client wrote to this socket.

Returns

[]string
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
Method

handlersBound reports whether the socket still has Go callbacks attached.

Returns

bool
func (fakeSocket) handlersBound() bool
{
	for _, event := range []string{"onopen", "onmessage", "onerror", "onclose"} {
		if f.value.Get(event).Truthy() {
			return true
		}
	}
	return false
}

Fields

Name Type Description
t *testing.T
value js.Value
F
function

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

hostclient/browser_socket_wasm_test.go:52-63
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")
	})
}
F
function

lastSocket

lastSocket returns the most recently constructed fake.

Parameters

Returns

hostclient/browser_socket_wasm_test.go:66-74
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)}
}
F
function

socketCount

Returns

int
hostclient/browser_socket_wasm_test.go:76-76
func socketCount() int

{ return js.Get("__fakeSockets").Get("length").Int() }
F
function

dialFake

dialFake opens a connection through the fake and completes the handshake.

Parameters

Returns

hostclient/browser_socket_wasm_test.go:123-152
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{}
	}
}
F
function

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

hostclient/browser_socket_wasm_test.go:156-166
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{}
}
F
function

resetDeliveryState

resetDeliveryState clears the package level delivery bookkeeping so tests do
not inherit sequence numbers or pending calls from each other.

Parameters

hostclient/browser_socket_wasm_test.go:170-192
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()
	})
}
F
function

readOnce

readOnce runs readLoop until it returns, so a test can assert on the error a
single frame produces.

Parameters

timeout

Returns

error
hostclient/browser_socket_wasm_test.go:196-200
func readOnce(c *hostConn, timeout time.Duration) error

{
	ctx, cancel := context.WithTimeout(context.Background(), timeout)
	defer cancel()
	return readLoop(ctx, c)
}
F
function

TestBrowserSocketCompletesTheHandshakeAndWrites

Parameters

hostclient/browser_socket_wasm_test.go:202-218
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)
	}
}
F
function

TestBrowserSocketDeliversAnActionReply

Parameters

hostclient/browser_socket_wasm_test.go:220-240
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")
	}
}
F
function

TestBrowserSocketDeliversASubscriptionPayload

Parameters

hostclient/browser_socket_wasm_test.go:242-266
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")
	}
}
F
function

TestBrowserSocketRejectsAMalformedFrame

Parameters

hostclient/browser_socket_wasm_test.go:268-286
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")
	}
}
F
function

TestBrowserSocketReportsAServerClosure

Parameters

hostclient/browser_socket_wasm_test.go:288-306
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")
	}
}
F
function

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

hostclient/browser_socket_wasm_test.go:310-329
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")
	}
}
F
function

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

hostclient/browser_socket_wasm_test.go:333-352
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)
	}
}
F
function

TestBrowserSocketReconnectsOnANewSocket

Parameters

hostclient/browser_socket_wasm_test.go:354-389
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")
	}
}
F
function

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

hostclient/browser_socket_wasm_test.go:393-412
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)
	}
}
F
function

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

url
string

Returns

error
hostclient/transport_selector.go:30-47
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())
	}
}
F
function

sscTransport

Returns

string
hostclient/transport_selector.go:49-55
func sscTransport() string

{
	configured := js.Get("RFW_SSC_TRANSPORT")
	if !configured.Truthy() {
		return sscTransportBrowser
	}
	return strings.TrimSpace(configured.String())
}
F
function

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

url
string

Returns

error
hostclient/transport_selector.go:60-143
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()
	}
}
F
function

capacitorSSCPlugin

Returns

error
hostclient/transport_selector.go:145-159
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
}
F
function

decodeNativeFrame

Parameters

event

Returns

[]byte
error
hostclient/transport_selector.go:161-175
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")
	}
}
F
function

nativeEventError

Parameters

event
fallback
string

Returns

error
hostclient/transport_selector.go:177-183
func nativeEventError(event js.Value, fallback string) error

{
	message := strings.TrimSpace(event.Get("message").String())
	if message == "" {
		message = fallback
	}
	return errors.New("hostclient: " + message)
}
F
function

callNativePlugin

Parameters

plugin
method
string
args
...any

Returns

err
error
hostclient/transport_selector.go:185-193
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
}
F
function

TestReadDrainsBufferedFramesBeforeReportingTheClose

Parameters

hostclient/transport_wasm_test.go:15-41
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")
	}
}
F
function

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

hostclient/transport_wasm_test.go:45-58
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)
	}
}
F
function

TestInboundQueueOverflowClosesTheConnection

Parameters

hostclient/transport_wasm_test.go:60-75
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")
	}
}
F
function

TestHeartbeatFailsWithoutInboundTraffic

Parameters

hostclient/transport_wasm_test.go:77-84
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")
	}
}
F
function

TestHeartbeatSurvivesAnsweredProbes

Parameters

hostclient/transport_wasm_test.go:86-98
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)
	}
}
F
function

TestControlFrameIsUnsequencedAndNotRetained

Parameters

hostclient/transport_wasm_test.go:100-151
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)
	}
}
F
function

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

hostclient/transport_wasm_test.go:155-174
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)
		}
	}
}
F
function

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

hostclient/transport_wasm_test.go:178-188
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)
	}
}
F
function

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.

hostclient/coreinit.go:10-10
func init()

{ core.SetHostRegistrar(RegisterComponentOwned) }