- internal/model: единый источник статусов задачи (Status, TraceStatus, машина переходов), без зависимости от storage. - internal/storage: совместимый мост (type Status = model.Status, re-export констант) — внешний код не меняется. - internal/events: шина событий (fan-out, блокирующий Publish с гарантией порядка), события задач/трейсов, отдельная логовая шина + LogWriter, Publisher/NilPublisher для внедрения в Core/Worker. - docs/ui-spec.md: спецификация десктопного UI (Fyne).
110 lines
3.6 KiB
Go
110 lines
3.6 KiB
Go
// Package events — шина событий для обмена UI ↔ Core.
|
||
//
|
||
// Две независимые шины: доменная (статусы задач, история, трейсы) и логовая
|
||
// (сырые строки лога для панели «Логи»). Обе построены на одном типе *Bus,
|
||
// DOMAIN шина блокирующая (гарантия доставки и порядка), логовая — та же,
|
||
// но с большим буфером, чтобы не тормозить логирование.
|
||
package events
|
||
|
||
import (
|
||
"sync"
|
||
)
|
||
|
||
// Event — доменное событие. Закрытый интерфейс: новые типы добавляются
|
||
// только внутри пакета.
|
||
type Event interface {
|
||
_event()
|
||
}
|
||
|
||
// Bus — широковещательная шина событий (fan-out).
|
||
//
|
||
// Publish блокирует вызывающую горутину до тех пор, пока все подписчики не
|
||
// получат событие (в копию их буфера). Порядок событий для каждого
|
||
// подписчика сохраняется. Удаление подписчика происходит горутиной-монтируется
|
||
// close(done), что снимает блокировку Publish.
|
||
type Bus struct {
|
||
bufSize int
|
||
|
||
mu sync.RWMutex
|
||
subs map[*subscriber]struct{}
|
||
}
|
||
|
||
type subscriber struct {
|
||
ch chan Event
|
||
done chan struct{}
|
||
}
|
||
|
||
// New создаёт шину с буфером bufSize на каждого подписчика.
|
||
func New(bufSize int) *Bus {
|
||
if bufSize < 1 {
|
||
bufSize = 1
|
||
}
|
||
return &Bus{
|
||
bufSize: bufSize,
|
||
subs: make(map[*subscriber]struct{}),
|
||
}
|
||
}
|
||
|
||
// Subscribe регистрирует нового подписчика и возвращает канал событий вместе
|
||
// с функцией отписки. Рекомендуемый паттерн потребления:
|
||
//
|
||
// ch, unsub := bus.Subscribe()
|
||
// defer unsub()
|
||
// for {
|
||
// select {
|
||
// case e := <-ch:
|
||
// switch ev := e.(type) { ... }
|
||
// case <-closeCh:
|
||
// return
|
||
// }
|
||
// }
|
||
//
|
||
// Канал не закрывается шиной: выход из горутины подписчика организуется через
|
||
// закрытие канала приложения либо другого сигнала в select.
|
||
func (b *Bus) Subscribe() (<-chan Event, func()) {
|
||
s := &subscriber{
|
||
ch: make(chan Event, b.bufSize),
|
||
done: make(chan struct{}),
|
||
}
|
||
|
||
b.mu.Lock()
|
||
b.subs[s] = struct{}{}
|
||
b.mu.Unlock()
|
||
|
||
var once sync.Once
|
||
unsubscribe := func() {
|
||
once.Do(func() {
|
||
close(s.done)
|
||
b.mu.Lock()
|
||
delete(b.subs, s)
|
||
b.mu.Unlock()
|
||
})
|
||
}
|
||
return s.ch, unsubscribe
|
||
}
|
||
|
||
// Publish рассылает событие всем подписчикам и блокируется, пока каждый
|
||
// подписчик либо примет событие (в свой буфер), либо отпишется. Порядок
|
||
// рассылки стабилен (по списку подписок). Безопасен для Concurrent вызовов.
|
||
func (b *Bus) Publish(e Event) {
|
||
b.mu.RLock()
|
||
subs := make([]*subscriber, 0, len(b.subs))
|
||
for s := range b.subs {
|
||
subs = append(subs, s)
|
||
}
|
||
b.mu.RUnlock()
|
||
|
||
for _, s := range subs {
|
||
select {
|
||
case s.ch <- e:
|
||
case <-s.done:
|
||
}
|
||
}
|
||
}
|
||
|
||
// SubscribersCount — число активных подписчиков (для юнит-тестов и отладки).
|
||
func (b *Bus) SubscribersCount() int {
|
||
b.mu.RLock()
|
||
defer b.mu.RUnlock()
|
||
return len(b.subs)
|
||
} |