Skip to content
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,14 @@ Part of the receiver API is compatible with DataDog, those parts are extracted h

Provide exactly one authentication source: `APIKey` or `ServiceAccountToken`. The token callback is read on every request, including after rotation; empty credentials and header delimiters are rejected. The existing `NewOpenAPIClient` signature, `Connect()` method and legacy authentication behavior remain available unchanged.

### Capability queries

Feature queries bound attempts, elapsed time and response bodies (1 MiB). Set `QueryOptions.BooleanCapabilities` to the keys the consumer requires as booleans, such as `otel-logs` or `k8s-rbac`. Missing keys are valid, unrelated values are preserved, and an empty list checks only the response object. A 404 means the endpoint was not found; it cannot distinguish an older Receiver from an incorrect base path. Query results expose status and outcome without raw URLs, errors or bodies.

Proxy authentication rejection (407), including HTTPS CONNECT, is a configuration failure. CONNECT 5xx responses remain transient; other permanent CONNECT rejections do not imply an unsupported Receiver.

`Client.StartPolling` replaces `StartFeaturesPoller`. Callers must migrate to the cancellable poller and handle classified results.

### Bumping the openapi version

- Change the version/branch/commit sha in the `stackstate_openapi/openapi_version` file
Expand Down
222 changes: 222 additions & 0 deletions pkg/openapiclient/features/features_client.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,222 @@
package features

import (
"context"
"crypto/tls"
"crypto/x509"
"errors"
"io"
"math/rand/v2"
"net"
"net/http"
"strconv"
"strings"
"time"

"github.com/StackVista/stackstate-receiver-go-client/generated/receiver_api"
"github.com/StackVista/stackstate-receiver-go-client/pkg/openapiclient"
)

// QueryOptions bounds a complete feature query, including attempts and retry waits.
type QueryOptions struct {
Timeout, AttemptTimeout time.Duration
MaxAttempts int
InitialBackoff, MaxBackoff time.Duration
// BooleanCapabilities validates only listed keys that are present in the response.
BooleanCapabilities []string
}

// Client queries the generated features API with bounded retries.
type Client struct {
api receiver_api.FeaturesAPI
opts QueryOptions
now func() time.Time
random func() float64
}

// NewClient validates query bounds and constructs a feature client.
func NewClient(api receiver_api.FeaturesAPI, opts QueryOptions) (*Client, error) {
if api == nil {
return nil, errors.New("features API is required")
}
if opts.Timeout <= 0 || opts.AttemptTimeout <= 0 || opts.AttemptTimeout > opts.Timeout || opts.MaxAttempts < 1 || opts.InitialBackoff <= 0 || opts.MaxBackoff < opts.InitialBackoff {
return nil, errors.New("invalid feature query timeout, attempt count or backoff bounds")
}
opts.BooleanCapabilities = append([]string(nil), opts.BooleanCapabilities...)
return &Client{api: api, opts: opts, now: time.Now, random: rand.Float64}, nil
}

// FetchFeatures returns one observation after the bounded query completes.
func (c *Client) FetchFeatures(authCtx context.Context) Result {
queryCtx, cancel := context.WithTimeout(authCtx, c.opts.Timeout)
defer cancel()
result := Result{}
backoff := c.opts.InitialBackoff
for {
if queryCtx.Err() != nil {
result.Class = contextClass(authCtx, queryCtx.Err())
break
}
attemptCtx, stop := context.WithTimeout(queryCtx, c.opts.AttemptTimeout)
values, response, err := c.api.GetFeaturesExecute(c.api.GetFeatures(attemptCtx))
result.Attempts++
result.StatusCode = 0
var proxyError *openapiclient.ProxyConnectError
if errors.As(err, &proxyError) {
result.StatusCode = proxyError.StatusCode
}
retryAfter := time.Duration(0)
if response != nil {
result.StatusCode = response.StatusCode
if response.StatusCode == http.StatusTooManyRequests || response.StatusCode == http.StatusServiceUnavailable {
retryAfter = parseRetryAfter(response.Header.Get("Retry-After"), c.now(), c.opts.Timeout)
}
if response.Body != nil {
response.Body.Close()
}
}
result.Class = c.classify(authCtx, attemptCtx, values, response, err)
stop()
if result.Class == Valid {
result.Features = values
}
if (result.Class != Transient && result.Class != Timeout) || result.Attempts >= c.opts.MaxAttempts {
break
}
if queryCtx.Err() != nil {
result.Class = contextClass(authCtx, queryCtx.Err())
break
}
delay := time.Duration(c.random() * float64(backoff))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is full jitter, so a run of small random values retries almost immediately and the retry bound degrades to MaxAttempts back-to-back requests. A floor of backoff/2 would keep the spacing.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Full jitter is the selected retry policy, so near-zero delays are allowed. MaxAttempts and the whole-query deadline still bound requests, and Retry-After remains a minimum when present. Existing retry/deadline/cancellation tests pass; no minimum spacing is promised for retries.

if retryAfter > delay {
delay = retryAfter
}
deadline, _ := queryCtx.Deadline()
if delay >= time.Until(deadline) {
break
}
timer := time.NewTimer(delay)
select {
case <-queryCtx.Done():
timer.Stop()
result.Class = contextClass(authCtx, queryCtx.Err())
result.FinishedAt = c.now()
return result
case <-timer.C:
}
if backoff > c.opts.MaxBackoff/2 {
backoff = c.opts.MaxBackoff
} else {
backoff *= 2
}
}
result.FinishedAt = c.now()
return result
}

func contextClass(parent context.Context, err error) Class {
if parent.Err() != nil {
return Canceled
}
if errors.Is(err, context.DeadlineExceeded) {
return Timeout
}
return Canceled
}

func (c *Client) classify(parent, attempt context.Context, values map[string]any, response *http.Response, err error) Class {
if parent.Err() != nil {
return Canceled
}
var proxyError *openapiclient.ProxyConnectError
if errors.As(err, &proxyError) {
switch status := proxyError.StatusCode; {
case status == http.StatusProxyAuthRequired:
return Configuration
case status == 408 || status == 429 || status >= 500 && status <= 599:
return Transient
default:
return Rejected
}
}
if response != nil {
switch status := response.StatusCode; {
case status == 401 || status == 403:
return Authentication
case status == http.StatusProxyAuthRequired:
return Configuration
case status == 404:
return Unsupported
Comment thread
craffit marked this conversation as resolved.
case status == 408 || status == 429 || status >= 500 && status <= 599:
return Transient
case status != 200:
return Rejected
}
}
if errors.Is(err, openapiclient.ErrMissingCredential) {
return Authentication
}
var verification *tls.CertificateVerificationError
var unknown x509.UnknownAuthorityError
var hostname x509.HostnameError
var invalid x509.CertificateInvalidError
if errors.As(err, &verification) || errors.As(err, &unknown) || errors.As(err, &hostname) || errors.As(err, &invalid) {
return Configuration
}
if attempt.Err() != nil || errors.Is(err, context.DeadlineExceeded) {
return Timeout
}
if errors.Is(err, context.Canceled) {
return Canceled
}
var network net.Error
if errors.As(err, &network) && network.Timeout() {
return Timeout
}
var operation *net.OpError
if errors.As(err, &operation) || errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) || (errors.As(err, &network) && network.Temporary()) {
return Transient
}
if response == nil {
Comment thread
craffit marked this conversation as resolved.
if err != nil {
return Transient
}
return Rejected
}
if err != nil || values == nil {
return Malformed
}
for _, key := range c.opts.BooleanCapabilities {
if capability, present := values[key]; present {
if _, ok := capability.(bool); !ok {
return Malformed
}
}
}
return Valid
}

func parseRetryAfter(value string, now time.Time, limit time.Duration) time.Duration {
value = strings.TrimSpace(value)
if value == "" {
return 0
}
if seconds, err := strconv.ParseUint(value, 10, 64); err == nil {
if seconds > uint64(limit/time.Second) {
return limit
}
return time.Duration(seconds) * time.Second
} else if errors.Is(err, strconv.ErrRange) && strings.Trim(value, "0123456789") == "" {
return limit
}
if deadline, err := http.ParseTime(value); err == nil {
delay := deadline.Sub(now)
if delay > limit {
return limit
}
if delay > 0 {
return delay
}
}
return 0
}
Loading
Loading