daemon

package
v0.4.0 Latest Latest
Warning

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

Go to latest
Published: Mar 8, 2026 License: MIT Imports: 44 Imported by: 0

Documentation

Index

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.

func (*AckTracker) Track

func (a *AckTracker) Track(env *messagepb.Envelope, peerAddr string)

Track registers a sent envelope as pending ACK.

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 NewActivityBus

func NewActivityBus() *ActivityBus

NewActivityBus creates a new activity bus.

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 (b *ActivityBus) EmitRoomMessagePosted(
	roomID string,
	roomSeq uint64,
	from string,
	members []string,
	traceID string,
	contentType string,
	contentKind string,
	targetHandle string,
	turnID string,
	status string,
	round uint32,
)

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 (*AgentServer) CreateRoom added in v0.4.0

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

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

GetNodeStatus returns a snapshot of the node's current state for the dashboard.

func (*AgentServer) GetTrace

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

IntrospectHandle returns the manifest for a handle.

func (*AgentServer) JoinRoom added in v0.4.0

func (*AgentServer) LeaveRoom added in v0.4.0

func (*AgentServer) ListHandles

ListHandles returns all known handles, optionally filtered by tags.

func (*AgentServer) ListRoomMembers added in v0.4.0

func (*AgentServer) ListRooms added in v0.4.0

func (*AgentServer) ListSessions

ListSessions lists sessions involving a handle.

func (*AgentServer) OpenSession

OpenSession opens a new session and sends the opening message.

func (*AgentServer) PostRoomMessage added in v0.4.0

func (*AgentServer) Register

Register registers an agent handle on this node.

func (*AgentServer) ReplayRoom added in v0.4.0

func (*AgentServer) ResolveSession

ResolveSession resolves (closes) a session with a final message.

func (*AgentServer) SendMessage

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

Shutdown gracefully shuts down the daemon.

func (*AgentServer) Subscribe

Subscribe opens a stream of incoming messages for a handle.

func (*AgentServer) WatchActivity

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) Close

func (c *CoordClient) Close() error

Close closes the connection.

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

func (c *CoordClient) UpsertRoom(ctx context.Context, room *messagepb.RoomInfo) error

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 New

func New(cfg *config.DaemonConfig, logger *slog.Logger) (*Daemon, error)

New creates a new daemon from config.

func (*Daemon) AgentServer

func (d *Daemon) AgentServer() *AgentServer

AgentServer returns the agent server (used for testing).

func (*Daemon) Resolver

func (d *Daemon) Resolver() *handle.Resolver

Resolver returns the handle resolver.

func (*Daemon) Run

func (d *Daemon) Run(ctx context.Context) error

Run starts the daemon and blocks until the context is cancelled.

func (*Daemon) Sessions

func (d *Daemon) Sessions() *session.Store

Sessions returns the session store.

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) Route

Route routes an envelope to the appropriate destination.

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

func (ms *MessageStore) UpdatePending(messageID string, sentAt time.Time, retries int) error

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

func (m *Metrics) Serve(ctx context.Context, addr string, logger *slog.Logger)

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

func (m *Metrics) SetReadyFunc(fn func() bool)

SetReadyFunc sets the readiness check function for /readyz.

type RoomDeliverer added in v0.4.0

type RoomDeliverer interface {
	DeliverRoomEventLocal(event *messagepb.RoomEvent) bool
}

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) CloseRoom added in v0.4.0

func (rm *RoomManager) CloseRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)

func (*RoomManager) CreateRoom added in v0.4.0

func (rm *RoomManager) CreateRoom(ctx context.Context, creator, title string, initialMembers []string) (*messagepb.RoomInfo, error)

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) JoinRoom added in v0.4.0

func (rm *RoomManager) JoinRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)

func (*RoomManager) LeaveRoom added in v0.4.0

func (rm *RoomManager) LeaveRoom(ctx context.Context, roomID, handle string) (*messagepb.RoomInfo, error)

func (*RoomManager) ListMembers added in v0.4.0

func (rm *RoomManager) ListMembers(ctx context.Context, roomID, handle string) ([]string, error)

func (*RoomManager) ListRooms added in v0.4.0

func (rm *RoomManager) ListRooms(ctx context.Context, handle string) ([]*messagepb.RoomInfo, error)

func (*RoomManager) PostMessage added in v0.4.0

func (rm *RoomManager) PostMessage(ctx context.Context, roomID, from string, payload []byte, contentType, traceID string) (*messagepb.RoomEvent, error)

func (*RoomManager) Replay added in v0.4.0

func (rm *RoomManager) Replay(ctx context.Context, roomID, handle string, sinceSeq uint64) ([]*messagepb.RoomEvent, error)

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 Router

type Router interface {
	Route(ctx context.Context, env *messagepb.Envelope) error
}

Router is the interface the agent server uses to send outbound messages.

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.

Jump to

Keyboard shortcuts

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