Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
65802ea
fix(health-checks): thread the JSON-RPC method into both heuristic sites
oten91 Aug 6, 2026
a722e32
docs(config): catch the schema up to the health-check keys that ship
oten91 Aug 6, 2026
933bccf
feat(admin): per-operator reputation drain endpoint
oten91 Aug 6, 2026
77fc889
feat(metrics): per-connection websocket frame-rate histogram
oten91 Aug 6, 2026
f5419e7
docs: reputation drain endpoint + WS frame-rate distribution
oten91 Aug 6, 2026
95e89c0
fix(admin): drain survives storage refresh; target by URL not node id
oten91 Aug 7, 2026
234d160
feat(admin): fleet-wide reputation drains with a hard expiry ceiling
oten91 Aug 7, 2026
3ef1e3d
test(admin): pin drain isolation from the scoring system
oten91 Aug 7, 2026
03de910
feat(metrics): report admin drains separately from earned cooldowns
oten91 Aug 7, 2026
017bc5c
fix(admin): key drains by operator, not by resolved endpoint keys
oten91 Aug 7, 2026
d79b126
fix(admin): apply the drain to the set the filter actually returns
oten91 Aug 7, 2026
502c523
test(shannon): selection-path harness for routing changes
oten91 Aug 7, 2026
57c074e
fix(websocket): close the endpoint connection cleanly on rebind
oten91 Aug 7, 2026
cacce39
test(qos): restore the evm build after the reputation interface grew
oten91 Aug 7, 2026
b4640ad
fix(websocket): close health-check probes with a handshake, not a soc…
oten91 Aug 7, 2026
efef6d4
fix(websocket): close live connections on SIGTERM instead of dropping…
oten91 Aug 7, 2026
1d6a29a
fix(websocket): send each peer a close code that fits its role
oten91 Aug 7, 2026
ca5cfbe
fix(admin): apply the drain to the endpoint a connection is already b…
oten91 Aug 7, 2026
9bde542
feat(protocol): ban operator domains from serving an RPC type on ever…
oten91 Aug 7, 2026
999fbf0
test(protocol): cover the NewProtocol domain-blocklist wiring
oten91 Aug 7, 2026
a842f00
docs: correct two stale claims about drains and the over-servicing brake
oten91 Aug 7, 2026
b8c0257
refactor(test): use placeholder domains in fixtures and examples
oten91 Aug 7, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
252 changes: 249 additions & 3 deletions CLAUDE.md

Large diffs are not rendered by default.

49 changes: 49 additions & 0 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,22 @@ var (
// the executable to get the full path to the config file.
const defaultConfigPath = "config/.config.yaml"

// websocketShutdownTimeout bounds the close-handshake sweep over live websocket
// connections at termination.
//
// The sweep runs the closes concurrently, so this is a ceiling on the slowest peer rather
// than a per-connection cost — five seconds is far past the one-second write deadline on
// each close frame, and comfortably inside the pod's 30s termination grace period. Being
// polite must never be the reason a pod gets SIGKILLed, which would produce exactly the
// abrupt teardown the sweep exists to prevent.
const websocketShutdownTimeout = 5 * time.Second

// websocketShutdowner is implemented by protocols that hold live websocket bridges and
// must close them explicitly at termination.
type websocketShutdowner interface {
ShutdownWebsockets(ctx context.Context, reason string) int
}

func main() {
log.Printf(`{"level":"info","message":"PATH 🌿 gateway starting..."}`)

Expand Down Expand Up @@ -347,6 +363,14 @@ func main() {
// failing startup.
websocketAdmin, _ := protocol.(router.WebsocketAdmin)

// Admin handler to temporarily bench one operator's endpoints for a service via
// POST /admin/reputation/drain/{serviceId}. nil when reputation is disabled, which
// leaves the endpoint reporting 503 rather than failing startup.
var reputationAdmin router.ReputationAdmin
if reputationSvc := protocol.GetReputationService(); reputationSvc != nil {
reputationAdmin = reputationSvc
}

// Initialize the API router to serve requests to the PATH API.
apiRouter := router.NewRouter(
logger,
Expand All @@ -357,6 +381,7 @@ func main() {
gtw.DomainCircuitBreaker,
chainStateAdmin,
websocketAdmin,
reputationAdmin,
unifiedServicesConfig,
)

Expand Down Expand Up @@ -433,6 +458,30 @@ func main() {

logger.Info().Msg("Shutting down PATH...")

// Close live websocket connections FIRST, and with a close handshake.
//
// server.Shutdown below cannot do this: it explicitly does not close hijacked
// connections, and every websocket is hijacked. Without this the process just exits,
// every socket dies with its TCP connection, and both peers see an abnormal closure
// (1006) — indistinguishable from a crash for the client, and logged as their own
// fault by the endpoint. Fleetwide that made every rollout emit a burst of 1006s
// across all services within the same second.
//
// Before backgroundCancel() so the bridges are torn down deliberately rather than
// racing a context cancellation that would reach them as a generic failure.
//
// Its own short budget, well inside the pod's termination grace period: the close
// frames are best-effort and a peer that will not answer must not delay exiting.
//
// Behind an assertion rather than a method on gateway.Protocol: holding live
// websocket bridges is a property of the protocol implementation, not of the
// interface, and a protocol that holds none needs nothing here.
if wsShutdowner, ok := protocol.(websocketShutdowner); ok {
wsCtx, wsCancel := context.WithTimeout(context.Background(), websocketShutdownTimeout)
wsShutdowner.ShutdownWebsockets(wsCtx, "gateway shutting down")
wsCancel()
}

// Cancel background context to stop all background services (pprof, health checks)
backgroundCancel()

Expand Down
37 changes: 35 additions & 2 deletions config/config.schema.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,19 @@ properties:
description: "Enable/disable health checks."
type: boolean
default: false
sync_allowance:
description: "Default blocks behind the latest an endpoint may be before it counts as out of sync. Per-service entries in local[] override it. 0 disables the sync check."
type: integer
minimum: 0
default: 0
max_workers:
description: "Maximum concurrent health check workers. Auto-sized from service and check counts when unset."
type: integer
minimum: 1
backend_dedup:
description: "Fire ONE relay per unique backend URL each cycle (a rotating representative supplier) and fan the result to the other suppliers sharing that URL. Cuts redundant health-check relay volume by the redundancy factor. WebSocket checks are never deduped. Default: true."
type: boolean
default: true
coordination:
description: "Leader election configuration for multi-instance deployments."
type: object
Expand Down Expand Up @@ -447,6 +460,10 @@ properties:
description: "Enable/disable health checks for this service."
type: boolean
default: true
sync_allowance:
description: "Blocks behind the latest an endpoint may be before it counts as out of sync. Overrides the global default for this service. 0 disables the sync check."
type: integer
minimum: 0
checks:
description: "List of health check configurations."
type: array
Expand All @@ -461,8 +478,17 @@ properties:
description: "Name of the health check."
type: string
type:
description: "Type of health check (e.g., 'jsonrpc')."
description: "RPC protocol type (must match service's rpc_types): json_rpc, rest, comet_bft, websocket, grpc."
type: string
enum: ["json_rpc", "rest", "comet_bft", "websocket", "grpc"]
enabled:
description: "Enable/disable this individual check."
type: boolean
headers:
description: "Request headers."
type: object
additionalProperties:
type: string
method:
description: "HTTP method (GET, POST)."
type: string
Expand All @@ -481,10 +507,17 @@ properties:
timeout:
description: "Request timeout."
type: string
archival:
description: "Archival-specific check: only runs on endpoints marked archival."
type: boolean
sync_check:
description: "Enable sync check validation. When true, extracts block height from the response and validates against the service's sync_allowance. Requires sync_allowance > 0."
type: boolean
default: false
reputation_signal:
description: "Reputation signal to emit on failure."
type: string
enum: ["minor_error", "major_error", "critical_error"]
enum: ["minor_error", "major_error", "critical_error", "fatal_error"]

# ====================================
# UNIFIED SERVICE CONFIGURATION
Expand Down
42 changes: 34 additions & 8 deletions gateway/health_check_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -902,7 +902,7 @@ func (e *HealthCheckExecutor) recordCheckResult(
// Over-servicing rejections — same no-penalty rule as the request path.
// Without this, every health-check probe to an exhausted supplier records a
// MajorError signal and pins their reputation at 0 even after the request
// path has stopped penalizing them. Production canary observed easy2stake's
// path has stopped penalizing them. Production canary observed operator-zeta's
// BSC supplier set stuck at score=0 with success-only request-path signals
// because the health-check executor was draining them in parallel.
if heuristic.IsOverServicedError(checkErr.Error()) {
Expand Down Expand Up @@ -1245,8 +1245,21 @@ func (e *HealthCheckExecutor) ExecuteCheckViaProtocol(
}

// Heuristic analysis - detect bad gateway, empty responses, and other error patterns
// This runs BEFORE QoS validation to catch issues the basic validation might miss
heuristicResult := heuristic.Analyze(responseBody, httpStatusCode, servicePayload.EffectiveRPCType(), "")
// This runs BEFORE QoS validation to catch issues the basic validation might miss.
//
// The method must be threaded through here exactly as the other four call sites do it
// (protocol/shannon/context.go, gateway/hedge.go, http_request_context_handle_request.go).
// This call passed a hardcoded "" - so even with JSONRPCMethod populated on the payload,
// every method-aware rule was blind at this site, and this is the site that decides the
// health check's outcome. A CometBFT `health` response
// ({"jsonrpc":"2.0","id":1,"result":{}}) is flagged jsonrpc_empty_object_result at
// confidence 0.95 unless the method is known, which maps to SignalMajorError below.
jsonrpcMethod := servicePayload.JSONRPCMethod
if jsonrpcMethod == "" {
// REST checks carry no JSON-RPC method; path-aware rules key off Path instead.
jsonrpcMethod = servicePayload.Path
}
heuristicResult := heuristic.Analyze(responseBody, httpStatusCode, servicePayload.EffectiveRPCType(), jsonrpcMethod)
if heuristicResult.ShouldRetry {
heuristicErr := fmt.Errorf("heuristic detected error: %s - %s", heuristicResult.Reason, heuristicResult.Details)
e.logger.Debug().
Expand Down Expand Up @@ -1628,12 +1641,25 @@ func (e *HealthCheckExecutor) buildServicePayload(check HealthCheckConfig) proto
rpcType = sharedtypes.RPCType_UNKNOWN_RPC
}

// JSONRPCMethod drives the method-aware heuristic checks. User traffic gets it
// from QoS parsing; a health check never goes through QoS, so without this it
// stayed empty and the heuristic fell back to Path — which for a JSON-RPC check
// is "/", matching no method at all.
//
// That broke the CometBFT carve-out in analyzeJSONRPC: `health` legitimately
// answers {"jsonrpc":"2.0","id":1,"result":{}}, and the empty-object rule is
// skipped only when rpcType is COMET_BFT or isCometBFTMethod(method) holds.
// Cosmos services declare the check as type json_rpc, so neither guard fired and
// every healthy node was scored jsonrpc_empty_object_result → retry + minor_error.
// Measured 2026-08-06: 16 cosmos services failing this check 100%, ~19% of all
// failing health-check signals fleet-wide.
return protocol.Payload{
Method: check.Method,
Path: check.Path,
Data: check.Body,
Headers: headers,
RPCType: rpcType, // Set from aligned health check type
Method: check.Method,
Path: check.Path,
Data: check.Body,
Headers: headers,
RPCType: rpcType, // Set from aligned health check type
JSONRPCMethod: extractJSONRPCMethod([]byte(check.Body)),
}
}

Expand Down
Loading
Loading