diff --git a/.env.example b/.env.example index 4bedd52..7c6dd33 100644 --- a/.env.example +++ b/.env.example @@ -121,15 +121,26 @@ EMBEDDING_DIMENSIONS=0 # Google Geocoding when set; OpenStreetMap's Nominatim otherwise. GEOCODER_API_KEY= -# Nearle Buddy. Any OpenAI-compatible endpoint; empty provider disables the -# assistant entirely and the console's composer stays disabled. -# Groq: https://api.groq.com/openai/v1 openai/gpt-oss-120b -# Ollama: http://localhost:11434/v1 (no key needed) +# ── Nearle Buddy ──────────────────────────────────────────────────────────── +# +# ONE variable. The provider, endpoint and model are defaults in config.go +# (openai / api.groq.com / openai/gpt-oss-120b) because each has one right +# answer for this product — and three variables that must be typed correctly +# into a hosting platform are three ways for the assistant to sit silently off, +# which is how it spent its first week. +# +# The key is the only one that differs per deployment and the only one that +# cannot live in this repository. Locally it goes in `.env.secrets`, which git +# ignores; in production it is set on the platform. +ASSISTANT_API_KEY= + +# Overrides, none of them needed for the shipped setup. +# Ollama on a laptop: ASSISTANT_BASE_URL=http://localhost:11434/v1 and +# ASSISTANT_MODEL=llama3 — a local endpoint needs no key. ASSISTANT_PROVIDER= ASSISTANT_BASE_URL= -ASSISTANT_API_KEY= ASSISTANT_MODEL= -# Optional per-tier overrides. ASSISTANT_MODEL alone sets all three. +# Per-tier overrides. ASSISTANT_MODEL alone sets all three. ASSISTANT_MODEL_FAST= ASSISTANT_MODEL_BALANCED= ASSISTANT_MODEL_DEEP= diff --git a/.env.secrets b/.env.secrets deleted file mode 100644 index 2f66c65..0000000 --- a/.env.secrets +++ /dev/null @@ -1,17 +0,0 @@ -# Secrets for LOCAL runs. Git ignores this file — that is the whole point of it. -# -# Every other .env file in this folder is tracked, so a key written there is a -# key pushed to the repository. This one is not, so it is where a key goes. -# -# It does NOT reach production. The container takes its environment from Dokploy -# and no .env file is copied into the image, so these same three names have to be -# set again in Dokploy → your backend application → Environment → Redeploy. - -# Nearle Buddy. Three variables and nothing else is needed; ASSISTANT_PROVIDER is -# worked out from the model name. -ASSISTANT_BASE_URL=https://api.groq.com/openai/v1 -ASSISTANT_MODEL=openai/gpt-oss-120b - -# Paste your Groq key here, from https://console.groq.com/keys -# The one used while building this is in the chat history and should be replaced. -ASSISTANT_API_KEY=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 diff --git a/.gitignore b/.gitignore index b172fdc..894861a 100644 --- a/.gitignore +++ b/.gitignore @@ -59,3 +59,4 @@ Thumbs.db # that true, and it is why this file exists separately from .env.local, which # IS tracked and therefore cannot hold a key. +.env.secrets diff --git a/config/assistant_test.go b/config/assistant_test.go index 63a22e1..d1a84dd 100644 --- a/config/assistant_test.go +++ b/config/assistant_test.go @@ -131,3 +131,78 @@ func TestTheEnvironmentsOwnFileBeatsTheSharedOne(t *testing.T) { t.Fatalf("the environment's own file does not take precedence: %v", order) } } + +// One variable, not four. +// +// The assistant sat switched off for days because `ASSISTANT_PROVIDER` had not +// been typed into a hosting platform's environment tab — a variable whose only +// correct value is "openai", because every endpoint this server speaks is +// OpenAI-compatible. The base URL and the model had one right answer too. +// +// So three of the four are constants now. The key is the only one that varies +// between deployments and the only one that cannot live in the repository. +func TestTheKeyAloneSwitchesTheAssistantOn(t *testing.T) { + for _, name := range []string{ + "ASSISTANT_PROVIDER", "ASSISTANT_BASE_URL", "ASSISTANT_MODEL", + "ASSISTANT_MODEL_BALANCED", "ASSISTANT_MODEL_FAST", "ASSISTANT_MODEL_DEEP", + } { + t.Setenv(name, "") + } + t.Setenv("ASSISTANT_API_KEY", "gsk_not-a-real-key") + + cfg := AssistantConfig{ + Provider: assistantProvider(), + BaseURL: env("ASSISTANT_BASE_URL", defaultAssistantBaseURL), + APIKey: env("ASSISTANT_API_KEY", ""), + Balanced: env("ASSISTANT_MODEL_BALANCED", env("ASSISTANT_MODEL", defaultAssistantModel)), + } + + if !cfg.Enabled() { + t.Fatalf("the key alone did not switch it on: %s", cfg.Why()) + } + if cfg.Provider != "openai" { + t.Fatalf("provider defaulted to %q", cfg.Provider) + } + if cfg.ModelFor("fast") != defaultAssistantModel { + t.Fatalf("the fast tier fell through to %q", cfg.ModelFor("fast")) + } +} + +func TestNoKeyIsStillOffAndSaysWhich(t *testing.T) { + // The defaults must not make an unconfigured deployment look ready. Without + // a key every question would reach Groq and come back 401, which reads as + // the assistant being broken rather than as not being set up. + cfg := AssistantConfig{ + Provider: defaultAssistantProvider, + BaseURL: defaultAssistantBaseURL, + Balanced: defaultAssistantModel, + } + if cfg.Enabled() { + t.Fatal("reported ready with no key") + } + if !strings.Contains(cfg.Why(), "ASSISTANT_API_KEY") { + t.Fatalf("did not name the one variable left to set: %q", cfg.Why()) + } +} + +func TestEachDefaultIsStillOverridable(t *testing.T) { + // Running against Ollama on a laptop must not need a code change. + t.Setenv("ASSISTANT_PROVIDER", "ollama") + t.Setenv("ASSISTANT_BASE_URL", "http://localhost:11434/v1") + t.Setenv("ASSISTANT_MODEL", "llama3") + + cfg := AssistantConfig{ + Provider: assistantProvider(), + BaseURL: env("ASSISTANT_BASE_URL", defaultAssistantBaseURL), + APIKey: env("ASSISTANT_API_KEY", ""), + Balanced: env("ASSISTANT_MODEL_BALANCED", env("ASSISTANT_MODEL", defaultAssistantModel)), + } + + if cfg.Provider != "ollama" || cfg.Balanced != "llama3" { + t.Fatalf("an override was ignored: %+v", cfg) + } + // Local endpoints need no key, so this must be on without one. + if !cfg.Enabled() { + t.Fatalf("a local model was refused: %s", cfg.Why()) + } +} diff --git a/config/config.go b/config/config.go index 150cd38..4d3b535 100644 --- a/config/config.go +++ b/config/config.go @@ -245,20 +245,59 @@ func (a AssistantConfig) ModelFor(tier string) string { return a.Balanced } +// What Nearle Buddy runs on unless a deployment says otherwise. +// +// These are in the code rather than in the environment because they are not +// secrets and not deployment-specific — they are what this product uses. Every +// variable that has one right answer is a variable somebody has to remember, +// get past a platform's UI, and then re-enter on the next environment; three of +// the four were exactly that, and the assistant sat switched off for days +// because one of them had not been typed. +// +// The API key is the one that genuinely varies and genuinely cannot live here. +const ( + defaultAssistantProvider = "openai" + defaultAssistantBaseURL = "https://api.groq.com/openai/v1" + defaultAssistantModel = "openai/gpt-oss-120b" +) + // assistantProvider reads the provider, defaulting to the one shape this // server speaks. // -// A deployment that names a model and a key has said what it wants; making it -// also name a protocol it has no choice about is a variable that exists only to -// be forgotten. +// Every endpoint here is OpenAI-compatible — Groq, Ollama, Together and OpenAI +// itself — so the base URL is what actually distinguishes them. Naming a +// protocol you have no choice about is a variable that exists only to be +// forgotten. func assistantProvider() string { if named := strings.ToLower(strings.TrimSpace(env("ASSISTANT_PROVIDER", ""))); named != "" { return named } - if strings.TrimSpace(env("ASSISTANT_MODEL_BALANCED", env("ASSISTANT_MODEL", ""))) != "" { - return "openai" + return defaultAssistantProvider +} + +// AssistantFromEnv reads the assistant's settings, defaults and all. +// +// Exported and used by `Load` rather than written inline there, because the +// live tests need the SAME reading. They used to build this struct by hand from +// `os.Getenv`, which meant they skipped silently the moment a default was +// introduced — they were testing a configuration production no longer uses, +// and the one time that mattered was the day the provider stopped being +// required and nothing noticed. +// +// Three of the four fields have one right answer and come from the constants +// above. The key varies between deployments and is the only one that cannot +// live in this repository. +func AssistantFromEnv() AssistantConfig { + return AssistantConfig{ + Provider: assistantProvider(), + BaseURL: env("ASSISTANT_BASE_URL", defaultAssistantBaseURL), + APIKey: env("ASSISTANT_API_KEY", ""), + Fast: env("ASSISTANT_MODEL_FAST", ""), + // ASSISTANT_MODEL alone still sets every tier, for a deployment that + // wants one model everywhere but not this one. + Balanced: env("ASSISTANT_MODEL_BALANCED", env("ASSISTANT_MODEL", defaultAssistantModel)), + Deep: env("ASSISTANT_MODEL_DEEP", ""), } - return "" } // IsProduction is true under APP_ENV=production. @@ -317,20 +356,7 @@ func Load() (*Config, error) { BaseURL: env("EMBEDDING_BASE_URL", ""), }, - Assistant: AssistantConfig{ - // Defaults to "openai" when a model is named, because every endpoint - // this speaks is OpenAI-compatible and the base URL is what actually - // distinguishes them. One less variable to set, and one less way to - // have the assistant silently off. - Provider: assistantProvider(), - BaseURL: env("ASSISTANT_BASE_URL", ""), - APIKey: env("ASSISTANT_API_KEY", ""), - Fast: env("ASSISTANT_MODEL_FAST", ""), - // ASSISTANT_MODEL alone sets every tier, for a deployment that has - // not thought about tiers yet. - Balanced: env("ASSISTANT_MODEL_BALANCED", env("ASSISTANT_MODEL", "")), - Deep: env("ASSISTANT_MODEL_DEEP", ""), - }, + Assistant: AssistantFromEnv(), POSTokenSecret: env("POS_TOKEN_SECRET", ""), JWTSecret: env("JWT_SECRET_KEY", ""), diff --git a/controllers/assistantController.go b/controllers/assistantController.go index b29a00f..4cff704 100644 --- a/controllers/assistantController.go +++ b/controllers/assistantController.go @@ -112,6 +112,13 @@ func (ctl *AssistantController) Ask(c *fiber.Ctx) error { if errors.As(err, &tooFast) { return assistantRefuse(c, http.StatusTooManyRequests, err.Error()) } + // The provider's quota, as opposed to our own limiter above. Same status + // for the same reason — it is not a bad question, it is a busy minute — + // and the message is ours rather than Groq's, which names our billing + // account and the tokens-per-minute arithmetic behind it. + if errors.Is(err, utils.ErrBusy) { + return assistantRefuse(c, http.StatusTooManyRequests, utils.ErrBusy.Error()) + } return assistantRefuse(c, http.StatusBadRequest, err.Error()) } diff --git a/controllers/assistantHTTP_test.go b/controllers/assistantHTTP_test.go index 499f742..cc21489 100644 --- a/controllers/assistantHTTP_test.go +++ b/controllers/assistantHTTP_test.go @@ -3,9 +3,9 @@ package controllers import ( "context" "encoding/json" + "hash/crc32" "io" "net/http/httptest" - "os" "strconv" "strings" "testing" @@ -179,9 +179,24 @@ func buildApp(t *testing.T, chat utils.Chat) (*fiber.App, *fakeShop) { return app, shop } +// webSession mints a session for THIS test's own user. +// +// One user id across the file put every test in one rate-limit bucket — six +// questions and then 429 for ten seconds — so the suite passed test by test and +// failed when run together, which is the worst way round: green locally, red in +// CI, and the failure blamed on the model. +// +// A per-test user is also the truthful shape. The limiter is per person, and +// two tests are two people. func webSession(t *testing.T) string { t.Helper() - token, _, err := utils.MintWebToken(testCaller, time.Now()) + + claims := testCaller + // Stable across runs and distinct per test, so a failure names the same + // user every time. The fakes key on tenant, never on this. + claims.Userid = testCaller.Userid + int(crc32.ChecksumIEEE([]byte(t.Name()))%10_000) + + token, _, err := utils.MintWebToken(claims, time.Now()) if err != nil { t.Fatalf("minting a session: %v", err) } @@ -306,20 +321,9 @@ func TestStatusNamesTheMissingVariable(t *testing.T) { func liveHTTPChat(t *testing.T) utils.Chat { t.Helper() - // Defaults exactly as production does, so this proves three variables are - // enough rather than working around the question. - provider := strings.ToLower(strings.TrimSpace(os.Getenv("ASSISTANT_PROVIDER"))) - model := strings.TrimSpace(os.Getenv("ASSISTANT_MODEL")) - if provider == "" && model != "" { - provider = "openai" - } - - cfg := config.AssistantConfig{ - Provider: provider, - BaseURL: os.Getenv("ASSISTANT_BASE_URL"), - APIKey: os.Getenv("ASSISTANT_API_KEY"), - Balanced: model, - } + // Read exactly as production reads it, so this proves the shipped defaults + // work rather than quietly testing a configuration of its own. + cfg := config.AssistantFromEnv() if !cfg.Enabled() { t.Skipf("no model configured: %s", cfg.Why()) } diff --git a/services/assistantLive_test.go b/services/assistantLive_test.go index 8bbeb30..b4b8cf2 100644 --- a/services/assistantLive_test.go +++ b/services/assistantLive_test.go @@ -2,7 +2,6 @@ package services import ( "context" - "os" "strings" "testing" "time" @@ -41,23 +40,14 @@ import ( func liveChat(t *testing.T) (utils.Chat, string) { t.Helper() - // Provider defaults exactly as production does, so this exercises the - // defaulting rather than working around it. The first version read the - // variable straight and skipped silently the moment `ASSISTANT_PROVIDER` - // was left unset — which is precisely the configuration this is meant to - // prove works. - provider := strings.ToLower(strings.TrimSpace(os.Getenv("ASSISTANT_PROVIDER"))) - model := os.Getenv("ASSISTANT_MODEL") - if provider == "" && strings.TrimSpace(model) != "" { - provider = "openai" - } - - cfg := config.AssistantConfig{ - Provider: provider, - BaseURL: os.Getenv("ASSISTANT_BASE_URL"), - APIKey: os.Getenv("ASSISTANT_API_KEY"), - Balanced: model, - } + // Read exactly as production reads it, defaults and all. + // + // This was a hand-built struct twice, and it was wrong both times: first it + // read ASSISTANT_PROVIDER straight and skipped silently once that stopped + // being required, then it kept demanding ASSISTANT_MODEL after that gained + // a default. A test that builds its own configuration is a test of a + // configuration nobody runs. + cfg := config.AssistantFromEnv() if !cfg.Enabled() { t.Skipf("no model configured, skipping the live test: %s", cfg.Why()) } diff --git a/utils/embedding.go b/utils/embedding.go index c8498f9..01a49ff 100644 --- a/utils/embedding.go +++ b/utils/embedding.go @@ -8,6 +8,8 @@ import ( "fmt" "io" "net/http" + "reflect" + "strconv" "strings" "time" @@ -161,7 +163,41 @@ func (e *geminiEmbedder) Embed(ctx context.Context, text string) ([]float32, err // so a rate-limited assistant told a shopkeeper "embedding: HTTP 429" — a // sentence about a subsystem they have never heard of, describing something // that was not involved. +// postJSON sends the request, and sends it a second time if the provider said +// it was over its quota and named a wait we are willing to hold for. +// +// One retry, not a loop: past that, a queue forms behind a limit that is not +// going to lift, and the person is better told to try again than left watching +// a spinner. Both the chat gateway and the embedder go through here, so neither +// can be the one that forgot. func postJSON(ctx context.Context, client *http.Client, what, url, auth string, body, out interface{}, headers ...string) error { + err := postJSONOnce(ctx, client, what, url, auth, body, out, headers...) + + var busy *tooManyRequests + if errors.As(err, &busy) && waitBeforeRetry(ctx, busy.after) { + // `out` must be emptied first. The first attempt decoded the provider's + // error body into it, and `encoding/json` leaves fields the second + // payload does not mention exactly as it found them — so a retry that + // SUCCEEDED came back carrying the 429's `error` object, and every + // caller here checks that field before the data. The result was a + // successful call reported as the failure it had just recovered from. + resetForRetry(out) + return postJSONOnce(ctx, client, what, url, auth, body, out, headers...) + } + return err +} + +// resetForRetry empties a decode target so a second attempt cannot inherit the +// first one's fields. +func resetForRetry(out interface{}) { + value := reflect.ValueOf(out) + if value.Kind() != reflect.Ptr || value.IsNil() { + return + } + value.Elem().Set(reflect.Zero(value.Elem().Type())) +} + +func postJSONOnce(ctx context.Context, client *http.Client, what, url, auth string, body, out interface{}, headers ...string) error { payload, err := json.Marshal(body) if err != nil { return err @@ -194,6 +230,20 @@ func postJSON(ctx context.Context, client *http.Client, what, url, auth string, return fmt.Errorf("%s: HTTP %d, unreadable body: %w", what, resp.StatusCode, err) } if resp.StatusCode/100 != 2 { + // Too many requests is the one status that is not about this request. + // The provider's own sentence is unusable here — Groq's reads + // + // "Rate limit reached for model `openai/gpt-oss-120b` in organization + // `org_01m38x8s72e759kn6g88ve2dhj` service tier `on_demand` on tokens + // per minute (TPM): Limit 8000, Used 7320…" + // + // which is shown to a shopkeeper as Buddy's answer, names our billing + // account, and tells them nothing they can act on. It is also usually + // over within a second, so the honest handling is to wait and try again + // rather than to report it at all. + if resp.StatusCode == http.StatusTooManyRequests { + return &tooManyRequests{what: what, after: retryAfter(resp, raw)} + } // The decoded body carries the provider's message where there is one; // this is the fallback for a bare status. if msg := extractMessage(raw); msg != "" { @@ -230,3 +280,106 @@ func VectorLiteral(v []float32) string { b.WriteByte(']') return b.String() } + +// ── Being rate limited ────────────────────────────────────────────────────── +// +// A shared provider quota is not a fault in the request that happened to hit +// it, and on Groq's free tier it is reached by ordinary use: 8,000 tokens a +// minute is three or four Buddy questions. The waits are short — the provider +// states them in milliseconds — so one retry turns almost all of them into a +// slightly slower answer instead of an error. + +// ErrBusy is what a caller sees when the wait did not help. +// +// Sentinel so the HTTP layer can answer 429 and the console can say "a moment" +// rather than rendering a provider's billing details as an answer. +var ErrBusy = errors.New("the assistant is busy right now — try again in a moment") + +type tooManyRequests struct { + what string + after time.Duration +} + +func (e *tooManyRequests) Error() string { return e.what + ": " + ErrBusy.Error() } +func (e *tooManyRequests) Unwrap() error { return ErrBusy } + +// maxRetryWait bounds how long a request may be held. Beyond this the honest +// answer is "busy" — a person watching a spinner has already decided something +// is broken, and the provider's own suggestion can be a minute on a hard quota. +const maxRetryWait = 3 * time.Second + +// retryAfter reads how long the provider asked us to wait. +// +// `Retry-After` first, because it is the standard and a proxy may add it where +// the body has nothing. Groq puts the number in prose instead — "Please try +// again in 840ms" — so that is read next. Zero means "no idea", and the caller +// uses its own floor rather than hammering immediately. +func retryAfter(resp *http.Response, raw []byte) time.Duration { + if header := strings.TrimSpace(resp.Header.Get("Retry-After")); header != "" { + // Seconds, as an integer, is the only form worth reading: the HTTP-date + // form is for caches and no model provider sends it. + if secs, err := strconv.ParseFloat(header, 64); err == nil && secs > 0 { + return time.Duration(secs * float64(time.Second)) + } + } + return waitFromMessage(extractMessage(raw)) +} + +// waitFromMessage pulls "try again in 840ms" or "try again in 1.5s" out of prose. +// +// Its own function because it is the part worth testing: the wording comes from +// somebody else's error strings and is the first thing that will change. +func waitFromMessage(message string) time.Duration { + lower := strings.ToLower(message) + marker := "try again in " + at := strings.Index(lower, marker) + if at < 0 { + return 0 + } + + rest := lower[at+len(marker):] + end := 0 + for end < len(rest) && (rest[end] == '.' || (rest[end] >= '0' && rest[end] <= '9')) { + end++ + } + if end == 0 { + return 0 + } + amount, err := strconv.ParseFloat(rest[:end], 64) + if err != nil || amount <= 0 { + return 0 + } + + switch { + case strings.HasPrefix(rest[end:], "ms"): + return time.Duration(amount * float64(time.Millisecond)) + case strings.HasPrefix(rest[end:], "s"): + return time.Duration(amount * float64(time.Second)) + } + return 0 +} + +// waitBeforeRetry sleeps for what the provider asked, bounded, and reports +// whether waiting is worth it at all. +// +// Returns false when the ask is longer than we are prepared to hold a request +// for, or when the caller's context is done — a retry after the browser has +// given up is work nobody will see. +func waitBeforeRetry(ctx context.Context, after time.Duration) bool { + if after <= 0 { + // No stated wait. A short one anyway: retrying instantly on a quota is + // how a burst becomes two bursts. + after = 250 * time.Millisecond + } + if after > maxRetryWait { + return false + } + timer := time.NewTimer(after) + defer timer.Stop() + select { + case <-ctx.Done(): + return false + case <-timer.C: + return true + } +} diff --git a/utils/embedding_test.go b/utils/embedding_test.go index 0a7653a..cfe00d3 100644 --- a/utils/embedding_test.go +++ b/utils/embedding_test.go @@ -3,6 +3,7 @@ package utils import ( "context" "encoding/json" + "errors" "net/http" "net/http/httptest" "strings" @@ -73,16 +74,70 @@ func TestGeminiEmbedderSendsTheRequestTheAPIExpects(t *testing.T) { } func TestEmbedderSurfacesProviderErrors(t *testing.T) { + // Used to assert this with a 429, which is now the one status that is + // deliberately NOT passed through — see the two tests below. The rule it was + // written for still holds for every other failure: a provider that explains + // itself should not have that explanation swallowed. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(429) - w.Write([]byte(`{"error":{"message":"Rate limit reached"}}`)) + w.WriteHeader(400) + w.Write([]byte(`{"error":{"message":"input exceeds the maximum token count"}}`)) })) defer srv.Close() e, _ := NewEmbedder(config.EmbeddingConfig{Provider: "openai", Model: "m", APIKey: "k", BaseURL: srv.URL}) _, err := e.Embed(context.Background(), "x") - if err == nil || !strings.Contains(err.Error(), "429") || !strings.Contains(err.Error(), "Rate limit") { - t.Fatalf("want a 429 with the provider's message, got %v", err) + if err == nil || !strings.Contains(err.Error(), "400") || !strings.Contains(err.Error(), "maximum token count") { + t.Fatalf("want a 400 with the provider's message, got %v", err) + } +} + +func TestARateLimitIsWaitedOutRatherThanReported(t *testing.T) { + // The common case on a free tier: the quota clears in under a second, so + // the right answer is a slightly slower success rather than an error. + var calls int + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + if calls == 1 { + w.WriteHeader(429) + w.Write([]byte(`{"error":{"message":"Rate limit reached. Please try again in 20ms."}}`)) + return + } + w.Write([]byte(`{"data":[{"embedding":[0.5,0.25]}]}`)) + })) + defer srv.Close() + + e, _ := NewEmbedder(config.EmbeddingConfig{Provider: "openai", Model: "m", APIKey: "k", BaseURL: srv.URL}) + vector, err := e.Embed(context.Background(), "x") + if err != nil { + t.Fatalf("a 20ms quota wait became an error: %v", err) + } + if len(vector) != 2 { + t.Fatalf("the retry did not return the answer: %v", vector) + } + if calls != 2 { + t.Fatalf("expected one retry, saw %d calls", calls) + } +} + +func TestAPersistentRateLimitIsReportedWithoutTheProvidersDetails(t *testing.T) { + // When waiting does not help, the person is told to try again — not handed + // our organisation id and the tokens-per-minute arithmetic behind it. + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(429) + w.Write([]byte(`{"error":{"message":"Rate limit reached for model X in organization org_01m38 on tokens per minute (TPM): Limit 8000. Please try again in 20ms."}}`)) + })) + defer srv.Close() + + e, _ := NewEmbedder(config.EmbeddingConfig{Provider: "openai", Model: "m", APIKey: "k", BaseURL: srv.URL}) + _, err := e.Embed(context.Background(), "x") + + if !errors.Is(err, ErrBusy) { + t.Fatalf("a persistent quota was not reported as busy: %v", err) + } + for _, leaked := range []string{"org_01m38", "TPM", "8000"} { + if strings.Contains(err.Error(), leaked) { + t.Fatalf("%q reached the caller: %s", leaked, err.Error()) + } } } diff --git a/utils/ratelimit_test.go b/utils/ratelimit_test.go new file mode 100644 index 0000000..91bc63d --- /dev/null +++ b/utils/ratelimit_test.go @@ -0,0 +1,131 @@ +package utils + +import ( + "context" + "errors" + "net/http" + "testing" + "time" +) + +/* +Being rate limited by the model provider. + +Groq's free tier is 8,000 tokens a minute, which is three or four Buddy +questions — so this is ordinary use, not an edge case. Its message reads: + + Rate limit reached for model `openai/gpt-oss-120b` in organization + `org_01m38x8s72e759kn6g88ve2dhj` service tier `on_demand` on tokens per + minute (TPM): Limit 8000, Used 7320, Requested 792. Please try again in 840ms. + +That sentence was being returned to the console as Buddy's ANSWER. It names our +billing account, it is about arithmetic the shopkeeper cannot influence, and the +condition it describes is usually over in under a second. +*/ + +func TestTheWaitIsReadOutOfTheProvidersProse(t *testing.T) { + // The part most likely to change, because it is somebody else's wording. + for _, tc := range []struct { + message string + want time.Duration + }{ + {"Please try again in 840ms.", 840 * time.Millisecond}, + {"please try again in 1.5s", 1500 * time.Millisecond}, + {"Try again in 2s. Need more tokens?", 2 * time.Second}, + {"Rate limit reached … Please try again in 397.499999ms.", 397499999 * time.Nanosecond}, + } { + got := waitFromMessage(tc.message) + // Milliseconds is the resolution that matters; the fractional tail of + // 397.499999ms is noise. + if (got - tc.want).Abs() > time.Millisecond { + t.Fatalf("%q read as %v, want %v", tc.message, got, tc.want) + } + } +} + +func TestProseWithNoWaitInItReadsAsUnknown(t *testing.T) { + // Zero means "no idea", and the caller uses its own floor. Guessing a + // number here would be inventing one. + for _, message := range []string{ + "", "Rate limit reached.", "try again later", "try again in soon", "try again in 0ms", + } { + if got := waitFromMessage(message); got != 0 { + t.Fatalf("%q produced a wait of %v", message, got) + } + } +} + +func TestTheRetryAfterHeaderWinsOverTheProse(t *testing.T) { + // It is the standard, and a proxy can add it where the body has nothing. + resp := &http.Response{Header: http.Header{}} + resp.Header.Set("Retry-After", "2") + + got := retryAfter(resp, []byte(`{"error":{"message":"try again in 840ms"}}`)) + if got != 2*time.Second { + t.Fatalf("header ignored: got %v", got) + } +} + +func TestAnUnwaitableLimitIsNotWaitedFor(t *testing.T) { + // A hard quota can suggest a minute. Holding a request that long is worse + // than saying "busy" — the person watching the spinner decided it was + // broken long before it returned. + if waitBeforeRetry(context.Background(), time.Minute) { + t.Fatal("agreed to hold the request for a minute") + } +} + +func TestNothingIsWaitedForOnceTheCallerHasGone(t *testing.T) { + // A retry after the browser gave up is work nobody will see. + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + if waitBeforeRetry(ctx, 10*time.Millisecond) { + t.Fatal("waited although the request was already abandoned") + } +} + +func TestAShortWaitIsHonoured(t *testing.T) { + start := time.Now() + if !waitBeforeRetry(context.Background(), 30*time.Millisecond) { + t.Fatal("refused a 30ms wait") + } + if time.Since(start) < 25*time.Millisecond { + t.Fatal("returned without waiting") + } +} + +func TestBeingBusyIsRecognisableWithoutReadingTheText(t *testing.T) { + // The HTTP layer answers 429 on this, and the console tells the person to + // try again. Matching on the provider's wording instead would break the + // moment Groq rephrases it. + err := error(&tooManyRequests{what: "assistant", after: time.Second}) + + if !errors.Is(err, ErrBusy) { + t.Fatal("a rate-limited call is not recognisable as busy") + } +} + +func TestTheProvidersBillingDetailsAreNotInTheMessage(t *testing.T) { + // The whole point. Whatever the provider said, this is what a shopkeeper + // reads. + err := error(&tooManyRequests{what: "assistant", after: 840 * time.Millisecond}) + + for _, leaked := range []string{"org_", "TPM", "8000", "tier", "billing"} { + if contains(err.Error(), leaked) { + t.Fatalf("%q reaches the console: %s", leaked, err.Error()) + } + } +} + +func contains(haystack, needle string) bool { + return len(needle) > 0 && len(haystack) >= len(needle) && + func() bool { + for i := 0; i+len(needle) <= len(haystack); i++ { + if haystack[i:i+len(needle)] == needle { + return true + } + } + return false + }() +}