Documentation
¶
Overview ¶
Package curo adds adaptive resilience to outbound HTTP calls.
A Transport is an http.RoundTripper that wraps a base transport, such as http.DefaultTransport. It observes the requests sent through it, learns how each dependency behaves, and in Enforce mode applies three controls that need no tuning:
- a retry of a replay-safe read that failed transiently, within budgets that allow about one retry per ten requests;
- a dependency breaker that fails requests fast with ErrBreakerOpen once a dependency is diagnosed as down; and
- an adaptive timeout that ends a read with ErrTimeout when its response headers take too long.
Wrap the base transport once, share the client, and close the Transport when the application shuts down:
transport, err := curo.New(http.DefaultTransport)
if err != nil {
return err
}
client := &http.Client{Transport: transport, Timeout: time.Minute}
Modes ¶
A new Transport starts in Observe: it sends every request unchanged and reports what it would do. Enforce applies the controls, and Off passes requests straight to the base transport. Transport.SetMode changes the mode at runtime, so a rollback needs no deploy. Apart from the timeout bounds, set with WithTimeoutBounds, every threshold is fixed.
Errors ¶
http.Client wraps the errors that Curo returns in a *url.Error, so match them with errors.Is. ErrTimeout also matches context.DeadlineExceeded, so check it first.
Visibility ¶
Transport.Stats returns aggregate counters, and Transport.Report returns each dependency's latest decision with its evidence and reasons. Both encode to JSON.
Safety ¶
Curo runs inside the application it protects. A failure inside Curo never fails a request, and repeated failures make the Transport pass requests straight through. Transport.Stats counts these failures, and WithLogger logs where each one happened. Curo never retries a request that is not replay-safe, never extends a caller's deadline, and keeps its state bounded however many hosts the application calls.
The operations guide explains how to run Curo in a service, and the behavior reference describes each mechanism in detail.
Example (ErrorHandling) ¶
package main
import (
"context"
"errors"
"fmt"
"net/url"
"github.com/raj1kshtz/curo"
)
func main() {
describe := func(err error) string {
switch {
case err == nil:
return "succeeded"
case errors.Is(err, curo.ErrBreakerOpen):
return "not sent: the dependency is down, so use a fallback"
case errors.Is(err, curo.ErrTimeout):
return "ended by Curo: the dependency may have received it"
case errors.Is(err, context.DeadlineExceeded):
return "the caller's deadline passed"
default:
return "failed"
}
}
// http.Client wraps transport errors in *url.Error, which errors.Is
// unwraps. ErrTimeout also matches context.DeadlineExceeded, so check it
// first.
for _, cause := range []error{
curo.ErrBreakerOpen,
curo.ErrTimeout,
context.DeadlineExceeded,
} {
err := &url.Error{
Op: "Get",
URL: "https://api.example/orders",
Err: cause,
}
fmt.Println(describe(err))
}
}
Output: not sent: the dependency is down, so use a fallback ended by Curo: the dependency may have received it the caller's deadline passed
Index ¶
Examples ¶
Constants ¶
This section is empty.
Variables ¶
var ( // ErrNilBaseTransport indicates that a Transport has no base RoundTripper. ErrNilBaseTransport = errors.New("curo: base transport is nil") // ErrInvalidMode indicates that a Mode is not one of Off, Observe, or Enforce. ErrInvalidMode = errors.New("curo: invalid mode") // ErrClosed indicates that an operation requires an open Transport. ErrClosed = errors.New("curo: transport is closed") // ErrBreakerOpen indicates that Enforce failed a request fast because the // dependency breaker of the request's target is open. The base transport // was not called, so the request was not sent. ErrBreakerOpen = errors.New("curo: dependency breaker is open") // ErrTimeout indicates that Enforce ended a request because the base // transport did not return response headers within the target's adaptive // timeout. The request may have reached the dependency. // // Like the timeout error of an http.Client, ErrTimeout reports true from // Timeout, and errors.Is(ErrTimeout, context.DeadlineExceeded) is true. ErrTimeout error = &timeoutError{} )
Functions ¶
This section is empty.
Types ¶
type Candidates ¶
type Candidates uint8
Candidates is a set of controls the policy considers eligible.
const ( // CandidateRetry marks transient dependency failures as eligible for // budgeted retries of replay-safe requests. Enforce mode applies it. CandidateRetry Candidates = 1 << iota // CandidateBreakerOpen marks a dependency that is down as eligible for // opening its dependency breaker. Enforce mode applies it, and requests // then fail fast with ErrBreakerOpen until a probe shows that the // dependency answers again. CandidateBreakerOpen // CandidateTimeout marks a target whose retained latency evidence selects // the adaptive timeout in Decision.Timeout. It does not depend on the // diagnosis. Enforce mode applies it to read requests other than breaker // probes, which fail with ErrTimeout when the base transport does not // return response headers in time. CandidateTimeout )
func (Candidates) Has ¶
func (candidates Candidates) Has(want Candidates) bool
Has reports whether candidates contains every control in want. It reports false when want is empty.
func (Candidates) MarshalText ¶
func (candidates Candidates) MarshalText() ([]byte, error)
MarshalText encodes the set as String does, such as "Retry|Timeout" or "None". It never returns an error.
func (Candidates) String ¶
func (candidates Candidates) String() string
String returns candidate names joined by "|", or "None" for an empty set.
type Change ¶
type Change struct {
// Decision is the evaluation that produced the change.
Decision Decision
// Sequence orders changes within one Transport. It starts at 1 and
// increases by one per change, so a gap between reads means older changes
// were overwritten.
Sequence uint64
}
Change is a retained evaluation that changed a target's candidates.
A change is recorded when the candidate set or the Timeout changes, or when the same CandidateRetry or CandidateBreakerOpen selection is justified by a different diagnosis. CandidateTimeout does not depend on the diagnosis, so a new diagnosis with no other candidate is not a change. Renewing an unchanged decision is not a change. Changes are not a target lifecycle log: target eviction and evidence going stale are not recorded.
type Decision ¶
type Decision struct {
// Target identifies the normalized dependency.
Target Target
// EvaluatedAt is when the evaluation ran.
EvaluatedAt time.Time
// ExpiresAt is when the evaluation stops describing current evidence.
// Report never re-evaluates, so an idle target keeps its last decision
// after ExpiresAt.
ExpiresAt time.Time
// Reasons explains readiness, diagnosis, and candidate selection in
// evaluation order. It is nil when there is no reason.
Reasons []Reason
// Recent summarizes the recent evidence window used by the evaluation.
Recent Evidence
// Historical summarizes the non-overlapping historical window used by the
// evaluation.
Historical Evidence
// Latency summarizes the latency samples the adaptive timeout was derived
// from.
Latency Latency
// Timeout is the adaptive timeout that Enforce applies to the target's
// read requests while the decision is current. It is non-zero exactly when
// Candidates has CandidateTimeout.
Timeout time.Duration
// PolicyVersion identifies the rules that selected Candidates.
PolicyVersion uint32
// Readiness reports whether evidence was sufficient for a diagnosis.
Readiness Readiness
// Diagnosis is the deterministic diagnosis class.
Diagnosis Diagnosis
// Candidates holds controls the policy considers eligible. Enforce mode
// applies CandidateRetry, subject to replay safety, the caller's deadline,
// and retry budgets. It applies CandidateBreakerOpen by opening the
// target's dependency breaker when a request fails because of the
// dependency. It applies CandidateTimeout by ending a read request that
// waits longer than Timeout for response headers.
Candidates Candidates
}
Decision is one published evaluation of a target.
A target without an evaluation has zero times, ReadinessCold, DiagnosisNone, no reasons, zero evidence and latency, no candidates, a zero Timeout, and a zero PolicyVersion.
type Diagnosis ¶
type Diagnosis uint8
Diagnosis is a closed diagnosis class.
const ( // DiagnosisNone means the evidence does not support a diagnosis. DiagnosisNone Diagnosis = iota // DiagnosisHealthy means ready evidence is within bounds. DiagnosisHealthy // DiagnosisTransient means isolated dependency failures are present. DiagnosisTransient // DiagnosisDependencyDown means dependency failures dominate recent // attempts. DiagnosisDependencyDown // DiagnosisSaturation means the dependency is persistently rate limiting. DiagnosisSaturation // DiagnosisClientError means request-side HTTP failures dominate recent // attempts. DiagnosisClientError // DiagnosisDegrading means recent failure rate or latency moved materially // above the historical baseline. DiagnosisDegrading )
func (Diagnosis) MarshalText ¶
MarshalText encodes the diagnosis as its name. It never returns an error.
type Evidence ¶
type Evidence struct {
// Attempts is the number of dependency-relevant attempts.
Attempts uint64
// DependencyFailures counts transport failures and HTTP 5xx responses.
DependencyFailures uint64
// RateLimited counts HTTP 429 responses.
RateLimited uint64
// ClientFailures counts HTTP 4xx responses other than 429.
ClientFailures uint64
// Span is the time between the first and last counted attempts.
Span time.Duration
}
Evidence summarizes the dependency-relevant attempts in one window.
Caller-owned cancellations and caller deadlines are not dependency relevant and are excluded from every count.
type Latency ¶
type Latency struct {
// Samples is the number of retained samples.
Samples uint64
// Slowest is the upper bound of the bucket that holds the slowest retained
// sample. It is the largest Duration when that sample took longer than 30
// seconds, and zero without samples.
Slowest time.Duration
}
Latency summarizes the latency samples a target retains for its adaptive timeout, from the recent window and the historical baseline, which together cover about the last 30 minutes.
Every initial attempt that returned a response is a sample of its time to response headers. An attempt that the adaptive timeout ended is a sample at the timeout. Samples are kept in fixed buckets, the slowest of which holds everything above 30 seconds.
type MethodClass ¶
type MethodClass uint8
MethodClass is a bounded class of HTTP request methods.
const ( // MethodRead covers GET, HEAD, OPTIONS, TRACE, and an empty method. MethodRead MethodClass = iota // MethodWrite covers POST, PUT, PATCH, and DELETE. MethodWrite // MethodConnect covers CONNECT. MethodConnect // MethodOther covers every other method, including non-canonical case. MethodOther )
func (MethodClass) MarshalText ¶
func (class MethodClass) MarshalText() ([]byte, error)
MarshalText encodes the method class as its name. It never returns an error.
func (MethodClass) String ¶
func (class MethodClass) String() string
String returns the method class name.
type Mode ¶
type Mode uint32
Mode controls how much authority a Transport has over a request.
A Mode prints and encodes as its name, and UnmarshalText accepts the name in any letter case, so flag.TextVar, encoding/json, and other text decoders can set a Mode from configuration or an admin control.
const ( // Off delegates requests without collecting evidence or applying actions. Off Mode = iota // Observe permits evidence collection but does not permit request changes. Observe // Enforce permits eligible actions within fixed safety bounds: at most one // budgeted retry of a replay-safe request whose target has a current Retry // candidate, failing requests fast with ErrBreakerOpen while the target's // dependency breaker is open, and ending a read request with ErrTimeout // when response headers do not arrive within its target's adaptive // timeout. Enforce )
func (Mode) MarshalText ¶
MarshalText encodes the mode as its name. It returns an error wrapping ErrInvalidMode for an invalid mode.
func (Mode) String ¶
String returns the mode name: "Off", "Observe", or "Enforce". An invalid mode returns "Mode(N)", where N is its number.
func (*Mode) UnmarshalText ¶
UnmarshalText sets the mode from its name in any letter case, such as "off", "Observe", or "ENFORCE". For any other text, it returns an error wrapping ErrInvalidMode and leaves the mode unchanged. A nil *Mode also returns an error wrapping ErrInvalidMode.
type Option ¶
type Option interface {
// contains filtered or unexported methods
}
Option configures a Transport during construction.
Options are created by this package's With functions.
func WithLogger ¶
WithLogger logs Curo's internal failures to logger. By default, Curo logs nothing.
Curo contains a panic or an error in its own work, including application code that the work calls, such as a body's Close method or an error's Unwrap method. It offers each failure to logger as a Warn record with these attributes:
- stage: the work that failed: preflight, postflight, retry, timeout, body, or report
- kind: panic, or error when the work returned an error
- type: the Go type of the panic value or error. Types that package reflect builds at run time can hold application data, so an unnamed struct, function, or channel type is written as struct, func, or chan, and an array type as [...]T, without its length. A name is cut after 16 levels of nested types or at 512 bytes, and then ends with "...".
- error: the message of a Go runtime error whose message is fixed text, such as a nil pointer dereference. Other messages can hold application data, as an index out of range error holds the index and a failed type assertion names types, so they are omitted.
- failures: the failure's number, from the count that Stats reports as InternalFailures
- skipped: the number of failures since the previous failure record that were not offered to logger, present only when some were
- stack: for a panic, the functions and source lines of the goroutine that panicked, from the panic outward, without argument values, and truncated to whole frames within 8 KiB
When three internal failures within a minute disable the Transport, it offers that once as an Error record with a failures attribute.
A Transport offers one record at a time, and at most three failure records a minute. A failure that happens while a record is offered is skipped, including one that the handler causes by calling the Transport, so the handler is never called reentrantly. The per-minute limit skips no failure before self-disable, only failures in the work that continues afterwards: settling the timeouts of requests in progress, closing bodies, and building reports.
Records are offered synchronously, outside Curo's locks, with context.Background. A failure record is offered on the goroutine where the failure happened, and the self-disable record follows a failure record on the same goroutine. The logger's level and handler decide which records are written. A panic from the handler is recovered and ignored, and does not count as an internal failure.
A nil logger is rejected.
Example ¶
package main
import (
"fmt"
"log/slog"
"net/http"
"github.com/raj1kshtz/curo"
)
func main() {
// Curo writes records only when it contains a failure in its own work, so
// a healthy Transport logs nothing.
transport, err := curo.New(
http.DefaultTransport,
curo.WithMode(curo.Enforce),
curo.WithLogger(slog.Default()),
)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
fmt.Println(transport.Stats().InternalFailures)
_, err = curo.New(http.DefaultTransport, curo.WithLogger(nil))
fmt.Println(err)
}
Output: 0 curo: apply option 1: nil logger
func WithTimeoutBounds ¶
WithTimeoutBounds sets the floor and ceiling of adaptive timeouts. The default floor is 2 seconds and the default ceiling is 30 seconds.
In Enforce mode, a read request to a target with enough latency evidence must receive response headers within the target's adaptive timeout: three times its slowest retained latency, clamped to these bounds. Reading the response body is not bounded. Set the ceiling above the longest time to response headers that a read may legitimately take, because no read may wait longer. Also set it below the deadlines that callers put on those reads: a read is timed only when its deadline leaves more time than its timeout, and only a timeout at the ceiling counts as a dependency failure. A read that its caller's deadline ends is never a dependency failure, even when the dependency hangs.
Passing zero for both disables adaptive timeouts. Otherwise minimum must be positive and must not exceed maximum.
Example ¶
package main
import (
"fmt"
"net/http"
"time"
"github.com/raj1kshtz/curo"
)
func main() {
// Keep the ceiling below the client's timeout. A hung read then ends at
// the ceiling, which counts as a dependency failure and can open the
// breaker, rather than at the client's deadline, which counts as the
// caller giving up.
transport, err := curo.New(
http.DefaultTransport,
curo.WithMode(curo.Enforce),
curo.WithTimeoutBounds(2*time.Second, 9*time.Second),
)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
client := &http.Client{
Transport: transport,
Timeout: 10 * time.Second,
}
fmt.Println(client.Transport == transport)
fmt.Println(client.Timeout)
fmt.Println(transport.Stats().Timeouts)
_, err = curo.New(
http.DefaultTransport,
curo.WithTimeoutBounds(10*time.Second, time.Second),
)
fmt.Println(err != nil)
}
Output: true 10s 0 true
type Readiness ¶
type Readiness uint8
Readiness reports whether retained evidence can support a diagnosis.
const ( // ReadinessCold means no dependency-relevant attempt has been observed. ReadinessCold Readiness = iota // ReadinessWarming means recent evidence exists but is not sufficient. ReadinessWarming // ReadinessReady means recent or historical evidence is sufficient. ReadinessReady // ReadinessStale means earlier evidence exists but no recent attempt // remains. ReadinessStale )
func (Readiness) MarshalText ¶
MarshalText encodes the readiness as its name. It never returns an error.
type Reason ¶
type Reason uint8
Reason is a stable explanation for one part of a Decision.
const ( // ReasonNoEvidence means no dependency-relevant attempt has been observed. ReasonNoEvidence Reason = iota + 1 // ReasonEvidenceStale means no recent dependency-relevant attempt remains. ReasonEvidenceStale // ReasonInsufficientEvidence means recent evidence has not reached the // readiness volume or span. ReasonInsufficientEvidence // ReasonRecentEvidenceReady means recent volume and span are sufficient. ReasonRecentEvidenceReady // ReasonHistoricalBaselineReady means historical evidence is sufficient. ReasonHistoricalBaselineReady // ReasonInsufficientRecentEvidence means recent evidence supports no // diagnosis class. ReasonInsufficientRecentEvidence // ReasonClientFailureRate means request-side HTTP failures dominate recent // attempts. ReasonClientFailureRate // ReasonRateLimitRate means HTTP 429 responses are sustained. ReasonRateLimitRate // ReasonDependencyFailureRate means dependency failures dominate recent // attempts. ReasonDependencyFailureRate // ReasonFailureRateIncrease means the recent dependency failure rate rose // materially above the historical baseline. ReasonFailureRateIncrease // ReasonLatencyIncrease means recent p95 latency rose materially above the // historical baseline. ReasonLatencyIncrease // ReasonIsolatedDependencyFailure means dependency failures are present // without dominating recent attempts. ReasonIsolatedDependencyFailure // ReasonWithinBaseline means ready evidence shows no failure pattern. ReasonWithinBaseline // ReasonReadinessRequired means the diagnosis maps to a candidate, but // evidence is not Ready, so no candidate was selected. ReasonReadinessRequired )
func (Reason) MarshalText ¶
MarshalText encodes the reason as its name, so encoding/json renders Decision.Reasons as an array of names. It never returns an error.
type Report ¶
type Report struct {
// Targets holds the latest decision for each tracked regular target,
// ordered by scheme, host, port, and method class. It holds at most 128
// entries. The overflow aggregate is never reported.
Targets []Decision
// Changes holds the most recent candidate changes in ascending Sequence
// order. It holds at most 256 entries.
Changes []Change
}
Report is a detached snapshot of published target decisions and recent candidate changes for one Transport.
Reports contain normalized target identities and bounded counters. They never include paths, queries, fragments, user information, headers, bodies, raw errors, or response metadata.
A Report encodes to JSON with names for its enumerated values, such as "Ready", and with Reasons as an array of names. The encoding is output only: the enumerated types implement encoding.TextMarshaler but not encoding.TextUnmarshaler. Later versions may add names and fields, so consumers should decode names as strings and tolerate unknown names and fields.
type Stats ¶
type Stats struct {
// ObservedRequests is the number of completed initial attempts recorded in
// Observe or Enforce mode. Retry attempts and requests that failed fast
// with ErrBreakerOpen are not included.
ObservedRequests uint64
// TrackedTargets is the current number of admitted regular targets. It does
// not include the fixed overflow aggregate.
TrackedTargets uint64
// OverflowRequests is the number of observations assigned to the bounded
// non-actionable overflow aggregate.
OverflowRequests uint64
// RetryAttempts is the number of retry attempts Curo started in Enforce
// mode. Each is one additional base transport call, and a request starts
// at most one. Retries the base transport performs internally are not
// counted.
RetryAttempts uint64
// RetrySuccesses is the number of retry attempts whose base transport call
// returned a 2xx or 3xx response without an error. Reading that response
// body can still fail.
RetrySuccesses uint64
// RetryBudgetDenials is the number of Enforce retries that met every other
// condition after the initial attempt but did not start because the target
// or instance retry budget was empty. Observe funds retry budgets but never
// reserves from them, so it records no denials.
RetryBudgetDenials uint64
// BreakerOpens is the number of times Enforce opened a closed dependency
// breaker. A failed probe keeps a breaker open without counting another
// open.
BreakerOpens uint64
// BreakerProbes is the number of requests that open dependency breakers
// sent to the base transport as probes after a cooldown.
BreakerProbes uint64
// BreakerRejections is the number of requests that failed fast with
// ErrBreakerOpen without calling the base transport.
BreakerRejections uint64
// Timeouts is the number of attempts, including retry attempts, that
// Enforce ended with ErrTimeout because the base transport did not return
// response headers within the adaptive timeout.
Timeouts uint64
// ShadowTimeouts is the number of Observe initial attempts that took
// longer than the adaptive timeout Enforce would have applied to them.
// Observe never ends an attempt.
ShadowTimeouts uint64
// InternalFailures is the number of contained Curo-owned failures in
// request stages and Report. [WithLogger] offers them to its logger,
// within the limits that it describes.
InternalFailures uint64
// SelfDisabled reports whether repeated internal failures permanently
// selected direct pass-through for this Transport. [WithLogger] offers the
// change to its logger once.
SelfDisabled bool
}
Stats is a privacy-safe aggregate snapshot of one Transport.
Stats never includes target identifiers, URLs, headers, bodies, or raw errors. Fields are sampled independently while requests are in flight. OverflowRequests never exceeds ObservedRequests, RetrySuccesses never exceeds RetryAttempts, RetryAttempts plus RetryBudgetDenials never exceeds ObservedRequests, ShadowTimeouts never exceeds ObservedRequests, and SelfDisabled implies at least three InternalFailures.
type Target ¶
type Target struct {
// Scheme is "http" or "https".
Scheme string
// Host is a lowercase DNS name without a trailing dot, or a canonical IP
// address.
Host string
// Port is the explicit port or the scheme default.
Port uint16
// Method is the request method class.
Method MethodClass
}
Target is a normalized dependency identity.
Host comes from the request URL and can be influenced by untrusted input, for example when an application calls user-supplied URLs. Do not use it as an unbounded metric label.
type Transport ¶
type Transport struct {
// contains filtered or unexported fields
}
Transport is an explicit, concurrency-safe wrapper around a host-owned http.RoundTripper.
Off delegates every request to the base transport exactly once without mutation. Observe and Enforce collect bounded result evidence behind narrow failure boundaries without wrapping the host-owned transport call, evaluate deterministic diagnoses and control candidates, and publish them through Report. Enforce also applies three candidates. For Retry, when every retry condition holds, it starts one budgeted retry of a replay-safe request with a clone of the original request. For BreakerOpen, it opens the target's dependency breaker, so later requests to that target fail fast with ErrBreakerOpen until a probe request gets an answer from the dependency. For Timeout, it ends a read request that waits longer than the target's adaptive timeout for response headers, and returns ErrTimeout.
A Transport must not be copied after first use.
func New ¶
func New(base http.RoundTripper, options ...Option) (*Transport, error)
New constructs a Transport around base.
New does not take ownership of base. Closing the returned Transport never closes base. A nil base is rejected so fallback ownership remains explicit.
Example ¶
package main
import (
"fmt"
"net/http"
"github.com/raj1kshtz/curo"
)
func main() {
transport, err := curo.New(
http.DefaultTransport,
curo.WithMode(curo.Observe),
)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
client := &http.Client{Transport: transport}
fmt.Println(transport.Mode() == curo.Observe)
fmt.Println(client.Transport == transport)
fmt.Println(transport.Stats().ObservedRequests)
}
Output: true true 0
func (*Transport) Close ¶
Close releases resources owned by the Transport.
Close is idempotent. Requests made after Close continue to delegate directly to the base transport, and retries that have not started and adaptive timeouts that have not expired are withdrawn. Close does not close the base transport or discard the bounded observation snapshot.
func (*Transport) CloseIdleConnections ¶
func (t *Transport) CloseIdleConnections()
CloseIdleConnections calls the base transport's CloseIdleConnections method, if it has one, so that http.Client.CloseIdleConnections reaches the base transport through Curo. Otherwise it does nothing.
CloseIdleConnections remains available after Close. The zero value and a nil *Transport do nothing.
func (*Transport) Mode ¶
Mode returns the configured operating mode.
The zero value and a nil *Transport report Off.
func (*Transport) Report ¶
Report returns a detached snapshot of published decisions and recent candidate changes.
Decisions are produced while Observe or Enforce requests complete. Report only copies them and never evaluates evidence, so an idle target keeps its last decision; compare ExpiresAt with the current time before acting on one. Targets are never older than Changes in the same Report for a target that is still tracked. A target admitted after an idle target is replaced starts without evidence.
Report is safe for concurrent use. The zero value and a nil *Transport return an empty Report. Report remains readable after Close, in Off mode, and after self-disable. If Curo fails while building a Report, Report returns an empty Report and the failure counts toward InternalFailures and self-disable.
Example ¶
package main
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"github.com/raj1kshtz/curo"
)
func main() {
server := httptest.NewServer(http.HandlerFunc(
func(response http.ResponseWriter, _ *http.Request) {
response.WriteHeader(http.StatusNoContent)
},
))
defer server.Close()
transport, err := curo.New(server.Client().Transport)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
client := &http.Client{Transport: transport}
request, err := http.NewRequestWithContext(
context.Background(),
http.MethodGet,
server.URL+"/orders/42?token=secret",
nil,
)
if err != nil {
fmt.Println("request setup failed")
return
}
response, err := client.Do(request)
if err != nil {
fmt.Println("request failed")
return
}
_ = response.Body.Close()
for _, decision := range transport.Report().Targets {
fmt.Printf(
"%s %s: readiness=%s diagnosis=%s candidates=%s\n",
decision.Target.Scheme,
decision.Target.Method,
decision.Readiness,
decision.Diagnosis,
decision.Candidates,
)
}
}
Output: http Read: readiness=Warming diagnosis=None candidates=None
func (*Transport) RoundTrip ¶
RoundTrip delegates req to the base transport, unless a dependency breaker rejects it.
Off, Observe, and Enforce requests that are neither rejected by a breaker nor retried reach the base transport exactly once, unmodified unless an adaptive timeout applies. In Enforce mode, a replay-safe request whose initial attempt failed in a retryable way may reach the base transport a second time, as a clone with a body from GetBody, after a short cancellable backoff. The retry's result is returned and the discarded first response body is closed without being drained. If the retry does not start, the first result is returned untouched. A done request context, or a closed deprecated Cancel channel, stops a retry that has not started.
In Enforce mode, while the dependency breaker of req's target is open, RoundTrip closes req.Body, if any, and returns ErrBreakerOpen without calling the base transport. After a cooldown, one request at a time is sent as a probe instead, and a probe is never retried.
In Enforce mode, a read request to a target whose decision selects CandidateTimeout, other than a probe, has an adaptive timeout unless its own deadline comes first. Each attempt reaches the base transport as a shallow copy of the request whose context Curo cancels if response headers do not arrive within the timeout. RoundTrip then returns ErrTimeout, closes any response that arrives late, and does not retry. Otherwise the response body is wrapped, so its concrete type is not preserved, and closing it or reading it to the end releases the timeout. A wrapped body keeps io.Writer when the original body has it, as for 101 Switching Protocols.
Curo-owned stages run only before an attempt starts or after its result has been captured, except that an expiring adaptive timeout cancels the attempt's context from the timer's goroutine. Base transport calls remain outside Curo's recovery boundary, so RoundTrip preserves their error and panic behavior, and their response apart from a timed body. RoundTrip remains available after Close.
The zero value and a nil *Transport close req.Body, if any, and return ErrNilBaseTransport.
func (*Transport) SetMode ¶
SetMode changes the operating mode for requests that start after the change.
A request never gains authority after it starts. Any mode change also withdraws a retry that an earlier Enforce request has not started yet, and an adaptive timeout that has not expired yet, even if Enforce is restored before either would act.
Example ¶
package main
import (
"errors"
"fmt"
"net/http"
"github.com/raj1kshtz/curo"
)
func main() {
transport, err := curo.New(http.DefaultTransport)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
fmt.Println(transport.Mode())
// Connect SetMode to a runtime control, such as an admin endpoint, so
// that a rollout or a rollback needs no deploy. Each change applies to
// requests that start after it. UnmarshalText accepts a mode name in any
// letter case and rejects any other text with ErrInvalidMode.
for _, name := range []string{"Enforce", "observe", "on", "OFF"} {
var mode curo.Mode
err = mode.UnmarshalText([]byte(name))
if err != nil {
fmt.Println("unknown mode:", name)
continue
}
err = transport.SetMode(mode)
if err != nil {
fmt.Println("mode change failed")
return
}
fmt.Println(transport.Mode())
}
_ = transport.Close()
fmt.Println(errors.Is(transport.SetMode(curo.Enforce), curo.ErrClosed))
}
Output: Observe Enforce Observe unknown mode: on Off true
func (*Transport) Stats ¶
Stats returns a concurrency-safe aggregate snapshot.
The zero value and a nil *Transport return an empty snapshot. Stats remains readable after Close.
Example ¶
package main
import (
"context"
"expvar"
"fmt"
"net/http"
"net/http/httptest"
"github.com/raj1kshtz/curo"
)
func main() {
server := httptest.NewServer(http.HandlerFunc(
func(response http.ResponseWriter, _ *http.Request) {
response.WriteHeader(http.StatusNoContent)
},
))
defer server.Close()
transport, err := curo.New(server.Client().Transport)
if err != nil {
fmt.Println("setup failed")
return
}
defer func() {
_ = transport.Close()
}()
client := &http.Client{Transport: transport}
request, err := http.NewRequestWithContext(
context.Background(),
http.MethodGet,
server.URL,
nil,
)
if err != nil {
fmt.Println("request setup failed")
return
}
response, err := client.Do(request)
if err != nil {
fmt.Println("request failed")
return
}
_ = response.Body.Close()
// expvar.Publish("curo", stats) serves the snapshot as JSON at
// /debug/vars. The counters are cumulative, so alert on their rates.
stats := expvar.Func(func() any {
return transport.Stats()
})
fmt.Println(stats.String())
}
Output: {"ObservedRequests":1,"TrackedTargets":1,"OverflowRequests":0,"RetryAttempts":0,"RetrySuccesses":0,"RetryBudgetDenials":0,"BreakerOpens":0,"BreakerProbes":0,"BreakerRejections":0,"Timeouts":0,"ShadowTimeouts":0,"InternalFailures":0,"SelfDisabled":false}
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
internal
|
|
|
breaker
Package breaker contains Curo's per-target dependency breaker.
|
Package breaker contains Curo's per-target dependency breaker. |
|
diagnose
Package diagnose evaluates bounded evidence using deterministic rules.
|
Package diagnose evaluates bounded evidence using deterministic rules. |
|
guard
Package guard contains the failure boundary for Curo-owned work.
|
Package guard contains the failure boundary for Curo-owned work. |
|
observe
Package observe contains Curo's bounded request observation state.
|
Package observe contains Curo's bounded request observation state. |
|
policy
Package policy maps diagnoses to bounded control candidates.
|
Package policy maps diagnoses to bounded control candidates. |
|
retry
Package retry contains Curo's fixed retry rules and lock-free retry budgets.
|
Package retry contains Curo's fixed retry rules and lock-free retry budgets. |
|
timeout
Package timeout contains Curo's adaptive timeout rule and the per-attempt mechanism that enforces it.
|
Package timeout contains Curo's adaptive timeout rule and the per-attempt mechanism that enforces it. |