plugin/bundle: Correct bundle delay behavior (#7812)

* plugin/bundle: Correct bundle delay behavior

I ran into an issue when testing an earlier change:

https://github.com/open-policy-agent/opa/actions/runs/16646981900/job/47110035165

I found that this test generated around 100,000 lines of errors showing
the bundle downloader running.

This can be be tested using:

```
go test -v ./v1/plugins/bundle -count=1 2>&1 | grep -c "request failed"
```

This commit closes managers and plugins correctly.

Signed-off-by: Charlie Egan <charlie@styra.com>

* download: Update stop to be idempotent

I had some race detector issues with TestStartStopWithLongPollNotSupported
https://github.com/open-policy-agent/opa/actions/runs/16722869930/job/47334690407?pr=7812

I think this is a deadlock around multiple calls to Stop dead locking
updating the stopped var.

Signed-off-by: Charlie Egan <charlie@styra.com>

---------

Signed-off-by: Charlie Egan <charlie@styra.com>
This commit is contained in:
Charlie Egan
2025-08-07 17:13:37 +00:00
committed by GitHub
parent 16158ba8c9
commit f746d8caa9
4 changed files with 87 additions and 24 deletions
+11 -12
View File
@@ -61,8 +61,8 @@ type Downloader struct {
respHdrTimeoutSec int64
wg sync.WaitGroup
logger logging.Logger
mtx sync.Mutex
stopped bool
stopOnce sync.Once
persist bool
longPollingEnabled bool
lazyLoadingMode bool
@@ -204,16 +204,11 @@ func (d *Downloader) Stop(context.Context) {
return
}
d.mtx.Lock()
defer d.mtx.Unlock()
if d.stopped {
return
}
done := make(chan struct{})
d.stop <- done
<-done
d.stopOnce.Do(func() {
done := make(chan struct{})
d.stop <- done
<-done
})
}
func (d *Downloader) loop(ctx context.Context) {
@@ -231,13 +226,18 @@ func (d *Downloader) loop(ctx context.Context) {
return
}
// Note: MaxDelaySeconds, MinDelaySeconds and LongPollingTimeoutSeconds
// are scaled from int seconds to ns in ValidateAndInjectDefaults.
if err != nil {
// when there was an error, use a delay that's based on the retry count
delay = util.DefaultBackoff(float64(minRetryDelay), float64(*d.config.Polling.MaxDelaySeconds), retry)
} else if !d.longPollingEnabled || d.config.Polling.LongPollingTimeoutSeconds == nil {
// revert the response header timeout value on the http client's transport
if *d.client.Config().ResponseHeaderTimeoutSeconds == 0 {
d.client = d.client.SetResponseHeaderTimeout(&d.respHdrTimeoutSec)
}
// when polling, use a jittered delay based on min and max delay config
min := float64(*d.config.Polling.MinDelaySeconds)
max := float64(*d.config.Polling.MaxDelaySeconds)
delay = time.Duration(((max - min) * rand.Float64()) + min)
@@ -263,7 +263,6 @@ func (d *Downloader) loop(ctx context.Context) {
func (d *Downloader) oneShot(ctx context.Context) error {
m := metrics.New()
resp, err := d.download(ctx, m)
if err != nil {
d.etag = ""
+64 -3
View File
@@ -51,6 +51,7 @@ func TestPluginOneShot(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName, Metrics: metrics.New()}
@@ -125,6 +126,7 @@ func TestPluginOneShotWithAstStore(t *testing.T) {
ctx := context.Background()
store := inmem.NewWithOpts(inmem.OptRoundTripOnWrite(false), inmem.OptReturnASTValuesOnRead(true))
manager := getTestManagerWithOpts(nil, store)
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName, Metrics: metrics.New()}
@@ -225,6 +227,7 @@ corge contains 1 if {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName, Metrics: metrics.New()}
@@ -559,6 +562,7 @@ func TestPluginOneShotWithAuthzSchemaVerification(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
info, err := runtime.Term(runtime.Params{Config: nil, IsAuthorizationEnabled: true})
if err != nil {
@@ -704,6 +708,7 @@ func TestPluginOneShotWithAuthzSchemaVerificationNonDefaultAuthzPath(t *testing.
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
s := "/foo/authz/allow"
manager.Config.DefaultAuthorizationDecision = &s
@@ -843,6 +848,7 @@ func TestPluginStartLazyLoadInMem(t *testing.T) {
t.Fatal(err)
}
}))
defer s1.Close()
mockBundle2 := bundle.Bundle{
Data: map[string]any{"q": "x2"},
@@ -858,6 +864,7 @@ func TestPluginStartLazyLoadInMem(t *testing.T) {
t.Fatal(err)
}
}))
defer s2.Close()
config := fmt.Appendf(nil, `{
"services": {
@@ -1099,6 +1106,7 @@ func TestPluginOneShotDeltaBundle(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName, Metrics: metrics.New()}
@@ -1203,6 +1211,7 @@ func TestPluginOneShotDeltaBundleWithAstStore(t *testing.T) {
ctx := context.Background()
store := inmem.NewWithOpts(inmem.OptRoundTripOnWrite(false), inmem.OptReturnASTValuesOnRead(true))
manager := getTestManagerWithOpts(nil, store)
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName, Metrics: metrics.New()}
@@ -1306,6 +1315,7 @@ func TestPluginStart(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
bundles := map[string]*Source{}
plugin := New(&Config{Bundles: bundles}, manager)
@@ -1313,6 +1323,7 @@ func TestPluginStart(t *testing.T) {
if err != nil {
t.Fatal("unexpected error:", err)
}
defer plugin.Stop(ctx)
}
func TestStop(t *testing.T) {
@@ -1337,10 +1348,11 @@ func TestStop(t *testing.T) {
time.Sleep(time.Duration(longPollTimeout) * time.Second)
fmt.Fprintln(w) // Note: this is an invalid bundle and will fail the download
}))
defer ts.Close()
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
defer ts.Close()
serviceName := "test-svc"
err := manager.Reconfigure(&config.Config{
@@ -1391,6 +1403,7 @@ func TestPluginOneShotBundlePersistence(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -1996,6 +2009,7 @@ func TestPluginOneShotSignedBundlePersistence(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -2093,6 +2107,7 @@ func TestLoadAndActivateBundlesFromDisk(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -2180,6 +2195,7 @@ func TestLoadAndActivateBundlesFromDisk(t *testing.T) {
func TestLoadAndActivateBundlesFromDiskReservedChars(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -2779,6 +2795,7 @@ func TestLoadAndActivateDepBundlesFromDisk(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -2887,6 +2904,7 @@ func TestLoadAndActivateDepBundlesFromDiskMaxAttempts(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
dir := t.TempDir()
@@ -2956,6 +2974,7 @@ func TestPluginOneShotCompileError(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName}
@@ -3049,6 +3068,7 @@ func TestPluginOneShotHTTPError(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName}
@@ -3090,6 +3110,7 @@ func TestPluginOneShotActivationRemovesOld(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName}
@@ -3168,6 +3189,7 @@ func TestPluginOneShotActivationConflictingRoots(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
ensurePluginState(t, plugin, plugins.StateNotReady)
@@ -3237,6 +3259,7 @@ func TestPluginOneShotActivationPrefixMatchingRoots(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -3293,6 +3316,7 @@ func TestPluginListener(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := New(&Config{}, manager)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName}
@@ -3407,6 +3431,7 @@ func TestPluginListenerErrorClearedOn304(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -3461,6 +3486,7 @@ func TestPluginBulkListener(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -3664,6 +3690,7 @@ func TestPluginBulkListenerStatusCopyOnly(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -3887,6 +3914,7 @@ func TestPluginSetCompilerOnContext(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -3971,10 +3999,11 @@ func TestPluginReconfigure(t *testing.T) {
}
fmt.Fprintln(w, "") // Note: this is an invalid bundle and will fail the download
}))
defer ts.Close()
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
defer ts.Close()
serviceName := "test-svc"
err := manager.Reconfigure(&config.Config{
@@ -3985,6 +4014,7 @@ func TestPluginReconfigure(t *testing.T) {
}
plugin := New(&Config{}, manager)
defer plugin.Stop(ctx)
var delay int64 = 10
@@ -4165,6 +4195,7 @@ func TestPluginRequestVsDownloadTimestamp(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -4537,6 +4568,7 @@ func TestUpgradeLegacyBundleToMuiltiBundleSameBundle(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
status: map[string]*Status{},
@@ -4634,6 +4666,7 @@ func TestUpgradeLegacyBundleToMultiBundleNewBundles(t *testing.T) {
ctx := context.Background()
manager := getTestManager()
defer manager.Stop(ctx)
plugin := Plugin{
manager: manager,
@@ -4641,6 +4674,7 @@ func TestUpgradeLegacyBundleToMultiBundleNewBundles(t *testing.T) {
etags: map[string]string{},
downloaders: map[string]Loader{},
}
defer plugin.Stop(ctx)
bundleName := "test-bundle"
plugin.status[bundleName] = &Status{Name: bundleName}
@@ -4881,6 +4915,7 @@ func TestSaveBundleToDiskNew(t *testing.T) {
t.Parallel()
manager := getTestManager()
defer manager.Stop(context.Background())
dir := t.TempDir()
@@ -4900,6 +4935,7 @@ func TestSaveBundleToDiskNewConfiguredPersistDir(t *testing.T) {
dir := t.TempDir()
manager := getTestManager()
defer manager.Stop(context.Background())
manager.Config.PersistenceDirectory = &dir
bundles := map[string]*Source{}
plugin := New(&Config{Bundles: bundles}, manager)
@@ -4908,6 +4944,7 @@ func TestSaveBundleToDiskNewConfiguredPersistDir(t *testing.T) {
if err != nil {
t.Fatalf("unexpected error %v", err)
}
defer plugin.Stop(context.Background())
err = plugin.saveBundleToDisk("foo", getTestRawBundle(t))
if err != nil {
@@ -4925,6 +4962,7 @@ func TestSaveBundleToDiskOverWrite(t *testing.T) {
t.Parallel()
manager := getTestManager()
defer manager.Stop(context.Background())
// test to check existing bundle is replaced
dir := t.TempDir()
@@ -5011,6 +5049,7 @@ func TestLoadBundleFromDisk(t *testing.T) {
t.Parallel()
manager := getTestManager()
defer manager.Stop(context.Background())
plugin := New(&Config{}, manager)
// no bundle on disk
@@ -5108,6 +5147,7 @@ func TestLoadSignedBundleFromDisk(t *testing.T) {
t.Parallel()
manager := getTestManager()
defer manager.Stop(context.Background())
plugin := New(&Config{}, manager)
// no bundle on disk
@@ -5150,7 +5190,9 @@ func TestLoadSignedBundleFromDisk(t *testing.T) {
func TestGetDefaultBundlePersistPath(t *testing.T) {
t.Parallel()
plugin := New(&Config{}, getTestManager())
manager := getTestManager()
defer manager.Stop(context.Background())
plugin := New(&Config{}, manager)
path, err := plugin.getBundlePersistPath()
if err != nil {
t.Fatalf("unexpected error %v", err)
@@ -5166,6 +5208,7 @@ func TestConfiguredBundlePersistPath(t *testing.T) {
persistPath := "/var/opa"
manager := getTestManager()
defer manager.Stop(context.Background())
manager.Config.PersistenceDirectory = &persistPath
plugin := New(&Config{}, manager)
@@ -6261,6 +6304,7 @@ func TestPluginReadBundleEtagFromDiskStore(t *testing.T) {
t.Fatal(err)
}
}))
defer s.Close()
test.WithTempFS(nil, func(dir string) {
ctx := context.Background()
@@ -6458,6 +6502,7 @@ func TestPluginStateReconciliationOnReconfigure(t *testing.T) {
t.Fatal(err)
}
}))
defer s.Close()
// setup plugin pointing at fake server
manager := getTestManagerWithOpts(fmt.Appendf(nil, `{
@@ -6496,6 +6541,7 @@ func TestPluginStateReconciliationOnReconfigure(t *testing.T) {
if err != nil {
t.Fatal(err)
}
defer plugin.Stop(ctx)
// manually trigger bundle download
go func() { _ = plugin.Loaders()["b1"].Trigger(ctx) }()
@@ -6628,6 +6674,7 @@ func TestPluginManualTrigger(t *testing.T) {
t.Fatal(err)
}
}))
defer s.Close()
// setup plugin pointing at fake server
manager := getTestManagerWithOpts(fmt.Appendf(nil, `{
@@ -6637,6 +6684,7 @@ func TestPluginManualTrigger(t *testing.T) {
}
}
}`, s.URL))
defer manager.Stop(ctx)
var mode plugins.TriggerMode = "manual"
@@ -6661,6 +6709,7 @@ func TestPluginManualTrigger(t *testing.T) {
if err != nil {
t.Fatal(err)
}
defer plugin.Stop(ctx)
// manually trigger bundle download
go func() {
@@ -6729,6 +6778,7 @@ func TestPluginManualTriggerMultipleDiskStorage(t *testing.T) {
t.Fatal(err)
}
}))
defer s1.Close()
mockBundle2 := bundle.Bundle{
Data: map[string]any{"q": "x2"},
@@ -6744,6 +6794,7 @@ func TestPluginManualTriggerMultipleDiskStorage(t *testing.T) {
t.Fatal(err)
}
}))
defer s2.Close()
test.WithTempFS(nil, func(dir string) {
store, err := disk.New(ctx, logging.NewNoOpLogger(), nil, disk.Options{
@@ -6782,6 +6833,7 @@ func TestPluginManualTriggerMultipleDiskStorage(t *testing.T) {
},
},
}, manager)
defer plugin.Stop(ctx)
statusCh := make(chan map[string]*Status)
@@ -6875,6 +6927,7 @@ func TestPluginManualTriggerMultiple(t *testing.T) {
t.Fatal(err)
}
}))
defer s1.Close()
mockBundle2 := bundle.Bundle{
Data: map[string]any{"q": "x2"},
@@ -6890,6 +6943,7 @@ func TestPluginManualTriggerMultiple(t *testing.T) {
t.Fatal(err)
}
}))
defer s2.Close()
// setup plugin pointing at fake server
manager := getTestManagerWithOpts(fmt.Appendf(nil, `{
@@ -6902,6 +6956,7 @@ func TestPluginManualTriggerMultiple(t *testing.T) {
}
}
}`, s1.URL, s2.URL))
defer manager.Stop(ctx)
var mode plugins.TriggerMode = "manual"
@@ -6931,6 +6986,7 @@ func TestPluginManualTriggerMultiple(t *testing.T) {
if err != nil {
t.Fatal(err)
}
defer plugin.Stop(ctx)
// manually trigger bundle download on all configured bundles
go func() {
@@ -6969,6 +7025,7 @@ func TestPluginManualTriggerWithTimeout(t *testing.T) {
s := httptest.NewServer(http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {
time.Sleep(3 * time.Second) // this should cause the context deadline to exceed
}))
defer s.Close()
// setup plugin pointing at fake server
manager := getTestManagerWithOpts(fmt.Appendf(nil, `{
@@ -6978,6 +7035,7 @@ func TestPluginManualTriggerWithTimeout(t *testing.T) {
}
}
}`, s.URL))
defer manager.Stop(ctx)
var mode plugins.TriggerMode = "manual"
@@ -7003,6 +7061,7 @@ func TestPluginManualTriggerWithTimeout(t *testing.T) {
if err != nil {
t.Fatal(err)
}
defer plugin.Stop(ctx)
// manually trigger bundle download
go func() {
@@ -7030,6 +7089,7 @@ func TestPluginManualTriggerWithServerError(t *testing.T) {
s := httptest.NewServer(http.HandlerFunc(func(resp http.ResponseWriter, _ *http.Request) {
resp.WriteHeader(500)
}))
defer s.Close()
// setup plugin pointing at fake server
manager := getTestManagerWithOpts(fmt.Appendf(nil, `{
@@ -7039,6 +7099,7 @@ func TestPluginManualTriggerWithServerError(t *testing.T) {
}
}
}`, s.URL))
defer manager.Stop(ctx)
var manual plugins.TriggerMode = "manual"
+10 -9
View File
@@ -4,6 +4,7 @@ import (
"bytes"
"strings"
"testing"
"time"
"github.com/open-policy-agent/opa/v1/logging"
test_sdk "github.com/open-policy-agent/opa/v1/sdk/test"
@@ -29,7 +30,7 @@ func TestEnablePrintStatementsForFilesystemPolicies(t *testing.T) {
params := e2e.NewAPIServerTestParams()
params.Paths = []string{dir}
buf := bytes.NewBuffer(nil)
buf := &test.BlockingWriter{}
logger := logging.New()
logger.SetOutput(buf)
@@ -69,7 +70,7 @@ func TestEnablePrintStatementsForHTTPAPIPushedPolicies(t *testing.T) {
params := e2e.NewAPIServerTestParams()
buf := bytes.NewBuffer(nil)
buf := &test.BlockingWriter{}
logger := logging.New()
logger.SetOutput(buf)
@@ -93,9 +94,9 @@ func TestEnablePrintStatementsForHTTPAPIPushedPolicies(t *testing.T) {
expContains := "hello world"
if !strings.Contains(buf.String(), expContains) {
t.Fatalf("expected logs to contain %q but got: %v", expContains, buf.String())
}
test.EventuallyOrFatal(t, 100*time.Millisecond, func() bool {
return strings.Contains(buf.String(), expContains)
})
})
}
@@ -117,7 +118,7 @@ func TestEnablePrintStatementsForBundles(t *testing.T) {
params := e2e.NewAPIServerTestParams()
buf := bytes.NewBuffer(nil)
buf := &test.BlockingWriter{}
logger := logging.New()
logger.SetOutput(buf)
@@ -142,8 +143,8 @@ func TestEnablePrintStatementsForBundles(t *testing.T) {
expContains := "hello world"
if !strings.Contains(buf.String(), expContains) {
t.Fatalf("expected logs to contain %q but got: %v", expContains, buf.String())
}
test.EventuallyOrFatal(t, 100*time.Millisecond, func() bool {
return strings.Contains(buf.String(), expContains)
})
})
}
+2
View File
@@ -17,6 +17,8 @@ func DefaultBackoff(base, maxNS float64, retries int) time.Duration {
// Backoff returns a delay with an exponential backoff based on the number of
// retries. Same algorithm used in gRPC.
// Note that if maxNS is smaller than base, the backoff will still be capped at
// maxNS.
func Backoff(base, maxNS, jitter, factor float64, retries int) time.Duration {
if retries == 0 {
return 0