experimental

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Sep 25, 2026 License: MIT Imports: 31 Imported by: 0

README

Experimental Radius transport

This package provides the opt-in Radius relay library from the pinned Pi sources experimental/radius-auth.ts and experimental/radius-relay.ts. It does not register commands, enable experimental mode, or select a hosted endpoint.

Selection and authentication

Supply the gateway explicitly to NewRadiusRelayAuthResolver. Supply either a token, a token file, or a lazy CreateRuntime callback. Configure that runtime without model-network refresh and with OAuth refresh directed to the selected gateway. Do not use the built-in Earendil Radius runtime as a substitute for a PiG-owned gateway (D64). The existing RequestAuthRuntime supports an explicit radius provider configuration with oauth: "radius" and the selected baseUrl. TestRadiusAuthRefreshesFiveMinuteCredentialAtLocalGateway exercises that composition against a local gateway.

The resolver checks context cancellation before PI_OFFLINE presence. Offline mode performs no token-file or runtime access. The resolver trims explicit tokens using JavaScript whitespace semantics and reads a token file anew for every attempt. The resolver creates a stored-auth runtime once and resolves credentials anew for every attempt. Stored OAuth requires five minutes of remaining validity, including after refresh.

Ownership and async contract

RadiusRelayHost.Start owns one reconnect loop. RadiusRelayHost.Close cancels and joins opening, serving, pending writes, and retry waits. The host accepts independent virtual connections in arrival order. Pong and unknown or rejected connection-close replies reserve their order and bytes immediately without waiting for socket writes. One owned writer drains control and data submissions in that order while the reader continues delivering inbound data, closes, and opens. Control-write failures close the socket with the transport-error code and report the failure on the reader. Shutdown rejects queued writes and joins the active write and all completion callbacks. A local virtual close writes its non-nil final chunk before its close control. A remote close or host failure reports one terminal callback per active connection in insertion order.

CreateRadiusClientTransportFactory awaits authentication and the WebSocket handshake. The returned transport delivers raw binary messages without the host multiplexing envelope. Send copies each submission and awaits its ordered write. Pending submissions share the upstream four-frame byte budget. Native socket writes provide backpressure instead of browser bufferedAmount polling. Close initiates local shutdown without synthesizing a remote callback. Wait on Done to join transport cleanup.

Callbacks execute on the transport reader, not the TUI loop. Do not wait for that reader's shutdown from its own callback. Marshal any UI mutation through the UI owner's executor.

RadiusClientReconnect observes an established client. It retries after one second, doubles failed-attempt delays to thirty seconds, and restores the last selected Session after reconnect succeeds. A disconnected attachment reset preserves the selection. An explicit detach while connected clears it. A failed reattachment disconnects the client and retries. Dispose removes listeners, cancels pending work, disconnects, and joins the retry loop.

Evidence

Run go test -race -count=3 ./internal/experimental. The suite includes the relevant pinned upstream relay cases, malformed controls and envelopes, token rotation, local OAuth refresh, a native WebSocket echo gateway, ordered bidirectional backpressure, reconnect selection, and cancellation/cleanup tests. TestPinnedUpstreamRadiusSource executes the actual pinned TypeScript modules with local dependency adapters and rejects implicit endpoints or live sockets. The compiler adapter uses the repository's existing TypeScript installation and installs nothing.

Run go test ./internal/experimental -run '^$' -bench 'Benchmark(RelayDataFrame|RadiusClientSend|RadiusHostControl)$' -benchmem to measure framing, client-send, and host-control allocation costs. These benchmarks do not claim a speedup or complete CLI/server parity. The server's accept adapter and the established client's reconnect adapter remain caller-owned library boundaries.

Documentation

Overview

Package experimental implements the opt-in local process transport used by Pi's experimental server.

Package experimental provides opt-in transports for the experimental Session runtime. It does not register CLI commands or select network endpoints.

Index

Constants

View Source
const CoordinatorProtocolVersion = 3

CoordinatorProtocolVersion is the exact upstream-owned coordinator wire version.

View Source
const EnvRadiusGateway = codingagent.EnvRadiusGateway

EnvRadiusGateway is the upstream experimental gateway override.

View Source
const InternalProcessEnv = "__PI_INTERNAL_SPAWN"
View Source
const MaxControlLineBytes = 128 * 1024 * 1024
View Source
const RadiusRelayClientSubprotocol = "pi-session-relay.client.v1"

RadiusRelayClientSubprotocol is the upstream raw client WebSocket protocol.

View Source
const RadiusRelayHostSubprotocol = "pi-session-relay.host.v1"

RadiusRelayHostSubprotocol is the upstream multiplexed host WebSocket protocol.

Variables

This section is empty.

Functions

func EncodeControlLine

func EncodeControlLine(message any) (string, error)

EncodeControlLine encodes one value with JSON.stringify number, string and object-key semantics and a newline, enforcing Pi's UTF-8 byte limit including the delimiter.

func EncodeRelayDataFrame

func EncodeRelayDataFrame(connectionID string, payload []byte) ([]byte, error)

EncodeRelayDataFrame copies payload into the Radius version/type/UUID envelope.

func RunCoordinatorEntry

func RunCoordinatorEntry(ctx context.Context, args []string) error

RunCoordinatorEntry validates and consumes the internal role before starting the coordinator. An opt-in executable calls this entry explicitly; importing the package never changes CLI dispatch.

func RunCoordinatorProcess

func RunCoordinatorProcess(ctx context.Context, args []string) error

RunCoordinatorProcess serves the stable endpoint until shutdown, cancellation or the empty grace period. Unlike Node's process-owned event loop, Go explicitly joins the owned socket work before returning.

func TerminateInternalProcess

func TerminateInternalProcess(child *InternalProcess) error

TerminateInternalProcess forces a child to exit and waits until it can no longer take ownership. A nonzero exit status is expected after a kill and is available separately through Wait.

Types

type AuthInput

type AuthInput struct {
	Type  string
	Token string
	Path  string
}

AuthInput selects an explicit token or a token file. Nil selects stored credentials.

type CoordinatorConnection

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

CoordinatorConnection is the server-side endpoint of the coordinator's opaque message router. Event callbacks run in receive order, outside the state lock. They may send, close or unsubscribe.

func NewCoordinatorConnection

func NewCoordinatorConnection(options CoordinatorConnectionOptions) *CoordinatorConnection

NewCoordinatorConnection creates an unconnected server generation with a random ID when none is supplied.

func (*CoordinatorConnection) Broadcast

func (c *CoordinatorConnection) Broadcast(payload any) error

Broadcast writes a payload to every connected peer after registration.

func (*CoordinatorConnection) Close

func (c *CoordinatorConnection) Close()

Close destroys the socket and clears membership without marking an intentional close as replacement.

func (*CoordinatorConnection) Connect

func (c *CoordinatorConnection) Connect(ctx context.Context) error

Connect waits for a validated registration. Cancellation closes the pending connection.

func (*CoordinatorConnection) ControlPath

func (c *CoordinatorConnection) ControlPath() string

func (*CoordinatorConnection) OnEvent

func (c *CoordinatorConnection) OnEvent(listener *CoordinatorConnectionListener) func()

OnEvent adds a listener reference once, in insertion order. Each cleanup removes that reference, including any later re-registration.

func (*CoordinatorConnection) PeerIDs

func (c *CoordinatorConnection) PeerIDs() []string

PeerIDs returns a detached snapshot in the coordinator's insertion order.

func (*CoordinatorConnection) Replaced

func (c *CoordinatorConnection) Replaced() <-chan struct{}

Replaced closes once on replacement or an unexpected disconnect, but not on an explicit Close.

func (*CoordinatorConnection) Send

func (c *CoordinatorConnection) Send(peerID string, payload any) error

Send writes a payload to a named peer after registration and waits for the socket write.

func (*CoordinatorConnection) ServerConnectionID

func (c *CoordinatorConnection) ServerConnectionID() string

func (*CoordinatorConnection) WasReplaced

func (c *CoordinatorConnection) WasReplaced() bool

type CoordinatorConnectionEvent

type CoordinatorConnectionEvent struct {
	Type    string
	PeerID  string
	From    string
	Payload json.RawMessage
}

CoordinatorConnectionEvent carries peer membership changes or an opaque routed payload. A nil Payload means absent; json.RawMessage("null") is an explicit JSON null.

type CoordinatorConnectionListener

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

CoordinatorConnectionListener is a callback with stable reference identity. Go functions cannot be compared; OnEvent identifies listeners by this pointer instead.

func NewCoordinatorConnectionListener

func NewCoordinatorConnectionListener(listener func(CoordinatorConnectionEvent)) *CoordinatorConnectionListener

NewCoordinatorConnectionListener creates a listener reference. Reuse it to register the same callback again.

type CoordinatorConnectionOptions

type CoordinatorConnectionOptions struct {
	ControlPath        string
	Endpoint           string
	ServerConnectionID *string
}

CoordinatorConnectionOptions identifies this server generation and its private forwarding endpoint.

type CoordinatorStartupLease

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

CoordinatorStartupLease keeps an unregistered control connection alive while a server starts.

func EnsureCoordinator

func EnsureCoordinator(ctx context.Context, publicPath, controlPath string, options ...InternalProcessSpawnOptions) (*CoordinatorStartupLease, error)

EnsureCoordinator connects to an existing router or launches one and waits for its control endpoint. The optional native spawn options select an explicit opt-in entry executable, without changing Stock CLI dispatch.

func (*CoordinatorStartupLease) Close

func (lease *CoordinatorStartupLease) Close()

Close releases startup demand. It does not stop a coordinator that has other demand.

type InternalProcess

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

InternalProcess owns the child's exit observation and reaps it exactly once. A detached child may outlive its launcher. Done closes only after Wait has reaped it.

func SpawnInternalProcess

func SpawnInternalProcess(role InternalProcessRole, args []string, options InternalProcessSpawnOptions) (*InternalProcess, error)

SpawnInternalProcess starts a detached native role with ignored stdio and the current working directory. The caller must select an executable that explicitly dispatches internal roles. Stock CLI dispatch is unchanged.

func (*InternalProcess) Done

func (p *InternalProcess) Done() <-chan struct{}

Done closes when the child can no longer take ownership of a socket or session.

func (*InternalProcess) PID

func (p *InternalProcess) PID() int

PID returns the spawned process ID.

func (*InternalProcess) ProcessState

func (p *InternalProcess) ProcessState() *os.ProcessState

ProcessState returns nil while the process runs and its final state after it exits.

func (*InternalProcess) Wait

func (p *InternalProcess) Wait() error

Wait joins exit observation and returns the child's exit error, if any.

type InternalProcessRole

type InternalProcessRole string

InternalProcessRole selects an internal entrypoint, not a public CLI command.

func ConsumeInternalProcessRole

func ConsumeInternalProcessRole() (InternalProcessRole, error)

ConsumeInternalProcessRole validates before removing the role so descendants cannot inherit it.

func GetInternalProcessRole

func GetInternalProcessRole() (InternalProcessRole, error)

GetInternalProcessRole reads and validates the role without consuming it.

type InternalProcessSpawnOptions

type InternalProcessSpawnOptions struct {
	EntryPath string
	Env       map[string]string
}

InternalProcessSpawnOptions supplies an optional native entry executable and environment overrides. EntryPath is the native counterpart of a source module URL; an empty path re-executes this executable.

type RadiusClientAttachment

type RadiusClientAttachment struct{ SessionID string }

RadiusClientAttachment identifies the currently selected Session.

type RadiusClientByteTransport

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

RadiusClientByteTransport carries raw binary messages and reports one remote terminal event. Close initiates local shutdown; Done joins reads, writes, and cancellation cleanup.

func (*RadiusClientByteTransport) Close

func (t *RadiusClientByteTransport) Close()

Close initiates shutdown without emitting a remote terminal callback. It is safe in a data callback.

func (*RadiusClientByteTransport) Done

func (t *RadiusClientByteTransport) Done() <-chan struct{}

Done closes when every operation owned by this transport has finished.

func (*RadiusClientByteTransport) Send

func (t *RadiusClientByteTransport) Send(chunk []byte) error

Send copies and writes a chunk, in submission order, with bounded pending bytes.

type RadiusClientReconnect

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

RadiusClientReconnect reconnects an established client and restores its last selected Session. Dispose removes listeners, cancels, disconnects, and joins the retry loop.

func NewRadiusClientReconnect

func NewRadiusClientReconnect(ctx context.Context, client RadiusReconnectClient, reattach func(context.Context, string) error) *RadiusClientReconnect

NewRadiusClientReconnect observes an already established client; it does not initiate a connection until a disconnection event.

func (*RadiusClientReconnect) Dispose

func (r *RadiusClientReconnect) Dispose()

Dispose is idempotent and joins all retry and reattachment work.

type RadiusClientTransportFactory

type RadiusClientTransportFactory func(context.Context, RelayByteConnectionHandler) (*RadiusClientByteTransport, error)

RadiusClientTransportFactory resolves fresh credentials and connects one client byte stream. The context owns the resulting transport until Close.

func CreateRadiusClientTransportFactory

func CreateRadiusClientTransportFactory(options RadiusClientTransportOptions) RadiusClientTransportFactory

CreateRadiusClientTransportFactory creates an inert connection factory.

type RadiusClientTransportOptions

type RadiusClientTransportOptions struct {
	ServerID         string
	Auth             *RadiusRelayAuthResolver
	WebSocketFactory RadiusRelayWebSocketFactory
}

RadiusClientTransportOptions selects one authenticated raw client relay.

type RadiusReconnectClient

type RadiusReconnectClient interface {
	Attachment() *RadiusClientAttachment
	Connected() bool
	ConnectionState() string
	Disconnect(reason string)
	OnAttachmentChange(func(*RadiusClientAttachment)) func()
	OnConnectionStateChange(func(string)) func()
	Reconnect(context.Context) error
}

RadiusReconnectClient is the established-client contract used by the reconnect owner. Listener removers must be safe during notification; methods and listeners may run concurrently.

type RadiusRelayAuth

type RadiusRelayAuth struct {
	Gateway string
	Token   string
}

RadiusRelayAuth is one connection attempt's resolved credential.

type RadiusRelayAuthOptions

type RadiusRelayAuthOptions struct {
	Input         *AuthInput
	Gateway       string
	CreateRuntime func() (*codingagent.RequestAuthRuntime, error)
}

RadiusRelayAuthOptions requires an explicit gateway. CreateRuntime supplies the stored-auth runtime without model-network refresh; it is called once, lazily, only without Input. Its owner must configure OAuth refresh for its selected gateway.

type RadiusRelayAuthResolver

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

RadiusRelayAuthResolver rereads explicit or stored credentials for every connection attempt.

func NewRadiusRelayAuthResolver

func NewRadiusRelayAuthResolver(options RadiusRelayAuthOptions) (*RadiusRelayAuthResolver, error)

NewRadiusRelayAuthResolver creates an inert resolver. Gateway and the stored-auth runtime are caller-owned so construction cannot select a hosted service.

func (*RadiusRelayAuthResolver) Gateway

func (r *RadiusRelayAuthResolver) Gateway() string

Gateway returns the normalized gateway selected by the caller.

func (*RadiusRelayAuthResolver) Resolve

func (r *RadiusRelayAuthResolver) Resolve(ctx context.Context, required bool) (*RadiusRelayAuth, error)

Resolve checks cancellation, then PI_OFFLINE presence, before reading credentials. Required turns missing auth into an actionable error. Stored OAuth must have at least five minutes of validity.

type RadiusRelayHost

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

RadiusRelayHost owns a reconnecting, multiplexed host connection. Incoming controls enqueue ordered replies without waiting for socket writes. Close cancels and joins every owned operation.

func NewRadiusRelayHost

func NewRadiusRelayHost(options RadiusRelayHostOptions) (*RadiusRelayHost, error)

NewRadiusRelayHost validates the collaborators without starting work.

func (*RadiusRelayHost) Close

func (h *RadiusRelayHost) Close()

Close is idempotent and joins shutdown, including a pending authentication or opening attempt.

func (*RadiusRelayHost) Start

func (h *RadiusRelayHost) Start(ctx context.Context)

Start begins the owned reconnect loop once. The context owns its lifetime.

type RadiusRelayHostOptions

type RadiusRelayHostOptions struct {
	ServerID         string
	Accept           func(*RelayServerByteConnection) RelayByteConnectionHandler
	Auth             *RadiusRelayAuthResolver
	WebSocketFactory RadiusRelayWebSocketFactory
	OnStatus         func(RadiusRelayHostStatus)
}

RadiusRelayHostOptions binds the relay to a server's accept operation. No command or default endpoint is activated by construction.

type RadiusRelayHostStatus

type RadiusRelayHostStatus struct {
	Status string
	Error  string
}

RadiusRelayHostStatus reports authentication, connection, and retry transitions.

type RadiusRelayWebSocket

type RadiusRelayWebSocket interface {
	Protocol() string
	Read() (binary bool, data []byte, err error)
	Send(binary bool, data []byte) error
	Close(code int, reason string) error
}

RadiusRelayWebSocket is a connected transport. Read has one owner; Send is serialized by the relay. Close must unblock both and may run concurrently with them.

type RadiusRelayWebSocketFactory

type RadiusRelayWebSocketFactory func(context.Context, RadiusRelayWebSocketOptions) (RadiusRelayWebSocket, error)

RadiusRelayWebSocketFactory completes the opening handshake or returns its error. The context owns cancellation of the opening attempt.

type RadiusRelayWebSocketOptions

type RadiusRelayWebSocketOptions struct {
	URL           string
	Protocol      string
	Authorization string
}

RadiusRelayWebSocketOptions carries the exact URL, subprotocol, and bearer header for an attempt.

type RelayByteConnectionHandler

type RelayByteConnectionHandler struct {
	OnData  func([]byte)
	OnClose func()
	OnError func(error)
}

RelayByteConnectionHandler receives ordered data and exactly one terminal callback. Callbacks run on the host's reader, not the TUI loop; they must not wait for host shutdown.

type RelayDataFrame

type RelayDataFrame struct {
	ConnectionID string
	Payload      []byte
}

RelayDataFrame contains a validated lowercase UUIDv4 and an independently owned payload.

func ParseRelayDataFrame

func ParseRelayDataFrame(frame []byte) (RelayDataFrame, bool)

ParseRelayDataFrame validates the envelope and returns a copied payload, or false for malformed frames.

type RelayServerByteConnection

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

RelayServerByteConnection is one server-side virtual byte stream. Send and Close await ordered writes; a final chunk precedes its close control.

func (*RelayServerByteConnection) Close

func (c *RelayServerByteConnection) Close(finalChunk []byte) error

Close removes the connection once, then sends a non-nil final chunk before the close control. It does not synthesize a remote-close callback.

func (*RelayServerByteConnection) Closed

func (c *RelayServerByteConnection) Closed() bool

Closed reports local or remote terminal state.

func (*RelayServerByteConnection) Send

func (c *RelayServerByteConnection) Send(chunk []byte) error

Send sends one copied payload on this connection.

Directories

Path Synopsis
mini
shared
Package shared implements mini's local service protocol and transports.
Package shared implements mini's local service protocol and transports.
Package services contains the opt-in experimental presentation service contracts.
Package services contains the opt-in experimental presentation service contracts.

Jump to

Keyboard shortcuts

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