Files
ITWorx Pulse release export bd774932d5
Public source validation / validate (push) Failing after 3m8s
Publish ITWorx Pulse source
2026-09-03 02:09:19 +02:00

244 lines
5.4 KiB
Go

package live
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"sync"
"time"
"github.com/itworx/pulse/internal/queryplan"
)
const (
defaultRegistryEntries = 256
defaultRegistryConcurrency = 16
)
var (
ErrRegistryLimit = errors.New("live subscription registry limit reached")
ErrLeaseReleased = errors.New("live subscription lease is released")
ErrOutboundBackpressure = errors.New("live outbound queue is full")
ErrSamplerUnavailable = errors.New("live sampler is not configured")
)
type RegistryOptions struct {
MaxEntries int
MaxConcurrent int
}
type Registry struct {
sampler Sampler
now func() time.Time
maxEntries int
sem chan struct{}
mu sync.Mutex
entries map[string]*registryEntry
}
type registryEntry struct {
key string
request queryplan.Request
references int
minInterval time.Duration
lastAt time.Time
lastSamples []Sample
inFlight *sampleCall
}
type sampleCall struct {
done chan struct{}
samples []Sample
err error
}
type Lease struct {
registry *Registry
key string
once sync.Once
}
// NewRegistry builds the shared live subscription registry. A nil sampler keeps
// subscription lifecycle traffic working but is not a usable data source: every
// Sample then fails with ErrSamplerUnavailable instead of silently reporting an
// empty successful result.
func NewRegistry(sampler Sampler, options RegistryOptions) *Registry {
maxEntries := options.MaxEntries
if maxEntries <= 0 {
maxEntries = defaultRegistryEntries
}
maxConcurrent := options.MaxConcurrent
if maxConcurrent <= 0 {
maxConcurrent = defaultRegistryConcurrency
}
return &Registry{
sampler: sampler,
now: func() time.Time { return time.Now().UTC() },
maxEntries: maxEntries,
sem: make(chan struct{}, maxConcurrent),
entries: make(map[string]*registryEntry),
}
}
func (r *Registry) Acquire(request queryplan.Request, interval time.Duration) (*Lease, error) {
if r == nil || interval <= 0 {
return nil, ErrLeaseReleased
}
key, err := normalizedKey(request)
if err != nil {
return nil, err
}
r.mu.Lock()
defer r.mu.Unlock()
entry, exists := r.entries[key]
if !exists {
if len(r.entries) >= r.maxEntries {
return nil, ErrRegistryLimit
}
entry = &registryEntry{key: key, request: request, minInterval: interval}
r.entries[key] = entry
} else if interval < entry.minInterval {
entry.minInterval = interval
}
entry.references++
return &Lease{registry: r, key: key}, nil
}
func (l *Lease) Release() {
if l == nil || l.registry == nil {
return
}
l.once.Do(func() {
l.registry.release(l.key)
})
}
func (r *Registry) release(key string) {
r.mu.Lock()
defer r.mu.Unlock()
entry, ok := r.entries[key]
if !ok {
return
}
if entry.references > 0 {
entry.references--
}
if entry.references == 0 && entry.inFlight == nil {
delete(r.entries, key)
}
}
func (r *Registry) Sample(ctx context.Context, request queryplan.Request) ([]Sample, error) {
if r == nil {
return nil, ErrLeaseReleased
}
if r.sampler == nil {
return nil, ErrSamplerUnavailable
}
if err := ctx.Err(); err != nil {
return nil, err
}
key, err := normalizedKey(request)
if err != nil {
return nil, err
}
r.mu.Lock()
entry, ok := r.entries[key]
if !ok || entry.references == 0 {
r.mu.Unlock()
return nil, ErrLeaseReleased
}
now := r.now()
if !entry.lastAt.IsZero() && now.Sub(entry.lastAt) < entry.minInterval {
samples := cloneSamples(entry.lastSamples)
r.mu.Unlock()
return samples, nil
}
if entry.inFlight != nil {
call := entry.inFlight
r.mu.Unlock()
select {
case <-call.done:
return cloneSamples(call.samples), call.err
case <-ctx.Done():
return nil, ctx.Err()
}
}
call := &sampleCall{done: make(chan struct{})}
entry.inFlight = call
requestCopy := entry.request
r.mu.Unlock()
select {
case r.sem <- struct{}{}:
case <-ctx.Done():
r.finish(key, entry, call, nil, ctx.Err())
return nil, ctx.Err()
}
samples, sampleErr := r.sampler.Sample(ctx, requestCopy)
<-r.sem
r.finish(key, entry, call, samples, sampleErr)
return cloneSamples(samples), sampleErr
}
func (r *Registry) finish(key string, entry *registryEntry, call *sampleCall, samples []Sample, err error) {
r.mu.Lock()
call.samples = cloneSamples(samples)
call.err = err
if err == nil {
entry.lastAt = r.now()
entry.lastSamples = cloneSamples(samples)
}
entry.inFlight = nil
if entry.references == 0 {
delete(r.entries, key)
}
close(call.done)
r.mu.Unlock()
}
func (r *Registry) Active() (entries, references int) {
if r == nil {
return 0, 0
}
r.mu.Lock()
defer r.mu.Unlock()
for _, entry := range r.entries {
entries++
references += entry.references
}
return entries, references
}
func normalizedKey(request queryplan.Request) (string, error) {
payload, err := json.Marshal(request)
if err != nil {
return "", err
}
digest := sha256.Sum256(payload)
return hex.EncodeToString(digest[:]), nil
}
func cloneSamples(samples []Sample) []Sample {
if samples == nil {
return nil
}
cloned := make([]Sample, len(samples))
for index, sample := range samples {
cloned[index] = sample
if sample.Value != nil {
value := *sample.Value
cloned[index].Value = &value
}
if sample.Labels != nil {
cloned[index].Labels = make(map[string]string, len(sample.Labels))
for key, value := range sample.Labels {
cloned[index].Labels[key] = value
}
}
}
return cloned
}