Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
22 changes: 20 additions & 2 deletions apps/edge-api/internal/errors/provider_blind.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,8 +41,26 @@ var (
func WriteProviderBlindUpstreamError(w http.ResponseWriter, alias string, httpStatus int, rawMessage string) {
errType := "api_error"
code := "upstream_error"
// responseStatus is what the CUSTOMER sees. It starts equal to httpStatus,
// the real upstream status, and only diverges for 402 below. httpStatus
// itself stays untouched so the operator log two lines down keeps the real
// value even after responseStatus is remapped.
responseStatus := httpStatus

switch httpStatus {
case http.StatusPaymentRequired:
// A 402 here is the upstream refusing for ITS OWN funding (our
// provider account's balance with OpenRouter, DeepSeek, or whoever),
// never the Hive customer's balance: CreateReservation already
// verified the customer's own credit before dispatch ever reached
// this point (D-034, fail closed). Relaying a literal 402 tells the
// customer the opposite, that THEY must pay, which is exactly the
// provider refusal being presented as a caller fault the way issue
// #1411 names. Treated as the same "temporarily unavailable" verdict
// a 503/504 already gets: it is an availability problem on Hive's
// side, not a request the caller can fix.
responseStatus = http.StatusServiceUnavailable
code = "upstream_unavailable"
case http.StatusTooManyRequests:
errType = "rate_limit_error"
code = "upstream_rate_limited"
Expand All @@ -57,9 +75,9 @@ func WriteProviderBlindUpstreamError(w http.ResponseWriter, alias string, httpSt
code = "invalid_request"
}

message := sanitizeProviderBlindMessage(alias, httpStatus, rawMessage)
message := sanitizeProviderBlindMessage(alias, responseStatus, rawMessage)
logProviderBlindUpstreamError(w, alias, httpStatus, rawMessage, message)
WriteError(w, httpStatus, errType, message, &code)
WriteError(w, responseStatus, errType, message, &code)
}

func sanitizeProviderBlindMessage(alias string, httpStatus int, raw string) string {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
package errors

import (
"net/http"
"net/http/httptest"
"strings"
"testing"
)

// TestWriteProviderBlindUpstreamErrorRemaps402ToUpstreamUnavailable is the
// regression guard for issue #1411: OpenRouter (or any upstream) answering
// 402 Payment Required is talking about ITS OWN account balance with the
// upstream vendor, never the Hive customer's balance: CreateReservation
// already verified the customer's own credit before dispatch reached this
// point. Forwarding a literal 402 to the customer says the opposite: it
// reads as "you must pay," which is backwards and is exactly the "provider
// refusal presented as a customer/model fault" shape the issue names.
//
// The fixture below is OpenRouter's real documented error envelope
// (https://openrouter.ai/docs/api_reference/errors-and-debugging.md,
// confirmed live 2026-08-29): {"error":{"code":402,"message":"...",
// "metadata":{"error_type":"payment_required"}}}, HTTP status equal to
// error.code. This is a SIMULATED upstream response (an in-process
// httptest.ResponseRecorder, no network call to any real provider); the
// shape is the part that is real, not the transport.
func TestWriteProviderBlindUpstreamErrorRemaps402ToUpstreamUnavailable(t *testing.T) {
w := httptest.NewRecorder()

raw := `{"error":{"code":402,"message":"Insufficient credits. Add more using https://openrouter.ai/settings/credits","metadata":{"error_type":"payment_required"}}}`
logs := captureProviderBlindLogs(t, func() {
WriteProviderBlindUpstreamError(w, "hive-auto", http.StatusPaymentRequired, raw)
})

if w.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d (a provider funding refusal must not read as the CUSTOMER owing money)", w.Code, http.StatusServiceUnavailable)
}

resp := decodeOpenAIError(t, w)
if resp.Error.Code == nil || *resp.Error.Code != "upstream_unavailable" {
t.Fatalf("code = %v, want %q", resp.Error.Code, "upstream_unavailable")
}
if resp.Error.Type != "api_error" {
t.Fatalf("type = %q, want %q", resp.Error.Type, "api_error")
}
if resp.Error.Message != "hive-auto is temporarily unavailable." {
t.Fatalf("message = %q, want the same verdict a 503/504 already gets", resp.Error.Message)
}
assertNoProviderLeak(t, resp.Error.Message)
for _, forbidden := range []string{"insufficient", "credit", "balance", "402", "payment"} {
if strings.Contains(strings.ToLower(resp.Error.Message), forbidden) {
t.Fatalf("upstream funding vocabulary reached the customer: %q in %q", forbidden, resp.Error.Message)
}
}

// The operator log must keep the REAL status (402) and the raw body, even
// though the customer-facing status changed to 503: this is a diagnosis
// aid, not a second leak surface.
if !strings.Contains(logs, "status=402") {
t.Fatalf("operator log lost the real upstream status: %q", logs)
}
if !strings.Contains(logs, "payment_required") {
t.Fatalf("operator log lost the raw upstream body: %q", logs)
}
}
135 changes: 135 additions & 0 deletions apps/edge-api/internal/inference/upstream_payment_required_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
package inference

import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"strings"
"testing"

apierrors "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/errors"
)

// --- issue #1411: upstream funds exhaustion ---
//
// These tests drive the REAL Orchestrator.executeSync / executeStreaming
// lifecycle end to end against in-process httptest stand-ins for the
// control-plane (routing + accounting) and LiteLLM, the same harness
// stream_disconnect_test.go and sync_settlement_test.go already use. Only
// the LiteLLM stand-in is faked here; the reservation, release, dispatch
// and sanitization code paths are the real production code.
//
// The upstream body below is OpenRouter's documented error envelope
// (https://openrouter.ai/docs/api_reference/errors-and-debugging.md,
// confirmed live 2026-08-29): the JSON shape is error.code plus
// error.message plus error.metadata, with the HTTP status equal to
// error.code. It is SIMULATED (an httptest.Server, never a call to a real
// provider): draining the real OpenRouter wallet to observe this would
// cause the exact outage this test exists to prevent.

const openRouterInsufficientCreditsBody = "{\"error\":{\"code\":402,\"message\":\"Insufficient credits. Add more using https://openrouter.ai/settings/credits\",\"metadata\":{\"error_type\":\"payment_required\"}}}"

// paymentRequiredServer stands in for LiteLLM/the upstream provider
// answering with OpenRouter's real 402 shape, and counts how many times it
// was hit so a test can confirm the bounded-retry path (429/5xx only) did
// not treat 402 as retryable.
func paymentRequiredServer(hits *int) *httptest.Server {
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
*hits++
w.WriteHeader(http.StatusPaymentRequired)
_, _ = w.Write([]byte(openRouterInsufficientCreditsBody))
}))
}

func assertUpstreamUnavailableBody(t *testing.T, w *httptest.ResponseRecorder) {
t.Helper()
if w.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d (upstream funding refusal must not read as the customer owing money)", w.Code, http.StatusServiceUnavailable)
}
var resp apierrors.OpenAIError
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
t.Fatalf("response is not a valid OpenAI error envelope: %v, body=%s", err, w.Body.String())
}
if resp.Error.Code == nil || *resp.Error.Code != "upstream_unavailable" {
t.Fatalf("code = %v, want upstream_unavailable", resp.Error.Code)
}
lower := strings.ToLower(resp.Error.Message)
for _, forbidden := range []string{"openrouter", "insufficient", "credit", "balance", "402", "payment", "settings/credits"} {
if strings.Contains(lower, forbidden) {
t.Fatalf("upstream funding/identity vocabulary reached the customer: %q in %q", forbidden, resp.Error.Message)
}
}
}

// TestExecuteSync_UpstreamPaymentRequired_ReleasesHoldAndDoesNotBill is the
// sync /v1/chat/completions and /v1/responses regression guard.
func TestExecuteSync_UpstreamPaymentRequired_ReleasesHoldAndDoesNotBill(t *testing.T) {
hits := 0
litellmSrv := paymentRequiredServer(&hits)
defer litellmSrv.Close()
routingSrv := newRoutingMock(litellmSrv.URL)
defer routingSrv.Close()
rec := &accountingRecorder{}
acctSrv := newAccountingMock(rec)
defer acctSrv.Close()

orch := newAuthorizedOrchestrator(acctSrv.URL, routingSrv.URL, litellmSrv.URL)
w := callSyncCtx(orch, context.Background())

assertUpstreamUnavailableBody(t, w)

if hits != 1 {
t.Fatalf("litellm hits = %d, want 1 (a 402 is neither 429 nor 5xx and must not be retried)", hits)
}

releaseBody, released := rec.find("/internal/accounting/reservations/release")
if !released {
t.Fatalf("expected the reservation to be released; calls seen: %+v", rec.calls)
}
if releaseBody["reason"] != "upstream_error" {
t.Errorf("release reason = %v, want upstream_error", releaseBody["reason"])
}
if rec.has("/internal/accounting/reservations/finalize") {
t.Fatalf("reservation must never be finalized (charged) on an upstream funding refusal; calls seen: %+v", rec.calls)
}
}

// TestExecuteStreaming_UpstreamPaymentRequired_ReleasesHoldAndDoesNotBill is
// the streaming twin: the 402 arrives as LiteLLM's initial HTTP response,
// before the SSE 200 is ever committed to the client, so this is a
// pre-stream refusal, not a mid-stream error frame.
func TestExecuteStreaming_UpstreamPaymentRequired_ReleasesHoldAndDoesNotBill(t *testing.T) {
hits := 0
litellmSrv := paymentRequiredServer(&hits)
defer litellmSrv.Close()
routingSrv := newRoutingMock(litellmSrv.URL)
defer routingSrv.Close()
rec := &accountingRecorder{}
acctSrv := newAccountingMock(rec)
defer acctSrv.Close()

orch := newAuthorizedOrchestrator(acctSrv.URL, routingSrv.URL, litellmSrv.URL)

req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader("{}"))
req.Header.Set("Authorization", "Bearer test-token")
w := httptest.NewRecorder()
_ = orch.executeStreaming(context.Background(), w, req, EndpointChatCompletions, []byte("{}"), "gpt-4o", "gpt-4o",
NeedFlags{NeedChatCompletions: true, NeedStreaming: true}, 10000, false, nil, orch.litellm.ChatCompletion)

assertUpstreamUnavailableBody(t, w)

if hits != 1 {
t.Fatalf("litellm hits = %d, want 1 (a 402 is neither 429 nor 5xx and must not be retried)", hits)
}
releaseBody, released := rec.find("/internal/accounting/reservations/release")
if !released {
t.Fatalf("expected the reservation to be released; calls seen: %+v", rec.calls)
}
if releaseBody["reason"] != "upstream_error" {
t.Errorf("release reason = %v, want upstream_error", releaseBody["reason"])
}
if rec.has("/internal/accounting/reservations/finalize") {
t.Fatalf("reservation must never be finalized (charged) on an upstream funding refusal; calls seen: %+v", rec.calls)
}
}
Loading