Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

bugfix(cli): CSV import via remote marketstore connection #515

Merged
merged 2 commits into from
Oct 6, 2021
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 10 additions & 7 deletions cmd/connect/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ package connect
import (
"errors"

"github.com/alpacahq/marketstore/v4/frontend/client"

"github.com/spf13/cobra"

"github.com/alpacahq/marketstore/v4/cmd/connect/session"
Expand Down Expand Up @@ -86,7 +88,13 @@ func executeConnect(cmd *cobra.Command, args []string) error {

// Attempt remote mode.
if len(url) != 0 {
conn, err = session.NewRemoteAPIClient(url)
// Attempt connection to remote host.
rpcClient, err := client.NewClient(url)
if err != nil {
return err
}

conn, err = session.NewRemoteAPIClient(url, rpcClient)
if err != nil {
return err
}
Expand All @@ -96,13 +104,8 @@ func executeConnect(cmd *cobra.Command, args []string) error {
utils.InstanceConfig.DisableVariableCompression = true
}

c = session.NewClient(conn)

// Initialize connection.
err = c.Connect()
if err != nil {
return err
}
c = session.NewClient(conn)

// Enter command loop
err = c.Read()
Expand Down
6 changes: 0 additions & 6 deletions cmd/connect/session/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,6 @@ type Client struct {
type APIClient interface {
// PrintConnectInfo prints connection information to stdout.
PrintConnectInfo()
// Connect initializes a client connection.
Connect() error
// Create creates a new bucket in the marketstore server
Create(reqs *frontend.MultiCreateRequest, responses *frontend.MultiServerResponse) error
// Write executes a write operation to the marketstore server.
Expand All @@ -60,10 +58,6 @@ type APIClient interface {
SQL(line string) (cs *dbio.ColumnSeries, err error)
}

func (c *Client) Connect() error {
return c.apiClient.Connect()
}

// RPCClient is a marketstore API client interface.
type RPCClient interface {
DoRPC(functionName string, args interface{}) (response interface{}, err error)
Expand Down
4 changes: 0 additions & 4 deletions cmd/connect/session/local_api_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,10 +49,6 @@ type LocalAPIClient struct {
func (lc *LocalAPIClient) PrintConnectInfo() {
fmt.Fprintf(os.Stderr, "Connected to local instance at path: %v\n", lc.dir)
}
func (lc *LocalAPIClient) Connect() error {
// Nothing to do here yet..
return nil
}

func (lc *LocalAPIClient) Write(reqs *frontend.MultiWriteRequest, responses *frontend.MultiServerResponse) error {
ds := frontend.NewDataService(lc.dir, lc.catalogDir, lc.aggRunner, lc.writer, lc.query)
Expand Down
53 changes: 19 additions & 34 deletions cmd/connect/session/mock/client.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

92 changes: 61 additions & 31 deletions cmd/connect/session/remote_api_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,10 @@ import (
"github.com/alpacahq/marketstore/v4/frontend"
"github.com/alpacahq/marketstore/v4/planner"
"github.com/alpacahq/marketstore/v4/utils/io"

"github.com/alpacahq/marketstore/v4/frontend/client"
)

// NewRemoteAPIClient generates a new client struct.
func NewRemoteAPIClient(url string) (rc *RemoteAPIClient, err error) {
func NewRemoteAPIClient(url string, client RPCClient) (rc *RemoteAPIClient, err error) {
// TODO: validate url using go core packages.
splits := strings.Split(url, ":")
if len(splits) != 2 {
Expand All @@ -24,7 +22,7 @@ func NewRemoteAPIClient(url string) (rc *RemoteAPIClient, err error) {
}
// build url.
url = "http://" + url
return &RemoteAPIClient{url: url}, nil
return &RemoteAPIClient{url: url, rpcClient: client}, nil
}

// RemoteAPIClient represents an agent that manages a database
Expand All @@ -40,25 +38,22 @@ type RemoteAPIClient struct {
func (rc *RemoteAPIClient) PrintConnectInfo() {
fmt.Fprintf(os.Stderr, "Connected to remote instance at: %v\n", rc.url)
}
func (rc *RemoteAPIClient) Connect() error {
// Attempt connection to remote host.
cli, err := client.NewClient(rc.url)
if err != nil {
return err
}
rc.rpcClient = cli

// Success.
return nil
}

func (rc *RemoteAPIClient) Write(reqs *frontend.MultiWriteRequest, responses *frontend.MultiServerResponse) error {
var respI interface{}
respI, err := rc.rpcClient.DoRPC("Write", reqs)
if err != nil {
return fmt.Errorf("DoRPC:Write error:%w", err)
}

if respI != nil {
responses = respI.(*frontend.MultiServerResponse)
if val, ok := respI.(*frontend.MultiServerResponse); ok {
*responses = *val
} else {
return fmt.Errorf("[bug] unexpected data type returned from DoRPC:Write func. resp=%v", respI)
}
}
return err
return nil
}

func (rc *RemoteAPIClient) Show(tbk *io.TimeBucketKey, start, end *time.Time) (csm io.ColumnSeriesMap, err error) {
Expand All @@ -79,41 +74,72 @@ func (rc *RemoteAPIClient) Show(tbk *io.TimeBucketKey, start, end *time.Time) (c
Requests: []frontend.QueryRequest{req},
}

resp, err := rc.rpcClient.DoRPC("Query", args)
respI, err := rc.rpcClient.DoRPC("Query", args)
if err != nil {
return nil, err
return nil, fmt.Errorf("DoRPC:Query error:%w", err)
}

if respI == nil {
return io.ColumnSeriesMap{}, nil
}

return *resp.(*io.ColumnSeriesMap), nil
if val, ok := respI.(*io.ColumnSeriesMap); ok {
return *val, nil
} else {
return nil, fmt.Errorf("[bug] unexpected data type returned from DoRPC:Query func. resp=%v",
respI,
)
}
}

func (rc *RemoteAPIClient) Create(reqs *frontend.MultiCreateRequest, responses *frontend.MultiServerResponse) error {
var respI interface{}
respI, err := rc.rpcClient.DoRPC("Create", reqs)
if err != nil {
return fmt.Errorf("DoRPC:Create error:%w", err)
}
if respI != nil {
responses = respI.(*frontend.MultiServerResponse)
if val, ok := respI.(*frontend.MultiServerResponse); ok {
*responses = *val
} else {
return fmt.Errorf("[bug] unexpected data type returned from DoRPC:Create func. resp=%v", respI)
}
}
return err

return nil
}

func (rc *RemoteAPIClient) Destroy(reqs *frontend.MultiKeyRequest, responses *frontend.MultiServerResponse) error {
var respI interface{}
respI, err := rc.rpcClient.DoRPC("Destroy", reqs)
if err != nil {
return fmt.Errorf("DoRPC:Destroy error:%w", err)
}
if respI != nil {
responses = respI.(*frontend.MultiServerResponse)
if val, ok := respI.(*frontend.MultiServerResponse); ok {
*responses = *val
} else {
return fmt.Errorf("[bug] unexpected data type returned from DoRPC:Destroy func. resp=%v", respI)
}
}
return err
return nil
}

func (rc *RemoteAPIClient) GetBucketInfo(reqs *frontend.MultiKeyRequest, responses *frontend.MultiGetInfoResponse,
) error {
var respI interface{}
respI, err := rc.rpcClient.DoRPC("GetInfo", reqs)
if err != nil {
return fmt.Errorf("DoRPC:GetBucketInfo error:%w", err)
}

if respI != nil {
responses = respI.(*frontend.MultiGetInfoResponse)
if val, ok := respI.(*frontend.MultiGetInfoResponse); ok {
*responses = *val
} else {
return fmt.Errorf("[bug] unexpected data type returned from DoRPC:GetBucketInfo func. resp=%v", respI)
}
}
return err
return nil
}

func (rc *RemoteAPIClient) SQL(line string) (cs *io.ColumnSeries, err error) {
Expand All @@ -125,12 +151,16 @@ func (rc *RemoteAPIClient) SQL(line string) (cs *io.ColumnSeries, err error) {

resp, err := rc.rpcClient.DoRPC("Query", args)
if err != nil {
return nil, err
return nil, fmt.Errorf("DoRPC:Query SQL error: %w", err)
}

for _, sub := range *resp.(*io.ColumnSeriesMap) {
cs = sub
break
if val, ok := resp.(*io.ColumnSeriesMap); ok {
for _, sub := range *val {
cs = sub
break
}
} else {
return nil, fmt.Errorf("[bug] unexpected data type returned from DoRPC:SQL func. resp=%v", resp)
}
return cs, err
}
Loading