mirror of
https://github.com/open-policy-agent/opa.git
synced 2026-08-12 19:32:48 -06:00
plugins: Surface decision log errors via status API
Currently errors encountered during decision log uploads are logged. This change now includes these errors as part of status updates which can be useful for the control plane to determine if OPA is having any issues while processing, uploading logs etc. Fixes: #5637 Signed-off-by: Ashutosh Narkar <anarkar4387@gmail.com>
This commit is contained in:
+48
-27
@@ -12,6 +12,7 @@ import (
|
||||
"net/http"
|
||||
"reflect"
|
||||
|
||||
lstat "github.com/open-policy-agent/opa/plugins/logs/status"
|
||||
prom "github.com/prometheus/client_golang/prometheus"
|
||||
|
||||
"github.com/open-policy-agent/opa/logging"
|
||||
@@ -31,32 +32,35 @@ type Logger interface {
|
||||
// UpdateRequestV1 represents the status update message that OPA sends to
|
||||
// remote HTTP endpoints.
|
||||
type UpdateRequestV1 struct {
|
||||
Labels map[string]string `json:"labels"`
|
||||
Bundle *bundle.Status `json:"bundle,omitempty"` // Deprecated: Use bulk `bundles` status updates instead
|
||||
Bundles map[string]*bundle.Status `json:"bundles,omitempty"`
|
||||
Discovery *bundle.Status `json:"discovery,omitempty"`
|
||||
Metrics map[string]interface{} `json:"metrics,omitempty"`
|
||||
Plugins map[string]*plugins.Status `json:"plugins,omitempty"`
|
||||
Labels map[string]string `json:"labels"`
|
||||
Bundle *bundle.Status `json:"bundle,omitempty"` // Deprecated: Use bulk `bundles` status updates instead
|
||||
Bundles map[string]*bundle.Status `json:"bundles,omitempty"`
|
||||
Discovery *bundle.Status `json:"discovery,omitempty"`
|
||||
DecisionLogs *lstat.Status `json:"decision_logs,omitempty"`
|
||||
Metrics map[string]interface{} `json:"metrics,omitempty"`
|
||||
Plugins map[string]*plugins.Status `json:"plugins,omitempty"`
|
||||
}
|
||||
|
||||
// Plugin implements status reporting. Updates can be triggered by the caller.
|
||||
type Plugin struct {
|
||||
manager *plugins.Manager
|
||||
config Config
|
||||
bundleCh chan bundle.Status // Deprecated: Use bulk bundle status updates instead
|
||||
lastBundleStatus *bundle.Status // Deprecated: Use bulk bundle status updates instead
|
||||
bulkBundleCh chan map[string]*bundle.Status
|
||||
lastBundleStatuses map[string]*bundle.Status
|
||||
discoCh chan bundle.Status
|
||||
lastDiscoStatus *bundle.Status
|
||||
pluginStatusCh chan map[string]*plugins.Status
|
||||
lastPluginStatuses map[string]*plugins.Status
|
||||
queryCh chan chan *UpdateRequestV1
|
||||
stop chan chan struct{}
|
||||
reconfig chan interface{}
|
||||
metrics metrics.Metrics
|
||||
logger logging.Logger
|
||||
trigger chan trigger
|
||||
manager *plugins.Manager
|
||||
config Config
|
||||
bundleCh chan bundle.Status // Deprecated: Use bulk bundle status updates instead
|
||||
lastBundleStatus *bundle.Status // Deprecated: Use bulk bundle status updates instead
|
||||
bulkBundleCh chan map[string]*bundle.Status
|
||||
lastBundleStatuses map[string]*bundle.Status
|
||||
discoCh chan bundle.Status
|
||||
lastDiscoStatus *bundle.Status
|
||||
pluginStatusCh chan map[string]*plugins.Status
|
||||
decisionLogsCh chan lstat.Status
|
||||
lastDecisionLogsStatus *lstat.Status
|
||||
lastPluginStatuses map[string]*plugins.Status
|
||||
queryCh chan chan *UpdateRequestV1
|
||||
stop chan chan struct{}
|
||||
reconfig chan interface{}
|
||||
metrics metrics.Metrics
|
||||
logger logging.Logger
|
||||
trigger chan trigger
|
||||
}
|
||||
|
||||
// Config contains configuration for the plugin.
|
||||
@@ -192,6 +196,7 @@ func New(parsedConfig *Config, manager *plugins.Manager) *Plugin {
|
||||
bundleCh: make(chan bundle.Status),
|
||||
bulkBundleCh: make(chan map[string]*bundle.Status),
|
||||
discoCh: make(chan bundle.Status),
|
||||
decisionLogsCh: make(chan lstat.Status),
|
||||
stop: make(chan chan struct{}),
|
||||
reconfig: make(chan interface{}),
|
||||
pluginStatusCh: make(chan map[string]*plugins.Status),
|
||||
@@ -279,6 +284,11 @@ func (p *Plugin) UpdateDiscoveryStatus(status bundle.Status) {
|
||||
p.discoCh <- status
|
||||
}
|
||||
|
||||
// UpdateDecisionLogsStatus notifies the plugin that status of a decision log upload event.
|
||||
func (p *Plugin) UpdateDecisionLogsStatus(status lstat.Status) {
|
||||
p.decisionLogsCh <- status
|
||||
}
|
||||
|
||||
// UpdatePluginStatus notifies the plugin that a plugin status was updated.
|
||||
func (p *Plugin) UpdatePluginStatus(status map[string]*plugins.Status) {
|
||||
p.pluginStatusCh <- status
|
||||
@@ -358,6 +368,16 @@ func (p *Plugin) loop() {
|
||||
p.logger.Info("Status update sent successfully in response to discovery update.")
|
||||
}
|
||||
}
|
||||
case status := <-p.decisionLogsCh:
|
||||
p.lastDecisionLogsStatus = &status
|
||||
if *p.config.Trigger == plugins.TriggerPeriodic {
|
||||
err := p.oneShot(ctx)
|
||||
if err != nil {
|
||||
p.logger.Error("%v.", err)
|
||||
} else {
|
||||
p.logger.Info("Status update sent successfully in response to discovery update.")
|
||||
}
|
||||
}
|
||||
case newConfig := <-p.reconfig:
|
||||
p.reconfigure(newConfig)
|
||||
case respCh := <-p.queryCh:
|
||||
@@ -437,11 +457,12 @@ func (p *Plugin) reconfigure(config interface{}) {
|
||||
func (p *Plugin) snapshot() *UpdateRequestV1 {
|
||||
|
||||
s := &UpdateRequestV1{
|
||||
Labels: p.manager.Labels(),
|
||||
Discovery: p.lastDiscoStatus,
|
||||
Bundle: p.lastBundleStatus,
|
||||
Bundles: p.lastBundleStatuses,
|
||||
Plugins: p.lastPluginStatuses,
|
||||
Labels: p.manager.Labels(),
|
||||
Discovery: p.lastDiscoStatus,
|
||||
DecisionLogs: p.lastDecisionLogsStatus,
|
||||
Bundle: p.lastBundleStatus,
|
||||
Bundles: p.lastBundleStatuses,
|
||||
Plugins: p.lastPluginStatuses,
|
||||
}
|
||||
|
||||
if p.metrics != nil {
|
||||
|
||||
@@ -23,6 +23,8 @@ import (
|
||||
inmem "github.com/open-policy-agent/opa/storage/inmem/test"
|
||||
"github.com/open-policy-agent/opa/util"
|
||||
"github.com/open-policy-agent/opa/version"
|
||||
|
||||
lstat "github.com/open-policy-agent/opa/plugins/logs/status"
|
||||
)
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
@@ -464,6 +466,49 @@ func TestPluginStartDiscovery(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPluginStartDecisionLogs(t *testing.T) {
|
||||
|
||||
fixture := newTestFixture(t, nil)
|
||||
fixture.server.ch = make(chan UpdateRequestV1)
|
||||
defer fixture.server.stop()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
err := fixture.plugin.Start(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer fixture.plugin.Stop(ctx)
|
||||
|
||||
// Ignore the plugin updating its status (tested elsewhere)
|
||||
<-fixture.server.ch
|
||||
|
||||
status := &lstat.Status{
|
||||
Code: "decision_log_error",
|
||||
Message: "Upload Failed",
|
||||
HTTPCode: "400",
|
||||
}
|
||||
|
||||
fixture.plugin.UpdateDecisionLogsStatus(*status)
|
||||
result := <-fixture.server.ch
|
||||
|
||||
exp := UpdateRequestV1{
|
||||
Labels: map[string]string{
|
||||
"id": "test-instance-id",
|
||||
"app": "example-app",
|
||||
"version": version.Version,
|
||||
},
|
||||
DecisionLogs: status,
|
||||
Plugins: map[string]*plugins.Status{
|
||||
"status": {State: plugins.StateOK},
|
||||
},
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(result, exp) {
|
||||
t.Fatalf("Expected: %+v but got: %+v", exp, result)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPluginBadAuth(t *testing.T) {
|
||||
fixture := newTestFixture(t, nil)
|
||||
ctx := context.Background()
|
||||
|
||||
Reference in New Issue
Block a user