2013-11-20 02:37:03 -05:00
|
|
|
package engine
|
|
|
|
|
|
|
|
import (
|
|
|
|
"bufio"
|
|
|
|
"container/ring"
|
|
|
|
"fmt"
|
|
|
|
"io"
|
2014-01-21 20:56:09 -05:00
|
|
|
"io/ioutil"
|
2013-11-20 02:37:03 -05:00
|
|
|
"sync"
|
|
|
|
)
|
|
|
|
|
|
|
|
type Output struct {
|
|
|
|
sync.Mutex
|
|
|
|
dests []io.Writer
|
|
|
|
tasks sync.WaitGroup
|
2014-01-10 17:54:54 -05:00
|
|
|
used bool
|
2013-11-20 02:37:03 -05:00
|
|
|
}
|
|
|
|
|
|
|
|
// NewOutput returns a new Output object with no destinations attached.
|
|
|
|
// Writing to an empty Output will cause the written data to be discarded.
|
|
|
|
func NewOutput() *Output {
|
|
|
|
return &Output{}
|
|
|
|
}
|
|
|
|
|
2014-01-10 17:54:54 -05:00
|
|
|
// Return true if something was written on this output
|
|
|
|
func (o *Output) Used() bool {
|
2014-01-23 19:20:51 -05:00
|
|
|
o.Lock()
|
|
|
|
defer o.Unlock()
|
2014-01-10 17:54:54 -05:00
|
|
|
return o.used
|
|
|
|
}
|
|
|
|
|
2013-11-20 02:37:03 -05:00
|
|
|
// Add attaches a new destination to the Output. Any data subsequently written
|
|
|
|
// to the output will be written to the new destination in addition to all the others.
|
|
|
|
// This method is thread-safe.
|
2014-01-22 18:54:22 -05:00
|
|
|
func (o *Output) Add(dst io.Writer) {
|
2014-01-23 19:20:51 -05:00
|
|
|
o.Lock()
|
|
|
|
defer o.Unlock()
|
2013-11-20 02:37:03 -05:00
|
|
|
o.dests = append(o.dests, dst)
|
2014-01-22 18:54:22 -05:00
|
|
|
}
|
|
|
|
|
|
|
|
// Set closes and remove existing destination and then attaches a new destination to
|
|
|
|
// the Output. Any data subsequently written to the output will be written to the new
|
|
|
|
// destination in addition to all the others. This method is thread-safe.
|
|
|
|
func (o *Output) Set(dst io.Writer) {
|
|
|
|
o.Close()
|
2014-01-23 19:20:51 -05:00
|
|
|
o.Lock()
|
|
|
|
defer o.Unlock()
|
2014-01-22 18:54:22 -05:00
|
|
|
o.dests = []io.Writer{dst}
|
2013-11-20 02:37:03 -05:00
|
|
|
}
|
|
|
|
|
|
|
|
// AddPipe creates an in-memory pipe with io.Pipe(), adds its writing end as a destination,
|
|
|
|
// and returns its reading end for consumption by the caller.
|
|
|
|
// This is a rough equivalent similar to Cmd.StdoutPipe() in the standard os/exec package.
|
|
|
|
// This method is thread-safe.
|
|
|
|
func (o *Output) AddPipe() (io.Reader, error) {
|
|
|
|
r, w := io.Pipe()
|
|
|
|
o.Add(w)
|
|
|
|
return r, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// AddTail starts a new goroutine which will read all subsequent data written to the output,
|
|
|
|
// line by line, and append the last `n` lines to `dst`.
|
|
|
|
func (o *Output) AddTail(dst *[]string, n int) error {
|
|
|
|
src, err := o.AddPipe()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
o.tasks.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer o.tasks.Done()
|
|
|
|
Tail(src, n, dst)
|
|
|
|
}()
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// AddString starts a new goroutine which will read all subsequent data written to the output,
|
|
|
|
// line by line, and store the last line into `dst`.
|
|
|
|
func (o *Output) AddString(dst *string) error {
|
|
|
|
src, err := o.AddPipe()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
o.tasks.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer o.tasks.Done()
|
|
|
|
lines := make([]string, 0, 1)
|
|
|
|
Tail(src, 1, &lines)
|
|
|
|
if len(lines) == 0 {
|
|
|
|
*dst = ""
|
|
|
|
} else {
|
|
|
|
*dst = lines[0]
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// Write writes the same data to all registered destinations.
|
|
|
|
// This method is thread-safe.
|
|
|
|
func (o *Output) Write(p []byte) (n int, err error) {
|
2014-01-23 19:20:51 -05:00
|
|
|
o.Lock()
|
|
|
|
defer o.Unlock()
|
2014-01-10 17:54:54 -05:00
|
|
|
o.used = true
|
2013-11-20 02:37:03 -05:00
|
|
|
var firstErr error
|
|
|
|
for _, dst := range o.dests {
|
|
|
|
_, err := dst.Write(p)
|
|
|
|
if err != nil && firstErr == nil {
|
|
|
|
firstErr = err
|
|
|
|
}
|
|
|
|
}
|
2013-11-21 21:28:56 -05:00
|
|
|
return len(p), firstErr
|
2013-11-20 02:37:03 -05:00
|
|
|
}
|
|
|
|
|
|
|
|
// Close unregisters all destinations and waits for all background
|
|
|
|
// AddTail and AddString tasks to complete.
|
|
|
|
// The Close method of each destination is called if it exists.
|
|
|
|
func (o *Output) Close() error {
|
2014-01-23 19:20:51 -05:00
|
|
|
o.Lock()
|
|
|
|
defer o.Unlock()
|
2013-11-20 02:37:03 -05:00
|
|
|
var firstErr error
|
|
|
|
for _, dst := range o.dests {
|
2014-05-08 15:57:19 -04:00
|
|
|
if closer, ok := dst.(io.Closer); ok {
|
2013-11-20 02:37:03 -05:00
|
|
|
err := closer.Close()
|
|
|
|
if err != nil && firstErr == nil {
|
|
|
|
firstErr = err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
o.tasks.Wait()
|
|
|
|
return firstErr
|
|
|
|
}
|
|
|
|
|
|
|
|
type Input struct {
|
|
|
|
src io.Reader
|
|
|
|
sync.Mutex
|
|
|
|
}
|
|
|
|
|
|
|
|
// NewInput returns a new Input object with no source attached.
|
|
|
|
// Reading to an empty Input will return io.EOF.
|
|
|
|
func NewInput() *Input {
|
|
|
|
return &Input{}
|
|
|
|
}
|
|
|
|
|
|
|
|
// Read reads from the input in a thread-safe way.
|
|
|
|
func (i *Input) Read(p []byte) (n int, err error) {
|
|
|
|
i.Mutex.Lock()
|
|
|
|
defer i.Mutex.Unlock()
|
|
|
|
if i.src == nil {
|
|
|
|
return 0, io.EOF
|
|
|
|
}
|
|
|
|
return i.src.Read(p)
|
|
|
|
}
|
|
|
|
|
2014-01-08 17:05:05 -05:00
|
|
|
// Closes the src
|
|
|
|
// Not thread safe on purpose
|
|
|
|
func (i *Input) Close() error {
|
|
|
|
if i.src != nil {
|
2014-05-08 15:57:19 -04:00
|
|
|
if closer, ok := i.src.(io.Closer); ok {
|
2014-01-08 17:05:05 -05:00
|
|
|
return closer.Close()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2013-11-20 02:37:03 -05:00
|
|
|
// Add attaches a new source to the input.
|
|
|
|
// Add can only be called once per input. Subsequent calls will
|
|
|
|
// return an error.
|
|
|
|
func (i *Input) Add(src io.Reader) error {
|
|
|
|
i.Mutex.Lock()
|
|
|
|
defer i.Mutex.Unlock()
|
|
|
|
if i.src != nil {
|
|
|
|
return fmt.Errorf("Maximum number of sources reached: 1")
|
|
|
|
}
|
|
|
|
i.src = src
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// Tail reads from `src` line per line, and returns the last `n` lines as an array.
|
|
|
|
// A ring buffer is used to only store `n` lines at any time.
|
|
|
|
func Tail(src io.Reader, n int, dst *[]string) {
|
|
|
|
scanner := bufio.NewScanner(src)
|
|
|
|
r := ring.New(n)
|
|
|
|
for scanner.Scan() {
|
|
|
|
if n == 0 {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
r.Value = scanner.Text()
|
|
|
|
r = r.Next()
|
|
|
|
}
|
|
|
|
r.Do(func(v interface{}) {
|
|
|
|
if v == nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
*dst = append(*dst, v.(string))
|
|
|
|
})
|
|
|
|
}
|
2013-12-08 01:16:10 -05:00
|
|
|
|
|
|
|
// AddEnv starts a new goroutine which will decode all subsequent data
|
|
|
|
// as a stream of json-encoded objects, and point `dst` to the last
|
|
|
|
// decoded object.
|
2013-11-20 02:37:03 -05:00
|
|
|
// The result `env` can be queried using the type-neutral Env interface.
|
2013-12-08 01:16:10 -05:00
|
|
|
// It is not safe to query `env` until the Output is closed.
|
|
|
|
func (o *Output) AddEnv() (dst *Env, err error) {
|
|
|
|
src, err := o.AddPipe()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
dst = &Env{}
|
|
|
|
o.tasks.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer o.tasks.Done()
|
|
|
|
decoder := NewDecoder(src)
|
|
|
|
for {
|
|
|
|
env, err := decoder.Decode()
|
|
|
|
if err != nil {
|
|
|
|
return
|
|
|
|
}
|
2013-11-20 02:37:03 -05:00
|
|
|
*dst = *env
|
2013-12-08 01:16:10 -05:00
|
|
|
}
|
|
|
|
}()
|
|
|
|
return dst, nil
|
|
|
|
}
|
2013-12-12 17:39:35 -05:00
|
|
|
|
2014-01-21 18:06:23 -05:00
|
|
|
func (o *Output) AddListTable() (dst *Table, err error) {
|
|
|
|
src, err := o.AddPipe()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
dst = NewTable("", 0)
|
|
|
|
o.tasks.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer o.tasks.Done()
|
2014-01-21 20:56:09 -05:00
|
|
|
content, err := ioutil.ReadAll(src)
|
|
|
|
if err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
if _, err := dst.ReadListFrom(content); err != nil {
|
2014-01-21 18:06:23 -05:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
return dst, nil
|
|
|
|
}
|
|
|
|
|
2013-12-12 17:39:35 -05:00
|
|
|
func (o *Output) AddTable() (dst *Table, err error) {
|
|
|
|
src, err := o.AddPipe()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
dst = NewTable("", 0)
|
|
|
|
o.tasks.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer o.tasks.Done()
|
|
|
|
if _, err := dst.ReadFrom(src); err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
return dst, nil
|
|
|
|
}
|