From d6b8eb7777aeed945bd11dfb26eecef741fd152e Mon Sep 17 00:00:00 2001 From: Andriy Oblivantsev Date: Wed, 16 Sep 2026 21:32:14 +0000 Subject: [PATCH] fix(retry): retry transient edge answers centrally, fix stale task test (#57) - Route HTTP helpers (getJSON/formRequest/deleteReq/postJSON/putJSON/ multipart upload), Query and AuthenticateContext/ensureToken through the deterministic 429/502/503/504 retry (retryRaw), so the integration suite no longer fails on the shared openresty rate limit under parallel runs. - Fix stale cmd/office/fetch integration test: loader.TaskFields was replaced by loader.DetailForm (broke go vet -tags=integration). - Add unit tests for Transient/DoRetry. --- AGENTS.md | 2 +- auth.go | 23 ++++-- cmd/office/fetch/task_integration_test.go | 10 +-- http.go | 34 ++++++++- request.go | 16 ++++- retry.go | 14 ++++ retry_test.go | 88 +++++++++++++++++++++++ 7 files changed, 169 insertions(+), 18 deletions(-) create mode 100644 retry_test.go diff --git a/AGENTS.md b/AGENTS.md index 3c0a26b..b37c558 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. **`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 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). - **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/auth.go b/auth.go index 98c0986..3d2a67a 100644 --- a/auth.go +++ b/auth.go @@ -43,6 +43,11 @@ func (c *Client) Authenticate() error { return c.ensureToken() } // cached token is still valid it returns immediately; otherwise it performs // a POST to /api/2.0/authentication.json that is cancellable via ctx. // +// Transient answers from the edge (openresty 429/502/503/504) are retried with +// the same deterministic policy as every other request (see retry.go), because +// the server rate-limits authentication and the integration suite otherwise +// fails with a raw HTML 429 page. +// // This is the recommended entry point for long-running syncs (cron, // watchers) because it guarantees that a stalled auth call will not block // the caller past its deadline. @@ -50,6 +55,14 @@ func (c *Client) AuthenticateContext(ctx context.Context) error { if c.tokenValid() { return nil } + return DoRetry(ctx, DefaultRetryPolicy(), func() error { + return c.authenticateOnce(ctx) + }) +} + +// authenticateOnce performs a single authentication POST. Callers must handle +// retries; use AuthenticateContext. +func (c *Client) authenticateOnce(ctx context.Context) error { body, err := json.Marshal(c.credentials) if err != nil { return fmt.Errorf("marshal credentials: %w", err) @@ -99,17 +112,13 @@ func (c *Client) tokenValid() bool { // ensureToken refreshes the authentication token when missing or expired. // Mirrors the logic inline in Query() but is safe to call from helpers that -// bypass the typed Request abstraction. +// bypass the typed Request abstraction. It shares AuthenticateContext so the +// transient-retry policy applies to every code path. func (c *Client) ensureToken() error { if c.tokenValid() { return nil } - tok, err := c.Auth(c.credentials) - if err != nil { - return err - } - c.token = tok - return nil + return c.AuthenticateContext(context.Background()) } // authHeader returns the value for the Authorization header, ensuring a token. diff --git a/cmd/office/fetch/task_integration_test.go b/cmd/office/fetch/task_integration_test.go index 94a60d0..00d6547 100644 --- a/cmd/office/fetch/task_integration_test.go +++ b/cmd/office/fetch/task_integration_test.go @@ -4,7 +4,6 @@ package fetch_test import ( "context" - "os" "testing" "github.com/eslider/go-onlyoffice/cmd/office/model" @@ -46,9 +45,6 @@ func TestIntegrationUpdateTaskTitleDescription(t *testing.T) { } func TestIntegrationTaskFieldsFromLiveAPI(t *testing.T) { - if os.Getenv("ONLYOFFICE_URL") == "" && os.Getenv("ONLYOFFICE_HOST") == "" { - t.Skip("ONLYOFFICE_URL not set") - } loader, ctx := liveLoader(t) items, err := loader.List(ctx, model.ListSpec{Subject: model.SubjectTasks}) if err != nil { @@ -57,12 +53,12 @@ func TestIntegrationTaskFieldsFromLiveAPI(t *testing.T) { if len(items) == 0 { t.Skip("no tasks") } - title, desc, err := loader.TaskFields(ctx, items[0]) + fields, err := loader.DetailForm(ctx, items[0]) if err != nil { t.Fatal(err) } - if title == "" { + if fields.Primary == "" { t.Fatal("empty title") } - _ = desc + _ = fields.Secondary } diff --git a/http.go b/http.go index 5cf8787..1452d3d 100644 --- a/http.go +++ b/http.go @@ -144,8 +144,13 @@ func unmarshalResponseObject(raw json.RawMessage) (map[string]any, error) { } } -// getJSON issues an authenticated GET and returns the raw response body. +// getJSON issues an authenticated GET and returns the raw response body, +// retrying transient answers (see retryRaw). func (c *Client) getJSON(ctx context.Context, path string) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { return c.getJSONOnce(ctx, path) }) +} + +func (c *Client) getJSONOnce(ctx context.Context, path string) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err @@ -187,6 +192,12 @@ func (c *Client) deleteForm(ctx context.Context, path string, fields url.Values) } func (c *Client) formRequest(ctx context.Context, method, path string, fields url.Values) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { + return c.formRequestOnce(ctx, method, path, fields) + }) +} + +func (c *Client) formRequestOnce(ctx context.Context, method, path string, fields url.Values) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err @@ -215,6 +226,10 @@ func (c *Client) formRequest(ctx context.Context, method, path string, fields ur // deleteReq issues an authenticated DELETE. func (c *Client) deleteReq(ctx context.Context, path string) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { return c.deleteReqOnce(ctx, path) }) +} + +func (c *Client) deleteReqOnce(ctx context.Context, path string) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err @@ -260,6 +275,10 @@ func (c *Client) postJSONObject(ctx context.Context, path string, body any) (map // postJSON issues an authenticated POST with application/json body. func (c *Client) postJSON(ctx context.Context, path string, body any) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { return c.postJSONOnce(ctx, path, body) }) +} + +func (c *Client) postJSONOnce(ctx context.Context, path string, body any) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err @@ -303,6 +322,10 @@ func (c *Client) postJSON(ctx context.Context, path string, body any) (json.RawM // putJSON issues an authenticated PUT with application/json body. func (c *Client) putJSON(ctx context.Context, path string, body any) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { return c.putJSONOnce(ctx, path, body) }) +} + +func (c *Client) putJSONOnce(ctx context.Context, path string, body any) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err @@ -352,8 +375,15 @@ func (c *Client) uploadMultipart(ctx context.Context, path, fieldName, filePath // uploadMultipartMethod sends a single-file multipart request with the given // HTTP method. The OnlyOffice Documents API needs PUT for /update (a new // version) and POST for /upload (a new file); sending POST to /update answers -// 500 on current servers. +// 500 on current servers. The file is re-opened per attempt, so transient +// answers are retried like every other request. func (c *Client) uploadMultipartMethod(ctx context.Context, method, path, fieldName, filePath string) (json.RawMessage, error) { + return retryRaw(ctx, func() (json.RawMessage, error) { + return c.uploadMultipartOnce(ctx, method, path, fieldName, filePath) + }) +} + +func (c *Client) uploadMultipartOnce(ctx context.Context, method, path, fieldName, filePath string) (json.RawMessage, error) { auth, err := c.authHeader() if err != nil { return nil, err diff --git a/request.go b/request.go index a4bb443..74adb30 100644 --- a/request.go +++ b/request.go @@ -8,6 +8,7 @@ package onlyoffice import ( "bytes" + "context" "encoding/json" "fmt" "io" @@ -49,6 +50,12 @@ func (r Request) GetMethod() string { // - If request.NoAuth is true then no token is fetched — the caller is // responsible for authenticating requests (used internally by Auth()). func (c *Client) Query(request Request, result interface{}) error { + return DoRetry(context.Background(), DefaultRetryPolicy(), func() error { + return c.queryOnce(request, result) + }) +} + +func (c *Client) queryOnce(request Request, result interface{}) error { url := c.credentials.Url + request.Uri if request.Params != nil { @@ -90,10 +97,17 @@ func (c *Client) Query(request Request, result interface{}) error { } defer resp.Body.Close() + raw, err := io.ReadAll(resp.Body) + if err != nil { + return err + } + if resp.StatusCode >= 400 { + return fmt.Errorf("%s %s: %d %s", request.GetMethod(), request.Uri, resp.StatusCode, truncate(string(raw), 400)) + } if result == nil { return nil } - return json.NewDecoder(resp.Body).Decode(result) + return json.Unmarshal(raw, result) } // requestBodyReader normalises Query() body input into an io.Reader. diff --git a/retry.go b/retry.go index 73210e9..bfa1555 100644 --- a/retry.go +++ b/retry.go @@ -2,6 +2,7 @@ package onlyoffice import ( "context" + "encoding/json" "regexp" "time" ) @@ -58,3 +59,16 @@ func DoRetry(ctx context.Context, p RetryPolicy, fn func() error) error { } return err } + +// retryRaw runs a transport attempt under the default transient-retry policy +// and returns its raw payload. All HTTP helpers and Query() go through it, so +// an openresty 429/502/503/504 is retried exactly like every bulk tool. +func retryRaw(ctx context.Context, fn func() (json.RawMessage, error)) (json.RawMessage, error) { + var raw json.RawMessage + err := DoRetry(ctx, DefaultRetryPolicy(), func() error { + var e error + raw, e = fn() + return e + }) + return raw, err +} diff --git a/retry_test.go b/retry_test.go new file mode 100644 index 0000000..c774dda --- /dev/null +++ b/retry_test.go @@ -0,0 +1,88 @@ +package onlyoffice + +import ( + "context" + "errors" + "testing" + "time" +) + +func TestTransient(t *testing.T) { + cases := []struct { + err error + want bool + }{ + {nil, false}, + {errors.New("boom"), false}, + {errors.New("GET /api/2.0/files/1.json: 429 Too Many Requests"), true}, + {errors.New("auth: 502 bad gateway"), true}, + {errors.New("POST /x: 503 service unavailable"), true}, + {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}, + } + for _, c := range cases { + if got := Transient(c.err); got != c.want { + t.Errorf("Transient(%v) = %v, want %v", c.err, got, c.want) + } + } +} + +func TestDoRetryRetriesTransientThenSucceeds(t *testing.T) { + p := RetryPolicy{Attempts: 5, Base: time.Millisecond, Max: 5 * time.Millisecond} + attempts := 0 + err := DoRetry(context.Background(), p, func() error { + attempts++ + if attempts < 3 { + return errors.New("GET /x: 429 Too Many Requests") + } + return nil + }) + if err != nil { + t.Fatalf("DoRetry: %v", err) + } + if attempts != 3 { + t.Fatalf("attempts = %d, want 3", attempts) + } +} + +func TestDoRetryStopsOnPermanentError(t *testing.T) { + p := RetryPolicy{Attempts: 5, Base: time.Millisecond, Max: time.Millisecond} + attempts := 0 + err := DoRetry(context.Background(), p, func() error { + attempts++ + return errors.New("GET /x: 404 not found") + }) + if err == nil { + t.Fatal("want error") + } + if attempts != 1 { + t.Fatalf("attempts = %d, want 1", attempts) + } +} + +func TestDoRetryExhaustsAttemptsOnPersistentTransient(t *testing.T) { + p := RetryPolicy{Attempts: 3, Base: time.Millisecond, Max: time.Millisecond} + attempts := 0 + err := DoRetry(context.Background(), p, func() error { + attempts++ + return errors.New("GET /x: 429 Too Many Requests") + }) + if !Transient(err) { + t.Fatalf("want transient error, got %v", err) + } + if attempts != 3 { + t.Fatalf("attempts = %d, want 3", attempts) + } +} + +func TestDoRetryHonoursContextCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + err := DoRetry(ctx, RetryPolicy{Attempts: 3, Base: time.Millisecond, Max: time.Millisecond}, func() error { + return errors.New("GET /x: 429 Too Many Requests") + }) + if !errors.Is(err, context.Canceled) { + t.Fatalf("want context.Canceled, got %v", err) + } +} -- 2.54.0