parser

package
v1.0.20 Latest Latest
Warning

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

Go to latest
Published: Aug 16, 2026 License: MIT Imports: 17 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Component

type Component struct {
	cfacade.Component
	// contains filtered or unexported fields
}

Component hosts AGP over one or more configured connectors and owns all accepted client connections.

func New

func New(name string, connectors []cfacade.IConnector, opts ...Option) *Component

New creates an AGP server component. Connectors are started after all application components and Actor methods have completed initialization.

func (*Component) Bind

func (c *Component) Bind(connectionID string, uid int64, data map[string]string) error

Bind associates an authenticated UID and server-owned session data with a connection for later requests and server push.

func (*Component) Broadcast

func (c *Component) Broadcast(methodID uint32, payload any) error

Broadcast pushes a notification to every active connection and joins errors.

func (*Component) Connections

func (c *Component) Connections() *ConnectionManager

Connections exposes the connection index for authentication and server push.

func (*Component) HandleConn

func (c *Component) HandleConn(conn net.Conn)

HandleConn adapts a connector connection to AGP packet boundaries.

func (*Component) Init

func (c *Component) Init()

Init verifies that the shared method table is available.

func (*Component) Kick

func (c *Component) Kick(connectionID string, reasonCode int32, reason string, reconnectable bool) error

Kick sends a terminal notification and closes the selected connection.

func (*Component) Name

func (c *Component) Name() string

Name returns the application-unique component name.

func (*Component) Notify

func (c *Component) Notify(connectionID string, methodID uint32, payload any) error

Notify pushes a notification to one connection ID.

func (*Component) NotifyUID

func (c *Component) NotifyUID(uid int64, methodID uint32, payload any) error

NotifyUID pushes a notification to the connection bound to uid.

func (*Component) OnAfterInit

func (c *Component) OnAfterInit()

OnAfterInit configures and starts every transport connector.

func (*Component) OnBeforeStop

func (c *Component) OnBeforeStop()

OnBeforeStop first stops accepting sockets, then drains a stable connection snapshot so no new client can appear after CloseAll.

func (*Component) OnStop

func (c *Component) OnStop()

OnStop closes all configured connectors.

type Connection

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

Connection owns AGP framing, request cancellation and the session associated with one TCP or WebSocket client.

func (*Connection) Close

func (c *Connection) Close()

Close is idempotent and cancels every in-flight request before unregistering the connection from its manager.

func (*Connection) GoAway

func (c *Connection) GoAway(reasonCode int32, retryAfter time.Duration) error

GoAway asks the client to reconnect later and then closes the connection.

func (*Connection) ID

func (c *Connection) ID() string

ID returns the server-assigned connection identifier.

func (*Connection) Kick

func (c *Connection) Kick(reasonCode int32, reason string, reconnectable bool) error

Kick sends the terminal kick notification synchronously before closing.

func (*Connection) Notify

func (c *Connection) Notify(methodID uint32, payload any) error

Notify pushes an application notification using the configured default codec.

func (*Connection) Run

func (c *Connection) Run()

Run performs the handshake/state validation loop. Packet handlers may run in separate goroutines, but all writes are serialized by writeLoop.

func (*Connection) Session

func (c *Connection) Session() *cproto.Session

Session returns a snapshot so handlers cannot mutate connection state without going through the manager.

func (*Connection) State

func (c *Connection) State() ConnectionState

State returns the current connection lifecycle state.

func (*Connection) UID

func (c *Connection) UID() int64

UID returns the currently bound user ID, or zero when unbound.

type ConnectionManager

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

ConnectionManager indexes live connections by both connection ID and bound UID.

func NewConnectionManager

func NewConnectionManager() *ConnectionManager

NewConnectionManager creates an empty connection index.

func (*ConnectionManager) Add

func (m *ConnectionManager) Add(connection *Connection)

Add indexes a newly accepted connection.

func (*ConnectionManager) Bind

func (m *ConnectionManager) Bind(id string, uid int64, data map[string]string) error

Bind enforces a one-to-one UID-to-connection mapping and replaces the server-owned session data exposed to subsequent Actor requests.

func (*ConnectionManager) CloseAll

func (m *ConnectionManager) CloseAll()

CloseAll idempotently closes every indexed connection.

func (*ConnectionManager) Count

func (m *ConnectionManager) Count() int

Count returns the number of active indexed connections.

func (*ConnectionManager) Get

func (m *ConnectionManager) Get(id string) (*Connection, bool)

Get looks up a connection by its server-assigned ID.

func (*ConnectionManager) GetUID

func (m *ConnectionManager) GetUID(uid int64) (*Connection, bool)

GetUID looks up the connection currently bound to uid.

func (*ConnectionManager) Range

func (m *ConnectionManager) Range(fn func(*Connection) bool)

Range visits a stable snapshot and stops when fn returns false.

func (*ConnectionManager) Remove

func (m *ConnectionManager) Remove(id string)

Remove deletes both the connection ID and any UID binding.

func (*ConnectionManager) Snapshot

func (m *ConnectionManager) Snapshot() []*Connection

Snapshot returns the current connections without holding the manager lock.

type ConnectionState

type ConnectionState int32

ConnectionState is the lifecycle state of one AGP client connection.

const (
	// A connection must complete the protobuf handshake before application calls.
	ConnectionHandshaking ConnectionState = iota + 1
	ConnectionReady
	ConnectionDraining
	ConnectionClosed
)

type Option

type Option func(*Options)

Option mutates AGP component options during construction.

func WithHandshakeTimeout

func WithHandshakeTimeout(timeout time.Duration) Option

WithHandshakeTimeout limits how long a new connection may remain uninitialized.

func WithHeartbeat

func WithHeartbeat(interval, idle time.Duration) Option

WithHeartbeat configures the advertised heartbeat interval and idle deadline.

func WithLimits

func WithLimits(limits cproto.Limits) Option

WithLimits replaces packet, body and metadata limits.

func WithMaxInflight

func WithMaxInflight(max int) Option

WithMaxInflight limits concurrent request handlers on each connection.

func WithMaxRequestTimeout

func WithMaxRequestTimeout(timeout time.Duration) Option

WithMaxRequestTimeout caps the timeout requested by an AGP client.

func WithOnDisconnect

func WithOnDisconnect(fn func(*Connection)) Option

WithOnDisconnect registers a callback invoked once per connection close.

func WithWriteQueue

func WithWriteQueue(size int) Option

WithWriteQueue sets the maximum buffered outbound packet count.

func WithWriteTimeout

func WithWriteTimeout(timeout time.Duration) Option

WithWriteTimeout limits one packet write.

type Options

type Options struct {
	Limits            cproto.Limits
	WriteQueueSize    int
	MaxInflight       int
	HandshakeTimeout  time.Duration
	HeartbeatInterval time.Duration
	IdleTimeout       time.Duration
	WriteTimeout      time.Duration
	MaxRequestTimeout time.Duration
	// OnDisconnect is invoked asynchronously once after a connection is removed
	// from the manager. Shutdown waits for it only up to WriteTimeout.
	OnDisconnect func(*Connection)
}

Options controls AGP limits, timeouts and connection backpressure.

type TCPPacketFramer

type TCPPacketFramer struct {
	MaxPacketSize uint32
}

TCPPacketFramer implements AGP's four-byte big-endian packet length prefix.

func NewTCPPacketFramer

func NewTCPPacketFramer(maxPacketSize uint32) *TCPPacketFramer

NewTCPPacketFramer creates a framer and applies the protocol default for zero.

func (*TCPPacketFramer) ReadPacketBytes

func (f *TCPPacketFramer) ReadPacketBytes(reader io.Reader) ([]byte, error)

ReadPacketBytes reads exactly one bounded AGP packet.

func (*TCPPacketFramer) WritePacketBytes

func (f *TCPPacketFramer) WritePacketBytes(writer io.Writer, packet []byte) error

WritePacketBytes writes one length-prefixed AGP packet.

Jump to

Keyboard shortcuts

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