Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Broadcaster ¶
type Broadcaster[T any] struct { // contains filtered or unexported fields }
Broadcaster distributes events to multiple subscribers with bounded buffers.
func NewBroadcaster ¶
func NewBroadcaster[T any](metrics ...BroadcasterMetrics) *Broadcaster[T]
NewBroadcaster creates a new bounded broadcaster.
func (*Broadcaster[T]) Publish ¶
func (b *Broadcaster[T]) Publish(event Event[T])
Publish sends an event to all subscribers. If a subscriber's buffer is full, it is dropped.
func (*Broadcaster[T]) Subscribe ¶
func (b *Broadcaster[T]) Subscribe(ctx context.Context) *Subscriber[T]
Subscribe creates a new subscriber with a bounded channel. When the buffer is full, the subscriber is dropped (closed).
func (*Broadcaster[T]) SubscriberCount ¶
func (b *Broadcaster[T]) SubscriberCount() int
SubscriberCount returns number of active subscribers.
func (*Broadcaster[T]) Unsubscribe ¶
func (b *Broadcaster[T]) Unsubscribe(id string)
Unsubscribe removes a subscriber.
type BroadcasterMetrics ¶ added in v0.5.0
type BroadcasterMetrics struct {
IncSubscribers func()
DecSubscribers func()
IncEvents func()
IncDrops func()
}
BroadcasterMetrics defines optional callbacks for observability.
type Event ¶
type Event[T any] struct { Type string // ADDED, MODIFIED, DELETED Object T Version string // resourceVersion }
Event represents a watch event.
type InformerManager ¶
type InformerManager struct {
EnvBroadcaster *Broadcaster[*divergev1alpha1.Environment]
PgBroadcaster *Broadcaster[*divergev1alpha1.PreviewGroup]
// contains filtered or unexported fields
}
InformerManager manages shared informers and broadcasts events.
func NewInformerManager ¶
func NewInformerManager(logger *slog.Logger, metrics ...BroadcasterMetrics) *InformerManager
NewInformerManager creates informer event handlers and connects them to broadcasters.
func (*InformerManager) HandleEnvironmentEvent ¶
func (m *InformerManager) HandleEnvironmentEvent(eventType string, obj client.Object)
HandleEnvironmentEvent handles an environment event from an informer.
func (*InformerManager) HandlePreviewGroupEvent ¶
func (m *InformerManager) HandlePreviewGroupEvent(eventType string, obj client.Object)
HandlePreviewGroupEvent handles a preview group event from an informer.
type LogStreamer ¶
type LogStreamer struct {
// contains filtered or unexported fields
}
LogStreamer streams logs from pods.
func NewLogStreamer ¶
func NewLogStreamer(clientset kubernetes.Interface) *LogStreamer
func (*LogStreamer) StreamPodLogs ¶
func (l *LogStreamer) StreamPodLogs(ctx context.Context, namespace, podName, container string, follow bool, tailLines int64, since *time.Time) (io.ReadCloser, error)
StreamPodLogs streams logs from a specific pod/container.
type Subscriber ¶
type Subscriber[T any] struct { // contains filtered or unexported fields }
Subscriber is a registered watch subscriber.
func (*Subscriber[T]) Events ¶
func (s *Subscriber[T]) Events() <-chan Event[T]
Events returns the subscriber's event channel.