Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
feat: groundwork for comleting events and routing hierarchy.
Signed-off-by: Andrew Neudegg <andrew.neudegg@finbourne.com>
- Loading branch information
1 parent
c87595a
commit 271f3f1
Showing
13 changed files
with
237 additions
and
5 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,3 @@ | ||
# Meta | ||
|
||
Meta is for when you need to modify the behaviour of a source, relay or distributor. |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,29 @@ | ||
package builder | ||
|
||
import ( | ||
"github.com/andrewneudegg/delta/pkg/meta" | ||
"github.com/andrewneudegg/delta/pkg/meta/chaos" | ||
"github.com/andrewneudegg/delta/pkg/meta/example" | ||
|
||
"github.com/mitchellh/mapstructure" | ||
) | ||
|
||
// Get will return the given source with its data values initialised. | ||
func Get(distributorName string, metaConfiguration interface{}) (meta.M, error) { | ||
switch distributorName { | ||
case "meta/example": | ||
m := example.Example{} | ||
err := mapstructure.Decode(metaConfiguration, &m) | ||
return m, err | ||
case "meta/example2": | ||
m := example.Example{} | ||
err := mapstructure.Decode(metaConfiguration, &m) | ||
return m, err | ||
case "meta/chaos/simple": | ||
m := chaos.ChaosSimple{} | ||
err := mapstructure.Decode(metaConfiguration, &m) | ||
return m, err | ||
} | ||
|
||
return nil, nil | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,95 @@ | ||
package chaos | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"math/rand" | ||
|
||
"github.com/andrewneudegg/delta/pkg/distributor" | ||
"github.com/andrewneudegg/delta/pkg/events" | ||
"github.com/andrewneudegg/delta/pkg/relay" | ||
"github.com/andrewneudegg/delta/pkg/source" | ||
) | ||
|
||
// ChaosSimple intercepts and randomly fails events. | ||
type ChaosSimple struct { | ||
FailChance float32 // 0.5 | ||
} | ||
|
||
func (m ChaosSimple) isChance(f float32) bool { | ||
// rand.Float64() == 0.1, 0.5, 0.8 | ||
// if f == 0.1 (10% chance) then rand has to be above 0.9. | ||
// if f == 0.90 (90% chance) then rand has to be above 0.1. | ||
return rand.Float32() > (1 - f) | ||
} | ||
|
||
// DoS will do S with some modification. | ||
func (m ChaosSimple) DoS(ctx context.Context, ch chan events.Event, s source.S) error { | ||
nCh := make(chan events.Event) | ||
go func() { | ||
for { | ||
select { | ||
case e := <-ch: | ||
if m.isChance(m.FailChance) { | ||
e.Fail(fmt.Errorf("event was unlucky")) | ||
continue | ||
} | ||
|
||
// if its lucky then continue... | ||
nCh <- e | ||
case _ = <-ctx.Done(): | ||
return | ||
} | ||
} | ||
}() | ||
|
||
return s.Do(ctx, nCh) | ||
} | ||
|
||
// DoR will do R with some modification. | ||
func (m ChaosSimple) DoR(ctx context.Context, chOut chan events.Event, chIn chan events.Event, r relay.R) error { | ||
nCh := make(chan events.Event) | ||
|
||
go func() { | ||
for { | ||
select { | ||
case e := <-chOut: | ||
if m.isChance(m.FailChance) { | ||
e.Fail(fmt.Errorf("event was unlucky")) | ||
continue | ||
} | ||
|
||
// if its lucky then continue... | ||
nCh <- e | ||
case _ = <-ctx.Done(): | ||
return | ||
} | ||
} | ||
}() | ||
|
||
return r.Do(ctx, nCh, chIn) | ||
} | ||
|
||
// DoD will do D with some modification. | ||
func (m ChaosSimple) DoD(ctx context.Context, ch chan events.Event, d distributor.D) error { | ||
nCh := make(chan events.Event) | ||
|
||
go func() { | ||
for { | ||
select { | ||
case e := <-ch: | ||
if m.isChance(m.FailChance) { | ||
e.Fail(fmt.Errorf("event was unlucky")) | ||
continue | ||
} | ||
|
||
// if its lucky then continue... | ||
nCh <- e | ||
case _ = <-ctx.Done(): | ||
return | ||
} | ||
} | ||
}() | ||
|
||
return d.Do(ctx, nCh) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,31 @@ | ||
package example | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/andrewneudegg/delta/pkg/distributor" | ||
"github.com/andrewneudegg/delta/pkg/events" | ||
"github.com/andrewneudegg/delta/pkg/meta" | ||
"github.com/andrewneudegg/delta/pkg/relay" | ||
"github.com/andrewneudegg/delta/pkg/source" | ||
) | ||
|
||
// Example doesn't really do anything here. | ||
type Example struct { | ||
meta.M | ||
} | ||
|
||
// DoS will do S with some modification. | ||
func (m Example) DoS(ctx context.Context, ch chan events.Event, s source.S) error { | ||
return s.Do(ctx, ch) | ||
} | ||
|
||
// DoR will do R with some modification. | ||
func (m Example) DoR(ctx context.Context, chOut chan events.Event, chIn chan events.Event, r relay.R) error { | ||
return r.Do(ctx, chOut, chIn) | ||
} | ||
|
||
// DoD will do D with some modification. | ||
func (m Example) DoD(ctx context.Context, ch chan events.Event, d distributor.D) error { | ||
return d.Do(ctx, ch) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,17 @@ | ||
package meta | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/andrewneudegg/delta/pkg/distributor" | ||
"github.com/andrewneudegg/delta/pkg/events" | ||
"github.com/andrewneudegg/delta/pkg/relay" | ||
"github.com/andrewneudegg/delta/pkg/source" | ||
) | ||
|
||
// M defines a block that can augment other functions. | ||
type M interface { | ||
DoS(context.Context, chan events.Event, source.S) error | ||
DoR(context.Context, chan events.Event, chan events.Event, relay.R) error | ||
DoD(context.Context, chan events.Event, distributor.D) error | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters