Files
releases/plugins/logs/plugin.go
T
Torin Sandall 2f5a0fe0a4 Update decision log event to include error
The error field from the server event was not being copied into the
decision log event. Also, we didn't have test cases to verify that the
error was being set correctly in the first place.

In the future, we should remove the duplication of the server event
and the decision log event (preferring the latter).

Signed-off-by: Torin Sandall <torinsandall@gmail.com>
2019-01-16 12:45:47 -08:00

428 lines
11 KiB
Go

// Copyright 2018 The OPA Authors. All rights reserved.
// Use of this source code is governed by an Apache2
// license that can be found in the LICENSE file.
// Package logs implements decision log buffering and uploading.
package logs
import (
"context"
"fmt"
"math/rand"
"net/http"
"reflect"
"strings"
"sync"
"time"
"github.com/open-policy-agent/opa/plugins"
"github.com/open-policy-agent/opa/plugins/rest"
"github.com/open-policy-agent/opa/server"
"github.com/open-policy-agent/opa/util"
"github.com/open-policy-agent/opa/version"
"github.com/pkg/errors"
"github.com/sirupsen/logrus"
)
// Logger defines the interface for decision logging plugins.
type Logger interface {
plugins.Plugin
Log(context.Context, EventV1)
}
// EventV1 represents a decision log event.
type EventV1 struct {
Labels map[string]string `json:"labels"`
DecisionID string `json:"decision_id"`
Revision string `json:"revision,omitempty"`
Path string `json:"path"`
Input *interface{} `json:"input,omitempty"`
Result *interface{} `json:"result,omitempty"`
Error error `json:"error,omitempty"`
RequestedBy string `json:"requested_by"`
Timestamp time.Time `json:"timestamp"`
Version string `json:"version"`
Metrics map[string]interface{} `json:"metrics,omitempty"`
}
const (
// min amount of time to wait following a failure
minRetryDelay = time.Millisecond * 100
defaultMinDelaySeconds = int64(300)
defaultMaxDelaySeconds = int64(600)
defaultUploadSizeLimitBytes = int64(32768) // 32KB limit
defaultBufferSizeLimitBytes = int64(0) // unlimited
)
// ReportingConfig represents configuration for the plugin's reporting behaviour.
type ReportingConfig struct {
BufferSizeLimitBytes *int64 `json:"buffer_size_limit_bytes,omitempty"` // max size of in-memory buffer
UploadSizeLimitBytes *int64 `json:"upload_size_limit_bytes,omitempty"` // max size of upload payload
MinDelaySeconds *int64 `json:"min_delay_seconds,omitempty"` // min amount of time to wait between successful poll attempts
MaxDelaySeconds *int64 `json:"max_delay_seconds,omitempty"` // max amount of time to wait between poll attempts
}
// Config represents the plugin configuration.
type Config struct {
Plugin *string `json:"plugin"`
Service string `json:"service"`
PartitionName string `json:"partition_name,omitempty"`
Reporting ReportingConfig `json:"reporting"`
}
func (c *Config) validateAndInjectDefaults(services []string, plugins []string) error {
if c.Plugin != nil {
var found bool
for _, other := range plugins {
if other == *c.Plugin {
found = true
break
}
}
if !found {
return fmt.Errorf("invalid plugin name %q in decision_logs", *c.Plugin)
}
} else if c.Service == "" && len(services) != 0 {
c.Service = services[0]
} else {
found := false
for _, svc := range services {
if svc == c.Service {
found = true
break
}
}
if !found {
return fmt.Errorf("invalid service name %q in decision_logs", c.Service)
}
}
min := defaultMinDelaySeconds
max := defaultMaxDelaySeconds
// reject bad min/max values
if c.Reporting.MaxDelaySeconds != nil && c.Reporting.MinDelaySeconds != nil {
if *c.Reporting.MaxDelaySeconds < *c.Reporting.MinDelaySeconds {
return fmt.Errorf("max reporting delay must be >= min reporting delay in decision_logs")
}
min = *c.Reporting.MinDelaySeconds
max = *c.Reporting.MaxDelaySeconds
} else if c.Reporting.MaxDelaySeconds == nil && c.Reporting.MinDelaySeconds != nil {
return fmt.Errorf("reporting configuration missing 'max_delay_seconds' in decision_logs")
} else if c.Reporting.MinDelaySeconds == nil && c.Reporting.MaxDelaySeconds != nil {
return fmt.Errorf("reporting configuration missing 'min_delay_seconds' in decision_logs")
}
// scale to seconds
minSeconds := int64(time.Duration(min) * time.Second)
c.Reporting.MinDelaySeconds = &minSeconds
maxSeconds := int64(time.Duration(max) * time.Second)
c.Reporting.MaxDelaySeconds = &maxSeconds
// default the upload size limit
uploadLimit := defaultUploadSizeLimitBytes
if c.Reporting.UploadSizeLimitBytes != nil {
uploadLimit = *c.Reporting.UploadSizeLimitBytes
}
c.Reporting.UploadSizeLimitBytes = &uploadLimit
// default the buffer size limit
bufferLimit := defaultBufferSizeLimitBytes
if c.Reporting.BufferSizeLimitBytes != nil {
bufferLimit = *c.Reporting.BufferSizeLimitBytes
}
c.Reporting.BufferSizeLimitBytes = &bufferLimit
return nil
}
// Plugin implements decision log buffering and uploading.
type Plugin struct {
manager *plugins.Manager
config Config
buffer *logBuffer
enc *chunkEncoder
mtx sync.Mutex
stop chan chan struct{}
reconfig chan interface{}
}
// ParseConfig validates the config and injects default values.
func ParseConfig(config []byte, services []string, plugins []string) (*Config, error) {
if config == nil {
return nil, nil
}
var parsedConfig Config
if err := util.Unmarshal(config, &parsedConfig); err != nil {
return nil, err
}
if err := parsedConfig.validateAndInjectDefaults(services, plugins); err != nil {
return nil, err
}
return &parsedConfig, nil
}
// New returns a new Plugin with the given config.
func New(parsedConfig *Config, manager *plugins.Manager) *Plugin {
plugin := &Plugin{
manager: manager,
config: *parsedConfig,
stop: make(chan chan struct{}),
buffer: newLogBuffer(*parsedConfig.Reporting.BufferSizeLimitBytes),
enc: newChunkEncoder(*parsedConfig.Reporting.UploadSizeLimitBytes),
reconfig: make(chan interface{}),
}
return plugin
}
// Name identifies the plugin on manager.
const Name = "decision_logs"
// Lookup returns the decision logs plugin registered with the manager.
func Lookup(manager *plugins.Manager) *Plugin {
if p := manager.Plugin(Name); p != nil {
return p.(*Plugin)
}
return nil
}
// Start starts the plugin.
func (p *Plugin) Start(ctx context.Context) error {
p.logInfo("Starting decision log uploader.")
go p.loop()
return nil
}
// Stop stops the plugin.
func (p *Plugin) Stop(ctx context.Context) {
p.logInfo("Stopping decision log uploader.")
done := make(chan struct{})
p.stop <- done
_ = <-done
}
// Log appends a decision log event to the buffer for uploading.
func (p *Plugin) Log(ctx context.Context, decision *server.Info) {
path := strings.Replace(strings.TrimPrefix(decision.Query, "data."), ".", "/", -1)
event := EventV1{
Labels: p.manager.Labels(),
DecisionID: decision.DecisionID,
Revision: decision.Revision,
Path: path,
Input: &decision.Input,
Result: decision.Results,
RequestedBy: decision.RemoteAddr,
Timestamp: decision.Timestamp,
Version: version.Version,
}
if decision.Metrics != nil {
event.Metrics = decision.Metrics.All()
}
if decision.Error != nil {
event.Error = decision.Error
}
if p.config.Plugin != nil {
proxy, ok := p.manager.Plugin(*p.config.Plugin).(Logger)
if !ok {
p.logError("Plugin does not implement Logger interface. Dropping event.")
return
}
proxy.Log(ctx, event)
return
}
p.mtx.Lock()
defer p.mtx.Unlock()
result, err := p.enc.Write(event)
if err != nil {
p.logError("Log encoding failed: %v.", err)
return
}
if result != nil {
p.bufferChunk(p.buffer, result)
}
}
// Reconfigure notifies the plugin with a new configuration.
func (p *Plugin) Reconfigure(_ context.Context, config interface{}) {
p.reconfig <- config
}
func (p *Plugin) loop() {
ctx, cancel := context.WithCancel(context.Background())
var retry int
for {
var err error
if p.config.Plugin == nil {
var uploaded bool
uploaded, err = p.oneShot(ctx)
if err != nil {
p.logError("%v.", err)
} else if uploaded {
p.logInfo("Logs uploaded successfully.")
} else {
p.logInfo("Log upload skipped.")
}
}
var delay time.Duration
if err == nil {
min := float64(*p.config.Reporting.MinDelaySeconds)
max := float64(*p.config.Reporting.MaxDelaySeconds)
delay = time.Duration(((max - min) * rand.Float64()) + min)
} else {
delay = util.DefaultBackoff(float64(minRetryDelay), float64(*p.config.Reporting.MaxDelaySeconds), retry)
}
if p.config.Plugin == nil {
p.logDebug("Waiting %v before next upload/retry.", delay)
}
timer := time.NewTimer(delay)
select {
case <-timer.C:
if err != nil {
retry++
} else {
retry = 0
}
case newConfig := <-p.reconfig:
p.reconfigure(newConfig)
case done := <-p.stop:
cancel()
done <- struct{}{}
return
}
}
}
func (p *Plugin) oneShot(ctx context.Context) (ok bool, err error) {
// Make a local copy of the plugins's encoder and buffer and create
// a new encoder and buffer. This is needed as locking the buffer for
// the upload duration will block policy evaluation and result in
// increased latency for OPA clients
p.mtx.Lock()
oldChunkEnc := p.enc
oldBuffer := p.buffer
p.buffer = newLogBuffer(*p.config.Reporting.BufferSizeLimitBytes)
p.enc = newChunkEncoder(*p.config.Reporting.UploadSizeLimitBytes)
p.mtx.Unlock()
// Along with uploading the compressed events in the buffer
// to the remote server, flush any pending compressed data to the
// underlying writer and add to the buffer.
chunk, err := oldChunkEnc.Flush()
if err != nil {
return false, err
} else if chunk != nil {
p.bufferChunk(oldBuffer, chunk)
}
if oldBuffer.Len() == 0 {
return false, nil
}
for bs := oldBuffer.Pop(); bs != nil; bs = oldBuffer.Pop() {
err := uploadChunk(ctx, p.manager.Client(p.config.Service), p.config.PartitionName, bs)
if err != nil {
// requeue the chunk
p.mtx.Lock()
p.bufferChunk(p.buffer, bs)
p.mtx.Unlock()
return false, err
}
}
return true, nil
}
func (p *Plugin) reconfigure(config interface{}) {
newConfig := config.(*Config)
if reflect.DeepEqual(p.config, *newConfig) {
p.logDebug("Decision log uploader configuration unchanged.")
return
}
p.logInfo("Decision log uploader configuration changed.")
p.config = *newConfig
}
func (p *Plugin) bufferChunk(buffer *logBuffer, bs []byte) {
dropped := buffer.Push(bs)
if dropped > 0 {
p.logError("Dropped %v chunks from buffer. Reduce reporting interval or increase buffer size.", dropped)
}
}
func uploadChunk(ctx context.Context, client rest.Client, partitionName string, data []byte) error {
resp, err := client.
WithHeader("Content-Type", "application/json").
WithHeader("Content-Encoding", "gzip").
WithBytes(data).
Do(ctx, "POST", fmt.Sprintf("/logs/%v", partitionName))
if err != nil {
return errors.Wrap(err, "Log upload failed")
}
defer util.Close(resp)
switch resp.StatusCode {
case http.StatusOK:
return nil
case http.StatusNotFound:
return fmt.Errorf("Log upload failed, server replied with not found")
case http.StatusUnauthorized:
return fmt.Errorf("Log upload failed, server replied with not authorized")
default:
return fmt.Errorf("Log upload failed, server replied with HTTP %v", resp.StatusCode)
}
}
func (p *Plugin) logError(fmt string, a ...interface{}) {
logrus.WithFields(p.logrusFields()).Errorf(fmt, a...)
}
func (p *Plugin) logInfo(fmt string, a ...interface{}) {
logrus.WithFields(p.logrusFields()).Infof(fmt, a...)
}
func (p *Plugin) logDebug(fmt string, a ...interface{}) {
logrus.WithFields(p.logrusFields()).Debugf(fmt, a...)
}
func (p *Plugin) logrusFields() logrus.Fields {
return logrus.Fields{
"plugin": Name,
}
}