diff --git a/api.go b/api.go index 59f4105bb..eeb81f0c7 100644 --- a/api.go +++ b/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 diff --git a/api/client/grpc.go b/api/client/grpc.go index f36f764e4..c9a5ff2bf 100644 --- a/api/client/grpc.go +++ b/api/client/grpc.go @@ -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() diff --git a/cluster.go b/cluster.go index ca95cbe67..d7d54edc9 100644 --- a/cluster.go +++ b/cluster.go @@ -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 } diff --git a/cmd.go b/cmd.go index f131eda49..b9e46a8a8 100644 --- a/cmd.go +++ b/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 } diff --git a/cmd/server.go b/cmd/server.go index 90efc7150..1fe9e8400 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -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") diff --git a/ctl/common.go b/ctl/common.go index 8cd22451c..49e35afc6 100644 --- a/ctl/common.go +++ b/ctl/common.go @@ -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 diff --git a/diagnostics.go b/diagnostics.go index 42ae83338..0f835b9d1 100644 --- a/diagnostics.go +++ b/diagnostics.go @@ -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 diff --git a/executor.go b/executor.go index 593be48dc..8bdc65f1d 100644 --- a/executor.go +++ b/executor.go @@ -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 diff --git a/field.go b/field.go index aebd93477..59570465f 100644 --- a/field.go +++ b/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 } diff --git a/fragment.go b/fragment.go index 565cc8aac..09c11d7ce 100644 --- a/fragment.go +++ b/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 { diff --git a/generation.go b/generation.go index 268585b1d..9c5d66e9f 100644 --- a/generation.go +++ b/generation.go @@ -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, diff --git a/holder.go b/holder.go index 4a81a953e..bc59af08f 100644 --- a/holder.go +++ b/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 } } diff --git a/http/handler.go b/http/handler.go index 05109f71a..e515e3a50 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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) } } diff --git a/logger/logger.go b/logger/logger.go index 46e274e4f..28c829be2 100644 --- a/logger/logger.go +++ b/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) +} diff --git a/pg/protocol.go b/pg/protocol.go index 55ed8b750..f7016ac9f 100644 --- a/pg/protocol.go +++ b/pg/protocol.go @@ -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, diff --git a/pg/server.go b/pg/server.go index 9cde6b374..f53cdda8b 100644 --- a/pg/server.go +++ b/pg/server.go @@ -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) } }() } diff --git a/prometheus/prometheus.go b/prometheus/prometheus.go index 551996fb9..3c322376a 100644 --- a/prometheus/prometheus.go +++ b/prometheus/prometheus.go @@ -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] diff --git a/server.go b/server.go index 018d13eec..6b8bbd34f 100644 --- a/server.go +++ b/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 diff --git a/server/grpc.go b/server/grpc.go index e95952688..d91409b7a 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -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 { diff --git a/server/pg.go b/server/pg.go index 8938cc44b..1db60ecca 100644 --- a/server/pg.go +++ b/server/pg.go @@ -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() } diff --git a/server/server.go b/server/server.go index 58827eca0..429dfc39c 100644 --- a/server/server.go +++ b/server/server.go @@ -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()) } } }() diff --git a/server/tlsconfig.go b/server/tlsconfig.go index 4e5f0cf57..d01c82e31 100644 --- a/server/tlsconfig.go +++ b/server/tlsconfig.go @@ -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 { diff --git a/snapshotqueue.go b/snapshotqueue.go index e4fa67db5..9c29a59cb 100644 --- a/snapshotqueue.go +++ b/snapshotqueue.go @@ -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") diff --git a/statsd/statsd.go b/statsd/statsd.go index ccd7afd39..f859f3b3f 100644 --- a/statsd/statsd.go +++ b/statsd/statsd.go @@ -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) } } diff --git a/test/cluster.go b/test/cluster.go index 8eaef7b6f..ec98f3e0f 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -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) } diff --git a/test/logger.go b/test/logger.go deleted file mode 100644 index ecdfbfc40..000000000 --- a/test/logger.go +++ /dev/null @@ -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) -} diff --git a/tracing/opentracing/opentracing.go b/tracing/opentracing/opentracing.go index 79daf5abc..465f051ef 100644 --- a/tracing/opentracing/opentracing.go +++ b/tracing/opentracing/opentracing.go @@ -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) } } } diff --git a/transaction.go b/transaction.go index ca0d5b39a..99926eb42 100644 --- a/transaction.go +++ b/transaction.go @@ -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) diff --git a/transaction_test.go b/transaction_test.go index 50998cd0d..d994564be 100644 --- a/transaction_test.go +++ b/transaction_test.go @@ -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 diff --git a/txfactory.go b/txfactory.go index ed93a1e41..4be8258b6 100644 --- a/txfactory.go +++ b/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 { diff --git a/view.go b/view.go index 093b0c2e9..732f309cf 100644 --- a/view.go +++ b/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()