mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-05 08:10:50 +00:00
Add log prefix levels
This commit is contained in:
parent
ebbb196a23
commit
285d0a0af8
31 changed files with 383 additions and 311 deletions
36
api.go
36
api.go
|
|
@ -273,7 +273,7 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error {
|
|||
Index: indexName,
|
||||
})
|
||||
if err != nil {
|
||||
api.server.logger.Printf("problem sending DeleteIndex message: %s", err)
|
||||
api.server.logger.Errorf("problem sending DeleteIndex message: %s", err)
|
||||
return errors.Wrap(err, "sending DeleteIndex message")
|
||||
}
|
||||
api.holder.Stats.Count(MetricDeleteIndex, 1, 1.0)
|
||||
|
|
@ -582,7 +582,7 @@ func (api *API) DeleteField(ctx context.Context, indexName string, fieldName str
|
|||
Field: fieldName,
|
||||
})
|
||||
if err != nil {
|
||||
api.server.logger.Printf("problem sending DeleteField message: %s", err)
|
||||
api.server.logger.Errorf("problem sending DeleteField message: %s", err)
|
||||
return errors.Wrap(err, "sending DeleteField message")
|
||||
}
|
||||
api.holder.Stats.CountWithCustomTags(MetricDeleteField, 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)})
|
||||
|
|
@ -614,7 +614,7 @@ func (api *API) DeleteAvailableShard(_ context.Context, indexName, fieldName str
|
|||
ShardID: shardID,
|
||||
})
|
||||
if err != nil {
|
||||
api.server.logger.Printf("problem sending DeleteAvailableShard message: %s", err)
|
||||
api.server.logger.Errorf("problem sending DeleteAvailableShard message: %s", err)
|
||||
return errors.Wrap(err, "sending DeleteAvailableShard message")
|
||||
}
|
||||
api.holder.Stats.CountWithCustomTags(MetricDeleteAvailableShard, 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)})
|
||||
|
|
@ -636,7 +636,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
|
||||
// Validate that this handler owns the shard.
|
||||
if !snap.OwnsShard(api.NodeID(), indexName, shard) {
|
||||
api.server.logger.Printf("node %s does not own shard %d of index %s", api.NodeID(), shard, indexName)
|
||||
api.server.logger.Errorf("node %s does not own shard %d of index %s", api.NodeID(), shard, indexName)
|
||||
return ErrClusterDoesNotOwnShard
|
||||
}
|
||||
|
||||
|
|
@ -907,16 +907,16 @@ func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, e
|
|||
si := api.server.systemInfo
|
||||
diskCapacity, err := si.DiskCapacity(api.holder.path)
|
||||
if err != nil {
|
||||
api.server.logger.Printf("couldn't read disk capacity: %s", err)
|
||||
api.server.logger.Infof("couldn't read disk capacity: %s", err)
|
||||
}
|
||||
|
||||
memoryCapacity, err := si.MemTotal()
|
||||
if err != nil {
|
||||
api.server.logger.Printf("couldn't read memory capacity: %s", err)
|
||||
api.server.logger.Infof("couldn't read memory capacity: %s", err)
|
||||
}
|
||||
memoryUse, err := si.MemUsed()
|
||||
if err != nil {
|
||||
api.server.logger.Printf("couldn't read memory usage: %s", err)
|
||||
api.server.logger.Infof("couldn't read memory usage: %s", err)
|
||||
}
|
||||
|
||||
// Insert into result.
|
||||
|
|
@ -1135,7 +1135,7 @@ func (api *API) DeleteView(ctx context.Context, indexName string, fieldName stri
|
|||
View: viewName,
|
||||
})
|
||||
if err != nil {
|
||||
api.server.logger.Printf("problem sending DeleteView message: %s", err)
|
||||
api.server.logger.Errorf("problem sending DeleteView message: %s", err)
|
||||
}
|
||||
|
||||
return errors.Wrap(err, "sending DeleteView message")
|
||||
|
|
@ -1483,7 +1483,7 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest,
|
|||
// so don't expect it to be invariant.
|
||||
if !options.Clear {
|
||||
if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil {
|
||||
api.server.logger.Printf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
return err
|
||||
}
|
||||
if err != nil {
|
||||
|
|
@ -1494,7 +1494,7 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest,
|
|||
// Import into fragment.
|
||||
err = field.Import(qcx, req.RowIDs, req.ColumnIDs, timestamps, opts...)
|
||||
if err != nil {
|
||||
api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
return errors.Wrap(err, "importing")
|
||||
}
|
||||
return errors.Wrap(err, "committing")
|
||||
|
|
@ -1602,7 +1602,7 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu
|
|||
// Import columnIDs into existence field.
|
||||
if !options.Clear {
|
||||
if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil {
|
||||
api.server.logger.Printf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
return errors.Wrap(err, "importing existence columns")
|
||||
}
|
||||
}
|
||||
|
|
@ -1611,17 +1611,17 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu
|
|||
if len(req.Values) > 0 {
|
||||
err = field.importValue(qcx, req.ColumnIDs, req.Values, options)
|
||||
if err != nil {
|
||||
api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
}
|
||||
} else if len(req.TimestampValues) > 0 {
|
||||
err = field.importTimestampValue(qcx, req.ColumnIDs, req.TimestampValues, options)
|
||||
if err != nil {
|
||||
api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
}
|
||||
} else if len(req.FloatValues) > 0 {
|
||||
err = field.importFloatValue(qcx, req.ColumnIDs, req.FloatValues, options)
|
||||
if err != nil {
|
||||
api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
|
||||
}
|
||||
}
|
||||
return errors.Wrap(err, "importing value")
|
||||
|
|
@ -1704,7 +1704,7 @@ func (api *API) ImportColumnAttrs(ctx context.Context, req *ImportColumnAttrsReq
|
|||
bulkAttrs[uint64(req.ColumnIDs[n])] = map[string]interface{}{req.AttrKey: req.AttrVals[n]}
|
||||
}
|
||||
if err := index.ColumnAttrStore().SetBulkAttrs(bulkAttrs); err != nil {
|
||||
api.server.logger.Printf("import error: index=%s, shard=%d, len(columns)=%d, err=%s", req.Index, req.Shard, len(req.ColumnIDs), err)
|
||||
api.server.logger.Errorf("import error: index=%s, shard=%d, len(columns)=%d, err=%s", req.Index, req.Shard, len(req.ColumnIDs), err)
|
||||
return errors.Wrap(err, "importing column attrs")
|
||||
}
|
||||
return nil
|
||||
|
|
@ -1776,7 +1776,7 @@ func (api *API) validateShardOwnership(indexName string, shard uint64) error {
|
|||
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
|
||||
// Validate that this handler owns the shard.
|
||||
if !snap.OwnsShard(api.NodeID(), indexName, shard) {
|
||||
api.server.logger.Printf("node %s does not own shard %d of index %s", api.NodeID(), shard, indexName)
|
||||
api.server.logger.Errorf("node %s does not own shard %d of index %s", api.NodeID(), shard, indexName)
|
||||
return ErrClusterDoesNotOwnShard
|
||||
}
|
||||
return nil
|
||||
|
|
@ -1788,14 +1788,14 @@ func (api *API) indexField(indexName string, fieldName string, shard uint64) (*I
|
|||
// Find the Index.
|
||||
index := api.holder.Index(indexName)
|
||||
if index == nil {
|
||||
api.server.logger.Printf("fragment error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrIndexNotFound.Error())
|
||||
api.server.logger.Errorf("fragment error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrIndexNotFound.Error())
|
||||
return nil, nil, newNotFoundError(ErrIndexNotFound, indexName)
|
||||
}
|
||||
|
||||
// Retrieve field.
|
||||
field := index.Field(fieldName)
|
||||
if field == nil {
|
||||
api.server.logger.Printf("field error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrFieldNotFound.Error())
|
||||
api.server.logger.Errorf("field error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrFieldNotFound.Error())
|
||||
return nil, nil, newNotFoundError(ErrFieldNotFound, fieldName)
|
||||
}
|
||||
return index, field, nil
|
||||
|
|
|
|||
|
|
@ -17,9 +17,9 @@ package client
|
|||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"log"
|
||||
"sync"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
pb "github.com/pilosa/pilosa/v2/proto"
|
||||
"github.com/pkg/errors"
|
||||
"google.golang.org/grpc"
|
||||
|
|
@ -33,6 +33,7 @@ const maxMsgSize = 1024 * 1024 * 100 // 100 megs ought to be enough for anybody!
|
|||
type GRPCClient struct {
|
||||
dialTargets []string
|
||||
tlsConfig *tls.Config
|
||||
logger logger.Logger
|
||||
|
||||
mu sync.RWMutex
|
||||
conn *grpc.ClientConn
|
||||
|
|
@ -40,10 +41,11 @@ type GRPCClient struct {
|
|||
}
|
||||
|
||||
// NewGRPCClient returns a new instance of GRPCClient.
|
||||
func NewGRPCClient(dialTargets []string, tlsConfig *tls.Config) (*GRPCClient, error) {
|
||||
func NewGRPCClient(dialTargets []string, tlsConfig *tls.Config, logger logger.Logger) (*GRPCClient, error) {
|
||||
c := &GRPCClient{
|
||||
dialTargets: dialTargets,
|
||||
tlsConfig: tlsConfig,
|
||||
logger: logger,
|
||||
}
|
||||
// resetConn sets GRPCClient.conn when it doesn't
|
||||
// exist yet.
|
||||
|
|
@ -124,8 +126,7 @@ func (c *GRPCClient) Conn() *grpc.ClientConn {
|
|||
c.mu.RUnlock()
|
||||
|
||||
if err := c.resetConn(); err != nil {
|
||||
// TODO: log this error with logger
|
||||
log.Printf("error resetting connection: %s", err)
|
||||
c.logger.Errorf("error resetting connection: %s", err)
|
||||
}
|
||||
|
||||
c.mu.RLock()
|
||||
|
|
|
|||
|
|
@ -1061,7 +1061,7 @@ func (c *cluster) followResizeInstruction(ctx context.Context, instr *ResizeInst
|
|||
// if we don't know about a field locally, log an error because
|
||||
// fields should be created and synced prior to shard creation
|
||||
if f == nil {
|
||||
c.logger.Printf("local field not found: %s/%s", is.Name, fs.Name)
|
||||
c.logger.Errorf("local field not found: %s/%s", is.Name, fs.Name)
|
||||
continue
|
||||
}
|
||||
|
||||
|
|
|
|||
9
cmd.go
9
cmd.go
|
|
@ -16,7 +16,8 @@ package pilosa
|
|||
|
||||
import (
|
||||
"io"
|
||||
"log"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
)
|
||||
|
||||
// CmdIO holds standard unix inputs and outputs.
|
||||
|
|
@ -24,7 +25,7 @@ type CmdIO struct {
|
|||
Stdin io.Reader
|
||||
Stdout io.Writer
|
||||
Stderr io.Writer
|
||||
logger *log.Logger
|
||||
logger logger.Logger
|
||||
}
|
||||
|
||||
// NewCmdIO returns a new instance of CmdIO with inputs and outputs set to the
|
||||
|
|
@ -34,10 +35,10 @@ func NewCmdIO(stdin io.Reader, stdout, stderr io.Writer) *CmdIO {
|
|||
Stdin: stdin,
|
||||
Stdout: stdout,
|
||||
Stderr: stderr,
|
||||
logger: log.New(stderr, "", log.LstdFlags),
|
||||
logger: logger.NewStandardLogger(stderr),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *CmdIO) Logger() *log.Logger {
|
||||
func (c *CmdIO) Logger() logger.Logger {
|
||||
return c.logger
|
||||
}
|
||||
|
|
|
|||
|
|
@ -89,7 +89,7 @@ on the configured port.`,
|
|||
return errors.Wrap(err, "initializing jaeger tracer")
|
||||
}
|
||||
defer closer.Close()
|
||||
tracing.GlobalTracer = opentracing.NewTracer(tracer)
|
||||
tracing.GlobalTracer = opentracing.NewTracer(tracer, Server.Logger())
|
||||
}
|
||||
|
||||
return errors.Wrap(Server.Wait(), "waiting on Server")
|
||||
|
|
|
|||
|
|
@ -15,9 +15,8 @@
|
|||
package ctl
|
||||
|
||||
import (
|
||||
"log"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/http"
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
"github.com/pilosa/pilosa/v2/server"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/spf13/pflag"
|
||||
|
|
@ -27,7 +26,7 @@ import (
|
|||
type CommandWithTLSSupport interface {
|
||||
TLSHost() string
|
||||
TLSConfiguration() server.TLSConfig
|
||||
Logger() *log.Logger
|
||||
Logger() logger.Logger
|
||||
}
|
||||
|
||||
// SetTLSConfig creates common TLS flags
|
||||
|
|
|
|||
|
|
@ -125,7 +125,7 @@ func (d *diagnosticsCollector) CheckVersion() error {
|
|||
|
||||
d.lastVersion = rsp.Version
|
||||
if err := d.compareVersion(rsp.Version); err != nil {
|
||||
d.Logger.Printf("%s\n", err.Error())
|
||||
d.Logger.Infof("%s\n", err.Error())
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -169,7 +169,7 @@ func (d *diagnosticsCollector) Set(name string, value interface{}) {
|
|||
// logErr logs the error and returns true if an error exists
|
||||
func (d *diagnosticsCollector) logErr(err error) bool {
|
||||
if err != nil {
|
||||
d.Logger.Printf("%v", err)
|
||||
d.Logger.Errorf("%v", err)
|
||||
return true
|
||||
}
|
||||
return false
|
||||
|
|
|
|||
|
|
@ -716,7 +716,7 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
|
|||
// See: https://github.com/pilosa/pilosa/issues/2009
|
||||
// TODO: Remove at version 2.0
|
||||
if e.detectRangeCall(c) {
|
||||
e.Holder.Logger.Printf("DEPRECATED: Range() is deprecated, please use Row() instead.")
|
||||
e.Holder.Logger.Infof("DEPRECATED: Range() is deprecated, please use Row() instead.")
|
||||
}
|
||||
|
||||
// If shards are specified, then use that value for shards. If shards aren't
|
||||
|
|
|
|||
19
field.go
19
field.go
|
|
@ -495,7 +495,7 @@ func (f *Field) loadAvailableShards() error {
|
|||
}
|
||||
// some other problem:
|
||||
if err != nil {
|
||||
f.holder.Logger.Printf("available shards file present but unreadable, discarding: %v", err)
|
||||
f.holder.Logger.Errorf("available shards file present but unreadable, discarding: %v", err)
|
||||
err = os.Remove(path)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "deleting corrupt available shards list")
|
||||
|
|
@ -504,7 +504,7 @@ func (f *Field) loadAvailableShards() error {
|
|||
}
|
||||
bm := roaring.NewBitmap()
|
||||
if err = bm.UnmarshalBinary(buf); err != nil {
|
||||
f.holder.Logger.Printf("available shards file corrupt, discarding: %v", err)
|
||||
f.holder.Logger.Errorf("available shards file corrupt, discarding: %v", err)
|
||||
err = os.Remove(path)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "deleting corrupt available shards list")
|
||||
|
|
@ -589,6 +589,7 @@ func (f *Field) Open() error {
|
|||
}
|
||||
|
||||
f.holder.Logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name)
|
||||
|
||||
if err := f.loadAvailableShards(); err != nil {
|
||||
return errors.Wrap(err, "loading available shards")
|
||||
}
|
||||
|
|
@ -637,8 +638,8 @@ func (f *Field) Open() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func blockingWriteAvailableShards(fieldPath string, availableShardBytes []byte) {
|
||||
path := filepath.Join(fieldPath, ".available.shards")
|
||||
func (f *Field) blockingWriteAvailableShards(availableShardBytes []byte) {
|
||||
path := filepath.Join(f.path, ".available.shards")
|
||||
|
||||
// Create a temporary file to save to.
|
||||
tempPath := path + tempExt
|
||||
|
|
@ -650,15 +651,15 @@ func blockingWriteAvailableShards(fieldPath string, availableShardBytes []byte)
|
|||
|
||||
// Move snapshot to data file location.
|
||||
if err := os.Rename(tempPath, path); err != nil {
|
||||
log.Printf("rename snapshot: %s", err)
|
||||
f.holder.Logger.Errorf("rename snapshot: %s", err)
|
||||
}
|
||||
}
|
||||
func nonBlockingWriteAvailableShards(fieldPath string, availableShardBytes []byte, done chan bool) {
|
||||
func (f *Field) nonBlockingWriteAvailableShards(availableShardBytes []byte, done chan bool) {
|
||||
if len(availableShardBytes) == 0 {
|
||||
return
|
||||
}
|
||||
go func() {
|
||||
blockingWriteAvailableShards(fieldPath, availableShardBytes)
|
||||
f.blockingWriteAvailableShards(availableShardBytes)
|
||||
done <- true
|
||||
}()
|
||||
}
|
||||
|
|
@ -678,7 +679,7 @@ func (f *Field) writeAvailableShards() {
|
|||
if len(data) > 0 {
|
||||
if !writing {
|
||||
writing = true
|
||||
nonBlockingWriteAvailableShards(f.path, data, tracker)
|
||||
f.nonBlockingWriteAvailableShards(data, tracker)
|
||||
data = nil
|
||||
}
|
||||
}
|
||||
|
|
@ -689,7 +690,7 @@ func (f *Field) writeAvailableShards() {
|
|||
<-tracker
|
||||
}
|
||||
if len(data) > 0 {
|
||||
blockingWriteAvailableShards(f.path, data)
|
||||
f.blockingWriteAvailableShards(data)
|
||||
}
|
||||
alive = false
|
||||
}
|
||||
|
|
|
|||
24
fragment.go
24
fragment.go
|
|
@ -372,7 +372,7 @@ func (f *fragment) importStorage(data []byte, file *os.File, newGen generation,
|
|||
}
|
||||
return false, fmt.Errorf("unmarshal storage: file=%s, err=%s", file.Name(), err)
|
||||
}
|
||||
f.holder.Logger.Printf("warning: unmarshal storage, file=%s, err=%v", file.Name(), err)
|
||||
f.holder.Logger.Warnf("unmarshal storage, file=%s, err=%v", file.Name(), err)
|
||||
trunc, ok := cause.(roaring.FileShouldBeTruncatedError)
|
||||
if ok && !f.holder.Opts.ReadOnly {
|
||||
// if the holder is ReadOnly, we silently ignore the "advisory"
|
||||
|
|
@ -401,7 +401,7 @@ func (f *fragment) applyStorage(data []byte, file *os.File, newGen generation, m
|
|||
if file != nil {
|
||||
fi, err := file.Stat()
|
||||
if err != nil {
|
||||
f.holder.Logger.Printf("trying to apply new storage to existing bitmap, stat failed: %v", err)
|
||||
f.holder.Logger.Errorf("trying to apply new storage to existing bitmap, stat failed: %v", err)
|
||||
}
|
||||
if err == nil && fi != nil && fi.Size() == 0 {
|
||||
return f.emptyStorage(file)
|
||||
|
|
@ -536,7 +536,7 @@ func (f *fragment) openCache() error {
|
|||
// Unmarshal cache data.
|
||||
var pb pb.Cache
|
||||
if err := proto.Unmarshal(buf, &pb); err != nil {
|
||||
f.holder.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
|
||||
f.holder.Logger.Errorf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -576,13 +576,13 @@ func (f *fragment) Close() error {
|
|||
func (f *fragment) close() error {
|
||||
// Flush cache if closing gracefully.
|
||||
if err := f.flushCache(); err != nil {
|
||||
f.holder.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path())
|
||||
f.holder.Logger.Errorf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path())
|
||||
return errors.Wrap(err, "flushing cache")
|
||||
}
|
||||
|
||||
// Close underlying storage.
|
||||
if err := f.closeStorage(); err != nil {
|
||||
f.holder.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path())
|
||||
f.holder.Logger.Errorf("fragment: error closing storage: err=%s, path=%s", err, f.path())
|
||||
return errors.Wrap(err, "closing storage")
|
||||
}
|
||||
|
||||
|
|
@ -2547,11 +2547,11 @@ func (f *fragment) importPositions(tx Tx, set, clear []uint64, rowSet map[uint64
|
|||
// we got an error. it's possible that the error indicates that something went wrong.
|
||||
mappedIn, mappedOut, unmappedIn, errs, e2 := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to)
|
||||
if errs != 0 {
|
||||
f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
|
||||
f.holder.Logger.Errorf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
|
||||
f.path(), mappedIn, mappedOut, unmappedIn, errs, e2)
|
||||
if f.prevdata.from != f.currdata.from {
|
||||
mappedIn, mappedOut, unmappedIn, errs, e2 = f.storage.SanityCheckMapping(f.prevdata.from, f.prevdata.to)
|
||||
f.holder.Logger.Printf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
|
||||
f.holder.Logger.Errorf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
|
||||
mappedIn, mappedOut, unmappedIn, errs, e2)
|
||||
}
|
||||
}
|
||||
|
|
@ -2675,8 +2675,8 @@ func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int
|
|||
}(); err != nil {
|
||||
errOpenStorage := f.openStorage(true)
|
||||
if errOpenStorage != nil {
|
||||
f.Logger.Printf("failed to import data into fragment: %v", err)
|
||||
f.Logger.Printf("recovery with openStorage failed for fragment: %v", errOpenStorage)
|
||||
f.Logger.Errorf("failed to import data into fragment: %v", err)
|
||||
f.Logger.Errorf("recovery with openStorage failed for fragment: %v", errOpenStorage)
|
||||
f.Logger.Debugf("%s", debug.Stack())
|
||||
os.Exit(1)
|
||||
}
|
||||
|
|
@ -2735,8 +2735,8 @@ func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDep
|
|||
}(); err != nil {
|
||||
errOpenStorage := f.openStorage(true)
|
||||
if errOpenStorage != nil {
|
||||
f.Logger.Printf("failed to import data into fragment: %v", err)
|
||||
f.Logger.Printf("recovery with openStorage failed for fragment: %v", errOpenStorage)
|
||||
f.Logger.Errorf("failed to import data into fragment: %v", err)
|
||||
f.Logger.Errorf("recovery with openStorage failed for fragment: %v", errOpenStorage)
|
||||
f.Logger.Debugf("%s", debug.Stack())
|
||||
os.Exit(1)
|
||||
}
|
||||
|
|
@ -2900,7 +2900,7 @@ func (f *fragment) snapshot() (err error) {
|
|||
// we can't see the actual values that were used to generate this, probably.
|
||||
if e2.Error() == "runtime error: invalid memory address or nil pointer dereference" {
|
||||
mappedIn, mappedOut, unmappedIn, errs, _ := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to)
|
||||
f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total",
|
||||
f.holder.Logger.Errorf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total",
|
||||
f.path(), mappedIn, mappedOut, unmappedIn, errs)
|
||||
}
|
||||
} else {
|
||||
|
|
|
|||
|
|
@ -139,7 +139,7 @@ func (m *mmapGeneration) Transaction(fileP *io.Writer, fn func() error) (transac
|
|||
// open.
|
||||
if m.dead {
|
||||
elapsed := time.Since(m.deadSince)
|
||||
m.logger.Printf("WARNING: transaction against %s, which has been dead for %v\n", m.id, elapsed)
|
||||
m.logger.Warnf("transaction against %s, which has been dead for %v\n", m.id, elapsed)
|
||||
}
|
||||
if fileP != nil {
|
||||
if m.file == nil {
|
||||
|
|
@ -217,7 +217,7 @@ func (m *mmapGeneration) Done() {
|
|||
m.deadSince = time.Now()
|
||||
err := m.closeFile()
|
||||
if err != nil {
|
||||
m.logger.Printf("error closing generation %s: %v", m.id, err)
|
||||
m.logger.Errorf("error closing generation %s: %v", m.id, err)
|
||||
}
|
||||
// If we're not debugging, the finalizer won't have been enabled
|
||||
// previously. Finalizers have non-zero cost, so having them not be
|
||||
|
|
@ -274,18 +274,18 @@ func (m *mmapGeneration) openFile() (shouldClose bool, err error) {
|
|||
func generationFinalizer(m *mmapGeneration) {
|
||||
m.mu.Lock()
|
||||
if !m.dead {
|
||||
m.logger.Printf("finalizing generation %s which isn't dead yet\n",
|
||||
m.logger.Infof("finalizing generation %s which isn't dead yet\n",
|
||||
m.id)
|
||||
}
|
||||
m.mu.Unlock()
|
||||
err := m.closeFile()
|
||||
if err != nil {
|
||||
m.logger.Printf("finalizing generation, closing file: %v\n", err)
|
||||
m.logger.Errorf("finalizing generation, closing file: %v\n", err)
|
||||
}
|
||||
if m.data != nil {
|
||||
err := syswrap.Munmap(m.data)
|
||||
if err != nil {
|
||||
m.logger.Printf("finalizing generation, munmap: %v\n", err)
|
||||
m.logger.Errorf("finalizing generation, munmap: %v\n", err)
|
||||
}
|
||||
m.data = nil
|
||||
}
|
||||
|
|
@ -308,7 +308,7 @@ func (m *mmapGeneration) Cancel() {
|
|||
}
|
||||
err := m.closeFile()
|
||||
if err != nil {
|
||||
m.logger.Printf("error cancelling generation %s: %v", m.id, err)
|
||||
m.logger.Errorf("error cancelling generation %s: %v", m.id, err)
|
||||
}
|
||||
runtime.SetFinalizer(m, nil)
|
||||
m.dead = true
|
||||
|
|
@ -362,7 +362,7 @@ func newGeneration(existing generation, path string, readData bool, setup func([
|
|||
data, err = syswrap.Mmap(int(m.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED)
|
||||
if err == syswrap.ErrMaxMapCountReached {
|
||||
// I have no idea where/how to display this message.
|
||||
m.logger.Printf("maximum number of maps reached, reading file '%s' instead", m.path)
|
||||
m.logger.Warnf("maximum number of maps reached, reading file '%s' instead", m.path)
|
||||
} else if err != nil {
|
||||
m.Cancel()
|
||||
return nil, errors.Wrap(err, "mmap failed")
|
||||
|
|
@ -389,12 +389,12 @@ func newGeneration(existing generation, path string, readData bool, setup func([
|
|||
// be truncated: For instance, if a bitmap has a corrupted
|
||||
// ops log, we could truncate that part of it and retry.
|
||||
if err, ok := err.(roaring.FileShouldBeTruncatedError); ok && m.retries < 1 {
|
||||
m.logger.Printf("file %s read partially, but should-be-truncated at %d bytes\n", m.path, err.SuggestedLength())
|
||||
m.logger.Infof("file %s read partially, but should-be-truncated at %d bytes\n", m.path, err.SuggestedLength())
|
||||
// close this generation, then try again. once.
|
||||
m.retries++
|
||||
err := os.Truncate(m.path, err.SuggestedLength())
|
||||
if err != nil {
|
||||
m.logger.Printf("truncating file failed [but retrying anyway]: %v\n", err)
|
||||
m.logger.Errorf("truncating file failed [but retrying anyway]: %v\n", err)
|
||||
}
|
||||
return newGeneration(&m, path, readData, setup, logger)
|
||||
}
|
||||
|
|
@ -417,7 +417,7 @@ func newGeneration(existing generation, path string, readData bool, setup func([
|
|||
// doesn't need to exist, yay.
|
||||
unmapErr := syswrap.Munmap(data)
|
||||
if unmapErr != nil {
|
||||
m.logger.Printf("error unmapping (probably harmless): %v", unmapErr)
|
||||
m.logger.Errorf("error unmapping (probably harmless): %v", unmapErr)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -427,7 +427,7 @@ func newGeneration(existing generation, path string, readData bool, setup func([
|
|||
if shouldClose {
|
||||
err := m.closeFile()
|
||||
if err != nil {
|
||||
m.logger.Printf("closing file to preserve open files failed: %v\n", err)
|
||||
m.logger.Errorf("closing file to preserve open files failed: %v\n", err)
|
||||
}
|
||||
}
|
||||
// It's possible that the generation has no actual data to track,
|
||||
|
|
|
|||
30
holder.go
30
holder.go
|
|
@ -683,7 +683,7 @@ func (h *Holder) Open() error {
|
|||
|
||||
index, err := h.newIndex(h.IndexPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
|
||||
if errors.Cause(err) == ErrName {
|
||||
h.Logger.Printf("ERROR opening index: %s, err=%s", fi.Name(), err)
|
||||
h.Logger.Errorf("opening index: %s, err=%s", fi.Name(), err)
|
||||
continue
|
||||
} else if err != nil {
|
||||
return errors.Wrap(err, "opening index")
|
||||
|
|
@ -702,7 +702,7 @@ func (h *Holder) Open() error {
|
|||
if err != nil {
|
||||
_ = h.txf.Close()
|
||||
if err == ErrName {
|
||||
h.Logger.Printf("ERROR opening index: %s, err=%s", index.Name(), err)
|
||||
h.Logger.Errorf("opening index: %s, err=%s", index.Name(), err)
|
||||
continue
|
||||
}
|
||||
return fmt.Errorf("open index: name=%s, err=%s", index.Name(), err)
|
||||
|
|
@ -1428,7 +1428,7 @@ func (h *Holder) flushCaches() {
|
|||
}
|
||||
|
||||
if err := fragment.FlushCache(); err != nil {
|
||||
h.Logger.Printf("ERROR flushing cache: err=%s, path=%s", err, fragment.cachePath())
|
||||
h.Logger.Errorf("flushing cache: err=%s, path=%s", err, fragment.cachePath())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1454,7 +1454,7 @@ func (h *Holder) setFileLimit() {
|
|||
newLimit := &syscall.Rlimit{}
|
||||
|
||||
if err := syscall.Getrlimit(syscall.RLIMIT_NOFILE, oldLimit); err != nil {
|
||||
h.Logger.Printf("ERROR checking open file limit: %s", err)
|
||||
h.Logger.Errorf("checking open file limit: %s", err)
|
||||
return
|
||||
}
|
||||
// If the soft limit is lower than the FileLimit constant, we will try to change it.
|
||||
|
|
@ -1478,26 +1478,26 @@ func (h *Holder) setFileLimit() {
|
|||
}
|
||||
// Try setting again with lowered Max (hard limit)
|
||||
if err := syscall.Setrlimit(syscall.RLIMIT_NOFILE, newLimit); err != nil {
|
||||
h.Logger.Printf("ERROR setting open file limit: %s", err)
|
||||
h.Logger.Errorf("setting open file limit: %s", err)
|
||||
}
|
||||
// If we weren't trying to change the hard limit, let the user know something is wrong.
|
||||
} else {
|
||||
h.Logger.Printf("ERROR setting open file limit: %s", err)
|
||||
h.Logger.Errorf("setting open file limit: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Check the limit after setting it. OS may not obey Setrlimit call.
|
||||
if err := syscall.Getrlimit(syscall.RLIMIT_NOFILE, oldLimit); err != nil {
|
||||
h.Logger.Printf("ERROR checking open file limit: %s", err)
|
||||
h.Logger.Errorf("checking open file limit: %s", err)
|
||||
} else {
|
||||
if oldLimit.Cur < fileLimit {
|
||||
h.Logger.Printf("WARNING: Tried to set open file limit to %d, but it is %d. You may consider running \"sudo ulimit -n %d\" before starting Pilosa to avoid \"too many open files\" error. See https://www.pilosa.com/docs/latest/administration/#open-file-limits for more information.", fileLimit, oldLimit.Cur, fileLimit)
|
||||
h.Logger.Warnf("Tried to set open file limit to %d, but it is %d. You may consider running \"sudo ulimit -n %d\" before starting Pilosa to avoid \"too many open files\" error. See https://www.pilosa.com/docs/latest/administration/#open-file-limits for more information.", fileLimit, oldLimit.Cur, fileLimit)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Log startup time and version to $DATA_DIR/startup.log
|
||||
// Log startup time and version to $DATA_DIR/.startup.log
|
||||
func (h *Holder) logStartup() error {
|
||||
RFC3339NanoFixedWidth := "2006-01-02T15:04:05.000000 07:00"
|
||||
time := time.Now().Format(RFC3339NanoFixedWidth)
|
||||
|
|
@ -2006,7 +2006,7 @@ func (s *holderSyncer) readBothTranslateReader(rd TranslateEntryReader, snap *to
|
|||
for {
|
||||
var entry TranslateEntry
|
||||
if err := rd.ReadEntry(&entry); err != nil {
|
||||
s.Holder.Logger.Printf("cannot read translate entry: %s", err)
|
||||
s.Holder.Logger.Errorf("cannot read translate entry: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -2015,30 +2015,30 @@ func (s *holderSyncer) readBothTranslateReader(rd TranslateEntryReader, snap *to
|
|||
// Find appropriate store.
|
||||
f := s.Holder.Field(entry.Index, entry.Field)
|
||||
if f == nil {
|
||||
s.Holder.Logger.Printf("field not found: %s/%s", entry.Index, entry.Field)
|
||||
s.Holder.Logger.Errorf("field not found: %s/%s", entry.Index, entry.Field)
|
||||
return
|
||||
}
|
||||
store = f.TranslateStore()
|
||||
if store == nil {
|
||||
s.Holder.Logger.Printf("no translate store suitable for index %q, field %q, key %q", entry.Index, entry.Field, entry.Key)
|
||||
s.Holder.Logger.Errorf("no translate store suitable for index %q, field %q, key %q", entry.Index, entry.Field, entry.Key)
|
||||
return
|
||||
}
|
||||
} else {
|
||||
// Find appropriate store.
|
||||
idx := s.Holder.Index(entry.Index)
|
||||
if idx == nil {
|
||||
s.Holder.Logger.Printf("index not found: %q", entry.Index)
|
||||
s.Holder.Logger.Errorf("index not found: %q", entry.Index)
|
||||
return
|
||||
}
|
||||
store = idx.TranslateStore(snap.KeyToKeyPartition(entry.Index, entry.Key))
|
||||
if store == nil {
|
||||
s.Holder.Logger.Printf("no translate store suitable for index %q, key %q", entry.Index, entry.Key)
|
||||
s.Holder.Logger.Errorf("no translate store suitable for index %q, key %q", entry.Index, entry.Key)
|
||||
return
|
||||
}
|
||||
}
|
||||
// Apply replication to store.
|
||||
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
|
||||
s.Holder.Logger.Printf("cannot force set field translation data: %d=%q", entry.ID, entry.Key)
|
||||
s.Holder.Logger.Errorf("cannot force set field translation data: %d=%q", entry.ID, entry.Key)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -202,7 +202,7 @@ func NewHandler(opts ...handlerOption) (*Handler, error) {
|
|||
func (h *Handler) Serve() error {
|
||||
err := h.server.Serve(h.ln)
|
||||
if err != nil && err.Error() != "http: Server closed" {
|
||||
h.logger.Printf("HTTP handler terminated with error: %s\n", err)
|
||||
h.logger.Errorf("HTTP handler terminated with error: %s\n", err)
|
||||
return errors.Wrap(err, "serve http")
|
||||
}
|
||||
return nil
|
||||
|
|
@ -295,7 +295,7 @@ func (h *Handler) queryArgValidator(next http.Handler) http.Handler {
|
|||
response := errorResponse{Error: errText}
|
||||
data, err := json.Marshal(response)
|
||||
if err != nil {
|
||||
h.logger.Printf("failed to encode error %q as JSON: %v", errText, err)
|
||||
h.logger.Errorf("failed to encode error %q as JSON: %v", errText, err)
|
||||
} else {
|
||||
errText = string(data)
|
||||
}
|
||||
|
|
@ -487,8 +487,8 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||
if err := recover(); err != nil {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
stack := debug.Stack()
|
||||
msg := "PANIC: %s\n%s"
|
||||
h.logger.Printf(msg, err, stack)
|
||||
msg := "%s\n%s"
|
||||
h.logger.Panicf(msg, err, stack)
|
||||
fmt.Fprintf(w, msg, err, stack)
|
||||
}
|
||||
}()
|
||||
|
|
@ -529,7 +529,7 @@ func (s statikHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
if s.statikFS == nil {
|
||||
msg := "Web UI is not available. Please run `make generate-statik` before building Pilosa with `make install`."
|
||||
s.handler.logger.Printf(msg)
|
||||
s.handler.logger.Infof(msg)
|
||||
http.Error(w, msg, http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
|
@ -601,12 +601,12 @@ func (r *successResponse) write(w http.ResponseWriter, err error) {
|
|||
w.Header().Set("Content-Type", "application/json")
|
||||
_, err := w.Write(msg)
|
||||
if err != nil {
|
||||
r.h.logger.Printf("error writing response: %v", err)
|
||||
r.h.logger.Errorf("error writing response: %v", err)
|
||||
return
|
||||
}
|
||||
_, err = w.Write([]byte("\n"))
|
||||
if err != nil {
|
||||
r.h.logger.Printf("error writing newline after response: %v", err)
|
||||
r.h.logger.Errorf("error writing newline after response: %v", err)
|
||||
return
|
||||
}
|
||||
} else {
|
||||
|
|
@ -681,7 +681,7 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
|
||||
if err := json.NewEncoder(w).Encode(pilosa.Schema{Indexes: schema}); err != nil {
|
||||
h.logger.Printf("write schema response error: %s", err)
|
||||
h.logger.Errorf("write schema response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -745,7 +745,7 @@ func (h *Handler) handleGetUsage(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(nodeUsages); err != nil {
|
||||
h.logger.Printf("write status response error: %s", err)
|
||||
h.logger.Errorf("write status response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -754,7 +754,7 @@ func (h *Handler) handleGetShardDistribution(w http.ResponseWriter, r *http.Requ
|
|||
dist := h.api.ShardDistribution(r.Context())
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(dist); err != nil {
|
||||
h.logger.Printf("write status response error: %s", err)
|
||||
h.logger.Errorf("write status response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -779,7 +779,7 @@ func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(status); err != nil {
|
||||
h.logger.Printf("write status response error: %s", err)
|
||||
h.logger.Errorf("write status response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -791,7 +791,7 @@ func (h *Handler) handleGetInfo(w http.ResponseWriter, r *http.Request) {
|
|||
info := h.api.Info()
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(info); err != nil {
|
||||
h.logger.Printf("write info response error: %s", err)
|
||||
h.logger.Errorf("write info response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -822,7 +822,7 @@ func (h *Handler) handleInspect(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(info); err != nil {
|
||||
h.logger.Printf("write inspect response error: %s", err)
|
||||
h.logger.Errorf("write inspect response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -883,7 +883,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
|
|||
w.WriteHeader(http.StatusBadRequest)
|
||||
e := h.writeQueryResponse(w, r, &pilosa.QueryResponse{Err: err})
|
||||
if e != nil {
|
||||
h.logger.Printf("write query response error: %v (while trying to write another error: %v)", e, err)
|
||||
h.logger.Errorf("write query response error: %v (while trying to write another error: %v)", e, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -905,7 +905,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
e := h.writeQueryResponse(w, r, &pilosa.QueryResponse{Err: err})
|
||||
if e != nil {
|
||||
h.logger.Printf("write query response error: %v (while trying to write another error: %v)", e, err)
|
||||
h.logger.Errorf("write query response error: %v (while trying to write another error: %v)", e, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -924,7 +924,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
// Write response back to client.
|
||||
if err := h.writeQueryResponse(w, r, &resp); err != nil {
|
||||
h.logger.Printf("write query response error: %s", err)
|
||||
h.logger.Errorf("write query response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -986,7 +986,7 @@ func (h *Handler) handleGetShardsMax(w http.ResponseWriter, r *http.Request) {
|
|||
if err := json.NewEncoder(w).Encode(getShardsMaxResponse{
|
||||
Standard: h.api.MaxShards(r.Context()),
|
||||
}); err != nil {
|
||||
h.logger.Printf("write shards-max response error: %s", err)
|
||||
h.logger.Errorf("write shards-max response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1018,7 +1018,7 @@ func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) {
|
|||
if idx.Name == indexName {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(idx); err != nil {
|
||||
h.logger.Printf("write response error: %s", err)
|
||||
h.logger.Errorf("write response error: %s", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -1186,7 +1186,7 @@ func (h *Handler) handlePostIndexAttrDiff(w http.ResponseWriter, r *http.Request
|
|||
if err := json.NewEncoder(w).Encode(postIndexAttrDiffResponse{
|
||||
Attrs: attrs,
|
||||
}); err != nil {
|
||||
h.logger.Printf("response encoding error: %s", err)
|
||||
h.logger.Errorf("response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1222,16 +1222,16 @@ func (h *Handler) handleGetActiveQueries(w http.ResponseWriter, r *http.Request)
|
|||
for i, q := range queries {
|
||||
_, err := fmt.Fprintf(w, "%*s%q\n", -(maxlen + 2), durations[i], q.PQL)
|
||||
if err != nil {
|
||||
h.logger.Printf("sending GetActiveQueries response: %s", err)
|
||||
h.logger.Errorf("sending GetActiveQueries response: %s", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
if _, err := w.Write([]byte{'\n'}); err != nil {
|
||||
h.logger.Printf("sending GetActiveQueries response: %s", err)
|
||||
h.logger.Errorf("sending GetActiveQueries response: %s", err)
|
||||
}
|
||||
case "application/json":
|
||||
if err := json.NewEncoder(w).Encode(queries); err != nil {
|
||||
h.logger.Printf("encoding GetActiveQueries response: %s", err)
|
||||
h.logger.Errorf("encoding GetActiveQueries response: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1252,7 +1252,7 @@ func (h *Handler) handleGetPastQueries(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(queries); err != nil {
|
||||
h.logger.Printf("encoding GetActiveQueries response: %s", err)
|
||||
h.logger.Errorf("encoding GetActiveQueries response: %s", err)
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1559,7 +1559,7 @@ func (h *Handler) handleGetTransactionList(w http.ResponseWriter, r *http.Reques
|
|||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(trnsList); err != nil {
|
||||
h.logger.Printf("encoding GetTransactionList response: %s", err)
|
||||
h.logger.Errorf("encoding GetTransactionList response: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1581,7 +1581,7 @@ func (h *Handler) handleGetTransactions(w http.ResponseWriter, r *http.Request)
|
|||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(trnsMap); err != nil {
|
||||
h.logger.Printf("encoding GetTransactions response: %s", err)
|
||||
h.logger.Errorf("encoding GetTransactions response: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1612,7 +1612,7 @@ func (h *Handler) doTransactionResponse(w http.ResponseWriter, err error, trns *
|
|||
err = json.NewEncoder(w).Encode(
|
||||
TransactionResponse{Error: errString, Transaction: trns})
|
||||
if err != nil {
|
||||
h.logger.Printf("encoding transaction response: %v", err)
|
||||
h.logger.Errorf("encoding transaction response: %v", err)
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1709,7 +1709,7 @@ func (h *Handler) handlePostFieldAttrDiff(w http.ResponseWriter, r *http.Request
|
|||
if err := json.NewEncoder(w).Encode(postFieldAttrDiffResponse{
|
||||
Attrs: attrs,
|
||||
}); err != nil {
|
||||
h.logger.Printf("response encoding error: %s", err)
|
||||
h.logger.Errorf("response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1862,7 +1862,7 @@ func (h *Handler) handleGetMetricsJSON(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
err := json.NewEncoder(w).Encode(metrics)
|
||||
if err != nil {
|
||||
h.logger.Printf("json write error: %s", err)
|
||||
h.logger.Errorf("json write error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1926,7 +1926,7 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request)
|
|||
// Write to response.
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(nodes); err != nil {
|
||||
h.logger.Printf("json write error: %s", err)
|
||||
h.logger.Errorf("json write error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1943,7 +1943,7 @@ func (h *Handler) handleGetNodes(w http.ResponseWriter, r *http.Request) {
|
|||
// Write to response.
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(nodes); err != nil {
|
||||
h.logger.Printf("json write error: %s", err)
|
||||
h.logger.Errorf("json write error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1966,7 +1966,7 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ
|
|||
w.Header().Set("Content-Length", strconv.Itoa(len(buf)))
|
||||
_, err = w.Write(buf)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing fragment/block/data response: %v", err)
|
||||
h.logger.Errorf("writing fragment/block/data response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1999,7 +1999,7 @@ func (h *Handler) handleGetFragmentBlocks(w http.ResponseWriter, r *http.Request
|
|||
if err := json.NewEncoder(w).Encode(getFragmentBlocksResponse{
|
||||
Blocks: blocks,
|
||||
}); err != nil {
|
||||
h.logger.Printf("block response encoding error: %s", err)
|
||||
h.logger.Errorf("block response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2024,7 +2024,7 @@ func (h *Handler) handleGetFragmentData(w http.ResponseWriter, r *http.Request)
|
|||
}
|
||||
// Stream fragment to response body.
|
||||
if _, err := f.WriteTo(w); err != nil {
|
||||
h.logger.Printf("error streaming fragment data: %s", err)
|
||||
h.logger.Errorf("error streaming fragment data: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2045,7 +2045,7 @@ func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request)
|
|||
}
|
||||
// Stream translate partition to response body.
|
||||
if _, err := p.WriteTo(w); err != nil {
|
||||
h.logger.Printf("error streaming translation data: %s", err)
|
||||
h.logger.Errorf("error streaming translation data: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2062,7 +2062,7 @@ func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) {
|
|||
Version: h.api.Version(),
|
||||
})
|
||||
if err != nil {
|
||||
h.logger.Printf("write version response error: %s", err)
|
||||
h.logger.Errorf("write version response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2121,7 +2121,7 @@ func (h *Handler) handlePostClusterResizeRemoveNode(w http.ResponseWriter, r *ht
|
|||
if err := json.NewEncoder(w).Encode(removeNodeResponse{
|
||||
Remove: removeNode,
|
||||
}); err != nil {
|
||||
h.logger.Printf("response encoding error: %s", err)
|
||||
h.logger.Errorf("response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2158,7 +2158,7 @@ func (h *Handler) handlePostClusterResizeAbort(w http.ResponseWriter, r *http.Re
|
|||
if err := json.NewEncoder(w).Encode(clusterResizeAbortResponse{
|
||||
Info: msg,
|
||||
}); err != nil {
|
||||
h.logger.Printf("response encoding error: %s", err)
|
||||
h.logger.Errorf("response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2199,7 +2199,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
|
|||
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if err := json.NewEncoder(w).Encode(defaultClusterMessageResponse{}); err != nil {
|
||||
h.logger.Printf("response encoding error: %s", err)
|
||||
h.logger.Errorf("response encoding error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2235,7 +2235,7 @@ func (h *Handler) handlePostTranslateData(w http.ResponseWriter, r *http.Request
|
|||
if err := rd.ReadEntry(&entry); err == io.EOF {
|
||||
return
|
||||
} else if err != nil {
|
||||
h.logger.Printf("http: translate store read error: %s", err)
|
||||
h.logger.Errorf("http: translate store read error: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -2363,7 +2363,7 @@ func (h *Handler) handlePostImportAtomicRecord(w http.ResponseWriter, r *http.Re
|
|||
// Write response.
|
||||
_, err = w.Write(importOk)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing import response: %v", err)
|
||||
h.logger.Errorf("writing import response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2468,7 +2468,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
|
|||
// Write response.
|
||||
_, err = w.Write(importOk)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing import response: %v", err)
|
||||
h.logger.Errorf("writing import response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2510,7 +2510,7 @@ func (h *Handler) handlePostImportColumnAttrs(w http.ResponseWriter, r *http.Req
|
|||
// Write response.
|
||||
_, err = w.Write(importOk)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing import-column-attrs response: %v", err)
|
||||
h.logger.Errorf("writing import-column-attrs response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -2587,7 +2587,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request
|
|||
// Write response.
|
||||
_, err = w.Write(buf)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing import-roaring response: %v", err)
|
||||
h.logger.Errorf("writing import-roaring response: %v", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
|
@ -2607,7 +2607,7 @@ func (h *Handler) handlePostTranslateKeys(w http.ResponseWriter, r *http.Request
|
|||
case nil:
|
||||
// Write response.
|
||||
if _, err = w.Write(buf); err != nil {
|
||||
h.logger.Printf("writing translate keys response: %v", err)
|
||||
h.logger.Errorf("writing translate keys response: %v", err)
|
||||
}
|
||||
|
||||
case pilosa.ErrTranslatingKeyNotFound:
|
||||
|
|
@ -2639,7 +2639,7 @@ func (h *Handler) handlePostTranslateIDs(w http.ResponseWriter, r *http.Request)
|
|||
// Write response.
|
||||
_, err = w.Write(buf)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing translate keys response: %v", err)
|
||||
h.logger.Errorf("writing translate keys response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
198
logger/logger.go
198
logger/logger.go
|
|
@ -15,9 +15,12 @@
|
|||
package logger
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
|
|
@ -28,8 +31,24 @@ var _ Logger = &nopLogger{}
|
|||
|
||||
// Logger represents an interface for a shared logger.
|
||||
type Logger interface {
|
||||
Printf(format string, v ...interface{})
|
||||
Printf(format string, v ...interface{}) // backward compatibility
|
||||
Debugf(format string, v ...interface{})
|
||||
Infof(format string, v ...interface{})
|
||||
Warnf(format string, v ...interface{})
|
||||
Errorf(format string, v ...interface{})
|
||||
Panicf(format string, v ...interface{})
|
||||
}
|
||||
|
||||
const (
|
||||
LevelPanic = iota
|
||||
LevelError
|
||||
LevelWarn
|
||||
LevelInfo
|
||||
LevelDebug
|
||||
)
|
||||
|
||||
func LevelPrefix(level int) string {
|
||||
return [...]string{"PANIC: ", "ERROR: ", "WARN: ", "INFO: ", "DEBUG: "}[level]
|
||||
}
|
||||
|
||||
// NopLogger represents a Logger that doesn't do anything.
|
||||
|
|
@ -43,9 +62,22 @@ func (n *nopLogger) Printf(format string, v ...interface{}) {}
|
|||
// Debugf is a no-op implementation of the Logger Debugf method.
|
||||
func (n *nopLogger) Debugf(format string, v ...interface{}) {}
|
||||
|
||||
// Infof is a no-op implementation of the Logger Printf method.
|
||||
func (n *nopLogger) Infof(format string, v ...interface{}) {}
|
||||
|
||||
// Warnf is a no-op implementation of the Logger Warnf method.
|
||||
func (n *nopLogger) Warnf(format string, v ...interface{}) {}
|
||||
|
||||
// Errorf is a no-op implementation of the Logger Errorf method.
|
||||
func (n *nopLogger) Errorf(format string, v ...interface{}) {}
|
||||
|
||||
// Panicf is a no-op implementation of the Logger Panicf method.
|
||||
func (n *nopLogger) Panicf(format string, v ...interface{}) {}
|
||||
|
||||
// standardLogger is a basic implementation of Logger based on log.Logger.
|
||||
type standardLogger struct {
|
||||
logger *log.Logger
|
||||
logger *log.Logger
|
||||
verbosity int
|
||||
}
|
||||
|
||||
// write in UTC with constant width and microsecond resolution.
|
||||
|
|
@ -57,51 +89,61 @@ func (fl formatLog) Write(bytes []byte) (int, error) {
|
|||
return fmt.Fprintf(fl.w, "%v %v", time.Now().UTC().Format(RFC3339UsecTz0), string(bytes))
|
||||
}
|
||||
|
||||
func NewStandardLogger(w io.Writer) *standardLogger {
|
||||
func newStandardLogger(w io.Writer, verbosity int) *standardLogger {
|
||||
logger := log.New(w, "", 0)
|
||||
logger.SetOutput(formatLog{w: w})
|
||||
return &standardLogger{
|
||||
logger: logger,
|
||||
logger: logger,
|
||||
verbosity: verbosity,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *standardLogger) Printf(format string, v ...interface{}) {
|
||||
s.logger.Printf(format, v...)
|
||||
func NewStandardLogger(w io.Writer) *standardLogger {
|
||||
return newStandardLogger(w, LevelInfo)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Debugf(format string, v ...interface{}) {}
|
||||
func NewVerboseLogger(w io.Writer) *standardLogger {
|
||||
return newStandardLogger(w, LevelDebug)
|
||||
}
|
||||
|
||||
func (s *standardLogger) printf(level int, format string, v ...interface{}) {
|
||||
if level > s.verbosity {
|
||||
return
|
||||
}
|
||||
|
||||
s.logger.Printf(LevelPrefix(level)+format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Printf(format string, v ...interface{}) {
|
||||
s.printf(LevelInfo, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Debugf(format string, v ...interface{}) {
|
||||
s.printf(LevelDebug, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Infof(format string, v ...interface{}) {
|
||||
s.printf(LevelInfo, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Warnf(format string, v ...interface{}) {
|
||||
s.printf(LevelWarn, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Errorf(format string, v ...interface{}) {
|
||||
s.printf(LevelError, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Panicf(format string, v ...interface{}) {
|
||||
s.printf(LevelPanic, format, v...)
|
||||
}
|
||||
|
||||
func (s *standardLogger) Logger() *log.Logger {
|
||||
return s.logger
|
||||
}
|
||||
|
||||
// verboseLogger is an implementation of Logger which includes debug messages.
|
||||
type verboseLogger struct {
|
||||
logger *log.Logger
|
||||
}
|
||||
|
||||
func NewVerboseLogger(w io.Writer) *verboseLogger {
|
||||
logger := log.New(w, "", 0)
|
||||
logger.SetOutput(formatLog{w: w})
|
||||
return &verboseLogger{
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
func (vb *verboseLogger) Printf(format string, v ...interface{}) {
|
||||
vb.logger.Printf(format, v...)
|
||||
}
|
||||
|
||||
func (vb *verboseLogger) Debugf(format string, v ...interface{}) {
|
||||
vb.logger.Printf(format, v...)
|
||||
}
|
||||
|
||||
func (vb *verboseLogger) Logger() *log.Logger {
|
||||
return vb.logger
|
||||
}
|
||||
|
||||
// CaptureLogger is a logger that stores all the print and debug messages
|
||||
// it sees, useful for testing.
|
||||
// CaptureLogger is a test logger that stores all the print and debug messages
|
||||
// it sees.
|
||||
type CaptureLogger struct {
|
||||
Prints []string
|
||||
Debugs []string
|
||||
|
|
@ -112,14 +154,34 @@ func NewCaptureLogger() *CaptureLogger {
|
|||
return &CaptureLogger{}
|
||||
}
|
||||
|
||||
// Printf formats a message and appends it to Prints.
|
||||
// Printf formats a message and appends it to Debugs.
|
||||
func (cl *CaptureLogger) Printf(format string, v ...interface{}) {
|
||||
cl.Prints = append(cl.Prints, fmt.Sprintf(format, v...))
|
||||
cl.Debugs = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelInfo)+format, v...))
|
||||
}
|
||||
|
||||
// Debugf formats a message and appends it to Debugs.
|
||||
func (cl *CaptureLogger) Debugf(format string, v ...interface{}) {
|
||||
cl.Debugs = append(cl.Debugs, fmt.Sprintf(format, v...))
|
||||
cl.Debugs = append(cl.Debugs, fmt.Sprintf(LevelPrefix(LevelDebug)+format, v...))
|
||||
}
|
||||
|
||||
// Infof formats a message and appends it to Prints.
|
||||
func (cl *CaptureLogger) Infof(format string, v ...interface{}) {
|
||||
cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelInfo)+format, v...))
|
||||
}
|
||||
|
||||
// Warnf formats a message and appends it to Prints.
|
||||
func (cl *CaptureLogger) Warnf(format string, v ...interface{}) {
|
||||
cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelWarn)+format, v...))
|
||||
}
|
||||
|
||||
// Errorf formats a message and appends it to Prints.
|
||||
func (cl *CaptureLogger) Errorf(format string, v ...interface{}) {
|
||||
cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelError)+format, v...))
|
||||
}
|
||||
|
||||
// Panicf formats a message and appends it to Prints.
|
||||
func (cl *CaptureLogger) Panicf(format string, v ...interface{}) {
|
||||
cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelPanic)+format, v...))
|
||||
}
|
||||
|
||||
// Logfer is a thing that has only a Logf() method, like for instance,
|
||||
|
|
@ -128,7 +190,7 @@ type Logfer interface {
|
|||
Logf(format string, v ...interface{})
|
||||
}
|
||||
|
||||
// LogfLogger is a logger that wraps something that has a Logf interface
|
||||
// LogfLogger is a test logger that wraps something that has a Logf interface
|
||||
// and makes it act like our logger.
|
||||
type LogfLogger struct {
|
||||
wrapped Logfer
|
||||
|
|
@ -142,6 +204,66 @@ func (ll *LogfLogger) Debugf(format string, v ...interface{}) {
|
|||
ll.wrapped.Logf(format, v...)
|
||||
}
|
||||
|
||||
func (ll *LogfLogger) Infof(format string, v ...interface{}) {
|
||||
ll.wrapped.Logf(format, v...)
|
||||
}
|
||||
|
||||
func (ll *LogfLogger) Warnf(format string, v ...interface{}) {
|
||||
ll.wrapped.Logf(format, v...)
|
||||
}
|
||||
|
||||
func (ll *LogfLogger) Errorf(format string, v ...interface{}) {
|
||||
ll.wrapped.Logf(format, v...)
|
||||
}
|
||||
|
||||
func (ll *LogfLogger) Panicf(format string, v ...interface{}) {
|
||||
ll.wrapped.Logf(format, v...)
|
||||
}
|
||||
|
||||
func NewLogfLogger(l Logfer) *LogfLogger {
|
||||
return &LogfLogger{wrapped: l}
|
||||
}
|
||||
|
||||
// bufferLogger represents a test Logger that holds log messages
|
||||
// in a buffer for review.
|
||||
type bufferLogger struct {
|
||||
buf *bytes.Buffer
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// NewBufferLogger returns a new instance of BufferLogger.
|
||||
func NewBufferLogger() *bufferLogger {
|
||||
return &bufferLogger{
|
||||
buf: &bytes.Buffer{},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *bufferLogger) Printf(format string, v ...interface{}) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
s := fmt.Sprintf(format, v...)
|
||||
_, err := b.buf.WriteString(s)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
||||
func (b *bufferLogger) Debugf(format string, v ...interface{}) {}
|
||||
func (b *bufferLogger) Infof(format string, v ...interface{}) {
|
||||
b.Printf(LevelPrefix(1)+format, v...)
|
||||
}
|
||||
func (b *bufferLogger) Warnf(format string, v ...interface{}) {
|
||||
b.Printf(LevelPrefix(2)+format, v...)
|
||||
}
|
||||
func (b *bufferLogger) Errorf(format string, v ...interface{}) {
|
||||
b.Printf(LevelPrefix(3)+format, v...)
|
||||
}
|
||||
func (b *bufferLogger) Panicf(format string, v ...interface{}) {
|
||||
b.Printf(LevelPrefix(4)+format, v...)
|
||||
}
|
||||
|
||||
func (b *bufferLogger) ReadAll() ([]byte, error) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
return ioutil.ReadAll(b.buf)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -423,7 +423,7 @@ func (s *Server) handleStandard(ctx context.Context, proto Protocol, conn net.Co
|
|||
default:
|
||||
// The message is not supported yet.
|
||||
// Send an error.
|
||||
s.Logger.Printf("unrecognized postgres packet %v", msg)
|
||||
s.Logger.Errorf("unrecognized postgres packet %v", msg)
|
||||
msg, err = encoder.Error(
|
||||
message.NoticeField{
|
||||
Type: message.NoticeFieldSeverity,
|
||||
|
|
|
|||
|
|
@ -108,7 +108,7 @@ func (s *Server) Serve(ctx context.Context, l net.Listener) (err error) {
|
|||
if limit != nil {
|
||||
// Wait for connection limit.
|
||||
if len(limit) == cap(limit) {
|
||||
s.Logger.Printf("postgres connection limit reached")
|
||||
s.Logger.Warnf("postgres connection limit reached")
|
||||
}
|
||||
select {
|
||||
case limit <- struct{}{}:
|
||||
|
|
@ -133,7 +133,7 @@ func (s *Server) Serve(ctx context.Context, l net.Listener) (err error) {
|
|||
}
|
||||
err := s.handle(ctx, conn)
|
||||
if err != nil {
|
||||
s.Logger.Printf("postgres connection terminated with error: %v", err)
|
||||
s.Logger.Errorf("postgres connection terminated with error: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -146,7 +146,7 @@ func (c *prometheusClient) Count(name string, value int64, rate float64) {
|
|||
var err error
|
||||
counter, err = counterVec.GetMetricWith(labels)
|
||||
if err != nil {
|
||||
c.logger.Printf("counterVec.GetMetricWith error: %s", err)
|
||||
c.logger.Errorf("counterVec.GetMetricWith error: %s", err)
|
||||
}
|
||||
}
|
||||
if value == 1 {
|
||||
|
|
@ -195,7 +195,7 @@ func (c *prometheusClient) Gauge(name string, value float64, rate float64) {
|
|||
var err error
|
||||
gauge, err = gaugeVec.GetMetricWith(labels)
|
||||
if err != nil {
|
||||
c.logger.Printf("gaugeVec.GetMetricWith error: %s", err)
|
||||
c.logger.Errorf("gaugeVec.GetMetricWith error: %s", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
|
@ -238,7 +238,7 @@ func (c *prometheusClient) Histogram(name string, value float64, rate float64) {
|
|||
var err error
|
||||
observer, err = summaryVec.GetMetricWith(labels)
|
||||
if err != nil {
|
||||
c.logger.Printf("summaryVec.GetMetricWith error: %s", err)
|
||||
c.logger.Errorf("summaryVec.GetMetricWith error: %s", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
|
@ -247,7 +247,7 @@ func (c *prometheusClient) Histogram(name string, value float64, rate float64) {
|
|||
|
||||
// Set tracks number of unique elements.
|
||||
func (c *prometheusClient) Set(name string, value string, rate float64) {
|
||||
c.logger.Printf("prometheusClient.Set unimplemented: %s=%s", name, value)
|
||||
c.logger.Infof("prometheusClient.Set unimplemented: %s=%s", name, value)
|
||||
}
|
||||
|
||||
// Timing tracks timing information for a metric.
|
||||
|
|
@ -301,7 +301,7 @@ func tagsToLabels(tags []string, logger logger.Logger) (labels prometheus.Labels
|
|||
tagParts := strings.SplitAfterN(tag, ":", 2)
|
||||
if len(tagParts) != 2 {
|
||||
// only process tags in "key:value" form
|
||||
logger.Printf("Error: invalid Prometheus label: %v\n", tag)
|
||||
logger.Errorf("invalid Prometheus label: %v\n", tag)
|
||||
continue
|
||||
}
|
||||
labels[tagParts[0][0:len(tagParts[0])-1]] = tagParts[1]
|
||||
|
|
|
|||
44
server.go
44
server.go
|
|
@ -225,7 +225,7 @@ func OptServerExecutorPoolSize(size int) ServerOption {
|
|||
// OptServerPrimaryTranslateStore has been deprecated.
|
||||
func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.logger.Printf("DEPRECATED: OptServerPrimaryTranslateStore")
|
||||
s.logger.Infof("DEPRECATED: OptServerPrimaryTranslateStore")
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
|
@ -472,13 +472,13 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
}
|
||||
s.holder = NewHolder(path, s.holderConfig)
|
||||
s.holder.Stats.SetLogger(s.logger)
|
||||
s.holder.Logger.Printf("RowCacheOn: %v", s.holderConfig.RowcacheOn)
|
||||
s.holder.Logger.Infof("RowCacheOn: %v", s.holderConfig.RowcacheOn)
|
||||
cwd, err := os.Getwd()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.holder.Logger.Printf("cwd: %v", cwd)
|
||||
s.holder.Logger.Printf("cmd line: %v", strings.Join(os.Args, " "))
|
||||
s.holder.Logger.Infof("cwd: %v", cwd)
|
||||
s.holder.Logger.Infof("cmd line: %v", strings.Join(os.Args, " "))
|
||||
|
||||
s.cluster.Path = path
|
||||
s.cluster.logger = s.logger
|
||||
|
|
@ -517,7 +517,7 @@ func (s *Server) GRPCURI() pnet.URI {
|
|||
// UpAndDown brings the server up minimally and shuts it down
|
||||
// again; basically, it exists for testing holder open and close.
|
||||
func (s *Server) UpAndDown() error {
|
||||
s.logger.Printf("open server. PID %v", os.Getpid())
|
||||
s.logger.Infof("open server. PID %v", os.Getpid())
|
||||
|
||||
// Log startup
|
||||
err := s.holder.logStartup()
|
||||
|
|
@ -540,7 +540,7 @@ func (s *Server) UpAndDown() error {
|
|||
|
||||
// Open opens and initializes the server.
|
||||
func (s *Server) Open() error {
|
||||
s.logger.Printf("open server. PID %v", os.Getpid())
|
||||
s.logger.Infof("open server. PID %v", os.Getpid())
|
||||
|
||||
if s.holder.NeedsSnapshot() {
|
||||
// Start background monitoring.
|
||||
|
|
@ -781,14 +781,14 @@ func (s *Server) SyncData() error {
|
|||
// listens for events indicating the need to reset the translation
|
||||
// sync processes.
|
||||
func (s *Server) monitorResetTranslationSync() {
|
||||
s.logger.Printf("holder translation sync monitor initializing")
|
||||
s.logger.Infof("holder translation sync monitor initializing")
|
||||
for {
|
||||
// Wait for a reset or a close.
|
||||
select {
|
||||
case <-s.closing:
|
||||
return
|
||||
case <-s.resetTranslationSyncCh:
|
||||
s.logger.Printf("holder translation sync beginning")
|
||||
s.logger.Infof("holder translation sync beginning")
|
||||
s.wg.Add(1)
|
||||
go func() {
|
||||
// Obtaining this lock ensures that there is only
|
||||
|
|
@ -798,7 +798,7 @@ func (s *Server) monitorResetTranslationSync() {
|
|||
defer s.syncer.mu.Unlock()
|
||||
defer s.wg.Done()
|
||||
if err := s.syncer.resetTranslationSync(); err != nil {
|
||||
s.logger.Printf("holder translation sync error: err=%s", err)
|
||||
s.logger.Errorf("holder translation sync error: err=%s", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
|
@ -814,7 +814,7 @@ func (s *Server) monitorAntiEntropy() {
|
|||
ticker := time.NewTicker(s.antiEntropyInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
s.logger.Printf("holder sync monitor initializing (%s interval)", s.antiEntropyInterval)
|
||||
s.logger.Infof("holder sync monitor initializing (%s interval)", s.antiEntropyInterval)
|
||||
|
||||
// Initialize syncer with local holder and remote client.
|
||||
for {
|
||||
|
|
@ -842,17 +842,17 @@ func (s *Server) monitorAntiEntropy() {
|
|||
}
|
||||
|
||||
// Sync holders.
|
||||
s.logger.Printf("holder sync beginning")
|
||||
s.logger.Infof("holder sync beginning")
|
||||
s.cluster.muAntiEntropy.Lock()
|
||||
if err := s.syncer.SyncHolder(); err != nil {
|
||||
s.cluster.muAntiEntropy.Unlock()
|
||||
s.logger.Printf("holder sync error: err=%s", err)
|
||||
s.logger.Errorf("holder sync error: err=%s", err)
|
||||
continue
|
||||
}
|
||||
s.cluster.muAntiEntropy.Unlock()
|
||||
|
||||
// Record successful sync in log.
|
||||
s.logger.Printf("holder sync complete")
|
||||
s.logger.Infof("holder sync complete")
|
||||
dif := time.Since(t)
|
||||
s.holder.Stats.Timing(MetricAntiEntropyDurationSeconds, dif, 1.0)
|
||||
|
||||
|
|
@ -1069,7 +1069,7 @@ func (s *Server) handleRemoteStatus(pb Message) {
|
|||
|
||||
err := s.mergeRemoteStatus(pb.(*NodeStatus))
|
||||
if err != nil {
|
||||
s.logger.Printf("merge remote status: %s", err)
|
||||
s.logger.Errorf("merge remote status: %s", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
|
@ -1093,7 +1093,7 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error {
|
|||
// if we don't know about a field locally, log an error because
|
||||
// fields should be created and synced prior to shard creation
|
||||
if f == nil {
|
||||
s.logger.Printf("local field not found: %s/%s", is.Name, fs.Name)
|
||||
s.logger.Errorf("local field not found: %s/%s", is.Name, fs.Name)
|
||||
continue
|
||||
}
|
||||
if err := f.AddRemoteAvailableShards(fs.AvailableShards); err != nil {
|
||||
|
|
@ -1114,10 +1114,10 @@ func (s *Server) IsPrimary() bool {
|
|||
func (s *Server) monitorDiagnostics() {
|
||||
// Do not send more than once a minute
|
||||
if s.diagnosticInterval < time.Minute {
|
||||
s.logger.Printf("diagnostics disabled")
|
||||
s.logger.Infof("diagnostics disabled")
|
||||
return
|
||||
}
|
||||
s.logger.Printf("Pilosa is currently configured to send small diagnostics reports to our team every %v. More information here: https://www.pilosa.com/docs/latest/administration/#diagnostics", s.diagnosticInterval)
|
||||
s.logger.Infof("Pilosa is currently configured to send small diagnostics reports to our team every %v. More information here: https://www.pilosa.com/docs/latest/administration/#diagnostics", s.diagnosticInterval)
|
||||
|
||||
s.diagnostics.Logger = s.logger
|
||||
s.diagnostics.SetVersion(Version)
|
||||
|
|
@ -1141,11 +1141,11 @@ func (s *Server) monitorDiagnostics() {
|
|||
s.diagnostics.EnrichWithSchemaProperties()
|
||||
err = s.diagnostics.CheckVersion()
|
||||
if err != nil {
|
||||
s.logger.Printf("can't check version: %v", err)
|
||||
s.logger.Errorf("can't check version: %v", err)
|
||||
}
|
||||
err = s.diagnostics.Flush()
|
||||
if err != nil {
|
||||
s.logger.Printf("diagnostics error: %s", err)
|
||||
s.logger.Errorf("diagnostics error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1176,7 +1176,7 @@ func (s *Server) monitorRuntime() {
|
|||
|
||||
defer s.gcNotifier.Close()
|
||||
|
||||
s.logger.Printf("runtime stats initializing (%s interval)", s.metricInterval)
|
||||
s.logger.Infof("runtime stats initializing (%s interval)", s.metricInterval)
|
||||
|
||||
for {
|
||||
// Wait for tick or a close.
|
||||
|
|
@ -1245,7 +1245,7 @@ func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time
|
|||
},
|
||||
)
|
||||
if errLocal != nil || errBroadcast != nil {
|
||||
srv.logger.Printf("error(s) while trying to clean up transaction which failed to start, local: %v, broadcast: %v",
|
||||
srv.logger.Errorf("error(s) while trying to clean up transaction which failed to start, local: %v, broadcast: %v",
|
||||
errLocal,
|
||||
errBroadcast,
|
||||
)
|
||||
|
|
@ -1279,7 +1279,7 @@ func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool
|
|||
},
|
||||
)
|
||||
if err != nil {
|
||||
srv.logger.Printf("error broadcasting transaction finish: %v", err)
|
||||
srv.logger.Errorf("error broadcasting transaction finish: %v", err)
|
||||
// TODO retry?
|
||||
}
|
||||
return trns, nil
|
||||
|
|
|
|||
|
|
@ -224,7 +224,7 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ
|
|||
}
|
||||
longQueryTime := h.api.LongQueryTime()
|
||||
if longQueryTime > 0 && durQuery > longQueryTime {
|
||||
h.logger.Printf("GRPC QueryPQL %v %s", durQuery, query.Query)
|
||||
h.logger.Infof("GRPC QueryPQL %v %s", durQuery, query.Query)
|
||||
}
|
||||
|
||||
rslt := resp.Results[0]
|
||||
|
|
@ -271,7 +271,7 @@ func (h *GRPCHandler) QueryPQLUnary(ctx context.Context, req *pb.QueryPQLRequest
|
|||
}
|
||||
longQueryTime := h.api.LongQueryTime()
|
||||
if longQueryTime > 0 && durQuery > longQueryTime {
|
||||
h.logger.Printf("GRPC QueryPQLUnary %v %s", durQuery, query.Query)
|
||||
h.logger.Infof("GRPC QueryPQLUnary %v %s", durQuery, query.Query)
|
||||
}
|
||||
|
||||
rslt := resp.Results[0]
|
||||
|
|
@ -569,7 +569,7 @@ func (h *GRPCHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSe
|
|||
const defaultLimit = 100000
|
||||
|
||||
h.inspectDeprecated.Do(func() {
|
||||
h.logger.Printf("DEPRECATED: Inspect is deprecated, please use Extract() instead.")
|
||||
h.logger.Infof("DEPRECATED: Inspect is deprecated, please use Extract() instead.")
|
||||
})
|
||||
|
||||
index, err := h.api.Index(stream.Context(), req.Index)
|
||||
|
|
@ -1358,7 +1358,7 @@ func OptGRPCServerStats(stats stats.StatsClient) grpcServerOption {
|
|||
}
|
||||
|
||||
func (s *grpcServer) Serve() error {
|
||||
s.logger.Printf("enabled grpc listening on %s", s.ln.Addr())
|
||||
s.logger.Infof("enabled grpc listening on %s", s.ln.Addr())
|
||||
|
||||
// and start...
|
||||
if err := s.grpcServer.Serve(s.ln); err != nil {
|
||||
|
|
|
|||
|
|
@ -82,7 +82,7 @@ func (s *PostgresServer) Start(addr string) error {
|
|||
return errors.Wrap(err, "creating listener")
|
||||
}
|
||||
|
||||
s.logger.Printf("serving postgres wire protocol on %s", l.Addr())
|
||||
s.logger.Infof("serving postgres wire protocol on %s", l.Addr())
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
s.stop = cancel
|
||||
|
|
@ -102,7 +102,7 @@ func (s *PostgresServer) Close() error {
|
|||
}
|
||||
s.stop()
|
||||
|
||||
s.logger.Printf("waiting for postgres connections to shut down")
|
||||
s.logger.Infof("waiting for postgres connections to shut down")
|
||||
|
||||
return s.eg.Wait()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -148,7 +148,6 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer, opts ...CommandOption
|
|||
func (m *Command) Start() (err error) {
|
||||
// Seed random number generator
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
|
||||
// SetupServer
|
||||
err = m.SetupServer()
|
||||
if err != nil {
|
||||
|
|
@ -158,13 +157,13 @@ func (m *Command) Start() (err error) {
|
|||
if runtime.GOOS == "linux" {
|
||||
result, err := ioutil.ReadFile("/proc/sys/vm/max_map_count")
|
||||
if err != nil {
|
||||
m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err)
|
||||
m.logger.Infof("Tried unsuccessfully to check system mmap limit: %v", err)
|
||||
} else {
|
||||
sysMmapLimit, err := strconv.ParseUint(strings.TrimSuffix(string(result), "\n"), 10, 64)
|
||||
if err != nil {
|
||||
m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err)
|
||||
m.logger.Infof("Tried unsuccessfully to check system mmap limit: %v", err)
|
||||
} else if m.Config.MaxMapCount > sysMmapLimit {
|
||||
m.logger.Printf("WARNING: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, sysMmapLimit)
|
||||
m.logger.Warnf("Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, sysMmapLimit)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -177,7 +176,7 @@ func (m *Command) Start() (err error) {
|
|||
// Initialize HTTP.
|
||||
go func() {
|
||||
if err := m.Handler.Serve(); err != nil {
|
||||
m.logger.Printf("handler serve error: %v", err)
|
||||
m.logger.Errorf("handler serve error: %v", err)
|
||||
}
|
||||
}()
|
||||
m.logger.Printf("listening as %s\n", m.listenURI)
|
||||
|
|
@ -185,7 +184,7 @@ func (m *Command) Start() (err error) {
|
|||
// Initialize gRPC.
|
||||
go func() {
|
||||
if err := m.grpcServer.Serve(); err != nil {
|
||||
m.logger.Printf("grpc server error: %v", err)
|
||||
m.logger.Errorf("grpc server error: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
|
|
@ -194,7 +193,7 @@ func (m *Command) Start() (err error) {
|
|||
if m.Config.Postgres.Bind != "" {
|
||||
var tlsConf *tls.Config
|
||||
if m.Config.Postgres.TLS.CertificatePath != "" {
|
||||
conf, err := GetTLSConfig(&m.Config.Postgres.TLS, m.logger.Logger())
|
||||
conf, err := GetTLSConfig(&m.Config.Postgres.TLS, m.logger)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "setting up postgres TLS")
|
||||
}
|
||||
|
|
@ -230,7 +229,7 @@ func (m *Command) UpAndDown() (err error) {
|
|||
go func() {
|
||||
err := m.Handler.Serve()
|
||||
if err != nil {
|
||||
m.logger.Printf("handler serve error: %v", err)
|
||||
m.logger.Errorf("handler serve error: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
|
|
@ -239,7 +238,7 @@ func (m *Command) UpAndDown() (err error) {
|
|||
return errors.Wrap(err, "bringing server up and down")
|
||||
}
|
||||
|
||||
m.logger.Printf("brought up and shut down again")
|
||||
m.logger.Errorf("brought up and shut down again")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
@ -251,13 +250,13 @@ func (m *Command) Wait() error {
|
|||
signal.Notify(c, os.Interrupt, syscall.SIGTERM)
|
||||
select {
|
||||
case sig := <-c:
|
||||
m.logger.Printf("received signal '%s', gracefully shutting down...\n", sig.String())
|
||||
m.logger.Infof("received signal '%s', gracefully shutting down...\n", sig.String())
|
||||
|
||||
// Second signal causes a hard shutdown.
|
||||
go func() { <-c; os.Exit(1) }()
|
||||
return errors.Wrap(m.Close(), "closing command")
|
||||
case <-m.done:
|
||||
m.logger.Printf("server closed externally")
|
||||
m.logger.Infof("server closed externally")
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
|
@ -275,7 +274,7 @@ func (m *Command) SetupServer() error {
|
|||
return errors.Wrap(err, "setting up logger")
|
||||
}
|
||||
|
||||
m.logger.Printf("%s", pilosa.VersionInfo())
|
||||
m.logger.Infof("%s", pilosa.VersionInfo())
|
||||
|
||||
handleTrialDeadline(m.logger)
|
||||
|
||||
|
|
@ -313,7 +312,7 @@ func (m *Command) SetupServer() error {
|
|||
|
||||
// Setup TLS
|
||||
if uri.Scheme == "https" {
|
||||
m.tlsConfig, err = GetTLSConfig(&m.Config.TLS, m.logger.Logger())
|
||||
m.tlsConfig, err = GetTLSConfig(&m.Config.TLS, m.logger)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "get tls config")
|
||||
}
|
||||
|
|
@ -364,13 +363,13 @@ func (m *Command) SetupServer() error {
|
|||
|
||||
// Primary store configuration is handled automatically now.
|
||||
if m.Config.Translation.PrimaryURL != "" {
|
||||
m.logger.Printf("DEPRECATED: The primary-url configuration option is no longer used.")
|
||||
m.logger.Infof("DEPRECATED: The primary-url configuration option is no longer used.")
|
||||
}
|
||||
// Handle renamed and deprecated config parameter
|
||||
longQueryTime := m.Config.LongQueryTime
|
||||
if m.Config.Cluster.LongQueryTime >= 0 {
|
||||
longQueryTime = m.Config.Cluster.LongQueryTime
|
||||
m.logger.Printf("DEPRECATED: Configuration parameter cluster.long-query-time has been renamed to long-query-time")
|
||||
m.logger.Infof("DEPRECATED: Configuration parameter cluster.long-query-time has been renamed to long-query-time")
|
||||
}
|
||||
|
||||
// Use other config parameters to set Etcd parameters which we don't want to
|
||||
|
|
@ -498,14 +497,14 @@ func (m *Command) setupLogger() error {
|
|||
// duplicate stderr onto log file
|
||||
err := m.dup(int(f.Fd()), int(os.Stderr.Fd()))
|
||||
if err != nil {
|
||||
m.logger.Printf("syscall dup: %s\n", err.Error())
|
||||
m.logger.Errorf("syscall dup: %s\n", err.Error())
|
||||
}
|
||||
|
||||
// reopen log file on SIGHUP
|
||||
<-sighup
|
||||
err = f.Reopen()
|
||||
if err != nil {
|
||||
m.logger.Printf("reopen: %s\n", err.Error())
|
||||
m.logger.Infof("reopen: %s\n", err.Error())
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
|
|
|||
|
|
@ -49,12 +49,12 @@ import (
|
|||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"os"
|
||||
"os/signal"
|
||||
"sync"
|
||||
"syscall"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
|
|
@ -65,7 +65,7 @@ type keypairReloader struct {
|
|||
keyPath string
|
||||
}
|
||||
|
||||
func NewKeypairReloader(certPath, keyPath string, logger *log.Logger) (*keypairReloader, error) {
|
||||
func NewKeypairReloader(certPath, keyPath string, logger logger.Logger) (*keypairReloader, error) {
|
||||
result := &keypairReloader{
|
||||
certPath: certPath,
|
||||
keyPath: keyPath,
|
||||
|
|
@ -79,9 +79,9 @@ func NewKeypairReloader(certPath, keyPath string, logger *log.Logger) (*keypairR
|
|||
c := make(chan os.Signal, 1)
|
||||
signal.Notify(c, syscall.SIGHUP)
|
||||
for range c {
|
||||
logger.Printf("Received SIGHUP, reloading TLS certificate and key from %q and %q", certPath, keyPath)
|
||||
logger.Infof("Received SIGHUP, reloading TLS certificate and key from %q and %q", certPath, keyPath)
|
||||
if err := result.maybeReload(); err != nil {
|
||||
logger.Printf("Keeping old TLS certificate because the new one could not be loaded: %v", err)
|
||||
logger.Printf("ERROR: Keeping old TLS certificate because the new one could not be loaded: %v", err)
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
|
@ -115,7 +115,7 @@ func (kpr *keypairReloader) GetClientCertificateFunc() func(*tls.CertificateRequ
|
|||
}
|
||||
}
|
||||
|
||||
func GetTLSConfig(tlsConfig *TLSConfig, logger *log.Logger) (TLSConfig *tls.Config, err error) {
|
||||
func GetTLSConfig(tlsConfig *TLSConfig, logger logger.Logger) (TLSConfig *tls.Config, err error) {
|
||||
if tlsConfig.CertificatePath != "" && tlsConfig.CertificateKeyPath != "" {
|
||||
kpr, err := NewKeypairReloader(tlsConfig.CertificatePath, tlsConfig.CertificateKeyPath, logger)
|
||||
if err != nil {
|
||||
|
|
|
|||
|
|
@ -145,7 +145,7 @@ func (sq *prioritySnapshotQueue) spawnWorkers(w int) {
|
|||
sq.mu.Lock()
|
||||
defer sq.mu.Unlock()
|
||||
if sq.ctx.Err() != nil {
|
||||
sq.logger.Printf("prioritySnapshotQueue worker: already done")
|
||||
sq.logger.Infof("prioritySnapshotQueue worker: already done")
|
||||
return
|
||||
}
|
||||
sq.workerWG.Add(w)
|
||||
|
|
@ -193,8 +193,8 @@ func (sq *prioritySnapshotQueue) process(req snapshotRequest) {
|
|||
if f.snapshotStamp.Before(req.when) {
|
||||
f.snapshotErr = f.snapshot()
|
||||
if f.snapshotErr != nil {
|
||||
fmt.Printf("snapshot error: %v\n", f.snapshotErr)
|
||||
sq.logger.Printf("snapshot error: %v", f.snapshotErr)
|
||||
fmt.Printf("ERROR: snapshot error: %v\n", f.snapshotErr)
|
||||
sq.logger.Errorf("snapshot error: %v", f.snapshotErr)
|
||||
}
|
||||
f.snapshotPending = false
|
||||
f.snapshotCond.Broadcast()
|
||||
|
|
@ -225,7 +225,7 @@ func (sq *prioritySnapshotQueue) Stop() {
|
|||
enqueued := atomic.LoadUint32(&sq.stats.enqueued)
|
||||
skipped := atomic.LoadUint32(&sq.stats.skipped)
|
||||
if skipped > 0 || enqueued > 1 {
|
||||
sq.logger.Printf("snapshot queue: enqueued %d, skipped %d\n", sq.stats.enqueued, sq.stats.skipped)
|
||||
sq.logger.Infof("snapshot queue: enqueued %d, skipped %d\n", sq.stats.enqueued, sq.stats.skipped)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -239,7 +239,7 @@ func (sq *prioritySnapshotQueue) Enqueue(f *fragment) {
|
|||
sq.mu.RLock()
|
||||
defer sq.mu.RUnlock()
|
||||
if sq.normal == nil {
|
||||
sq.logger.Printf("requested snapshot after snapshot queue was closed")
|
||||
sq.logger.Infof("requested snapshot after snapshot queue was closed")
|
||||
return
|
||||
}
|
||||
// we have to set this before enqueing, because it's
|
||||
|
|
@ -286,7 +286,7 @@ func (sq *prioritySnapshotQueue) Immediate(f *fragment) error {
|
|||
// *don't* need this lock anymore so someone else should have it.
|
||||
if sq.urgent == nil {
|
||||
sq.mu.RUnlock()
|
||||
sq.logger.Printf("requested immediate snapshot after snapshot queue was closed")
|
||||
sq.logger.Errorf("requested immediate snapshot after snapshot queue was closed")
|
||||
return errors.New("requested immediate snapshot after snapshot queue was closed")
|
||||
}
|
||||
f.snapshotPending = true
|
||||
|
|
@ -376,7 +376,7 @@ func (sq *prioritySnapshotQueue) computeMaxOpN() {
|
|||
sq.maxOpN--
|
||||
}
|
||||
if prevMaxOpN != sq.maxOpN {
|
||||
sq.logger.Printf("background scan: %d/%d fragments considered have opN %d or higher\n",
|
||||
sq.logger.Infof("background scan: %d/%d fragments considered have opN %d or higher\n",
|
||||
subTotal, total, sq.maxOpN)
|
||||
}
|
||||
break
|
||||
|
|
@ -490,7 +490,7 @@ func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan sna
|
|||
}
|
||||
|
||||
if scanner.hits > 0 {
|
||||
sq.logger.Printf("background scan: %d/%d fragments needed snapshots\n", scanner.hits, scanner.seen)
|
||||
sq.logger.Infof("background scan: %d/%d fragments needed snapshots\n", scanner.hits, scanner.seen)
|
||||
scanner.hits = 0
|
||||
} else {
|
||||
sq.logger.Debugf("background scan: no fragments needed snapshots, waiting\n")
|
||||
|
|
|
|||
|
|
@ -82,7 +82,7 @@ func (c *statsClient) WithTags(tags ...string) stats.StatsClient {
|
|||
// Count tracks the number of times something occurs per second.
|
||||
func (c *statsClient) Count(name string, value int64, rate float64) {
|
||||
if err := c.client.Count(prefix+name, value, c.tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Count error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Count error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -90,35 +90,35 @@ func (c *statsClient) Count(name string, value int64, rate float64) {
|
|||
func (c *statsClient) CountWithCustomTags(name string, value int64, rate float64, t []string) {
|
||||
tags := append(c.tags, t...)
|
||||
if err := c.client.Count(prefix+name, value, tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Count error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Count error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Gauge sets the value of a metric.
|
||||
func (c *statsClient) Gauge(name string, value float64, rate float64) {
|
||||
if err := c.client.Gauge(prefix+name, value, c.tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Gauge error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Gauge error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Histogram tracks statistical distribution of a metric.
|
||||
func (c *statsClient) Histogram(name string, value float64, rate float64) {
|
||||
if err := c.client.Histogram(prefix+name, value, c.tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Histogram error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Histogram error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Set tracks number of unique elements.
|
||||
func (c *statsClient) Set(name string, value string, rate float64) {
|
||||
if err := c.client.Set(prefix+name, value, c.tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Set error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Set error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Timing tracks timing information for a metric.
|
||||
func (c *statsClient) Timing(name string, value time.Duration, rate float64) {
|
||||
if err := c.client.Timing(prefix+name, value, c.tags, rate); err != nil {
|
||||
c.logger.Printf("statsd.StatsClient.Timing error: %s", err)
|
||||
c.logger.Errorf("statsd.StatsClient.Timing error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ import (
|
|||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/api/client"
|
||||
"github.com/pilosa/pilosa/v2/disco"
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
"github.com/pilosa/pilosa/v2/proto"
|
||||
"github.com/pilosa/pilosa/v2/server"
|
||||
"github.com/pilosa/pilosa/v2/storage"
|
||||
|
|
@ -77,7 +78,7 @@ func (c *Cluster) QueryGRPC(t testing.TB, index, query string) *proto.TableRespo
|
|||
t.Fatal("must have at least one node in cluster to query")
|
||||
}
|
||||
|
||||
grpcClient, err := client.NewGRPCClient([]string{fmt.Sprintf("%s:%d", c.GetPrimary().Server.GRPCURI().Host, c.GetPrimary().Server.GRPCURI().Port)}, nil)
|
||||
grpcClient, err := client.NewGRPCClient([]string{fmt.Sprintf("%s:%d", c.GetPrimary().Server.GRPCURI().Host, c.GetPrimary().Server.GRPCURI().Port)}, nil, logger.NopLogger)
|
||||
if err != nil {
|
||||
t.Fatalf("getting GRPC client: %v", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,54 +0,0 @@
|
|||
// Copyright 2017 Pilosa Corp.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// bufferLogger represents a test Logger that holds log messages
|
||||
// in a buffer for review.
|
||||
type bufferLogger struct {
|
||||
buf *bytes.Buffer
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// NewBufferLogger returns a new instance of BufferLogger.
|
||||
func NewBufferLogger() *bufferLogger {
|
||||
return &bufferLogger{
|
||||
buf: &bytes.Buffer{},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *bufferLogger) Printf(format string, v ...interface{}) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
s := fmt.Sprintf(format, v...)
|
||||
_, err := b.buf.WriteString(s)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
||||
func (b *bufferLogger) Debugf(format string, v ...interface{}) {}
|
||||
|
||||
func (b *bufferLogger) ReadAll() ([]byte, error) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
return ioutil.ReadAll(b.buf)
|
||||
}
|
||||
|
|
@ -16,11 +16,11 @@ package opentracing
|
|||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
|
||||
"github.com/opentracing/opentracing-go"
|
||||
"github.com/opentracing/opentracing-go/ext"
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
"github.com/pilosa/pilosa/v2/tracing"
|
||||
)
|
||||
|
||||
|
|
@ -30,11 +30,12 @@ var _ tracing.Tracer = (*Tracer)(nil)
|
|||
// Tracer represents a wrapper for OpenTracing that implements tracing.Tracer.
|
||||
type Tracer struct {
|
||||
tracer opentracing.Tracer
|
||||
logger logger.Logger
|
||||
}
|
||||
|
||||
// NewTracer returns a new instance of Tracer.
|
||||
func NewTracer(tracer opentracing.Tracer) *Tracer {
|
||||
return &Tracer{tracer: tracer}
|
||||
func NewTracer(tracer opentracing.Tracer, logger logger.Logger) *Tracer {
|
||||
return &Tracer{tracer: tracer, logger: logger}
|
||||
}
|
||||
|
||||
// StartSpanFromContext returns a new child span and context from a given context.
|
||||
|
|
@ -55,7 +56,7 @@ func (t *Tracer) InjectHTTPHeaders(r *http.Request) {
|
|||
opentracing.HTTPHeaders,
|
||||
opentracing.HTTPHeadersCarrier(r.Header),
|
||||
); err != nil {
|
||||
log.Printf("opentracing inject error: %s", err)
|
||||
t.logger.Errorf("opentracing inject error: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -175,7 +175,7 @@ func (tm *TransactionManager) finish(id string) (*Transaction, error) {
|
|||
// After removing, check to see if we need to activate an exclusive transaction
|
||||
trnsMap, err := tm.store.List()
|
||||
if err != nil {
|
||||
tm.log().Printf("error listing transactions in Finish: %v", err)
|
||||
tm.log().Errorf("error listing transactions in Finish: %v", err)
|
||||
return trns, nil
|
||||
}
|
||||
|
||||
|
|
@ -188,7 +188,7 @@ func (tm *TransactionManager) finish(id string) (*Transaction, error) {
|
|||
etrans.Active = true
|
||||
etrans.Deadline = time.Now().Add(etrans.Timeout)
|
||||
if err := tm.store.Put(etrans); err != nil {
|
||||
tm.log().Printf("activating exclusive transaction after finishing last transaction: %v", err)
|
||||
tm.log().Errorf("activating exclusive transaction after finishing last transaction: %v", err)
|
||||
return trns, nil
|
||||
}
|
||||
}
|
||||
|
|
@ -262,7 +262,7 @@ func (tm *TransactionManager) checkDeadlines() time.Duration {
|
|||
|
||||
trnsMap, err := tm.store.List()
|
||||
if err != nil {
|
||||
tm.log().Printf("transaction deadline checker couldn't list transactions: %v", err)
|
||||
tm.log().Errorf("transaction deadline checker couldn't list transactions: %v", err)
|
||||
return 0
|
||||
}
|
||||
|
||||
|
|
@ -287,9 +287,9 @@ func (tm *TransactionManager) checkDeadlines() time.Duration {
|
|||
if !now.Before(trns.Deadline) {
|
||||
trnsF, err := tm.finish(id)
|
||||
if err != nil {
|
||||
tm.log().Printf("error finishing expired transaction '%s': %+v: %v", id, trnsF, err)
|
||||
tm.log().Errorf("error finishing expired transaction '%s': %+v: %v", id, trnsF, err)
|
||||
} else {
|
||||
tm.log().Printf("cleared expired transaction: %+v", trnsF)
|
||||
tm.log().Infof("cleared expired transaction: %+v", trnsF)
|
||||
}
|
||||
} else {
|
||||
interval := trns.Deadline.Sub(now)
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/logger"
|
||||
"github.com/pilosa/pilosa/v2/test"
|
||||
)
|
||||
|
||||
|
|
@ -32,7 +33,7 @@ func TestTransactionManager(t *testing.T) {
|
|||
store := pilosa.NewInMemTransactionStore()
|
||||
|
||||
tm := pilosa.NewTransactionManager(store)
|
||||
tm.Log = test.NewBufferLogger()
|
||||
tm.Log = logger.NewBufferLogger()
|
||||
ctx := context.Background()
|
||||
|
||||
// can add a non-exclusive transaction
|
||||
|
|
|
|||
16
txfactory.go
16
txfactory.go
|
|
@ -503,7 +503,7 @@ func NewTxFactory(backend string, holderDir string, holder *Holder) (f *TxFactor
|
|||
f.dbPerShard = f.NewDBPerShard(types, holderDir, holder)
|
||||
|
||||
if f.hasRBF() {
|
||||
holder.Logger.Printf("rbf config = %#v", holder.cfg.RBFConfig)
|
||||
holder.Logger.Infof("rbf config = %#v", holder.cfg.RBFConfig)
|
||||
}
|
||||
|
||||
return f, err
|
||||
|
|
@ -1399,7 +1399,7 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) {
|
|||
return nil
|
||||
}
|
||||
|
||||
holder.Logger.Printf("green2blue analysis begins.")
|
||||
holder.Logger.Infof("green2blue analysis begins.")
|
||||
|
||||
blueDest := f.types[0]
|
||||
greenSrc := f.types[1]
|
||||
|
|
@ -1422,12 +1422,12 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) {
|
|||
return errors.Wrap(err, "TxFactory.green2blue f.greenHasData()")
|
||||
}
|
||||
if !blueHasData && !greenHasData {
|
||||
holder.Logger.Printf("no data in blue or green. No migration or verification to do.")
|
||||
holder.Logger.Infof("no data in blue or green. No migration or verification to do.")
|
||||
return nil
|
||||
}
|
||||
// INVAR: blue has data.
|
||||
if !greenHasData {
|
||||
holder.Logger.Printf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc)
|
||||
holder.Logger.Errorf("cannot migrate from green '%v' because it has no data in it.", greenSrc)
|
||||
return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc)
|
||||
}
|
||||
|
||||
|
|
@ -1441,11 +1441,11 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) {
|
|||
action := "verify"
|
||||
if blueHasData {
|
||||
verifyInsteadOfCopy = true
|
||||
defer holder.Logger.Printf("bitmap-backend verification done : %v compared to %v", blueDest, greenSrc)
|
||||
defer holder.Logger.Infof("bitmap-backend verification done : %v compared to %v", blueDest, greenSrc)
|
||||
} else {
|
||||
action = "migrate"
|
||||
holder.Logger.Printf("bitmap-backend migration starting: populating %v from %v with %v threads", blueDest, greenSrc, nGoro)
|
||||
defer holder.Logger.Printf("bitmap-backend migration done : populated %v from %v", blueDest, greenSrc)
|
||||
holder.Logger.Infof("bitmap-backend migration starting: populating %v from %v with %v threads", blueDest, greenSrc, nGoro)
|
||||
defer holder.Logger.Infof("bitmap-backend migration done : populated %v from %v", blueDest, greenSrc)
|
||||
}
|
||||
firstPjobStarted := false
|
||||
|
||||
|
|
@ -1496,7 +1496,7 @@ indexloop:
|
|||
return errors.Wrap(err, fmt.Sprintf("GetDBShard(index='%v', shard='%v')", idx.name, int(shard)))
|
||||
}
|
||||
|
||||
holder.Logger.Printf("%v progress on index '%v' (%v of %v): on shard '%v' (%v of %v) [worker %v]",
|
||||
holder.Logger.Infof("%v progress on index '%v' (%v of %v): on shard '%v' (%v of %v) [worker %v]",
|
||||
action, idx.name, k+1, len(idxs), shard, shnum, len(greenShards), worker)
|
||||
|
||||
if verifyInsteadOfCopy {
|
||||
|
|
|
|||
4
view.go
4
view.go
|
|
@ -361,7 +361,7 @@ func (v *view) notifyIfNewShard(shard uint64) {
|
|||
Shard: shard,
|
||||
})
|
||||
if err != nil {
|
||||
v.holder.Logger.Printf("broadcasting create shard: %v", err)
|
||||
v.holder.Logger.Errorf("broadcasting create shard: %v", err)
|
||||
}
|
||||
close(broadcastChan)
|
||||
}()
|
||||
|
|
@ -409,7 +409,7 @@ func (v *view) deleteFragment(shard uint64) error {
|
|||
return ErrFragmentNotFound
|
||||
}
|
||||
|
||||
v.holder.Logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard)
|
||||
v.holder.Logger.Infof("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard)
|
||||
|
||||
idx := f.holder.Index(v.index)
|
||||
f.Close()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue