diff --git a/cli.go b/cli.go index 0db0a69..1e8ad59 100644 --- a/cli.go +++ b/cli.go @@ -10,8 +10,7 @@ import ( "slices" "strings" - dlyrs "github.com/orca-telemetry/core/internal/datalayers" - envs "github.com/orca-telemetry/core/internal/envs" + "github.com/orca-telemetry/core/internal" ) type cliFlags struct { @@ -138,7 +137,7 @@ func validateFlags(flags cliFlags) error { return nil } -func validateConfig(config *envs.Config) error { +func validateConfig(config *internal.Config) error { if config.Platform == "" { return fmt.Errorf("platform cannot be determined from connection string") } @@ -175,8 +174,8 @@ func runCLI(flags cliFlags) { return } - // get singleton configuration - config := envs.GetConfig() + // get singleton configuration + config := internal.GetConfig() // validate configuration if err := validateConfig(config); err != nil { @@ -201,7 +200,7 @@ func runCLI(flags cliFlags) { slog.Info("premigration") if flags.migrate { slog.Info("migrating datalayer", "platform", config.Platform) - err := dlyrs.MigrateDatalayer(config.Platform, config.ConnectionString) + err := MigrateDatalayer(config.Platform, config.ConnectionString) if err != nil { slog.Error("could not migrate the datalayer, exiting", "error", err) os.Exit(1) diff --git a/internal/datalayers/export_test.go b/export_test.go similarity index 91% rename from internal/datalayers/export_test.go rename to export_test.go index 1b79f1e..4193737 100644 --- a/internal/datalayers/export_test.go +++ b/export_test.go @@ -1,7 +1,7 @@ // Bridge package that implements useful mocking functionality for // Orca processors -package datalayers +package main import ( "context" @@ -27,7 +27,7 @@ func (s *mockOrcaProcessorServer) ExecuteDagPart(req *pb.ExecutionRequest, strea slog.Debug("Received ExecuteDagPart request", "exec_id", req.GetExecId()) // simulate processing each algorithm in the request - for i, execution := range req.GetAlgorithmExecutions() { + for _, execution := range req.GetAlgorithmExecutions() { // create a mock result for this algorithm result := &pb.ExecutionResult{ @@ -49,7 +49,6 @@ func (s *mockOrcaProcessorServer) ExecuteDagPart(req *pb.ExecutionRequest, strea return status.Errorf(codes.Internal, "failed to send result: %v", err) } - slog.Debug("sent result for algorithm", "result_num", i+1, "algorithm_num", len(req.GetAlgorithms()), "algorithm_name", execution.GetAlgorithm().GetName()) } slog.Debug("completed ExecuteDagPart", "exec_id", req.GetExecId()) diff --git a/internal/datalayers/main.go b/internal/datalayers/main.go deleted file mode 100644 index 2aaa6a5..0000000 --- a/internal/datalayers/main.go +++ /dev/null @@ -1,54 +0,0 @@ -// Package datalayers provides a factory function for generating a -// datalayer client. Current supported datalayers are: -// - PostgreSQL -package datalayers - -import ( - "context" - "fmt" - "log/slog" - - psql "github.com/orca-telemetry/core/internal/datalayers/postgresql" - types "github.com/orca-telemetry/core/internal/types" -) - -// Platform resprents a database storage platform (e.g. PostgreSQL) -type Platform string - -const ( - // PostgreSQL is the postgresql platform - PostgreSQL Platform = "postgresql" -) - -// check if the platform is supported -func (p Platform) isValid() bool { - switch p { - case PostgreSQL: - return true - default: - return false - } -} - -// NewDatalayerClient generates a new datalayer client of the specificed type. -func NewDatalayerClient( - ctx context.Context, - platform Platform, - connStr string, -) (types.Datalayer, error) { - if !platform.isValid() { - return nil, fmt.Errorf("unsupported platform: %s", platform) - } - - switch platform { - case PostgreSQL: - return psql.NewClient(ctx, connStr) - default: - slog.Error( - "attempted to access unsuported platform", - "platform", - platform, - ) - return nil, fmt.Errorf("platform not implemented: %s", platform) - } -} diff --git a/internal/datalayers/postgresql/docker-compose.yml b/internal/datalayers/postgresql/docker-compose.yml deleted file mode 100755 index f85a95f..0000000 --- a/internal/datalayers/postgresql/docker-compose.yml +++ /dev/null @@ -1,29 +0,0 @@ -services: - postgres: - image: postgres:17-alpine - container_name: orca_postgres - environment: - POSTGRES_USER: orca - POSTGRES_PASSWORD: orca_password - POSTGRES_DB: orca - ports: - - "4747:5432" - volumes: - - postgres_data:/var/lib/postgresql/data - healthcheck: - test: ["CMD-SHELL", "pg_isready -U orca"] - interval: 5s - timeout: 5s - retries: 5 - restart: unless-stopped - networks: - - orca_network - -volumes: - postgres_data: - name: orca_postgres_data - -networks: - orca_network: - name: orca_network - driver: bridge diff --git a/internal/datalayers/postgresql/db.go b/internal/db/db.go similarity index 96% rename from internal/datalayers/postgresql/db.go rename to internal/db/db.go index ae5ac4c..9d485b5 100644 --- a/internal/datalayers/postgresql/db.go +++ b/internal/db/db.go @@ -2,7 +2,7 @@ // versions: // sqlc v1.30.0 -package postgresql +package db import ( "context" diff --git a/internal/types/main.go b/internal/db/errors.go similarity index 53% rename from internal/types/main.go rename to internal/db/errors.go index cefb40b..8c5c8eb 100644 --- a/internal/types/main.go +++ b/internal/db/errors.go @@ -1,26 +1,7 @@ -package datalayers +package db import ( - "context" "fmt" - - pb "github.com/orca-telemetry/contract/go" -) - -// the interface that all datalayers must implement to be compatible with Orca -type ( - Tx interface { - Rollback(ctx context.Context) - Commit(ctx context.Context) error - } - Datalayer interface { - WithTx(ctx context.Context) (Tx, error) - - // Core level operations - RegisterProcessor(ctx context.Context, proc *pb.ProcessorRegistration) error - EmitWindow(ctx context.Context, window *pb.Window) (pb.WindowEmitStatus, error) - Expose(ctx context.Context, settings *pb.ExposeSettings) (*pb.InternalState, error) - } ) // custom errors diff --git a/internal/datalayers/postgresql/funcs.go b/internal/db/funcs.go similarity index 90% rename from internal/datalayers/postgresql/funcs.go rename to internal/db/funcs.go index a50937d..ff43b1a 100644 --- a/internal/datalayers/postgresql/funcs.go +++ b/internal/db/funcs.go @@ -1,4 +1,4 @@ -package postgresql +package db import ( "context" @@ -13,7 +13,6 @@ import ( "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" pb "github.com/orca-telemetry/contract/go" - types "github.com/orca-telemetry/core/internal/types" ) type Datalayer struct { @@ -53,23 +52,13 @@ func NewClient(ctx context.Context, connStr string) (*Datalayer, error) { }, nil } -func (d *Datalayer) WithTx(ctx context.Context) (types.Tx, error) { - tx, err := d.conn.Begin(ctx) - if err != nil { - slog.Error("could not start transaction", "error", err) - return nil, err - } - return &PgTx{tx: tx}, nil -} - func (d *Datalayer) createProcessor( ctx context.Context, - tx types.Tx, + tx pgx.Tx, proc *pb.ProcessorRegistration, ) error { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) err := qtx.CreateProcessor(ctx, CreateProcessorParams{ Name: proc.GetName(), @@ -86,11 +75,10 @@ func (d *Datalayer) createProcessor( func (d *Datalayer) createMetadataField( ctx context.Context, - tx types.Tx, + tx pgx.Tx, metadataField *pb.MetadataField, ) (int64, error) { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) metadataFieldId, err := qtx.CreateMetadataField(ctx, CreateMetadataFieldParams{ Name: metadataField.GetName(), Description: metadataField.GetDescription(), @@ -108,11 +96,10 @@ func (d *Datalayer) createMetadataField( func (d *Datalayer) readMetadataFieldsByWindowType( ctx context.Context, - tx types.Tx, + tx pgx.Tx, windowTypeId int64, ) ([]*pb.MetadataField, error) { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) metadataFields, err := qtx.ReadMetadataFieldsByWindowType(ctx, windowTypeId) @@ -133,11 +120,10 @@ func (d *Datalayer) readMetadataFieldsByWindowType( func (d *Datalayer) createWindowType( ctx context.Context, - tx types.Tx, + tx pgx.Tx, windowType *pb.WindowType, ) (int64, error) { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) windowTypeId, err := qtx.CreateWindowType(ctx, CreateWindowTypeParams{ Name: windowType.GetName(), Version: windowType.GetVersion(), @@ -152,12 +138,11 @@ func (d *Datalayer) createWindowType( func (d *Datalayer) createMetadataFieldBridge( ctx context.Context, - tx types.Tx, + tx pgx.Tx, windowTypeId int64, metadataFieldId int64, ) error { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) err := qtx.CreateWindowTypeMetadataFieldBridge(ctx, CreateWindowTypeMetadataFieldBridgeParams{ WindowTypeID: windowTypeId, MetadataFieldsID: metadataFieldId, @@ -171,12 +156,11 @@ func (d *Datalayer) createMetadataFieldBridge( func (d *Datalayer) addAlgorithm( ctx context.Context, - tx types.Tx, + tx pgx.Tx, algo *pb.Algorithm, proc *pb.ProcessorRegistration, ) error { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) // create algos var resultType ResultType @@ -218,12 +202,12 @@ func (d *Datalayer) addAlgorithm( func (d *Datalayer) addOverwriteAlgorithmDependency( ctx context.Context, - tx types.Tx, + tx pgx.Tx, algo *pb.Algorithm, proc *pb.ProcessorRegistration, ) error { - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) + // get algorithm id algoId, err := qtx.ReadAlgorithmId(ctx, ReadAlgorithmIdParams{ AlgorithmName: algo.GetName(), @@ -264,7 +248,7 @@ func (d *Datalayer) addOverwriteAlgorithmDependency( "to_algo", algo, ) - return &types.CircularDependencyError{ + return &CircularDependencyError{ FromAlgoName: algoDependentOn.GetName(), FromAlgoVersion: algoDependentOn.GetVersion(), FromAlgoProcessor: algoDependentOn.GetProcessorName(), diff --git a/internal/datalayers/postgresql/main.go b/internal/db/main.go similarity index 96% rename from internal/datalayers/postgresql/main.go rename to internal/db/main.go index 6c90b23..fb9b1aa 100644 --- a/internal/datalayers/postgresql/main.go +++ b/internal/db/main.go @@ -1,4 +1,4 @@ -package postgresql +package db import ( "context" @@ -9,6 +9,7 @@ import ( "strconv" "strings" + "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgtype" "google.golang.org/protobuf/types/known/structpb" @@ -16,20 +17,26 @@ import ( "github.com/orca-telemetry/core/internal/dag" ) +func rollbackTransaction(tx pgx.Tx, err *error) { + if p := recover(); p != nil { + tx.Rollback(context.Background()) + *err = fmt.Errorf("panic: %v", p) + return + } + if err != nil && *err != nil { + tx.Rollback(context.Background()) + } +} + // RegisterProcessor with Orca Core func (d *Datalayer) RegisterProcessor( ctx context.Context, proc *pb.ProcessorRegistration, -) (txErr error) { +) (retErr error) { slog.Debug("registering processor", "processor", proc) - tx, err := d.WithTx(ctx) - - defer func() { - if txErr != nil { - tx.Rollback(ctx) - } - }() + tx, err := d.conn.Begin(ctx) + defer rollbackTransaction(tx, &retErr) if err != nil { slog.Error("could not start a transaction", "error", err) @@ -130,29 +137,19 @@ func (d *Datalayer) RegisterProcessor( func (d *Datalayer) EmitWindow( ctx context.Context, window *pb.Window, + useTls bool, ) (_ pb.WindowEmitStatus, retErr error) { slog.Debug("recieved emitted window", "window", window) - tx, err := d.WithTx(ctx) - - defer func() { - if r := recover(); r != nil { - tx.Rollback(ctx) - return - } - - if retErr != nil { - tx.Rollback(ctx) - } - }() + tx, err := d.conn.Begin(ctx) + defer rollbackTransaction(tx, &retErr) if err != nil { slog.Error("could not start a transaction", "error", err) return pb.WindowEmitStatus{}, err } - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) metadata := window.GetMetadata() @@ -367,7 +364,7 @@ func (d *Datalayer) EmitWindow( } go func() { - err := processTasks(d, executionPlan, window, insertedWindow, metadataFilterBytes) + err := processTasks(d, executionPlan, window, insertedWindow, metadataFilterBytes, useTls) if err != nil { err = setWindowStateToFailed(ctx, d.queries, err, insertedWindow.ID) slog.Error("issue processing tasks", "error", err) @@ -388,25 +385,18 @@ func (d *Datalayer) EmitWindow( func (d *Datalayer) Expose( ctx context.Context, settings *pb.ExposeSettings, -) (*pb.InternalState, error) { +) (_ *pb.InternalState, retErr error) { // settings not handled for now - tx, err := d.WithTx(ctx) - - defer func() { - if tx != nil { - tx.Rollback(ctx) - } - }() + tx, err := d.conn.Begin(ctx) + defer rollbackTransaction(tx, &retErr) if err != nil { slog.Error("could not start a transaction", "error", err) return nil, err } - pgTx := tx.(*PgTx) - - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) var processors []Processor if len(settings.ExcludeProject) > 0 { processors, err = qtx.ReadProcessorExcludeProject(ctx, pgtype.Text{ diff --git a/internal/datalayers/postgresql/models.go b/internal/db/models.go similarity index 99% rename from internal/datalayers/postgresql/models.go rename to internal/db/models.go index ad1d3d5..0e2b251 100644 --- a/internal/datalayers/postgresql/models.go +++ b/internal/db/models.go @@ -2,7 +2,7 @@ // versions: // sqlc v1.30.0 -package postgresql +package db import ( "database/sql/driver" diff --git a/internal/datalayers/postgresql/query.sql.go b/internal/db/query.sql.go similarity index 99% rename from internal/datalayers/postgresql/query.sql.go rename to internal/db/query.sql.go index 46c100e..30f56b0 100644 --- a/internal/datalayers/postgresql/query.sql.go +++ b/internal/db/query.sql.go @@ -3,7 +3,7 @@ // sqlc v1.30.0 // source: query.sql -package postgresql +package db import ( "context" diff --git a/internal/datalayers/postgresql/utils.go b/internal/db/utils.go similarity index 97% rename from internal/datalayers/postgresql/utils.go rename to internal/db/utils.go index 61d7c17..2b0f683 100644 --- a/internal/datalayers/postgresql/utils.go +++ b/internal/db/utils.go @@ -1,4 +1,4 @@ -package postgresql +package db import ( "context" @@ -18,7 +18,6 @@ import ( "github.com/orca-telemetry/core/internal/dag" pb "github.com/orca-telemetry/contract/go" - "github.com/orca-telemetry/core/internal/envs" "google.golang.org/grpc" "google.golang.org/grpc/credentials" @@ -103,24 +102,6 @@ func setResultStateToProcessing(ctx context.Context, qtx *Queries, err error, re return oldErr } -func setResultStateToSucceeded(ctx context.Context, qtx *Queries, err error, resultId int64) error { - oldErr := err - err = qtx.UpdateResultState(ctx, UpdateResultStateParams{ - State: AlgorithmStateSUCCEEDED, - ResultID: resultId, - }) - - if err != nil { - slog.Error("could not update result state", "error", err) - if oldErr != nil { - return fmt.Errorf("%w & could not update result state to failed - %w", oldErr, err) - } else { - return fmt.Errorf("could not update result state to failed - %w", err) - } - } - return oldErr -} - // converts a structpb struct to a form that can be used to filter metadata // in the db type MetadataFilter struct { @@ -172,7 +153,8 @@ func processTasks( window *pb.Window, insertedWindow RegisterWindowRow, metadataFilter []byte, -) error { + useTls bool, +) (retErr error) { ctx := context.Background() slog.Info("calculated execution paths", "execution_paths", executionPlan) @@ -214,19 +196,16 @@ func processTasks( algorithmMap[algo.ID] = algo } - // get the environment - config := envs.GetConfig() - // go through the executionPlan and preallocate all the results in the Db - tx, err := d.WithTx(ctx) - defer tx.Rollback(ctx) + tx, err := d.conn.Begin(ctx) + defer rollbackTransaction(tx, &retErr) + if err != nil { slog.Error("could not start a transaction", "error", err) return err } - pgTx := tx.(*PgTx) - qtx := d.queries.WithTx(pgTx.tx) + qtx := d.queries.WithTx(tx) algoIdResultMap := make(map[int64]int64, executionPlan.NumAffectedAlgos) for _, stage := range executionPlan.Stages { @@ -269,7 +248,7 @@ func processTasks( } var conn *grpc.ClientConn - if config.IsProduction { + if useTls { host, _, err := net.SplitHostPort(proc.ConnectionString) if err != nil { host = proc.ConnectionString diff --git a/internal/envs/main.go b/internal/env.go similarity index 99% rename from internal/envs/main.go rename to internal/env.go index 0977d47..46d9f54 100644 --- a/internal/envs/main.go +++ b/internal/env.go @@ -1,4 +1,4 @@ -package envs +package internal import ( "os" diff --git a/internal/main.go b/internal/main.go index 974acd1..ec90893 100644 --- a/internal/main.go +++ b/internal/main.go @@ -6,15 +6,14 @@ import ( "buf.build/go/protovalidate" pb "github.com/orca-telemetry/contract/go" - dlyr "github.com/orca-telemetry/core/internal/datalayers" - types "github.com/orca-telemetry/core/internal/types" + "github.com/orca-telemetry/core/internal/db" "google.golang.org/protobuf/proto" ) type ( OrcaCoreServer struct { pb.UnimplementedOrcaCoreServer - client types.Datalayer + client *db.Datalayer } ) @@ -25,15 +24,12 @@ var ( // NewServer produces a new ORCA gRPC server func NewServer( ctx context.Context, - platform dlyr.Platform, connStr string, ) (*OrcaCoreServer, error) { - client, err := dlyr.NewDatalayerClient(ctx, platform, connStr) + client, err := db.NewClient(ctx, connStr) if err != nil { slog.Error( - "Could not initialise new platform client whilst initialising server", - "platform", - platform, + "could not initialise client", "error", err, ) @@ -94,7 +90,8 @@ func (o *OrcaCoreServer) EmitWindow( return nil, err } slog.Info("emitting window", "window", window) - windowEmitStatus, err := o.client.EmitWindow(ctx, window) + config := GetConfig() + windowEmitStatus, err := o.client.EmitWindow(ctx, window, config.IsProduction) return &windowEmitStatus, err } diff --git a/main.go b/main.go index 201a474..47a2535 100644 --- a/main.go +++ b/main.go @@ -12,7 +12,6 @@ import ( pb "github.com/orca-telemetry/contract/go" orca "github.com/orca-telemetry/core/internal" - dlyr "github.com/orca-telemetry/core/internal/datalayers" ) func startGRPCServer( @@ -21,7 +20,7 @@ func startGRPCServer( port int, _ string, ) { - orcaServer, err := orca.NewServer(context.Background(), dlyr.Platform(platform), dbConnString) + orcaServer, err := orca.NewServer(context.Background(), dbConnString) if err != nil { slog.Error("issue launching Orca Server", "error", err) os.Exit(1) diff --git a/internal/datalayers/main_test.go b/main_test.go similarity index 93% rename from internal/datalayers/main_test.go rename to main_test.go index 6ec53a2..f9ffe11 100644 --- a/internal/datalayers/main_test.go +++ b/main_test.go @@ -1,4 +1,4 @@ -package datalayers +package main import ( "context" @@ -7,7 +7,7 @@ import ( "time" pb "github.com/orca-telemetry/contract/go" - types "github.com/orca-telemetry/core/internal/types" + "github.com/orca-telemetry/core/internal/db" "google.golang.org/protobuf/types/known/structpb" "google.golang.org/protobuf/types/known/timestamppb" @@ -85,8 +85,7 @@ func TestAddProcessor(t *testing.T) { processorConnStr_1 := mockListener_1.Addr().String() processorConnStr_2 := mockListener_2.Addr().String() - // TODO: paramaterise if we have more datalayers (e.g. MySQL, SQLite) - high level function should be the same between them - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) asset_id := pb.MetadataField{Name: "asset_id", Description: "Unique ID of the asset"} @@ -157,7 +156,7 @@ func TestAddProcessor(t *testing.T) { }, }, } - emitStatus, err := dlyr.EmitWindow(testCtx, &window) + emitStatus, err := dlyr.EmitWindow(testCtx, &window, false) assert.Equal(t, emitStatus.GetStatus(), pb.WindowEmitStatus_PROCESSING_TRIGGERED) } @@ -182,8 +181,7 @@ func TestLookbackDependenciesBetweenAlgorithms(t *testing.T) { processorConnStr_1 := mockListener_1.Addr().String() processorConnStr_2 := mockListener_2.Addr().String() - // TODO: paramaterise if we have more datalayers (e.g. MySQL, SQLite) - high level function should be the same between them - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) asset_id := pb.MetadataField{Name: "asset_id", Description: "Unique ID of the asset", Filter: true} @@ -267,7 +265,7 @@ func TestLookbackDependenciesBetweenAlgorithms(t *testing.T) { }, }, } - emitStatus, err := dlyr.EmitWindow(testCtx, &window) + emitStatus, err := dlyr.EmitWindow(testCtx, &window, false) assert.NoError(t, err) assert.Equal(t, emitStatus.GetStatus(), pb.WindowEmitStatus_PROCESSING_TRIGGERED) } @@ -289,8 +287,7 @@ func TestMetadataFieldsChangeable(t *testing.T) { // get the actual address the mock server is listening on processorConnStr := mockListener.Addr().String() - // TODO: paramaterise if we have more datalayers (e.g. MySQL, SQLite) - high level function should be the same between them - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) asset_id := pb.MetadataField{Name: "asset_id", Description: "Unique ID of the asset"} @@ -387,7 +384,7 @@ func TestMetadataFieldsChangeable(t *testing.T) { }, }, } - emitStatus, err := dlyr.EmitWindow(testCtx, &window) + emitStatus, err := dlyr.EmitWindow(testCtx, &window, false) assert.NoError(t, err) assert.Equal(t, pb.WindowEmitStatus_PROCESSING_TRIGGERED, emitStatus.GetStatus()) @@ -413,7 +410,7 @@ func TestMetadataFieldsChangeable(t *testing.T) { }, } // 5. Confirm that the window could not be emitted becuase it is missing the field ID - emitStatus, err = dlyr.EmitWindow(testCtx, &window) + emitStatus, err = dlyr.EmitWindow(testCtx, &window, false) assert.Error(t, err) assert.Equal(t, pb.WindowEmitStatus_TRIGGERING_FAILED, emitStatus.GetStatus()) @@ -440,7 +437,7 @@ func TestMetadataFieldsChangeable(t *testing.T) { } // 7. Confirm that this window is totally different and so not bound by the old metadata fields - emitStatus, err = dlyr.EmitWindow(testCtx, &window) + emitStatus, err = dlyr.EmitWindow(testCtx, &window, false) assert.NoError(t, err) assert.Equal(t, pb.WindowEmitStatus_PROCESSING_TRIGGERED, emitStatus.GetStatus()) } @@ -461,8 +458,7 @@ func TestWindowTypeDefintion(t *testing.T) { // get the actual address the mock server is listening on processorConnStr := mockListener.Addr().String() - // TODO: paramaterise if we have more datalayers (e.g. MySQL, SQLite) - high level function should be the same between them - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) asset_id := pb.MetadataField{Name: "asset_id", Description: "Unique ID of the asset"} @@ -510,7 +506,7 @@ func TestWindowTypeDefintion(t *testing.T) { } func TestCircularDependency(t *testing.T) { - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) windowType := pb.WindowType{ @@ -568,7 +564,7 @@ func TestCircularDependency(t *testing.T) { err = dlyr.RegisterProcessor(testCtx, &proc) - var circularError *types.CircularDependencyError + var circularError *db.CircularDependencyError assert.ErrorAs(t, err, &circularError) assert.Equal(t, algo1.GetName(), circularError.FromAlgoName) @@ -580,7 +576,7 @@ func TestCircularDependency(t *testing.T) { } func TestValidDependenciesBetweenProcessors(t *testing.T) { - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) windowType := pb.WindowType{ @@ -671,13 +667,13 @@ func TestValidDependenciesBetweenProcessors(t *testing.T) { WindowTypeVersion: windowType.GetVersion(), Origin: "Test", } - emitStatus, err := dlyr.EmitWindow(testCtx, window) + emitStatus, err := dlyr.EmitWindow(testCtx, window, false) assert.NoError(t, err) assert.Equal(t, emitStatus.GetStatus(), pb.WindowEmitStatus_PROCESSING_TRIGGERED) } func TestAlgosSameNamesDifferentProcessors(t *testing.T) { - dlyr, err := NewDatalayerClient(testCtx, "postgresql", testConnStr) + dlyr, err := db.NewClient(testCtx, testConnStr) assert.NoError(t, err) windowType := pb.WindowType{ diff --git a/makefile b/makefile index ea0876d..4f8e8e8 100644 --- a/makefile +++ b/makefile @@ -21,8 +21,8 @@ BINARY_NAME = orca export CGO_ENABLED = 0 .datalayer: - sqlc vet -f internal/datalayers/postgresql/sqlc.yaml - sqlc generate -f internal/datalayers/postgresql/sqlc.yaml + sqlc vet -f sqlc.yaml + sqlc generate -f sqlc.yaml .stop_datalayer: cd local_storage && docker-compose stop @@ -59,7 +59,7 @@ export CGO_ENABLED = 0 sudo chmod 640 local_storage/_ca/server.key .test_all: - go test ./internal/... -v + go test ./... -v # ------------- BUILD ------------- diff --git a/internal/datalayers/migrate.go b/migrate.go similarity index 86% rename from internal/datalayers/migrate.go rename to migrate.go index 5da8265..3744cea 100644 --- a/internal/datalayers/migrate.go +++ b/migrate.go @@ -1,4 +1,4 @@ -package datalayers +package main import ( "embed" @@ -10,13 +10,13 @@ import ( "github.com/golang-migrate/migrate/v4/source/iofs" ) -//go:embed postgresql/migrations/*.sql +//go:embed migrations/*.sql var PostgresqlMigrations embed.FS func MigrateDatalayer(platform string, connStr string) error { switch platform { case "postgresql": - d, err := iofs.New(PostgresqlMigrations, "postgresql/migrations") + d, err := iofs.New(PostgresqlMigrations, "migrations") if err != nil { return fmt.Errorf("failed to load embedded migrations: %w", err) } diff --git a/internal/datalayers/postgresql/migrations/000001_initial_migration.down.sql b/migrations/000001_initial_migration.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000001_initial_migration.down.sql rename to migrations/000001_initial_migration.down.sql diff --git a/internal/datalayers/postgresql/migrations/000001_initial_migration.up.sql b/migrations/000001_initial_migration.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000001_initial_migration.up.sql rename to migrations/000001_initial_migration.up.sql diff --git a/internal/datalayers/postgresql/migrations/000002_added_annotations.down.sql b/migrations/000002_added_annotations.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000002_added_annotations.down.sql rename to migrations/000002_added_annotations.down.sql diff --git a/internal/datalayers/postgresql/migrations/000002_added_annotations.up.sql b/migrations/000002_added_annotations.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000002_added_annotations.up.sql rename to migrations/000002_added_annotations.up.sql diff --git a/internal/datalayers/postgresql/migrations/000003_added_metadata_fields.down.sql b/migrations/000003_added_metadata_fields.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000003_added_metadata_fields.down.sql rename to migrations/000003_added_metadata_fields.down.sql diff --git a/internal/datalayers/postgresql/migrations/000003_added_metadata_fields.up.sql b/migrations/000003_added_metadata_fields.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000003_added_metadata_fields.up.sql rename to migrations/000003_added_metadata_fields.up.sql diff --git a/internal/datalayers/postgresql/migrations/000004_modified_algorithm_constraint.down.sql b/migrations/000004_modified_algorithm_constraint.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000004_modified_algorithm_constraint.down.sql rename to migrations/000004_modified_algorithm_constraint.down.sql diff --git a/internal/datalayers/postgresql/migrations/000004_modified_algorithm_constraint.up.sql b/migrations/000004_modified_algorithm_constraint.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000004_modified_algorithm_constraint.up.sql rename to migrations/000004_modified_algorithm_constraint.up.sql diff --git a/internal/datalayers/postgresql/migrations/000005_add_description_to_algorithm.down.sql b/migrations/000005_add_description_to_algorithm.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000005_add_description_to_algorithm.down.sql rename to migrations/000005_add_description_to_algorithm.down.sql diff --git a/internal/datalayers/postgresql/migrations/000005_add_description_to_algorithm.up.sql b/migrations/000005_add_description_to_algorithm.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000005_add_description_to_algorithm.up.sql rename to migrations/000005_add_description_to_algorithm.up.sql diff --git a/internal/datalayers/postgresql/migrations/000006_added_project_name_to_processors.down.sql b/migrations/000006_added_project_name_to_processors.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000006_added_project_name_to_processors.down.sql rename to migrations/000006_added_project_name_to_processors.down.sql diff --git a/internal/datalayers/postgresql/migrations/000006_added_project_name_to_processors.up.sql b/migrations/000006_added_project_name_to_processors.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000006_added_project_name_to_processors.up.sql rename to migrations/000006_added_project_name_to_processors.up.sql diff --git a/internal/datalayers/postgresql/migrations/000007_add_lookback_requirements_to_dependencies.down.sql b/migrations/000007_add_lookback_requirements_to_dependencies.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000007_add_lookback_requirements_to_dependencies.down.sql rename to migrations/000007_add_lookback_requirements_to_dependencies.down.sql diff --git a/internal/datalayers/postgresql/migrations/000007_add_lookback_requirements_to_dependencies.up.sql b/migrations/000007_add_lookback_requirements_to_dependencies.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000007_add_lookback_requirements_to_dependencies.up.sql rename to migrations/000007_add_lookback_requirements_to_dependencies.up.sql diff --git a/internal/datalayers/postgresql/migrations/000008_modify_dependency_execution_paths_view_for_lookback.down.sql b/migrations/000008_modify_dependency_execution_paths_view_for_lookback.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000008_modify_dependency_execution_paths_view_for_lookback.down.sql rename to migrations/000008_modify_dependency_execution_paths_view_for_lookback.down.sql diff --git a/internal/datalayers/postgresql/migrations/000008_modify_dependency_execution_paths_view_for_lookback.up.sql b/migrations/000008_modify_dependency_execution_paths_view_for_lookback.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000008_modify_dependency_execution_paths_view_for_lookback.up.sql rename to migrations/000008_modify_dependency_execution_paths_view_for_lookback.up.sql diff --git a/internal/datalayers/postgresql/migrations/000009_add_lookback_requirements_to_algorithms.down.sql b/migrations/000009_add_lookback_requirements_to_algorithms.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000009_add_lookback_requirements_to_algorithms.down.sql rename to migrations/000009_add_lookback_requirements_to_algorithms.down.sql diff --git a/internal/datalayers/postgresql/migrations/000009_add_lookback_requirements_to_algorithms.up.sql b/migrations/000009_add_lookback_requirements_to_algorithms.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000009_add_lookback_requirements_to_algorithms.up.sql rename to migrations/000009_add_lookback_requirements_to_algorithms.up.sql diff --git a/internal/datalayers/postgresql/migrations/000010_add_lookback_gap.down.sql b/migrations/000010_add_lookback_gap.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000010_add_lookback_gap.down.sql rename to migrations/000010_add_lookback_gap.down.sql diff --git a/internal/datalayers/postgresql/migrations/000010_add_lookback_gap.up.sql b/migrations/000010_add_lookback_gap.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000010_add_lookback_gap.up.sql rename to migrations/000010_add_lookback_gap.up.sql diff --git a/internal/datalayers/postgresql/migrations/000011_move_metadata_to_column.down.sql b/migrations/000011_move_metadata_to_column.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000011_move_metadata_to_column.down.sql rename to migrations/000011_move_metadata_to_column.down.sql diff --git a/internal/datalayers/postgresql/migrations/000011_move_metadata_to_column.up.sql b/migrations/000011_move_metadata_to_column.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000011_move_metadata_to_column.up.sql rename to migrations/000011_move_metadata_to_column.up.sql diff --git a/internal/datalayers/postgresql/migrations/000012_add_filter_keys_to_window_type.down.sql b/migrations/000012_add_filter_keys_to_window_type.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000012_add_filter_keys_to_window_type.down.sql rename to migrations/000012_add_filter_keys_to_window_type.down.sql diff --git a/internal/datalayers/postgresql/migrations/000012_add_filter_keys_to_window_type.up.sql b/migrations/000012_add_filter_keys_to_window_type.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000012_add_filter_keys_to_window_type.up.sql rename to migrations/000012_add_filter_keys_to_window_type.up.sql diff --git a/internal/datalayers/postgresql/migrations/000013_add_window_status.down.sql b/migrations/000013_add_window_status.down.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000013_add_window_status.down.sql rename to migrations/000013_add_window_status.down.sql diff --git a/internal/datalayers/postgresql/migrations/000013_add_window_status.up.sql b/migrations/000013_add_window_status.up.sql similarity index 100% rename from internal/datalayers/postgresql/migrations/000013_add_window_status.up.sql rename to migrations/000013_add_window_status.up.sql diff --git a/internal/datalayers/postgresql/query.sql b/query.sql similarity index 100% rename from internal/datalayers/postgresql/query.sql rename to query.sql diff --git a/internal/datalayers/postgresql/sqlc.yaml b/sqlc.yaml similarity index 73% rename from internal/datalayers/postgresql/sqlc.yaml rename to sqlc.yaml index 2b52251..a30b009 100755 --- a/internal/datalayers/postgresql/sqlc.yaml +++ b/sqlc.yaml @@ -5,6 +5,6 @@ sql: schema: "migrations" gen: go: - package: "postgresql" - out: "./" + package: "db" + out: "./internal/db/" sql_package: "pgx/v5"