Documentation
¶
Overview ¶
Package runtime implements the nightseam.duplex/1 profile over a frames duplex connection (duplex): JSON text frames carrying requests, responses, events and cancellations. It never touches a WebSocket; Dial and Accept open one and hand it over as a connection, and a Peer over any other transport speaks the same profile byte for byte. It has no application authorization, replay, retries, or persistence policy.
Index ¶
- Constants
- Variables
- func CallWire(ctx context.Context, wire duplex.Wire, path []string, params, result any, ...) error
- func CheckIdentity(ctx context.Context, call func(context.Context, string, any, any) error, ...) error
- func EmitWire(ctx context.Context, wire duplex.Wire, path []string, data any, ...) error
- func ForwardWire(inbound, outbound duplex.Endpoint) (func(), error)
- func HandleWire(wire HandlerRegistry, path []string, handler WireHandler) (func(), error)
- func MarshalJSON(value any) ([]byte, error)
- func MarshalObject(order []string, members map[string]json.RawMessage) ([]byte, error)
- func MustTypeExpression(encoded string) any
- func NewHandler(options ServerOptions) (http.Handler, error)
- func NewWirePair(options Options) (left, right duplex.Endpoint, err error)
- func RegisterWire(wire HandlerRegistry, path []string, handlers WireHandlers) (func(), error)
- func RelayInvocationControl(message duplex.Message) error
- func Unpublished(err error) error
- func ValidateDrawnType(binding TypeBinding, member string, requireObject bool) error
- func ValidateUnicodeJSON(data []byte) error
- func WithMeta(ctx context.Context, meta Meta) context.Context
- func WithoutUnpublishedProof(err error) error
- type AdapterContext
- type Backpressure
- type ConnectionClosed
- type ConnectionOpened
- type DeclarationIdentity
- type DialOptions
- type Dispatcher
- func (d *Dispatcher) Close(code duplex.Code, reason string) error
- func (d *Dispatcher) Register(path []string, receiver duplex.Receiver) (func(), error)
- func (d *Dispatcher) RegisterPrefix(path []string, receiver duplex.Receiver) (func(), error)
- func (d *Dispatcher) Select(path []string) *SelectedEndpoint
- func (d *Dispatcher) Send(path []string, message duplex.Message) error
- type DispatcherOptions
- type Event
- type EventDelivered
- type EventEmitted
- type EventHandler
- type FrameReceived
- type FrameSent
- type Handler
- type HandlerPanic
- type HandlerRegistry
- type IdentityPreparation
- type Invocation
- type InvocationBodyHandle
- type InvocationCaptureHandle
- type InvocationLimits
- type Meta
- type Nullable
- type Observer
- type ObserverEvent
- type Of
- type Opaque
- type Optional
- type Options
- type Outcome
- type Peer
- func (p *Peer) Call(ctx context.Context, method string, params, result any) error
- func (p *Peer) Close() error
- func (p *Peer) Context() context.Context
- func (p *Peer) Done() <-chan struct{}
- func (p *Peer) Emit(ctx context.Context, event string, data any) error
- func (p *Peer) Err() error
- func (p *Peer) Handle(method string, handler Handler) error
- func (p *Peer) HandleEvent(name string, handler EventHandler) error
- func (p *Peer) MaxFrameBytes() int64
- func (p *Peer) Observe(event ObserverEvent)
- func (p *Peer) Observer() Observer
- func (p *Peer) OnEvent(listener func(context.Context, Event)) func()
- func (p *Peer) Role() Role
- func (p *Peer) Subprotocol() string
- func (p *Peer) Wire() duplex.Endpoint
- type Propagator
- type PublicError
- type Raw
- type RequestEnded
- type RequestStarted
- type Role
- type Schema
- func (s *Schema) Bind(types map[string]any, families map[string]*Schema) *Schema
- func (s *Schema) BoundDeclaration() (string, error)
- func (s *Schema) DeclarationDigest() (string, error)
- func (s *Schema) Fields(name string) []string
- func (s *Schema) MustWithDeclaration(document string) *Schema
- func (s *Schema) ValidateExpressionRaw(expression any, data []byte, at ...string) error
- func (s *Schema) ValidateRaw(name string, data []byte, at ...string) error
- func (s *Schema) ValidateValue(expression any, value any) error
- func (s *Schema) WithDeclaration(document string) (*Schema, error)
- type SelectedEndpoint
- type ServerOptions
- type Trace
- type TypeBinding
- type UnpublishedError
- type ValueAdapter
- type ValueEnvironment
- type WireCallOptions
- type WireEmitOptions
- type WireEventHandler
- type WireHandler
- type WireHandlers
Constants ¶
const ( // InvocationCapture claims one immutable routing decision for this // traversal. Its message carries the capture's control sink as the return // address. Admission is the capture; a refusal is an error from Send. InvocationCapture = "invocation.capture" // InvocationReady says the captured request has been delivered. A control // latched before this arrives is pushed to the sink now. InvocationReady = "invocation.ready" // InvocationRelease drops a capture whose traversal wants no more controls. InvocationRelease = "invocation.release" // InvocationBegin takes one execution lease. Admission is the lease. InvocationBegin = "invocation.begin" // InvocationDone reports that an executing body actually finished. InvocationDone = "invocation.done" // InvocationControl relays a cancellation into the invocation, which // latches it and pushes it to every ready capture exactly once. InvocationControl = "invocation.control" )
An admitted request's return capability is the invocation, presented as a Wire. The empty path carries its outcome, as it always has; these operations carry its lifecycle. They are ordinary events of the profile — a layer's own vocabulary, as `channel.` is the tunnel's — and a participant needs nothing of Nightseam's to speak them but the Wire it was already handed.
The capture or body a verb is about is one opaque segment after the operation, because a path is what addresses a thing. The verbs never reach a peer root and never cross a physical hop, so they take no built-in family and reserve no namespace there; a return capability's path space is the invocation's alone.
const DefaultInvocationBodies = 64
DefaultInvocationBodies bounds the execution leases one admitted invocation may take.
const DefaultInvocationCaptures = 64
DefaultInvocationCaptures bounds the captures one admitted invocation may take. It bounds traversal depth and shallow fan-out together, since a capture is taken once per routing boundary crossed.
const IdentityMethod = "identity.check"
IdentityMethod is the ordinary request used before interpreting a wire through a declaration. It does not participate in transport negotiation.
const Profile = "nightseam.duplex/1"
Profile names what this package speaks: JSON text frames carrying requests, responses, events and cancellations both ways over a connection of the seam. docs/wire/profile.md is its specification.
Variables ¶
var ( // ErrInvocationUnsupported is the refusal a participant receives from a // return capability that does not speak this vocabulary. A dispatcher // answers it explicitly rather than routing with weaker guarantees. ErrInvocationUnsupported = errors.New("return capability does not carry an invocation lifecycle") // ErrInvocationEnded refuses participation once the invocation retired. ErrInvocationEnded = errors.New("invocation no longer admits participation") // ErrInvocationLimit refuses participation beyond a bound. ErrInvocationLimit = errors.New("invocation participation limit reached") // ErrInvocationDuplicate refuses a participant identifier already in use. ErrInvocationDuplicate = errors.New("invocation participant already exists") )
var ( ErrClosed = errors.New("duplex connection closed") ErrBackpressure = errors.New("duplex consumer is stalled") )
The errors a peer ends with: ErrClosed when Close was called or the connection ended, ErrBackpressure when a queue stayed full past its write deadline — a consumer that does not drain is disconnected rather than allowed to hold the connection up. Both reach Err and every pending call.
Functions ¶
func CallWire ¶ added in v0.6.0
func CallWire(ctx context.Context, wire duplex.Wire, path []string, params, result any, options ...WireCallOptions) error
CallWire calls a relative operation through the peer's request primitive. Its local return address is independent of every other call's identifier.
func CheckIdentity ¶ added in v0.6.0
func CheckIdentity(ctx context.Context, call func(context.Context, string, any, any) error, expected DeclarationIdentity) error
CheckIdentity checks the remote declaration before the caller exposes its model. The supplied call observes its context; a context without a deadline receives a 30-second limit. Only method_not_found means absent identity. Refusals neither retry the exchange nor close the underlying carrier.
func EmitWire ¶ added in v0.6.0
func EmitWire(ctx context.Context, wire duplex.Wire, path []string, data any, options ...WireEmitOptions) error
EmitWire admits one event at a relative path. The return says only that the destination accepted it; processing and transport remain asynchronous.
func ForwardWire ¶ added in v0.6.0
func ForwardWire(inbound, outbound duplex.Endpoint) (func(), error)
ForwardWire joins two existing origins without allocating a peer or channel. Detach removes only the forwarding registrations; both wires remain owned by their callers. Each root remains responsible for ending its failed carrier.
func HandleWire ¶ added in v0.6.0
func HandleWire(wire HandlerRegistry, path []string, handler WireHandler) (func(), error)
HandleWire registers one relative operation. The receiver returns before running application code, and request cancellation uses its return address.
func MarshalJSON ¶ added in v0.5.0
MarshalJSON encodes a wire value without replacing malformed Unicode. Generated codecs use it before validation, while invalid UTF-8 in Go strings and unpaired surrogate escapes in raw encodings are still visible.
func MarshalObject ¶ added in v0.5.0
MarshalObject writes an object whose members are already encoded, in the order given. It is what a generated live codec builds its value with: a live value is assembled member by member, because each callable in it has to become a binding first, and a map would write the members in another order than the record declares them.
func MustTypeExpression ¶
MustTypeExpression decodes a type expression the generator wrote, and panics on one it cannot read, as MustSchema does: its caller is generated code passing a constant, for which a malformed expression is a generator bug and not a consumer's error. A hand-written expression is decoded with json.Unmarshal.
func NewHandler ¶
func NewHandler(options ServerOptions) (http.Handler, error)
NewHandler serves the profile at an HTTP endpoint: each request that passes CheckOrigin and Authenticate is upgraded to a WebSocket and becomes a server-role peer with the options given. The generated binding's NewHandler wraps this with the family's handlers installed.
func NewWirePair ¶ added in v0.6.0
NewWirePair constructs a bounded local carrier with two relative origins. Sending on either endpoint delivers to receivers on the other. It allocates no Peer and preserves structured frames and verified local request context. The limits, propagator and observer in options apply to both directions.
func RegisterWire ¶ added in v0.6.0
func RegisterWire(wire HandlerRegistry, path []string, handlers WireHandlers) (func(), error)
RegisterWire installs a single receiver for a declared method, event, or both. The one detach removes the group; an event-only path refuses requests.
func RelayInvocationControl ¶ added in v0.6.0
func RelayInvocationControl(message duplex.Message) error
RelayInvocationControl hands a cancellation to the invocation it names, which latches it and pushes it to the traversals that captured it. A router that receives a control frame relays it here rather than resolving a route of its own: the capture, not the current registration, decides where it goes.
func Unpublished ¶ added in v0.5.0
Unpublished marks an error only at a boundary that can prove its own payload was neither queued nor dispatched locally. It preserves a nil error. Never use it to classify a received error code or an uncertain transport outcome.
func ValidateDrawnType ¶ added in v0.6.0
func ValidateDrawnType(binding TypeBinding, member string, requireObject bool) error
ValidateDrawnType checks the original declaration selected by a supplied family interpretation. It does not normalize an alias into an allowed kind.
func ValidateUnicodeJSON ¶ added in v0.5.0
ValidateUnicodeJSON checks original JSON text for malformed Unicode, including strings a decoder would discard as duplicate members. It does not replace JSON syntax or envelope validation performed by its caller.
func WithMeta ¶
WithMeta says what the requests and events sent from ctx carry. The map is copied, so a later write to the caller's does not reach a frame already sent; a nil or empty map carries nothing. Keys under the reserved prefix are the profile's and are dropped rather than sent, since the peer at the far end refuses a frame carrying one.
func WithoutUnpublishedProof ¶ added in v0.5.0
WithoutUnpublishedProof preserves the ordinary cause of an error crossing a dispatch boundary, but removes evidence that belongs to a nested send attempt. Transport adapters and local implementations can return another call's error; that error cannot prove that the surrounding request was never published.
Types ¶
type AdapterContext ¶ added in v0.6.0
type AdapterContext struct {
Options Options
ValueEnvironment ValueEnvironment
}
AdapterContext carries runtime options used when constructing model wires.
type Backpressure ¶
Backpressure is a queue that could not take a frame: Queued is its depth, Deadline the write deadline the producer then waited out where it waited, and Stalled the peer giving up on the consumer and closing the connection.
func (Backpressure) ObserverEvent ¶
func (Backpressure) ObserverEvent()
type ConnectionClosed ¶
ConnectionClosed is the connection ending, once and whatever ended it. Local says this side ended it, which is every case but the one a peer can lay at the other's door: a close the remote sent, with its code and its reason.
func (ConnectionClosed) ObserverEvent ¶
func (ConnectionClosed) ObserverEvent()
type ConnectionOpened ¶
ConnectionOpened is the peer taking the connection over, before it has read or written anything on it.
func (ConnectionOpened) ObserverEvent ¶
func (ConnectionOpened) ObserverEvent()
type DeclarationIdentity ¶ added in v0.6.0
type DeclarationIdentity struct {
Path string `json:"path"`
Digest string `json:"digest,omitempty"`
}
DeclarationIdentity names a declaration and, when specified, its generated SHA-256 digest. An empty Digest is omitted from the exchange.
func CallableIdentity ¶ added in v0.6.0
func CallableIdentity(binding TypeBinding) (DeclarationIdentity, error)
CallableIdentity identifies a closed nominal callable interpretation. Its application path and digest come from the same canonical declaration graph; a nongeneric callable retains its declaring family's revision digest.
type DialOptions ¶
type DialOptions struct {
Options Options
HTTPHeader http.Header
HTTPClient *http.Client
// ConnectTimeout bounds the handshake alone, as TypeScript's
// connectTimeoutMs does, and takes the same default of 30 seconds where
// it is zero. A dial that has not become a connection by then is refused
// with the code connect_timeout and nothing is opened; the connection's
// own lifetime is the context's, as it is without one.
ConnectTimeout time.Duration
// Subprotocols are offered to the server in order of preference; the
// default offers none. A server that selects none leaves the connection
// with none and the profile is spoken over it either way — but a browser
// refuses a handshake whose offer went unselected, so a client that
// offers must be met by a server that selects (docs/wire/profile.md).
Subprotocols []string
}
DialOptions is what Dial opens a WebSocket with: the peer's Options, the headers and client of the HTTP upgrade, the bound on the handshake, and the subprotocols to offer.
type Dispatcher ¶ added in v0.6.0
type Dispatcher struct {
// contains filtered or unexported fields
}
Dispatcher owns one endpoint attachment and an explicit exact/longest-prefix routing policy. It captures each request's traversal on the invocation its return capability carries, through the public vocabulary alone, and refuses a request whose return capability carries none rather than routing it with weaker detach and cancellation guarantees. An opaque wrapper is therefore as good as a native endpoint: the lifecycle travels with the unchanged return capability, and nothing here recognizes a concrete type.
func NewDispatcher ¶ added in v0.6.0
func NewDispatcher(root duplex.Endpoint, options ...DispatcherOptions) (*Dispatcher, error)
func (*Dispatcher) Close ¶ added in v0.6.0
func (d *Dispatcher) Close(code duplex.Code, reason string) error
func (*Dispatcher) Register ¶ added in v0.6.0
func (d *Dispatcher) Register(path []string, receiver duplex.Receiver) (func(), error)
func (*Dispatcher) RegisterPrefix ¶ added in v0.6.0
func (d *Dispatcher) RegisterPrefix(path []string, receiver duplex.Receiver) (func(), error)
func (*Dispatcher) Select ¶ added in v0.6.0
func (d *Dispatcher) Select(path []string) *SelectedEndpoint
Select returns a receiving view of this shared dispatcher. The view owns its prefix route and never acquires closure authority over the root endpoint.
func (*Dispatcher) Send ¶ added in v0.6.0
func (d *Dispatcher) Send(path []string, message duplex.Message) error
type DispatcherOptions ¶ added in v0.6.0
type DispatcherOptions struct{ OwnEndpoint bool }
DispatcherOptions explicitly transfers closure authority for an endpoint the caller owns. Borrowed endpoints remain the default.
type Event ¶
type Event struct {
Name string `json:"event"`
Data json.RawMessage `json:"data"`
}
Event is one event of the profile as a handler or an emitter sees it: its name and its data.
type EventDelivered ¶
EventDelivered is an event reaching the application, before its handler and the listeners run.
func (EventDelivered) ObserverEvent ¶
func (EventDelivered) ObserverEvent()
type EventEmitted ¶
EventEmitted is the application emitting an event, before the frame carrying it is queued. Bytes is the size of the event's data, as EventDelivered's is.
func (EventEmitted) ObserverEvent ¶
func (EventEmitted) ObserverEvent()
type EventHandler ¶
type EventHandler func(context.Context, *Peer, json.RawMessage)
EventHandler takes one event's data; events have no answer. Handlers run one at a time in the order the events arrived.
type FrameReceived ¶
type FrameReceived struct {
At time.Time
Kind string
Name string
Bytes int
ID string
Trace Trace
Family string
}
FrameReceived is one frame the decoder accepted, before anything routed it. A frame the decoder refused closes the connection and is no frame at all.
func (FrameReceived) ObserverEvent ¶
func (FrameReceived) ObserverEvent()
type FrameSent ¶
type FrameSent struct {
At time.Time
Kind string
Name string
Bytes int
ID string
Trace Trace
Family string
}
FrameSent is one frame handed to the transport, which is as far as this peer carries it: told by the writer immediately before the bytes leave, so that a frame is observed sent before anything it draws can be received. Name is what the frame names — the method of a request, the name of an event — and a response and a cancel name nothing, a response's method being no member of the wire.
func (FrameSent) ObserverEvent ¶
func (FrameSent) ObserverEvent()
type Handler ¶
Handler answers one request: it takes the params as they arrived and returns the result, or an error — a *PublicError crosses the wire with its code; any other error reaches the caller as internal. The context is cancelled when the caller withdraws the request or its deadline passes, and carries the trace and the meta the frame brought.
func IdentityHandler ¶ added in v0.6.0
func IdentityHandler(expected DeclarationIdentity) (Handler, error)
IdentityHandler validates the supplied declaration once and answers identity requests without running model code. Register it before exposing that model.
type HandlerPanic ¶
HandlerPanic is a method handler that gave up. Value is the panic value as %v renders it and nothing else: what the handler was given is the handler's, and reaches no observer here.
func (HandlerPanic) ObserverEvent ¶
func (HandlerPanic) ObserverEvent()
type HandlerRegistry ¶ added in v0.6.0
type HandlerRegistry interface {
duplex.Wire
Register([]string, duplex.Receiver) (func(), error)
Close(duplex.Code, string) error
}
HandlerRegistry is the explicit registration capability used by generated bindings. Closing it releases its registrations, not its borrowed carrier.
type IdentityPreparation ¶ added in v0.6.0
type IdentityPreparation struct {
// contains filtered or unexported fields
}
IdentityPreparation installs model receivers before a carrier starts reading, checks its declaration after attachment, and releases dispatch after binding. It owns its registrations, never the source carrier or another receive queue.
func PrepareIdentity ¶ added in v0.6.0
func PrepareIdentity(wire duplex.Endpoint, expected DeclarationIdentity, options Options) (*IdentityPreparation, error)
PrepareIdentity is synchronous: callers install receivers through Wire before attaching their carrier, then call Check and bind the model before Ready. RequestTimeout bounds the whole preparation, including a factory never bound.
func (*IdentityPreparation) Check ¶ added in v0.6.0
func (p *IdentityPreparation) Check(ctx context.Context) error
Check makes the interpretation's first request. It may be started only once. A failed or cancelled check releases waiting deliveries without dispatching.
func (*IdentityPreparation) Close ¶ added in v0.6.0
func (p *IdentityPreparation) Close()
Close abandons only this interpretation, including its deferred deliveries. The source wire and unrelated registrations remain usable.
func (*IdentityPreparation) Ready ¶ added in v0.6.0
func (p *IdentityPreparation) Ready() error
Ready releases dispatch after a successful check and model binding.
func (*IdentityPreparation) Wire ¶ added in v0.6.0
func (p *IdentityPreparation) Wire() HandlerRegistry
Wire registers on the source directly and holds only model dispatch. Its outgoing sends require readiness; Check uses the ungated source separately.
type Invocation ¶ added in v0.6.0
type Invocation struct {
// contains filtered or unexported fields
}
Invocation is the lifecycle an admitting runtime keeps for one admitted request, and the answer its return capability gives to the vocabulary above. A runtime that is not Nightseam's composes it — or answers the same paths itself — and the same participants work against either.
func NewInvocation ¶ added in v0.6.0
func NewInvocation(limits InvocationLimits, onRetired func()) *Invocation
NewInvocation creates the lifecycle of one admitted request. onRetired runs once, outside every lifecycle lock, when no permitted future participation and no already admitted control can still need the captures.
func (*Invocation) Deliver ¶ added in v0.6.0
func (v *Invocation) Deliver(path []string, message duplex.Message) error
Deliver answers one operation of the invocation vocabulary. A return capability routes every nonempty path here; the empty path stays its own.
func (*Invocation) DispatchDone ¶ added in v0.6.0
func (v *Invocation) DispatchDone()
DispatchDone reports that the admitted request's own delivery has returned.
func (*Invocation) Retired ¶ added in v0.6.0
func (v *Invocation) Retired() bool
Retired reports whether the invocation has released its captures.
func (*Invocation) Settle ¶ added in v0.6.0
func (v *Invocation) Settle()
Settle fixes the outcome. It neither finishes a body nor drains a control.
type InvocationBodyHandle ¶ added in v0.6.0
type InvocationBodyHandle struct {
// contains filtered or unexported fields
}
InvocationBodyHandle is one execution lease of one admitted invocation.
func BeginInvocationBody ¶ added in v0.6.0
func BeginInvocationBody(message duplex.Message) (*InvocationBodyHandle, error)
BeginInvocationBody takes an execution lease for work this participant owns. The invocation does not retire while the lease is held, so an early answer to the caller — a deadline, a withdrawal — never retires an invocation whose body is still running.
func (*InvocationBodyHandle) Done ¶ added in v0.6.0
func (b *InvocationBodyHandle) Done()
Done reports that the body actually finished. It is idempotent and releases only the lease it took: a participant cannot finish another owner's work.
type InvocationCaptureHandle ¶ added in v0.6.0
type InvocationCaptureHandle struct {
// contains filtered or unexported fields
}
InvocationCaptureHandle is one immutable routing decision a participant took for one traversal of one admitted invocation. Repeated traversal of the same dispatcher takes a fresh handle, so no two traversals share a slot.
func CaptureInvocation ¶ added in v0.6.0
func CaptureInvocation(message duplex.Message, control func(duplex.Message)) (*InvocationCaptureHandle, error)
CaptureInvocation claims a routing decision for this traversal and supplies the sink the invocation pushes its cancellation to. The sink receives the control at most once, after Ready and never before.
A return capability that does not speak the vocabulary refuses, and the refusal is the caller's to answer: routing an invocation-aware request with weaker cancellation guarantees is exactly what this reports instead.
func (*InvocationCaptureHandle) Ready ¶ added in v0.6.0
func (c *InvocationCaptureHandle) Ready()
Ready says the captured request has been delivered. A cancellation latched while the capture was being installed reaches the sink now.
func (*InvocationCaptureHandle) Release ¶ added in v0.6.0
func (c *InvocationCaptureHandle) Release()
Release drops the capture. Ready and Release are each idempotent, and a release after Ready gives up only this traversal's remaining controls.
type InvocationLimits ¶ added in v0.6.0
type InvocationLimits struct{ Captures, Bodies int }
InvocationLimits bounds the total captures and execution leases of one admitted invocation. Both are totals, so neither depth nor fan-out can grow the state an invocation retains.
func DefaultInvocationLimits ¶ added in v0.6.0
func DefaultInvocationLimits() InvocationLimits
DefaultInvocationLimits are the bounds an admitting runtime uses when it states none of its own.
type Meta ¶
Meta is what a frame carries about a call rather than of it: a flat map of strings — a tenant, an idempotency key, a credential that is per request — which the profile carries verbatim and reads nothing into.
type Nullable ¶
Nullable distinguishes JSON null from a non-null value. Use Optional[Nullable[T]] for an optional nullable field.
func (Nullable[T]) MarshalJSON ¶
MarshalJSON writes the value, or null.
func (*Nullable[T]) UnmarshalJSON ¶
UnmarshalJSON reads null as Null and anything else as the value.
type Observer ¶
type Observer interface{ Observe(event ObserverEvent) }
Observer is what a peer tells about the traffic it carries. It emits and never aggregates, and it chooses no backend: an observer sees names, ids, sizes, durations, outcomes and close codes — never a payload. A params object reaches no observer by any path, and diagnostic logging is one observer among others; the runtime writes to no logger of its own.
Observe is called on whichever goroutine the event happened on — the reader, a handler, a caller — and so from several at once: an observer that keeps anything keeps it under a lock of its own, and one that blocks holds up the connection it is watching. One that panics panics alone: the peer recovers it, loses that event and carries on, a diagnostic being no reason for a connection to end on a goroutine the consumer does not own.
type ObserverEvent ¶
type ObserverEvent interface{ ObserverEvent() }
ObserverEvent is one thing a peer did, or one thing a layer running over a peer did: it is implemented by the events of this package, by those of the tunnel and of the live layer, which reach an observer through the peer they run over, and by nothing else. A type switch is the consumer's dispatch, and a later profile or a later layer may add a case to it. It is not named Event because an Event of this package is one event of the profile, which an observer only ever hears about.
type Of ¶
type Of[T any] interface{ Of() T }
Of is what an entry point of a package generic in a family holds its type arguments to. Every record and enum a family's protocol package declares returns the package's Tag from Of, and an entry point that takes a parameter S drawn at several types constrains each to Of[STag]: they are then types of one family, or the call does not compile, whether the caller spells them or the compiler infers them from the handler. Go has no associated types; this is the pairing they would have given.
type Optional ¶
Optional distinguishes an absent member from a present value. Generated Go fields use the omitzero JSON option; a standalone absent value encodes as null.
func (Optional[T]) MarshalJSON ¶
MarshalJSON writes the value, or null when absent; omitzero is what keeps an absent member off the wire.
func (*Optional[T]) UnmarshalJSON ¶
UnmarshalJSON reads a member that is present; an absent one never reaches it.
type Options ¶
type Options struct {
Handlers map[string]Handler
Events map[string]EventHandler
// Prepare runs on the peer once it is built and before it reads its
// first frame: what it installs — Handle, HandleEvent, a tunnel over the
// peer — is there before anything can arrive, so the other side's first
// request cannot be refused method_not_found by a peer whose handlers
// are still on their way. Install in Prepare, use in OnConnect, which
// runs on a peer that is already live. An error fails the construction:
// the peer never runs and the constructor answers with it.
Prepare func(*Peer) error
MaxConcurrentHandlers int
// MaxPendingRequests bounds the calls this peer may have outstanding at
// once; the one past it is refused busy without reaching the wire. It is
// the caller's own bound, as MaxConcurrentHandlers is the receiver's.
MaxPendingRequests int
QueueCapacity int
MaxFrameBytes int64
RequestTimeout time.Duration
WriteTimeout time.Duration
Propagator Propagator
// Observer is told what the peer does with the traffic it carries; nil
// observes nothing and costs nothing. Families labels a method or event
// name with the family it belongs to, which the generated install fills:
// an unlabelled name has no family, and the runtime parses none.
Observer Observer
Families map[string]string
}
Options limits are per connection. Zero values select the documented defaults. Handlers may run concurrently; event callbacks run serially in receive order.
type Outcome ¶
type Outcome int
Outcome is how a request ended.
const ( // OutcomeOK is a request answered with a result. OutcomeOK Outcome = iota // OutcomeErrored is a request answered with an error. OutcomeErrored // OutcomeCancelled is a request the caller withdrew before it was answered. OutcomeCancelled // OutcomeTimedOut is a request that outlived its deadline. OutcomeTimedOut )
type Peer ¶
type Peer struct {
// contains filtered or unexported fields
}
Peer owns a frames duplex connection until Close or transport failure. Never send on or receive from the connection after handing it to NewPeer. The connection's receive limit is its maker's to set to MaxFrameBytes; the peer refuses a larger frame it is nonetheless handed.
func Accept ¶
func Accept(w http.ResponseWriter, r *http.Request, options ServerOptions) (*Peer, error)
Accept upgrades an authenticated request to a WebSocket and speaks the profile over it. The caller must keep its HTTP handler alive until Peer.Done; NewHandler implements that lifetime contract.
func Dial ¶
Dial opens a WebSocket and speaks the profile over it. It uses ctx for the handshake and the connection lifetime. Use a separate context for each Call; cancelling the dialing context disconnects the peer.
func NewPeer ¶
NewPeer speaks the profile over any connection of the seam — a pipe, a tunnel channel, a socket already accepted — as the given role. The peer owns the connection from here and closes it when it ends; ctx ending ends the peer. Dial and Accept are this over a WebSocket.
func (*Peer) Call ¶
Call sends one request and waits for its response. It never retries. A caller cancellation also sends best-effort cancellation to the remote handler.
func (*Peer) Close ¶
Close ends the peer with ErrClosed, closing the connection beneath it with 1000 and failing every pending call. It is safe to call more than once.
func (*Peer) Context ¶
Context is the peer's own, cancelled when it ends: what a handler or a caller derives its own from to be released with the connection.
func (*Peer) Done ¶
func (p *Peer) Done() <-chan struct{}
Done is closed when the peer has ended, for whatever reason; Err says which.
func (*Peer) Emit ¶
Emit queues an event. Success means queued for this connection, not persisted or processed by the remote application.
func (*Peer) Err ¶
Err is why the peer ended, or nil while it runs: ErrClosed, ErrBackpressure, the context's error, or the connection's own.
func (*Peer) HandleEvent ¶
func (p *Peer) HandleEvent(name string, handler EventHandler) error
HandleEvent registers the handler for the event of that name, replacing any before it; an event with no handler is dropped.
func (*Peer) MaxFrameBytes ¶
MaxFrameBytes is the largest frame this peer sends or receives.
func (*Peer) Observe ¶
func (p *Peer) Observe(event ObserverEvent)
Observe tells this peer's observer one event, the runtime's own or one of a layer running over the peer; a peer given no observer does nothing. It is how a tunnel and a live scope observe — through the peer they run over, which is the observer they were never given a second way to choose.
func (*Peer) Observer ¶
Observer is what this peer was given, or nil. Whatever runs over a peer — a tunnel, a live scope — observes through this one rather than taking its own.
func (*Peer) OnEvent ¶
OnEvent observes every event and returns an idempotent unsubscribe function. Slow callbacks consume the bounded event queue and can disconnect the peer.
func (*Peer) Role ¶
Role is the side of the connection this peer is: it prefixes the request ids it mints, and a tunnel over it chooses channel ids by it.
func (*Peer) Subprotocol ¶
Subprotocol is what the WebSocket handshake beneath this peer selected, and "" when it selected none or the peer does not run over a WebSocket. It is fixed for the peer's life; the profile reads nothing into it.
type Propagator ¶
type Propagator interface {
Extract(ctx context.Context, trace Trace) context.Context
Inject(ctx context.Context) Trace
}
Propagator moves a trace between a context and a frame. Extract places an incoming frame's trace in the context its handler runs under; Inject says what an outgoing frame carries — a child of the trace the context holds, or a new trace where it holds none. Both members reach it verbatim, and a peer sends what it returns as it returns it.
var DefaultPropagator Propagator = w3cPropagator{}
DefaultPropagator is what a peer with no Options.Propagator propagates by: it keeps an incoming trace verbatim and mints an outgoing one's ids itself, so that frames correlate across hops with no tracing library installed.
type PublicError ¶
type PublicError struct {
Code string `json:"code"`
Message string `json:"message"`
Data json.RawMessage `json:"data,omitempty"`
}
PublicError is safe to send to the remote caller. Other handler errors are replaced by a generic internal error; their messages are not disclosed.
func (*PublicError) Error ¶
func (e *PublicError) Error() string
Error is the code and the message, as a log line shows them.
type Raw ¶
type Raw []byte
Raw fills a slot with JSON passed through unexamined, which is what a relay wants: it carries the bytes and satisfies Of, so a generic package instantiated with it compiles, and validates nothing of what it carries.
func (Raw) MarshalJSON ¶
func (*Raw) UnmarshalJSON ¶
type RequestEnded ¶
type RequestEnded struct {
At time.Time
ID string
Method string
Incoming bool
Duration time.Duration
Outcome Outcome
ErrorCode string
Trace Trace
Family string
}
RequestEnded pairs with every RequestStarted. Duration is the span between the two, and ErrorCode names the local cause: request_timeout for a deadline, cancelled for a cancellation, or the public error's code for a refusal. These codes hold for incoming and outgoing requests; success names no code.
func (RequestEnded) ObserverEvent ¶
func (RequestEnded) ObserverEvent()
type RequestStarted ¶
type RequestStarted struct {
At time.Time
ID string
Method string
Incoming bool
Trace Trace
Family string
}
RequestStarted is a request beginning, whichever side raised it: Incoming is one this peer serves, and a request it refuses for want of a method or of a slot begins and ends like any other.
func (RequestStarted) ObserverEvent ¶
func (RequestStarted) ObserverEvent()
type Role ¶
type Role string
Role is the side of a connection a peer takes. It decides the prefix of the request ids the peer mints — c: for a client, s: for a server — so the two sides never mint the same id, and a tunnel over the peer chooses its channel ids by it.
type Schema ¶
type Schema struct {
// contains filtered or unexported fields
}
Schema validates a family's descriptor: {"types": {...}, "parameters": [...]}. Its imports retain their own descriptors, so an application can fill type and family parameters without losing the caller's scope.
Expressions include primitives, named and imported types, family draws, arrays, maps, entity references, applications, literals, nullable values, inline shapes and the empty object. Records and adjacently tagged unions inherit their fields or variants. Field presence is separate from nullness.
Bind supplies parameters to raw validation. An explicitly applied argument is always validated. A generated Go generic codec may instead leave a family parameter unbound: the schema checks the enclosing shape and the instantiated Go codec checks that slot during marshal or unmarshal.
A refusal gives the location and the fact about it. Both runtimes read conformance/tables/validator.json, including the same refusal strings.
func MustSchema ¶
MustSchema is NewSchema for a description the generator wrote.
func NewSchema ¶
NewSchema reads a family's wire description and its generated digest; imported maps each family it refers to to that family's Schema. An empty digest leaves declaration identity unspecified.
func (*Schema) Bind ¶ added in v0.5.0
Bind returns a schema with type and family parameters filled. Type arguments are interpreted in the receiver's scope, before this binding.
func (*Schema) BoundDeclaration ¶ added in v0.6.0
BoundDeclaration returns this family's closed canonical application. Every declared parameter must be supplied; an incomplete interpretation never silently acquires the identity of its template.
func (*Schema) DeclarationDigest ¶ added in v0.6.0
DeclarationDigest identifies the closed family interpretation retained by Bind.
func (*Schema) MustWithDeclaration ¶ added in v0.6.0
MustWithDeclaration is WithDeclaration for a graph the generator wrote.
func (*Schema) ValidateExpressionRaw ¶
ValidateExpressionRaw verifies a type expression's value and rejects trailing values. at roots the diagnostic, as ValidateRaw's does.
func (*Schema) ValidateRaw ¶
ValidateRaw verifies a named type's value, including null and field presence. at roots the diagnostic where a family that imports this one holds the value; absent, the value is its own root.
func (*Schema) ValidateValue ¶
ValidateValue validates a typed value before publishing it on the wire.
func (*Schema) WithDeclaration ¶ added in v0.6.0
WithDeclaration returns a schema carrying the canonical declaration graph emitted beside its validator descriptor. It leaves the explicit digest and validation behavior unchanged. Bind retains this provenance.
type SelectedEndpoint ¶ added in v0.6.0
type SelectedEndpoint struct {
// contains filtered or unexported fields
}
func (*SelectedEndpoint) Close ¶ added in v0.6.0
func (s *SelectedEndpoint) Close(code duplex.Code, reason string) error
func (*SelectedEndpoint) Receive ¶ added in v0.6.0
func (s *SelectedEndpoint) Receive(receiver duplex.Receiver) (func(), error)
func (*SelectedEndpoint) Select ¶ added in v0.6.0
func (s *SelectedEndpoint) Select(path []string) *SelectedEndpoint
func (*SelectedEndpoint) Send ¶ added in v0.6.0
func (s *SelectedEndpoint) Send(path []string, message duplex.Message) error
type ServerOptions ¶
type ServerOptions struct {
Options Options
Authenticate func(*http.Request) (context.Context, error)
CheckOrigin func(*http.Request) bool
// OnConnect is called with a peer that is already live: it has read
// frames and may have answered them. It is where a server uses the peer
// — calls it, keeps it, waits on it. What a peer must serve is installed
// in Options.Prepare, which runs before it reads anything; a handler
// installed here can be too late for the other side's first request.
OnConnect func(*Peer)
// Subprotocols are what the server will select, in its own order of
// preference, from what a client offers; empty selects none, which is
// the default and what every consumer that sets nothing keeps. The
// profile names itself here nowhere and refuses nothing on this ground
// (docs/wire/profile.md).
Subprotocols []string
// SelectSubprotocol answers with the one subprotocol to select out of
// what this request offered, "" for none. It is the selection, not a
// filter over Subprotocols: a browser's ticket travels in the offer and
// is accepted only if it is selected back unchanged, which no fixed
// list can do. Nil selects the first offered that Subprotocols names.
SelectSubprotocol func(r *http.Request, offered []string) string
}
ServerOptions requires an explicit authentication and origin policy. The authenticated context is the base context for every incoming invocation.
type Trace ¶
type Trace struct{ Parent, State string }
Trace is the W3C Trace Context a frame carries: the two members verbatim, neither of them read by the runtime beyond the form the decoder holds a traceparent to. The runtime imports no tracing library; one binds to it as a Propagator and nowhere else.
type TypeBinding ¶ added in v0.5.0
TypeBinding retains the schema where a supplied type expression belongs. Generated WireType methods return it so aliases, literals and constraints survive use as Go type arguments, including inside another family.
func TypeArgument ¶ added in v0.5.0
func TypeArgument[T any]() TypeBinding
TypeArgument supplies an instantiated Go type without a consumer codec argument or registry. Discovery is lazy: recursive generated types may return bindings for their arguments without recursively describing them.
func (TypeBinding) Declaration ¶ added in v0.6.0
func (b TypeBinding) Declaration() (string, error)
Declaration selects the argument's reachable declaration content, retaining nested applications and the schemas in which their arguments were supplied.
type UnpublishedError ¶ added in v0.5.0
type UnpublishedError struct {
// contains filtered or unexported fields
}
UnpublishedError reports a local refusal before a frame entered the peer's outbound queue or a local implementation dispatched. Its underlying cause retains the ordinary public error or cancellation identity. A queued write failure or a remote response never supplies this proof, even if it has the same error code or message.
The proof belongs to this send attempt. A handler must not treat a nested call's refusal as proof that its own already-dispatched request was unsent.
func (*UnpublishedError) Error ¶ added in v0.5.0
func (e *UnpublishedError) Error() string
func (*UnpublishedError) Unwrap ¶ added in v0.5.0
func (e *UnpublishedError) Unwrap() error
Unwrap preserves errors.Is and errors.As for the original refusal.
type ValueAdapter ¶ added in v0.6.0
type ValueAdapter[T any] struct { Binding TypeBinding NeedsContext bool Export func(context.Context, T) (json.RawMessage, error) Import func(context.Context, json.RawMessage) (T, error) }
ValueAdapter pairs a declaration with its two primitive conversions. A reusable adapter retains no operation context or lifetime. When NeedsContext is true, conversion runs inside the supplied ValueEnvironment's matching Export, Import or Publish boundary, using the active context it provides.
func JSONAdapter ¶ added in v0.6.0
func JSONAdapter[T any]() ValueAdapter[T]
JSONAdapter validates ordinary data in both directions. It neither consults the context nor requires a conversion environment.
type ValueEnvironment ¶ added in v0.6.0
type ValueEnvironment interface {
Select(context.Context) (context.Context, error)
Child(context.Context) (context.Context, error)
Export(context.Context, func(context.Context) (json.RawMessage, error)) (json.RawMessage, error)
Import(context.Context, func(context.Context) error) error
Publish(context.Context, func(context.Context) (json.RawMessage, error), func(json.RawMessage) (json.RawMessage, error)) (json.RawMessage, error)
}
ValueEnvironment supplies operation-local conversion effects. It belongs to the consumer's chosen component; the runtime never interprets its context. Select chooses an outgoing context, and Child derives an incoming lifetime. Build callbacks are synchronous and must pass their active context through every nested conversion. Publish completes the build before attempting send.
type WireCallOptions ¶ added in v0.6.0
type WireCallOptions struct {
Propagator Propagator
RequestTimeout time.Duration
Observer Observer
Family string
}
WireCallOptions configures one model call, independently of its carrier. A zero RequestTimeout uses the runtime default; Propagator nil uses DefaultPropagator.
type WireEmitOptions ¶ added in v0.6.0
type WireEmitOptions struct {
Propagator Propagator
Observer Observer
Family string
}
WireEmitOptions configures one model event emission, independently of its carrier.
type WireEventHandler ¶ added in v0.6.0
type WireEventHandler func(context.Context, json.RawMessage) error
WireEventHandler receives an event body beside its carried context.
type WireHandler ¶ added in v0.6.0
WireHandler is a typed adapter's decoded request body, independent of the concrete carrier. The runtime supplies cancellation and response routing.
type WireHandlers ¶ added in v0.6.0
type WireHandlers struct {
Request WireHandler
Event WireEventHandler
Observer Observer
Family string
}
WireHandlers groups a method and event that share one declared name.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package slogobserver writes what a peer tells an observer to a slog.Logger: one record per event, at a level per kind, with the event's fields as attributes.
|
Package slogobserver writes what a peer tells an observer to a slog.Logger: one record per event, at a level per kind, with the event's fields as attributes. |