diff --git a/.github/workflows/go.yml b/.github/workflows/go.yml index b1973199c..8c5726a84 100644 --- a/.github/workflows/go.yml +++ b/.github/workflows/go.yml @@ -65,7 +65,7 @@ jobs: PKGS=$(go list ./pkg/... | grep -v '/monitoring$' | tr '\n' ' ') # Run tests with coverage # Use -coverpkg to measure client package coverage through server tests - go test -p 1 -race -coverprofile=coverage.raw -covermode=atomic -coverpkg=./pkg/... $PKGS + go test -race -coverprofile=coverage.raw -covermode=atomic -coverpkg=./pkg/... $PKGS # Filter out test utilities, mocks, and monitoring from coverage report grep -v 'pkg/server/test_utils.go' coverage.raw | \ grep -v 'pkg/monitoring/' | \ diff --git a/CLAUDE.md b/CLAUDE.md index 409246a47..ca46caa27 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -26,10 +26,16 @@ The deployment docker-compose files expect `colonyos/colonies:latest`. Using a d ### Testing ```bash -make test # Run all tests (requires grc for colored output) -make github_test # Run tests without grc (for CI) +make test # Run all tests: needs Postgres on localhost:5432 (make startdb) + # and an S3 server on localhost:9000 for pkg/fs +make github_test # Alias for test (used by CI) ``` +Set a timezone when running tests, e.g. `export TZ=Europe/Stockholm` (the +database layer refuses to connect without TZ, and one core test expects a +non-UTC zone). Test packages use per-process databases and dynamic ports, so +they run in parallel. + ### Development Environment ```bash docker-compose up -d # Start Colonies server with dependencies (TimescaleDB, MinIO) diff --git a/Makefile b/Makefile index a644d3787..19490b69f 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,5 @@ all: build -.PHONY: all build +.PHONY: all build container container-multiplatform container-multiplatform-push push coverage test test-all test-compat github_test install startdb nukedb BUILD_IMAGE ?= colonyos/colonies PUSH_IMAGE ?= colonyos/colonies:v1.9.10 @@ -34,55 +34,18 @@ push: docker push $(PUSH_IMAGE) coverage: - ./buildtools/coverage.sh - ./buildtools/codecov + @go test -coverprofile=coverage.txt -covermode=atomic ./... build_cryptolib_ubuntu_2020: cd buildtools; ./build_cryptolib_ubuntu.sh +# Runs all tests: needs Postgres on localhost:5432 (make startdb) and an S3 +# server on localhost:9000 for pkg/fs. Test packages use per-process databases +# and dynamic ports, so they run in parallel. test: - @cd tests/reliability; go test -v --race - @cd internal/crypto; go test -v --race - @cd pkg/core; go test -v --race - @cd pkg/database/postgresql; go test -v --race - @cd pkg/rpc; go test -v --race - @cd pkg/security; go test -v --race - @cd pkg/security/crypto; go test -v --race - @cd pkg/security/validator; go test -v --race - @cd pkg/backends/gin; go test -v --race - @cd pkg/client; go test -v --race - @cd pkg/client/backends; go test -v --race - @cd pkg/client/gin; go test -v --race - @cd pkg/server; go test -v --race - @cd pkg/server/controllers; go test -v --race - @cd pkg/server/handlers/attribute; go test -v --race - @cd pkg/server/handlers/blueprint; go test -v --race - @cd pkg/server/handlers/channel; go test -v --race - @cd pkg/server/handlers/colony; go test -v --race - @cd pkg/server/handlers/cron; go test -v --race - @cd pkg/server/handlers/executor; go test -v --race - @cd pkg/server/handlers/file; go test -v --race - @cd pkg/server/handlers/function; go test -v --race - @cd pkg/server/handlers/generator; go test -v --race - @cd pkg/server/handlers/location; go test -v --race - @cd pkg/server/handlers/log; go test -v --race - @cd pkg/server/handlers/process; go test -v --race - @cd pkg/server/handlers/processgraph; go test -v --race - @cd pkg/server/handlers/security; go test -v --race - @cd pkg/server/handlers/server; go test -v --race - @cd pkg/server/handlers/snapshot; go test -v --race - @cd pkg/server/handlers/user; go test -v --race - @cd pkg/server/handlers/realtime; go test -v --race - @cd pkg/server/registry; go test -v --race - @cd pkg/server/utils; go test -v --race - @cd pkg/scheduler; go test -v --race - @cd pkg/parsers; go test -v --race - @cd pkg/utils; go test -v --race - @cd pkg/validate; go test -v --race - @cd pkg/channel; go test -v --race - @cd pkg/cluster; go test -v --race - @cd pkg/cron; go test -v --race - @cd pkg/fs; go test -v --race + @go test -race ./... + +github_test: test install: cp ./bin/colonies /usr/local/bin diff --git a/internal/cryptolib.wasm/cryptolib.go b/internal/cryptolib.wasm/cryptolib.go index 59a694e3f..9e12e9390 100644 --- a/internal/cryptolib.wasm/cryptolib.go +++ b/internal/cryptolib.wasm/cryptolib.go @@ -1,3 +1,5 @@ +//go:build js && wasm + package main import ( diff --git a/pkg/cluster/config.go b/pkg/cluster/config.go index 792e4cd78..2420a9463 100644 --- a/pkg/cluster/config.go +++ b/pkg/cluster/config.go @@ -73,7 +73,7 @@ func (config *Config) Equals(config2 *Config) bool { } } - if counter == len(config.Nodes) && counter == len(config.Nodes) { + if counter == len(config.Nodes) && counter == len(config2.Nodes) { return true } diff --git a/pkg/cluster/etcd_test.go b/pkg/cluster/etcd_test.go index 9ebaf0211..efdb29540 100644 --- a/pkg/cluster/etcd_test.go +++ b/pkg/cluster/etcd_test.go @@ -1,6 +1,7 @@ package cluster import ( + "fmt" "os" "testing" @@ -8,10 +9,8 @@ import ( ) func TestCreateEtcdCluster(t *testing.T) { - node1 := Node{Name: "etcd1", Host: "localhost", EtcdClientPort: 24100, EtcdPeerPort: 23100, RelayPort: 25100, APIPort: 26100} - node2 := Node{Name: "etcd2", Host: "localhost", EtcdClientPort: 24200, EtcdPeerPort: 23200, RelayPort: 25200, APIPort: 26200} - node3 := Node{Name: "etcd3", Host: "localhost", EtcdClientPort: 24300, EtcdPeerPort: 23300, RelayPort: 25300, APIPort: 26300} - node4 := Node{Name: "etcd4", Host: "localhost", EtcdClientPort: 24400, EtcdPeerPort: 23400, RelayPort: 25400, APIPort: 26400} + nodes := testClusterNodes(t, "etcd1", "etcd2", "etcd3", "etcd4") + node1, node2, node3, node4 := nodes[0], nodes[1], nodes[2], nodes[3] config := Config{} config.AddNode(node1) @@ -19,12 +18,15 @@ func TestCreateEtcdCluster(t *testing.T) { config.AddNode(node3) config.AddNode(node4) - server1 := CreateEtcdServer(node1, config, ".") - server2 := CreateEtcdServer(node2, config, ".") - server3 := CreateEtcdServer(node3, config, ".") - server4 := CreateEtcdServer(node4, config, ".") + dataPath := t.TempDir() + server1 := CreateEtcdServer(node1, config, dataPath) + server2 := CreateEtcdServer(node2, config, dataPath) + server3 := CreateEtcdServer(node3, config, dataPath) + server4 := CreateEtcdServer(node4, config, dataPath) - assert.Equal(t, server1.buildInitialClusterStr(), "etcd1=http://localhost:23100,etcd2=http://localhost:23200,etcd3=http://localhost:23300,etcd4=http://localhost:23400") + expectedClusterStr := fmt.Sprintf("etcd1=http://localhost:%d,etcd2=http://localhost:%d,etcd3=http://localhost:%d,etcd4=http://localhost:%d", + node1.EtcdPeerPort, node2.EtcdPeerPort, node3.EtcdPeerPort, node4.EtcdPeerPort) + assert.Equal(t, server1.buildInitialClusterStr(), expectedClusterStr) server1.Start() server2.Start() @@ -70,11 +72,11 @@ func TestCreateEtcdCluster(t *testing.T) { } func TestEtcdAssignmentsPauseResume(t *testing.T) { - node := Node{Name: "etcd1", Host: "localhost", EtcdClientPort: 24500, EtcdPeerPort: 23500, RelayPort: 25500, APIPort: 26500} + node := testClusterNodes(t, "etcd1")[0] config := Config{} config.AddNode(node) - server := CreateEtcdServer(node, config, ".") + server := CreateEtcdServer(node, config, t.TempDir()) server.Start() server.WaitToStart() @@ -123,11 +125,11 @@ func TestEtcdAssignmentsPauseResume(t *testing.T) { } func TestEtcdAssignmentsPauseResumeWithoutClient(t *testing.T) { - node := Node{Name: "etcd2", Host: "localhost", EtcdClientPort: 24600, EtcdPeerPort: 23600, RelayPort: 25600, APIPort: 26600} + node := testClusterNodes(t, "etcd2")[0] config := Config{} config.AddNode(node) - server := CreateEtcdServer(node, config, ".") + server := CreateEtcdServer(node, config, t.TempDir()) colonyName := "test_colony" // Test methods fail when etcd client is not initialized @@ -146,15 +148,16 @@ func TestEtcdAssignmentsPauseResumeWithoutClient(t *testing.T) { } func TestEtcdAssignmentsPauseResumeMultiNode(t *testing.T) { - node1 := Node{Name: "etcd1", Host: "localhost", EtcdClientPort: 24700, EtcdPeerPort: 23700, RelayPort: 25700, APIPort: 26700} - node2 := Node{Name: "etcd2", Host: "localhost", EtcdClientPort: 24800, EtcdPeerPort: 23800, RelayPort: 25800, APIPort: 26800} + nodes := testClusterNodes(t, "etcd1", "etcd2") + node1, node2 := nodes[0], nodes[1] config := Config{} config.AddNode(node1) config.AddNode(node2) - server1 := CreateEtcdServer(node1, config, ".") - server2 := CreateEtcdServer(node2, config, ".") + dataPath := t.TempDir() + server1 := CreateEtcdServer(node1, config, dataPath) + server2 := CreateEtcdServer(node2, config, dataPath) server1.Start() server2.Start() diff --git a/pkg/cluster/relay_server_test.go b/pkg/cluster/relay_server_test.go index 11986e420..53ab4b213 100644 --- a/pkg/cluster/relay_server_test.go +++ b/pkg/cluster/relay_server_test.go @@ -12,9 +12,8 @@ func TestRelayServer(t *testing.T) { gin.SetMode(gin.ReleaseMode) gin.DefaultWriter = ioutil.Discard - node1 := Node{Name: "etcd1", Host: "localhost", EtcdClientPort: 24100, EtcdPeerPort: 23100, RelayPort: 25100, APIPort: 26100} - node2 := Node{Name: "etcd2", Host: "localhost", EtcdClientPort: 24200, EtcdPeerPort: 23200, RelayPort: 25200, APIPort: 26200} - node3 := Node{Name: "etcd3", Host: "localhost", EtcdClientPort: 24300, EtcdPeerPort: 23300, RelayPort: 25300, APIPort: 26300} + nodes := testClusterNodes(t, "etcd1", "etcd2", "etcd3") + node1, node2, node3 := nodes[0], nodes[1], nodes[2] config := Config{} config.AddNode(node1) diff --git a/pkg/cluster/test_helpers_test.go b/pkg/cluster/test_helpers_test.go new file mode 100644 index 000000000..8ef133b86 --- /dev/null +++ b/pkg/cluster/test_helpers_test.go @@ -0,0 +1,30 @@ +package cluster + +import ( + "testing" + + "github.com/colonyos/colonies/pkg/utils" +) + +// testClusterNodes returns n cluster nodes with kernel-assigned free ports so +// tests never collide on fixed port numbers. +func testClusterNodes(t *testing.T, names ...string) []Node { + t.Helper() + ports, err := utils.FreePorts(4 * len(names)) + if err != nil { + t.Fatalf("failed to allocate ports: %v", err) + } + + nodes := make([]Node, len(names)) + for i, name := range names { + nodes[i] = Node{ + Name: name, + Host: "localhost", + EtcdClientPort: ports[4*i], + EtcdPeerPort: ports[4*i+1], + RelayPort: ports[4*i+2], + APIPort: ports[4*i+3], + } + } + return nodes +} diff --git a/pkg/constants/constants.go b/pkg/constants/constants.go index 20927eb8a..212a0da68 100644 --- a/pkg/constants/constants.go +++ b/pkg/constants/constants.go @@ -8,7 +8,6 @@ const MAX_LOG_COUNT = 500 // Maximum number of log entries that can be reque // Test Configuration - Default values used in test environments const TESTHOST = "localhost" // Default hostname for test servers -const TESTPORT = 28088 // Default port for test servers // Background Processing Periods - How frequently various system tasks run const RELEASE_PERIOD = 1 // Period in seconds when processes are checked for max exec time or max wait time diff --git a/pkg/database/postgresql/test_utils.go b/pkg/database/postgresql/test_utils.go index 31465ba2e..cca908104 100644 --- a/pkg/database/postgresql/test_utils.go +++ b/pkg/database/postgresql/test_utils.go @@ -1,21 +1,138 @@ package postgresql import ( - "io/ioutil" + "database/sql" + "fmt" + "io" "log" - "math/rand" "os" - "time" + "strings" + "sync" ) +// testTables lists every table created by Initialize. Kept in sync with +// Initialize so test cleanup can truncate instead of recreating the schema. +var testTables = []string{ + "USERS", + "SERVER", + "COLONIES", + "EXECUTORS", + "NODES", + "FUNCTIONS", + "PROCESSES", + "LOGS", + "FILES", + "SNAPSHOTS", + "ATTRIBUTES", + "PROCESSGRAPHS", + "GENERATORS", + "GENERATORARGS", + "CRONS", + "BLUEPRINTDEFINITIONS", + "BLUEPRINTS", + "LOCATIONS", + "BLUEPRINT_HISTORY", +} + +var testSchemas = struct { + sync.Mutex + ready map[string]bool +}{ready: make(map[string]bool)} + +var testDatabase struct { + once sync.Once + name string + err error +} + +// pinnedTestConn holds one open connection to this process's test database for +// the lifetime of the process, so the database shows up in pg_stat_activity +// and is never considered stale by other test processes. +var pinnedTestConn *sql.DB + +// testDatabaseName returns the name of a database dedicated to this test +// process, creating it on first use. Per-process databases let test packages +// run in parallel against one Postgres instance without interfering. Stale +// databases from finished or crashed runs are dropped opportunistically: a +// database is stale only if it has no active connections AND the process +// that created it (encoded in the name) is no longer alive on this host. +func testDatabaseName(host string, port int, user string, password string) (string, error) { + testDatabase.once.Do(func() { + adminDSN := fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=postgres sslmode=disable", host, port, user, password) + admin, err := sql.Open("postgres", adminDSN) + if err != nil { + testDatabase.err = err + return + } + defer admin.Close() + + rows, err := admin.Query(`SELECT datname FROM pg_database WHERE datname LIKE 'colonies_test_%' AND datname NOT IN (SELECT datname FROM pg_stat_activity WHERE datname IS NOT NULL)`) + if err == nil { + var stale []string + for rows.Next() { + var name string + if rows.Scan(&name) == nil && !testDatabaseOwnerAlive(name) { + stale = append(stale, name) + } + } + rows.Close() + for _, name := range stale { + // Best effort; a database that just became active is skipped + admin.Exec(`DROP DATABASE IF EXISTS ` + name) + } + } + + name := fmt.Sprintf("colonies_test_%d", os.Getpid()) + if _, err := admin.Exec(`DROP DATABASE IF EXISTS ` + name); err != nil { + testDatabase.err = err + return + } + if _, err := admin.Exec(`CREATE DATABASE ` + name); err != nil { + testDatabase.err = err + return + } + + // Pin a connection so the new database is visible as in-use + pinnedDSN := fmt.Sprintf("host=%s port=%d user=%s password=%s dbname=%s sslmode=disable", host, port, user, password, name) + pinnedTestConn, err = sql.Open("postgres", pinnedDSN) + if err == nil { + pinnedTestConn.SetMaxIdleConns(1) + pinnedTestConn.SetConnMaxLifetime(0) + pinnedTestConn.SetConnMaxIdleTime(0) + testDatabase.err = pinnedTestConn.Ping() + if testDatabase.err != nil { + return + } + } + + testDatabase.name = name + }) + + return testDatabase.name, testDatabase.err +} + +// testDatabaseOwnerAlive reports whether the process that created a test +// database (pid encoded in the database name) is still running on this host. +func testDatabaseOwnerAlive(datname string) bool { + var pid int + if _, err := fmt.Sscanf(datname, "colonies_test_%d", &pid); err != nil { + // Unrecognized name; leave it alone + return true + } + _, err := os.Stat(fmt.Sprintf("/proc/%d", pid)) + return err == nil +} + func PrepareTests() (*PQDatabase, error) { return PrepareTestsWithPrefix("TEST_") } +// PrepareTestsWithPrefix returns a connected database with an empty schema. +// The schema is dropped and recreated only on the first call per process and +// prefix; subsequent calls truncate all tables, which is roughly an order of +// magnitude faster than Drop+Initialize. func PrepareTestsWithPrefix(prefix string) (*PQDatabase, error) { - log.SetOutput(ioutil.Discard) - - rand.Seed(time.Now().UTC().UnixNano()) + log.SetOutput(io.Discard) dbHost := os.Getenv("COLONIES_DB_HOST") if dbHost == "" { @@ -30,18 +147,59 @@ func PrepareTestsWithPrefix(prefix string) (*PQDatabase, error) { if dbPassword == "" { dbPassword = "rFcLGNkgsNtksg6Pgtn9CumL4xXBQ7" } - dbName := "postgres" - dbPrefix := prefix + dbName, err := testDatabaseName(dbHost, dbPort, dbUser, dbPassword) + if err != nil { + return nil, err + } - db := CreatePQDatabase(dbHost, dbPort, dbUser, dbPassword, dbName, dbPrefix, false) + db := CreatePQDatabase(dbHost, dbPort, dbUser, dbPassword, dbName, prefix, false) - err := db.Connect() + // Keep pools small so parallel test packages do not exhaust the Postgres + // server's connection limit + if os.Getenv("COLONIES_DB_MAX_OPEN_CONNS") == "" { + os.Setenv("COLONIES_DB_MAX_OPEN_CONNS", "25") + } + if os.Getenv("COLONIES_DB_MAX_IDLE_CONNS") == "" { + os.Setenv("COLONIES_DB_MAX_IDLE_CONNS", "5") + } + + err = db.Connect() if err != nil { return nil, err } + testSchemas.Lock() + defer testSchemas.Unlock() + + if testSchemas.ready[prefix] { + return db, db.clearTestData() + } + db.Drop() err = db.Initialize() + if err != nil { + return nil, err + } + + testSchemas.ready[prefix] = true + + return db, nil +} + +// clearTestData empties all tables and resets the file sequence, giving each +// test a clean database without paying for schema recreation. +func (db *PQDatabase) clearTestData() error { + tables := make([]string, len(testTables)) + for i, table := range testTables { + tables[i] = db.dbPrefix + table + } + + _, err := db.postgresql.Exec(`TRUNCATE TABLE ` + strings.Join(tables, ", ")) + if err != nil { + return err + } + + _, err = db.postgresql.Exec(`ALTER SEQUENCE ` + db.dbPrefix + `FILE_SEQ RESTART WITH 1`) - return db, err + return err } diff --git a/pkg/server/channel_integration_test.go b/pkg/server/channel_integration_test.go index ea498328e..c52d133c3 100644 --- a/pkg/server/channel_integration_test.go +++ b/pkg/server/channel_integration_test.go @@ -30,14 +30,15 @@ func TestChannelEndToEndIntegration(t *testing.T) { defer db.Close() // Create server - port := 8081 + ports := utils.FreePortsOrPanic(4) + port := ports[0] thisNode := cluster.Node{ Name: "test-node", Host: "localhost", APIPort: port, - EtcdClientPort: 2379, - EtcdPeerPort: 2380, - RelayPort: 25100, + EtcdClientPort: ports[1], + EtcdPeerPort: ports[2], + RelayPort: ports[3], } clusterConfig := cluster.Config{ Nodes: []cluster.Node{thisNode}, @@ -51,7 +52,7 @@ func TestChannelEndToEndIntegration(t *testing.T) { "", thisNode, clusterConfig, - "/tmp/test-etcd-"+time.Now().Format("20060102150405"), // etcd path in /tmp + t.TempDir(), 10, // generator period 10, // cron period false, // exclusive assign @@ -219,14 +220,15 @@ func TestChannelCleanupOnProcessFail(t *testing.T) { defer db.Close() // Create server - port := 8082 + ports := utils.FreePortsOrPanic(4) + port := ports[0] thisNode := cluster.Node{ Name: "test-node", Host: "localhost", APIPort: port, - EtcdClientPort: 2379, - EtcdPeerPort: 2380, - RelayPort: 25101, + EtcdClientPort: ports[1], + EtcdPeerPort: ports[2], + RelayPort: ports[3], } clusterConfig := cluster.Config{ Nodes: []cluster.Node{thisNode}, @@ -240,7 +242,7 @@ func TestChannelCleanupOnProcessFail(t *testing.T) { "", thisNode, clusterConfig, - "/tmp/test-etcd-fail-"+time.Now().Format("20060102150405"), + t.TempDir(), 10, // generator period 10, // cron period false, // exclusive assign diff --git a/pkg/server/controllers/mock_test.go b/pkg/server/controllers/mock_test.go index 3fb809244..5ef01b33d 100644 --- a/pkg/server/controllers/mock_test.go +++ b/pkg/server/controllers/mock_test.go @@ -3,620 +3,639 @@ package controllers import ( "errors" "fmt" + "os" "sync/atomic" "time" - "github.com/colonyos/colonies/pkg/backends" "github.com/colonyos/colonies/pkg/cluster" "github.com/colonyos/colonies/pkg/constants" "github.com/colonyos/colonies/pkg/core" + "github.com/colonyos/colonies/pkg/database" "github.com/colonyos/colonies/pkg/database/postgresql" + "github.com/colonyos/colonies/pkg/utils" ) // portCounter is used to allocate unique ports for each test to avoid port conflicts var portCounter int32 = 0 -// ControllerMock implements the Controller interface for testing -type ControllerMock struct { +// DatabaseMock implements database interfaces for testing +type DatabaseMock struct { ReturnError string ReturnValue string } -func (v *ControllerMock) GetCronPeriod() int { - return -1 +// Implement all database interfaces as no-ops for testing +// ColonyDatabase interface +func (db *DatabaseMock) AddColony(colony *core.Colony) error { + if db.ReturnError == "AddColony" { + return errors.New("mock error") + } + return nil } - -func (v *ControllerMock) GetGeneratorPeriod() int { - return -1 +func (db *DatabaseMock) GetColonies() ([]*core.Colony, error) { + if db.ReturnError == "GetColonies" { + return nil, errors.New("mock error") + } + return nil, nil } - -func (v *ControllerMock) GetEtcdServer() *cluster.EtcdServer { +func (db *DatabaseMock) GetColonyByID(id string) (*core.Colony, error) { + if db.ReturnError == "GetColonyByID" { + return nil, errors.New("mock error") + } + return nil, nil +} +func (db *DatabaseMock) GetColonyByName(name string) (*core.Colony, error) { + if db.ReturnError == "GetColonyByName" { + return nil, errors.New("mock error") + } + return nil, nil +} +func (db *DatabaseMock) RenameColony(colonyName string, newColonyName string) error { + if db.ReturnError == "RenameColony" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) GetEventHandler() backends.RealtimeEventHandler { +func (db *DatabaseMock) RemoveColonyByName(colonyName string) error { + if db.ReturnError == "RemoveColonyByName" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) GetThisNode() cluster.Node { - return cluster.Node{} +func (db *DatabaseMock) CountColonies() (int, error) { + if db.ReturnError == "CountColonies" { + return 0, errors.New("mock error") + } + return 0, nil } -func (v *ControllerMock) SubscribeProcesses(executorID string, subscription *backends.RealtimeSubscription) error { +// ExecutorDatabase interface +func (db *DatabaseMock) AddExecutor(executor *core.Executor) error { + if db.ReturnError == "AddExecutor" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) SubscribeProcess(executorID string, subscription *backends.RealtimeSubscription) error { +func (db *DatabaseMock) SetAllocations(colonyName string, executorName string, allocations core.Allocations) error { + if db.ReturnError == "SetAllocations" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) AddProcessToDB(process *core.Process) (*core.Process, error) { +func (db *DatabaseMock) GetExecutors() ([]*core.Executor, error) { + if db.ReturnError == "GetExecutors" { + return nil, errors.New("mock error") + } return nil, nil } - -func (v *ControllerMock) AddProcess(process *core.Process) (*core.Process, error) { +func (db *DatabaseMock) GetExecutorByID(executorID string) (*core.Executor, error) { + if db.ReturnError == "GetExecutorByID" { + return nil, errors.New("mock error") + } + // Return a dummy executor when no error is set + return &core.Executor{ID: executorID, ColonyName: "test-colony"}, nil +} +func (db *DatabaseMock) GetExecutorsByColonyName(colonyName string) ([]*core.Executor, error) { + if db.ReturnError == "GetExecutorsByColonyName" || db.ReturnError == "GetExecutorByColonyName" { + return nil, errors.New("mock error") + } return nil, nil } - -func (v *ControllerMock) AddChild(processGraphID string, parentProcessID string, childProcessID string, process *core.Process, executorID string, insert bool) (*core.Process, error) { +func (db *DatabaseMock) GetExecutorByName(colonyName string, executorName string) (*core.Executor, error) { + if db.ReturnError == "GetExecutorByName" { + return nil, errors.New("mock error") + } + if db.ReturnValue == "GetExecutorByName" { + return &core.Executor{}, nil + } return nil, nil } - -func (v *ControllerMock) UpdateProcessGraph(graph *core.ProcessGraph) error { +func (db *DatabaseMock) ApproveExecutor(executor *core.Executor) error { return nil } +func (db *DatabaseMock) RejectExecutor(executor *core.Executor) error { return nil } +func (db *DatabaseMock) MarkAlive(executor *core.Executor) error { return nil } +func (db *DatabaseMock) RemoveExecutorByName(colonyName string, executorName string) error { return nil } - -func (v *ControllerMock) CreateProcessGraph(workflowSpec *core.WorkflowSpec, args []interface{}, kwargs map[string]interface{}, rootInput []interface{}, recoveredID string) (*core.ProcessGraph, error) { - return nil, nil +func (db *DatabaseMock) RemoveExecutorsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) CountExecutors() (int, error) { return 0, nil } +func (db *DatabaseMock) CountExecutorsByColonyName(colonyName string) (int, error) { return 0, nil } +func (db *DatabaseMock) CountExecutorsByColonyNameAndState(colonyName string, state int) (int, error) { + return 0, nil } - -func (v *ControllerMock) SubmitWorkflowSpec(workflowSpec *core.WorkflowSpec, recoveredID string) (*core.ProcessGraph, error) { +func (db *DatabaseMock) GetExecutorsByBlueprintID(blueprintID string) ([]*core.Executor, error) { return nil, nil } - -func (v *ControllerMock) GetProcessGraphByID(processGraphID string) (*core.ProcessGraph, error) { - return nil, nil +func (db *DatabaseMock) UpdateExecutorCapabilities(colonyName string, executorName string, capabilities core.Capabilities) error { + return nil } -func (v *ControllerMock) FindWaitingProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { +// LocationDatabase interface +func (db *DatabaseMock) AddLocation(location *core.Location) error { return nil } +func (db *DatabaseMock) GetLocationsByColonyName(colonyName string) ([]*core.Location, error) { return nil, nil } - -func (v *ControllerMock) FindRunningProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { +func (db *DatabaseMock) GetLocationByID(locationID string) (*core.Location, error) { return nil, nil } +func (db *DatabaseMock) GetLocationByName(colonyName string, name string) (*core.Location, error) { return nil, nil } +func (db *DatabaseMock) RemoveLocationByID(locationID string) error { return nil } +func (db *DatabaseMock) RemoveLocationByName(colonyName string, name string) error { return nil } +func (db *DatabaseMock) RemoveLocationsByColonyName(colonyName string) error { return nil } -func (v *ControllerMock) FindSuccessfulProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { - return nil, nil +// ProcessDatabase interface +func (db *DatabaseMock) AddProcess(process *core.Process) error { + if db.ReturnError == "AddProcess" { + return errors.New("mock error") + } + return nil } - -func (v *ControllerMock) FindFailedProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { +func (db *DatabaseMock) GetProcesses() ([]*core.Process, error) { + if db.ReturnError == "GetProcesses" { + return nil, errors.New("mock error") + } return nil, nil } - -func (v *ControllerMock) CloseSuccessful(processID string, executorID string, output []interface{}) error { - return nil +func (db *DatabaseMock) GetProcessByID(processID string) (*core.Process, error) { + if db.ReturnError == "GetProcessByID" { + return nil, errors.New("mock error") + } + return nil, nil } - -func (v *ControllerMock) NotifyChildren(process *core.Process) error { +func (db *DatabaseMock) SetProcessState(processID string, state int) error { + if db.ReturnError == "SetProcessState" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) CloseFailed(processID string, errs []string) error { +func (db *DatabaseMock) SetWaitForParents(processID string, waitForParent bool) error { + if db.ReturnError == "SetWaitForParents" { + return errors.New("mock error") + } return nil } - -func (v *ControllerMock) HandleDefunctProcessgraph(processGraphID string, processID string, err error) error { - return nil +func (db *DatabaseMock) MarkSuccessful(processID string) (float64, float64, error) { + if db.ReturnError == "MarkSuccessful" { + return 0, 0, errors.New("mock error") + } + return 1.0, 1.0, nil } - -func (v *ControllerMock) Assign(executorID string, colonyName string, cpu int64, memory int64) (*AssignResult, error) { +func (db *DatabaseMock) FindProcessesByColonyName(colonyName string, seconds int, state int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) DistributedAssign(executor *core.Executor, colonyName string, cpu int64, memory int64, storage int64) (*AssignResult, error) { +func (db *DatabaseMock) FindProcessesByExecutorID(colonyName string, executorID string, seconds int, state int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) UnassignExecutor(processID string) error { - return nil +func (db *DatabaseMock) FindWaitingProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { + return nil, nil } - -func (v *ControllerMock) ResetProcess(processID string) error { - return nil +func (db *DatabaseMock) FindRunningProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { + return nil, nil } - -func (v *ControllerMock) AddGenerator(generator *core.Generator) (*core.Generator, error) { +func (db *DatabaseMock) FindSuccessfulProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) GetGenerator(generatorID string) (*core.Generator, error) { +func (db *DatabaseMock) FindFailedProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) ResolveGenerator(colonyName string, generatorName string) (*core.Generator, error) { +func (db *DatabaseMock) FindAllRunningProcesses() ([]*core.Process, error) { return nil, nil } +func (db *DatabaseMock) FindAllWaitingProcesses() ([]*core.Process, error) { return nil, nil } +func (db *DatabaseMock) FindCandidates(colonyName string, executorType string, executorLocationName string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) GetGenerators(colonyName string, count int) ([]*core.Generator, error) { +func (db *DatabaseMock) FindCandidatesByName(colonyName string, executorName string, executorType string, executorLocationName string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) ([]*core.Process, error) { return nil, nil } - -func (v *ControllerMock) PackGenerator(generatorID string, colonyName, arg string) error { +func (db *DatabaseMock) RemoveProcessByID(processID string) error { return nil } +func (db *DatabaseMock) RemoveAllProcesses() error { return nil } +func (db *DatabaseMock) RemoveAllWaitingProcessesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllRunningProcessesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllSuccessfulProcessesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllFailedProcessesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllProcessesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllProcessesByProcessGraphID(processGraphID string) error { return nil } +func (db *DatabaseMock) RemoveAllProcessesInProcessGraphsByColonyName(colonyName string) error { return nil } - -func (v *ControllerMock) GeneratorTriggerLoop() { -} - -func (v *ControllerMock) TriggerGenerators() { +func (db *DatabaseMock) ResetProcess(process *core.Process) error { return nil } +func (db *DatabaseMock) SetInput(processID string, input []interface{}) error { return nil } +func (db *DatabaseMock) SetOutput(processID string, output []interface{}) error { return nil } +func (db *DatabaseMock) SetErrors(processID string, errs []string) error { return nil } +func (db *DatabaseMock) SetParents(processID string, parents []string) error { return nil } +func (db *DatabaseMock) SetChildren(processID string, children []string) error { return nil } +func (db *DatabaseMock) Assign(executorID string, process *core.Process) error { return nil } +func (db *DatabaseMock) SelectAndAssign(colonyName string, executorID string, executorName string, executorType string, executorLocation string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) (*core.Process, error) { + return nil, nil } - -func (v *ControllerMock) SubmitWorkflow(generator *core.Generator, counter int, recoveredID string) { +func (db *DatabaseMock) Unassign(process *core.Process) error { return nil } +func (db *DatabaseMock) MarkFailed(processID string, errs []string) error { return nil } +func (db *DatabaseMock) CountProcesses() (int, error) { return 0, nil } +func (db *DatabaseMock) CountWaitingProcesses() (int, error) { return 0, nil } +func (db *DatabaseMock) CountRunningProcesses() (int, error) { return 0, nil } +func (db *DatabaseMock) CountSuccessfulProcesses() (int, error) { return 0, nil } +func (db *DatabaseMock) CountFailedProcesses() (int, error) { return 0, nil } +func (db *DatabaseMock) CountWaitingProcessesByColonyName(colonyName string) (int, error) { + return 0, nil } - -func (v *ControllerMock) AddCron(cron *core.Cron) (*core.Cron, error) { - return nil, nil +func (db *DatabaseMock) CountRunningProcessesByColonyName(colonyName string) (int, error) { + return 0, nil } - -func (v *ControllerMock) RemoveGenerator(generatorID string) error { - return nil +func (db *DatabaseMock) CountSuccessfulProcessesByColonyName(colonyName string) (int, error) { + return 0, nil } - -func (v *ControllerMock) GetCron(cronID string) (*core.Cron, error) { - return nil, nil +func (db *DatabaseMock) CountFailedProcessesByColonyName(colonyName string) (int, error) { + return 0, nil } -func (v *ControllerMock) GetCrons(colonyName string, count int) ([]*core.Cron, error) { +// UserDatabase interface +func (db *DatabaseMock) AddUser(user *core.User) error { return nil } +func (db *DatabaseMock) GetUserByName(colonyName string, name string) (*core.User, error) { return nil, nil } - -func (v *ControllerMock) GetCronByName(colonyName string, cronName string) (*core.Cron, error) { +func (db *DatabaseMock) GetUserByID(colonyName string, userID string) (*core.User, error) { return nil, nil } - -func (v *ControllerMock) RunCron(cronID string) (*core.Cron, error) { +func (db *DatabaseMock) GetUsersByColonyName(colonyName string) ([]*core.User, error) { return nil, nil } +func (db *DatabaseMock) RemoveUserByName(colonyName string, name string) error { return nil } +func (db *DatabaseMock) RemoveUserByID(colonyName string, userID string) error { return nil } +func (db *DatabaseMock) RemoveUsersByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) CountUsers() (int, error) { return 0, nil } -func (v *ControllerMock) RemoveCron(cronID string) error { - return nil -} - -func (v *ControllerMock) CalcNextRun(cron *core.Cron) time.Time { - return time.Time{} +// AttributeDatabase interface +func (db *DatabaseMock) AddAttribute(attribute core.Attribute) error { return nil } +func (db *DatabaseMock) AddAttributes(attributes []core.Attribute) error { return nil } +func (db *DatabaseMock) GetAttributeByID(attributeID string) (core.Attribute, error) { + return core.Attribute{}, nil } - -func (v *ControllerMock) StartCron(cron *core.Cron) { +func (db *DatabaseMock) GetAttributesByColonyName(colonyName string) ([]core.Attribute, error) { + return nil, nil } - -func (v *ControllerMock) TriggerCrons() { +func (db *DatabaseMock) GetAttribute(targetID string, key string, attributeType int) (core.Attribute, error) { + return core.Attribute{}, nil } - -func (v *ControllerMock) CronTriggerLoop() { +func (db *DatabaseMock) GetAttributes(targetID string) ([]core.Attribute, error) { return nil, nil } +func (db *DatabaseMock) GetAttributesByType(targetID string, attributeType int) ([]core.Attribute, error) { + return nil, nil } - -func (v *ControllerMock) ResetDatabase() error { +func (db *DatabaseMock) UpdateAttribute(attribute core.Attribute) error { return nil } +func (db *DatabaseMock) RemoveAttributeByID(attributeID string) error { return nil } +func (db *DatabaseMock) RemoveAllAttributesByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllAttributesByColonyNameWithState(colonyName string, state int) error { return nil } - -func (v *ControllerMock) PauseColonyAssignments(colonyName string) error { +func (db *DatabaseMock) RemoveAllAttributesByProcessGraphID(processGraphID string) error { return nil } +func (db *DatabaseMock) RemoveAllAttributesInProcessGraphsByColonyName(colonyName string) error { return nil } - -func (v *ControllerMock) ResumeColonyAssignments(colonyName string) error { +func (db *DatabaseMock) RemoveAllAttributesInProcessGraphsByColonyNameWithState(colonyName string, state int) error { return nil } - -func (v *ControllerMock) AreColonyAssignmentsPaused(colonyName string) (bool, error) { - return false, nil -} - -func (v *ControllerMock) Stop() { -} - -func (v *ControllerMock) IsLeader() bool { - return false -} - -func (v *ControllerMock) TryBecomeLeader() bool { - return false -} - -func (v *ControllerMock) TimeoutLoop() { -} - -func (v *ControllerMock) BlockingCmdQueueWorker() { -} - -func (v *ControllerMock) RetentionWorker() { -} - -func (v *ControllerMock) CmdQueueWorker() { -} - -// DatabaseMock implements database interfaces for testing -type DatabaseMock struct { - ReturnError string - ReturnValue string -} - -// Implement all database interfaces as no-ops for testing -// ColonyDatabase interface -func (db *DatabaseMock) AddColony(colony *core.Colony) error { - if db.ReturnError == "AddColony" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveAttributesByTargetID(targetID string, attributeType int) error { return nil } -func (db *DatabaseMock) GetColonies() ([]*core.Colony, error) { - if db.ReturnError == "GetColonies" { return nil, errors.New("mock error") } - return nil, nil -} -func (db *DatabaseMock) GetColonyByID(id string) (*core.Colony, error) { - if db.ReturnError == "GetColonyByID" { return nil, errors.New("mock error") } +func (db *DatabaseMock) RemoveAllAttributesByTargetID(targetID string) error { return nil } +func (db *DatabaseMock) RemoveAllAttributes() error { return nil } + +// FunctionDatabase interface +func (db *DatabaseMock) AddFunction(function *core.Function) error { return nil } +func (db *DatabaseMock) GetFunctionByID(functionID string) (*core.Function, error) { return nil, nil } +func (db *DatabaseMock) GetFunctionsByExecutorName(colonyName string, executorName string) ([]*core.Function, error) { return nil, nil } -func (db *DatabaseMock) GetColonyByName(name string) (*core.Colony, error) { - if db.ReturnError == "GetColonyByName" { return nil, errors.New("mock error") } +func (db *DatabaseMock) GetFunctionsByColonyName(colonyName string) ([]*core.Function, error) { return nil, nil } -func (db *DatabaseMock) RenameColony(colonyName string, newColonyName string) error { - if db.ReturnError == "RenameColony" { return errors.New("mock error") } - return nil -} -func (db *DatabaseMock) RemoveColonyByName(colonyName string) error { - if db.ReturnError == "RemoveColonyByName" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveFunctionByID(functionID string) error { return nil } +func (db *DatabaseMock) RemoveFunctionsByExecutorName(colonyName string, executorName string) error { return nil } -func (db *DatabaseMock) CountColonies() (int, error) { - if db.ReturnError == "CountColonies" { return 0, errors.New("mock error") } - return 0, nil +func (db *DatabaseMock) RemoveFunctionsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) GetFunctionsByExecutorAndName(colonyName string, executorName string, name string) (*core.Function, error) { + return nil, nil } - -// ExecutorDatabase interface -func (db *DatabaseMock) AddExecutor(executor *core.Executor) error { - if db.ReturnError == "AddExecutor" { return errors.New("mock error") } +func (db *DatabaseMock) UpdateFunctionStats(colonyName string, executorName string, name string, counter int, minWaitTime float64, maxWaitTime float64, minExecTime float64, maxExecTime float64, avgWaitTime float64, avgExecTime float64) error { return nil } -func (db *DatabaseMock) SetAllocations(colonyName string, executorName string, allocations core.Allocations) error { - if db.ReturnError == "SetAllocations" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveFunctionByName(colonyName string, executorName string, name string) error { return nil } -func (db *DatabaseMock) GetExecutors() ([]*core.Executor, error) { - if db.ReturnError == "GetExecutors" { return nil, errors.New("mock error") } +func (db *DatabaseMock) RemoveFunctions() error { return nil } +func (db *DatabaseMock) CountFunctions() (int, error) { return 0, nil } + +// ProcessGraphDatabase interface +func (db *DatabaseMock) AddProcessGraph(processGraph *core.ProcessGraph) error { return nil } +func (db *DatabaseMock) GetProcessGraphByID(processGraphID string) (*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) GetExecutorByID(executorID string) (*core.Executor, error) { - if db.ReturnError == "GetExecutorByID" { return nil, errors.New("mock error") } - // Return a dummy executor when no error is set - return &core.Executor{ID: executorID, ColonyName: "test-colony"}, nil +func (db *DatabaseMock) SetProcessGraphState(processGraphID string, state int) error { + if db.ReturnError == "SetProcessGraphState" { + return errors.New("mock error") + } + return nil } -func (db *DatabaseMock) GetExecutorsByColonyName(colonyName string) ([]*core.Executor, error) { - if db.ReturnError == "GetExecutorsByColonyName" || db.ReturnError == "GetExecutorByColonyName" { return nil, errors.New("mock error") } +func (db *DatabaseMock) FindWaitingProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) GetExecutorByName(colonyName string, executorName string) (*core.Executor, error) { - if db.ReturnError == "GetExecutorByName" { return nil, errors.New("mock error") } - if db.ReturnValue == "GetExecutorByName" { return &core.Executor{}, nil } +func (db *DatabaseMock) FindRunningProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) ApproveExecutor(executor *core.Executor) error { return nil } -func (db *DatabaseMock) RejectExecutor(executor *core.Executor) error { return nil } -func (db *DatabaseMock) MarkAlive(executor *core.Executor) error { return nil } -func (db *DatabaseMock) RemoveExecutorByName(colonyName string, executorName string) error { return nil } -func (db *DatabaseMock) RemoveExecutorsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) CountExecutors() (int, error) { return 0, nil } -func (db *DatabaseMock) CountExecutorsByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountExecutorsByColonyNameAndState(colonyName string, state int) (int, error) { return 0, nil } -func (db *DatabaseMock) GetExecutorsByBlueprintID(blueprintID string) ([]*core.Executor, error) { return nil, nil } -func (db *DatabaseMock) UpdateExecutorCapabilities(colonyName string, executorName string, capabilities core.Capabilities) error { return nil } - -// LocationDatabase interface -func (db *DatabaseMock) AddLocation(location *core.Location) error { return nil } -func (db *DatabaseMock) GetLocationsByColonyName(colonyName string) ([]*core.Location, error) { return nil, nil } -func (db *DatabaseMock) GetLocationByID(locationID string) (*core.Location, error) { return nil, nil } -func (db *DatabaseMock) GetLocationByName(colonyName string, name string) (*core.Location, error) { return nil, nil } -func (db *DatabaseMock) RemoveLocationByID(locationID string) error { return nil } -func (db *DatabaseMock) RemoveLocationByName(colonyName string, name string) error { return nil } -func (db *DatabaseMock) RemoveLocationsByColonyName(colonyName string) error { return nil } - -// ProcessDatabase interface -func (db *DatabaseMock) AddProcess(process *core.Process) error { - if db.ReturnError == "AddProcess" { return errors.New("mock error") } - return nil -} -func (db *DatabaseMock) GetProcesses() ([]*core.Process, error) { - if db.ReturnError == "GetProcesses" { return nil, errors.New("mock error") } +func (db *DatabaseMock) FindSuccessfulProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) GetProcessByID(processID string) (*core.Process, error) { - if db.ReturnError == "GetProcessByID" { return nil, errors.New("mock error") } +func (db *DatabaseMock) FindFailedProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) SetProcessState(processID string, state int) error { - if db.ReturnError == "SetProcessState" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveProcessGraphByID(processGraphID string) error { return nil } +func (db *DatabaseMock) RemoveAllProcessGraphsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllWaitingProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) SetWaitForParents(processID string, waitForParent bool) error { - if db.ReturnError == "SetWaitForParents" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveAllRunningProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) MarkSuccessful(processID string) (float64, float64, error) { - if db.ReturnError == "MarkSuccessful" { return 0, 0, errors.New("mock error") } - return 1.0, 1.0, nil -} -func (db *DatabaseMock) FindProcessesByColonyName(colonyName string, seconds int, state int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindProcessesByExecutorID(colonyName string, executorID string, seconds int, state int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindWaitingProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindRunningProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindSuccessfulProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindFailedProcesses(colonyName string, executorType string, label string, initiator string, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindAllRunningProcesses() ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindAllWaitingProcesses() ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindCandidates(colonyName string, executorType string, executorLocationName string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) FindCandidatesByName(colonyName string, executorName string, executorType string, executorLocationName string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) ([]*core.Process, error) { return nil, nil } -func (db *DatabaseMock) RemoveProcessByID(processID string) error { return nil } -func (db *DatabaseMock) RemoveAllProcesses() error { return nil } -func (db *DatabaseMock) RemoveAllWaitingProcessesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllRunningProcessesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllSuccessfulProcessesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllFailedProcessesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllProcessesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllProcessesByProcessGraphID(processGraphID string) error { return nil } -func (db *DatabaseMock) RemoveAllProcessesInProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) ResetProcess(process *core.Process) error { return nil } -func (db *DatabaseMock) SetInput(processID string, input []interface{}) error { return nil } -func (db *DatabaseMock) SetOutput(processID string, output []interface{}) error { return nil } -func (db *DatabaseMock) SetErrors(processID string, errs []string) error { return nil } -func (db *DatabaseMock) SetParents(processID string, parents []string) error { return nil } -func (db *DatabaseMock) SetChildren(processID string, children []string) error { return nil } -func (db *DatabaseMock) Assign(executorID string, process *core.Process) error { return nil } -func (db *DatabaseMock) SelectAndAssign(colonyName string, executorID string, executorName string, executorType string, executorLocation string, cpu int64, memory int64, storage int64, nodes int, processes int, processesPerNode int, count int) (*core.Process, error) { return nil, nil } -func (db *DatabaseMock) Unassign(process *core.Process) error { return nil } -func (db *DatabaseMock) MarkFailed(processID string, errs []string) error { return nil } -func (db *DatabaseMock) CountProcesses() (int, error) { return 0, nil } -func (db *DatabaseMock) CountWaitingProcesses() (int, error) { return 0, nil } -func (db *DatabaseMock) CountRunningProcesses() (int, error) { return 0, nil } -func (db *DatabaseMock) CountSuccessfulProcesses() (int, error) { return 0, nil } -func (db *DatabaseMock) CountFailedProcesses() (int, error) { return 0, nil } -func (db *DatabaseMock) CountWaitingProcessesByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountRunningProcessesByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountSuccessfulProcessesByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountFailedProcessesByColonyName(colonyName string) (int, error) { return 0, nil } - -// UserDatabase interface -func (db *DatabaseMock) AddUser(user *core.User) error { return nil } -func (db *DatabaseMock) GetUserByName(colonyName string, name string) (*core.User, error) { return nil, nil } -func (db *DatabaseMock) GetUserByID(colonyName string, userID string) (*core.User, error) { return nil, nil } -func (db *DatabaseMock) GetUsersByColonyName(colonyName string) ([]*core.User, error) { return nil, nil } -func (db *DatabaseMock) RemoveUserByName(colonyName string, name string) error { return nil } -func (db *DatabaseMock) RemoveUserByID(colonyName string, userID string) error { return nil } -func (db *DatabaseMock) RemoveUsersByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) CountUsers() (int, error) { return 0, nil } - -// AttributeDatabase interface -func (db *DatabaseMock) AddAttribute(attribute core.Attribute) error { return nil } -func (db *DatabaseMock) AddAttributes(attributes []core.Attribute) error { return nil } -func (db *DatabaseMock) GetAttributeByID(attributeID string) (core.Attribute, error) { return core.Attribute{}, nil } -func (db *DatabaseMock) GetAttributesByColonyName(colonyName string) ([]core.Attribute, error) { return nil, nil } -func (db *DatabaseMock) GetAttribute(targetID string, key string, attributeType int) (core.Attribute, error) { return core.Attribute{}, nil } -func (db *DatabaseMock) GetAttributes(targetID string) ([]core.Attribute, error) { return nil, nil } -func (db *DatabaseMock) GetAttributesByType(targetID string, attributeType int) ([]core.Attribute, error) { return nil, nil } -func (db *DatabaseMock) UpdateAttribute(attribute core.Attribute) error { return nil } -func (db *DatabaseMock) RemoveAttributeByID(attributeID string) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesByColonyNameWithState(colonyName string, state int) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesByProcessGraphID(processGraphID string) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesInProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesInProcessGraphsByColonyNameWithState(colonyName string, state int) error { return nil } -func (db *DatabaseMock) RemoveAttributesByTargetID(targetID string, attributeType int) error { return nil } -func (db *DatabaseMock) RemoveAllAttributesByTargetID(targetID string) error { return nil } -func (db *DatabaseMock) RemoveAllAttributes() error { return nil } - -// FunctionDatabase interface -func (db *DatabaseMock) AddFunction(function *core.Function) error { return nil } -func (db *DatabaseMock) GetFunctionByID(functionID string) (*core.Function, error) { return nil, nil } -func (db *DatabaseMock) GetFunctionsByExecutorName(colonyName string, executorName string) ([]*core.Function, error) { return nil, nil } -func (db *DatabaseMock) GetFunctionsByColonyName(colonyName string) ([]*core.Function, error) { return nil, nil } -func (db *DatabaseMock) RemoveFunctionByID(functionID string) error { return nil } -func (db *DatabaseMock) RemoveFunctionsByExecutorName(colonyName string, executorName string) error { return nil } -func (db *DatabaseMock) RemoveFunctionsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) GetFunctionsByExecutorAndName(colonyName string, executorName string, name string) (*core.Function, error) { return nil, nil } -func (db *DatabaseMock) UpdateFunctionStats(colonyName string, executorName string, name string, counter int, minWaitTime float64, maxWaitTime float64, minExecTime float64, maxExecTime float64, avgWaitTime float64, avgExecTime float64) error { return nil } -func (db *DatabaseMock) RemoveFunctionByName(colonyName string, executorName string, name string) error { return nil } -func (db *DatabaseMock) RemoveFunctions() error { return nil } -func (db *DatabaseMock) CountFunctions() (int, error) { return 0, nil } - -// ProcessGraphDatabase interface -func (db *DatabaseMock) AddProcessGraph(processGraph *core.ProcessGraph) error { return nil } -func (db *DatabaseMock) GetProcessGraphByID(processGraphID string) (*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) SetProcessGraphState(processGraphID string, state int) error { - if db.ReturnError == "SetProcessGraphState" { return errors.New("mock error") } +func (db *DatabaseMock) RemoveAllSuccessfulProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) FindWaitingProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) FindRunningProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) FindSuccessfulProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) FindFailedProcessGraphs(colonyName string, count int) ([]*core.ProcessGraph, error) { return nil, nil } -func (db *DatabaseMock) RemoveProcessGraphByID(processGraphID string) error { return nil } -func (db *DatabaseMock) RemoveAllProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllWaitingProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllRunningProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) RemoveAllSuccessfulProcessGraphsByColonyName(colonyName string) error { return nil } func (db *DatabaseMock) RemoveAllFailedProcessGraphsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) CountWaitingProcessGraphs() (int, error) { return 0, nil } -func (db *DatabaseMock) CountRunningProcessGraphs() (int, error) { return 0, nil } -func (db *DatabaseMock) CountSuccessfulProcessGraphs() (int, error) { return 0, nil } -func (db *DatabaseMock) CountFailedProcessGraphs() (int, error) { return 0, nil } -func (db *DatabaseMock) CountWaitingProcessGraphsByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountRunningProcessGraphsByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountSuccessfulProcessGraphsByColonyName(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountFailedProcessGraphsByColonyName(colonyName string) (int, error) { return 0, nil } +func (db *DatabaseMock) CountWaitingProcessGraphs() (int, error) { return 0, nil } +func (db *DatabaseMock) CountRunningProcessGraphs() (int, error) { return 0, nil } +func (db *DatabaseMock) CountSuccessfulProcessGraphs() (int, error) { return 0, nil } +func (db *DatabaseMock) CountFailedProcessGraphs() (int, error) { return 0, nil } +func (db *DatabaseMock) CountWaitingProcessGraphsByColonyName(colonyName string) (int, error) { + return 0, nil +} +func (db *DatabaseMock) CountRunningProcessGraphsByColonyName(colonyName string) (int, error) { + return 0, nil +} +func (db *DatabaseMock) CountSuccessfulProcessGraphsByColonyName(colonyName string) (int, error) { + return 0, nil +} +func (db *DatabaseMock) CountFailedProcessGraphsByColonyName(colonyName string) (int, error) { + return 0, nil +} // GeneratorDatabase interface func (db *DatabaseMock) AddGenerator(generator *core.Generator) error { - if db.ReturnError == "AddGenerator" { return errors.New("mock error") } + if db.ReturnError == "AddGenerator" { + return errors.New("mock error") + } return nil } func (db *DatabaseMock) GetGeneratorByID(generatorID string) (*core.Generator, error) { - if db.ReturnError == "GetGeneratorByID" { return nil, errors.New("mock error") } + if db.ReturnError == "GetGeneratorByID" { + return nil, errors.New("mock error") + } return &core.Generator{ID: generatorID, ColonyName: "test-colony"}, nil } func (db *DatabaseMock) GetGeneratorByName(colonyName string, generatorName string) (*core.Generator, error) { - if db.ReturnError == "GetGeneratorByName" { return nil, errors.New("mock error") } + if db.ReturnError == "GetGeneratorByName" { + return nil, errors.New("mock error") + } return &core.Generator{ID: generatorName, ColonyName: colonyName}, nil } func (db *DatabaseMock) FindAllGenerators() ([]*core.Generator, error) { return nil, nil } -func (db *DatabaseMock) FindGeneratorsByColonyName(colonyName string, count int) ([]*core.Generator, error) { return nil, nil } -func (db *DatabaseMock) RemoveGeneratorByID(generatorID string) error { return nil } +func (db *DatabaseMock) FindGeneratorsByColonyName(colonyName string, count int) ([]*core.Generator, error) { + return nil, nil +} +func (db *DatabaseMock) RemoveGeneratorByID(generatorID string) error { return nil } func (db *DatabaseMock) AddGeneratorArg(generatorArg *core.GeneratorArg) error { return nil } -func (db *DatabaseMock) GetGeneratorArgs(generatorID string, count int) ([]*core.GeneratorArg, error) { return nil, nil } -func (db *DatabaseMock) CountGeneratorArgs(generatorID string) (int, error) { return 0, nil } -func (db *DatabaseMock) RemoveGeneratorArgByID(generatorArgID string) error { return nil } -func (db *DatabaseMock) SetGeneratorLastRun(generatorID string) error { return nil } -func (db *DatabaseMock) SetGeneratorFirstPack(generatorID string) error { return nil } -func (db *DatabaseMock) RemoveAllGeneratorsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) GetGeneratorArgs(generatorID string, count int) ([]*core.GeneratorArg, error) { + return nil, nil +} +func (db *DatabaseMock) CountGeneratorArgs(generatorID string) (int, error) { return 0, nil } +func (db *DatabaseMock) RemoveGeneratorArgByID(generatorArgID string) error { return nil } +func (db *DatabaseMock) SetGeneratorLastRun(generatorID string) error { return nil } +func (db *DatabaseMock) SetGeneratorFirstPack(generatorID string) error { return nil } +func (db *DatabaseMock) RemoveAllGeneratorsByColonyName(colonyName string) error { return nil } func (db *DatabaseMock) RemoveAllGeneratorArgsByGeneratorID(generatorID string) error { return nil } -func (db *DatabaseMock) RemoveAllGeneratorArgsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) RemoveAllGeneratorArgsByColonyName(colonyName string) error { return nil } // CronDatabase interface func (db *DatabaseMock) AddCron(cron *core.Cron) error { - if db.ReturnError == "AddCron" { return errors.New("mock error") } + if db.ReturnError == "AddCron" { + return errors.New("mock error") + } return nil } func (db *DatabaseMock) GetCronByID(cronID string) (*core.Cron, error) { - if db.ReturnError == "GetCronByID" { return nil, errors.New("mock error") } + if db.ReturnError == "GetCronByID" { + return nil, errors.New("mock error") + } return &core.Cron{ID: cronID, ColonyName: "test-colony"}, nil } func (db *DatabaseMock) FindAllCrons() ([]*core.Cron, error) { return nil, nil } -func (db *DatabaseMock) FindCronsByColonyName(colonyName string, count int) ([]*core.Cron, error) { return nil, nil } +func (db *DatabaseMock) FindCronsByColonyName(colonyName string, count int) ([]*core.Cron, error) { + return nil, nil +} func (db *DatabaseMock) RemoveCronByID(cronID string) error { - if db.ReturnError == "RemoveCronByID" { return errors.New("mock error") } + if db.ReturnError == "RemoveCronByID" { + return errors.New("mock error") + } return nil } -func (db *DatabaseMock) UpdateCron(cronID string, nextRun time.Time, lastRun time.Time, processGraphID string) error { return nil } -func (db *DatabaseMock) GetCronByName(colonyName string, cronName string) (*core.Cron, error) { return nil, nil } +func (db *DatabaseMock) UpdateCron(cronID string, nextRun time.Time, lastRun time.Time, processGraphID string) error { + return nil +} +func (db *DatabaseMock) GetCronByName(colonyName string, cronName string) (*core.Cron, error) { + return nil, nil +} func (db *DatabaseMock) RemoveAllCronsByColonyName(colonyName string) error { return nil } -// LogDatabase interface -func (db *DatabaseMock) AddLog(processID string, colonyName string, executorName string, timestamp int64, msg string) error { return nil } -func (db *DatabaseMock) GetLogsByProcessID(processID string, limit int) ([]*core.Log, error) { return nil, nil } -func (db *DatabaseMock) GetLogsByProcessIDSince(processID string, limit int, since int64) ([]*core.Log, error) { return nil, nil } -func (db *DatabaseMock) GetLogsByExecutor(executorName string, limit int) ([]*core.Log, error) { return nil, nil } -func (db *DatabaseMock) GetLogsByExecutorSince(executorName string, limit int, since int64) ([]*core.Log, error) { return nil, nil } -func (db *DatabaseMock) GetLogsByProcessIDLatest(processID string, limit int) ([]*core.Log, error) { return nil, nil } -func (db *DatabaseMock) GetLogsByExecutorLatest(executorName string, limit int) ([]*core.Log, error) { return nil, nil } +// LogDatabase interface +func (db *DatabaseMock) AddLog(processID string, colonyName string, executorName string, timestamp int64, msg string) error { + return nil +} +func (db *DatabaseMock) GetLogsByProcessID(processID string, limit int) ([]*core.Log, error) { + return nil, nil +} +func (db *DatabaseMock) GetLogsByProcessIDSince(processID string, limit int, since int64) ([]*core.Log, error) { + return nil, nil +} +func (db *DatabaseMock) GetLogsByExecutor(executorName string, limit int) ([]*core.Log, error) { + return nil, nil +} +func (db *DatabaseMock) GetLogsByExecutorSince(executorName string, limit int, since int64) ([]*core.Log, error) { + return nil, nil +} +func (db *DatabaseMock) GetLogsByProcessIDLatest(processID string, limit int) ([]*core.Log, error) { + return nil, nil +} +func (db *DatabaseMock) GetLogsByExecutorLatest(executorName string, limit int) ([]*core.Log, error) { + return nil, nil +} func (db *DatabaseMock) RemoveLogsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) CountLogs(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) SearchLogs(colonyName string, text string, days int, count int) ([]*core.Log, error) { return nil, nil } +func (db *DatabaseMock) CountLogs(colonyName string) (int, error) { return 0, nil } +func (db *DatabaseMock) SearchLogs(colonyName string, text string, days int, count int) ([]*core.Log, error) { + return nil, nil +} // FileDatabase interface func (db *DatabaseMock) AddFile(file *core.File) error { return nil } -func (db *DatabaseMock) GetFileByID(colonyName string, fileID string) (*core.File, error) { return nil, nil } -func (db *DatabaseMock) GetLatestFileByName(colonyName string, label string, name string) ([]*core.File, error) { return nil, nil } -func (db *DatabaseMock) GetFileByName(colonyName string, label string, name string) ([]*core.File, error) { return nil, nil } -func (db *DatabaseMock) GetFilenamesByLabel(colonyName string, label string) ([]string, error) { return nil, nil } -func (db *DatabaseMock) GetFileDataByLabel(colonyName string, label string) ([]*core.FileData, error) { return nil, nil } +func (db *DatabaseMock) GetFileByID(colonyName string, fileID string) (*core.File, error) { + return nil, nil +} +func (db *DatabaseMock) GetLatestFileByName(colonyName string, label string, name string) ([]*core.File, error) { + return nil, nil +} +func (db *DatabaseMock) GetFileByName(colonyName string, label string, name string) ([]*core.File, error) { + return nil, nil +} +func (db *DatabaseMock) GetFilenamesByLabel(colonyName string, label string) ([]string, error) { + return nil, nil +} +func (db *DatabaseMock) GetFileDataByLabel(colonyName string, label string) ([]*core.FileData, error) { + return nil, nil +} func (db *DatabaseMock) GetFileLabels(colonyName string) ([]*core.Label, error) { return nil, nil } -func (db *DatabaseMock) GetFileLabelsByName(colonyName string, name string, exact bool) ([]*core.Label, error) { return nil, nil } -func (db *DatabaseMock) GetFilesByColonyName(colonyName string) ([]*core.File, error) { return nil, nil } -func (db *DatabaseMock) GetFilesByProcessGraphID(processGraphID string) ([]*core.File, error) { return nil, nil } -func (db *DatabaseMock) GetFiles() ([]*core.File, error) { return nil, nil } -func (db *DatabaseMock) UpdateFile(file *core.File) error { return nil } +func (db *DatabaseMock) GetFileLabelsByName(colonyName string, name string, exact bool) ([]*core.Label, error) { + return nil, nil +} +func (db *DatabaseMock) GetFilesByColonyName(colonyName string) ([]*core.File, error) { + return nil, nil +} +func (db *DatabaseMock) GetFilesByProcessGraphID(processGraphID string) ([]*core.File, error) { + return nil, nil +} +func (db *DatabaseMock) GetFiles() ([]*core.File, error) { return nil, nil } +func (db *DatabaseMock) UpdateFile(file *core.File) error { return nil } func (db *DatabaseMock) RemoveFileByID(colonyName string, fileID string) error { return nil } -func (db *DatabaseMock) RemoveFileByName(colonyName string, label string, name string) error { return nil } +func (db *DatabaseMock) RemoveFileByName(colonyName string, label string, name string) error { + return nil +} func (db *DatabaseMock) CountFiles(colonyName string) (int, error) { return 0, nil } -func (db *DatabaseMock) CountFilesWithLabel(colonyName string, label string) (int, error) { return 0, nil } - +func (db *DatabaseMock) CountFilesWithLabel(colonyName string, label string) (int, error) { + return 0, nil +} // SnapshotDatabase interface -func (db *DatabaseMock) CreateSnapshot(colonyName string, label string, name string) (*core.Snapshot, error) { return nil, nil } -func (db *DatabaseMock) GetSnapshotByID(colonyName string, snapshotID string) (*core.Snapshot, error) { return nil, nil } -func (db *DatabaseMock) GetSnapshotsByColonyName(colonyName string) ([]*core.Snapshot, error) { return nil, nil } +func (db *DatabaseMock) CreateSnapshot(colonyName string, label string, name string) (*core.Snapshot, error) { + return nil, nil +} +func (db *DatabaseMock) GetSnapshotByID(colonyName string, snapshotID string) (*core.Snapshot, error) { + return nil, nil +} +func (db *DatabaseMock) GetSnapshotsByColonyName(colonyName string) ([]*core.Snapshot, error) { + return nil, nil +} func (db *DatabaseMock) RemoveSnapshotByID(colonyName string, snapshotID string) error { return nil } -func (db *DatabaseMock) GetSnapshotByName(colonyName string, name string) (*core.Snapshot, error) { return nil, nil } +func (db *DatabaseMock) GetSnapshotByName(colonyName string, name string) (*core.Snapshot, error) { + return nil, nil +} func (db *DatabaseMock) RemoveSnapshotByName(colonyName string, name string) error { return nil } -func (db *DatabaseMock) RemoveSnapshotsByColonyName(colonyName string) error { return nil } -func (db *DatabaseMock) CountSnapshots() (int, error) { return 0, nil } - +func (db *DatabaseMock) RemoveSnapshotsByColonyName(colonyName string) error { return nil } +func (db *DatabaseMock) CountSnapshots() (int, error) { return 0, nil } // SecurityDatabase interface func (db *DatabaseMock) SetServerID(oldServerID, newServerID string) error { return nil } -func (db *DatabaseMock) GetServerID() (string, error) { return "", nil } -func (db *DatabaseMock) ChangeColonyID(colonyName string, oldColonyID, newColonyID string) error { return nil } -func (db *DatabaseMock) ChangeUserID(colonyName string, oldUserID, newUserID string) error { return nil } -func (db *DatabaseMock) ChangeExecutorID(colonyName string, oldExecutorID, newExecutorID string) error { return nil } +func (db *DatabaseMock) GetServerID() (string, error) { return "", nil } +func (db *DatabaseMock) ChangeColonyID(colonyName string, oldColonyID, newColonyID string) error { + return nil +} +func (db *DatabaseMock) ChangeUserID(colonyName string, oldUserID, newUserID string) error { + return nil +} +func (db *DatabaseMock) ChangeExecutorID(colonyName string, oldExecutorID, newExecutorID string) error { + return nil +} // BlueprintDatabase interface - BlueprintDefinition methods func (db *DatabaseMock) AddBlueprintDefinition(sd *core.BlueprintDefinition) error { return nil } -func (db *DatabaseMock) GetBlueprintDefinitionByID(id string) (*core.BlueprintDefinition, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintDefinitionByName(namespace, name string) (*core.BlueprintDefinition, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintDefinitions() ([]*core.BlueprintDefinition, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintDefinitionsByNamespace(namespace string) ([]*core.BlueprintDefinition, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintDefinitionsByGroup(group string) ([]*core.BlueprintDefinition, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintDefinitionByKind(kind string) (*core.BlueprintDefinition, error) { return nil, nil } +func (db *DatabaseMock) GetBlueprintDefinitionByID(id string) (*core.BlueprintDefinition, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintDefinitionByName(namespace, name string) (*core.BlueprintDefinition, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintDefinitions() ([]*core.BlueprintDefinition, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintDefinitionsByNamespace(namespace string) ([]*core.BlueprintDefinition, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintDefinitionsByGroup(group string) ([]*core.BlueprintDefinition, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintDefinitionByKind(kind string) (*core.BlueprintDefinition, error) { + return nil, nil +} func (db *DatabaseMock) UpdateBlueprintDefinition(sd *core.BlueprintDefinition) error { return nil } -func (db *DatabaseMock) RemoveBlueprintDefinitionByID(id string) error { return nil } +func (db *DatabaseMock) RemoveBlueprintDefinitionByID(id string) error { return nil } func (db *DatabaseMock) RemoveBlueprintDefinitionByName(namespace, name string) error { return nil } -func (db *DatabaseMock) CountBlueprintDefinitions() (int, error) { return 0, nil } +func (db *DatabaseMock) CountBlueprintDefinitions() (int, error) { return 0, nil } // BlueprintDatabase interface - Blueprint methods -func (db *DatabaseMock) AddBlueprint(blueprint *core.Blueprint) error { return nil } +func (db *DatabaseMock) AddBlueprint(blueprint *core.Blueprint) error { return nil } func (db *DatabaseMock) GetBlueprintByID(id string) (*core.Blueprint, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintByName(namespace, name string) (*core.Blueprint, error) { return nil, nil } +func (db *DatabaseMock) GetBlueprintByName(namespace, name string) (*core.Blueprint, error) { + return nil, nil +} func (db *DatabaseMock) GetBlueprints() ([]*core.Blueprint, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintsByNamespace(namespace string) ([]*core.Blueprint, error) { return nil, nil } +func (db *DatabaseMock) GetBlueprintsByNamespace(namespace string) ([]*core.Blueprint, error) { + return nil, nil +} func (db *DatabaseMock) GetBlueprintsByKind(kind string) ([]*core.Blueprint, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintsByNamespaceAndKind(namespace, kind string) ([]*core.Blueprint, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintsByNamespaceKindAndLocation(namespace, kind, locationName string) ([]*core.Blueprint, error) { return nil, nil } +func (db *DatabaseMock) GetBlueprintsByNamespaceAndKind(namespace, kind string) ([]*core.Blueprint, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintsByNamespaceKindAndLocation(namespace, kind, locationName string) ([]*core.Blueprint, error) { + return nil, nil +} func (db *DatabaseMock) UpdateBlueprint(blueprint *core.Blueprint) error { return nil } -func (db *DatabaseMock) UpdateBlueprintStatus(id string, status map[string]interface{}) error { return nil } -func (db *DatabaseMock) RemoveBlueprintByID(id string) error { return nil } -func (db *DatabaseMock) RemoveBlueprintByName(namespace, name string) error { return nil } -func (db *DatabaseMock) RemoveBlueprintsByNamespace(namespace string) error { return nil } -func (db *DatabaseMock) CountBlueprints() (int, error) { return 0, nil } +func (db *DatabaseMock) UpdateBlueprintStatus(id string, status map[string]interface{}) error { + return nil +} +func (db *DatabaseMock) RemoveBlueprintByID(id string) error { return nil } +func (db *DatabaseMock) RemoveBlueprintByName(namespace, name string) error { return nil } +func (db *DatabaseMock) RemoveBlueprintsByNamespace(namespace string) error { return nil } +func (db *DatabaseMock) CountBlueprints() (int, error) { return 0, nil } func (db *DatabaseMock) CountBlueprintsByNamespace(namespace string) (int, error) { return 0, nil } func (db *DatabaseMock) AddBlueprintHistory(history *core.BlueprintHistory) error { return nil } -func (db *DatabaseMock) GetBlueprintHistory(blueprintID string, limit int) ([]*core.BlueprintHistory, error) { return nil, nil } -func (db *DatabaseMock) GetBlueprintHistoryByGeneration(blueprintID string, generation int64) (*core.BlueprintHistory, error) { return nil, nil } +func (db *DatabaseMock) GetBlueprintHistory(blueprintID string, limit int) ([]*core.BlueprintHistory, error) { + return nil, nil +} +func (db *DatabaseMock) GetBlueprintHistoryByGeneration(blueprintID string, generation int64) (*core.BlueprintHistory, error) { + return nil, nil +} func (db *DatabaseMock) RemoveBlueprintHistory(blueprintID string) error { return nil } // Implement the database.Database interface -func (db *DatabaseMock) CreateTables() error { return nil } -func (db *DatabaseMock) DropTables() error { return nil } -func (db *DatabaseMock) Close() { } -func (db *DatabaseMock) Initialize() error { return nil } -func (db *DatabaseMock) Drop() error { return nil } -func (db *DatabaseMock) Lock(timeout int) error { return nil } -func (db *DatabaseMock) Unlock() error { return nil } +func (db *DatabaseMock) CreateTables() error { return nil } +func (db *DatabaseMock) DropTables() error { return nil } +func (db *DatabaseMock) Close() {} +func (db *DatabaseMock) Initialize() error { return nil } +func (db *DatabaseMock) Drop() error { return nil } +func (db *DatabaseMock) Lock(timeout int) error { return nil } +func (db *DatabaseMock) Unlock() error { return nil } func (db *DatabaseMock) ApplyRetentionPolicy(retentionPeriod int64) error { return nil } // Test utility functions -func createFakeColoniesController() (*ColoniesController, *DatabaseMock) { - // Use atomic counter to get unique ports for each test to avoid "address already in use" errors - portOffset := atomic.AddInt32(&portCounter, 1) - etcdClientPort := 24100 + int(portOffset)*10 - etcdPeerPort := 23100 + int(portOffset)*10 - relayPort := 25100 + int(portOffset)*10 - nodeName := fmt.Sprintf("etcd-%d", portOffset) - dataPath := fmt.Sprintf("/tmp/colonies/etcd-%d", portOffset) - - node := cluster.Node{Name: nodeName, Host: "localhost", EtcdClientPort: etcdClientPort, EtcdPeerPort: etcdPeerPort, RelayPort: relayPort, APIPort: constants.TESTPORT} +func newTestColoniesController(db database.Database, name string) *ColoniesController { + // Unique node names keep etcd data directories apart; dynamic ports and + // temp dirs let test packages run in parallel + offset := atomic.AddInt32(&portCounter, 1) + ports := utils.FreePortsOrPanic(4) + nodeName := fmt.Sprintf("%s-%d", name, offset) + dataPath, err := os.MkdirTemp("", "colonies-etcd-") + if err != nil { + panic(err) + } + + node := cluster.Node{Name: nodeName, Host: "localhost", EtcdClientPort: ports[0], EtcdPeerPort: ports[1], RelayPort: ports[2], APIPort: ports[3]} clusterConfig := cluster.Config{} clusterConfig.AddNode(node) + return CreateColoniesController(db, node, clusterConfig, dataPath, constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) +} + +func createFakeColoniesController() (*ColoniesController, *DatabaseMock) { dbMock := &DatabaseMock{} - return CreateColoniesController(dbMock, node, clusterConfig, dataPath, constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second), dbMock + return newTestColoniesController(dbMock, "etcd"), dbMock } func createTestColoniesController(db *postgresql.PQDatabase) *ColoniesController { - node := cluster.Node{Name: "test", Host: "localhost", EtcdClientPort: 24101, EtcdPeerPort: 23101, RelayPort: 25101, APIPort: constants.TESTPORT} - clusterConfig := cluster.Config{} - clusterConfig.AddNode(node) - return CreateColoniesController(db, node, clusterConfig, "/tmp/colonies/etcd_test", constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) + return newTestColoniesController(db, "test") } func createTestColoniesController2(db *postgresql.PQDatabase) *ColoniesController { - node := cluster.Node{Name: "test2", Host: "localhost", EtcdClientPort: 24102, EtcdPeerPort: 23102, RelayPort: 25102, APIPort: constants.TESTPORT} - clusterConfig := cluster.Config{} - clusterConfig.AddNode(node) - return CreateColoniesController(db, node, clusterConfig, "/tmp/colonies/etcd_test2", constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) -} \ No newline at end of file + return newTestColoniesController(db, "test2") +} diff --git a/pkg/server/handlers/realtime/handler_test.go b/pkg/server/handlers/realtime/handler_test.go index 895381436..8710157c9 100644 --- a/pkg/server/handlers/realtime/handler_test.go +++ b/pkg/server/handlers/realtime/handler_test.go @@ -190,7 +190,9 @@ func TestSubscribeChangeStateProcess2(t *testing.T) { err = client.Close(assignedProcess.ID, env.Executor1PrvKey) assert.Nil(t, err) - time.Sleep(5 * time.Second) + // The process state is committed once Close returns; a short pause lets + // any in-flight events drain before subscribing after the fact + time.Sleep(500 * time.Millisecond) subscription, err := client.SubscribeProcess(env.Colony1Name, addedProcess.ID, diff --git a/pkg/server/server_manager_test.go b/pkg/server/server_manager_test.go index 1b921075a..2a23896be 100644 --- a/pkg/server/server_manager_test.go +++ b/pkg/server/server_manager_test.go @@ -9,22 +9,31 @@ import ( "github.com/colonyos/colonies/pkg/constants" "github.com/colonyos/colonies/pkg/database/postgresql" "github.com/colonyos/colonies/pkg/security/crypto" + "github.com/colonyos/colonies/pkg/utils" "github.com/stretchr/testify/assert" ) +// testServerManagerNode returns a single-node cluster config with +// kernel-assigned ports so tests can run in parallel. +func testServerManagerNode(t *testing.T) cluster.Node { + ports, err := utils.FreePorts(4) + assert.Nil(t, err) + return cluster.Node{ + Name: "test-node", + Host: "localhost", + EtcdClientPort: ports[0], + EtcdPeerPort: ports[1], + RelayPort: ports[2], + APIPort: ports[3], + } +} + func TestServerManagerCreation(t *testing.T) { db, err := postgresql.PrepareTests() assert.Nil(t, err) defer db.Close() - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -32,7 +41,7 @@ func TestServerManagerCreation(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -46,14 +55,7 @@ func TestServerManagerBackendFactoryRegistration(t *testing.T) { assert.Nil(t, err) defer db.Close() - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -61,7 +63,7 @@ func TestServerManagerBackendFactoryRegistration(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -88,14 +90,7 @@ func TestServerManagerConfigManagement(t *testing.T) { assert.Nil(t, err) defer db.Close() - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -103,7 +98,7 @@ func TestServerManagerConfigManagement(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -111,16 +106,16 @@ func TestServerManagerConfigManagement(t *testing.T) { // Add gin server config ginConfig := &ServerConfig{ BackendType: GinBackendType, - Port: constants.TESTPORT + 100, - TLS: false, - TLSPrivateKeyPath: "", - TLSCertPath: "", - ExclusiveAssign: true, + Port: node.APIPort, + TLS: false, + TLSPrivateKeyPath: "", + TLSCertPath: "", + ExclusiveAssign: true, AllowExecutorReregister: false, - Retention: false, - RetentionPolicy: -1, - RetentionPeriod: 500, - Enabled: true, + Retention: false, + RetentionPolicy: -1, + RetentionPeriod: 500, + Enabled: true, } err = sm.AddServerConfig(ginConfig) @@ -135,7 +130,7 @@ func TestServerManagerConfigManagement(t *testing.T) { sm.running = true err = sm.AddServerConfig(&ServerConfig{ BackendType: "test", - Enabled: true, + Enabled: true, }) assert.NotNil(t, err) assert.Contains(t, err.Error(), "cannot add server config while server manager is running") @@ -155,14 +150,7 @@ func TestServerManagerLifecycle(t *testing.T) { err = db.SetServerID("", serverID) assert.Nil(t, err) - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT + 300, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -170,7 +158,7 @@ func TestServerManagerLifecycle(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -182,16 +170,16 @@ func TestServerManagerLifecycle(t *testing.T) { ginConfig := &ServerConfig{ BackendType: GinBackendType, - Port: constants.TESTPORT + 300, - TLS: false, - TLSPrivateKeyPath: "", - TLSCertPath: "", - ExclusiveAssign: true, + Port: node.APIPort, + TLS: false, + TLSPrivateKeyPath: "", + TLSCertPath: "", + ExclusiveAssign: true, AllowExecutorReregister: false, - Retention: false, - RetentionPolicy: -1, - RetentionPeriod: 500, - Enabled: true, + Retention: false, + RetentionPolicy: -1, + RetentionPeriod: 500, + Enabled: true, } err = sm.AddServerConfig(ginConfig) @@ -217,7 +205,7 @@ func TestServerManagerLifecycle(t *testing.T) { assert.True(t, exists) assert.Equal(t, GinBackendType, ginStatus.BackendType) assert.True(t, ginStatus.Running) - assert.Equal(t, constants.TESTPORT+300, ginStatus.Port) + assert.Equal(t, node.APIPort, ginStatus.Port) // Get specific server ginServer, exists := sm.GetServer(GinBackendType) @@ -245,14 +233,7 @@ func TestServerManagerMissingFactory(t *testing.T) { assert.Nil(t, err) defer db.Close() - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT + 400, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -260,7 +241,7 @@ func TestServerManagerMissingFactory(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -268,8 +249,8 @@ func TestServerManagerMissingFactory(t *testing.T) { // Add config without registering factory ginConfig := &ServerConfig{ BackendType: GinBackendType, - Port: constants.TESTPORT + 400, - Enabled: true, + Port: node.APIPort, + Enabled: true, } err = sm.AddServerConfig(ginConfig) @@ -295,14 +276,7 @@ func TestServerManagerStopTimeout(t *testing.T) { err = db.SetServerID("", serverID) assert.Nil(t, err) - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT + 500, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -310,7 +284,7 @@ func TestServerManagerStopTimeout(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -322,16 +296,16 @@ func TestServerManagerStopTimeout(t *testing.T) { ginConfig := &ServerConfig{ BackendType: GinBackendType, - Port: constants.TESTPORT + 500, - TLS: false, - TLSPrivateKeyPath: "", - TLSCertPath: "", - ExclusiveAssign: true, + Port: node.APIPort, + TLS: false, + TLSPrivateKeyPath: "", + TLSCertPath: "", + ExclusiveAssign: true, AllowExecutorReregister: false, - Retention: false, - RetentionPolicy: -1, - RetentionPeriod: 500, - Enabled: true, + Retention: false, + RetentionPolicy: -1, + RetentionPeriod: 500, + Enabled: true, } err = sm.AddServerConfig(ginConfig) @@ -364,14 +338,7 @@ func TestServerManagerHealthCheck(t *testing.T) { err = db.SetServerID("", serverID) assert.Nil(t, err) - node := cluster.Node{ - Name: "test-node", - Host: "localhost", - EtcdClientPort: 24100, - EtcdPeerPort: 23100, - RelayPort: 25100, - APIPort: constants.TESTPORT + 600, - } + node := testServerManagerNode(t) clusterConfig := cluster.Config{} clusterConfig.AddNode(node) @@ -379,7 +346,7 @@ func TestServerManagerHealthCheck(t *testing.T) { db, node, clusterConfig, - "/tmp/colonies/etcd", + t.TempDir(), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, ) @@ -397,16 +364,16 @@ func TestServerManagerHealthCheck(t *testing.T) { ginConfig := &ServerConfig{ BackendType: GinBackendType, - Port: constants.TESTPORT + 600, - TLS: false, - TLSPrivateKeyPath: "", - TLSCertPath: "", - ExclusiveAssign: true, + Port: node.APIPort, + TLS: false, + TLSPrivateKeyPath: "", + TLSCertPath: "", + ExclusiveAssign: true, AllowExecutorReregister: false, - Retention: false, - RetentionPolicy: -1, - RetentionPeriod: 500, - Enabled: true, + Retention: false, + RetentionPolicy: -1, + RetentionPeriod: 500, + Enabled: true, } err = sm.AddServerConfig(ginConfig) @@ -424,4 +391,4 @@ func TestServerManagerHealthCheck(t *testing.T) { // Cleanup err = sm.StopAll(5 * time.Second) assert.Nil(t, err) -} \ No newline at end of file +} diff --git a/pkg/server/test_utils.go b/pkg/server/test_utils.go index fc21fba9c..23ca016b4 100644 --- a/pkg/server/test_utils.go +++ b/pkg/server/test_utils.go @@ -2,10 +2,11 @@ package server import ( "context" + "errors" "fmt" "io/ioutil" "math/rand" - "os" + "net/http" "strconv" "testing" "time" @@ -18,7 +19,6 @@ import ( "github.com/colonyos/colonies/pkg/database/postgresql" "github.com/colonyos/colonies/pkg/rpc" "github.com/colonyos/colonies/pkg/security/crypto" - "github.com/colonyos/colonies/pkg/server/controllers" "github.com/colonyos/colonies/pkg/utils" "github.com/gin-gonic/gin" log "github.com/sirupsen/logrus" @@ -206,8 +206,16 @@ func prepareTests(t *testing.T) (*client.ColoniesClient, *Server, string, chan b } func prepareTestsWithRetention(t *testing.T, retention bool) (*client.ColoniesClient, *Server, string, chan bool) { - os.RemoveAll("/tmp/colonies") - client := client.CreateColoniesClient(constants.TESTHOST, constants.TESTPORT, Insecure, SkipTLSVerify) + // Dynamic ports and a per-test etcd data directory allow test packages to + // run in parallel without colliding on fixed ports or /tmp/colonies. The + // ports stay reserved until just before each component binds, so parallel + // test processes cannot be handed the same port during the slow parts of + // the setup (database preparation, etcd startup). + reserved := utils.ReservePortsOrPanic(4) + apiPort, etcdClientPort, etcdPeerPort, relayPort := reserved[0].Port(), reserved[1].Port(), reserved[2].Port(), reserved[3].Port() + etcdDataPath := t.TempDir() + + client := client.CreateColoniesClient(constants.TESTHOST, apiPort, Insecure, SkipTLSVerify) db, err := postgresql.PrepareTests() assert.Nil(t, err) @@ -221,14 +229,23 @@ func prepareTestsWithRetention(t *testing.T, retention bool) (*client.ColoniesCl err = db.SetServerID("", serverID) assert.Nil(t, err) - node := cluster.Node{Name: "etcd", Host: "localhost", EtcdClientPort: 24100, EtcdPeerPort: 23100, RelayPort: 25100, APIPort: constants.TESTPORT} + node := cluster.Node{Name: "etcd", Host: "localhost", EtcdClientPort: etcdClientPort, EtcdPeerPort: etcdPeerPort, RelayPort: relayPort, APIPort: apiPort} clusterConfig := cluster.Config{} clusterConfig.AddNode(node) - server := CreateServer(db, constants.TESTPORT, EnableTLS, "", "", node, clusterConfig, "/tmp/colonies/etcd", constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, false, retention, 1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) + reserved[1].Release() + reserved[2].Release() + reserved[3].Release() + server := CreateServer(db, apiPort, EnableTLS, "", "", node, clusterConfig, etcdDataPath, constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, false, retention, 1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) done := make(chan bool) + reserved[0].Release() go func() { - server.ServeForever() + err := server.ServeForever() + if err != nil && !errors.Is(err, http.ErrServerClosed) { + // A silent bind failure would leave the client talking to some + // other process's server; fail loudly instead + panic("test server failed to serve: " + err.Error()) + } db.Close() done <- true }() @@ -236,20 +253,6 @@ func prepareTestsWithRetention(t *testing.T, retention bool) (*client.ColoniesCl return client, server, serverPrvKey, done } -func createTestColoniesController(db database.Database) *controllers.ColoniesController { - node := cluster.Node{Name: "etcd", Host: "localhost", EtcdClientPort: 24100, EtcdPeerPort: 23100, RelayPort: 25100, APIPort: constants.TESTPORT} - clusterConfig := cluster.Config{} - clusterConfig.AddNode(node) - return controllers.CreateColoniesController(db, node, clusterConfig, "/tmp/colonies/etcd", constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) -} - -func createTestColoniesController2(db database.Database) *controllers.ColoniesController { - node := cluster.Node{Name: "etcd2", Host: "localhost", EtcdClientPort: 26100, EtcdPeerPort: 27100, RelayPort: 28100, APIPort: constants.TESTPORT} - clusterConfig := cluster.Config{} - clusterConfig.AddNode(node) - return controllers.CreateColoniesController(db, node, clusterConfig, "/tmp/colonies/etcd", constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) -} - func GenerateDiamondtWorkflowSpec(colonyName string) *core.WorkflowSpec { // task1 // / \ @@ -371,71 +374,30 @@ type ServerInfo struct { } func StartCluster(t *testing.T, db database.Database, size int) []ServerInfo { - os.RemoveAll("/tmp/colonies") - gin.SetMode(gin.ReleaseMode) - gin.DefaultWriter = ioutil.Discard - - clusterConfig := cluster.Config{} - for i := 0; i < size; i++ { - node := cluster.Node{ - Name: "etcd" + strconv.Itoa(i), - Host: "localhost", - EtcdClientPort: 21000 + i, - EtcdPeerPort: 22000 + i, - RelayPort: 23000 + i, - APIPort: 24000 + i} - clusterConfig.AddNode(node) - } - - crypto := crypto.CreateCrypto() - serverPrvKey, err := crypto.GeneratePrivateKey() - assert.Nil(t, err) - serverID, err := crypto.GenerateID(serverPrvKey) - assert.Nil(t, err) - - db.SetServerID("", serverID) - - sChan := make(chan ServerInfo) - for i, node := range clusterConfig.Nodes { - go func(i int, node cluster.Node) { - log.WithFields(log.Fields{"APIPort": node.APIPort}).Info("Starting ColoniesServer") - server := CreateServer(db, node.APIPort, false, "", "", node, clusterConfig, "/tmp/colonies/etcd"+strconv.Itoa(i), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, true, false, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) - done := make(chan struct{}) - s := ServerInfo{ServerID: serverID, ServerPrvKey: serverPrvKey, Server: server, Node: node, Done: done} - go func(i int) { - log.Info("ColoniesServer serving") - server.ServeForever() - log.Info("ColoniesServer stopped") - done <- struct{}{} - }(i) - sChan <- s - }(i, node) - } - - var servers []ServerInfo - for range clusterConfig.Nodes { - s := <-sChan - servers = append(servers, s) - } - - return servers + return startCluster(t, db, size, true) } // StartClusterDistributed creates a cluster with ExclusiveAssign=false for testing distributed assignment func StartClusterDistributed(t *testing.T, db database.Database, size int) []ServerInfo { - os.RemoveAll("/tmp/colonies") + return startCluster(t, db, size, false) +} + +func startCluster(t *testing.T, db database.Database, size int, exclusiveAssign bool) []ServerInfo { gin.SetMode(gin.ReleaseMode) gin.DefaultWriter = ioutil.Discard + etcdDataPath := t.TempDir() + ports := utils.FreePortsOrPanic(4 * size) + clusterConfig := cluster.Config{} for i := 0; i < size; i++ { node := cluster.Node{ Name: "etcd" + strconv.Itoa(i), Host: "localhost", - EtcdClientPort: 21000 + i, - EtcdPeerPort: 22000 + i, - RelayPort: 23000 + i, - APIPort: 24000 + i} + EtcdClientPort: ports[4*i], + EtcdPeerPort: ports[4*i+1], + RelayPort: ports[4*i+2], + APIPort: ports[4*i+3]} clusterConfig.AddNode(node) } @@ -451,8 +413,7 @@ func StartClusterDistributed(t *testing.T, db database.Database, size int) []Ser for i, node := range clusterConfig.Nodes { go func(i int, node cluster.Node) { log.WithFields(log.Fields{"APIPort": node.APIPort}).Info("Starting ColoniesServer") - // ExclusiveAssign=false for distributed assignment - server := CreateServer(db, node.APIPort, false, "", "", node, clusterConfig, "/tmp/colonies/etcd"+strconv.Itoa(i), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, false, false, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) + server := CreateServer(db, node.APIPort, false, "", "", node, clusterConfig, etcdDataPath+"/etcd"+strconv.Itoa(i), constants.GENERATOR_TRIGGER_PERIOD, constants.CRON_TRIGGER_PERIOD, exclusiveAssign, false, false, -1, 500, time.Duration(constants.DEFAULT_STALE_EXECUTOR_DURATION)*time.Second) done := make(chan struct{}) s := ServerInfo{ServerID: serverID, ServerPrvKey: serverPrvKey, Server: server, Node: node, Done: done} go func(i int) { @@ -508,15 +469,15 @@ func WaitForServerToDie(t *testing.T, s ServerInfo) { func WaitForProcessGraphs(t *testing.T, c *client.ColoniesClient, colonyName string, generatorID string, executorPrvKey string, threshold int) int { var graphs []*core.ProcessGraph var err error - retries := 40 - for i := 0; i < retries; i++ { + deadline := time.Now().Add(40 * time.Second) + for { graphs, err = c.GetWaitingProcessGraphs(colonyName, 100, executorPrvKey) assert.Nil(t, err) - if len(graphs) >= threshold { + if len(graphs) >= threshold || time.Now().After(deadline) { break } - time.Sleep(1 * time.Second) + time.Sleep(50 * time.Millisecond) } return len(graphs) diff --git a/pkg/utils/free_ports.go b/pkg/utils/free_ports.go new file mode 100644 index 000000000..c3e9e5042 --- /dev/null +++ b/pkg/utils/free_ports.go @@ -0,0 +1,105 @@ +package utils + +import ( + "net" + "strconv" +) + +// FreePorts asks the kernel for n distinct free TCP ports. The listeners are +// closed before returning, so there is a small window in which another process +// could grab a port, but ports are not reused immediately in practice. Used by +// tests to avoid fixed port numbers so packages can run in parallel. +func FreePorts(n int) ([]int, error) { + listeners := make([]net.Listener, 0, n) + ports := make([]int, 0, n) + + for i := 0; i < n; i++ { + l, err := net.Listen("tcp", "localhost:0") + if err != nil { + for _, open := range listeners { + open.Close() + } + return nil, err + } + listeners = append(listeners, l) + ports = append(ports, l.Addr().(*net.TCPAddr).Port) + } + + for _, l := range listeners { + l.Close() + } + + return ports, nil +} + +// FreePort returns a single free TCP port. +func FreePort() (int, error) { + ports, err := FreePorts(1) + if err != nil { + return 0, err + } + return ports[0], nil +} + +// FreePortsOrPanic is a convenience wrapper for test setup code without error +// returns. +func FreePortsOrPanic(n int) []int { + ports, err := FreePorts(n) + if err != nil { + panic("failed to allocate free ports: " + err.Error() + " (n=" + strconv.Itoa(n) + ")") + } + return ports +} + +// ReservedPort is a free TCP port that stays bound until Release is called, so +// no other process can be handed the same port in the meantime. Test setups +// that do slow work between allocating a port and binding it (starting etcd, +// preparing a database) should hold a reservation and release it just before +// the component binds, keeping the unbound window to microseconds. +type ReservedPort struct { + listener net.Listener + port int +} + +// Port returns the reserved port number. The port stays bound until Release. +func (r *ReservedPort) Port() int { + return r.port +} + +// Release closes the reservation so the port can be bound. Safe to call more +// than once. +func (r *ReservedPort) Release() { + if r.listener != nil { + r.listener.Close() + r.listener = nil + } +} + +// ReservePorts asks the kernel for n distinct free TCP ports and keeps them +// bound until each reservation is released. +func ReservePorts(n int) ([]*ReservedPort, error) { + reservations := make([]*ReservedPort, 0, n) + + for i := 0; i < n; i++ { + l, err := net.Listen("tcp", "localhost:0") + if err != nil { + for _, r := range reservations { + r.Release() + } + return nil, err + } + reservations = append(reservations, &ReservedPort{listener: l, port: l.Addr().(*net.TCPAddr).Port}) + } + + return reservations, nil +} + +// ReservePortsOrPanic is a convenience wrapper for test setup code without +// error returns. +func ReservePortsOrPanic(n int) []*ReservedPort { + reservations, err := ReservePorts(n) + if err != nil { + panic("failed to reserve free ports: " + err.Error() + " (n=" + strconv.Itoa(n) + ")") + } + return reservations +} diff --git a/tests/smoke/main.go b/tests/smoke/main.go deleted file mode 100644 index 4040c959c..000000000 --- a/tests/smoke/main.go +++ /dev/null @@ -1,191 +0,0 @@ -package main - -import ( - "fmt" - "math/rand" - "os" - "strconv" - "time" - - "github.com/colonyos/colonies/pkg/client" - "github.com/colonyos/colonies/pkg/core" - "github.com/colonyos/colonies/pkg/security/crypto" -) - -func init() { - rand.Seed(time.Now().UnixNano()) -} - -var letterRunes = []rune("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ") - -func randStringRunes(n int) string { - b := make([]rune, n) - for i := range b { - b[i] = letterRunes[rand.Intn(len(letterRunes))] - } - return string(b) -} - -func checkError(err error) { - if err != nil { - fmt.Println(err) - os.Exit(-1) - } -} - -func submitProcess(client *client.ColoniesClient, colonyID string, executorPrvKey string) { - funcSpec := core.FunctionSpec{ - Func: "test_func", - Args: []string{"arg1"}, - MaxWaitTime: 100, - MaxExecTime: 2, - MaxRetries: 10, - Conditions: core.Conditions{ColonyName: colonyID, ExecutorType: "bemisexecutor"}, - Env: make(map[string]string)} - - client.Submit(&funcSpec, executorPrvKey) -} - -func startCron(client *client.ColoniesClient, colonyID string, executorPrvKey string) { - funcSpec1 := core.FunctionSpec{ - Name: "cron_task1", - Func: "cron_test_func", - Args: []string{"arg1"}, - MaxWaitTime: -1, - MaxExecTime: 2, - MaxRetries: 10, - Conditions: core.Conditions{ColonyName: colonyID, ExecutorType: "bemisexecutor"}, - Env: make(map[string]string)} - - funcSpec2 := core.FunctionSpec{ - Name: "cron_task2", - Func: "cron_test_func", - Args: []string{"arg1"}, - MaxWaitTime: -1, - MaxExecTime: 2, - MaxRetries: 30, - Conditions: core.Conditions{ColonyName: colonyID, ExecutorType: "bemisexecutor"}, - Env: make(map[string]string)} - - workflowSpec := core.CreateWorkflowSpec(colonyID) - funcSpec2.AddDependency("cron_task1") - workflowSpec.AddFunctionSpec(&funcSpec1) - workflowSpec.AddFunctionSpec(&funcSpec2) - jsonStr, err := workflowSpec.ToJSON() - checkError(err) - - cron := core.CreateCron(colonyID, "test_cron1"+core.GenerateRandomID(), "1 * * * * *", -1, false, jsonStr) - _, err = client.AddCron(cron, executorPrvKey) - checkError(err) -} - -func startGenerator(client *client.ColoniesClient, colonyID string, executorPrvKey string) { - funcSpec1 := core.FunctionSpec{ - Name: "gen_task1", - Func: "gen_test_func", - Args: []string{"arg1"}, - MaxWaitTime: -1, - MaxExecTime: 2, - MaxRetries: 10, - Conditions: core.Conditions{ColonyName: colonyID, ExecutorType: "bemisexecutor"}, - Env: make(map[string]string)} - - funcSpec2 := core.FunctionSpec{ - Name: "gen_task2", - Func: "gen_test_func", - Args: []string{"arg1"}, - MaxWaitTime: -1, - MaxExecTime: 2, - MaxRetries: 30, - Conditions: core.Conditions{ColonyName: colonyID, ExecutorType: "bemisexecutor"}, - Env: make(map[string]string)} - - workflowSpec := core.CreateWorkflowSpec(colonyID) - funcSpec2.AddDependency("gen_task1") - workflowSpec.AddFunctionSpec(&funcSpec1) - workflowSpec.AddFunctionSpec(&funcSpec2) - jsonStr, err := workflowSpec.ToJSON() - checkError(err) - generator := core.CreateGenerator(colonyID, "test_genname"+core.GenerateRandomID(), jsonStr, 10) - - generator, err = client.AddGenerator(generator, executorPrvKey) - checkError(err) - - go func() { - for { - client.PackGenerator(generator.ID, randStringRunes(10), executorPrvKey) - } - }() -} - -func startExecutor(client *client.ColoniesClient, colonyID string, colonyPrvKey string) { - crypto := crypto.CreateCrypto() - executorPrvKey, err := crypto.GeneratePrivateKey() - checkError(err) - executorID, err := crypto.GenerateID(executorPrvKey) - checkError(err) - - executor := core.CreateExecutor(executorID, "bemisexecutor", core.GenerateRandomID(), colonyID, time.Now(), time.Now()) - - _, err = client.AddExecutor(executor, colonyPrvKey) - checkError(err) - - err = client.ApproveExecutor(executorID, colonyPrvKey) - checkError(err) - - go func() { - for { - assignedProcess, err := client.Assign(colonyID, 10, executorPrvKey) - if err == nil { - time.Sleep(time.Duration(rand.Intn(300)) * time.Millisecond) - client.Close(assignedProcess.ID, executorPrvKey) - } - } - }() -} - -func main() { - colonyID := os.Getenv("COLONIES_COLONY_ID") - colonyPrvKey := os.Getenv("COLONIES_COLONY_PRVKEY") - - serverHost := os.Getenv("COLONIES_SERVER_HOST") - serverPortEnvStr := os.Getenv("COLONIES_SERVER_PORT") - serverPort := -1 - var err error - if serverPortEnvStr != "" { - serverPort, err = strconv.Atoi(serverPortEnvStr) - checkError(err) - } - - tlsEnv := os.Getenv("COLONIES_TLS") - insecure := true - if tlsEnv == "true" { - insecure = false - } else if tlsEnv == "false" { - insecure = true - } - - executorPrvKey := os.Getenv("COLONIES_EXECUTOR_PRVKEY") - - client := client.CreateColoniesClient(serverHost, serverPort, insecure, true) - - for i := 0; i < 20; i++ { - startExecutor(client, colonyID, colonyPrvKey) - } - - startGenerator(client, colonyID, executorPrvKey) - startCron(client, colonyID, executorPrvKey) - - //start := time.Now().UnixNano() - for i := 0; i < 100000000; i++ { - //now := time.Now().UnixNano() - // delta := now - start - // start = time.Now().UnixNano() - //fmt.Println(i, " ", delta) - - submitProcess(client, colonyID, executorPrvKey) - } - - done := make(chan struct{}) - <-done -}