fix(retry): retry transient edge answers centrally, fix stale task test (#57)
Release Please / Release Please (push) Skipped
Release / GoReleaser (push) Skipped
Tests / Secret scan (gitleaks) (push) Skipped
Tests / Test (Go 1.25) (push) Skipped
Tests / Test (Go stable) (push) Skipped
Tests / Secret scan (gitleaks) (pull_request) Successful in 4s
Tests / Test (Go 1.25) (pull_request) Successful in 21s
Tests / Test (Go stable) (pull_request) Successful in 21s
Release Please / Release Please (push) Skipped
Release / GoReleaser (push) Skipped
Tests / Secret scan (gitleaks) (push) Skipped
Tests / Test (Go 1.25) (push) Skipped
Tests / Test (Go stable) (push) Skipped
Tests / Secret scan (gitleaks) (pull_request) Successful in 4s
Tests / Test (Go 1.25) (pull_request) Successful in 21s
Tests / Test (Go stable) (pull_request) Successful in 21s
- 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.
This commit is contained in:
@@ -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):
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
+15
-1
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user