/
pyo
/
updater
Обзор
Документация
Войти
/
pyo
/
updater
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
develop
internal/postgres/pg.go
429 строк
10 KB
ilovepitsa
restore from another repo
16 июн 2025, 18:36
16 июн 2025, 18:36
91b9a34
Код
Авторство
О чём код?
package postgres import ( "context" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" "github.com/prometheus/client_golang/prometheus" "gitverse.ru/pyo/updater/pkg/metrics" "gitverse.ru/pyo/updater/pkg/models" ) const ( PgMetricType = "postgres" MetricLabelOperationValueDBSelect = "select" MetricLabelOperationValueDBExec = "exec" MetricLabelOperationValueDBQuery = "query" MetricLabelOperationValueDBQueryRow = "query_row" MetricLabelOperationValueDBCopyFrom = "copy_from" MetricLabelOperationValueDBCollectOneRow = "collect_one_row" MetricLabelOperationValueDBCollectRows = "collect_rows" MetricLabelOperationValueDBRowScan = "row_scan" MetricLabelOperationValueDBRowsScan = "rows_scan" ) func handleDBExecutorSelectRowErrorMetrics(err *error, operation, dbType string) { if err != nil && *err != nil { metrics.MetricDBExecutorRowSelectErrors.With(prometheus.Labels{ metrics.MetricLabelOperation: operation, metrics.MetricLabelDBType: dbType, }).Inc() } } type pgxRows struct { next pgx.Rows } func (r *pgxRows) Close() { r.next.Close() } func (r *pgxRows) Err() error { return r.next.Err() } func (r *pgxRows) CommandTag() pgconn.CommandTag { return r.next.CommandTag() } func (r *pgxRows) FieldDescriptions() []pgconn.FieldDescription { return r.next.FieldDescriptions() } func (r *pgxRows) Next() bool { return r.next.Next() } func (r *pgxRows) Scan(dest ...any) (err error) { defer handleDBExecutorSelectRowErrorMetrics(&err, MetricLabelOperationValueDBRowsScan, PgMetricType) return r.next.Scan(dest...) } func (r *pgxRows) Values() ([]any, error) { return r.next.Values() } func (r *pgxRows) RawValues() [][]byte { return r.next.RawValues() } func (r *pgxRows) Conn() *pgx.Conn { return r.next.Conn() } type pgxRow struct { next pgx.Row } func (r *pgxRow) Scan(dest ...any) (err error) { defer handleDBExecutorSelectRowErrorMetrics(&err, MetricLabelOperationValueDBRowScan, PgMetricType) return r.next.Scan(dest...) } func collectOneRowWithMetrics[T any](rows pgx.Rows, fn pgx.RowToFunc[T]) (result T, err error) { defer handleDBExecutorSelectRowErrorMetrics(&err, MetricLabelOperationValueDBCollectOneRow, PgMetricType) return pgx.CollectOneRow[T](rows, fn) } func collectRowsWithMetrics[T any](rows pgx.Rows, fn pgx.RowToFunc[T]) (results []T, err error) { defer handleDBExecutorSelectRowErrorMetrics(&err, MetricLabelOperationValueDBCollectRows, PgMetricType) return pgx.CollectRows[T](rows, fn) } type PgExecutorMetricsWrapper struct { pgExecutor PgxExecutor } func NewPgExecutorMetricsWrapper(pgExecutor PgxExecutor) PgxExecutor { return PgExecutorMetricsWrapper{ pgExecutor: pgExecutor, } } func (w PgExecutorMetricsWrapper) Exec(ctx context.Context, sql string, arguments ...any) (commandTag pgconn.CommandTag, err error) { metricsHandler := metrics.NewDBExecutorBaseMetricsHandler(PgMetricType, MetricLabelOperationValueDBExec) metricsHandler.Before() defer metricsHandler.After(&err) return w.pgExecutor.Exec(ctx, sql, arguments...) } func (w PgExecutorMetricsWrapper) Query(ctx context.Context, sql string, args ...any) (rows pgx.Rows, err error) { metricsHandler := metrics.NewDBExecutorBaseMetricsHandler(PgMetricType, MetricLabelOperationValueDBQuery) metricsHandler.Before() defer metricsHandler.After(&err) rows, err = w.pgExecutor.Query(ctx, sql, args...) return &pgxRows{next: rows}, err } func (w PgExecutorMetricsWrapper) QueryRow(ctx context.Context, sql string, args ...any) (row pgx.Row) { metricsHandler := metrics.NewDBExecutorBaseMetricsHandler(PgMetricType, MetricLabelOperationValueDBQueryRow) metricsHandler.Before() defer metricsHandler.After(nil) row = w.pgExecutor.QueryRow(ctx, sql, args...) return &pgxRow{next: row} } func (w PgExecutorMetricsWrapper) CopyFrom(ctx context.Context, tableName pgx.Identifier, columnNames []string, rowSrc pgx.CopyFromSource) (amount int64, err error) { metricsHandler := metrics.NewDBExecutorBaseMetricsHandler(PgMetricType, MetricLabelOperationValueDBCopyFrom) metricsHandler.Before() defer metricsHandler.After(&err) return w.pgExecutor.CopyFrom(ctx, tableName, columnNames, rowSrc) } type PgxExecutor interface { Exec(ctx context.Context, sql string, arguments ...any) (commandTag pgconn.CommandTag, err error) Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error) QueryRow(ctx context.Context, sql string, args ...any) pgx.Row CopyFrom(ctx context.Context, tableName pgx.Identifier, columnNames []string, rowSrc pgx.CopyFromSource) (int64, error) } type PgxConnProvider func(ctx context.Context) (PgxExecutor, error) type PG struct { conn PgxConnProvider } func NewPG(conn PgxConnProvider) *PG { return &PG{conn: conn} } func (pg *PG) GetComponents(ctx context.Context) ([]models.DBComponent, error) { const q = ` SELECT id, name, display_name, state, type, author, icon, price, homepage, description, publisher FROM component; ` conn, err := pg.conn(ctx) if err != nil { return nil, err } result, err := conn.Query(ctx, q) if err != nil { return nil, err } return collectRowsWithMetrics[models.DBComponent](result, pgx.RowToStructByName[models.DBComponent]) } func (pg *PG) GetNewerComponentsVersion(ctx context.Context, current []models.DBComponentWithVersion) ([]models.DBComponentWithVersion, error) { return nil, nil } func (pg *PG) GetComponentVersions(ctx context.Context, name string) ([]models.DBVersion, error) { const q = ` SELECT version, depend, build_options FROM versions WHERE component_name = $1 ` conn, err := pg.conn(ctx) if err != nil { return nil, err } result, err := conn.Query(ctx, q, name) if err != nil { return nil, err } return collectRowsWithMetrics[models.DBVersion](result, pgx.RowToStructByName[models.DBVersion]) } func (pg *PG) GetComponentByName(ctx context.Context, name string) (models.DBComponent, error) { const q = ` SELECT id, name, display_name, state, type, author, icon, price, homepage, description, publisher FROM component WHERE name = $1; ` conn, err := pg.conn(ctx) if err != nil { return models.DBComponent{}, err } row, err := conn.Query(ctx, q, name) if err != nil { return models.DBComponent{}, err } return collectOneRowWithMetrics(row, pgx.RowToStructByName[models.DBComponent]) } func (pg *PG) GetComponentWithVersionByNameAndVersion(ctx context.Context, name, version string) (models.DBComponentWithVersion, error) { const q = ` SELECT c.id, c.name, c.display_name, c.state, c.type, c.author, c.icon, c.price, c.homepage, c.description, c.publisher, v.version, v.depend, v.build_options FROM component as c JOIN versions as v ON c.name = v.component_name WHERE c.name = $1 AND v.version = $2; ` conn, err := pg.conn(ctx) if err != nil { return models.DBComponentWithVersion{}, err } row, err := conn.Query(ctx, q, name, version) if err != nil { return models.DBComponentWithVersion{}, err } return collectOneRowWithMetrics(row, pgx.RowToStructByName[models.DBComponentWithVersion]) } func (pg *PG) AddComponent(ctx context.Context, comp models.DBComponent) (int64, error) { const q = ` INSERT INTO component ( name, display_name, state, type, author, icon, price, homepage, description, publisher ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) RETURNING id; ` conn, err := pg.conn(ctx) if err != nil { return 0, err } id := int64(0) err = conn.QueryRow( ctx, q, comp.Name, comp.DisplayName, comp.State, comp.ComponentType, comp.Author, comp.Icon, comp.Price, comp.Homepage, comp.Description, comp.Publisher, ).Scan(&id) return id, err } func (pg *PG) AddComponentVersion(ctx context.Context, compName string, version models.DBVersion) (int64, error) { const q = ` INSERT INTO versions ( component_name, version, depend, build_options ) VALUES ($1, $2, $3, $4) RETURNING id; ` conn, err := pg.conn(ctx) if err != nil { return 0, err } id := int64(0) err = conn.QueryRow( ctx, q, compName, version.Version, version.DevDepend, version.BuildOptions, ).Scan(&id) return id, err } func (pg *PG) AddBuildOperation(ctx context.Context, op models.BuildOperation) (int, error) { const q = ` INSERT INTO build_operations ( component_id, version, created_at, updated_at, path, status, operation_type ) VALUES ($1, $2, $3, $4, $5, $6, $7) RETURNING id; ` conn, err := pg.conn(ctx) if err != nil { return 0, err } id := 0 err = conn.QueryRow(ctx, q, op.ComponentID, op.Version, op.CreatedAt, op.UpdatedAt, op.ComponentPath, op.Status, op.OperationType, ).Scan(&id) return id, err } func (pg *PG) AddComponentVersionByName(ctx context.Context, c models.Component, name string) error { return nil } func (pg *PG) UpdateComponent(ctx context.Context, comp models.DBComponent) error { return nil } func (pg *PG) UpdateVersion(ctx context.Context, compName string, version models.DBVersion) error { return nil } func (pg *PG) GetNewOperation(ctx context.Context) ([]models.BuildOperation, error) { const q = ` SELECT bo.id, bo.created_at, bo.updated_at, bo.deleted_at, bo.version, bo.component_id as component_id, c.name as component_name, c.display_name as component_display_name, c.state as component_state, c.type as component_type, c.author as component_author, c.icon as component_icon, c.price as component_price, c.homepage as component_homepage, c.description as component_description, c.publisher as component_publisher, v.depend as component_depend, v.build_options as component_build_options, bo.path as path, bo.status as status, bo.operation_type as operation_type FROM build_operations as bo JOIN component as c ON bo.component_id = c.id JOIN versions as v ON bo.version = v.version and c.name = v.component_name WHERE bo.status = 0; ` conn, err := pg.conn(ctx) if err != nil { return nil, err } result, err := conn.Query(ctx, q) if err != nil { return nil, err } return collectRowsWithMetrics[models.BuildOperation](result, pgx.RowToStructByName[models.BuildOperation]) } func (pg *PG) UpdateOperationStatus(ctx context.Context, operationID int, status int) error { const q = ` UPDATE build_operations SET status = $1, updated_at = $2 WHERE id = $3; ` now := time.Now() conn, err := pg.conn(ctx) if err != nil { return err } _, err = conn.Exec(ctx, q, status, now, operationID) return err }