iop/apps/edge/internal/events/bus.go
toki 6e2a1fc2f2 feat: edge runtime, opsconsole, CLI target terminology, and e2e validation updates
- Add edge runtime config and opsconsole package
- Refactor edge console to use new runtime config
- Add service source metadata support
- Update CLI adapter with target terminology
- Add edge operation contract and event bus replay
- Update node label and command ops surface
- Add E2E smoke tests and full validation
- Update proto runtime definitions
- Update documentation and agent-ops rules
2026-05-17 14:57:49 +09:00

265 lines
5.7 KiB
Go

package events
import (
"sync"
iop "iop/proto/gen/iop"
)
// DefaultReplayCapacity bounds the per-bus ring buffer of recent run/node
// events kept for late-joining ops surfaces.
const DefaultReplayCapacity = 64
// BusStats is a lock-safe snapshot of bus counters. Counters are cumulative
// since the Bus was created.
type BusStats struct {
DroppedRunEvents uint64
DroppedNodeEvents uint64
}
// Bus is the in-process edge event fanout used by temporary consoles, future
// API handlers, and later event persistence.
type Bus struct {
mu sync.Mutex
runSubs map[string]map[chan *iop.RunEvent]struct{}
nodeSubs map[string]map[chan *iop.EdgeNodeEvent]struct{}
allRuns map[chan *iop.RunEvent]struct{}
allNodes map[chan *iop.EdgeNodeEvent]struct{}
replayCap int
recentRuns []*iop.RunEvent
recentNodes []*iop.EdgeNodeEvent
droppedRun uint64
droppedNode uint64
}
func NewBus() *Bus {
return NewBusWithReplayCapacity(DefaultReplayCapacity)
}
// NewBusWithReplayCapacity creates a Bus whose replay ring keeps up to cap
// most-recent events of each kind. A cap <= 0 disables replay retention.
func NewBusWithReplayCapacity(cap int) *Bus {
if cap < 0 {
cap = 0
}
return &Bus{
runSubs: make(map[string]map[chan *iop.RunEvent]struct{}),
nodeSubs: make(map[string]map[chan *iop.EdgeNodeEvent]struct{}),
allRuns: make(map[chan *iop.RunEvent]struct{}),
allNodes: make(map[chan *iop.EdgeNodeEvent]struct{}),
replayCap: cap,
}
}
func (b *Bus) SubscribeRun(runID string, buffer int) (<-chan *iop.RunEvent, func()) {
if buffer <= 0 {
buffer = 1
}
ch := make(chan *iop.RunEvent, buffer)
b.mu.Lock()
if b.runSubs[runID] == nil {
b.runSubs[runID] = make(map[chan *iop.RunEvent]struct{})
}
b.runSubs[runID][ch] = struct{}{}
b.mu.Unlock()
return ch, func() {
b.mu.Lock()
if subs := b.runSubs[runID]; subs != nil {
delete(subs, ch)
if len(subs) == 0 {
delete(b.runSubs, runID)
}
}
close(ch)
b.mu.Unlock()
}
}
func (b *Bus) SubscribeNode(nodeID string, buffer int) (<-chan *iop.EdgeNodeEvent, func()) {
if buffer <= 0 {
buffer = 1
}
ch := make(chan *iop.EdgeNodeEvent, buffer)
b.mu.Lock()
if b.nodeSubs[nodeID] == nil {
b.nodeSubs[nodeID] = make(map[chan *iop.EdgeNodeEvent]struct{})
}
b.nodeSubs[nodeID][ch] = struct{}{}
b.mu.Unlock()
return ch, func() {
b.mu.Lock()
if subs := b.nodeSubs[nodeID]; subs != nil {
delete(subs, ch)
if len(subs) == 0 {
delete(b.nodeSubs, nodeID)
}
}
close(ch)
b.mu.Unlock()
}
}
func (b *Bus) SubscribeAllRuns(buffer int) (<-chan *iop.RunEvent, func()) {
if buffer <= 0 {
buffer = 1
}
ch := make(chan *iop.RunEvent, buffer)
b.mu.Lock()
b.allRuns[ch] = struct{}{}
b.mu.Unlock()
return ch, func() {
b.mu.Lock()
delete(b.allRuns, ch)
close(ch)
b.mu.Unlock()
}
}
func (b *Bus) SubscribeAllNodes(buffer int) (<-chan *iop.EdgeNodeEvent, func()) {
if buffer <= 0 {
buffer = 1
}
ch := make(chan *iop.EdgeNodeEvent, buffer)
b.mu.Lock()
b.allNodes[ch] = struct{}{}
b.mu.Unlock()
return ch, func() {
b.mu.Lock()
delete(b.allNodes, ch)
close(ch)
b.mu.Unlock()
}
}
func (b *Bus) PublishRun(event *iop.RunEvent) {
if event == nil {
return
}
b.mu.Lock()
defer b.mu.Unlock()
for ch := range b.allRuns {
if !offerRun(ch, event) {
b.droppedRun++
}
}
for ch := range b.runSubs[event.GetRunId()] {
if !offerRun(ch, event) {
b.droppedRun++
}
}
b.recordRunLocked(event)
}
func (b *Bus) PublishNode(event *iop.EdgeNodeEvent) {
if event == nil {
return
}
b.mu.Lock()
defer b.mu.Unlock()
for ch := range b.allNodes {
if !offerNode(ch, event) {
b.droppedNode++
}
}
for ch := range b.nodeSubs[event.GetNodeId()] {
if !offerNode(ch, event) {
b.droppedNode++
}
}
b.recordNodeLocked(event)
}
// Stats returns a snapshot of cumulative bus counters.
func (b *Bus) Stats() BusStats {
b.mu.Lock()
defer b.mu.Unlock()
return BusStats{
DroppedRunEvents: b.droppedRun,
DroppedNodeEvents: b.droppedNode,
}
}
// ReplayRun returns up to replayCap most-recent run events for runID, oldest
// first. Returns nil when replay is disabled or no matching events are buffered.
func (b *Bus) ReplayRun(runID string) []*iop.RunEvent {
b.mu.Lock()
defer b.mu.Unlock()
if len(b.recentRuns) == 0 {
return nil
}
out := make([]*iop.RunEvent, 0, len(b.recentRuns))
for _, e := range b.recentRuns {
if runID == "" || e.GetRunId() == runID {
out = append(out, e)
}
}
if len(out) == 0 {
return nil
}
return out
}
// ReplayNode returns up to replayCap most-recent node events for nodeID,
// oldest first. Returns nil when replay is disabled or no matching events are
// buffered.
func (b *Bus) ReplayNode(nodeID string) []*iop.EdgeNodeEvent {
b.mu.Lock()
defer b.mu.Unlock()
if len(b.recentNodes) == 0 {
return nil
}
out := make([]*iop.EdgeNodeEvent, 0, len(b.recentNodes))
for _, e := range b.recentNodes {
if nodeID == "" || e.GetNodeId() == nodeID {
out = append(out, e)
}
}
if len(out) == 0 {
return nil
}
return out
}
func (b *Bus) recordRunLocked(event *iop.RunEvent) {
if b.replayCap <= 0 {
return
}
if len(b.recentRuns) >= b.replayCap {
b.recentRuns = append(b.recentRuns[:0], b.recentRuns[1:]...)
}
b.recentRuns = append(b.recentRuns, event)
}
func (b *Bus) recordNodeLocked(event *iop.EdgeNodeEvent) {
if b.replayCap <= 0 {
return
}
if len(b.recentNodes) >= b.replayCap {
b.recentNodes = append(b.recentNodes[:0], b.recentNodes[1:]...)
}
b.recentNodes = append(b.recentNodes, event)
}
func offerRun(ch chan *iop.RunEvent, event *iop.RunEvent) bool {
select {
case ch <- event:
return true
default:
return false
}
}
func offerNode(ch chan *iop.EdgeNodeEvent, event *iop.EdgeNodeEvent) bool {
select {
case ch <- event:
return true
default:
return false
}
}