diff --git a/.golangci.yml b/.golangci.yml index 67090e9a8..28bf84791 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -50,6 +50,8 @@ linters: - wrapcheck - zerologlint - modernize + - goconst + - prealloc settings: gocyclo: min-complexity: 15 diff --git a/cmd/api/config.go b/cmd/api/config.go index 27787b695..df97309f1 100644 --- a/cmd/api/config.go +++ b/cmd/api/config.go @@ -15,6 +15,7 @@ import ( type config struct { Redis cmd.RedisConfig `yaml:"redis"` + Database cmd.DatabaseConfig `yaml:"database"` Logger cmd.LoggerConfig `yaml:"log"` API apiConfig `yaml:"api"` Web webConfig `yaml:"web"` diff --git a/cmd/cli/notification_history.go b/cmd/cli/notification_history.go index a50d40dee..9f15f5ba6 100644 --- a/cmd/cli/notification_history.go +++ b/cmd/cli/notification_history.go @@ -115,6 +115,7 @@ func mergeNotificationHistory(logger moira.Logger, database moira.Database) erro Member: eventBytes, }) } + return nil }) if err != nil { @@ -128,6 +129,7 @@ func mergeNotificationHistory(logger moira.Logger, database moira.Database) erro for _, id := range contactIDs { pipe.Del(connector.Context(), id) } + return nil }) if err != nil { @@ -135,6 +137,7 @@ func mergeNotificationHistory(logger moira.Logger, database moira.Database) erro } var totalDelCount int64 + for i, cmd := range cmds { deleted, err := cmd.(*goredis.IntCmd).Result() if err != nil { @@ -143,6 +146,7 @@ func mergeNotificationHistory(logger moira.Logger, database moira.Database) erro Error(err). Msg("failed to delete") } + totalDelCount += deleted } diff --git a/cmd/config.go b/cmd/config.go index ce5ffe2e5..de8d560a7 100644 --- a/cmd/config.go +++ b/cmd/config.go @@ -114,6 +114,20 @@ func (config *RedisConfig) GetSettings() redis.DatabaseConfig { } } +type PostgresqlConfig struct { + Master struct { + ConnectionString string + } + Replicas []struct { + ConnectionString string + } +} + +type DatabaseConfig struct { + Redis RedisConfig + Postgresql PostgresqlConfig +} + // NotificationHistoryConfig is the config which coordinates interaction with notification statistics. // E.g. how much time should we store it, or how many history items can we request from database. type NotificationHistoryConfig struct { diff --git a/database/postgresql/config.go b/database/postgresql/config.go new file mode 100644 index 000000000..989803ee8 --- /dev/null +++ b/database/postgresql/config.go @@ -0,0 +1,10 @@ +package postgresql + +type DatabaseConfig struct { + Master ReplicaConfig + Replicas []ReplicaConfig +} + +type ReplicaConfig struct { + ConnectionString string +} diff --git a/database/postgresql/contact.go b/database/postgresql/contact.go new file mode 100644 index 000000000..e94243973 --- /dev/null +++ b/database/postgresql/contact.go @@ -0,0 +1,311 @@ +package postgresql + +import ( + "fmt" + + "github.com/lib/pq" + "github.com/moira-alert/moira" +) + +func (connector *DbConnector) GetContact(contactId string) (moira.ContactData, error) { + query := ` +SELECT + c.type, + c.name, + c.value, + c.extra_message, + + u.login, + t.team_id + +FROM contacts c +LEFT JOIN users u + ON c.user_id = u.id +LEFT JOIN teams t + ON c.team_id = t.id +WHERE c.contact_id = $1 LIMIT 1; + ` + requester := func(db RDB) (moira.ContactData, error) { + row := db.QueryRowContext(connector.ctx, query, contactId) + var ( + contactType string + name string + value string + extraMessage string + userLogin *string + teamId *string + ) + err := row.Scan(&contactType, &name, &value, &extraMessage, &userLogin, &teamId) + return moira.ContactData{ + Type: contactType, + Name: name, + Value: value, + ID: contactId, + User: moira.UseString(userLogin), + Team: moira.UseString(teamId), + ExtraMessage: extraMessage, + }, err + } + + respFromReplica, err := requester(connector.db.Replica()) + if err == nil { + return respFromReplica, nil + } + + respFromMaster, err := requester(connector.db.Master()) + return respFromMaster, err +} + +func (connector *DbConnector) GetContacts(contactIDs []string) ([]*moira.ContactData, error) { + query := ` +SELECT + c.contact_id, + c.type, + c.name, + c.value, + c.extra_message, + + u.login, + t.team_id + +FROM contacts c +LEFT JOIN users u + ON c.user_id = u.id +LEFT JOIN teams t + ON c.team_id = t.id +WHERE c.contact_id = ANY($1); + ` + + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query, pq.Array(contactIDs)) + if err != nil { + return nil, fmt.Errorf("get contacts by ids error: %w", err) + } + + contacts := make([]*moira.ContactData, 0) + for responseFromReplica.Next() { + var ( + contactId string + contactType string + name string + value string + extraMessage string + userLogin *string + teamId *string + ) + err := responseFromReplica.Scan(&contactId, &contactType, &name, &value, &extraMessage, &userLogin, &teamId) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get contacts: %w", err) + } + contacts = append(contacts, &moira.ContactData{ + Type: contactType, + Name: name, + Value: value, + ID: contactId, + User: moira.UseString(userLogin), + Team: moira.UseString(teamId), + ExtraMessage: extraMessage, + }) + } + return contacts, nil +} + +func (connector *DbConnector) GetAllContacts() ([]*moira.ContactData, error) { + query := ` +SELECT + c.contact_id, + c.type, + c.name, + c.value, + c.extra_message, + + u.login, + t.team_id + +FROM contacts c +LEFT JOIN users u + ON c.user_id = u.id +LEFT JOIN teams t + ON c.team_id = t.id; + ` + + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query) + if err != nil { + return nil, fmt.Errorf("get contacts by ids error: %w", err) + } + + contacts := make([]*moira.ContactData, 0) + for responseFromReplica.Next() { + var ( + contactId string + contactType string + name string + value string + extraMessage string + userLogin *string + teamId *string + ) + err := responseFromReplica.Scan(&contactId, &contactType, &name, &value, &extraMessage, &userLogin, &teamId) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get all contacts: %w", err) + } + contacts = append(contacts, &moira.ContactData{ + Type: contactType, + Name: name, + Value: value, + ID: contactId, + User: moira.UseString(userLogin), + Team: moira.UseString(teamId), + ExtraMessage: extraMessage, + }) + } + return contacts, nil +} + +func (connector *DbConnector) SaveContact(contact *moira.ContactData) error { + if contact == nil { + return nil + } + + queryWithUser := ` +WITH selected_user AS ( + SELECT id + FROM users + WHERE login = $1 +) +INSERT INTO contacts ( + type, + name, + value, + contact_id, + user_id, + team_id, + extra_message +) +SELECT + $2, + $3, + $4, + $5, + selected_user.id, + NULL, + $6 +FROM selected_user; + ` + queryWithTeam := ` +INSERT INTO contacts ( + type, + name, + value, + contact_id, + user_id, + team_id, + extra_message +) +SELECT + $2, + $3, + $4, + $5, + NULL, + t.id, + $6 +FROM teams t +WHERE t.team_id = $1; + ` + + query := queryWithUser + ownerId := contact.User + if contact.Team != "" { + query = queryWithTeam + ownerId = contact.Team + } + _, err := connector.db.Master().ExecContext(connector.ctx, query, + ownerId, + contact.Type, + contact.Name, + contact.Value, + contact.ID, + contact.ExtraMessage, + ) + if err != nil { + return fmt.Errorf("save contact error: %w", err) + } + return nil +} + +func (connector *DbConnector) RemoveContact(contactID string) error { + //TODO: maybe we should use removing by setting flag instead of hard delete + query := ` +DELETE FROM contacts WHERE contact_id = $1 + ` + _, err := connector.db.Master().ExecContext(connector.ctx, query, contactID) + if err != nil { + return fmt.Errorf("remove contact error: %w", err) + } + return nil +} + +func (connector *DbConnector) GetUserContactIDs(login string) ([]string, error) { + query := ` +WITH selected_user AS ( + SELECT id + FROM users + WHERE login = $1 + LIMIT 1 +) +SELECT contact_id +FROM contacts +WHERE user_id = ( + SELECT id + FROM selected_user +) + ` + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query, login) + if err != nil { + return nil, fmt.Errorf("get contacts by user login error: %w", err) + } + + contacts := make([]string, 0) + for responseFromReplica.Next() { + var contactId string + err := responseFromReplica.Scan(&contactId) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get contacts ids: %w", err) + } + contacts = append(contacts, contactId) + } + return contacts, nil +} + +func (connector *DbConnector) GetTeamContactIDs(teamId string) ([]string, error) { + query := ` +WITH selected_team AS ( + SELECT id + FROM teams + WHERE team_id = $1 + LIMIT 1 +) +SELECT contact_id +FROM contacts +WHERE team_id = ( + SELECT id + FROM selected_team +) + ` + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query, teamId) + if err != nil { + return nil, fmt.Errorf("get contacts by team id error: %w", err) + } + + contacts := make([]string, 0) + for responseFromReplica.Next() { + var contactId string + err := responseFromReplica.Scan(&contactId) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get contacts ids: %w", err) + } + contacts = append(contacts, contactId) + } + return contacts, nil +} + diff --git a/database/postgresql/contacts_test.go b/database/postgresql/contacts_test.go new file mode 100644 index 000000000..5a79ad52a --- /dev/null +++ b/database/postgresql/contacts_test.go @@ -0,0 +1,111 @@ +package postgresql_test + +import ( + "testing" + + "github.com/moira-alert/moira" + "github.com/moira-alert/moira/logging/zerolog_adapter" + "github.com/stretchr/testify/require" +) + +func TestSaveContact(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + err = db.SaveContact(&moira.ContactData{ + Type: "email", + Name: "contact-1-name", + Value: "test@mail.com", + ID: "contact-1", + User: "user-1", + ExtraMessage: "", + }) + require.NoError(t, err) +} + +func TestGetContact(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + contact, err := db.GetContact("contact-1") + require.NoError(t, err) + require.Equal(t, contact, moira.ContactData{ + Type: "email", + Name: "contact-1-name", + Value: "test@mail.com", + ID: "contact-1", + User: "user-1", + ExtraMessage: "", + }) +} + +func TestGetContacts(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + contacts, err := db.GetContacts([]string{"contact-1"}) + require.NoError(t, err) + require.Equal(t, []*moira.ContactData{ + { + Type: "email", + Name: "contact-1-name", + Value: "test@mail.com", + ID: "contact-1", + User: "user-1", + ExtraMessage: "", + }, + }, contacts) +} + +func TestGetAllContacts(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + contacts, err := db.GetAllContacts() + require.NoError(t, err) + require.Equal(t, []*moira.ContactData{ + { + Type: "email", + Name: "contact-1-name", + Value: "test@mail.com", + ID: "contact-1", + User: "user-1", + ExtraMessage: "", + }, + }, contacts) +} + +func TestGetUserContactIDs(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + contacts, err := db.GetUserContactIDs("user-1") + require.NoError(t, err) + require.Equal(t, []string{"contact-1"}, contacts) +} + +func TestRemoveContact(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + err = db.RemoveContact("contact-1") + require.NoError(t, err) +} diff --git a/database/postgresql/database.go b/database/postgresql/database.go new file mode 100644 index 000000000..e38740085 --- /dev/null +++ b/database/postgresql/database.go @@ -0,0 +1,58 @@ +package postgresql + +import ( + "context" + "database/sql" + + "github.com/moira-alert/moira" + _ "github.com/jackc/pgx/v5/stdlib" +) + +const DriverName = "pgx" + +type RDB interface { + QueryRowContext(context.Context, string, ...any) *sql.Row + QueryContext(context.Context, string, ...any) (*sql.Rows, error) +} + +type RWDB interface { + RDB + ExecContext(context.Context, string, ...any) (sql.Result, error) +} + +type Database interface { + Master() RWDB + Replica() RDB +} + +type DefaultDatabase struct { + conn *sql.DB +} + +func (d *DefaultDatabase) Master() RWDB { + return d.conn +} + +func (d *DefaultDatabase) Replica() RDB { + return d.conn +} + +type DbConnector struct { + db Database + ctx context.Context +} + +func NewDatabase( + ctx context.Context, + logger moira.Logger, + config DatabaseConfig, +) (*DbConnector, error) { + p, err := sql.Open(DriverName, config.Master.ConnectionString) + + return &DbConnector{ + db: &DefaultDatabase{ + conn: p, + }, + ctx: ctx, + }, err +} diff --git a/database/postgresql/migrations.go b/database/postgresql/migrations.go new file mode 100644 index 000000000..3a24b8ac6 --- /dev/null +++ b/database/postgresql/migrations.go @@ -0,0 +1,93 @@ +package postgresql + +import ( + "context" + "database/sql" + "errors" + "fmt" + "slices" + + "github.com/moira-alert/moira/database/postgresql/migrations" +) + +func (conn *DbConnector) ApplyMigrations(ctx context.Context) error { + migrationsToApply := []migrations.Migration{ + migrations.UsersAndTeams_01(), + } + + if err := conn.createMigrationsIfNeeded(ctx); err != nil { + return err + } + + lastMigrationNumber, err := conn.getLastMigrationNumber(ctx) + if err != nil { + return err + } + + slices.SortFunc(migrationsToApply, func(a, b migrations.Migration) int { + return int(a.Number - b.Number) + }) + + lastAppliedMigrationIndex := -1 + if lastMigrationNumber != 0 { + lastAppliedMigrationIndex = slices.IndexFunc(migrationsToApply, func(m migrations.Migration) bool { + return m.Number == lastMigrationNumber + }) + if lastAppliedMigrationIndex == -1 { + return fmt.Errorf("last applied migration in database with number %d not found", lastMigrationNumber) + } + } + + for i := range migrationsToApply[lastAppliedMigrationIndex + 1:] { + // TODO: add logging + migration := migrationsToApply[i] + + err := conn.applyMigration(ctx, migration) + if err != nil { + return fmt.Errorf("error on applying migration %d: %w", migration.Number, err) + } + } + + return nil +} + +func (conn *DbConnector) applyMigration(ctx context.Context, migration migrations.Migration) error { + _, err := conn.db.Master().ExecContext(ctx, migration.ForwardSQL) + if err != nil { + return err + } + + query := ` +INSERT INTO migrations (number, applied_at) +VALUES ($1, NOW()) + ` + _, err = conn.db.Master().ExecContext(ctx, query, migration.Number) + return err +} + +func (conn *DbConnector) getLastMigrationNumber(ctx context.Context) (int64, error) { + query := "SELECT number FROM migrations ORDER BY applied_at LIMIT 1" + + var lastMigrationNumber int64 + + err := conn.db.Master().QueryRowContext(ctx, query).Scan(&lastMigrationNumber) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return -1, fmt.Errorf("error on getting last migration number: %w", err) + } + + return lastMigrationNumber, nil +} + +func (conn *DbConnector) createMigrationsIfNeeded(ctx context.Context) error { + query := ` +CREATE TABLE IF NOT EXISTS migrations ( + id SERIAL PRIMARY KEY, + number INTEGER NOT NULL, + applied_at TIMESTAMPTZ NOT NULL +) + ` + if _, err := conn.db.Master().ExecContext(ctx, query); err != nil { + return fmt.Errorf("error on creating migrations if needed: %w", err) + } + return nil +} diff --git a/database/postgresql/migrations/1-init.go b/database/postgresql/migrations/1-init.go new file mode 100644 index 000000000..7acd35975 --- /dev/null +++ b/database/postgresql/migrations/1-init.go @@ -0,0 +1,44 @@ +package migrations + +type Migration struct { + Number int64 + ForwardSQL string +} + +func UsersAndTeams_01() Migration { + return Migration{ + Number: 1, + ForwardSQL: ` +CREATE TABLE IF NOT EXISTS users ( + id SERIAL PRIMARY KEY, + login VARCHAR NOT NULL +); +CREATE TABLE IF NOT EXISTS teams ( + id SERIAL PRIMARY KEY, + team_id VARCHAR NOT NULL, + name VARCHAR NOT NULL, + description VARCHAR +); +CREATE TABLE IF NOT EXISTS team_members ( + team_id INTEGER NOT NULL REFERENCES teams(id) ON UPDATE CASCADE ON DELETE CASCADE, + user_id INTEGER NOT NULL REFERENCES users(id) ON UPDATE CASCADE ON DELETE CASCADE, + PRIMARY KEY (team_id, user_id) +); +CREATE TABLE IF NOT EXISTS contacts ( + id SERIAL PRIMARY KEY, + contact_id VARCHAR NOT NULL, + type VARCHAR NOT NULL, + name VARCHAR NOT NULL, + value VARCHAR NOT NULL, + user_id INTEGER REFERENCES users(id) ON UPDATE CASCADE ON DELETE CASCADE, + team_id INTEGER REFERENCES teams(id) ON UPDATE CASCADE ON DELETE CASCADE, + extra_message VARCHAR NOT NULL, + CHECK ( + (user_id IS NOT NULL AND team_id IS NULL) + OR + (user_id IS NULL AND team_id IS NOT NULL) + ) +); + `, + } +} diff --git a/database/postgresql/teams.go b/database/postgresql/teams.go new file mode 100644 index 000000000..580f65523 --- /dev/null +++ b/database/postgresql/teams.go @@ -0,0 +1,180 @@ +package postgresql + +/* +Scheme: teams -> team_members <- users +*/ + +import ( + "fmt" + + "github.com/moira-alert/moira" +) + +// SaveTeam saves team into postgresql. +func (connector *DbConnector) SaveTeam(teamID string, team moira.Team) error { + query := ` +INSERT INTO teams (team_id, name, description) +VALUES ($1, $2, $3);` + _, err := connector.db.Master().ExecContext(connector.ctx, query, + teamID, + team.Name, + team.Description, + ) + if err != nil { + return fmt.Errorf("save team error: %w", err) + } + return nil +} + +func (connector *DbConnector) GetAllTeams() ([]moira.Team, error) { + query := "SELECT team_id, name, description FROM teams;" + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query) + if err != nil { + return nil, fmt.Errorf("get all teams error: %w", err) + } + + teams := make([]moira.Team, 0) + for responseFromReplica.Next() { + var ( + id string + name string + description string + ) + err := responseFromReplica.Scan(&id, &name, &description) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get all teams: %w", err) + } + teams = append(teams, moira.Team{ + ID: id, + Name: name, + Description: description, + }) + } + return teams, nil +} + +func (connector *DbConnector) GetTeam(teamID string) (moira.Team, error) { + query := "SELECT id, name, description FROM teams WHERE id=$1 LIMIT 1;" + requester := func(db RDB) (moira.Team, error) { + row := db.QueryRowContext(connector.ctx, query, teamID) + var ( + id string + name string + description string + ) + err := row.Scan(&id, &name, &description) + return moira.Team{ + ID: id, + Name: name, + Description: description, + }, err + } + + respFromReplica, err := requester(connector.db.Replica()) + if err == nil { + return respFromReplica, nil + } + + respFromMaster, err := requester(connector.db.Master()) + return respFromMaster, err +} + +func (connector *DbConnector) GetTeamByName(name string) (moira.Team, error) { + query := "SELECT id, name, description FROM teams WHERE name=$1 LIMIT 1;" + requester := func(db RDB) (moira.Team, error) { + row := db.QueryRowContext(connector.ctx, query, name) + var ( + id string + name string + description string + ) + err := row.Scan(&id, &name, &description) + return moira.Team{ + ID: id, + Name: name, + Description: description, + }, err + } + + respFromReplica, err := requester(connector.db.Replica()) + if err == nil { + return respFromReplica, nil + } + + respFromMaster, err := requester(connector.db.Master()) + return respFromMaster, err +} + +func (connector *DbConnector) SaveTeamsAndUsers(teamID string, users []string, teams map[string][]string) error { + clearQuery := ` +DELETE FROM team_members WHERE team_id = $1 + ` + insertQuery := ` +INSERT INTO team_members (team_id, user_id) VALUES ($1, $2) + ` + + //TODO: interact via transactions manager + if _, err := connector.db.Master().ExecContext(connector.ctx, clearQuery, teamID); err != nil { + return fmt.Errorf("clear team members error: %w", err) + } + + for _, userID := range users { + if _, err := connector.db.Master().ExecContext(connector.ctx, insertQuery, teamID, userID); err != nil { + return fmt.Errorf("error on save user %s to group %s: %w", userID, teamID, err) + } + } + + return nil +} + +func (connector *DbConnector) GetTeamUsers(teamID string) ([]string, error) { + query := ` +SELECT user_id FROM team_members WHERE team_id = $1 + ` + + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query) + if err != nil { + return nil, fmt.Errorf("get users of team %s error: %w", teamID, err) + } + + users := make([]string, 0) + for responseFromReplica.Next() { + var userId string + err := responseFromReplica.Scan(&userId) + if err != nil { + return nil, fmt.Errorf("unmarshaling error on get user: %w", err) + } + users = append(users, userId) + } + return users, nil +} + +func (connector *DbConnector) IsTeamContainUser(teamID, userID string) (bool, error) { + query := ` +SELECT EXISTS ( + SELECT 1 FROM team_members WHERE team_id = $1 AND user_id = $2 +) + ` + responseFromReplica, err := connector.db.Replica().QueryContext(connector.ctx, query) + if err != nil { + return false, fmt.Errorf("check user %s existance in team %s error: %w", userID, teamID, err) + } + + var exists bool + if err := responseFromReplica.Scan(&exists); err != nil { + return false, fmt.Errorf("unmarshaling error on check user existance: %w", err) + } + + return exists, nil +} + +func (connector *DbConnector) DeleteTeam(teamID, userID string) error { + query := ` +DELETE FROM teams WHERE team_id = $1 + ` + if _, err := connector.db.Master().ExecContext(connector.ctx, query, teamID); err != nil { + return fmt.Errorf("clear team members error: %w", err) + } + + return nil +} diff --git a/database/postgresql/teams_test.go b/database/postgresql/teams_test.go new file mode 100644 index 000000000..b4a212a09 --- /dev/null +++ b/database/postgresql/teams_test.go @@ -0,0 +1,57 @@ +package postgresql_test + +import ( + "context" + "testing" + + "github.com/moira-alert/moira" + "github.com/moira-alert/moira/database/postgresql" + "github.com/moira-alert/moira/logging/zerolog_adapter" + "github.com/stretchr/testify/require" +) + +func newDatabase(ctx context.Context, logger moira.Logger) (*postgresql.DbConnector, error) { + db, err := postgresql.NewDatabase(ctx, logger, postgresql.DatabaseConfig{ + Master: postgresql.ReplicaConfig{ + ConnectionString: "postgresql://user:password@localhost/postgres?connect_timeout=10", + }, + }) + return db, err +} + +func TestSaveTeam(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + err = db.SaveTeam("team-1", moira.Team{ + Name: "team-1-name", + Description: "team-1-description", + }) + require.NoError(t, err) +} + +func TestGetAllTeams(t *testing.T) { + logger, err := zerolog_adapter.GetLogger("postgresql") + require.NoError(t, err) + db, err := newDatabase(t.Context(), logger) + require.NoError(t, err) + + require.NoError(t, db.ApplyMigrations(t.Context())) + // err = db.SaveTeam("team-1", moira.Team{ + // Name: "team-1-name", + // Description: "team-1-description", + // }) + teams, err := db.GetAllTeams() + require.NoError(t, err) + require.Equal(t, []moira.Team{ + { + ID: "team-1", + Name: "team-1-name", + Description: "team-1-description", + }, + }, teams) +} + diff --git a/datatypes.go b/datatypes.go index 6be9965e6..7cd8e7883 100644 --- a/datatypes.go +++ b/datatypes.go @@ -201,7 +201,7 @@ func (trigger TriggerData) GetTriggerURI(frontURI string) string { func (trigger *TriggerData) GetTags() string { var buffer bytes.Buffer for _, tag := range trigger.Tags { - buffer.WriteString(fmt.Sprintf("[%s]", tag)) + fmt.Fprintf(&buffer, "[%s]", tag) } return buffer.String() diff --git a/go.mod b/go.mod index f165fcf21..c5ef2e633 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module github.com/moira-alert/moira -go 1.24.0 +go 1.25.0 require ( github.com/Knetic/govaluate v3.0.1-0.20171022003610-9aa49832a739+incompatible @@ -51,6 +51,7 @@ require ( github.com/cenkalti/backoff/v4 v4.3.0 github.com/go-playground/validator/v10 v10.6.0 github.com/hashicorp/golang-lru/v2 v2.0.7 + github.com/jackc/pgx/v5 v5.9.2 github.com/mattermost/mattermost/server/public v0.1.9 github.com/moira-alert/blackfriday-slack v0.1.2 github.com/swaggo/http-swagger v1.3.4 @@ -200,6 +201,9 @@ require ( github.com/hashicorp/go-plugin v1.6.1 // indirect github.com/hashicorp/yamux v0.1.1 // indirect github.com/huandu/xstrings v1.5.0 // indirect + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/josharian/intern v1.0.0 // indirect github.com/leodido/go-urn v1.2.0 // indirect github.com/mailru/easyjson v0.7.7 // indirect @@ -220,6 +224,7 @@ require ( go.opentelemetry.io/auto/sdk v1.1.0 // indirect go.opentelemetry.io/otel/trace v1.38.0 // indirect go.opentelemetry.io/proto/otlp v1.7.0 // indirect + golang.org/x/sync v0.18.0 // indirect golang.org/x/tools v0.38.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20250603155806-513f23925822 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20250603155806-513f23925822 // indirect diff --git a/go.sum b/go.sum index 09d3b8c64..914b1652f 100644 --- a/go.sum +++ b/go.sum @@ -489,6 +489,14 @@ github.com/huandu/xstrings v1.5.0 h1:2ag3IFq9ZDANvthTwTiqSSZLjDc+BedvHPAp5tJy2TI github.com/huandu/xstrings v1.5.0/go.mod h1:y5/lhBue+AyNmUVz9RLU9xbLR0o4KIIExikq4ovT0aE= github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= github.com/ianlancetaylor/demangle v0.0.0-20200824232613-28f6c0f3b639/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.9.2 h1:3ZhOzMWnR4yJ+RW1XImIPsD1aNSz4T4fyP7zlQb56hw= +github.com/jackc/pgx/v5 v5.9.2/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/jcmturner/aescts/v2 v2.0.0/go.mod h1:AiaICIRyfYg35RUkr8yESTqvSy7csK90qZ5xfvvsoNs= github.com/jcmturner/dnsutils/v2 v2.0.0/go.mod h1:b0TnjGOvI/n42bZa+hmXL+kFJZsFT7G4t3HTlQ184QM= github.com/jcmturner/gofork v1.0.0/go.mod h1:MK8+TM0La+2rjBD4jE12Kj1pCCxK7d2LK/UM3ncEo0o= diff --git a/interfaces.go b/interfaces.go index 03224ba44..0af713160 100644 --- a/interfaces.go +++ b/interfaces.go @@ -76,23 +76,6 @@ type Database interface { GetTeamSubscriptionIDs(teamID string) ([]string, error) GetTagsSubscriptions(tags []string) ([]*SubscriptionData, error) - // Patterns and metrics storing - GetPatterns() ([]string, error) - AddPatternMetric(pattern, metric string) error - GetPatternMetrics(pattern string) ([]string, error) - RemovePattern(pattern string) error - RemovePatternsMetrics(pattern []string) error - RemovePatternWithMetrics(pattern string) error - - SubscribeMetricEvents(tomb *tomb.Tomb, params *SubscribeMetricEventsParams) (<-chan *MetricEvent, error) - SaveMetrics(buffer []*MatchedMetric) error - GetMetricRetention(metric string) (int64, error) - GetMetricsValues(metrics []string, from int64, until int64) (map[string][]*MetricValue, error) - RemoveMetricRetention(metric string) error - RemoveMetricValues(metric string, from, to string) (int64, error) - RemoveMetricsValues(metrics []string, toTime int64) error - GetMetricsTTLSeconds() int64 - AddTriggersToCheck(clusterKey ClusterKey, triggerIDs []string) error GetTriggersToCheck(clusterKey ClusterKey, count int) ([]string, error) GetTriggersToCheckCount(clusterKey ClusterKey) (int64, error) @@ -131,14 +114,6 @@ type Database interface { IsTeamContainUser(teamID, userID string) (bool, error) DeleteTeam(teamID, userID string) error - // Metrics management - CleanUpOutdatedMetrics(duration time.Duration) error - CleanUpFutureMetrics(duration time.Duration) error - CleanupOutdatedPatternMetrics() (int64, error) - CleanUpAbandonedRetentions() error - RemoveMetricsByPrefix(pattern string) error - RemoveAllMetrics() error - // Delivery checks DeliveryCheckerDatabase @@ -150,6 +125,34 @@ type Database interface { // ScheduledNotification storing ScheduledNotificationsDatabase + + // Patterns and metrics storing + PatternsAndMetricsDatabase +} + +type PatternsAndMetricsDatabase interface { + GetPatterns() ([]string, error) + AddPatternMetric(pattern, metric string) error + GetPatternMetrics(pattern string) ([]string, error) + RemovePattern(pattern string) error + RemovePatternsMetrics(pattern []string) error + RemovePatternWithMetrics(pattern string) error + + SubscribeMetricEvents(tomb *tomb.Tomb, params *SubscribeMetricEventsParams) (<-chan *MetricEvent, error) + SaveMetrics(buffer []*MatchedMetric) error + GetMetricRetention(metric string) (int64, error) + GetMetricsValues(metrics []string, from int64, until int64) (map[string][]*MetricValue, error) + RemoveMetricRetention(metric string) error + RemoveMetricValues(metric string, from, to string) (int64, error) + RemoveMetricsValues(metrics []string, toTime int64) error + GetMetricsTTLSeconds() int64 + + CleanUpOutdatedMetrics(duration time.Duration) error + CleanUpFutureMetrics(duration time.Duration) error + CleanupOutdatedPatternMetrics() (int64, error) + CleanUpAbandonedRetentions() error + RemoveMetricsByPrefix(pattern string) error + RemoveAllMetrics() error } // ScheduledNotificationsDatabase is used to schedule and fetch notifications, as well as to view list of all notifications and delete some when needed. diff --git a/notifier/selfstate/check.go b/notifier/selfstate/check.go index b304aa9bb..0ae846930 100644 --- a/notifier/selfstate/check.go +++ b/notifier/selfstate/check.go @@ -227,7 +227,7 @@ func (selfCheck *SelfCheckWorker) buildTriggersTableForSubscription(subscription for _, link := range triggersTable { builder.WriteString("- ") - builder.WriteString(fmt.Sprintf("[%s](%s)", strings.Join(link.Tags, "|"), link.Link)) + fmt.Fprintf(&builder, "[%s](%s)", strings.Join(link.Tags, "|"), link.Link) builder.WriteString("\n") } diff --git a/senders/opsgenie/send.go b/senders/opsgenie/send.go index 57e87a3b3..d1183f8f8 100644 --- a/senders/opsgenie/send.go +++ b/senders/opsgenie/send.go @@ -140,7 +140,7 @@ func (sender *Sender) buildTitle(events moira.NotificationEvents, trigger moira. for len([]rune(title)) > titleLimit { var tagBuffer bytes.Buffer for i := 0; i < len(trigger.Tags)-tags; i++ { - tagBuffer.WriteString(fmt.Sprintf("[%s]", trigger.Tags[i])) + fmt.Fprintf(&tagBuffer, "[%s]", trigger.Tags[i]) } title = fmt.Sprintf("%s %s %s.... (%d)", state, trigger.Name, tagBuffer.String(), len(events)) diff --git a/senders/pushover/pushover.go b/senders/pushover/pushover.go index e2c3bfe16..c14d0a22e 100644 --- a/senders/pushover/pushover.go +++ b/senders/pushover/pushover.go @@ -105,17 +105,17 @@ func (sender *Sender) buildMessage(events moira.NotificationEvents, throttled bo break } - message.WriteString(fmt.Sprintf("%s: %s = %s (%s to %s)", event.FormatTimestamp(sender.location, moira.DefaultTimeFormat), event.Metric, event.GetMetricsValues(moira.DefaultNotificationSettings), event.OldState, event.State)) + fmt.Fprintf(&message, "%s: %s = %s (%s to %s)", event.FormatTimestamp(sender.location, moira.DefaultTimeFormat), event.Metric, event.GetMetricsValues(moira.DefaultNotificationSettings), event.OldState, event.State) if msg := event.CreateMessage(sender.location); len(msg) > 0 { - message.WriteString(fmt.Sprintf(". %s\n", msg)) + fmt.Fprintf(&message, ". %s\n", msg) } else { message.WriteString("\n") } } if len(events) > printEventsCount { - message.WriteString(fmt.Sprintf("\n...and %d more events.", len(events)-printEventsCount)) + fmt.Fprintf(&message, "\n...and %d more events.", len(events)-printEventsCount) } if throttled { @@ -133,7 +133,7 @@ func (sender *Sender) buildTitle(events moira.NotificationEvents, trigger moira. for len([]rune(title)) > titleLimit { var tagBuffer bytes.Buffer for i := 0; i < len(trigger.Tags)-tags; i++ { - tagBuffer.WriteString(fmt.Sprintf("[%s]", trigger.Tags[i])) + fmt.Fprintf(&tagBuffer, "[%s]", trigger.Tags[i]) } title = fmt.Sprintf("%s %s %s.... (%d)", state, trigger.Name, tagBuffer.String(), len(events)) diff --git a/senders/twilio/sms.go b/senders/twilio/sms.go index 96716c5da..da72a6053 100644 --- a/senders/twilio/sms.go +++ b/senders/twilio/sms.go @@ -38,22 +38,22 @@ func (sender *twilioSenderSms) buildMessage(events moira.NotificationEvents, tri state := events.GetCurrentState(throttled) - message.WriteString(fmt.Sprintf("%s %s %s (%d)\n", state, trigger.Name, trigger.GetTags(), len(events))) + fmt.Fprintf(&message, "%s %s %s (%d)\n", state, trigger.Name, trigger.GetTags(), len(events)) for i, event := range events { if i > printEventsCount-1 { break } - message.WriteString(fmt.Sprintf("\n%s: %s = %s (%s to %s)", event.FormatTimestamp(sender.location, moira.DefaultTimeFormat), event.Metric, event.GetMetricsValues(moira.DefaultNotificationSettings), event.OldState, event.State)) + fmt.Fprintf(&message, "\n%s: %s = %s (%s to %s)", event.FormatTimestamp(sender.location, moira.DefaultTimeFormat), event.Metric, event.GetMetricsValues(moira.DefaultNotificationSettings), event.OldState, event.State) if msg := event.CreateMessage(sender.location); len(msg) > 0 { - message.WriteString(fmt.Sprintf(". %s", msg)) + fmt.Fprintf(&message, ". %s", msg) } } if len(events) > printEventsCount { - message.WriteString(fmt.Sprintf("\n\n...and %d more events.", len(events)-printEventsCount)) + fmt.Fprintf(&message, "\n\n...and %d more events.", len(events)-printEventsCount) } if throttled {