From c9e16c7169f92143bec88c664653bbe0258ba122 Mon Sep 17 00:00:00 2001 From: Andriy Oblivantsev Date: Thu, 17 Sep 2026 12:39:25 +0000 Subject: [PATCH] =?UTF-8?q?fix(retry):=20=D0=B3=D0=BB=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=BB=D1=8C=D0=BD=D1=8B=D0=B9=20rate-limit=20+=20Retry-After?= =?UTF-8?q?=20+=20cooldown=20(#70)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- AGENTS.md | 2 +- CHANGELOG.md | 8 ++ README.md | 20 ++-- auth.go | 2 +- client.go | 8 +- docs/README.md | 3 + docs/rate-limiting.md | 55 +++++++++++ files_webdav.go | 4 +- http.go | 12 +-- mails.go | 2 +- ratelimit.go | 210 ++++++++++++++++++++++++++++++++++++++++++ ratelimit_test.go | 160 ++++++++++++++++++++++++++++++++ request.go | 2 +- retry.go | 140 +++++++++++++++++++++++++--- retry_test.go | 132 ++++++++++++++++++++++++++ storage_fallback.go | 2 +- 16 files changed, 729 insertions(+), 33 deletions(-) create mode 100644 docs/rate-limiting.md create mode 100644 ratelimit.go create mode 100644 ratelimit_test.go diff --git a/AGENTS.md b/AGENTS.md index e322a39..b69277f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -9,7 +9,7 @@ Canonical Go client for OnlyOffice Workspace (Projects + Calendar + CRM) and the - `request.go` — `Request`, `Query`, `Time`, `Token`, `MetaResponse`, `Permissions`. - `auth.go` — `Authenticate`, `AuthenticateContext`, `InvalidateToken`, `Auth`, token lifecycle. - `http.go` — transport + DRY response decoders (`ResponseArray`/`ResponseObject`/`postFormObject`/`putFormObject`/`deleteObject`). - - `projects.go`, `tasks.go`, `users.go`, `calendar.go`, `crm.go`, `files.go`, `files_webdav.go`, `files_stem.go`, `retry.go`, `mails.go`, `invoices.go` — typed / untyped domain methods. **`files.go`** — CRM opportunity upload plus **project/task Documents** (`UpdateFile`, `UploadToFolderReplacing`). **`files_webdav.go`** — Documents module by id (`ListDavFolder`, `MoveDavItems`/`CopyDavItems` with per-operation error surfacing, `ListFileOps`). **`retry.go`** — `DoRetry`: deterministic linear backoff (no jitter) on 429/502/503/504; every bulk tool routes API calls through it, and the HTTP transport + auth (`retryRaw`, `AuthenticateContext`) retry transient answers centrally. **`mails.go`** — OnlyOffice Workspace Mail. **`invoices.go`** — CRM invoices, PDF regen/cleanup, status. Association rules: [`docs/crm-associations.md`](docs/crm-associations.md). + - `projects.go`, `tasks.go`, `users.go`, `calendar.go`, `crm.go`, `files.go`, `files_webdav.go`, `files_stem.go`, `retry.go`, `mails.go`, `invoices.go` — typed / untyped domain methods. **`files.go`** — CRM opportunity upload plus **project/task Documents** (`UpdateFile`, `UploadToFolderReplacing`). **`files_webdav.go`** — Documents module by id (`ListDavFolder`, `MoveDavItems`/`CopyDavItems` with per-operation error surfacing, `ListFileOps`). **`retry.go`** — `DoRetry`: deterministic exponential backoff (no jitter) on 429/502/503/504; every bulk tool routes API calls through it, and the HTTP transport + auth (`retryRaw`, `AuthenticateContext`) retry transient answers centrally. `ratelimit.go` adds a process-wide token bucket (`OO_RATE_LIMIT`/`OO_BURST`) and a shared 429 cooldown gate, installed via `pacedTransport` in `NewClient`; `Retry-After` is parsed into `*TransientError` and honoured. See [`docs/rate-limiting.md`](docs/rate-limiting.md). **`mails.go`** — OnlyOffice Workspace Mail. **`invoices.go`** — CRM invoices, PDF regen/cleanup, status. Association rules: [`docs/crm-associations.md`](docs/crm-associations.md). - **Unified file client (epic #34) — `file_core.go`, `file_rest.go`, `file_dav.go`, `file_pg.go`, `file_es.go`, `file_es_text.go`, `file_text_index.go`, `file_facade.go`.** `file_core.go` — model (`Entry`, `Kind`) + `FileStore`/`Searcher`; `file_rest.go`/`file_dav.go` — REST/WebDAV adapters; `file_pg.go` — **read-only** SQL store (PostgreSQL/MySQL, `ErrReadOnly` on writes); `file_es.go` — OnlyOffice Elasticsearch searcher; `file_es_text.go`/`file_text_index.go` — own PDF/scan index (`oo_docs_text`, PDF attachments via pdfdetach); `file_facade.go` — `FileClient` with read/write/search order and transient fallback. Use `c.Files()` (facade), `c.FileStore("rest"|"dav"|"pg"|"sql")` or `c.SQLFileStore()`; contract and how to add a backend: [`docs/unified-file-client.md`](docs/unified-file-client.md). - Pure stdlib + `google/go-querystring`; no UI, no dotenv. - **CLI — `cmd/oo/` as `package main`.** Cobra wrapper that loads `.env` via `godotenv` at startup. **Subject-based command tree** mirroring [`tea`](https://gitea.com/gitea/tea): diff --git a/CHANGELOG.md b/CHANGELOG.md index ed506c5..1b1c6d1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,14 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## Unreleased + +### Features + +* **retry:** global token-bucket rate limit (`OO_RATE_LIMIT`/`OO_BURST`), typed + `TransientError` with `Retry-After`, exponential backoff + (`OO_RETRY_ATTEMPTS`/`_BASE`/`_MAX`) and a process-wide 429 cooldown gate. + All HTTP paths are paced via `pacedTransport`. + ## [0.18.0](https://github.com/eSlider/go-onlyoffice/compare/v0.17.0...v0.18.0) (2026-09-04) diff --git a/README.md b/README.md index 53de82e..42badd6 100644 --- a/README.md +++ b/README.md @@ -522,9 +522,9 @@ type Task struct { | `ListFileOps(ctx)` | Active file operations (move/copy status polling) | | `FolderFiles(ctx, folderID)` | Flat file list of a folder (stem helpers) | | `DeleteFilesByStem(ctx, folderID, stem)` | Remove `stem\|ext` copies | -| `DoRetry(ctx, policy, fn)` | Deterministic linear backoff (N·Base, no jitter) on 429/502/503/504 | -| `DefaultRetryPolicy()` | 5 attempts, 1s·2s·3s·4s waits, 30s cap | -| `Transient(err)` | True for retriable OnlyOffice answers | +| `DoRetry(ctx, policy, fn)` | Deterministic exponential backoff (`Base·2^(N-1)`, no jitter) on 429/502/503/504; honours `Retry-After` and the process-wide cooldown gate | +| `DefaultRetryPolicy()` | From env: 7 attempts, 2s base, 2m cap (`OO_RETRY_ATTEMPTS/_BASE/_MAX`) | +| `Transient(err)` | True for retriable OnlyOffice answers (`*TransientError` or HTTP 429/502/503/504 text) | ### Helper Types @@ -732,9 +732,12 @@ fallback rules, env names and how to add a backend: ### Bulk tools (`cmd/`) Small single-purpose binaries for bulk Documents work. All of them pace -requests and retry transient OnlyOffice answers (429/502/503/504) with a -deterministic linear backoff — no jitter, same waits on every run -(see `DoRetry` below). Build with `go build ./cmd/`. +requests through a process-wide token bucket (default ~4 req/s, `OO_RATE_LIMIT`/ +`OO_BURST`) and retry transient OnlyOffice answers (429/502/503/504) with a +deterministic exponential backoff — no jitter, same waits on every run. A 429 +opens a shared cooldown gate and `Retry-After` is honoured (see `DoRetry` and +[`docs/rate-limiting.md`](docs/rate-limiting.md)). Build with +`go build ./cmd/`. ```bash ooscan 659 # recursive index → TSV: file_id, folder_id, path, title @@ -999,6 +1002,11 @@ oo projects files list 33 | `ONLYOFFICE_CALENDAR_ID` | Default calendar id used when omitted (default `1`) | | `ONLYOFFICE_PROJECT_ID` | Default project id used when omitted (default `33`) | | `OO_URL`, `OO_USER`, `OO_PASS` | Optional CLI-only aliases for `ONLYOFFICE_*` | +| `OO_RATE_LIMIT` | Process-wide request pacing, req/s (default `4`; `0` disables) | +| `OO_BURST` | Token-bucket burst (default `1`) | +| `OO_RETRY_ATTEMPTS` | Transient retries, total attempts (default `7`) | +| `OO_RETRY_BASE` | Exponential backoff base (default `2s`) | +| `OO_RETRY_MAX` | Backoff cap (default `2m`) | | `ONLYOFFICE_ES_URL` | Elasticsearch base URL (`oo search`, own index); see [`docs/elasticsearch.md`](docs/elasticsearch.md) | | `ONLYOFFICE_ES_INDEX` | OnlyOffice index (default `files_file`) | | `ONLYOFFICE_ES_TEXT_INDEX` | Own PDF/scan index (default `oo_docs_text`) | diff --git a/auth.go b/auth.go index 3d2a67a..ed897ad 100644 --- a/auth.go +++ b/auth.go @@ -83,7 +83,7 @@ func (c *Client) authenticateOnce(ctx context.Context) error { return err } if resp.StatusCode >= 400 { - return fmt.Errorf("auth: %d %s", resp.StatusCode, truncate(string(raw), 400)) + return statusError(resp.StatusCode, retryAfterOf(resp), "auth: %d %s", resp.StatusCode, truncate(string(raw), 400)) } var env struct { Response *Token `json:"response"` diff --git a/client.go b/client.go index 4997f27..cb9e5ba 100644 --- a/client.go +++ b/client.go @@ -38,11 +38,15 @@ type Client struct { folderTitlesMu sync.Mutex } -// NewClient returns a new Client backed by http.DefaultClient. +// NewClient returns a new Client whose transport is paced by the process-wide +// rate limiter and 429 cooldown gate (OO_RATE_LIMIT/OO_BURST; see ratelimit.go). func NewClient(c Credentials) *Client { jar, _ := cookiejar.New(nil) return &Client{ - client: &http.Client{Jar: jar}, + client: &http.Client{ + Jar: jar, + Transport: &pacedTransport{base: http.DefaultTransport}, + }, credentials: &c, } } diff --git a/docs/README.md b/docs/README.md index f8d4826..56030ec 100644 --- a/docs/README.md +++ b/docs/README.md @@ -22,6 +22,9 @@ related: - [rclone-webdav.md](rclone-webdav.md) — rclone-монтирование Documents (`deploy/docker-compose.rclone-webdav.yml`), smoke, ограничения. - [crm-associations.md](crm-associations.md) — правила ассоциаций CRM. +- [rate-limiting.md](rate-limiting.md) — rate limit, exponential backoff, + `Retry-After`, общий cooldown против 429; env `OO_RATE_LIMIT`/`OO_BURST`/ + `OO_RETRY_*`. ## Тесты Команды и туннели — раздел Testing в [README.md](../README.md#testing). diff --git a/docs/rate-limiting.md b/docs/rate-limiting.md new file mode 100644 index 0000000..aa35eda --- /dev/null +++ b/docs/rate-limiting.md @@ -0,0 +1,55 @@ +--- +type: reference +status: current +related: + - README.md + - ../AGENTS.md +--- + +# Rate limit, backoff и cooldown + +Устойчивость к 429 (openresty). Всё встроено в библиотеку — отдельный пакет не +нужен. Реализация: `ratelimit.go`, `retry.go`. + +## Что происходит с каждым запросом + +1. **Cooldown-гейт** — общий на процесс. Если недавно пришёл 429, все запросы + ждут конца окна. +2. **Rate limiter** — token bucket на процесс. Пейсит все HTTP-пути: листинг, + создание папок, загрузку, `get project`, auth. +3. Запрос уходит. +4. Ответ ≥400 → `*TransientError` (для 429/502/503/504) с `Retry-After`. +5. `DoRetry` — экспоненциальный backoff, без jitter. +6. `Retry-After` длиннее backoff → ждём его; окно уходит в общий cooldown. + +Установлено в `NewClient` через `pacedTransport`; отдельный код трогать не надо. + +## Env + +| Переменная | Default | Смысл | +|---|---|---| +| `OO_RATE_LIMIT` | `4` | запросов/с на процесс; `0` — лимитер выключен | +| `OO_BURST` | `1` | запас токенов token bucket | +| `OO_RETRY_ATTEMPTS` | `7` | всего попыток, включая первую | +| `OO_RETRY_BASE` | `2s` | база экспоненты: ждать перед попыткой N = `Base*2^(N-1)` | +| `OO_RETRY_MAX` | `2m` | потолок ожидания | + +Битые значения → default. `OO_RETRY_*` — формат `time.ParseDuration` +(`2s`, `30s`, `2m`). + +## Правила + +- Детерминированно, без jitter — повторный прогон ждёт столько же. +- `Retry-After` — секунды (`120`) или HTTP-date. +- Cooldown общий: параллельные и последовательные вызовы не бьют в стену. +- Backoff cap не ограничивает `Retry-After` — серверу верим больше. +- Только stdlib. + +## Когда руками снять нагрузку + +`OO_RATE_LIMIT` ниже (`2`), `OO_BURST=1`; при массовом apply — батчами. + +## Тесты + +`retry_test.go` — `Retry-After`, экспонента, cap; `ratelimit_test.go` — burst, +`OO_RATE_LIMIT=0`, cooldown. Фейковый сервер отдаёт 429 с заголовком. diff --git a/files_webdav.go b/files_webdav.go index a38b732..f7f31ec 100644 --- a/files_webdav.go +++ b/files_webdav.go @@ -325,7 +325,7 @@ func (c *Client) deleteJSON(ctx context.Context, path string, body any) (json.Ra return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("DELETE %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "DELETE %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -365,7 +365,7 @@ func (c *Client) uploadReader(ctx context.Context, path, fieldName, fileName str return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("upload %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "upload %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } diff --git a/http.go b/http.go index 1452d3d..29d5f21 100644 --- a/http.go +++ b/http.go @@ -171,7 +171,7 @@ func (c *Client) getJSONOnce(ctx context.Context, path string) (json.RawMessage, return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("GET %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "GET %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -219,7 +219,7 @@ func (c *Client) formRequestOnce(ctx context.Context, method, path string, field return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("%s form %s: %d %s", method, path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "%s form %s: %d %s", method, path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -250,7 +250,7 @@ func (c *Client) deleteReqOnce(ctx context.Context, path string) (json.RawMessag return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("DELETE %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "DELETE %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -315,7 +315,7 @@ func (c *Client) postJSONOnce(ctx context.Context, path string, body any) (json. return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("POST JSON %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "POST JSON %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -362,7 +362,7 @@ func (c *Client) putJSONOnce(ctx context.Context, path string, body any) (json.R return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("PUT JSON %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "PUT JSON %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } @@ -423,7 +423,7 @@ func (c *Client) uploadMultipartOnce(ctx context.Context, method, path, fieldNam return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("upload %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "upload %s: %d %s", path, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } diff --git a/mails.go b/mails.go index d15f821..cc17ec9 100644 --- a/mails.go +++ b/mails.go @@ -125,7 +125,7 @@ func (c *Client) DownloadMailAttachment(ctx context.Context, attachmentID string return nil, err } if resp.StatusCode >= 400 { - return nil, fmt.Errorf("DownloadMailAttachment %s: %d %s", id, resp.StatusCode, truncate(string(raw), 400)) + return nil, statusError(resp.StatusCode, retryAfterOf(resp), "DownloadMailAttachment %s: %d %s", id, resp.StatusCode, truncate(string(raw), 400)) } return raw, nil } diff --git a/ratelimit.go b/ratelimit.go new file mode 100644 index 0000000..f71acd9 --- /dev/null +++ b/ratelimit.go @@ -0,0 +1,210 @@ +package onlyoffice + +// Client-side pacing and the process-wide 429 gate. A token bucket applies to +// every HTTP request a *Client makes (via pacedTransport), so listing, folder +// creation, uploads and even `get project` share one budget and cannot burst +// past openresty. A typed transient answer carrying Retry-After opens the +// cooldown gate for the whole process. + +import ( + "context" + "net/http" + "os" + "strconv" + "strings" + "sync" + "time" +) + +// Default pacing: ~4 requests/second with a single-request burst. This keeps +// bulk syncs under openresty's threshold without noticeably slowing one-shot +// tools. +const ( + defaultRateLimit = 4.0 + defaultBurst = 1 +) + +// rateLimiter is a deterministic token bucket. A nil *rateLimiter is a no-op, +// which is how OO_RATE_LIMIT=0 disables pacing. +type rateLimiter struct { + mu sync.Mutex + rate float64 + burst float64 + tokens float64 + last time.Time +} + +// newRateLimiter returns a limiter at rate requests/second with the given +// burst. rate <= 0 disables pacing and returns nil. +func newRateLimiter(rate float64, burst int) *rateLimiter { + if rate <= 0 { + return nil + } + if burst < 1 { + burst = 1 + } + return &rateLimiter{ + rate: rate, + burst: float64(burst), + tokens: float64(burst), + last: time.Now(), + } +} + +// wait consumes one token, blocking until one is available or ctx is done. +func (l *rateLimiter) wait(ctx context.Context) error { + if l == nil { + return nil + } + for { + l.mu.Lock() + now := time.Now() + l.tokens += now.Sub(l.last).Seconds() * l.rate + if l.tokens > l.burst { + l.tokens = l.burst + } + l.last = now + if l.tokens >= 1 { + l.tokens-- + l.mu.Unlock() + return nil + } + need := time.Duration((1 - l.tokens) / l.rate * float64(time.Second)) + l.mu.Unlock() + if need < time.Millisecond { + need = time.Millisecond + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(need): + } + } +} + +// cooldownGate is the process-wide 429 gate: after a throttled answer every +// request waits until the gate opens, so concurrent and sequential callers do +// not hammer the server. +type cooldownGate struct { + mu sync.Mutex + until time.Time +} + +// note extends the gate to at least now+d. +func (g *cooldownGate) note(d time.Duration) { + if d <= 0 { + return + } + g.mu.Lock() + if u := time.Now().Add(d); u.After(g.until) { + g.until = u + } + g.mu.Unlock() +} + +// wait blocks until the gate opens or ctx is done. +func (g *cooldownGate) wait(ctx context.Context) error { + g.mu.Lock() + d := time.Until(g.until) + g.mu.Unlock() + if d <= 0 { + return nil + } + t := time.NewTimer(d) + defer t.Stop() + select { + case <-ctx.Done(): + return ctx.Err() + case <-t.C: + return nil + } +} + +var ( + globalLimiter *rateLimiter + globalLimiterOn bool + globalLimiterMu sync.Mutex + globalCooldown = &cooldownGate{} +) + +// limiter returns the process-wide rate limiter, initialising it from the +// environment on first use. OO_RATE_LIMIT=0 yields a nil (disabled) limiter. +func limiter() *rateLimiter { + globalLimiterMu.Lock() + defer globalLimiterMu.Unlock() + if !globalLimiterOn { + globalLimiter = newRateLimiterFromEnv() + globalLimiterOn = true + } + return globalLimiter +} + +// newRateLimiterFromEnv builds a limiter from OO_RATE_LIMIT (requests/second, +// default 4; 0 disables) and OO_BURST (default 1). Malformed values fall back +// to the defaults. +func newRateLimiterFromEnv() *rateLimiter { + rate := defaultRateLimit + if v, ok := os.LookupEnv("OO_RATE_LIMIT"); ok { + if f, err := strconv.ParseFloat(strings.TrimSpace(v), 64); err == nil { + rate = f + } + } + burst := defaultBurst + if v, ok := os.LookupEnv("OO_BURST"); ok { + if n, err := strconv.Atoi(strings.TrimSpace(v)); err == nil { + burst = n + } + } + return newRateLimiter(rate, burst) +} + +// pacedTransport applies the process-wide cooldown gate and rate limiter to +// every request, and arms the gate on a 429. Installed by NewClient, so all +// HTTP paths — typed Query, the http.go helpers, auth and WebDAV — are paced. +type pacedTransport struct { + base http.RoundTripper +} + +func (t *pacedTransport) RoundTrip(req *http.Request) (*http.Response, error) { + if err := globalCooldown.wait(req.Context()); err != nil { + return nil, err + } + if err := limiter().wait(req.Context()); err != nil { + return nil, err + } + resp, err := t.base.RoundTrip(req) + if err != nil { + return nil, err + } + if resp.StatusCode == http.StatusTooManyRequests { + globalCooldown.note(retryAfterOf(resp)) + } + return resp, nil +} + +// envInt reads a positive integer env var, falling back to def. +func envInt(name string, def int) int { + v, ok := os.LookupEnv(name) + if !ok { + return def + } + n, err := strconv.Atoi(strings.TrimSpace(v)) + if err != nil || n <= 0 { + return def + } + return n +} + +// envDuration reads a duration env var (Go syntax, e.g. "2s", "2m"), +// falling back to def. +func envDuration(name string, def time.Duration) time.Duration { + v, ok := os.LookupEnv(name) + if !ok { + return def + } + d, err := time.ParseDuration(strings.TrimSpace(v)) + if err != nil || d <= 0 { + return def + } + return d +} diff --git a/ratelimit_test.go b/ratelimit_test.go new file mode 100644 index 0000000..9582ebd --- /dev/null +++ b/ratelimit_test.go @@ -0,0 +1,160 @@ +package onlyoffice + +import ( + "context" + "net/http" + "net/http/httptest" + "os" + "testing" + "time" +) + +// TestMain disables the process-wide rate limiter for the suite so existing +// fake-server tests stay fast; the pacing tests configure their own limiters. +func TestMain(m *testing.M) { + _ = os.Setenv("OO_RATE_LIMIT", "0") + os.Exit(m.Run()) +} + +// resetPacingForTest restores the package-level limiter and cooldown gate to a +// clean state. White-box helper for tests that arm the global 429 gate. +func resetPacingForTest() { + globalLimiterMu.Lock() + globalLimiter = nil + globalLimiterOn = false + globalLimiterMu.Unlock() + globalCooldown = &cooldownGate{} +} + +func TestNewClientUsesPacedTransport(t *testing.T) { + c := NewClient(Credentials{Url: "https://example.test"}) + if _, ok := c.client.Transport.(*pacedTransport); !ok { + t.Fatalf("transport = %T, want *pacedTransport", c.client.Transport) + } +} + +func TestNewRateLimiterFromEnv(t *testing.T) { + t.Setenv("OO_RATE_LIMIT", "") + t.Setenv("OO_BURST", "") + if l := newRateLimiterFromEnv(); l == nil || l.rate != defaultRateLimit || l.burst != defaultBurst { + t.Fatalf("defaults: got %+v, want rate=%v burst=%d", l, defaultRateLimit, defaultBurst) + } + + t.Setenv("OO_RATE_LIMIT", "0") + if l := newRateLimiterFromEnv(); l != nil { + t.Fatalf("OO_RATE_LIMIT=0: got %+v, want nil (off)", l) + } + + t.Setenv("OO_RATE_LIMIT", "10") + t.Setenv("OO_BURST", "3") + if l := newRateLimiterFromEnv(); l == nil || l.rate != 10 || l.burst != 3 { + t.Fatalf("explicit: got %+v, want rate=10 burst=3", l) + } +} + +func TestNilRateLimiterIsNoOp(t *testing.T) { + var l *rateLimiter + start := time.Now() + if err := l.wait(context.Background()); err != nil { + t.Fatalf("nil limiter wait: %v", err) + } + if elapsed := time.Since(start); elapsed > 20*time.Millisecond { + t.Fatalf("nil limiter blocked for %v", elapsed) + } +} + +func TestRateLimiterPacesBurst(t *testing.T) { + l := newRateLimiter(50, 1) // 20ms per token, no burst headroom + start := time.Now() + for i := 0; i < 5; i++ { + if err := l.wait(context.Background()); err != nil { + t.Fatalf("wait %d: %v", i, err) + } + } + // 5 calls with burst 1 -> 4 gaps * 20ms = 80ms. + if elapsed := time.Since(start); elapsed < 60*time.Millisecond { + t.Fatalf("elapsed %v, want >= 60ms (not paced)", elapsed) + } +} + +func TestRateLimiterHonoursContext(t *testing.T) { + l := newRateLimiter(1, 1) + if err := l.wait(context.Background()); err != nil { + t.Fatalf("first wait: %v", err) + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond) + defer cancel() + if err := l.wait(ctx); err == nil { + t.Fatal("want context error while waiting for a token") + } +} + +func TestCooldownGateBlocksAndExpires(t *testing.T) { + g := &cooldownGate{} + g.note(50 * time.Millisecond) + start := time.Now() + if err := g.wait(context.Background()); err != nil { + t.Fatalf("wait: %v", err) + } + if elapsed := time.Since(start); elapsed < 45*time.Millisecond { + t.Fatalf("gate did not block: %v", elapsed) + } + start = time.Now() + if err := g.wait(context.Background()); err != nil { + t.Fatalf("second wait: %v", err) + } + if elapsed := time.Since(start); elapsed > 20*time.Millisecond { + t.Fatalf("expired gate still blocked: %v", elapsed) + } +} + +func TestCooldownGateNoteExtendsNotShortens(t *testing.T) { + g := &cooldownGate{} + g.note(80 * time.Millisecond) + g.note(10 * time.Millisecond) // must not shorten the window + if d := time.Until(g.until); d < 70*time.Millisecond { + t.Fatalf("note shortened gate: %v left", d) + } +} + +// TestPacedTransportArmsCooldownOn429 verifies that a throttled answer opens the +// process-wide gate and that the next request is blocked for the Retry-After +// window even though it targets a healthy endpoint. +func TestPacedTransportArmsCooldownOn429(t *testing.T) { + resetPacingForTest() + defer resetPacingForTest() + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/throttle" { + w.Header().Set("Retry-After", "1") + w.WriteHeader(http.StatusTooManyRequests) + return + } + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("{}")) + })) + defer srv.Close() + + tr := &pacedTransport{base: http.DefaultTransport} + + req1, _ := http.NewRequestWithContext(context.Background(), http.MethodGet, srv.URL+"/throttle", nil) + resp1, err := tr.RoundTrip(req1) + if err != nil { + t.Fatalf("throttled round trip: %v", err) + } + resp1.Body.Close() + if resp1.StatusCode != http.StatusTooManyRequests { + t.Fatalf("status = %d, want 429", resp1.StatusCode) + } + + start := time.Now() + req2, _ := http.NewRequestWithContext(context.Background(), http.MethodGet, srv.URL+"/ok", nil) + resp2, err := tr.RoundTrip(req2) + if err != nil { + t.Fatalf("second round trip: %v", err) + } + resp2.Body.Close() + if elapsed := time.Since(start); elapsed < 900*time.Millisecond { + t.Fatalf("cooldown did not block next call: %v", elapsed) + } +} diff --git a/request.go b/request.go index 74adb30..4c8d83d 100644 --- a/request.go +++ b/request.go @@ -102,7 +102,7 @@ func (c *Client) queryOnce(request Request, result interface{}) error { return err } if resp.StatusCode >= 400 { - return fmt.Errorf("%s %s: %d %s", request.GetMethod(), request.Uri, resp.StatusCode, truncate(string(raw), 400)) + return statusError(resp.StatusCode, retryAfterOf(resp), "%s %s: %d %s", request.GetMethod(), request.Uri, resp.StatusCode, truncate(string(raw), 400)) } if result == nil { return nil diff --git a/retry.go b/retry.go index bfa1555..587b333 100644 --- a/retry.go +++ b/retry.go @@ -3,38 +3,149 @@ package onlyoffice import ( "context" "encoding/json" + "errors" + "fmt" + "net/http" "regexp" + "strconv" + "strings" "time" ) -// RetryPolicy controls deterministic retries against OnlyOffice: fixed linear +// RetryPolicy controls deterministic retries against OnlyOffice: exponential // backoff without jitter, so repeated runs wait exactly the same schedule. // OnlyOffice throttles bulk reads/writes with 429 (and occasional 502/503/504 // from openresty), so every bulk tool routes API calls through DoRetry. type RetryPolicy struct { Attempts int // total attempts, including the first try - Base time.Duration // wait before retry N is N*Base - Max time.Duration // per-wait cap + Base time.Duration // wait before retry N is Base*2^(N-1) + Max time.Duration // per-wait cap; <=0 means no cap } -// DefaultRetryPolicy retries up to 5 times with 1s, 2s, 3s, 4s waits. +// DefaultRetryPolicy builds the policy from the environment, falling back to +// 7 attempts, a 2s base and a 2m cap: +// +// - OO_RETRY_ATTEMPTS (default 7) +// - OO_RETRY_BASE (duration, default 2s) +// - OO_RETRY_MAX (duration, default 2m) func DefaultRetryPolicy() RetryPolicy { - return RetryPolicy{Attempts: 5, Base: time.Second, Max: 30 * time.Second} + return RetryPolicy{ + Attempts: envInt("OO_RETRY_ATTEMPTS", 7), + Base: envDuration("OO_RETRY_BASE", 2*time.Second), + Max: envDuration("OO_RETRY_MAX", 2*time.Minute), + } } var transientRe = regexp.MustCompile(`:\s*(429|502|503|504)\b`) -// Transient reports whether err looks like a transient OnlyOffice answer -// (an HTTP 429/502/503/504 surfaced as "...: ..."). +// isTransientStatus reports whether an HTTP status is retriable at the edge. +func isTransientStatus(status int) bool { + switch status { + case http.StatusTooManyRequests, http.StatusBadGateway, + http.StatusServiceUnavailable, http.StatusGatewayTimeout: + return true + } + return false +} + +// TransientError is a typed transient answer from the HTTP layer. It carries +// the status code and, when present, the server's Retry-After delay so DoRetry +// can wait at least that long. +type TransientError struct { + Code int + RetryAfter time.Duration + Msg string +} + +func (e *TransientError) Error() string { return e.Msg } + +// Transient reports whether err looks like a transient OnlyOffice answer: a +// *TransientError with a retriable code, or an error whose text carries an +// HTTP 429/502/503/504. func Transient(err error) bool { if err == nil { return false } + var te *TransientError + if errors.As(err, &te) { + return isTransientStatus(te.Code) + } return transientRe.MatchString(err.Error()) } +// statusError wraps a non-2xx answer, tagging transient statuses so DoRetry +// recognises them and honours Retry-After. +func statusError(status int, retryAfter time.Duration, format string, args ...any) error { + msg := fmt.Sprintf(format, args...) + if isTransientStatus(status) { + return &TransientError{Code: status, RetryAfter: retryAfter, Msg: msg} + } + return errors.New(msg) +} + +// retryAfterOf parses the Retry-After header of a response (integer seconds or +// an HTTP-date). Returns 0 when absent or malformed. +func retryAfterOf(resp *http.Response) time.Duration { + if resp == nil { + return 0 + } + return parseRetryAfter(resp.Header.Get("Retry-After")) +} + +// parseRetryAfter parses a Retry-After value: delay-seconds (RFC 9110) or an +// HTTP-date. Zero, negative and malformed values yield 0. +func parseRetryAfter(v string) time.Duration { + v = strings.TrimSpace(v) + if v == "" { + return 0 + } + if secs, err := strconv.Atoi(v); err == nil { + if secs <= 0 { + return 0 + } + return time.Duration(secs) * time.Second + } + if t, err := http.ParseTime(v); err == nil { + if d := time.Until(t); d > 0 { + return d + } + } + return 0 +} + +// retryAfterOfError extracts Retry-After from a typed transient error. +func retryAfterOfError(err error) time.Duration { + var te *TransientError + if errors.As(err, &te) { + return te.RetryAfter + } + return 0 +} + +// backoffDelay returns the deterministic wait before retry `attempt` +// (counting from 1): Base*2^(attempt-1), capped at Max. +func backoffDelay(p RetryPolicy, attempt int) time.Duration { + if p.Base <= 0 { + return 0 + } + wait := p.Base + for i := 1; i < attempt; i++ { + if p.Max > 0 && wait >= p.Max { + return p.Max + } + wait *= 2 + } + if p.Max > 0 && wait > p.Max { + wait = p.Max + } + return wait +} + // DoRetry runs fn until it succeeds, fails non-transiently, or attempts run -// out. Waits are deterministic: N*Base capped at Max, no jitter. +// out. Waits are deterministic: Base*2^(N-1) capped at Max, no jitter. A +// transient error's Retry-After wins when it is longer than the backoff, and +// every wait arms the process-wide cooldown gate so concurrent and sequential +// callers back off too. func DoRetry(ctx context.Context, p RetryPolicy, fn func() error) error { if p.Attempts < 1 { p.Attempts = 1 @@ -44,13 +155,18 @@ func DoRetry(ctx context.Context, p RetryPolicy, fn func() error) error { if ctx.Err() != nil { return ctx.Err() } - if err = fn(); err == nil || !Transient(err) || attempt == p.Attempts { + err = fn() + if err == nil || !Transient(err) || attempt == p.Attempts { return err } - wait := time.Duration(attempt) * p.Base - if wait > p.Max { - wait = p.Max + wait := backoffDelay(p, attempt) + if ra := retryAfterOfError(err); ra > wait { + wait = ra } + if wait <= 0 { + continue + } + globalCooldown.note(wait) select { case <-ctx.Done(): return ctx.Err() diff --git a/retry_test.go b/retry_test.go index c774dda..a5a412f 100644 --- a/retry_test.go +++ b/retry_test.go @@ -3,6 +3,8 @@ package onlyoffice import ( "context" "errors" + "net/http" + "net/http/httptest" "testing" "time" ) @@ -20,6 +22,8 @@ func TestTransient(t *testing.T) { {errors.New("PUT /x: 504 gateway timeout"), true}, {errors.New("GET /x: 404 not found"), false}, {errors.New("GET /x: 4290 not a code"), false}, + {&TransientError{Code: 429, Msg: "throttled"}, true}, + {&TransientError{Code: 404, Msg: "missing"}, false}, } for _, c := range cases { if got := Transient(c.err); got != c.want { @@ -86,3 +90,131 @@ func TestDoRetryHonoursContextCancellation(t *testing.T) { t.Fatalf("want context.Canceled, got %v", err) } } + +func TestBackoffDelayExponentialAndCapped(t *testing.T) { + p := RetryPolicy{Base: 100 * time.Millisecond, Max: 250 * time.Millisecond} + want := []time.Duration{ + 100 * time.Millisecond, // 2^0 + 200 * time.Millisecond, // 2^1 + 250 * time.Millisecond, // 2^2 capped + 250 * time.Millisecond, + 250 * time.Millisecond, + } + for i, w := range want { + if got := backoffDelay(p, i+1); got != w { + t.Errorf("backoffDelay(attempt=%d) = %v, want %v", i+1, got, w) + } + } +} + +func TestDoRetryWaitsRetryAfterLongerThanBackoff(t *testing.T) { + const ra = 60 * time.Millisecond + p := RetryPolicy{Attempts: 2, Base: time.Millisecond, Max: time.Millisecond} + attempts := 0 + start := time.Now() + err := DoRetry(context.Background(), p, func() error { + attempts++ + if attempts == 1 { + return &TransientError{Code: http.StatusTooManyRequests, RetryAfter: ra, Msg: "GET /x: 429"} + } + return nil + }) + if err != nil { + t.Fatalf("DoRetry: %v", err) + } + if elapsed := time.Since(start); elapsed < ra { + t.Fatalf("waited %v, want >= %v", elapsed, ra) + } +} + +func TestParseRetryAfter(t *testing.T) { + cases := []struct { + in string + want time.Duration + }{ + {"", 0}, + {" 0 ", 0}, + {"-3", 0}, + {"bogus", 0}, + {"120", 120 * time.Second}, + } + for _, c := range cases { + if got := parseRetryAfter(c.in); got != c.want { + t.Errorf("parseRetryAfter(%q) = %v, want %v", c.in, got, c.want) + } + } + future := time.Now().Add(90 * time.Second).UTC().Format(http.TimeFormat) + if got := parseRetryAfter(future); got < 80*time.Second || got > 95*time.Second { + t.Errorf("HTTP-date Retry-After = %v, want ~90s", got) + } +} + +// TestGetJSON429CarriesRetryAfter checks that the HTTP layer surfaces a 429 as +// a typed transient error carrying the parsed Retry-After header. +func TestGetJSON429CarriesRetryAfter(t *testing.T) { + t.Setenv("OO_RETRY_ATTEMPTS", "1") + resetPacingForTest() + defer resetPacingForTest() + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Retry-After", "1") + w.WriteHeader(http.StatusTooManyRequests) + _, _ = w.Write([]byte("busy")) + })) + defer srv.Close() + + c := NewClient(Credentials{Url: srv.URL}) + c.token = &Token{Value: "tok", Expires: Time(time.Now().Add(time.Hour))} + + _, err := c.getJSON(context.Background(), "/api/2.0/project/4.json") + if err == nil { + t.Fatal("want error") + } + var te *TransientError + if !errors.As(err, &te) { + t.Fatalf("want *TransientError, got %T: %v", err, err) + } + if te.Code != http.StatusTooManyRequests { + t.Fatalf("code = %d, want 429", te.Code) + } + if te.RetryAfter != time.Second { + t.Fatalf("RetryAfter = %v, want 1s", te.RetryAfter) + } +} + +// TestGetJSONRetries429ThenSucceeds verifies the end-to-end path: a throttled +// answer with Retry-After is retried and the eventual success is returned. +func TestGetJSONRetries429ThenSucceeds(t *testing.T) { + t.Setenv("OO_RETRY_ATTEMPTS", "3") + t.Setenv("OO_RETRY_BASE", "10ms") + t.Setenv("OO_RETRY_MAX", "50ms") + resetPacingForTest() + defer resetPacingForTest() + + var calls int + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + if calls == 1 { + w.Header().Set("Retry-After", "0") + w.WriteHeader(http.StatusTooManyRequests) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"response":{}}`)) + })) + defer srv.Close() + + c := NewClient(Credentials{Url: srv.URL}) + c.token = &Token{Value: "tok", Expires: Time(time.Now().Add(time.Hour))} + + raw, err := c.getJSON(context.Background(), "/api/2.0/project/4.json") + if err != nil { + t.Fatalf("getJSON: %v", err) + } + if string(raw) != `{"response":{}}` { + t.Fatalf("raw = %s", raw) + } + if calls != 2 { + t.Fatalf("calls = %d, want 2", calls) + } +} diff --git a/storage_fallback.go b/storage_fallback.go index e29e178..be829af 100644 --- a/storage_fallback.go +++ b/storage_fallback.go @@ -150,7 +150,7 @@ func (c *Client) downloadFileEntry(ctx context.Context, f *FileEntry, dst io.Wri } return 0, fmt.Errorf("GET viewUrl: %d (stale S3) and minio fallback: %w", resp.StatusCode, merr) } - return 0, fmt.Errorf("GET viewUrl: %d %s", resp.StatusCode, truncate(string(b), 400)) + return 0, statusError(resp.StatusCode, retryAfterOf(resp), "GET viewUrl: %d %s", resp.StatusCode, truncate(string(b), 400)) } return io.Copy(dst, resp.Body) } -- 2.54.0