Merge pull request #1178 from travisturner/logger

Clean up logger; make it honor --log-path flag.
This commit is contained in:
Travis Turner 2018-04-02 11:48:17 -05:00 committed by GitHub
commit 083b730536
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
28 changed files with 467 additions and 379 deletions

View file

@ -20,9 +20,7 @@ import (
"errors"
"fmt"
"hash/fnv"
"io"
"io/ioutil"
"log"
"math/rand"
"net/http"
"os"
@ -268,8 +266,7 @@ type Cluster struct {
closing chan struct{}
prefect SecurityManager
// The writer for any logging.
LogOutput io.Writer
Logger Logger
//
RemoteClient *http.Client
@ -289,16 +286,11 @@ func NewCluster() *Cluster {
closing: make(chan struct{}),
joining: make(chan struct{}),
LogOutput: os.Stderr,
prefect: &NopSecurityManager{},
Logger: NopLogger,
prefect: &NopSecurityManager{},
}
}
// logger returns a logger for the cluster.
func (c *Cluster) logger() *log.Logger {
return log.New(c.LogOutput, "", log.LstdFlags)
}
// Coordinator returns the coordinator node.
func (c *Cluster) CoordinatorNode() *Node {
return c.nodeByID(c.Coordinator)
@ -359,7 +351,7 @@ func (c *Cluster) UpdateCoordinator(n *Node) bool {
// AddNode adds a node to the Cluster and updates and saves the
// new topology.
func (c *Cluster) AddNode(node *Node) error {
c.logger().Printf("add node %s to cluster on %s", node, c.Node)
c.Logger.Printf("add node %s to cluster on %s", node, c.Node)
// If the node being added is the coordinator, set it for this node.
if node.IsCoordinator {
@ -437,7 +429,7 @@ func (c *Cluster) setState(state string) {
return
}
c.logger().Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID)
c.Logger.Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID)
var doCleanup bool
@ -471,7 +463,7 @@ func (c *Cluster) setState(state string) {
// Clean holder.
if err := cleaner.CleanHolder(); err != nil {
c.logger().Printf("holder clean error: err=%s", err)
c.Logger.Printf("holder clean error: err=%s", err)
}
}
}
@ -487,7 +479,7 @@ func (c *Cluster) SetNodeState(state string) error {
State: state,
}
c.logger().Printf("Sending State %s (%s)", state, c.Coordinator)
c.Logger.Printf("Sending State %s (%s)", state, c.Coordinator)
if err := c.sendTo(c.CoordinatorNode(), ns); err != nil {
return fmt.Errorf("sending node state error: err=%s", err)
}
@ -509,7 +501,7 @@ func (c *Cluster) ReceiveNodeState(nodeID string, state string) error {
}
c.Topology.nodeStates[nodeID] = state
c.logger().Printf("received state %s (%s)", state, nodeID)
c.Logger.Printf("received state %s (%s)", state, nodeID)
// Set cluster state to NORMAL.
if c.haveTopologyAgreement() && c.allNodesReady() {
@ -951,9 +943,9 @@ func (c *Cluster) Open() error {
return fmt.Errorf("sending restart NodeJoin: %v", err)
}
c.logger().Printf("wait for joining to complete")
c.Logger.Printf("wait for joining to complete")
<-c.joining
c.logger().Printf("joining has completed")
c.Logger.Printf("joining has completed")
}
return nil
@ -968,7 +960,7 @@ func (c *Cluster) Close() error {
}
func (c *Cluster) markAsJoined() {
c.logger().Printf("mark node as joined (received coordinator update)")
c.Logger.Printf("mark node as joined (received coordinator update)")
if !c.joined {
c.joined = true
close(c.joining)
@ -1001,9 +993,9 @@ func (c *Cluster) allNodesReady() bool {
func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
j, err := c.generateResizeJob(nodeAction)
if err != nil {
c.logger().Printf("generateResizeJob error: err=%s", err)
c.Logger.Printf("generateResizeJob error: err=%s", err)
if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil {
c.logger().Printf("setStateAndBroadcast error: err=%s", err)
c.Logger.Printf("setStateAndBroadcast error: err=%s", err)
}
return err
}
@ -1017,7 +1009,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
})
// Wait for the ResizeJob to finish or be aborted.
c.logger().Printf("wait for jobResult")
c.Logger.Printf("wait for jobResult")
jobResult := <-j.result
// Make sure j.Run() didn't return an error.
@ -1025,7 +1017,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
return err
}
c.logger().Printf("received jobResult: %s", jobResult)
c.Logger.Printf("received jobResult: %s", jobResult)
switch jobResult {
case ResizeJobStateDone:
if err := c.CompleteCurrentJob(ResizeJobStateDone); err != nil {
@ -1048,7 +1040,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
func (c *Cluster) setStateAndBroadcast(state string) error {
c.SetState(state)
// Broadcast cluster status changes to the cluster.
c.logger().Printf("broadcasting ClusterStatus: %s", state)
c.Logger.Printf("broadcasting ClusterStatus: %s", state)
return c.Broadcaster.SendSync(c.Status())
}
@ -1081,7 +1073,7 @@ func (c *Cluster) listenForJoins() {
case nodeAction := <-c.joiningLeavingNodes:
err := c.handleNodeAction(nodeAction)
if err != nil {
c.logger().Printf("handleNodeAction error: err=%s", err)
c.Logger.Printf("handleNodeAction error: err=%s", err)
continue
}
setNormal = true
@ -1093,7 +1085,7 @@ func (c *Cluster) listenForJoins() {
if setNormal {
// Put the cluster back to state NORMAL and broadcast.
if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil {
c.logger().Printf("setStateAndBroadcast error: err=%s", err)
c.Logger.Printf("setStateAndBroadcast error: err=%s", err)
}
}
@ -1104,7 +1096,7 @@ func (c *Cluster) listenForJoins() {
case nodeAction := <-c.joiningLeavingNodes:
err := c.handleNodeAction(nodeAction)
if err != nil {
c.logger().Printf("handleNodeAction error: err=%s", err)
c.Logger.Printf("handleNodeAction error: err=%s", err)
continue
}
setNormal = true
@ -1117,7 +1109,7 @@ func (c *Cluster) listenForJoins() {
// added/removed. It also saves a reference to the ResizeJob in the `jobs` map
// for future lookup by JobID.
func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*ResizeJob, error) {
c.logger().Printf("generateResizeJob: %v", nodeAction)
c.Logger.Printf("generateResizeJob: %v", nodeAction)
c.mu.Lock()
defer c.mu.Unlock()
@ -1125,7 +1117,7 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*ResizeJob, error) {
if err != nil {
return nil, err
}
c.logger().Printf("generated ResizeJob: %d", j.ID)
c.Logger.Printf("generated ResizeJob: %d", j.ID)
// Save job in jobs map for future reference.
c.jobs[j.ID] = j
@ -1215,14 +1207,14 @@ func (c *Cluster) CompleteCurrentJob(state string) error {
// FollowResizeInstruction is run by any node that receives a ResizeInstruction.
func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) error {
c.logger().Printf("follow resize instruction on %s", c.Node.ID)
c.Logger.Printf("follow resize instruction on %s", c.Node.ID)
// Make sure the cluster status on this node agrees with the Coordinator
// before attempting a resize.
if err := c.MergeClusterStatus(instr.ClusterStatus); err != nil {
return err
}
c.logger().Printf("MergeClusterStatus done, start goroutine")
c.Logger.Printf("MergeClusterStatus done, start goroutine")
// The actual resizing runs in a goroutine because we don't want to block
// the distribution of other ResizeInstructions to the rest of the cluster.
@ -1242,7 +1234,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err
if err := func() error {
// Sync the schema received in the resize instruction.
c.logger().Printf("Holder ApplySchema")
c.Logger.Printf("Holder ApplySchema")
if err := c.Holder.ApplySchema(instr.Schema); err != nil {
return err
}
@ -1252,7 +1244,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err
// Request each source file in ResizeSources.
for _, src := range instr.Sources {
c.logger().Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
c.Logger.Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
srcURI := decodeURI(src.Node.URI)
@ -1275,7 +1267,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err
}
// Stream slice from remote node.
c.logger().Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
c.Logger.Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI)
rd, err := client.RetrieveSliceFromURI(context.Background(), src.Index, src.Frame, src.View, src.Slice, srcURI)
if err != nil {
// For now it is an acceptable error if the fragment is not found
@ -1309,7 +1301,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err
}
if err := c.sendTo(DecodeNode(instr.Coordinator), complete); err != nil {
c.logger().Printf("sending resizeInstructionComplete error: err=%s", err)
c.Logger.Printf("sending resizeInstructionComplete error: err=%s", err)
}
}()
return nil
@ -1363,13 +1355,7 @@ type ResizeJob struct {
mu sync.RWMutex
state string
// The writer for any logging.
LogOutput io.Writer
}
// logger returns a logger for the resize job.
func (j *ResizeJob) logger() *log.Logger {
return log.New(j.LogOutput, "", log.LstdFlags)
Logger Logger
}
// NewResizeJob returns a new instance of ResizeJob.
@ -1397,11 +1383,11 @@ func NewResizeJob(existingNodes []*Node, node *Node, action string) *ResizeJob {
}
return &ResizeJob{
ID: rand.Int63(),
IDs: ids,
action: action,
result: make(chan string),
LogOutput: os.Stderr,
ID: rand.Int63(),
IDs: ids,
action: action,
result: make(chan string),
Logger: NopLogger,
}
}
@ -1425,18 +1411,18 @@ func (j *ResizeJob) setState(state string) {
// Run distributes ResizeInstructions.
func (j *ResizeJob) Run() error {
j.logger().Printf("run ResizeJob")
j.Logger.Printf("run ResizeJob")
// Set job state to RUNNING.
j.SetState(ResizeJobStateRunning)
// Job can be considered done in the case where it doesn't require any action.
if !j.nodesArePending() {
j.logger().Printf("ResizeJob contains no pending tasks; mark as done")
j.Logger.Printf("ResizeJob contains no pending tasks; mark as done")
j.result <- ResizeJobStateDone
return nil
}
j.logger().Printf("distribute tasks for ResizeJob")
j.Logger.Printf("distribute tasks for ResizeJob")
err := j.distributeResizeInstructions()
if err != nil {
j.result <- ResizeJobStateAborted
@ -1466,7 +1452,7 @@ func (j *ResizeJob) nodesArePending() bool {
}
func (j *ResizeJob) distributeResizeInstructions() error {
j.logger().Printf("distributeResizeInstructions for job %d", j.ID)
j.Logger.Printf("distributeResizeInstructions for job %d", j.ID)
// Loop through the ResizeInstructions in ResizeJob and send to each host.
for _, instr := range j.Instructions {
// Because the node may not be in the cluster yet, create
@ -1475,7 +1461,7 @@ func (j *ResizeJob) distributeResizeInstructions() error {
ID: instr.Node.ID,
URI: decodeURI(instr.Node.URI),
}
j.logger().Printf("send resize instructions: %v", instr)
j.Logger.Printf("send resize instructions: %v", instr)
if err := j.Broadcaster.SendTo(node, instr); err != nil {
return err
}
@ -1681,7 +1667,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error {
switch e.Event {
case NodeJoin:
c.logger().Printf("received NodeJoin event: %v", e)
c.Logger.Printf("received NodeJoin event: %v", e)
// Ignore the event if this is not the coordinator.
if !c.IsCoordinator() {
return nil
@ -1701,7 +1687,7 @@ func (c *Cluster) nodeJoin(node *Node) error {
// A host that is not part of the topology can't be added to the STARTING cluster.
if !c.Topology.ContainsID(node.ID) {
err := fmt.Sprintf("host is not in topology: %s", node.ID)
c.logger().Print(err)
c.Logger.Printf("%v", err)
return errors.New(err)
}
@ -1816,7 +1802,7 @@ func (c *Cluster) nodeLeave(node *Node) error {
}
func (c *Cluster) MergeClusterStatus(cs *internal.ClusterStatus) error {
c.logger().Printf("merge cluster status: %v", cs)
c.Logger.Printf("merge cluster status: %v", cs)
// Ignore status updates from self (coordinator).
if c.IsCoordinator() {
return nil

View file

@ -17,7 +17,6 @@ package cmd
import (
"fmt"
"io"
"log"
"os"
"os/signal"
"runtime/pprof"
@ -43,14 +42,14 @@ func NewServeCmd(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
Long: `pilosa server runs Pilosa.
It will load existing data from the configured
directory, and start listening client connections
directory and start listening for client connections
on the configured port.`,
RunE: func(cmd *cobra.Command, args []string) error {
logOutput, err := server.GetLogWriter(Server.Config.LogPath, stderr)
if err != nil {
return err
// Set up the logger.
if err := Server.SetupLogger(); err != nil {
return fmt.Errorf("error setting up the logger: %v", err)
}
logger := log.New(logOutput, "", log.LstdFlags)
logger := Server.Server.Logger
logger.Printf("Pilosa %s, build time %s\n", pilosa.Version, pilosa.BuildTime)
// Start CPU profiling.

View file

@ -133,6 +133,7 @@ type Config struct {
MaxWritesPerRequest int `toml:"max-writes-per-request"`
LogPath string `toml:"log-path"`
Verbose bool `toml:"verbose"`
// TLS
TLS TLSConfig
@ -178,6 +179,7 @@ func NewConfig() *Config {
Bind: ":" + DefaultPort,
MaxWritesPerRequest: DefaultMaxWritesPerRequest,
// LogPath: "",
// Verbose: false,
TLS: TLSConfig{},
}

View file

@ -28,6 +28,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
flags.StringVarP(&srv.Config.Bind, "bind", "b", srv.Config.Bind, "Default URI on which pilosa should listen.")
flags.IntVarP(&srv.Config.MaxWritesPerRequest, "max-writes-per-request", "", srv.Config.MaxWritesPerRequest, "Number of write commands per request.")
flags.StringVar(&srv.Config.LogPath, "log-path", srv.Config.LogPath, "Log path")
flags.BoolVar(&srv.Config.Verbose, "verbose", srv.Config.Verbose, "Enable verbose logging")
// TLS
SetTLSConfig(flags, &srv.Config.TLS.CertificatePath, &srv.Config.TLS.CertificateKeyPath, &srv.Config.TLS.SkipVerify)

View file

@ -18,9 +18,6 @@ import (
"bytes"
"encoding/json"
"fmt"
"io"
"io/ioutil"
"log"
"net/http"
"strconv"
"strings"
@ -52,7 +49,7 @@ type DiagnosticsCollector struct {
client *http.Client
logOutput io.Writer
Logger Logger
server *Server
}
@ -66,7 +63,7 @@ func NewDiagnosticsCollector(host string) *DiagnosticsCollector {
start: time.Now(),
client: &http.Client{Timeout: 10 * time.Second},
metrics: make(map[string]interface{}),
logOutput: ioutil.Discard,
Logger: NopLogger,
}
}
@ -119,7 +116,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.Printf("%s\n", err.Error())
}
return nil
@ -160,20 +157,10 @@ func (d *DiagnosticsCollector) Set(name string, value interface{}) {
d.metrics[name] = value
}
// SetLogger Set the logger output type.
func (d *DiagnosticsCollector) SetLogger(logger io.Writer) {
d.logOutput = logger
}
// logger returns a logger that writes to LogOutput.
func (d *DiagnosticsCollector) logger() *log.Logger {
return log.New(d.logOutput, "", log.LstdFlags)
}
// logErr logs the error and returns true if an error exists
func (d *DiagnosticsCollector) logErr(err error) bool {
if err != nil {
d.logOutput.Write([]byte(err.Error()))
d.Logger.Printf("%v", err)
return true
}
return false

View file

@ -16,7 +16,6 @@ package pilosa
import (
"encoding/json"
"io/ioutil"
"net/http"
"net/http/httptest"
"reflect"
@ -31,7 +30,6 @@ func TestDiagnosticsClient(t *testing.T) {
// Create a new client.
d := NewDiagnosticsCollector(server.URL)
d.SetLogger(ioutil.Discard)
d.Set("gg", 10)
d.Set("ss", "ss")
@ -146,7 +144,6 @@ func BenchmarkDiagnostics(b *testing.B) {
// Create a new client.
d := NewDiagnosticsCollector(server.URL)
d.SetLogger(ioutil.Discard)
prev := runtime.GOMAXPROCS(4)
defer runtime.GOMAXPROCS(prev)

View file

@ -81,6 +81,17 @@ The config file is in the [toml format](https://github.com/toml-lang/toml) and h
log-path = "/path/to/logfile"
```
#### Verbose
* Description: Enable verbose logging.
* Flag: `--verbose`
* Env: `PILOSA_VERBOSE`
* Config:
```toml
verbose = true
```
#### Max Writes Per Request
* Description: Maximum number of mutating commands allowed per request. This includes SetBit, ClearBit, SetRowAttrs, SetColumnAttrs, and SetFieldValue.

View file

@ -26,7 +26,6 @@ import (
"hash"
"io"
"io/ioutil"
"log"
"net/http"
"os"
"sort"
@ -103,8 +102,8 @@ type Fragment struct {
// so that they can be mmapped and heap utilization can be kept low.
MaxOpN int
// Writer used for out-of-band log entries.
LogOutput io.Writer
// Logger used for out-of-band log entries.
Logger Logger
// Row attribute storage.
// This is set by the parent frame unless overridden for testing.
@ -124,8 +123,8 @@ func NewFragment(path, index, frame, view string, slice uint64) *Fragment {
CacheType: DefaultCacheType,
CacheSize: DefaultCacheSize,
LogOutput: ioutil.Discard,
MaxOpN: DefaultFragmentMaxOpN,
Logger: NopLogger,
MaxOpN: DefaultFragmentMaxOpN,
stats: NopStatsClient,
}
@ -273,7 +272,7 @@ func (f *Fragment) openCache() error {
// Unmarshal cache data.
var pb internal.Cache
if err := proto.Unmarshal(buf, &pb); err != nil {
f.logger().Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
f.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
return nil
}
@ -298,13 +297,13 @@ func (f *Fragment) Close() error {
func (f *Fragment) close() error {
// Flush cache if closing gracefully.
if err := f.flushCache(); err != nil {
f.logger().Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path)
f.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path)
return err
}
// Close underlying storage.
if err := f.closeStorage(); err != nil {
f.logger().Printf("fragment: error closing storage: err=%s, path=%s", err, f.path)
f.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path)
return err
}
@ -343,9 +342,6 @@ func (f *Fragment) closeStorage() error {
return nil
}
// logger returns a logger instance for the fragment.nt.
func (f *Fragment) logger() *log.Logger { return log.New(f.LogOutput, "", log.LstdFlags) }
// Row returns a row by ID.
func (f *Fragment) Row(rowID uint64) *Bitmap {
f.mu.Lock()
@ -1385,18 +1381,17 @@ func (f *Fragment) Snapshot() error {
defer f.mu.Unlock()
return f.snapshot()
}
func track(start time.Time, message string, stats StatsClient, logger *log.Logger) {
func track(start time.Time, message string, stats StatsClient, logger Logger) {
elapsed := time.Since(start)
logger.Printf("%s took %s", message, elapsed)
stats.Histogram("snapshot", elapsed.Seconds(), 1.0)
}
func (f *Fragment) snapshot() error {
logger := f.logger()
logger.Printf("fragment: snapshotting %s/%s/%s/%d", f.index, f.frame, f.view, f.slice)
f.Logger.Printf("fragment: snapshotting %s/%s/%s/%d", f.index, f.frame, f.view, f.slice)
completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.frame, f.view, f.slice)
start := time.Now()
defer track(start, completeMessage, f.stats, logger)
defer track(start, completeMessage, f.stats, f.Logger)
// Create a temporary file to snapshot to.
snapshotPath := f.path + SnapshotExt

View file

@ -17,7 +17,6 @@ package pilosa
import (
"errors"
"fmt"
"io"
"io/ioutil"
"os"
"path/filepath"
@ -65,7 +64,7 @@ type Frame struct {
rangeEnabled bool
fields []*Field
LogOutput io.Writer
Logger Logger
}
// NewFrame returns a new instance of frame.
@ -95,7 +94,7 @@ func NewFrame(path, index, name string) (*Frame, error) {
rangeEnabled: DefaultRangeEnabled,
//fields
LogOutput: ioutil.Discard,
Logger: NopLogger,
}, nil
}
@ -624,7 +623,7 @@ func (f *Frame) createViewIfNotExistsBase(name string) (*View, bool, error) {
func (f *Frame) newView(path, name string) *View {
view := NewView(path, f.index, f.name, name, f.cacheSize)
view.cacheType = f.cacheType
view.LogOutput = f.LogOutput
view.Logger = f.Logger
view.RowAttrStore = f.rowAttrStore
view.stats = f.Stats.WithTags(fmt.Sprintf("view:%s", name))
view.broadcaster = f.broadcaster

View file

@ -16,10 +16,8 @@ package gossip
import (
"fmt"
"io"
"io/ioutil"
"log"
"os"
"strconv"
"strings"
"sync"
@ -51,8 +49,10 @@ type GossipMemberSet struct {
statusHandler pilosa.StatusHandler
config *gossipConfig
// The writer for any logging.
LogOutput io.Writer
Logger pilosa.Logger
logger *log.Logger
transport *Transport
}
// Start implements the BroadcastReceiver interface and sets the BroadcastHandler.
@ -140,11 +140,6 @@ func retry(attempts int, sleep time.Duration, fn func() error) (err error) {
return fmt.Errorf("after %d attempts, last error: %s", attempts, err)
}
// logger returns a logger for the GossipMemberSet.
func (g *GossipMemberSet) logger() *log.Logger {
return log.New(g.LogOutput, "", log.LstdFlags)
}
////////////////////////////////////////////////////////////////
type gossipConfig struct {
@ -152,14 +147,61 @@ type gossipConfig struct {
memberlistConfig *memberlist.Config
}
// NewGossipMemberSetWithTransport returns a new instance of GossipMemberSet given a Transport.
func NewGossipMemberSetWithTransport(name string, cfg *pilosa.Config, transport *Transport, server *pilosa.Server) (*GossipMemberSet, error) {
// GossipMemberSetOption describes a functional option for GossipMemberSet.
type GossipMemberSetOption func(*GossipMemberSet) error
// WithTransport is a functional option for providing a transport to NewGossipMemberSet.
func WithTransport(transport *Transport) func(*GossipMemberSet) error {
return func(g *GossipMemberSet) error {
g.transport = transport
return nil
}
}
// WithLogger is a functional option for providing a logger to NewGossipMemberSet.
func WithLogger(logger *log.Logger) func(*GossipMemberSet) error {
return func(g *GossipMemberSet) error {
g.logger = logger
return nil
}
}
// NewGossipMemberSet returns a new instance of GossipMemberSet based on options.
func NewGossipMemberSet(name string, cfg *pilosa.Config, server *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
g := &GossipMemberSet{
LogOutput: server.LogOutput,
Logger: server.Logger,
}
port := transport.Net.GetAutoBindPort()
// options
for _, opt := range options {
if err := opt(g); err != nil {
return nil, err
}
}
if g.transport == nil {
port, err := strconv.Atoi(cfg.Gossip.Port)
if err != nil {
return nil, fmt.Errorf("convert port: %s", err)
}
bindURI, err := pilosa.NewURIFromAddress(cfg.Bind)
if err != nil {
return nil, fmt.Errorf("getting uri from bind address: %s", err)
}
host := bindURI.Host()
// Set up the transport.
transport, err := NewTransport(host, port, g.logger)
if err != nil {
return nil, fmt.Errorf("new tranport: %s", err)
}
g.transport = transport
}
port := g.transport.Net.GetAutoBindPort()
bindURI, err := pilosa.NewURIFromAddress(cfg.Bind)
if err != nil {
@ -177,7 +219,7 @@ func NewGossipMemberSetWithTransport(name string, cfg *pilosa.Config, transport
// memberlist config
conf := memberlist.DefaultWANConfig()
conf.Transport = transport.Net
conf.Transport = g.transport.Net
conf.Name = name
conf.BindAddr = host
conf.BindPort = port
@ -196,6 +238,7 @@ func NewGossipMemberSetWithTransport(name string, cfg *pilosa.Config, transport
conf.Delegate = g
conf.SecretKey = gossipKey
conf.Events = server.Cluster.EventReceiver.(memberlist.EventDelegate)
conf.Logger = g.logger
g.config = &gossipConfig{
memberlistConfig: conf,
@ -207,28 +250,6 @@ func NewGossipMemberSetWithTransport(name string, cfg *pilosa.Config, transport
return g, nil
}
// NewGossipMemberSet returns a new instance of GossipMemberSet given a gossip port.
func NewGossipMemberSet(name string, cfg *pilosa.Config, server *pilosa.Server) (*GossipMemberSet, error) {
port, err := strconv.Atoi(cfg.Gossip.Port)
if err != nil {
return nil, fmt.Errorf("convert port: %s", err)
}
bindURI, err := pilosa.NewURIFromAddress(cfg.Bind)
if err != nil {
return nil, fmt.Errorf("getting uri from bind address: %s", err)
}
host := bindURI.Host()
// Set up the transport.
transport, err := NewTransport(host, port)
if err != nil {
return nil, fmt.Errorf("new tranport: %s", err)
}
return NewGossipMemberSetWithTransport(name, cfg, transport, server)
}
// SendSync implementation of the Broadcaster interface.
func (g *GossipMemberSet) SendSync(pb proto.Message) error {
msg, err := pilosa.MarshalMessage(pb)
@ -276,7 +297,7 @@ func (g *GossipMemberSet) SendAsync(pb proto.Message) error {
func (g *GossipMemberSet) NodeMeta(limit int) []byte {
buf, err := proto.Marshal(pilosa.EncodeNode(g.node))
if err != nil {
g.logger().Printf("marshal message error: %s", err)
g.Logger.Printf("marshal message error: %s", err)
return []byte{}
}
return buf
@ -287,11 +308,11 @@ func (g *GossipMemberSet) NodeMeta(limit int) []byte {
func (g *GossipMemberSet) NotifyMsg(b []byte) {
m, err := pilosa.UnmarshalMessage(b)
if err != nil {
g.logger().Printf("unmarshal message error: %s", err)
g.Logger.Printf("unmarshal message error: %s", err)
return
}
if err := g.handler.ReceiveMessage(m); err != nil {
g.logger().Printf("receive message error: %s", err)
g.Logger.Printf("receive message error: %s", err)
return
}
}
@ -307,14 +328,14 @@ func (g *GossipMemberSet) GetBroadcasts(overhead, limit int) [][]byte {
func (g *GossipMemberSet) LocalState(join bool) []byte {
pb, err := g.statusHandler.LocalStatus()
if err != nil {
g.logger().Printf("error getting local state, err=%s", err)
g.Logger.Printf("error getting local state, err=%s", err)
return []byte{}
}
// Marshal nodestate data to bytes.
buf, err := proto.Marshal(pb)
if err != nil {
g.logger().Printf("error marshalling nodestate data, err=%s", err)
g.Logger.Printf("error marshalling nodestate data, err=%s", err)
return []byte{}
}
return buf
@ -326,12 +347,12 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) {
// Unmarshal nodestate data.
var pb internal.NodeStatus
if err := proto.Unmarshal(buf, &pb); err != nil {
g.logger().Printf("error unmarshalling nodestate data, err=%s", err)
g.Logger.Printf("error unmarshalling nodestate data, err=%s", err)
return
}
err := g.statusHandler.HandleRemoteStatus(&pb)
if err != nil {
g.logger().Printf("merge state error: %s", err)
g.Logger.Printf("merge state error: %s", err)
}
}
@ -344,15 +365,14 @@ type GossipEventReceiver struct {
ch chan memberlist.NodeEvent
eventHandler pilosa.EventHandler
// The writer for any logging.
LogOutput io.Writer
Logger pilosa.Logger
}
// NewGossipEventReceiver returns a new instance of GossipEventReceiver.
func NewGossipEventReceiver(logOutput io.Writer) *GossipEventReceiver {
func NewGossipEventReceiver(logger pilosa.Logger) *GossipEventReceiver {
return &GossipEventReceiver{
ch: make(chan memberlist.NodeEvent, 1),
LogOutput: logOutput,
ch: make(chan memberlist.NodeEvent, 1),
Logger: logger,
}
}
@ -375,11 +395,6 @@ func (g *GossipEventReceiver) Start(h pilosa.EventHandler) error {
return nil
}
// logger returns a logger for the GossipEventReceiver.
func (g *GossipEventReceiver) logger() *log.Logger {
return log.New(g.LogOutput, "", log.LstdFlags)
}
func (g *GossipEventReceiver) listen() {
var nodeEventType pilosa.NodeEventType
for {
@ -407,7 +422,7 @@ func (g *GossipEventReceiver) listen() {
Node: node,
}
if err := g.eventHandler.ReceiveEvent(ne); err != nil {
g.logger().Printf("receive event error: %s", err)
g.Logger.Printf("receive event error: %s", err)
}
}
}
@ -443,12 +458,13 @@ type Transport struct {
// It will dynamically bind to a port if port is 0.
// This is useful for test cases where specifiying a port is not reasonable.
//func NewTransport(host string, port int) (*memberlist.NetTransport, error) {
func NewTransport(host string, port int) (*Transport, error) {
func NewTransport(host string, port int, logger *log.Logger) (*Transport, error) {
// memberlist config
conf := memberlist.DefaultWANConfig()
conf.BindAddr = host
conf.BindPort = port
conf.AdvertisePort = port
conf.Logger = logger
net, err := newTransport(conf)
if err != nil {
@ -469,24 +485,10 @@ func NewTransport(host string, port int) (*Transport, error) {
// newTransport returns a NetTransport based on the memberlist configuration.
// It will dynamically bind to a port if conf.BindPort is 0.
func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) {
if conf.LogOutput != nil && conf.Logger != nil {
return nil, fmt.Errorf("Cannot specify both LogOutput and Logger. Please choose a single log configuration setting.")
}
logDest := conf.LogOutput
if logDest == nil {
logDest = os.Stderr
}
logger := conf.Logger
if logger == nil {
logger = log.New(logDest, "", log.LstdFlags)
}
nc := &memberlist.NetTransportConfig{
BindAddrs: []string{conf.BindAddr},
BindPort: conf.BindPort,
Logger: logger,
Logger: conf.Logger,
}
// See comment below for details about the retry in here.
@ -498,7 +500,7 @@ func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) {
return nt, nil
}
if strings.Contains(err.Error(), "address already in use") {
logger.Printf("[DEBUG] Got bind error: %v", err)
conf.Logger.Printf("[DEBUG] Got bind error: %v", err)
continue
}
}

View file

@ -23,12 +23,10 @@ import (
"fmt"
"io"
"io/ioutil"
"log"
"net/http"
"net/url"
// Imported for its side-effect of registering pprof endpoints with the server.
_ "net/http/pprof"
"os"
"runtime/debug"
"strconv"
"strings"
@ -67,8 +65,7 @@ type Handler struct {
Execute(context context.Context, index string, query *pql.Query, slices []uint64, opt *ExecOptions) ([]interface{}, error)
}
// The writer for any logging.
LogOutput io.Writer
Logger Logger
// Keeps the query argument validators for each handler
validators map[string]*queryValidationSpec
@ -98,8 +95,7 @@ func NewHandler() *Handler {
//BroadcastHandler: NopBroadcastHandler, // TODO: implement the nop
//StatusHandler: NopStatusHandler, // TODO: implement the nop
FileSystem: NopFileSystem,
LogOutput: os.Stderr,
Logger: NopLogger,
}
BuildRouters(handler)
handler.populateValidators()
@ -247,7 +243,7 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
stack := debug.Stack()
msg := "PANIC: %s\n%s"
fmt.Fprintf(h.LogOutput, msg, err, stack)
h.Logger.Printf(msg, err, stack)
fmt.Fprintf(w, msg, err, stack)
}
}()
@ -261,7 +257,7 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
statsTags := make([]string, 0, 3)
if h.Cluster.LongQueryTime > 0 && dif > h.Cluster.LongQueryTime {
h.logger().Printf("%s %s %v", r.Method, r.URL.String(), dif)
h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif)
statsTags = append(statsTags, "slow_query")
}
@ -289,7 +285,7 @@ func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) {
filesystem, err := h.FileSystem.New()
if err != nil {
h.writeQueryResponse(w, r, &QueryResponse{Err: err})
h.logger().Println("Pilosa WebUI is not available. Please run `make generate-statik` before building Pilosa with `make install`.")
h.Logger.Printf("Pilosa WebUI is not available. Please run `make generate-statik` before building Pilosa with `make install`.")
return
}
http.FileServer(filesystem).ServeHTTP(w, r)
@ -300,7 +296,7 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
if err := json.NewEncoder(w).Encode(getSchemaResponse{
Indexes: h.Holder.Schema(),
}); err != nil {
h.logger().Printf("write schema response error: %s", err)
h.Logger.Printf("write schema response error: %s", err)
}
}
@ -308,7 +304,7 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
pb, err := h.StatusHandler.ClusterStatus()
if err != nil {
h.logger().Printf("cluster status error: %s", err)
h.Logger.Printf("cluster status error: %s", err)
return
}
@ -317,7 +313,7 @@ func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
State: cs.State,
Nodes: DecodeNodes(cs.Nodes),
}); err != nil {
h.logger().Printf("write status response error: %s", err)
h.Logger.Printf("write status response error: %s", err)
}
}
@ -395,7 +391,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.Printf("write query response error: %s", err)
}
}
@ -405,7 +401,7 @@ func (h *Handler) handleGetSlicesMax(w http.ResponseWriter, r *http.Request) {
Standard: h.Holder.MaxSlices(),
Inverse: h.Holder.MaxInverseSlices(),
}); err != nil {
h.logger().Printf("write slices-max response error: %s", err)
h.Logger.Printf("write slices-max response error: %s", err)
}
}
@ -431,7 +427,7 @@ func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) {
if err := json.NewEncoder(w).Encode(getIndexResponse{
map[string]string{"name": index.Name()},
}); err != nil {
h.logger().Printf("write response error: %s", err)
h.Logger.Printf("write response error: %s", err)
}
}
@ -519,12 +515,12 @@ func (h *Handler) handleDeleteIndex(w http.ResponseWriter, r *http.Request) {
Index: indexName,
})
if err != nil {
h.logger().Printf("problem sending DeleteIndex message: %s", err)
h.Logger.Printf("problem sending DeleteIndex message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
h.Holder.Stats.Count("deleteIndex", 1, 1.0)
@ -564,14 +560,14 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) {
Meta: req.Options.Encode(),
})
if err != nil {
h.logger().Printf("problem sending CreateIndex message: %s", err)
h.Logger.Printf("problem sending CreateIndex message: %s", err)
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Encode response.
if err := json.NewEncoder(w).Encode(postIndexResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
h.Holder.Stats.Count("createIndex", 1, 1.0)
@ -610,7 +606,7 @@ func (h *Handler) handlePatchIndexTimeQuantum(w http.ResponseWriter, r *http.Req
// Encode response.
if err := json.NewEncoder(w).Encode(patchIndexTimeQuantumResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -665,7 +661,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.Printf("response encoding error: %s", err)
}
}
@ -718,12 +714,12 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) {
Meta: req.Options.Encode(),
})
if err != nil {
h.logger().Printf("problem sending CreateFrame message: %s", err)
h.Logger.Printf("problem sending CreateFrame message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(postFrameResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
h.Holder.Stats.CountWithCustomTags("createFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)})
@ -783,7 +779,7 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) {
index := h.Holder.Index(indexName)
if index == nil {
if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
return
}
@ -801,12 +797,12 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) {
Frame: frameName,
})
if err != nil {
h.logger().Printf("problem sending DeleteFrame message: %s", err)
h.Logger.Printf("problem sending DeleteFrame message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(deleteFrameResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
h.Holder.Stats.CountWithCustomTags("deleteFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)})
@ -848,7 +844,7 @@ func (h *Handler) handlePatchFrameTimeQuantum(w http.ResponseWriter, r *http.Req
// Encode response.
if err := json.NewEncoder(w).Encode(patchFrameTimeQuantumResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -899,12 +895,12 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) {
Field: encodeField(field),
})
if err != nil {
h.logger().Printf("problem sending CreateField message: %s", err)
h.Logger.Printf("problem sending CreateField message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(postFrameFieldResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -943,12 +939,12 @@ func (h *Handler) handleDeleteFrameField(w http.ResponseWriter, r *http.Request)
Field: fieldName,
})
if err != nil {
h.logger().Printf("problem sending DeleteField message: %s", err)
h.Logger.Printf("problem sending DeleteField message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(deleteFrameFieldResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -979,7 +975,7 @@ func (h *Handler) handleGetFrameFields(w http.ResponseWriter, r *http.Request) {
// Encode response.
if err := json.NewEncoder(w).Encode(getFrameFieldsResponse{Fields: fields}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -1012,7 +1008,7 @@ func (h *Handler) handleGetFrameViews(w http.ResponseWriter, r *http.Request) {
// Encode response.
if err := json.NewEncoder(w).Encode(getFrameViewsResponse{Views: names}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -1046,12 +1042,12 @@ func (h *Handler) handleDeleteView(w http.ResponseWriter, r *http.Request) {
View: viewName,
})
if err != nil {
h.logger().Printf("problem sending DeleteView message: %s", err)
h.Logger.Printf("problem sending DeleteView message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(deleteViewResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -1107,7 +1103,7 @@ func (h *Handler) handlePostFrameAttrDiff(w http.ResponseWriter, r *http.Request
if err := json.NewEncoder(w).Encode(postFrameAttrDiffResponse{
Attrs: attrs,
}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -1276,10 +1272,10 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
}
// Find the Index.
h.logger().Println("importing:", req.Index, req.Frame, req.Slice)
h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice)
index := h.Holder.Index(req.Index)
if index == nil {
h.logger().Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error())
h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error())
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
@ -1287,7 +1283,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
// Retrieve frame.
f := index.Frame(req.Frame)
if f == nil {
h.logger().Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error())
h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error())
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
@ -1295,7 +1291,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
// Import into fragment.
err = f.Import(req.RowIDs, req.ColumnIDs, timestamps)
if err != nil {
h.logger().Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err)
h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err)
return
}
@ -1346,10 +1342,10 @@ func (h *Handler) handlePostImportValue(w http.ResponseWriter, r *http.Request)
}
// Find the Index.
h.logger().Println("importing:", req.Index, req.Frame, req.Slice)
h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice)
index := h.Holder.Index(req.Index)
if index == nil {
h.logger().Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error())
h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error())
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
@ -1357,7 +1353,7 @@ func (h *Handler) handlePostImportValue(w http.ResponseWriter, r *http.Request)
// Retrieve frame.
f := index.Frame(req.Frame)
if f == nil {
h.logger().Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error())
h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error())
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
@ -1365,7 +1361,7 @@ func (h *Handler) handlePostImportValue(w http.ResponseWriter, r *http.Request)
// Import into fragment.
err = f.ImportValue(req.Field, req.ColumnIDs, req.Values)
if err != nil {
h.logger().Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err)
h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err)
return
}
@ -1452,7 +1448,7 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request)
// Write to response.
if err := json.NewEncoder(w).Encode(nodes); err != nil {
h.logger().Printf("json write error: %s", err)
h.Logger.Printf("json write error: %s", err)
}
}
@ -1475,7 +1471,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("fragment backup error: %s", err)
h.Logger.Printf("fragment backup error: %s", err)
}
}
@ -1545,7 +1541,7 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ
// Encode response.
buf, err := proto.Marshal(&resp)
if err != nil {
h.logger().Printf("merge block response encoding error: %s", err)
h.Logger.Printf("merge block response encoding error: %s", err)
return
}
@ -1579,7 +1575,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.Printf("block response encoding error: %s", err)
}
}
@ -1680,7 +1676,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request)
// handleGetHosts handles /hosts requests.
func (h *Handler) handleGetHosts(w http.ResponseWriter, r *http.Request) {
if err := json.NewEncoder(w).Encode(h.Cluster.Nodes); err != nil {
h.logger().Printf("write version response error: %s", err)
h.Logger.Printf("write version response error: %s", err)
}
}
@ -1696,7 +1692,7 @@ func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) {
}{
Version: version,
}); err != nil {
h.logger().Printf("write version response error: %s", err)
h.Logger.Printf("write version response error: %s", err)
}
}
@ -1716,11 +1712,6 @@ func (h *Handler) handleExpvar(w http.ResponseWriter, r *http.Request) {
fmt.Fprintf(w, "\n}\n")
}
// logger returns a logger for the handler.
func (h *Handler) logger() *log.Logger {
return log.New(h.LogOutput, "", log.LstdFlags)
}
// QueryResult types.
const (
QueryResultTypeNil uint32 = iota
@ -1908,11 +1899,11 @@ func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Reque
Definition: def,
})
if err != nil {
h.logger().Printf("problem sending CreateInputDefinition message: %s", err)
h.Logger.Printf("problem sending CreateInputDefinition message: %s", err)
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -1938,7 +1929,7 @@ func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Reques
Frames: inputDef.frames,
Fields: inputDef.fields,
}); err != nil {
h.logger().Printf("write status response error: %s", err)
h.Logger.Printf("write status response error: %s", err)
}
}
@ -1967,11 +1958,11 @@ func (h *Handler) handleDeleteInputDefinition(w http.ResponseWriter, r *http.Req
Name: inputDefName,
})
if err != nil {
h.logger().Printf("problem sending DeleteInputDefinition message: %s", err)
h.Logger.Printf("problem sending DeleteInputDefinition message: %s", err)
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -2013,7 +2004,7 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) {
}
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -2061,7 +2052,7 @@ func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r
Old: oldNode,
New: newNode,
}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}
@ -2101,7 +2092,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.Printf("response encoding error: %s", err)
}
}
@ -2134,7 +2125,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.Printf("response encoding error: %s", err)
}
}
@ -2271,7 +2262,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
}
if err := json.NewEncoder(w).Encode(defaultClusterMessageResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
h.Logger.Printf("response encoding error: %s", err)
}
}

View file

@ -37,13 +37,13 @@ import (
func TestHandlerPanics(t *testing.T) {
h := test.NewHandler()
buf := &bytes.Buffer{}
h.Handler.LogOutput = buf
bufLogger := test.NewBufferLogger()
h.Handler.Logger = bufLogger
w := httptest.NewRecorder()
// will panic since Handler has no Holder set up
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/index/taxi", nil))
bufbytes, err := ioutil.ReadAll(buf)
bufbytes, err := bufLogger.ReadAll()
if err != nil {
t.Fatalf("reading all logoutput: %v", err)
}

View file

@ -18,9 +18,7 @@ import (
"context"
"errors"
"fmt"
"io"
"io/ioutil"
"log"
"net/http"
"os"
"path"
@ -71,7 +69,7 @@ type Holder struct {
// The interval at which the cached row ids are persisted to disk.
CacheFlushInterval time.Duration
LogOutput io.Writer
Logger Logger
}
// NewHolder returns a new instance of Holder.
@ -89,7 +87,7 @@ func NewHolder() *Holder {
CacheFlushInterval: DefaultCacheFlushInterval,
LogOutput: os.Stderr,
Logger: NopLogger,
}
}
@ -97,7 +95,7 @@ func NewHolder() *Holder {
// without actually loading any data into memory.
// HasData is returned, and h.hasData is set.
func (h *Holder) Peek() bool {
h.logger().Printf("peek at holder path: %s", h.Path)
h.Logger.Printf("peek at holder path: %s", h.Path)
h.hasData = false
// Open path to read all index directories.
@ -127,7 +125,7 @@ func (h *Holder) Peek() bool {
func (h *Holder) Open() error {
h.setFileLimit()
h.logger().Printf("open holder path: %s", h.Path)
h.Logger.Printf("open holder path: %s", h.Path)
if err := os.MkdirAll(h.Path, 0777); err != nil {
return err
}
@ -149,18 +147,18 @@ func (h *Holder) Open() error {
continue
}
h.logger().Printf("opening index: %s", filepath.Base(fi.Name()))
h.Logger.Printf("opening index: %s", filepath.Base(fi.Name()))
index, err := h.newIndex(h.IndexPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
if err == ErrName {
h.logger().Printf("ERROR opening index: %s, err=%s", fi.Name(), err)
h.Logger.Printf("ERROR opening index: %s, err=%s", fi.Name(), err)
continue
} else if err != nil {
return err
}
if err := index.Open(); err != nil {
if err == ErrName {
h.logger().Printf("ERROR opening index: %s, err=%s", index.Name(), err)
h.Logger.Printf("ERROR opening index: %s, err=%s", index.Name(), err)
continue
}
return fmt.Errorf("open index: name=%s, err=%s", index.Name(), err)
@ -169,7 +167,7 @@ func (h *Holder) Open() error {
h.indexes[index.Name()] = index
h.mu.Unlock()
}
h.logger().Printf("open holder: complete")
h.Logger.Printf("open holder: complete")
// Periodically flush cache.
h.wg.Add(1)
@ -374,7 +372,7 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
if err != nil {
return nil, err
}
index.LogOutput = h.LogOutput
index.Logger = h.Logger
index.Stats = h.Stats.WithTags(fmt.Sprintf("index:%s", index.Name()))
index.broadcaster = h.Broadcaster
index.NewAttrStore = h.NewAttrStore
@ -464,7 +462,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.Printf("error flushing cache: err=%s, path=%s", err, fragment.CachePath())
}
}
}
@ -488,7 +486,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.Printf("ERROR checking open file limit: %s", err)
return
}
// If the soft limit is lower than the FileLimit constant, we will try to change it.
@ -512,32 +510,30 @@ 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.Printf("ERROR 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.Printf("ERROR 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.Printf("ERROR 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/administration/#open-file-limits for more information.", FileLimit, 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/administration/#open-file-limits for more information.", FileLimit, oldLimit.Cur, FileLimit)
}
}
}
}
func (h *Holder) logger() *log.Logger { return log.New(h.LogOutput, "", log.LstdFlags) }
func (h *Holder) loadNodeID() (string, error) {
idPath := path.Join(h.Path, "ID")
nodeID := ""
h.logger().Printf("load NodeID: %s", idPath)
h.Logger.Printf("load NodeID: %s", idPath)
if err := os.MkdirAll(h.Path, 0777); err != nil {
return "", err
}

View file

@ -15,6 +15,7 @@
package pilosa_test
import (
"bytes"
"context"
"os"
"path/filepath"
@ -30,6 +31,10 @@ import (
func TestHolder_Open(t *testing.T) {
t.Run("ErrIndexName", func(t *testing.T) {
h := test.MustOpenHolder()
bufLogger := test.NewBufferLogger()
h.Holder.Logger = bufLogger
defer h.Close()
if err := os.Mkdir(h.IndexPath("!"), 0777); err != nil {
@ -39,8 +44,12 @@ func TestHolder_Open(t *testing.T) {
}
if err := h.Reopen(); err != nil {
t.Fatal(err)
} else if logOutput := h.LogOutput.String(); !strings.Contains(logOutput, `ERROR opening index: !`) {
t.Fatalf("expected log error:\n%s", logOutput)
}
if bufbytes, err := bufLogger.ReadAll(); err != nil {
t.Fatal(err)
} else if !bytes.Contains(bufbytes, []byte("ERROR opening index: !")) {
t.Fatalf("expected log error:\n%s", bufbytes)
}
})

View file

@ -17,7 +17,6 @@ package pilosa
import (
"errors"
"fmt"
"io"
"io/ioutil"
"os"
"path/filepath"
@ -66,7 +65,7 @@ type Index struct {
broadcaster Broadcaster
Stats StatsClient
LogOutput io.Writer
Logger Logger
}
// NewIndex returns a new instance of Index.
@ -92,7 +91,7 @@ func NewIndex(path, name string) (*Index, error) {
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
LogOutput: ioutil.Discard,
Logger: NopLogger,
}, nil
}
@ -528,7 +527,7 @@ func (i *Index) newFrame(path, name string) (*Frame, error) {
if err != nil {
return nil, err
}
f.LogOutput = i.LogOutput
f.Logger = i.Logger
f.Stats = i.Stats.WithTags(fmt.Sprintf("frame:%s", name))
f.broadcaster = i.broadcaster
f.rowAttrStore = i.NewAttrStore(filepath.Join(f.path, ".data"))

88
logger.go Normal file
View file

@ -0,0 +1,88 @@
// 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 pilosa
import (
"io"
"log"
)
// Ensure nopLogger implements interface.
var _ Logger = &nopLogger{}
// Logger represents an interface for a shared logger.
type Logger interface {
Printf(format string, v ...interface{})
Debugf(format string, v ...interface{})
}
func init() {
NopLogger = &nopLogger{}
}
// NopLogger represents a Logger that doesn't do anything.
var NopLogger Logger
type nopLogger struct{}
// Printf is a no-op implementation of the Logger Printf method.
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{}) {}
// StandardLogger is a basic implementation of pilosa.Logger based on log.Logger.
type StandardLogger struct {
logger *log.Logger
}
func NewStandardLogger(w io.Writer) *StandardLogger {
return &StandardLogger{
logger: log.New(w, "", log.LstdFlags),
}
}
func (s *StandardLogger) Printf(format string, v ...interface{}) {
s.logger.Printf(format, v...)
}
func (s *StandardLogger) Debugf(format string, v ...interface{}) {}
func (s *StandardLogger) Logger() *log.Logger {
return s.logger
}
// VerboseLogger is an implementation of pilosa.Logger which includes debug messages.
type VerboseLogger struct {
logger *log.Logger
}
func NewVerboseLogger(w io.Writer) *VerboseLogger {
return &VerboseLogger{
logger: log.New(w, "", log.LstdFlags),
}
}
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
}

View file

@ -19,8 +19,6 @@ import (
"crypto/tls"
"errors"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
@ -87,8 +85,7 @@ type Server struct {
// Misc options.
MaxWritesPerRequest int
LogOutput io.Writer
logger *log.Logger
Logger Logger
defaultClient InternalClient
}
@ -115,9 +112,8 @@ func NewServer() *Server {
MetricInterval: 0,
DiagnosticInterval: 0,
LogOutput: os.Stderr,
Logger: NopLogger,
}
s.logger = log.New(s.LogOutput, "", log.LstdFlags)
s.Handler.Holder = s.Holder
s.diagnostics.server = s
@ -126,7 +122,7 @@ func NewServer() *Server {
// Open opens and initializes the server.
func (s *Server) Open() error {
s.Logger().Printf("open server")
s.Logger.Printf("open server")
// s.ln can be configured prior to Open() via s.OpenListener().
if s.ln == nil {
if err := s.OpenListener(); err != nil {
@ -151,7 +147,6 @@ func (s *Server) Open() error {
// Peek at the holder to determine if there is data on disk.
// Don't actually load the data until after the Cluster
// management starts.
s.Holder.LogOutput = s.LogOutput
s.Holder.Peek()
// Create default HTTP client
@ -175,7 +170,6 @@ func (s *Server) Open() error {
s.Handler.Node = node
s.Handler.Cluster = s.Cluster
s.Handler.Executor = e
s.Handler.LogOutput = s.LogOutput
s.Cluster.prefect = s.Handler
@ -186,7 +180,7 @@ func (s *Server) Open() error {
go func() {
err := http.Serve(s.ln, s.Handler)
if err != nil {
s.Logger().Printf("HTTP handler terminated with error: %s\n", err)
s.Logger.Printf("HTTP handler terminated with error: %s\n", err)
}
}()
@ -226,7 +220,7 @@ func (s *Server) Open() error {
// OpenListener opens a listener for the Server.
func (s *Server) OpenListener() error {
s.Logger().Printf("open server listener: %s", s.URI)
s.Logger.Printf("open server listener: %s", s.URI)
if s.ln != nil {
return fmt.Errorf("a listener already exists for server: %s", s.URI)
}
@ -288,7 +282,7 @@ func (s *Server) LoadNodeID() string {
}
nodeID, err := s.Holder.loadNodeID()
if err != nil {
s.Logger().Printf("loading NodeID: %v", err)
s.Logger.Printf("loading NodeID: %v", err)
return s.NodeID
}
return nodeID
@ -321,14 +315,11 @@ func GetHTTPClient(t *tls.Config) *http.Client {
return &http.Client{Transport: transport}
}
// Logger returns a logger that writes to LogOutput
func (s *Server) Logger() *log.Logger { return s.logger }
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.Printf("holder sync monitor initializing (%s interval)", s.AntiEntropyInterval)
for {
// Wait for tick or a close.
@ -339,7 +330,7 @@ func (s *Server) monitorAntiEntropy() {
s.Holder.Stats.Count("AntiEntropy", 1, 1.0)
}
t := time.Now()
s.Logger().Printf("holder sync beginning")
s.Logger.Printf("holder sync beginning")
// Initialize syncer with local holder and remote client.
var syncer HolderSyncer
@ -352,12 +343,12 @@ func (s *Server) monitorAntiEntropy() {
// Sync holders.
if err := syncer.SyncHolder(); err != nil {
s.Logger().Printf("holder sync error: err=%s", err)
s.Logger.Printf("holder sync error: err=%s", err)
continue
}
// Record successful sync in log.
s.Logger().Printf("holder sync complete")
s.Logger.Printf("holder sync complete")
dif := time.Since(t)
s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0)
}
@ -482,7 +473,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
func (s *Server) SendSync(pb proto.Message) error {
var eg errgroup.Group
for _, node := range s.Cluster.Nodes {
s.Logger().Printf("SendSync to: %s", node.URI)
s.Logger.Printf("SendSync to: %s", node.URI)
// Don't forward the message to ourselves.
if s.URI == node.URI {
continue
@ -504,7 +495,7 @@ func (s *Server) SendAsync(pb proto.Message) error {
// SendTo represents an implementation of Broadcaster.
func (s *Server) SendTo(to *Node, pb proto.Message) error {
s.Logger().Printf("SendTo: %s", to.URI)
s.Logger.Printf("SendTo: %s", to.URI)
ctx := context.WithValue(context.Background(), "uri", &to.URI)
return s.defaultClient.SendMessage(ctx, pb)
}
@ -554,7 +545,7 @@ func (s *Server) HandleRemoteStatus(pb proto.Message) error {
err := s.mergeRemoteStatus(pb.(*internal.NodeStatus))
if err != nil {
s.Logger().Printf("merge remote status: %s", err)
s.Logger.Printf("merge remote status: %s", err)
}
}()
@ -579,7 +570,7 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error {
// if we don't know about an index locally, log an error because
// indexes should be created and synced prior to slice creation
if localIndex == nil {
s.Logger().Printf("Local Index not found: %s", index)
s.Logger.Printf("Local Index not found: %s", index)
continue
}
if newMax > oldmaxslices[index] {
@ -595,7 +586,7 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error {
// if we don't know about an index locally, log an error because
// indexes should be created and synced prior to slice creation
if localIndex == nil {
s.Logger().Printf("Local Index not found: %s", index)
s.Logger.Printf("Local Index not found: %s", index)
continue
}
if newMaxInverse > oldMaxInverseSlices[index] {
@ -611,13 +602,13 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error {
func (s *Server) monitorDiagnostics() {
// Do not send more than once a minute
if s.DiagnosticInterval < time.Minute {
s.Logger().Printf("diagnostics disabled")
s.Logger.Printf("diagnostics disabled")
return
} else {
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.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.diagnostics.SetLogger(s.LogOutput)
s.diagnostics.Logger = s.Logger
s.diagnostics.SetVersion(Version)
s.diagnostics.Set("Host", s.URI.host)
s.diagnostics.Set("Cluster", strings.Join(s.Cluster.NodeIDs(), ","))
@ -639,7 +630,7 @@ func (s *Server) monitorDiagnostics() {
s.diagnostics.CheckVersion()
err = s.diagnostics.Flush()
if err != nil {
s.Logger().Printf("Diagnostics error: %s", err)
s.Logger.Printf("Diagnostics error: %s", err)
}
}
@ -670,7 +661,7 @@ func (s *Server) monitorRuntime() {
defer s.GCNotifier.Close()
s.Logger().Printf("runtime stats initializing (%s interval)", s.MetricInterval)
s.Logger.Printf("runtime stats initializing (%s interval)", s.MetricInterval)
for {
// Wait for tick or a close.

View file

@ -54,7 +54,7 @@ func TestMain_SendReceiveMessage(t *testing.T) {
m0.Server.Cluster.Coordinator = m0.Server.NodeID
m0.Server.Cluster.Topology = &pilosa.Topology{NodeIDs: []string{m0.Server.NodeID, m1.Server.NodeID}}
m0.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m0.Server.LogOutput)
m0.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m0.Server.Logger)
gossipMemberSet0, err := gossip.NewGossipMemberSet(m0.Server.URI.HostPort(), m0.Config, m0.Server)
if err != nil {
t.Fatal(err)
@ -81,7 +81,7 @@ func TestMain_SendReceiveMessage(t *testing.T) {
m1.Config.Gossip.Seeds = []string{gossipMemberSet0.GetBindAddr()}
m1.Server.Cluster.Coordinator = m0.Server.NodeID
m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m1.Server.LogOutput)
m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m1.Server.Logger)
gossipMemberSet1, err := gossip.NewGossipMemberSet(m1.Server.URI.HostPort(), m1.Config, m1.Server)
if err != nil {
t.Fatal(err)

View file

@ -23,6 +23,7 @@ import (
"errors"
"fmt"
"io"
"log"
"math/rand"
"os"
"path/filepath"
@ -71,6 +72,10 @@ type Command struct {
Started chan struct{}
// Done will be closed when Command.Close() is called
Done chan struct{}
// Passed to the Gossip implementation.
logOutput io.Writer
logger *log.Logger
}
// NewCommand returns a new instance of Main.
@ -115,7 +120,31 @@ func (m *Command) Run(args ...string) (err error) {
return fmt.Errorf("server.Open: %v", err)
}
m.Server.Logger().Printf("Listening as %s\n", m.Server.URI)
m.Server.Logger.Printf("Listening as %s\n", m.Server.URI)
return nil
}
// SetupLogger sets up the logger based on the configuration.
func (m *Command) SetupLogger() error {
var err error
if m.Config.LogPath == "" {
m.logOutput = m.Stderr
} else {
m.logOutput, err = os.OpenFile(m.Config.LogPath, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600)
if err != nil {
return err
}
}
if m.Config.Verbose {
vbl := pilosa.NewVerboseLogger(m.logOutput)
m.logger = vbl.Logger()
m.Server.Logger = vbl
} else {
sl := pilosa.NewStandardLogger(m.logOutput)
m.logger = sl.Logger()
m.Server.Logger = sl
}
return nil
}
@ -126,6 +155,10 @@ func (m *Command) SetupServer() error {
return err
}
m.Server.Handler.Logger = m.Server.Logger
m.Server.Holder.Logger = m.Server.Logger
m.Server.Holder.Stats.SetLogger(m.Server.Logger)
uri, err := pilosa.AddressWithDefaults(m.Config.Bind)
if err != nil {
@ -136,15 +169,10 @@ func (m *Command) SetupServer() error {
cluster := pilosa.NewCluster()
cluster.ReplicaN = m.Config.Cluster.ReplicaN
cluster.Holder = m.Server.Holder
cluster.Logger = m.Server.Logger
m.Server.Cluster = cluster
// Setup logging output.
m.Server.LogOutput, err = GetLogWriter(m.Config.LogPath, m.Stderr)
if err != nil {
return err
}
// Configure data directory (for Cluster .topology)
m.Server.Cluster.Path = m.Config.DataDir
@ -152,7 +180,7 @@ func (m *Command) SetupServer() error {
m.Server.Holder.NewAttrStore = boltdb.NewAttrStore
// Configure holder.
m.Server.Logger().Printf("Using data from: %s\n", m.Config.DataDir)
m.Server.Logger.Printf("Using data from: %s\n", m.Config.DataDir)
m.Server.Holder.Path = m.Config.DataDir
m.Server.MetricInterval = time.Duration(m.Config.Metric.PollInterval)
if m.Config.Metric.Diagnostics {
@ -165,8 +193,6 @@ func (m *Command) SetupServer() error {
return err
}
m.Server.Holder.Stats.SetLogger(m.Server.LogOutput)
// Copy configuration flags.
m.Server.MaxWritesPerRequest = m.Config.MaxWritesPerRequest
@ -246,7 +272,7 @@ func (m *Command) SetupNetworking() error {
if m.GossipTransport != nil {
transport = m.GossipTransport
} else {
transport, err = gossip.NewTransport(gossipHost, gossipPort)
transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger)
if err != nil {
return err
}
@ -257,8 +283,8 @@ func (m *Command) SetupNetworking() error {
m.Server.Cluster.Coordinator = m.Server.NodeID
}
m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.Server.LogOutput)
gossipMemberSet, err := gossip.NewGossipMemberSetWithTransport(m.Server.NodeID, m.Config, transport, m.Server)
m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.Server.Logger)
gossipMemberSet, err := gossip.NewGossipMemberSet(m.Server.NodeID, m.Config, m.Server, gossip.WithLogger(m.logger), gossip.WithTransport(transport))
if err != nil {
return err
}
@ -269,26 +295,11 @@ func (m *Command) SetupNetworking() error {
return nil
}
// GetLogWriter opens a file for logging, or a default io.Writer (such as stderr) for an empty path.
func GetLogWriter(path string, defaultWriter io.Writer) (io.Writer, error) {
// This is split out so it can be used in NewServeCmd as well as SetupServer
if path == "" {
return defaultWriter, nil
} else {
logFile, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600)
if err != nil {
return nil, err
}
return logFile, nil
}
}
// Close shuts down the server.
func (m *Command) Close() error {
var logErr error
serveErr := m.Server.Close()
logOutput := m.Server.LogOutput
if closer, ok := logOutput.(io.Closer); ok {
if closer, ok := m.logOutput.(io.Closer); ok {
logErr = closer.Close()
}
close(m.Done)

View file

@ -16,7 +16,6 @@ package pilosa
import (
"expvar"
"io"
"sort"
"strings"
"sync"
@ -57,7 +56,7 @@ type StatsClient interface {
Timing(name string, value time.Duration, rate float64)
// SetLogger Set the logger output type
SetLogger(logger io.Writer)
SetLogger(logger Logger)
// Starts the service
Open()
@ -79,7 +78,7 @@ func (c *nopStatsClient) Gauge(name string, value float64, rate float64)
func (c *nopStatsClient) Histogram(name string, value float64, rate float64) {}
func (c *nopStatsClient) Set(name string, value string, rate float64) {}
func (c *nopStatsClient) Timing(name string, value time.Duration, rate float64) {}
func (c *nopStatsClient) SetLogger(logger io.Writer) {}
func (c *nopStatsClient) SetLogger(logger Logger) {}
func (c *nopStatsClient) Open() {}
func (c *nopStatsClient) Close() error { return nil }
@ -154,7 +153,7 @@ func (c *ExpvarStatsClient) Timing(name string, value time.Duration, rate float6
}
// SetLogger has no logger.
func (c *ExpvarStatsClient) SetLogger(logger io.Writer) {
func (c *ExpvarStatsClient) SetLogger(logger Logger) {
}
// Open no-op.
@ -226,7 +225,7 @@ func (a MultiStatsClient) Timing(name string, value time.Duration, rate float64)
}
// SetLogger Sets the StatsD logger output type.
func (a MultiStatsClient) SetLogger(logger io.Writer) {
func (a MultiStatsClient) SetLogger(logger Logger) {
for _, c := range a {
c.SetLogger(logger)
}

View file

@ -16,8 +16,6 @@ package pilosa_test
import (
"context"
"io"
"io/ioutil"
"net/http"
"strings"
"testing"
@ -38,7 +36,6 @@ func TestMultiStatClient_Expvar(t *testing.T) {
ms[0] = c
hldr.Stats = ms
hldr.Stats.SetLogger(ioutil.Discard)
hldr.MustCreateFragmentIfNotExists("d", "f", pilosa.ViewStandard, 0).SetBit(0, 0)
hldr.MustCreateFragmentIfNotExists("d", "f", pilosa.ViewStandard, 0).SetBit(0, 1)
hldr.MustCreateFragmentIfNotExists("d", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth)
@ -357,6 +354,6 @@ func (c *MockStats) Gauge(name string, value float64, rate float64) {}
func (c *MockStats) Histogram(name string, value float64, rate float64) {}
func (c *MockStats) Set(name string, value string, rate float64) {}
func (c *MockStats) Timing(name string, value time.Duration, rate float64) {}
func (c *MockStats) SetLogger(logger io.Writer) {}
func (c *MockStats) SetLogger(logger pilosa.Logger) {}
func (c *MockStats) Open() {}
func (c *MockStats) Close() error { return nil }

View file

@ -15,9 +15,6 @@
package statsd
import (
"io"
"io/ioutil"
"log"
"time"
"github.com/DataDog/datadog-go/statsd"
@ -40,9 +37,9 @@ var _ pilosa.StatsClient = &StatsClient{}
// StatsClient represents a StatsD implementation of pilosa.StatsClient.
type StatsClient struct {
client *statsd.Client
tags []string
logOutput io.Writer
client *statsd.Client
tags []string
logger pilosa.Logger
}
// NewStatsClient returns a new instance of StatsClient.
@ -53,8 +50,8 @@ func NewStatsClient(host string) (*StatsClient, error) {
}
return &StatsClient{
client: c,
logOutput: ioutil.Discard,
client: c,
logger: pilosa.NopLogger,
}, nil
}
@ -74,16 +71,16 @@ func (c *StatsClient) Tags() []string {
// WithTags returns a new client with additional tags appended.
func (c *StatsClient) WithTags(tags ...string) pilosa.StatsClient {
return &StatsClient{
client: c.client,
tags: pilosa.UnionStringSlice(c.tags, tags),
logOutput: c.logOutput,
client: c.client,
tags: pilosa.UnionStringSlice(c.tags, tags),
logger: c.logger,
}
}
// 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.Printf("statsd.StatsClient.Count error: %s", err)
}
}
@ -91,44 +88,39 @@ 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.Printf("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.Printf("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.Printf("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.Printf("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.Printf("statsd.StatsClient.Timing error: %s", err)
}
}
// SetLogger has no logger
func (c *StatsClient) SetLogger(logger io.Writer) {
c.logOutput = logger
}
// logger returns a logger that writes to LogOutput
func (c *StatsClient) logger() *log.Logger {
return log.New(c.logOutput, "", log.LstdFlags)
// SetLogger sets the logger for client.
func (c *StatsClient) SetLogger(logger pilosa.Logger) {
c.logger = logger
}

View file

@ -15,7 +15,6 @@
package statsd_test
import (
"io/ioutil"
"reflect"
"testing"
"time"
@ -31,7 +30,6 @@ func TestStatsClient_WithTags(t *testing.T) {
t.Fatal(err)
}
defer c.Close()
c.SetLogger(ioutil.Discard)
// Create a new client with additional tags.
c1 := c.WithTags("foo", "bar")

View file

@ -41,7 +41,6 @@ func NewHandler() *Handler {
Handler: pilosa.NewHandler(),
}
h.Handler.Executor = &h.Executor
h.Handler.LogOutput = ioutil.Discard
// Handler test messages can no-op.
h.Broadcaster = pilosa.NopBroadcaster

View file

@ -15,7 +15,6 @@
package test
import (
"bytes"
"io/ioutil"
"os"
@ -26,7 +25,6 @@ import (
// Holder is a test wrapper for pilosa.Holder.
type Holder struct {
*pilosa.Holder
LogOutput bytes.Buffer
}
// NewHolder returns a new instance of Holder with a temporary path.
@ -38,7 +36,6 @@ func NewHolder() *Holder {
h := &Holder{Holder: pilosa.NewHolder()}
h.Path = path
h.Holder.LogOutput = &h.LogOutput
h.Holder.NewAttrStore = boltdb.NewAttrStore
return h
@ -62,10 +59,10 @@ func (h *Holder) Close() error {
// Reopen instantiates and opens a new holder.
// Note that the holder must be Closed first.
func (h *Holder) Reopen() error {
path, logOutput := h.Path, h.Holder.LogOutput
path, logger := h.Path, h.Holder.Logger
h.Holder = pilosa.NewHolder()
h.Holder.Path = path
h.Holder.LogOutput = logOutput
h.Holder.Logger = logger
h.Holder.NewAttrStore = boltdb.NewAttrStore
if err := h.Holder.Open(); err != nil {
return err

48
test/logger.go Normal file
View file

@ -0,0 +1,48 @@
// 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"
)
// BufferLogger represents a test Logger that holds log messages
// in a buffer for review.
type BufferLogger struct {
buf *bytes.Buffer
}
// NewBufferLogger returns a new instance of BufferLogger.
func NewBufferLogger() *BufferLogger {
return &BufferLogger{
buf: &bytes.Buffer{},
}
}
func (b *BufferLogger) Printf(format string, v ...interface{}) {
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) {
return ioutil.ReadAll(b.buf)
}

View file

@ -181,7 +181,7 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (
}
// Open gossip transport to use in SetupServer.
transport, err := gossip.NewTransport(host, bindPort)
transport, err := gossip.NewTransport(host, bindPort, nil)
if err != nil {
return seed, err
}

16
view.go
View file

@ -16,9 +16,6 @@ package pilosa
import (
"fmt"
"io"
"io/ioutil"
"log"
"os"
"path/filepath"
"strconv"
@ -64,7 +61,7 @@ type View struct {
stats StatsClient
RowAttrStore AttrStore
LogOutput io.Writer
Logger Logger
}
// NewView returns a new instance of View.
@ -81,7 +78,7 @@ func NewView(path, index, frame, name string, cacheSize uint32) *View {
broadcaster: NopBroadcaster,
stats: NopStatsClient,
LogOutput: ioutil.Discard,
Logger: NopLogger,
}
}
@ -126,9 +123,6 @@ func (v *View) Open() error {
return nil
}
// logger returns a logger instance for the view.
func (v *View) logger() *log.Logger { return log.New(v.LogOutput, "", log.LstdFlags) }
// openFragments opens and initializes the fragments inside the view.
func (v *View) openFragments() error {
file, err := os.Open(filepath.Join(v.path, "fragments"))
@ -275,7 +269,7 @@ func (v *View) newFragment(path string, slice uint64) *Fragment {
frag := NewFragment(path, v.index, v.frame, v.name, slice)
frag.CacheType = v.cacheType
frag.CacheSize = v.cacheSize
frag.LogOutput = v.LogOutput
frag.Logger = v.Logger
frag.stats = v.stats.WithTags(fmt.Sprintf("slice:%d", slice))
return frag
}
@ -288,7 +282,7 @@ func (v *View) DeleteFragment(slice uint64) error {
return ErrFragmentNotFound
}
v.logger().Printf("delete fragment: (%s/%s/%s) %d", v.index, v.frame, v.name, slice)
v.Logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.frame, v.name, slice)
// Close data files before deletion.
if err := fragment.Close(); err != nil {
@ -302,7 +296,7 @@ func (v *View) DeleteFragment(slice uint64) error {
// Delete fragment cache file.
if err := os.Remove(fragment.CachePath()); err != nil {
v.logger().Printf("no cache file to delete for slice %d", slice)
v.Logger.Printf("no cache file to delete for slice %d", slice)
}
delete(v.fragments, slice)