MM-25943: Upgrade dependencies for server (#14932)
* MM-25943: Upgrade dependencies for server * tmp
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
eff2209a7e
Коммит
d3156395a1
71
vendor/github.com/hashicorp/go-hclog/exclude.go
сгенерированный
поставляемый
Обычный файл
71
vendor/github.com/hashicorp/go-hclog/exclude.go
сгенерированный
поставляемый
Обычный файл
@@ -0,0 +1,71 @@
|
||||
package hclog
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// ExcludeByMessage provides a simple way to build a list of log messages that
|
||||
// can be queried and matched. This is meant to be used with the Exclude
|
||||
// option on Options to suppress log messages. This does not hold any mutexs
|
||||
// within itself, so normal usage would be to Add entries at setup and none after
|
||||
// Exclude is going to be called. Exclude is called with a mutex held within
|
||||
// the Logger, so that doesn't need to use a mutex. Example usage:
|
||||
//
|
||||
// f := new(ExcludeByMessage)
|
||||
// f.Add("Noisy log message text")
|
||||
// appLogger.Exclude = f.Exclude
|
||||
type ExcludeByMessage struct {
|
||||
messages map[string]struct{}
|
||||
}
|
||||
|
||||
// Add a message to be filtered. Do not call this after Exclude is to be called
|
||||
// due to concurrency issues.
|
||||
func (f *ExcludeByMessage) Add(msg string) {
|
||||
if f.messages == nil {
|
||||
f.messages = make(map[string]struct{})
|
||||
}
|
||||
|
||||
f.messages[msg] = struct{}{}
|
||||
}
|
||||
|
||||
// Return true if the given message should be included
|
||||
func (f *ExcludeByMessage) Exclude(level Level, msg string, args ...interface{}) bool {
|
||||
_, ok := f.messages[msg]
|
||||
return ok
|
||||
}
|
||||
|
||||
// ExcludeByPrefix is a simple type to match a message string that has a common prefix.
|
||||
type ExcludeByPrefix string
|
||||
|
||||
// Matches an message that starts with the prefix.
|
||||
func (p ExcludeByPrefix) Exclude(level Level, msg string, args ...interface{}) bool {
|
||||
return strings.HasPrefix(msg, string(p))
|
||||
}
|
||||
|
||||
// ExcludeByRegexp takes a regexp and uses it to match a log message string. If it matches
|
||||
// the log entry is excluded.
|
||||
type ExcludeByRegexp struct {
|
||||
Regexp *regexp.Regexp
|
||||
}
|
||||
|
||||
// Exclude the log message if the message string matches the regexp
|
||||
func (e ExcludeByRegexp) Exclude(level Level, msg string, args ...interface{}) bool {
|
||||
return e.Regexp.MatchString(msg)
|
||||
}
|
||||
|
||||
// ExcludeFuncs is a slice of functions that will called to see if a log entry
|
||||
// should be filtered or not. It stops calling functions once at least one returns
|
||||
// true.
|
||||
type ExcludeFuncs []func(level Level, msg string, args ...interface{}) bool
|
||||
|
||||
// Calls each function until one of them returns true
|
||||
func (ff ExcludeFuncs) Exclude(level Level, msg string, args ...interface{}) bool {
|
||||
for _, f := range ff {
|
||||
if f(level, msg, args...) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
25
vendor/github.com/hashicorp/go-hclog/intlogger.go
сгенерированный
поставляемый
25
vendor/github.com/hashicorp/go-hclog/intlogger.go
сгенерированный
поставляемый
@@ -65,6 +65,8 @@ type intLogger struct {
|
||||
level *int32
|
||||
|
||||
implied []interface{}
|
||||
|
||||
exclude func(level Level, msg string, args ...interface{}) bool
|
||||
}
|
||||
|
||||
// New returns a configured logger.
|
||||
@@ -106,11 +108,14 @@ func newLogger(opts *LoggerOptions) *intLogger {
|
||||
mutex: mutex,
|
||||
writer: newWriter(output, opts.Color),
|
||||
level: new(int32),
|
||||
exclude: opts.Exclude,
|
||||
}
|
||||
|
||||
l.setColorization(opts)
|
||||
|
||||
if opts.TimeFormat != "" {
|
||||
if opts.DisableTime {
|
||||
l.timeFormat = ""
|
||||
} else if opts.TimeFormat != "" {
|
||||
l.timeFormat = opts.TimeFormat
|
||||
}
|
||||
|
||||
@@ -131,6 +136,10 @@ func (l *intLogger) log(name string, level Level, msg string, args ...interface{
|
||||
l.mutex.Lock()
|
||||
defer l.mutex.Unlock()
|
||||
|
||||
if l.exclude != nil && l.exclude(level, msg, args...) {
|
||||
return
|
||||
}
|
||||
|
||||
if l.json {
|
||||
l.logJSON(t, name, level, msg, args...)
|
||||
} else {
|
||||
@@ -169,12 +178,14 @@ func trimCallerPath(path string) string {
|
||||
return path[idx+1:]
|
||||
}
|
||||
|
||||
var logImplFile = regexp.MustCompile(`github.com/hashicorp/go-hclog/.+logger.go$`)
|
||||
var logImplFile = regexp.MustCompile(`.+intlogger.go|.+interceptlogger.go$`)
|
||||
|
||||
// Non-JSON logging format function
|
||||
func (l *intLogger) logPlain(t time.Time, name string, level Level, msg string, args ...interface{}) {
|
||||
l.writer.WriteString(t.Format(l.timeFormat))
|
||||
l.writer.WriteByte(' ')
|
||||
if len(l.timeFormat) > 0 {
|
||||
l.writer.WriteString(t.Format(l.timeFormat))
|
||||
l.writer.WriteByte(' ')
|
||||
}
|
||||
|
||||
s, ok := _levelToBracket[level]
|
||||
if ok {
|
||||
@@ -260,6 +271,12 @@ func (l *intLogger) logPlain(t time.Time, name string, level Level, msg string,
|
||||
val = strconv.FormatUint(uint64(st), 10)
|
||||
case uint8:
|
||||
val = strconv.FormatUint(uint64(st), 10)
|
||||
case Hex:
|
||||
val = "0x" + strconv.FormatUint(uint64(st), 16)
|
||||
case Octal:
|
||||
val = "0" + strconv.FormatUint(uint64(st), 8)
|
||||
case Binary:
|
||||
val = "0b" + strconv.FormatUint(uint64(st), 2)
|
||||
case CapturedStacktrace:
|
||||
stacktrace = st
|
||||
continue FOR
|
||||
|
||||
22
vendor/github.com/hashicorp/go-hclog/logger.go
сгенерированный
поставляемый
22
vendor/github.com/hashicorp/go-hclog/logger.go
сгенерированный
поставляемый
@@ -52,6 +52,18 @@ func Fmt(str string, args ...interface{}) Format {
|
||||
return append(Format{str}, args...)
|
||||
}
|
||||
|
||||
// A simple shortcut to format numbers in hex when displayed with the normal
|
||||
// text output. For example: L.Info("header value", Hex(17))
|
||||
type Hex int
|
||||
|
||||
// A simple shortcut to format numbers in octal when displayed with the normal
|
||||
// text output. For example: L.Info("perms", Octal(17))
|
||||
type Octal int
|
||||
|
||||
// A simple shortcut to format numbers in binary when displayed with the normal
|
||||
// text output. For example: L.Info("bits", Binary(17))
|
||||
type Binary int
|
||||
|
||||
// ColorOption expresses how the output should be colored, if at all.
|
||||
type ColorOption uint8
|
||||
|
||||
@@ -218,9 +230,19 @@ type LoggerOptions struct {
|
||||
// The time format to use instead of the default
|
||||
TimeFormat string
|
||||
|
||||
// Control whether or not to display the time at all. This is required
|
||||
// because setting TimeFormat to empty assumes the default format.
|
||||
DisableTime bool
|
||||
|
||||
// Color the output. On Windows, colored logs are only avaiable for io.Writers that
|
||||
// are concretely instances of *os.File.
|
||||
Color ColorOption
|
||||
|
||||
// A function which is called with the log information and if it returns true the value
|
||||
// should not be logged.
|
||||
// This is useful when interacting with a system that you wish to suppress the log
|
||||
// message for (because it's too noisy, etc)
|
||||
Exclude func(level Level, msg string, args ...interface{}) bool
|
||||
}
|
||||
|
||||
// InterceptLogger describes the interface for using a logger
|
||||
|
||||
21
vendor/github.com/hashicorp/go-hclog/stdlog.go
сгенерированный
поставляемый
21
vendor/github.com/hashicorp/go-hclog/stdlog.go
сгенерированный
поставляемый
@@ -2,6 +2,7 @@ package hclog
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"log"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -72,3 +73,23 @@ func (s *stdlogAdapter) pickLevel(str string) (Level, string) {
|
||||
return Info, str
|
||||
}
|
||||
}
|
||||
|
||||
type logWriter struct {
|
||||
l *log.Logger
|
||||
}
|
||||
|
||||
func (l *logWriter) Write(b []byte) (int, error) {
|
||||
l.l.Println(string(bytes.TrimRight(b, " \n\t")))
|
||||
return len(b), nil
|
||||
}
|
||||
|
||||
// Takes a standard library logger and returns a Logger that will write to it
|
||||
func FromStandardLogger(l *log.Logger, opts *LoggerOptions) Logger {
|
||||
var dl LoggerOptions = *opts
|
||||
|
||||
// Use the time format that log.Logger uses
|
||||
dl.DisableTime = true
|
||||
dl.Output = &logWriter{l}
|
||||
|
||||
return New(&dl)
|
||||
}
|
||||
|
||||
16
vendor/github.com/hashicorp/go-plugin/client.go
сгенерированный
поставляемый
16
vendor/github.com/hashicorp/go-plugin/client.go
сгенерированный
поставляемый
@@ -212,6 +212,12 @@ type ReattachConfig struct {
|
||||
Protocol Protocol
|
||||
Addr net.Addr
|
||||
Pid int
|
||||
|
||||
// Test is set to true if this is reattaching to to a plugin in "test mode"
|
||||
// (see ServeConfig.Test). In this mode, client.Kill will NOT kill the
|
||||
// process and instead will rely on the plugin to terminate itself. This
|
||||
// should not be used in non-test environments.
|
||||
Test bool
|
||||
}
|
||||
|
||||
// SecureConfig is used to configure a client to verify the integrity of an
|
||||
@@ -825,15 +831,21 @@ func (c *Client) reattach() (net.Addr, error) {
|
||||
c.exited = true
|
||||
}(p.Pid)
|
||||
|
||||
// Set the address and process
|
||||
// Set the address and protocol
|
||||
c.address = c.config.Reattach.Addr
|
||||
c.process = p
|
||||
c.protocol = c.config.Reattach.Protocol
|
||||
if c.protocol == "" {
|
||||
// Default the protocol to net/rpc for backwards compatibility
|
||||
c.protocol = ProtocolNetRPC
|
||||
}
|
||||
|
||||
// If we're in test mode, we do NOT set the process. This avoids the
|
||||
// process being killed (the only purpose we have for c.process), since
|
||||
// in test mode the process is responsible for exiting on its own.
|
||||
if !c.config.Reattach.Test {
|
||||
c.process = p
|
||||
}
|
||||
|
||||
return c.address, nil
|
||||
}
|
||||
|
||||
|
||||
285
vendor/github.com/hashicorp/go-plugin/server.go
сгенерированный
поставляемый
285
vendor/github.com/hashicorp/go-plugin/server.go
сгенерированный
поставляемый
@@ -1,11 +1,13 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"net"
|
||||
@@ -84,20 +86,50 @@ type ServeConfig struct {
|
||||
// server will create a default logger.
|
||||
Logger hclog.Logger
|
||||
|
||||
// Listener is the listener that the plugin server will listen for
|
||||
// plugin connections. THIS DOES NOT NORMALLY NEED TO BE SET. If this
|
||||
// isn't set, the plugin chooses a listener. This is exposed in case you
|
||||
// want to carefully control how a plugin is served.
|
||||
// Test, if non-nil, will put plugin serving into "test mode". This is
|
||||
// meant to be used as part of `go test` within a plugin's codebase to
|
||||
// launch the plugin in-process and output a ReattachConfig.
|
||||
//
|
||||
// If TLSProvider is set, this listener will be wrapped with a TLS
|
||||
// listener. If you want to manually control TLS you should set
|
||||
// TLSProvider to nil but be aware that the client side will need to be
|
||||
// manually made aware of the certificate used.
|
||||
// This changes the behavior of the server in a number of ways to
|
||||
// accomodate the expectation of running in-process:
|
||||
//
|
||||
// Serve will take ownership of this listener and close it when it is
|
||||
// complete. The caller should NOT close this listener once `Serve` is
|
||||
// called.
|
||||
Listener net.Listener
|
||||
// * The handshake cookie is not validated.
|
||||
// * Stdout/stderr will receive plugin reads and writes
|
||||
// * Connection information will not be sent to stdout
|
||||
//
|
||||
Test *ServeTestConfig
|
||||
}
|
||||
|
||||
// ServeTestConfig configures plugin serving for test mode. See ServeConfig.Test.
|
||||
type ServeTestConfig struct {
|
||||
// Context, if set, will force the plugin serving to end when cancelled.
|
||||
// This is only a test configuration because the non-test configuration
|
||||
// expects to take over the process and therefore end on an interrupt or
|
||||
// kill signal. For tests, we need to kill the plugin serving routinely
|
||||
// and this provides a way to do so.
|
||||
//
|
||||
// If you want to wait for the plugin process to close before moving on,
|
||||
// you can wait on CloseCh.
|
||||
Context context.Context
|
||||
|
||||
// If this channel is non-nil, we will send the ReattachConfig via
|
||||
// this channel. This can be encoded (via JSON recommended) to the
|
||||
// plugin client to attach to this plugin.
|
||||
ReattachConfigCh chan<- *ReattachConfig
|
||||
|
||||
// CloseCh, if non-nil, will be closed when serving exits. This can be
|
||||
// used along with Context to determine when the server is fully shut down.
|
||||
// If this is not set, you can still use Context on its own, but note there
|
||||
// may be a period of time between canceling the context and the plugin
|
||||
// server being shut down.
|
||||
CloseCh chan<- struct{}
|
||||
|
||||
// SyncStdio, if true, will enable the client side "SyncStdout/Stderr"
|
||||
// functionality to work. This defaults to false because the implementation
|
||||
// of making this work within test environments is particularly messy
|
||||
// and SyncStdio functionality is fairly rare, so we default to the simple
|
||||
// scenario.
|
||||
SyncStdio bool
|
||||
}
|
||||
|
||||
// protocolVersion determines the protocol version and plugin set to be used by
|
||||
@@ -182,42 +214,46 @@ func protocolVersion(opts *ServeConfig) (int, Protocol, PluginSet) {
|
||||
// Serve serves the plugins given by ServeConfig.
|
||||
//
|
||||
// Serve doesn't return until the plugin is done being executed. Any
|
||||
// errors will be outputted to os.Stderr.
|
||||
// fixable errors will be output to os.Stderr and the process will
|
||||
// exit with a status code of 1. Serve will panic for unexpected
|
||||
// conditions where a user's fix is unknown.
|
||||
//
|
||||
// This is the method that plugins should call in their main() functions.
|
||||
func Serve(opts *ServeConfig) {
|
||||
// We use this to trigger an `os.Exit` so that we can execute our other
|
||||
// deferred functions.
|
||||
exitCode := -1
|
||||
// We use this to trigger an `os.Exit` so that we can execute our other
|
||||
// deferred functions. In test mode, we just output the err to stderr
|
||||
// and return.
|
||||
defer func() {
|
||||
if exitCode >= 0 {
|
||||
if opts.Test == nil && exitCode >= 0 {
|
||||
os.Exit(exitCode)
|
||||
}
|
||||
|
||||
if opts.Test != nil && opts.Test.CloseCh != nil {
|
||||
close(opts.Test.CloseCh)
|
||||
}
|
||||
}()
|
||||
|
||||
// If our listener is not nil, then we want to close that on exit.
|
||||
if opts.Listener != nil {
|
||||
defer opts.Listener.Close()
|
||||
}
|
||||
if opts.Test == nil {
|
||||
// Validate the handshake config
|
||||
if opts.MagicCookieKey == "" || opts.MagicCookieValue == "" {
|
||||
fmt.Fprintf(os.Stderr,
|
||||
"Misconfigured ServeConfig given to serve this plugin: no magic cookie\n"+
|
||||
"key or value was set. Please notify the plugin author and report\n"+
|
||||
"this as a bug.\n")
|
||||
exitCode = 1
|
||||
return
|
||||
}
|
||||
|
||||
// Validate the handshake config
|
||||
if opts.MagicCookieKey == "" || opts.MagicCookieValue == "" {
|
||||
fmt.Fprintf(os.Stderr,
|
||||
"Misconfigured ServeConfig given to serve this plugin: no magic cookie\n"+
|
||||
"key or value was set. Please notify the plugin author and report\n"+
|
||||
"this as a bug.\n")
|
||||
exitCode = 1
|
||||
return
|
||||
}
|
||||
|
||||
// First check the cookie
|
||||
if os.Getenv(opts.MagicCookieKey) != opts.MagicCookieValue {
|
||||
fmt.Fprintf(os.Stderr,
|
||||
"This binary is a plugin. These are not meant to be executed directly.\n"+
|
||||
"Please execute the program that consumes these plugins, which will\n"+
|
||||
"load any plugins automatically\n")
|
||||
exitCode = 1
|
||||
return
|
||||
// First check the cookie
|
||||
if os.Getenv(opts.MagicCookieKey) != opts.MagicCookieValue {
|
||||
fmt.Fprintf(os.Stderr,
|
||||
"This binary is a plugin. These are not meant to be executed directly.\n"+
|
||||
"Please execute the program that consumes these plugins, which will\n"+
|
||||
"load any plugins automatically\n")
|
||||
exitCode = 1
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// negotiate the version and plugins
|
||||
@@ -237,34 +273,18 @@ func Serve(opts *ServeConfig) {
|
||||
})
|
||||
}
|
||||
|
||||
// Create our new stdout, stderr files. These will override our built-in
|
||||
// stdout/stderr so that it works across the stream boundary.
|
||||
stdout_r, stdout_w, err := os.Pipe()
|
||||
// Register a listener so we can accept a connection
|
||||
listener, err := serverListener()
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "Error preparing plugin: %s\n", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
stderr_r, stderr_w, err := os.Pipe()
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "Error preparing plugin: %s\n", err)
|
||||
os.Exit(1)
|
||||
logger.Error("plugin init error", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
listener := opts.Listener
|
||||
if listener == nil {
|
||||
// Register a listener so we can accept a connection
|
||||
listener, err = serverListener()
|
||||
if err != nil {
|
||||
logger.Error("plugin init error", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Close the listener on return. We wrap this in a func() on purpose
|
||||
// because the "listener" reference may change to TLS.
|
||||
defer func() {
|
||||
listener.Close()
|
||||
}()
|
||||
}
|
||||
// Close the listener on return. We wrap this in a func() on purpose
|
||||
// because the "listener" reference may change to TLS.
|
||||
defer func() {
|
||||
listener.Close()
|
||||
}()
|
||||
|
||||
var tlsConfig *tls.Config
|
||||
if opts.TLSProvider != nil {
|
||||
@@ -313,6 +333,33 @@ func Serve(opts *ServeConfig) {
|
||||
// Create the channel to tell us when we're done
|
||||
doneCh := make(chan struct{})
|
||||
|
||||
// Create our new stdout, stderr files. These will override our built-in
|
||||
// stdout/stderr so that it works across the stream boundary.
|
||||
var stdout_r, stderr_r io.Reader
|
||||
stdout_r, stdout_w, err := os.Pipe()
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "Error preparing plugin: %s\n", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
stderr_r, stderr_w, err := os.Pipe()
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "Error preparing plugin: %s\n", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
// If we're in test mode, we tee off the reader and write the data
|
||||
// as-is to our normal Stdout and Stderr so that they continue working
|
||||
// while stdio works. This is because in test mode, we assume we're running
|
||||
// in `go test` or some equivalent and we want output to go to standard
|
||||
// locations.
|
||||
if opts.Test != nil {
|
||||
// TODO(mitchellh): This isn't super ideal because a TeeReader
|
||||
// only works if the reader side is actively read. If we never
|
||||
// connect via a plugin client, the output still gets swallowed.
|
||||
stdout_r = io.TeeReader(stdout_r, os.Stdout)
|
||||
stderr_r = io.TeeReader(stderr_r, os.Stderr)
|
||||
}
|
||||
|
||||
// Build the server type
|
||||
var server ServerProtocol
|
||||
switch protoType {
|
||||
@@ -355,40 +402,96 @@ func Serve(opts *ServeConfig) {
|
||||
|
||||
logger.Debug("plugin address", "network", listener.Addr().Network(), "address", listener.Addr().String())
|
||||
|
||||
// Output the address and service name to stdout so that the client can bring it up.
|
||||
fmt.Printf("%d|%d|%s|%s|%s|%s\n",
|
||||
CoreProtocolVersion,
|
||||
protoVersion,
|
||||
listener.Addr().Network(),
|
||||
listener.Addr().String(),
|
||||
protoType,
|
||||
serverCert)
|
||||
os.Stdout.Sync()
|
||||
|
||||
// Eat the interrupts
|
||||
ch := make(chan os.Signal, 1)
|
||||
signal.Notify(ch, os.Interrupt)
|
||||
go func() {
|
||||
count := 0
|
||||
for {
|
||||
<-ch
|
||||
count++
|
||||
logger.Trace("plugin received interrupt signal, ignoring", "count", count)
|
||||
// Output the address and service name to stdout so that the client can
|
||||
// bring it up. In test mode, we don't do this because clients will
|
||||
// attach via a reattach config.
|
||||
if opts.Test == nil {
|
||||
fmt.Printf("%d|%d|%s|%s|%s|%s\n",
|
||||
CoreProtocolVersion,
|
||||
protoVersion,
|
||||
listener.Addr().Network(),
|
||||
listener.Addr().String(),
|
||||
protoType,
|
||||
serverCert)
|
||||
os.Stdout.Sync()
|
||||
} else if ch := opts.Test.ReattachConfigCh; ch != nil {
|
||||
// Send back the reattach config that can be used. This isn't
|
||||
// quite ready if they connect immediately but the client should
|
||||
// retry a few times.
|
||||
ch <- &ReattachConfig{
|
||||
Protocol: protoType,
|
||||
Addr: listener.Addr(),
|
||||
Pid: os.Getpid(),
|
||||
Test: true,
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// Set our new out, err
|
||||
os.Stdout = stdout_w
|
||||
os.Stderr = stderr_w
|
||||
// Eat the interrupts. In test mode we disable this so that go test
|
||||
// can be cancelled properly.
|
||||
if opts.Test == nil {
|
||||
ch := make(chan os.Signal, 1)
|
||||
signal.Notify(ch, os.Interrupt)
|
||||
go func() {
|
||||
count := 0
|
||||
for {
|
||||
<-ch
|
||||
count++
|
||||
logger.Trace("plugin received interrupt signal, ignoring", "count", count)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// Set our stdout, stderr to the stdio stream that clients can retrieve
|
||||
// using ClientConfig.SyncStdout/err. We only do this for non-test mode
|
||||
// or if the test mode explicitly requests it.
|
||||
//
|
||||
// In test mode, we use a multiwriter so that the data continues going
|
||||
// to the normal stdout/stderr so output can show up in test logs. We
|
||||
// also send to the stdio stream so that clients can continue working
|
||||
// if they depend on that.
|
||||
if opts.Test == nil || opts.Test.SyncStdio {
|
||||
if opts.Test != nil {
|
||||
// In test mode we need to maintain the original values so we can
|
||||
// reset it.
|
||||
defer func(out, err *os.File) {
|
||||
os.Stdout = out
|
||||
os.Stderr = err
|
||||
}(os.Stdout, os.Stderr)
|
||||
}
|
||||
os.Stdout = stdout_w
|
||||
os.Stderr = stderr_w
|
||||
}
|
||||
|
||||
// Accept connections and wait for completion
|
||||
go server.Serve(listener)
|
||||
|
||||
// Note that given the documentation of Serve we should probably be
|
||||
// setting exitCode = 0 and using os.Exit here. That's how it used to
|
||||
// work before extracting this library. However, for years we've done
|
||||
// this so we'll keep this functionality.
|
||||
<-doneCh
|
||||
ctx := context.Background()
|
||||
if opts.Test != nil && opts.Test.Context != nil {
|
||||
ctx = opts.Test.Context
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
// Cancellation. We can stop the server by closing the listener.
|
||||
// This isn't graceful at all but this is currently only used by
|
||||
// tests and its our only way to stop.
|
||||
listener.Close()
|
||||
|
||||
// If this is a grpc server, then we also ask the server itself to
|
||||
// end which will kill all connections. There isn't an easy way to do
|
||||
// this for net/rpc currently but net/rpc is more and more unused.
|
||||
if s, ok := server.(*GRPCServer); ok {
|
||||
s.Stop()
|
||||
}
|
||||
|
||||
// Wait for the server itself to shut down
|
||||
<-doneCh
|
||||
|
||||
case <-doneCh:
|
||||
// Note that given the documentation of Serve we should probably be
|
||||
// setting exitCode = 0 and using os.Exit here. That's how it used to
|
||||
// work before extracting this library. However, for years we've done
|
||||
// this so we'll keep this functionality.
|
||||
}
|
||||
}
|
||||
|
||||
func serverListener() (net.Listener, error) {
|
||||
|
||||
6
vendor/github.com/hashicorp/yamux/stream.go
сгенерированный
поставляемый
6
vendor/github.com/hashicorp/yamux/stream.go
сгенерированный
поставляемый
@@ -446,15 +446,17 @@ func (s *Stream) SetDeadline(t time.Time) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetReadDeadline sets the deadline for future Read calls.
|
||||
// SetReadDeadline sets the deadline for blocked and future Read calls.
|
||||
func (s *Stream) SetReadDeadline(t time.Time) error {
|
||||
s.readDeadline.Store(t)
|
||||
asyncNotify(s.recvNotifyCh)
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetWriteDeadline sets the deadline for future Write calls
|
||||
// SetWriteDeadline sets the deadline for blocked and future Write calls
|
||||
func (s *Stream) SetWriteDeadline(t time.Time) error {
|
||||
s.writeDeadline.Store(t)
|
||||
asyncNotify(s.sendNotifyCh)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
Ссылка в новой задаче
Block a user