- Add console emitter interface for event-driven console output - Implement persistent cancel reason and explicit completion tracking - Add profile proto message definitions - Update edge and node transport layers with adapter execution terminology - Add new packages/events module - Update architecture documentation and README files
152 lines
3 KiB
Go
152 lines
3 KiB
Go
package events
|
|
|
|
import (
|
|
"sync"
|
|
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
// 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{}
|
|
}
|
|
|
|
func NewBus() *Bus {
|
|
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{}),
|
|
}
|
|
}
|
|
|
|
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 {
|
|
offerRun(ch, event)
|
|
}
|
|
for ch := range b.runSubs[event.GetRunId()] {
|
|
offerRun(ch, event)
|
|
}
|
|
}
|
|
|
|
func (b *Bus) PublishNode(event *iop.EdgeNodeEvent) {
|
|
if event == nil {
|
|
return
|
|
}
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
for ch := range b.allNodes {
|
|
offerNode(ch, event)
|
|
}
|
|
for ch := range b.nodeSubs[event.GetNodeId()] {
|
|
offerNode(ch, event)
|
|
}
|
|
}
|
|
|
|
func offerRun(ch chan *iop.RunEvent, event *iop.RunEvent) {
|
|
select {
|
|
case ch <- event:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func offerNode(ch chan *iop.EdgeNodeEvent, event *iop.EdgeNodeEvent) {
|
|
select {
|
|
case ch <- event:
|
|
default:
|
|
}
|
|
}
|