Files
ITWorx Pulse release export bd774932d5
Public source validation / validate (push) Failing after 3m8s
Publish ITWorx Pulse source
2026-09-03 02:09:19 +02:00

285 lines
12 KiB
Go

package dashboard
import (
"context"
"crypto/rand"
"crypto/sha256"
"encoding/json"
"errors"
"fmt"
"reflect"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var ErrConflict = errors.New("dashboard revision conflict")
var ErrNotFound = errors.New("dashboard not found")
var ErrForbidden = errors.New("dashboard access denied")
type Summary struct {
ID, Slug, Name, Description, OwnerID, Scope string
ArchivedAt *time.Time
Revision int64
CurrentVersion int
}
type Version struct {
ID, DashboardID string
Number, SchemaVersion int
Document Document
ChangeSummary, CreatedBy string
CreatedAt time.Time
}
type Repository struct{ Pool *pgxpool.Pool }
func (r Repository) CanAccess(ctx context.Context, id, actor string) (bool, error) {
if r.Pool == nil {
return false, errors.New("dashboard repository is not configured")
}
var allowed bool
err := r.Pool.QueryRow(ctx, `SELECT EXISTS (
SELECT 1 FROM dashboards
WHERE id=$1 AND archived_at IS NULL
AND (scope IN ('shared','system') OR owner_user_id=(SELECT id FROM users WHERE external_subject=$2))
)`, id, actor).Scan(&allowed)
if err != nil {
return false, fmt.Errorf("check dashboard access: %w", err)
}
return allowed, nil
}
func (r Repository) Create(ctx context.Context, actor string, document Document, changeSummary string) (Summary, Version, error) {
if err := Validate(document); err != nil {
return Summary{}, Version{}, err
}
id, ok := document["id"].(string)
if !ok || id == "" {
return Summary{}, Version{}, errors.New("dashboard id is required")
}
slug, _ := document["slug"].(string)
name, _ := document["name"].(string)
scope, _ := document["scope"].(string)
description, _ := document["description"].(string)
tx, err := r.Pool.BeginTx(ctx, pgx.TxOptions{})
if err != nil {
return Summary{}, Version{}, fmt.Errorf("begin dashboard create: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
docJSON, _ := json.Marshal(document)
versionID := idForVersion(id, 1)
if _, err = tx.Exec(ctx, `INSERT INTO dashboards (id,slug,name,description,owner_user_id,scope) VALUES ($1,$2,$3,$4,(SELECT id FROM users WHERE external_subject=$5),$6)`, id, slug, name, description, actor, scope); err != nil {
err = mapDatabaseError(err)
return Summary{}, Version{}, fmt.Errorf("create dashboard: %w", err)
}
if _, err = tx.Exec(ctx, `INSERT INTO dashboard_versions (id,dashboard_id,version_number,schema_version,document,change_summary,created_by) VALUES ($1,$2,1,$3,$4,$5,(SELECT id FROM users WHERE external_subject=$6))`, versionID, id, CurrentSchemaVersion, docJSON, changeSummary, actor); err != nil {
return Summary{}, Version{}, fmt.Errorf("create dashboard version: %w", err)
}
if _, err = tx.Exec(ctx, `UPDATE dashboards SET current_version_id=$1 WHERE id=$2`, versionID, id); err != nil {
return Summary{}, Version{}, fmt.Errorf("set current dashboard version: %w", err)
}
if err = tx.Commit(ctx); err != nil {
return Summary{}, Version{}, fmt.Errorf("commit dashboard create: %w", err)
}
createdSummary, createdVersion, err := r.Get(ctx, id, actor)
if err != nil {
return Summary{}, Version{}, err
}
return createdSummary, createdVersion, nil
}
func (r Repository) UpdateDocument(ctx context.Context, id, actor string, expected int64, document Document, summary string) (Summary, error) {
if err := Validate(document); err != nil {
return Summary{}, err
}
tx, err := r.Pool.BeginTx(ctx, pgx.TxOptions{})
if err != nil {
return Summary{}, err
}
defer func() { _ = tx.Rollback(ctx) }()
var currentVersionID string
var currentVersion int
var currentJSON []byte
var current Summary
var canEdit bool
err = tx.QueryRow(ctx, `SELECT id,slug,name,description,COALESCE(owner_user_id::text,''),scope,archived_at,revision,current_version_id,(owner_user_id=(SELECT id FROM users WHERE external_subject=$2) OR scope IN ('shared','system')) FROM dashboards WHERE id=$1 FOR UPDATE`, id, actor).Scan(&current.ID, &current.Slug, &current.Name, &current.Description, &current.OwnerID, &current.Scope, &current.ArchivedAt, &current.Revision, &currentVersionID, &canEdit)
if err != nil {
return Summary{}, err
}
if !canEdit {
return Summary{}, ErrForbidden
}
if current.Revision != expected {
return Summary{}, ErrConflict
}
if err = tx.QueryRow(ctx, `SELECT version_number,document FROM dashboard_versions WHERE id=$1`, currentVersionID).Scan(&currentVersion, &currentJSON); err != nil {
return Summary{}, err
}
newJSON, _ := json.Marshal(document)
var stored, normalized Document
_ = json.Unmarshal(currentJSON, &stored)
_ = json.Unmarshal(newJSON, &normalized)
if reflect.DeepEqual(stored, normalized) {
current.CurrentVersion = currentVersion
return current, nil
}
next := currentVersion + 1
versionID := idForVersion(id, next)
if _, err = tx.Exec(ctx, `INSERT INTO dashboard_versions (id,dashboard_id,version_number,schema_version,document,change_summary,created_by) VALUES ($1,$2,$3,$4,$5,$6,(SELECT id FROM users WHERE external_subject=$7))`, versionID, id, next, CurrentSchemaVersion, newJSON, summary, actor); err != nil {
return Summary{}, err
}
tag, err := tx.Exec(ctx, `UPDATE dashboards SET current_version_id=$1,revision=revision+1,updated_at=now() WHERE id=$2 AND revision=$3`, versionID, id, expected)
if err != nil {
return Summary{}, err
}
if tag.RowsAffected() != 1 {
return Summary{}, ErrConflict
}
current.Revision++
current.CurrentVersion = next
if err = tx.Commit(ctx); err != nil {
return Summary{}, err
}
return current, nil
}
func (r Repository) Restore(ctx context.Context, id, actor string, expected int64, versionNumber int) (Summary, error) {
version, err := r.GetVersion(ctx, id, actor, versionNumber)
if err != nil {
return Summary{}, err
}
return r.UpdateDocument(ctx, id, actor, expected, version.Document, fmt.Sprintf("restore version %d", versionNumber))
}
func (r Repository) Get(ctx context.Context, id, actor string) (Summary, Version, error) {
var s Summary
var versionID string
err := r.Pool.QueryRow(ctx, `SELECT id,slug,name,description,COALESCE(owner_user_id::text,''),scope,archived_at,revision,current_version_id FROM dashboards WHERE id=$1 AND (scope IN ('shared','system') OR owner_user_id=(SELECT id FROM users WHERE external_subject=$2))`, id, actor).Scan(&s.ID, &s.Slug, &s.Name, &s.Description, &s.OwnerID, &s.Scope, &s.ArchivedAt, &s.Revision, &versionID)
if errors.Is(err, pgx.ErrNoRows) {
return Summary{}, Version{}, ErrNotFound
}
if err != nil {
return Summary{}, Version{}, err
}
var raw []byte
var v Version
err = r.Pool.QueryRow(ctx, `SELECT id,dashboard_id,version_number,schema_version,document,change_summary,COALESCE(created_by::text,''),created_at FROM dashboard_versions WHERE id=$1`, versionID).Scan(&v.ID, &v.DashboardID, &v.Number, &v.SchemaVersion, &raw, &v.ChangeSummary, &v.CreatedBy, &v.CreatedAt)
if err != nil {
return Summary{}, Version{}, err
}
if err = json.Unmarshal(raw, &v.Document); err != nil {
return Summary{}, Version{}, errors.New("invalid stored dashboard document")
}
s.CurrentVersion = v.Number
return s, v, nil
}
func (r Repository) List(ctx context.Context, actor string, limit int) ([]Summary, error) {
if limit < 1 || limit > 100 {
return nil, errors.New("dashboard page limit is invalid")
}
rows, err := r.Pool.Query(ctx, `SELECT d.id,d.slug,d.name,d.description,COALESCE(d.owner_user_id::text,''),d.scope,d.archived_at,d.revision,v.version_number FROM dashboards d JOIN dashboard_versions v ON v.id=d.current_version_id WHERE d.scope <> 'personal' OR d.owner_user_id=(SELECT id FROM users WHERE external_subject=$1) ORDER BY d.name ASC,d.id ASC LIMIT $2`, actor, limit)
if err != nil {
return nil, fmt.Errorf("list dashboards: %w", err)
}
defer rows.Close()
result := make([]Summary, 0, limit)
for rows.Next() {
var s Summary
if err := rows.Scan(&s.ID, &s.Slug, &s.Name, &s.Description, &s.OwnerID, &s.Scope, &s.ArchivedAt, &s.Revision, &s.CurrentVersion); err != nil {
return nil, fmt.Errorf("scan dashboard: %w", err)
}
result = append(result, s)
}
return result, rows.Err()
}
func (r Repository) UpdateMetadata(ctx context.Context, id, actor string, expected int64, name, description string) (Summary, error) {
tag, err := r.Pool.Exec(ctx, `UPDATE dashboards SET name=$1,description=$2,revision=revision+1,updated_at=now() WHERE id=$3 AND revision=$4 AND (owner_user_id=(SELECT id FROM users WHERE external_subject=$5) OR scope IN ('shared','system'))`, name, description, id, expected, actor)
if err != nil {
return Summary{}, fmt.Errorf("update dashboard metadata: %w", err)
}
if tag.RowsAffected() != 1 {
return Summary{}, ErrConflict
}
s, _, err := r.Get(ctx, id, actor)
return s, err
}
func (r Repository) Archive(ctx context.Context, id, actor string, expected int64) (Summary, error) {
tag, err := r.Pool.Exec(ctx, `UPDATE dashboards SET archived_at=now(),revision=revision+1,updated_at=now() WHERE id=$1 AND revision=$2 AND (owner_user_id=(SELECT id FROM users WHERE external_subject=$3) OR scope IN ('shared','system'))`, id, expected, actor)
if err != nil {
return Summary{}, fmt.Errorf("archive dashboard: %w", err)
}
if tag.RowsAffected() != 1 {
return Summary{}, ErrConflict
}
s, _, err := r.Get(ctx, id, actor)
return s, err
}
func (r Repository) Versions(ctx context.Context, id, actor string, limit int) ([]Version, error) {
if limit < 1 || limit > 100 {
return nil, errors.New("version page limit is invalid")
}
rows, err := r.Pool.Query(ctx, `SELECT v.id,v.dashboard_id,v.version_number,v.schema_version,v.document,v.change_summary,COALESCE(v.created_by::text,''),v.created_at FROM dashboard_versions v JOIN dashboards d ON d.id=v.dashboard_id WHERE v.dashboard_id=$1 AND (d.scope IN ('shared','system') OR d.owner_user_id=(SELECT id FROM users WHERE external_subject=$2)) ORDER BY v.version_number DESC LIMIT $3`, id, actor, limit)
if err != nil {
return nil, fmt.Errorf("list dashboard versions: %w", err)
}
defer rows.Close()
result := make([]Version, 0, limit)
for rows.Next() {
var v Version
var raw []byte
if err := rows.Scan(&v.ID, &v.DashboardID, &v.Number, &v.SchemaVersion, &raw, &v.ChangeSummary, &v.CreatedBy, &v.CreatedAt); err != nil {
return nil, err
}
if err := json.Unmarshal(raw, &v.Document); err != nil {
return nil, errors.New("invalid stored dashboard document")
}
result = append(result, v)
}
return result, rows.Err()
}
func (r Repository) Clone(ctx context.Context, actor, sourceID, slug, name string) (Summary, Version, error) {
summary, version, err := r.Get(ctx, sourceID, actor)
if err != nil {
return Summary{}, Version{}, err
}
var allowed bool
if err := r.Pool.QueryRow(ctx, `SELECT EXISTS (SELECT 1 FROM dashboards WHERE id=$1 AND (scope <> 'personal' OR owner_user_id=(SELECT id FROM users WHERE external_subject=$2)))`, sourceID, actor).Scan(&allowed); err != nil {
return Summary{}, Version{}, err
}
if !allowed {
return Summary{}, Version{}, ErrForbidden
}
if slug == "" {
slug = summary.Slug + "-copy"
}
if name == "" {
name = summary.Name + " (kopie)"
}
copyDoc := clone(version.Document)
copyDoc["id"] = newUUID()
copyDoc["slug"] = slug
copyDoc["name"] = name
copyDoc["scope"] = "personal"
return r.Create(ctx, actor, copyDoc, "cloned dashboard")
}
func newUUID() string {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
return "00000000-0000-4000-8000-000000000000"
}
b[6] = (b[6] & 0x0f) | 0x40
b[8] = (b[8] & 0x3f) | 0x80
return fmt.Sprintf("%08x-%04x-%04x-%04x-%012x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16])
}
func idForVersion(id string, number int) string {
digest := sha256.Sum256([]byte(fmt.Sprintf("%s:%d", id, number)))
digest[6] = (digest[6] & 0x0f) | 0x50
digest[8] = (digest[8] & 0x3f) | 0x80
return fmt.Sprintf("%08x-%04x-%04x-%04x-%012x", digest[0:4], digest[4:6], digest[6:8], digest[8:10], digest[10:16])
}