WIP
This commit is contained in:
parent
5aae6ed383
commit
bb80029ef9
17
goldsmith.go
17
goldsmith.go
@ -35,9 +35,9 @@ import (
|
|||||||
type goldsmith struct {
|
type goldsmith struct {
|
||||||
srcDir, dstDir string
|
srcDir, dstDir string
|
||||||
|
|
||||||
stages []*stage
|
stages []*stage
|
||||||
active int64
|
active int64
|
||||||
stalled int64
|
idle int64
|
||||||
|
|
||||||
refs map[string]bool
|
refs map[string]bool
|
||||||
refMtx sync.Mutex
|
refMtx sync.Mutex
|
||||||
@ -57,7 +57,7 @@ func (gs *goldsmith) queueFiles(target uint) {
|
|||||||
|
|
||||||
for path := range files {
|
for path := range files {
|
||||||
for {
|
for {
|
||||||
if gs.active-gs.stalled >= int64(target) {
|
if gs.active-gs.idle >= int64(target) {
|
||||||
time.Sleep(time.Millisecond)
|
time.Sleep(time.Millisecond)
|
||||||
} else {
|
} else {
|
||||||
break
|
break
|
||||||
@ -149,14 +149,15 @@ func (gs *goldsmith) referenceFile(path string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (gs *goldsmith) fault(s *stage, step string, f *file, err error) {
|
func (gs *goldsmith) fault(name string, f *file, err error) {
|
||||||
gs.faultMtx.Lock()
|
gs.faultMtx.Lock()
|
||||||
defer gs.faultMtx.Unlock()
|
defer gs.faultMtx.Unlock()
|
||||||
|
|
||||||
color.Red("Fault Detected\n")
|
color.Red("Fault Detected\n")
|
||||||
color.Yellow("\tPlugin:\t%s\n", color.WhiteString(s.name))
|
color.Yellow("\tPlugin:\t%s\n", color.WhiteString(name))
|
||||||
color.Yellow("\tStep:\t%s\n", color.WhiteString(step))
|
if f != nil {
|
||||||
color.Yellow("\tFile:\t%s\n", color.WhiteString(f.path))
|
color.Yellow("\tFile:\t%s\n", color.WhiteString(f.path))
|
||||||
|
}
|
||||||
color.Yellow("\tError:\t%s\n\n", color.WhiteString(err.Error()))
|
color.Yellow("\tError:\t%s\n\n", color.WhiteString(err.Error()))
|
||||||
|
|
||||||
gs.tainted = true
|
gs.tainted = true
|
||||||
|
39
stage.go
39
stage.go
@ -28,7 +28,6 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type stage struct {
|
type stage struct {
|
||||||
name string
|
|
||||||
gs *goldsmith
|
gs *goldsmith
|
||||||
input, output chan *file
|
input, output chan *file
|
||||||
}
|
}
|
||||||
@ -45,9 +44,13 @@ func newStage(gs *goldsmith) *stage {
|
|||||||
|
|
||||||
func (s *stage) chain(p Plugin) {
|
func (s *stage) chain(p Plugin) {
|
||||||
defer close(s.output)
|
defer close(s.output)
|
||||||
s.name = p.Name()
|
|
||||||
|
|
||||||
init, _ := p.(Initializer)
|
name, flags, err := p.Initialize(s)
|
||||||
|
if err != nil {
|
||||||
|
s.gs.fault(name, nil, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
accept, _ := p.(Accepter)
|
accept, _ := p.(Accepter)
|
||||||
proc, _ := p.(Processor)
|
proc, _ := p.(Processor)
|
||||||
fin, _ := p.(Finalizer)
|
fin, _ := p.(Finalizer)
|
||||||
@ -59,21 +62,13 @@ func (s *stage) chain(p Plugin) {
|
|||||||
)
|
)
|
||||||
|
|
||||||
dispatch := func(f *file) {
|
dispatch := func(f *file) {
|
||||||
if fin == nil {
|
if flags&PLUGIN_FLAG_BATCH == PLUGIN_FLAG_BATCH {
|
||||||
s.output <- f
|
atomic.AddInt64(&s.gs.idle, 1)
|
||||||
} else {
|
|
||||||
mtx.Lock()
|
mtx.Lock()
|
||||||
batch = append(batch, f)
|
batch = append(batch, f)
|
||||||
mtx.Unlock()
|
mtx.Unlock()
|
||||||
|
} else {
|
||||||
atomic.AddInt64(&s.gs.stalled, 1)
|
s.output <- f
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if init != nil {
|
|
||||||
if err := init.Initialize(s); err != nil {
|
|
||||||
s.gs.fault(s, "Initialization", nil, err)
|
|
||||||
return
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -89,7 +84,7 @@ func (s *stage) chain(p Plugin) {
|
|||||||
f.rewind()
|
f.rewind()
|
||||||
keep, err := proc.Process(s, f)
|
keep, err := proc.Process(s, f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.gs.fault(s, "Processing", f, err)
|
s.gs.fault(name, f, err)
|
||||||
} else if keep {
|
} else if keep {
|
||||||
dispatch(f)
|
dispatch(f)
|
||||||
} else {
|
} else {
|
||||||
@ -102,14 +97,14 @@ func (s *stage) chain(p Plugin) {
|
|||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
if fin != nil {
|
if fin != nil {
|
||||||
if err := fin.Finalize(s, batch); err != nil {
|
if err := fin.Finalize(s); err != nil {
|
||||||
s.gs.fault(s, "Finalization", nil, err)
|
s.gs.fault(name, nil, err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
for _, f := range batch {
|
for _, f := range batch {
|
||||||
atomic.AddInt64(&s.gs.stalled, -1)
|
atomic.AddInt64(&s.gs.idle, -1)
|
||||||
s.output <- f.(*file)
|
s.output <- f.(*file)
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
12
types.go
12
types.go
@ -79,20 +79,20 @@ type Context interface {
|
|||||||
DstDir() string
|
DstDir() string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
PLUGIN_FLAG_BATCH = 1 << iota
|
||||||
|
)
|
||||||
|
|
||||||
type Plugin interface {
|
type Plugin interface {
|
||||||
Name() string
|
Initialize(ctx Context) (string, uint, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
type Accepter interface {
|
type Accepter interface {
|
||||||
Accept(file File) bool
|
Accept(file File) bool
|
||||||
}
|
}
|
||||||
|
|
||||||
type Initializer interface {
|
|
||||||
Initialize(ctx Context) error
|
|
||||||
}
|
|
||||||
|
|
||||||
type Finalizer interface {
|
type Finalizer interface {
|
||||||
Finalize(ctx Context, fs []File) error
|
Finalize(ctx Context) error
|
||||||
}
|
}
|
||||||
|
|
||||||
type Processor interface {
|
type Processor interface {
|
||||||
|
Loading…
Reference in New Issue
Block a user