fix(retry): глобальный rate-limit + Retry-After + cooldown (#70) #71

Merged
eSlider merged 1 commits from fix/oo-backoff#70 into main 2026-09-17 13:39:57 +01:00
16 changed files with 729 additions and 33 deletions
Showing only changes of commit c9e16c7169 - Show all commits
+1 -1
View File
@@ -9,7 +9,7 @@ Canonical Go client for OnlyOffice Workspace (Projects + Calendar + CRM) and the
- `request.go` — `Request`, `Query`, `Time`, `Token`, `MetaResponse`, `Permissions`. - `request.go` — `Request`, `Query`, `Time`, `Token`, `MetaResponse`, `Permissions`.
- `auth.go` — `Authenticate`, `AuthenticateContext`, `InvalidateToken`, `Auth`, token lifecycle. - `auth.go` — `Authenticate`, `AuthenticateContext`, `InvalidateToken`, `Auth`, token lifecycle.
- `http.go` — transport + DRY response decoders (`ResponseArray`/`ResponseObject`/`postFormObject`/`putFormObject`/`deleteObject`). - `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). - **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. - 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): - **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):
+8
View File
@@ -6,6 +6,14 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
## Unreleased ## 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) ## [0.18.0](https://github.com/eSlider/go-onlyoffice/compare/v0.17.0...v0.18.0) (2026-09-04)
+14 -6
View File
@@ -522,9 +522,9 @@ type Task struct {
| `ListFileOps(ctx)` | Active file operations (move/copy status polling) | | `ListFileOps(ctx)` | Active file operations (move/copy status polling) |
| `FolderFiles(ctx, folderID)` | Flat file list of a folder (stem helpers) | | `FolderFiles(ctx, folderID)` | Flat file list of a folder (stem helpers) |
| `DeleteFilesByStem(ctx, folderID, stem)` | Remove `stem\|ext` copies | | `DeleteFilesByStem(ctx, folderID, stem)` | Remove `stem\|ext` copies |
| `DoRetry(ctx, policy, fn)` | Deterministic linear backoff (N·Base, no jitter) on 429/502/503/504 | | `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()` | 5 attempts, 1s·2s·3s·4s waits, 30s cap | | `DefaultRetryPolicy()` | From env: 7 attempts, 2s base, 2m cap (`OO_RETRY_ATTEMPTS/_BASE/_MAX`) |
| `Transient(err)` | True for retriable OnlyOffice answers | | `Transient(err)` | True for retriable OnlyOffice answers (`*TransientError` or HTTP 429/502/503/504 text) |
### Helper Types ### Helper Types
@@ -732,9 +732,12 @@ fallback rules, env names and how to add a backend:
### Bulk tools (`cmd/`) ### Bulk tools (`cmd/`)
Small single-purpose binaries for bulk Documents work. All of them pace Small single-purpose binaries for bulk Documents work. All of them pace
requests and retry transient OnlyOffice answers (429/502/503/504) with a requests through a process-wide token bucket (default ~4 req/s, `OO_RATE_LIMIT`/
deterministic linear backoff — no jitter, same waits on every run `OO_BURST`) and retry transient OnlyOffice answers (429/502/503/504) with a
(see `DoRetry` below). Build with `go build ./cmd/<tool>`. 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/<tool>`.
```bash ```bash
ooscan 659 # recursive index → TSV: file_id, folder_id, path, title 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_CALENDAR_ID` | Default calendar id used when omitted (default `1`) |
| `ONLYOFFICE_PROJECT_ID` | Default project id used when omitted (default `33`) | | `ONLYOFFICE_PROJECT_ID` | Default project id used when omitted (default `33`) |
| `OO_URL`, `OO_USER`, `OO_PASS` | Optional CLI-only aliases for `ONLYOFFICE_*` | | `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_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_INDEX` | OnlyOffice index (default `files_file`) |
| `ONLYOFFICE_ES_TEXT_INDEX` | Own PDF/scan index (default `oo_docs_text`) | | `ONLYOFFICE_ES_TEXT_INDEX` | Own PDF/scan index (default `oo_docs_text`) |
+1 -1
View File
@@ -83,7 +83,7 @@ func (c *Client) authenticateOnce(ctx context.Context) error {
return err return err
} }
if resp.StatusCode >= 400 { 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 { var env struct {
Response *Token `json:"response"` Response *Token `json:"response"`
+6 -2
View File
@@ -38,11 +38,15 @@ type Client struct {
folderTitlesMu sync.Mutex 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 { func NewClient(c Credentials) *Client {
jar, _ := cookiejar.New(nil) jar, _ := cookiejar.New(nil)
return &Client{ return &Client{
client: &http.Client{Jar: jar}, client: &http.Client{
Jar: jar,
Transport: &pacedTransport{base: http.DefaultTransport},
},
credentials: &c, credentials: &c,
} }
} }
+3
View File
@@ -22,6 +22,9 @@ related:
- [rclone-webdav.md](rclone-webdav.md) — rclone-монтирование Documents - [rclone-webdav.md](rclone-webdav.md) — rclone-монтирование Documents
(`deploy/docker-compose.rclone-webdav.yml`), smoke, ограничения. (`deploy/docker-compose.rclone-webdav.yml`), smoke, ограничения.
- [crm-associations.md](crm-associations.md) — правила ассоциаций CRM. - [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). Команды и туннели — раздел Testing в [README.md](../README.md#testing).
+55
View File
@@ -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 с заголовком.
+2 -2
View File
@@ -325,7 +325,7 @@ func (c *Client) deleteJSON(ctx context.Context, path string, body any) (json.Ra
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -365,7 +365,7 @@ func (c *Client) uploadReader(ctx context.Context, path, fieldName, fileName str
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
+6 -6
View File
@@ -171,7 +171,7 @@ func (c *Client) getJSONOnce(ctx context.Context, path string) (json.RawMessage,
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -219,7 +219,7 @@ func (c *Client) formRequestOnce(ctx context.Context, method, path string, field
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -250,7 +250,7 @@ func (c *Client) deleteReqOnce(ctx context.Context, path string) (json.RawMessag
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -315,7 +315,7 @@ func (c *Client) postJSONOnce(ctx context.Context, path string, body any) (json.
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -362,7 +362,7 @@ func (c *Client) putJSONOnce(ctx context.Context, path string, body any) (json.R
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
@@ -423,7 +423,7 @@ func (c *Client) uploadMultipartOnce(ctx context.Context, method, path, fieldNam
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
+1 -1
View File
@@ -125,7 +125,7 @@ func (c *Client) DownloadMailAttachment(ctx context.Context, attachmentID string
return nil, err return nil, err
} }
if resp.StatusCode >= 400 { 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 return raw, nil
} }
+210
View File
@@ -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
}
+160
View File
@@ -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)
}
}
+1 -1
View File
@@ -102,7 +102,7 @@ func (c *Client) queryOnce(request Request, result interface{}) error {
return err return err
} }
if resp.StatusCode >= 400 { 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 { if result == nil {
return nil return nil
+128 -12
View File
@@ -3,38 +3,149 @@ package onlyoffice
import ( import (
"context" "context"
"encoding/json" "encoding/json"
"errors"
"fmt"
"net/http"
"regexp" "regexp"
"strconv"
"strings"
"time" "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. // backoff without jitter, so repeated runs wait exactly the same schedule.
// OnlyOffice throttles bulk reads/writes with 429 (and occasional 502/503/504 // OnlyOffice throttles bulk reads/writes with 429 (and occasional 502/503/504
// from openresty), so every bulk tool routes API calls through DoRetry. // from openresty), so every bulk tool routes API calls through DoRetry.
type RetryPolicy struct { type RetryPolicy struct {
Attempts int // total attempts, including the first try Attempts int // total attempts, including the first try
Base time.Duration // wait before retry N is N*Base Base time.Duration // wait before retry N is Base*2^(N-1)
Max time.Duration // per-wait cap 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 { 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`) var transientRe = regexp.MustCompile(`:\s*(429|502|503|504)\b`)
// Transient reports whether err looks like a transient OnlyOffice answer // isTransientStatus reports whether an HTTP status is retriable at the edge.
// (an HTTP 429/502/503/504 surfaced as "...: <code> ..."). 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 { func Transient(err error) bool {
if err == nil { if err == nil {
return false return false
} }
var te *TransientError
if errors.As(err, &te) {
return isTransientStatus(te.Code)
}
return transientRe.MatchString(err.Error()) 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 // 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 { func DoRetry(ctx context.Context, p RetryPolicy, fn func() error) error {
if p.Attempts < 1 { if p.Attempts < 1 {
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 { if ctx.Err() != nil {
return ctx.Err() 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 return err
} }
wait := time.Duration(attempt) * p.Base wait := backoffDelay(p, attempt)
if wait > p.Max { if ra := retryAfterOfError(err); ra > wait {
wait = p.Max wait = ra
} }
if wait <= 0 {
continue
}
globalCooldown.note(wait)
select { select {
case <-ctx.Done(): case <-ctx.Done():
return ctx.Err() return ctx.Err()
+132
View File
@@ -3,6 +3,8 @@ package onlyoffice
import ( import (
"context" "context"
"errors" "errors"
"net/http"
"net/http/httptest"
"testing" "testing"
"time" "time"
) )
@@ -20,6 +22,8 @@ func TestTransient(t *testing.T) {
{errors.New("PUT /x: 504 gateway timeout"), true}, {errors.New("PUT /x: 504 gateway timeout"), true},
{errors.New("GET /x: 404 not found"), false}, {errors.New("GET /x: 404 not found"), false},
{errors.New("GET /x: 4290 not a code"), 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 { for _, c := range cases {
if got := Transient(c.err); got != c.want { 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) 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)
}
}
+1 -1
View File
@@ -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 (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) return io.Copy(dst, resp.Body)
} }