session

package
v0.2.2 Latest Latest
Warning

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

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

Documentation

Overview

Package session implements the service that maintains one WebSocket connection's widget state.

Actor, in actor.go, is the existing Go type for this service. Its parent in internal/wsx starts the service and waits for its Run method during cleanup. The service lives for one connection and processes events on one goroutine. All widgets on that connection share its event loop and transport. Asynchronous effects may add goroutines.

The loop serializes reducer and renderer access to its state values. Pointers inside those values may still alias objects elsewhere: shared mutable objects require synchronization or exclusive ownership by agreement. The loop selects over three typed, bounded inputs: events and effect results, client acknowledgements, and heartbeat ticks. Browser events pass authorization before entering the mailbox.

One step

A step stamps the event with wall time and any generated identifiers, so reducers never call a clock or a random source; reduces under a panic guard; marks the fragments the transition may have touched; renders those fragments; emits one patch through the outbound validation boundary and the framer; and finally hands any effects to the actor boundary, which runs them in goroutines scoped to the session context and returns their results as ordinary events. Effects never run inside a reducer.

Backpressure

Patches in flight are bounded by an acknowledgement window that retains metadata, never frame bytes. When the window fills, the actor keeps reducing but stops rendering, then renders once from current state when an acknowledgement re-opens it. Because rendering is pure, the skipped frames were never needed — memory under a slow client is proportional to the number of fragments, not to the number of pending patches. Sustained pressure reaches the application as a synthesized event rather than as a transport call, so the reducer stays deterministic and replayable.

Panics

Effect goroutines start in the session's runtime scope through one helper that installs recovery, a counter, and shutdown wait-group registration. The transport owns the loop separately. At shutdown the actor joins every effect before the teardown hook runs; the drain timeout only reports an effect that overran it. Repeated panics at one site close that connection while other connections continue serving.

Status

Implemented: the actor and its three bounded inputs, the mailbox and its pool, the acknowledgement window, the coalescing flush, the effect boundary, the panic guard, the rate budgets, and the provenance record.

Index

Constants

View Source
const (
	EffectFailedSourceField    = "source"
	EffectFailedErrorField     = "error"
	EffectFailedRetryableField = "retryable"
)

The fields an EffectFailedEvent carries: which effect failed, what it said, and whether the failure was classified as transient.

View Source
const EffectFailedEvent = "gotth.effect_failed"

EffectFailedEvent is the name of the event a failed effect is turned into, so that a reducer sees a deterministic failure rather than silence.

It is spelled again in package live, which is where an application can reach it; live's own suite asserts the two are equal, because a reducer that matches on the wrong string handles nothing and looks like it handles something.

MaxCoalesceFlushAt is the largest CoalesceFlushAt a session can honour.

It is below the schema ceiling, and the difference is not a safety margin — it is arithmetic. The derivation below is stated in ONE set of terms, which it previously was not: the prose described the second term as "MaxEventContributing plus the scheduledBy edge, bounded together" (65) while the formula used MaxEventContributing (64) and carried the library's edge in a separate "+ 1". Both readings land on 959 only because the two spare terms cancel, and a derivation that is right by cancellation is one edit from being wrong (U-3).

The flush trigger is evaluated against the union the frame will actually carry (Actor.unionReaches), so a transition is deferred only while that union is strictly below CoalesceFlushAt — at most F - 1 identifiers. The next emission adds at most:

1 identifier    the deferred transition's own event id. unionEdges excludes
                an origin's own id from the union it is the origin of, so it
                is not in the F - 1 above; takePending folds it in one
                emission later.
1 identifier    the scheduledBy edge this library prepends to an effect's
                emission (effects.go, scheduledEdge).
64 identifiers  MaxEventContributing, the largest Contributing an
                application may put on the event being emitted. The emit
                path refuses more, which is what D-18 closed.

so the widest frame carries (F - 1) + 1 + 1 + MaxEventContributing = F + 1 + MaxEventContributing, and F may go up to CoalesceFlushCeiling - 1 - MaxEventContributing. At the maximum that is exactly 1024, which checkListBounds accepts because it refuses only a list LONGER than the bound. The headroom is one element, and backpressure_test.go's "carries exactly the widest union the schema permits" is what drives it there — a margin no test touches is a margin the next edit spends.

The "- 1" alone is D-14's constant, and this is that constant with the term D-14 did not have: before D-18 the library was the only contributor, the application's half was effectively 0, and 1023 was correct. It stopped being correct the moment an application was allowed to contribute, which is what C-31(b) is about — a legal per-event bound plus a legal CoalesceFlushAt could still overflow at flush, because the trigger was reading a proxy (len(pendingIDs)) and not the frame.

Measured over the D-14 repro, 4,000 unacknowledged transitions and a resync, with the library the only contributor:

CoalesceFlushAt   largest union on the wire   provenance
           512                          512   3,978 of 3,978 carried
           959                          959   3,982 of 3,982 carried

The union column used to read one higher at each row, because the old trigger counted the deferred set and the frame carried one more than it counted. Counting the frame removed the discrepancy rather than the identifier: the same provenance reaches the wire, one flush earlier.

View Source
const MaxEventContributing = 64

MaxEventContributing is the largest Event.Contributing an application may put on one event it emits through an Emitter. It is enforced at the emit path, before the event reaches a mailbox, so an over-long list is a deterministic failure of that effect rather than a frame that dies later on the actor goroutine (D-18).

64 is H-4's bound on every other repeated field in the schema — Event.fields, Patch.updates, Snapshot.updates. Origin.contributing_event_ids is 1024 not because one event may name a thousand causes but because it is an accumulator this library fills by coalescing, and a per-event claim is an ordinary per-message list. The bound also has to be small for a second reason, which the arithmetic below makes explicit: every identifier an application may add to one event is an identifier the library may not coalesce, so this number is subtracted from the coalescing headroom.

It is documented in live.Event.Contributing's godoc and derived here. There is no exported constant and no Limits field: making it configurable would let an operator configure their way back into D-18.

View Source
const ResyncEventName = "gotth.resync"

ResyncEventName is the reserved name a resync request is authorized under.

A resync reaches the actor and triggers work proportional to the whole state, so it is not exempt from authorization the way the pure plumbing frames are; it is authorized as a distinguished event kind instead. An application that authorizes by name must permit this one or its clients cannot recover from a gap.

Variables

View Source
var ErrSessionClosing = errors.New(
	"the session is closing and will accept no further events: return from the effect rather than retrying, " +
		"because nothing this session emits from here reaches a reducer")

ErrSessionClosing is the sentinel behind the error returned to an effect that emits after the session has begun shutting down. It is wrapped exactly as ErrSessionSaturated is, and for the same reason.

View Source
var ErrSessionSaturated = errors.New(
	"the session mailbox is full: back off and emit again, or raise Config.Limits.MailboxDepth")

ErrSessionSaturated is the sentinel behind the error returned to an effect whose emitted event could not be accepted, so that the effect learns about backpressure instead of the event vanishing.

It is never returned bare, and that is FR-58's doing: an error handed to application code has to name the session it belongs to and the event that caused the work, and a package-level sentinel can name neither. [Actor.emitter] wraps it with both; this value is what survives the wrapping for errors.Is, and it carries the half of the message that is the same every time — why the event was dropped, and what to do about it.

Functions

func IsRetryable

func IsRetryable(err error) bool

IsRetryable reports whether a failure was explicitly classified as transient.

Unclassified is terminal, and that direction is the whole point. An effect may have committed externally before it failed, so re-running one nobody classified risks duplicating a side effect that already happened; retrying is a claim about idempotence, and only the code that performed the effect can make it. A failure never retried costs a change that does not happen. A failure retried blindly costs one that happens twice.

errors.As rather than a type assertion, so the mark survives being wrapped by helpers between the effect that set it and the actor that reads it. It is exported out of this package because live.IsRetryable is this function and not a second copy of it: a predicate an application can call and a predicate the library decides on must not be two implementations that can disagree.

Types

type Actor

type Actor[I IIdentity] struct {
	// contains filtered or unexported fields
}

Actor implements the service owning one connection's widget state and effects. Its Run method processes events on one goroutine, serializing reducers and renderers. Its parent starts Run and waits for it at cleanup. Shared mutable objects accessed elsewhere still require synchronization or ownership transfer.

It selects over three typed, bounded inputs: a mailbox of events, effect results and synthesized backpressure signals; a channel of client acknowledgements; and a heartbeat tick. Only the mailbox can reach a reducer, and exactly one function writes to it from the wire, which is what makes the per-event authorization hook impossible to route around.

func New

func New[I IIdentity](o Options[I]) *Actor[I]

New builds an actor. It allocates the mailbox, the acknowledgement channel and the window, which is the moment a session's per-connection memory comes into existence — after authentication, never before.

func (*Actor[I]) Close

func (a *Actor[I]) Close(code protocol.CloseCode, reason string)

Close ends the session with an enumerated code. It is safe to call from any goroutine and is idempotent.

func (*Actor[I]) Done

func (a *Actor[I]) Done() <-chan struct{}

Done reports a channel closed when the actor has finished.

func (*Actor[I]) ID

func (a *Actor[I]) ID() ID

ID returns the session's identifier.

func (*Actor[I]) Ingress

func (a *Actor[I]) Ingress(ctx context.Context, in protocol.IInbound)

Ingress is the only path from the wire into this actor.

It runs on the connection's read pump, not on the actor goroutine, and that placement is the security property: an event is rate limited, checked against the registered names, and authorized before it occupies a mailbox slot. A new client frame kind cannot skip the hook, because this switch is exhaustive over a sum type closed in another package and there is no other way into the mailbox from the wire.

The frames that do not pass through Authorize are exactly three, and each is accounted for: an acknowledgement and a heartbeat are transport plumbing that no reducer can observe, and client telemetry is a report about a patch this session already sent. None of them can reach application state.

func (*Actor[I]) Ready

func (a *Actor[I]) Ready(ctx context.Context) error

Ready blocks until the mount transition has emitted its snapshot. The read pump waits on it so that a client cannot have a frame accepted before the snapshot that establishes the sequence it must reference.

func (*Actor[I]) Run

func (a *Actor[I]) Run(ctx context.Context)

Run drives the session until its context is cancelled or the session closes. It returns when the actor goroutine is finished, having cancelled and joined every in-flight effect and run the teardown hook exactly once.

func (*Actor[I]) TrackedBytes

func (a *Actor[I]) TrackedBytes() int64

TrackedBytes is the exactly-sized cost of the structures this actor owns: the window, the two channel backing arrays, and the fragment hashes. It is what the per-session memory gauge reports, and it deliberately does not pretend to know the heap cost of application state, because Go has no per-goroutine heap attribution.

type DenyError

type DenyError struct {
	// Reason is operator-facing. The client is told a generic denial, because
	// the reason an event was refused describes the rule that refused it.
	Reason string
}

DenyError rejects one event without closing the connection.

func (*DenyError) Error

func (e *DenyError) Error() string

Error renders that operator-facing reason, for the log.

type Effect

type Effect[I IIdentity] struct {
	// Source names the effect for provenance and metrics. It is refined before
	// the effect runs, because it becomes half of an Origin.source.
	Source string

	// Run performs the effect, on the goroutine spawn starts for it.
	//
	// scheduledBy is the identifier of the event whose transition returned this
	// effect, or zero when the server started the transition itself. It is
	// handed over rather than re-derived because FR-58 requires every
	// library-produced error to name the causal identifier where one exists,
	// and the errors the public adapter raises against an emitted event are
	// raised before any identifier of their own is minted.
	//
	// A nil Run is a mistake rather than a no-op, and execute refuses it with a
	// failure event: an effect that never runs is a change that never happens.
	Run func(ctx context.Context, p Peer[I], scheduledBy uint64, emit Emit) error
}

Effect is one unit of I/O the actor performs at its boundary.

It mirrors the public live.Effect, which is a concrete struct by operator ruling of 2026-09-03, and it is a second declaration rather than an alias because Run is in THIS package's vocabulary: the public signature speaks of a live.Session and a live.Emitter, neither of which an internal package can name. live's own adapter re-expresses one as the other, once, and that is the only translation between the two.

type Emit

type Emit func(event Event) error

Emit injects an event into the session that spawned an effect.

type Event

type Event struct {
	// ID is the server-minted causal root. It is zero for the transitions the
	// server started on its own, where the origin source names the cause.
	ID uint64
	// ClientRef is the client's own correlation handle, echoed back so the
	// browser can match a patch to the interaction that caused it.
	ClientRef uint64
	// SeenServerSeq is the causation edge: the sequence number of the last
	// patch the user had applied when they acted.
	SeenServerSeq uint64

	// Name is the event name the application registered, such as
	// "cart.add". An unregistered name never reaches here: it is refused at
	// ingress and counted.
	Name string

	// FragmentID is the live region the interaction happened in, empty when
	// the server started the transition itself.
	FragmentID string

	// At is stamped at the actor boundary, not read by the reducer from a
	// clock. That is what lets an event log replay to a byte-identical result.
	At time.Time

	// Fields are the form values the interaction carried, in the order they
	// arrived.
	Fields []Field

	// Contributing lists events in this session whose state changes this
	// event carries. It is only ever set on an event an effect emits, where
	// the application is the only party that knows the edge — the library
	// knows which event scheduled an effect, but not which event caused the
	// value an asynchronous fan-out is delivering.
	Contributing []uint64
}

Event is one input to a reducer, already past the refinement boundary and past authorization.

ID and At are stamped at the actor boundary, never read by a reducer from a clock or a random source: that is what makes an event log replay to a byte-identical result.

type FatalDenyError

type FatalDenyError struct {
	// Reason is operator-facing, as for DenyError.
	Reason string
}

FatalDenyError rejects an event and closes the connection.

func (*FatalDenyError) Error

func (e *FatalDenyError) Error() string

Error renders the reason and says the connection is going with it, so a log line distinguishes this from the survivable denial without the reader having to know which type produced it.

type Field

type Field struct {
	// Key is the form field's name, as the browser sent it.
	Key string

	// Value is its value, already past the refinement boundary — length
	// bounded and valid UTF-8 — and past authorization, but otherwise
	// untrusted application input.
	Value string
}

Field is one form value carried by an event.

type IApp

type IApp[I IIdentity] interface {
	// Init produces the session's initial state and any startup effects. It
	// runs once, as the first transition, before the first snapshot.
	Init(ctx context.Context, p Peer[I]) (state any, effects []Effect[I], err error)

	// Authorize runs before the reducer for every event, at the single
	// mailbox ingress. It is the one method called from the read pump rather
	// than from the actor goroutine, because refusing an event before it
	// occupies a mailbox slot is the entire point of where it sits.
	Authorize(ctx context.Context, p Peer[I], ev Event) error

	// Reduce is the pure state transition.
	Reduce(state any, ev Event) (any, []Effect[I])

	// Teardown runs after the actor exits, with final state.
	Teardown(ctx context.Context, p Peer[I], state any)

	// Registry returns the application's fragments.
	Registry() *render.Registry

	// Registered reports whether an event name is declared. An unregistered
	// name is refused and counted, never dispatched and never ignored.
	Registered(name string) bool

	// StateComparable reports whether the application's state type may be
	// compared with == to decide whether a transition changed anything.
	//
	// It is false for a type Go cannot compare at all, and — the part that is
	// not obvious — for one Go compares only by IDENTITY: a pointer, map,
	// slice, channel, function or interface. Those are comparable in Go's
	// sense, so a naive type switch takes the fast path on them and asks "is
	// this the same object", which a reducer that mutates in place and returns
	// the same pointer answers "yes" to. That is the ordinary Go mistake the
	// purity rule exists to forbid, and answering it wrongly freezes
	// state_version and makes P4 false (BR-7).
	//
	// It is a method rather than a value the actor derives because deriving it
	// costs reflection over the state type, and the type does not change for
	// the life of an application: the answer is computed once, where the type
	// is still a type parameter and not an opaque any.
	StateComparable() bool
}

IApp is the application behaviour the actor drives.

It is an interface with one implementation, which this library otherwise refuses, and the reason it earns its place is the type parameter: the public package is generic over the application's state type and the actor holds state as an opaque value, so something has to be the seam where the type assertion happens exactly once. That seam is this interface, and every method of it is called only from the actor goroutine except Authorize.

type ID

type ID [16]byte

ID is a session identifier: sixteen random bytes minted by the server at handshake. It is server-minted so that untrusted input can never name another session, and sixteen bytes wide so a patch frame captured in isolation is resolvable.

func (ID) String

func (id ID) String() string

String returns the lower-case hex form used in logs and span attributes.

type IIdentity

type IIdentity interface {
	// Subject returns a stable, non-secret identifier used for logging and
	// per-identity session limits.
	Subject() string
}

IIdentity is what this package needs from an application's identity, and it appears only as a CONSTRAINT and in the admission bookkeeping that calls Subject(). Since 2026-09-03 the identity itself travels as a type parameter, so nothing here returns one or hands one back through a field.

type Limits

type Limits struct {
	// MaxInboundFrameBytes caps a decoded frame, applied to the connection
	// before any payload is allocated.
	MaxInboundFrameBytes int

	// MaxEventsPerSecond and EventBurst are the inbound event token bucket.
	MaxEventsPerSecond float64

	// EventBurst is that bucket's depth: how far a flurry of interactions may
	// run ahead of the sustained rate before the limiter refuses.
	EventBurst int

	// MailboxDepth bounds the actor's mailbox. It is also a memory parameter:
	// a Go buffered channel allocates its whole backing array at make time,
	// for the life of the channel, occupied or not. The mailbox holds
	// pointers for that reason.
	MailboxDepth int

	// AckChannelDepth bounds the acknowledgement channel. A full channel
	// drops, which is lossless because an acknowledgement is a cumulative
	// high-water mark: the next one supersedes the one dropped and the window
	// re-opens a round trip later.
	AckChannelDepth int

	// AckWindow is the number of unacknowledged patches allowed in flight.
	AckWindow int

	// CoalesceFlushAt is the size of the contributing-event union at which a
	// coalesced patch is emitted immediately rather than coalesced further.
	// It is the union the frame will carry, not a proxy for it: the trigger
	// counts what emitPatch is about to build, including the identifiers the
	// application contributed to the event being emitted.
	//
	// It is a flush trigger, so the schema's ceiling is unreachable and no
	// contributing event is ever dropped — but only while it stays at or below
	// MaxCoalesceFlushAt. Above that the frame the flush constructs is one the
	// schema refuses, and the trigger becomes the loss. The exported boundary
	// enforces the range; this package documents it.
	CoalesceFlushAt int

	// MinResyncInterval and ResyncBurst are the resync bucket, deliberately
	// independent of the event bucket. A resync is the one client frame that
	// triggers work proportional to the whole state, so sharing a budget with
	// ordinary events would make it an amplification vector.
	MinResyncInterval time.Duration

	// ResyncBurst is that bucket's depth. It is small deliberately: a client
	// that legitimately needs a snapshot needs one.
	ResyncBurst int

	// WriteDeadline bounds a single socket write. Exceeding it with a full
	// outbound window is what evicts a client rather than blocking the actor
	// behind it.
	WriteDeadline time.Duration

	// SlowClientGrace is how long a full window is tolerated before the
	// session is closed as a slow client. It is the gap between "stop
	// emitting" and "give up".
	SlowClientGrace time.Duration

	// HeartbeatInterval is the value announced to the client in the mount
	// snapshot: how often it should send a heartbeat.
	HeartbeatInterval time.Duration

	// HeartbeatTimeout is how long the server waits for one before treating
	// the connection as dead. It is deliberately larger than
	// HeartbeatInterval, because equality would close a session on one late
	// frame.
	HeartbeatTimeout time.Duration

	// IdleTimeout closes a session that has exchanged no application traffic,
	// heartbeats excluded — a tab left open forever is a session holding
	// memory for nobody.
	IdleTimeout time.Duration

	// EffectDrainTimeout is how long teardown waits for in-flight effects to
	// return before it reports the overrun. Teardown keeps waiting after it:
	// every effect is joined before the session ends, so an effect that
	// ignores its cancelled context holds the session's shutdown open.
	EffectDrainTimeout time.Duration

	// PanicBudget is how many times one site may panic in a session before
	// the session closes. Other sessions are unaffected either way.
	PanicBudget int
}

Limits are a session's resource bounds. Every zero field takes its documented default through Normalize, so a caller can set one field without restating the rest.

func DefaultLimits

func DefaultLimits() Limits

DefaultLimits returns the documented defaults.

func (Limits) Normalize

func (l Limits) Normalize() Limits

Normalize fills every zero field from the defaults and returns the result.

type Options

type Options[I IIdentity] struct {
	// Peer is the identity and session identifier this actor is bound to for
	// its whole life. Neither changes: a re-authentication is a new session.
	Peer Peer[I]

	// App is the application behaviour — mount, reduce, render, execute — as
	// the type-erased interface this package can hold without knowing the
	// state type.
	App IApp[I]

	// Limits are the resource bounds. Zero fields are filled by Normalize, so
	// a caller may set one and leave the rest.
	Limits Limits

	// Framer is the only way out to the socket: it validates and marshals, so
	// there is no path from this actor to the transport that skips the
	// outbound boundary.
	Framer *protocol.Framer

	// Close ends the connection with an enumerated code. It is the transport's
	// close, called from whichever goroutine notices, and it is idempotent.
	Close func(code protocol.CloseCode, reason string)

	// Metrics, Tracer and Logger are the instrumentation triple, and each may
	// be nil, which is the disabled configuration rather than a missing
	// dependency: every method on all three is nil-receiver safe.
	Metrics *obs.Metrics

	// Tracer starts the session and transition spans. Nil disables tracing.
	Tracer *obs.Tracer

	// Logger writes the event, patch and provenance records. Nil disables
	// logging, including the provenance stream.
	Logger *obs.Logger

	// Dev is developer mode (FR-23). Its whole effect is on the message of the
	// Error frame a contained panic produces: see Actor.devMessage. It must be
	// false in production.
	Dev bool

	// Now and Ticks are the actor's only sources of time. They are injected
	// so that a test can drive a thirty-minute idle timeout without waiting
	// thirty minutes, and so that nothing on the transition path reads a clock
	// the test cannot see.
	Now func() time.Time

	// Ticks drives the heartbeat, the idle timeout and the slow-client grace.
	// A test supplies its own channel and delivers a tick when it wants one;
	// nothing here calls time.NewTicker.
	Ticks <-chan time.Time

	// Scope is the runtime scope every effect goroutine starts in. The
	// transport passes the live UI service's sessions scope, so the service's
	// join covers every effect; the actor itself also joins its own effects
	// before Run returns. It is required.
	Scope *runtime.Scope
}

Options configure one session actor.

type Peer

type Peer[I IIdentity] struct {
	// ID is the sixteen server-minted bytes naming this session, carried on
	// every frame in both directions.
	ID ID

	// Identity is the authenticated principal, derived once from the upgrade
	// request, as the application's own type. It never changes for the life of
	// the session.
	Identity I
}

Peer is the immutable pair a session is bound to.

type RetryableError

type RetryableError struct {
	// Err is the failure being classified. It is never nil: Retryable(nil) is
	// nil rather than a marked nil.
	Err error
}

RetryableError marks an effect's failure as transient, so the failure event says so and a reducer can decide to schedule another attempt.

It carries no message of its own. The classification travels as its own field on the event, and prefixing the error text would put the same fact in two places — one of which a reducer would then be tempted to parse.

func (*RetryableError) Error

func (e *RetryableError) Error() string

Error is the wrapped error's message, unchanged. The classification is not prefixed onto it deliberately — it travels as its own event field, and a message that also carried it would invite a reducer to parse the string.

func (*RetryableError) Unwrap

func (e *RetryableError) Unwrap() error

Unwrap exposes the underlying failure, so the mark survives errors.Is and errors.As in both directions.

Jump to

Keyboard shortcuts

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