mirror of
https://github.com/d0zingcat/gocryptotrader.git
synced 2026-05-13 23:16:45 +00:00
* Adds lovely initial concept for historical data doer
* Adds ability to save tasks. Adds config. Adds startStop to engine
* Has a database microservice without use of globals! Further infrastructure design. Adds readme
* Commentary to help design
* Adds migrations for database
* readme and adds database models
* Some modelling that doesn't work end of day
* Completes datahistoryjob sql.Begins datahistoryjobresult
* Adds datahistoryjob functions to retreive job results. Adapts subsystem
* Adds process for upserting jobs and job results to the database
* Broken end of day weird sqlboiler crap
* Fixes issue with SQL generation.
* RPC generation and addition of basic upsert command
* Renames types
* Adds rpc functions
* quick commit before context swithc. Exchanges aren't being populated
* Begin the tests!
* complete sql tests. stop failed jobs. CLI command creation
* Defines rpc commands
* Fleshes out RPC implementation
* Expands testing
* Expands testing, removes double remove
* Adds coverage of data history subsystem, expands errors and nil checks
* Minor logic improvement
* streamlines datahistory test setup
* End of day minor linting
* Lint, convert simplify, rpc expansion, type expansion, readme expansion
* Documentation update
* Renames for consistency
* Completes RPC server commands
* Fixes tests
* Speeds up testing by reducing unnecessary actions. Adds maxjobspercycle config
* Comments for everything
* Adds missing result string. checks interval supported. default start end cli
* Fixes ID problem. Improves binance trade fetch. job ranges are processed
* adds dbservice coverage. adds rpcserver coverage
* docs regen, uses dbcon interface, reverts binance, fixes races, toggle manager
* Speed up tests, remove bad global usage, fix uuid check
* Adds verbose. Updates docs. Fixes postgres
* Minor changes to logging and start stop
* Fixes postgres db tests, fixes postgres column typo
* Fixes old string typo,removes constraint,error parsing for nonreaders
* prevents dhm running when table doesn't exist. Adds prereq documentation
* Adds parallel, rmlines, err fix, comment fix, minor param fixes
* doc regen, common time range check and test updating
* Fixes job validation issues. Updates candle range checker.
* Ensures test cannot fail due to time.Now() shenanigans
* Fixes oopsie, adds documentation and a warn
* Fixes another time test, adjusts copy
* Drastically speeds up data history manager tests via function overrides
* Fixes summary bug and better logs
* Fixes local time test, fixes websocket tests
* removes defaults and comment,updates error messages,sets cli command args
* Fixes FTX trade processing
* Fixes issue where jobs got stuck if data wasn't returned but retrieval was successful
* Improves test speed. Simplifies trade verification SQL. Adds command help
* Fixes the oopsies
* Fixes use of query within transaction. Fixes trade err
* oopsie, not needed
* Adds missing data status. Properly ends job even when data is missing
* errors are more verbose and so have more words to describe them
* Doc regen for new status
* tiny test tinkering
* str := string("Removes .String()").String()
* Merge fixups
* Fixes a data race discovered during github actions
* Allows websocket test to pass consistently
* Fixes merge issue preventing datahistorymanager from starting via config
* Niterinos cmd defaults and explanations
* fixes default oopsie
* Fixes lack of nil protection
* Additional oopsie
* More detailed error for validating job exchange
115 lines
3.1 KiB
Go
115 lines
3.1 KiB
Go
package engine
|
|
|
|
import (
|
|
"fmt"
|
|
"sync/atomic"
|
|
|
|
"github.com/thrasher-corp/gocryptotrader/communications"
|
|
"github.com/thrasher-corp/gocryptotrader/communications/base"
|
|
"github.com/thrasher-corp/gocryptotrader/log"
|
|
)
|
|
|
|
// CommunicationsManagerName is an exported subsystem name
|
|
const CommunicationsManagerName = "communications"
|
|
|
|
// CommunicationManager ensures operations of communications
|
|
type CommunicationManager struct {
|
|
started int32
|
|
shutdown chan struct{}
|
|
relayMsg chan base.Event
|
|
comms *communications.Communications
|
|
}
|
|
|
|
// SetupCommunicationManager creates a communications manager
|
|
func SetupCommunicationManager(cfg *base.CommunicationsConfig) (*CommunicationManager, error) {
|
|
if cfg == nil {
|
|
return nil, errNilConfig
|
|
}
|
|
manager := &CommunicationManager{
|
|
shutdown: make(chan struct{}),
|
|
relayMsg: make(chan base.Event),
|
|
}
|
|
var err error
|
|
manager.comms, err = communications.NewComm(cfg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return manager, nil
|
|
}
|
|
|
|
// IsRunning safely checks whether the subsystem is running
|
|
func (m *CommunicationManager) IsRunning() bool {
|
|
if m == nil {
|
|
return false
|
|
}
|
|
return atomic.LoadInt32(&m.started) == 1
|
|
}
|
|
|
|
// Start runs the subsystem
|
|
func (m *CommunicationManager) Start() error {
|
|
if m == nil {
|
|
return fmt.Errorf("communications manager server %w", ErrNilSubsystem)
|
|
}
|
|
if !atomic.CompareAndSwapInt32(&m.started, 0, 1) {
|
|
return fmt.Errorf("communications manager %w", ErrSubSystemAlreadyStarted)
|
|
}
|
|
log.Debugf(log.CommunicationMgr, "Communications manager %s", MsgSubSystemStarting)
|
|
m.shutdown = make(chan struct{})
|
|
go m.run()
|
|
return nil
|
|
}
|
|
|
|
// GetStatus returns the status of communications
|
|
func (m *CommunicationManager) GetStatus() (map[string]base.CommsStatus, error) {
|
|
if !m.IsRunning() {
|
|
return nil, fmt.Errorf("communications manager %w", ErrSubSystemNotStarted)
|
|
}
|
|
return m.comms.GetStatus(), nil
|
|
}
|
|
|
|
// Stop attempts to shutdown the subsystem
|
|
func (m *CommunicationManager) Stop() error {
|
|
if m == nil {
|
|
return fmt.Errorf("communications manager server %w", ErrNilSubsystem)
|
|
}
|
|
if atomic.LoadInt32(&m.started) == 0 {
|
|
return fmt.Errorf("communications manager %w", ErrSubSystemNotStarted)
|
|
}
|
|
defer func() {
|
|
atomic.CompareAndSwapInt32(&m.started, 1, 0)
|
|
}()
|
|
close(m.shutdown)
|
|
log.Debugf(log.CommunicationMgr, "Communications manager %s", MsgSubSystemShuttingDown)
|
|
return nil
|
|
}
|
|
|
|
// PushEvent pushes an event to the communications relay
|
|
func (m *CommunicationManager) PushEvent(evt base.Event) {
|
|
if !m.IsRunning() {
|
|
return
|
|
}
|
|
select {
|
|
case m.relayMsg <- evt:
|
|
default:
|
|
log.Errorf(log.CommunicationMgr, "Failed to send, no receiver when pushing event [%v]", evt)
|
|
}
|
|
}
|
|
|
|
// run takes awaiting messages and pushes them to be handled by communications
|
|
func (m *CommunicationManager) run() {
|
|
log.Debugf(log.Global, "Communications manager %s", MsgSubSystemStarted)
|
|
defer func() {
|
|
// TO-DO shutdown comms connections for connected services (Slack etc)
|
|
log.Debugf(log.CommunicationMgr, "Communications manager %s", MsgSubSystemShutdown)
|
|
}()
|
|
|
|
for {
|
|
select {
|
|
case msg := <-m.relayMsg:
|
|
m.comms.PushEvent(msg)
|
|
case <-m.shutdown:
|
|
return
|
|
}
|
|
}
|
|
}
|