Documentation
¶
Index ¶
- type AckTracker
- func (a *AckTracker) Acknowledge(messageID string) bool
- func (a *AckTracker) PendingCount() int
- func (a *AckTracker) Restore(pm *pendingMessage)
- func (a *AckTracker) SetStore(store *MessageStore)
- func (a *AckTracker) StartRetryLoop(ctx context.Context, interval time.Duration)
- func (a *AckTracker) Track(env *messagepb.Envelope, peerAddr string)
- type ActivityBus
- func (b *ActivityBus) Counters() *agentpb.Counters
- func (b *ActivityBus) Emit(event *agentpb.ActivityEvent)
- func (b *ActivityBus) EmitHandleRegistered(handle string)
- func (b *ActivityBus) EmitMessageRouted(sessionID, from, to string, remote bool, traceID, messageID string)
- func (b *ActivityBus) EmitRoomClosed(roomID, closedBy string, members []string)
- func (b *ActivityBus) EmitRoomCreated(roomID, title, createdBy string, members []string)
- func (b *ActivityBus) EmitRoomMemberJoined(roomID, handle string, members []string)
- func (b *ActivityBus) EmitRoomMemberLeft(roomID, handle string, members []string)
- func (b *ActivityBus) EmitRoomMessagePosted(roomID string, roomSeq uint64, from string, members []string, traceID string, ...)
- func (b *ActivityBus) EmitSessionOpened(sessionID, from, to string)
- func (b *ActivityBus) EmitSessionResolved(sessionID, from string)
- func (b *ActivityBus) Subscribe() chan *agentpb.ActivityEvent
- func (b *ActivityBus) Unsubscribe(ch chan *agentpb.ActivityEvent)
- type AgentServer
- func (s *AgentServer) CloseRoom(ctx context.Context, req *agentpb.CloseRoomRequest) (*agentpb.CloseRoomResponse, error)
- func (s *AgentServer) CreateRoom(ctx context.Context, req *agentpb.CreateRoomRequest) (*agentpb.CreateRoomResponse, error)
- func (s *AgentServer) DeliverRoomEventLocal(event *messagepb.RoomEvent) bool
- func (s *AgentServer) DeliverToLocal(env *messagepb.Envelope) bool
- func (s *AgentServer) FindHandles(_ context.Context, req *agentpb.FindHandlesRequest) (*agentpb.FindHandlesResponse, error)
- func (s *AgentServer) ForceStop()
- func (s *AgentServer) GetHandles() []string
- func (s *AgentServer) GetManifests() map[string]*messagepb.ServiceManifest
- func (s *AgentServer) GetNodeStatus(_ context.Context, _ *agentpb.GetNodeStatusRequest) (*agentpb.GetNodeStatusResponse, error)
- func (s *AgentServer) GetTrace(_ context.Context, req *agentpb.GetTraceRequest) (*agentpb.GetTraceResponse, error)
- func (s *AgentServer) GracefulStop()
- func (s *AgentServer) HasHandle(h string) bool
- func (s *AgentServer) IntrospectHandle(_ context.Context, req *agentpb.IntrospectHandleRequest) (*agentpb.IntrospectHandleResponse, error)
- func (s *AgentServer) JoinRoom(ctx context.Context, req *agentpb.JoinRoomRequest) (*agentpb.JoinRoomResponse, error)
- func (s *AgentServer) LeaveRoom(ctx context.Context, req *agentpb.LeaveRoomRequest) (*agentpb.LeaveRoomResponse, error)
- func (s *AgentServer) ListHandles(_ context.Context, req *agentpb.ListHandlesRequest) (*agentpb.ListHandlesResponse, error)
- func (s *AgentServer) ListRoomMembers(ctx context.Context, req *agentpb.ListRoomMembersRequest) (*agentpb.ListRoomMembersResponse, error)
- func (s *AgentServer) ListRooms(ctx context.Context, req *agentpb.ListRoomsRequest) (*agentpb.ListRoomsResponse, error)
- func (s *AgentServer) ListSessions(_ context.Context, req *agentpb.ListSessionsRequest) (*agentpb.ListSessionsResponse, error)
- func (s *AgentServer) OpenSession(ctx context.Context, req *agentpb.OpenSessionRequest) (*agentpb.OpenSessionResponse, error)
- func (s *AgentServer) PostRoomMessage(ctx context.Context, req *agentpb.PostRoomMessageRequest) (*agentpb.PostRoomMessageResponse, error)
- func (s *AgentServer) Register(ctx context.Context, req *agentpb.RegisterRequest) (*agentpb.RegisterResponse, error)
- func (s *AgentServer) ReplayRoom(ctx context.Context, req *agentpb.ReplayRoomRequest) (*agentpb.ReplayRoomResponse, error)
- func (s *AgentServer) ResolveSession(ctx context.Context, req *agentpb.ResolveSessionRequest) (*agentpb.ResolveSessionResponse, error)
- func (s *AgentServer) SendMessage(ctx context.Context, req *agentpb.SendMessageRequest) (*agentpb.SendMessageResponse, error)
- func (s *AgentServer) ServeTCP(lis net.Listener) error
- func (s *AgentServer) ServeUnix(socketPath string) error
- func (s *AgentServer) SetAckTracker(at *AckTracker)
- func (s *AgentServer) SetAuthToken(token string)
- func (s *AgentServer) SetDashboardDeps(nodeID string, resolver *handle.Resolver, tp *transport.GRPCTransport)
- func (s *AgentServer) SetOnHandleChange(fn func(handles []string, manifests map[string]*messagepb.ServiceManifest))
- func (s *AgentServer) SetRoomManager(rm RoomService)
- func (s *AgentServer) SetRouter(r Router)
- func (s *AgentServer) SetShutdownFunc(fn func())
- func (s *AgentServer) SetTracing(ts *TraceStore, m *Metrics)
- func (s *AgentServer) Shutdown(_ context.Context, _ *agentpb.ShutdownRequest) (*agentpb.ShutdownResponse, error)
- func (s *AgentServer) Subscribe(req *agentpb.SubscribeRequest, stream agentpb.AgentAPI_SubscribeServer) error
- func (s *AgentServer) WatchActivity(_ *agentpb.WatchActivityRequest, stream agentpb.AgentAPI_WatchActivityServer) error
- type CoordClient
- func (c *CoordClient) Close() error
- func (c *CoordClient) Heartbeat(ctx context.Context, getHandles func() []string, ...)
- func (c *CoordClient) ListRoomsForHandle(ctx context.Context, handle string) ([]*messagepb.RoomInfo, error)
- func (c *CoordClient) Register(ctx context.Context, handles []string, ...) error
- func (c *CoordClient) SetIsRelay(isRelay bool)
- func (c *CoordClient) SetOnRelayUpdate(fn func())
- func (c *CoordClient) SetTeamID(teamID string)
- func (c *CoordClient) SetTokenRefresh(fn func(ctx context.Context) (string, error))
- func (c *CoordClient) UpsertRoom(ctx context.Context, room *messagepb.RoomInfo) error
- func (c *CoordClient) WatchPeerMap(ctx context.Context) error
- type Daemon
- type LocalDeliverer
- type MessageRouter
- type MessageStore
- func (ms *MessageStore) Close() error
- func (ms *MessageStore) LoadPending() ([]*pendingMessage, error)
- func (ms *MessageStore) LoadRoom(roomID string) (*messagepb.RoomInfo, error)
- func (ms *MessageStore) LoadRooms() ([]*messagepb.RoomInfo, error)
- func (ms *MessageStore) LoadSessions() ([]*session.Session, error)
- func (ms *MessageStore) PendingCount() int
- func (ms *MessageStore) RemovePending(messageID string) error
- func (ms *MessageStore) RemoveSession(sessionID string) error
- func (ms *MessageStore) ReplayRoomEvents(roomID string, sinceSeq uint64) ([]*messagepb.RoomEvent, error)
- func (ms *MessageStore) StorePending(env *messagepb.Envelope, peerAddr string) error
- func (ms *MessageStore) StoreRoom(room *messagepb.RoomInfo) error
- func (ms *MessageStore) StoreRoomEvent(event *messagepb.RoomEvent, maxEvents int) error
- func (ms *MessageStore) StoreSession(sess *session.Session) error
- func (ms *MessageStore) UpdatePending(messageID string, sentAt time.Time, retries int) error
- type Metrics
- type RoomDeliverer
- type RoomManager
- func (rm *RoomManager) CloseRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
- func (rm *RoomManager) CreateRoom(ctx context.Context, creator, title string, initialMembers []string) (*messagepb.RoomInfo, error)
- func (rm *RoomManager) DashboardRooms(handles []string) ([]*messagepb.RoomInfo, error)
- func (rm *RoomManager) HandleTransportMessage(ctx context.Context, msg *transportpb.TransportMessage)
- func (rm *RoomManager) JoinRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
- func (rm *RoomManager) LeaveRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
- func (rm *RoomManager) ListMembers(ctx context.Context, roomID, handle string) ([]string, error)
- func (rm *RoomManager) ListRooms(ctx context.Context, handle string) ([]*messagepb.RoomInfo, error)
- func (rm *RoomManager) PostMessage(ctx context.Context, roomID, from string, payload []byte, ...) (*messagepb.RoomEvent, error)
- func (rm *RoomManager) Replay(ctx context.Context, roomID, handle string, sinceSeq uint64) ([]*messagepb.RoomEvent, error)
- func (rm *RoomManager) SetCoordClient(cc *CoordClient)
- type RoomService
- type Router
- type TraceStore
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AckTracker ¶
type AckTracker struct {
// contains filtered or unexported fields
}
AckTracker tracks pending messages awaiting ACK and retries on timeout.
func NewAckTracker ¶
func NewAckTracker(sendFn func(addr string, env *messagepb.Envelope) error, logger *slog.Logger) *AckTracker
NewAckTracker creates a new ACK tracker.
func (*AckTracker) Acknowledge ¶
func (a *AckTracker) Acknowledge(messageID string) bool
Acknowledge removes a message from pending. Returns true if it was pending.
func (*AckTracker) PendingCount ¶
func (a *AckTracker) PendingCount() int
PendingCount returns the number of messages awaiting ACK.
func (*AckTracker) Restore ¶
func (a *AckTracker) Restore(pm *pendingMessage)
Restore adds a pending message loaded from durable storage (no re-persist).
func (*AckTracker) SetStore ¶
func (a *AckTracker) SetStore(store *MessageStore)
SetStore attaches a durable message store for persistence.
func (*AckTracker) StartRetryLoop ¶
func (a *AckTracker) StartRetryLoop(ctx context.Context, interval time.Duration)
StartRetryLoop runs a background loop that retries expired pending messages. It checks every interval and exits when ctx is cancelled.
type ActivityBus ¶
type ActivityBus struct {
MessagesRouted atomic.Int64
MessagesDeliveredLocal atomic.Int64
MessagesSentRemote atomic.Int64
MessagesReceivedRemote atomic.Int64
SessionsOpened atomic.Int64
SessionsResolved atomic.Int64
RoomsCreated atomic.Int64
RoomMessagesPosted atomic.Int64
RoomMembersJoined atomic.Int64
RoomMembersLeft atomic.Int64
RoomsClosed atomic.Int64
// contains filtered or unexported fields
}
ActivityBus is an in-process pub/sub for dashboard activity events. It uses non-blocking sends so subscribers never slow down the daemon.
func (*ActivityBus) Counters ¶
func (b *ActivityBus) Counters() *agentpb.Counters
Counters returns a snapshot of all atomic counters.
func (*ActivityBus) Emit ¶
func (b *ActivityBus) Emit(event *agentpb.ActivityEvent)
Emit fans out an event to all subscribers without blocking.
func (*ActivityBus) EmitHandleRegistered ¶
func (b *ActivityBus) EmitHandleRegistered(handle string)
func (*ActivityBus) EmitMessageRouted ¶
func (b *ActivityBus) EmitMessageRouted(sessionID, from, to string, remote bool, traceID, messageID string)
func (*ActivityBus) EmitRoomClosed ¶ added in v0.4.0
func (b *ActivityBus) EmitRoomClosed(roomID, closedBy string, members []string)
func (*ActivityBus) EmitRoomCreated ¶ added in v0.4.0
func (b *ActivityBus) EmitRoomCreated(roomID, title, createdBy string, members []string)
func (*ActivityBus) EmitRoomMemberJoined ¶ added in v0.4.0
func (b *ActivityBus) EmitRoomMemberJoined(roomID, handle string, members []string)
func (*ActivityBus) EmitRoomMemberLeft ¶ added in v0.4.0
func (b *ActivityBus) EmitRoomMemberLeft(roomID, handle string, members []string)
func (*ActivityBus) EmitRoomMessagePosted ¶ added in v0.4.0
func (*ActivityBus) EmitSessionOpened ¶
func (b *ActivityBus) EmitSessionOpened(sessionID, from, to string)
func (*ActivityBus) EmitSessionResolved ¶
func (b *ActivityBus) EmitSessionResolved(sessionID, from string)
func (*ActivityBus) Subscribe ¶
func (b *ActivityBus) Subscribe() chan *agentpb.ActivityEvent
Subscribe returns a channel that receives activity events. The caller must call Unsubscribe when done.
func (*ActivityBus) Unsubscribe ¶
func (b *ActivityBus) Unsubscribe(ch chan *agentpb.ActivityEvent)
Unsubscribe removes a subscriber channel.
type AgentServer ¶
type AgentServer struct {
agentpb.UnimplementedAgentAPIServer
// contains filtered or unexported fields
}
AgentServer is the local gRPC server that agent programs connect to via Unix socket.
func NewAgentServer ¶
func NewAgentServer(sessions *session.Store, router Router, activity *ActivityBus, logger *slog.Logger) *AgentServer
NewAgentServer creates a new agent server.
func (*AgentServer) CloseRoom ¶ added in v0.4.0
func (s *AgentServer) CloseRoom(ctx context.Context, req *agentpb.CloseRoomRequest) (*agentpb.CloseRoomResponse, error)
func (*AgentServer) CreateRoom ¶ added in v0.4.0
func (s *AgentServer) CreateRoom(ctx context.Context, req *agentpb.CreateRoomRequest) (*agentpb.CreateRoomResponse, error)
func (*AgentServer) DeliverRoomEventLocal ¶ added in v0.4.0
func (s *AgentServer) DeliverRoomEventLocal(event *messagepb.RoomEvent) bool
DeliverRoomEventLocal delivers a room event to all local room members except the sender.
func (*AgentServer) DeliverToLocal ¶
func (s *AgentServer) DeliverToLocal(env *messagepb.Envelope) bool
DeliverToLocal delivers an envelope to local subscribers. Returns true if delivered.
func (*AgentServer) FindHandles ¶ added in v0.4.0
func (s *AgentServer) FindHandles(_ context.Context, req *agentpb.FindHandlesRequest) (*agentpb.FindHandlesResponse, error)
FindHandles returns ranked handles matching structured discovery constraints.
func (*AgentServer) ForceStop ¶ added in v0.4.0
func (s *AgentServer) ForceStop()
ForceStop immediately stops the gRPC server without waiting for streams.
func (*AgentServer) GetHandles ¶
func (s *AgentServer) GetHandles() []string
GetHandles returns the list of currently registered handles.
func (*AgentServer) GetManifests ¶
func (s *AgentServer) GetManifests() map[string]*messagepb.ServiceManifest
GetManifests returns a snapshot of all handle manifests.
func (*AgentServer) GetNodeStatus ¶
func (s *AgentServer) GetNodeStatus(_ context.Context, _ *agentpb.GetNodeStatusRequest) (*agentpb.GetNodeStatusResponse, error)
GetNodeStatus returns a snapshot of the node's current state for the dashboard.
func (*AgentServer) GetTrace ¶
func (s *AgentServer) GetTrace(_ context.Context, req *agentpb.GetTraceRequest) (*agentpb.GetTraceResponse, error)
GetTrace returns trace spans for a given trace ID from the local node.
func (*AgentServer) GracefulStop ¶
func (s *AgentServer) GracefulStop()
GracefulStop stops the server gracefully and removes the auth token file.
func (*AgentServer) HasHandle ¶
func (s *AgentServer) HasHandle(h string) bool
HasHandle returns true if the given handle is registered locally.
func (*AgentServer) IntrospectHandle ¶
func (s *AgentServer) IntrospectHandle(_ context.Context, req *agentpb.IntrospectHandleRequest) (*agentpb.IntrospectHandleResponse, error)
IntrospectHandle returns the manifest for a handle.
func (*AgentServer) JoinRoom ¶ added in v0.4.0
func (s *AgentServer) JoinRoom(ctx context.Context, req *agentpb.JoinRoomRequest) (*agentpb.JoinRoomResponse, error)
func (*AgentServer) LeaveRoom ¶ added in v0.4.0
func (s *AgentServer) LeaveRoom(ctx context.Context, req *agentpb.LeaveRoomRequest) (*agentpb.LeaveRoomResponse, error)
func (*AgentServer) ListHandles ¶
func (s *AgentServer) ListHandles(_ context.Context, req *agentpb.ListHandlesRequest) (*agentpb.ListHandlesResponse, error)
ListHandles returns all known handles, optionally filtered by tags.
func (*AgentServer) ListRoomMembers ¶ added in v0.4.0
func (s *AgentServer) ListRoomMembers(ctx context.Context, req *agentpb.ListRoomMembersRequest) (*agentpb.ListRoomMembersResponse, error)
func (*AgentServer) ListRooms ¶ added in v0.4.0
func (s *AgentServer) ListRooms(ctx context.Context, req *agentpb.ListRoomsRequest) (*agentpb.ListRoomsResponse, error)
func (*AgentServer) ListSessions ¶
func (s *AgentServer) ListSessions(_ context.Context, req *agentpb.ListSessionsRequest) (*agentpb.ListSessionsResponse, error)
ListSessions lists sessions involving a handle.
func (*AgentServer) OpenSession ¶
func (s *AgentServer) OpenSession(ctx context.Context, req *agentpb.OpenSessionRequest) (*agentpb.OpenSessionResponse, error)
OpenSession opens a new session and sends the opening message.
func (*AgentServer) PostRoomMessage ¶ added in v0.4.0
func (s *AgentServer) PostRoomMessage(ctx context.Context, req *agentpb.PostRoomMessageRequest) (*agentpb.PostRoomMessageResponse, error)
func (*AgentServer) Register ¶
func (s *AgentServer) Register(ctx context.Context, req *agentpb.RegisterRequest) (*agentpb.RegisterResponse, error)
Register registers an agent handle on this node.
func (*AgentServer) ReplayRoom ¶ added in v0.4.0
func (s *AgentServer) ReplayRoom(ctx context.Context, req *agentpb.ReplayRoomRequest) (*agentpb.ReplayRoomResponse, error)
func (*AgentServer) ResolveSession ¶
func (s *AgentServer) ResolveSession(ctx context.Context, req *agentpb.ResolveSessionRequest) (*agentpb.ResolveSessionResponse, error)
ResolveSession resolves (closes) a session with a final message.
func (*AgentServer) SendMessage ¶
func (s *AgentServer) SendMessage(ctx context.Context, req *agentpb.SendMessageRequest) (*agentpb.SendMessageResponse, error)
SendMessage sends a message within an existing session.
func (*AgentServer) ServeTCP ¶
func (s *AgentServer) ServeTCP(lis net.Listener) error
ServeTCP starts the gRPC server on a TCP address (for testing).
func (*AgentServer) ServeUnix ¶
func (s *AgentServer) ServeUnix(socketPath string) error
ServeUnix starts the gRPC server on a Unix socket. It generates a random auth token and writes it to <socketPath>.token (mode 0600).
func (*AgentServer) SetAckTracker ¶
func (s *AgentServer) SetAckTracker(at *AckTracker)
SetAckTracker sets the ACK tracker for delivery acknowledgement.
func (*AgentServer) SetAuthToken ¶
func (s *AgentServer) SetAuthToken(token string)
SetAuthToken sets the auth token for testing (bypasses file-based token generation).
func (*AgentServer) SetDashboardDeps ¶
func (s *AgentServer) SetDashboardDeps(nodeID string, resolver *handle.Resolver, tp *transport.GRPCTransport)
SetDashboardDeps sets dependencies needed for dashboard RPCs.
func (*AgentServer) SetOnHandleChange ¶
func (s *AgentServer) SetOnHandleChange(fn func(handles []string, manifests map[string]*messagepb.ServiceManifest))
SetOnHandleChange sets a callback invoked when local handles change.
func (*AgentServer) SetRoomManager ¶ added in v0.4.0
func (s *AgentServer) SetRoomManager(rm RoomService)
SetRoomManager sets the room manager used by room RPCs.
func (*AgentServer) SetRouter ¶
func (s *AgentServer) SetRouter(r Router)
SetRouter sets the router (used for breaking circular dependency during setup).
func (*AgentServer) SetShutdownFunc ¶ added in v0.3.6
func (s *AgentServer) SetShutdownFunc(fn func())
SetShutdownFunc sets the function called when a Shutdown RPC is received.
func (*AgentServer) SetTracing ¶
func (s *AgentServer) SetTracing(ts *TraceStore, m *Metrics)
SetTracing sets the trace store and metrics for tracing support.
func (*AgentServer) Shutdown ¶ added in v0.3.6
func (s *AgentServer) Shutdown(_ context.Context, _ *agentpb.ShutdownRequest) (*agentpb.ShutdownResponse, error)
Shutdown gracefully shuts down the daemon.
func (*AgentServer) Subscribe ¶
func (s *AgentServer) Subscribe(req *agentpb.SubscribeRequest, stream agentpb.AgentAPI_SubscribeServer) error
Subscribe opens a stream of incoming messages for a handle.
func (*AgentServer) WatchActivity ¶
func (s *AgentServer) WatchActivity(_ *agentpb.WatchActivityRequest, stream agentpb.AgentAPI_WatchActivityServer) error
WatchActivity streams activity events to the dashboard.
type CoordClient ¶
type CoordClient struct {
// contains filtered or unexported fields
}
CoordClient connects to the coordination server for registration and peer map updates.
func NewCoordClient ¶
func NewCoordClient(coordAddr, nodeID string, pubKey []byte, advertiseAddr string, resolver *handle.Resolver, logger *slog.Logger, kp *identity.Keypair, coordFPFile string, authToken string) (*CoordClient, error)
NewCoordClient creates a new coordination client. If kp is non-nil, mTLS is enabled with TOFU verification of the coord cert using the fingerprint file at coordFPFile. authToken is sent with RegisterNode requests for admission control (empty = no token).
func (*CoordClient) Heartbeat ¶
func (c *CoordClient) Heartbeat(ctx context.Context, getHandles func() []string, getManifests func() map[string]*messagepb.ServiceManifest, interval time.Duration, reRegister func(ctx context.Context) error)
Heartbeat sends periodic heartbeats to the coordination server. If the coord responds with Ok: false (e.g. "node not found" after a coord restart), it calls reRegister to re-register the node. Blocks until the context is cancelled.
func (*CoordClient) ListRoomsForHandle ¶ added in v0.4.0
func (c *CoordClient) ListRoomsForHandle(ctx context.Context, handle string) ([]*messagepb.RoomInfo, error)
ListRoomsForHandle returns coord-scoped room metadata for a handle.
func (*CoordClient) Register ¶
func (c *CoordClient) Register(ctx context.Context, handles []string, manifests map[string]*messagepb.ServiceManifest) error
Register registers this node with the coordination server.
func (*CoordClient) SetIsRelay ¶
func (c *CoordClient) SetIsRelay(isRelay bool)
SetIsRelay marks this client as a relay node for registration.
func (*CoordClient) SetOnRelayUpdate ¶
func (c *CoordClient) SetOnRelayUpdate(fn func())
SetOnRelayUpdate sets a callback invoked after relay info is updated from the peer map.
func (*CoordClient) SetTeamID ¶
func (c *CoordClient) SetTeamID(teamID string)
SetTeamID sets the team ID for registration and peer map requests.
func (*CoordClient) SetTokenRefresh ¶
func (c *CoordClient) SetTokenRefresh(fn func(ctx context.Context) (string, error))
SetTokenRefresh sets a callback that refreshes the auth token when auth errors occur.
func (*CoordClient) UpsertRoom ¶ added in v0.4.0
UpsertRoom stores room discovery metadata in the coordination server.
func (*CoordClient) WatchPeerMap ¶
func (c *CoordClient) WatchPeerMap(ctx context.Context) error
WatchPeerMap starts watching for peer map updates and updating the resolver. Blocks until the context is cancelled.
type Daemon ¶
type Daemon struct {
// contains filtered or unexported fields
}
Daemon is the main node daemon that ties together all components.
func (*Daemon) AgentServer ¶
func (d *Daemon) AgentServer() *AgentServer
AgentServer returns the agent server (used for testing).
func (*Daemon) TraceStore ¶
func (d *Daemon) TraceStore() *TraceStore
TraceStore returns the trace store (used for testing).
type LocalDeliverer ¶
type LocalDeliverer interface {
DeliverToLocal(env *messagepb.Envelope) bool
HasHandle(h string) bool
}
LocalDeliverer delivers messages to local agents.
type MessageRouter ¶
type MessageRouter struct {
// contains filtered or unexported fields
}
MessageRouter routes envelopes to local agents or remote daemons.
func NewMessageRouter ¶
func NewMessageRouter(resolver *handle.Resolver, transport transport.Transport, local LocalDeliverer, activity *ActivityBus, logger *slog.Logger) *MessageRouter
NewMessageRouter creates a new message router.
func (*MessageRouter) SetAckTracker ¶
func (r *MessageRouter) SetAckTracker(at *AckTracker)
SetAckTracker sets the ACK tracker for the router.
func (*MessageRouter) SetTracing ¶
func (r *MessageRouter) SetTracing(ts *TraceStore, m *Metrics, nodeID string)
SetTracing sets the trace store, metrics, and node ID for tracing support.
type MessageStore ¶
type MessageStore struct {
// contains filtered or unexported fields
}
MessageStore provides durable storage for pending messages and sessions using bbolt. Messages are stored until ACKed. Sessions survive daemon restarts.
func NewMessageStore ¶
func NewMessageStore(path string, logger *slog.Logger) (*MessageStore, error)
NewMessageStore opens (or creates) a bbolt database at the given path.
func (*MessageStore) Close ¶
func (ms *MessageStore) Close() error
Close closes the underlying database.
func (*MessageStore) LoadPending ¶
func (ms *MessageStore) LoadPending() ([]*pendingMessage, error)
LoadPending returns all pending messages from the store.
func (*MessageStore) LoadRoom ¶ added in v0.4.0
func (ms *MessageStore) LoadRoom(roomID string) (*messagepb.RoomInfo, error)
LoadRoom returns a persisted room, if any.
func (*MessageStore) LoadRooms ¶ added in v0.4.0
func (ms *MessageStore) LoadRooms() ([]*messagepb.RoomInfo, error)
LoadRooms returns all persisted rooms.
func (*MessageStore) LoadSessions ¶
func (ms *MessageStore) LoadSessions() ([]*session.Session, error)
LoadSessions returns all persisted sessions.
func (*MessageStore) PendingCount ¶
func (ms *MessageStore) PendingCount() int
PendingCount returns the number of pending messages in the store.
func (*MessageStore) RemovePending ¶
func (ms *MessageStore) RemovePending(messageID string) error
RemovePending deletes a pending message (called on ACK).
func (*MessageStore) RemoveSession ¶
func (ms *MessageStore) RemoveSession(sessionID string) error
RemoveSession deletes a session from the store.
func (*MessageStore) ReplayRoomEvents ¶ added in v0.4.0
func (ms *MessageStore) ReplayRoomEvents(roomID string, sinceSeq uint64) ([]*messagepb.RoomEvent, error)
ReplayRoomEvents returns retained room events with sequence greater than sinceSeq.
func (*MessageStore) StorePending ¶
func (ms *MessageStore) StorePending(env *messagepb.Envelope, peerAddr string) error
StorePending persists a sent envelope that is awaiting ACK.
func (*MessageStore) StoreRoom ¶ added in v0.4.0
func (ms *MessageStore) StoreRoom(room *messagepb.RoomInfo) error
StoreRoom persists room metadata.
func (*MessageStore) StoreRoomEvent ¶ added in v0.4.0
func (ms *MessageStore) StoreRoomEvent(event *messagepb.RoomEvent, maxEvents int) error
StoreRoomEvent appends a room event and enforces a max retained event count.
func (*MessageStore) StoreSession ¶
func (ms *MessageStore) StoreSession(sess *session.Session) error
StoreSession persists a session to disk.
func (*MessageStore) UpdatePending ¶ added in v0.4.0
UpdatePending updates retry metadata for a pending message.
type Metrics ¶
type Metrics struct {
// Histograms (observed explicitly)
MessageRoutingDuration *prometheus.Histogram
SessionLifetime *prometheus.Histogram
// contains filtered or unexported fields
}
Metrics exports Prometheus metrics for the daemon. It reads counters directly from ActivityBus atomics at scrape time (no double-counting), and provides histograms for routing duration and session lifetime.
func NewMetrics ¶
func NewMetrics(activity *ActivityBus) *Metrics
NewMetrics creates a new Metrics instance backed by the given ActivityBus.
func (*Metrics) Collect ¶
func (m *Metrics) Collect(ch chan<- prometheus.Metric)
Collect implements prometheus.Collector.
func (*Metrics) Describe ¶
func (m *Metrics) Describe(ch chan<- *prometheus.Desc)
Describe implements prometheus.Collector.
func (*Metrics) Serve ¶
Serve starts an HTTP server exposing /metrics on the given address. It blocks until the context is cancelled, then shuts down gracefully.
func (*Metrics) SetReadyFunc ¶
SetReadyFunc sets the readiness check function for /readyz.
type RoomDeliverer ¶ added in v0.4.0
RoomDeliverer delivers room events to local agent subscribers.
type RoomManager ¶ added in v0.4.0
type RoomManager struct {
// contains filtered or unexported fields
}
func NewRoomManager ¶ added in v0.4.0
func NewRoomManager(nodeID string, resolver *handle.Resolver, tp transport.Transport, deliverer RoomDeliverer, store *MessageStore, activity *ActivityBus, logger *slog.Logger) *RoomManager
func (*RoomManager) CreateRoom ¶ added in v0.4.0
func (*RoomManager) DashboardRooms ¶ added in v0.4.0
func (rm *RoomManager) DashboardRooms(handles []string) ([]*messagepb.RoomInfo, error)
func (*RoomManager) HandleTransportMessage ¶ added in v0.4.0
func (rm *RoomManager) HandleTransportMessage(ctx context.Context, msg *transportpb.TransportMessage)
func (*RoomManager) ListMembers ¶ added in v0.4.0
func (*RoomManager) PostMessage ¶ added in v0.4.0
func (*RoomManager) SetCoordClient ¶ added in v0.4.0
func (rm *RoomManager) SetCoordClient(cc *CoordClient)
type RoomService ¶ added in v0.4.0
type RoomService interface {
CreateRoom(ctx context.Context, creator, title string, initialMembers []string) (*messagepb.RoomInfo, error)
JoinRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
LeaveRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
PostMessage(ctx context.Context, roomID, from string, payload []byte, contentType, traceID string) (*messagepb.RoomEvent, error)
ListRooms(ctx context.Context, handle string) ([]*messagepb.RoomInfo, error)
ListMembers(ctx context.Context, roomID, handle string) ([]string, error)
Replay(ctx context.Context, roomID, handle string, sinceSeq uint64) ([]*messagepb.RoomEvent, error)
CloseRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)
DashboardRooms(handles []string) ([]*messagepb.RoomInfo, error)
}
type TraceStore ¶
type TraceStore struct {
// contains filtered or unexported fields
}
TraceStore is a concurrent ring buffer that stores trace spans indexed by trace ID.
func NewTraceStore ¶
func NewTraceStore(capacity int) *TraceStore
NewTraceStore creates a new trace store with the given capacity.
func (*TraceStore) GetTrace ¶
func (ts *TraceStore) GetTrace(traceID string) []*agentpb.TraceSpan
GetTrace returns all spans for a given trace ID, ordered by their insertion order.
func (*TraceStore) Record ¶
func (ts *TraceStore) Record(span *agentpb.TraceSpan)
Record stores a trace span in the ring buffer.
func (*TraceStore) RecordSpan ¶
func (ts *TraceStore) RecordSpan(traceID, messageID, nodeID string, action agentpb.TraceAction, metadata map[string]string)
RecordSpan is a convenience method that creates and records a TraceSpan.