Public source validation / validate (push) Failing after 3m8s
244 lines
5.4 KiB
Go
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 = ®istryEntry{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
|
|
}
|