feat: add unified storage data migration step for playlists (#114582)
* fix: add type * feat: register step * feat: add playlist support * test: add test case * fix: gen mock * fix: go gen * fix: lint * fix: lint * fix: tests * fix: add resource * fix: readd * fix: address comments * fix: independent playlist query for migrations * fix: remove lock logic for sqlite * fix: handle creation and update datetimes * fix: query templating * fix: simply resources and address comments
This commit is contained in:
@@ -24,6 +24,7 @@ import (
|
||||
dashboardV1 "github.com/grafana/grafana/apps/dashboard/pkg/apis/dashboard/v1beta1"
|
||||
"github.com/grafana/grafana/apps/dashboard/pkg/migration/schemaversion"
|
||||
folders "github.com/grafana/grafana/apps/folder/pkg/apis/folder/v1beta1"
|
||||
playlistv0 "github.com/grafana/grafana/apps/playlist/pkg/apis/playlist/v0alpha1"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/apis/common/v0alpha1"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/identity"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
@@ -138,7 +139,16 @@ func NewDashboardSQLAccess(sql legacysql.LegacyDatabaseProvider,
|
||||
}
|
||||
}
|
||||
|
||||
func (a *dashboardSqlAccess) getRows(ctx context.Context, sql *legacysql.LegacyDatabaseHelper, query *DashboardQuery) (*rowsWrapper, error) {
|
||||
func (a *dashboardSqlAccess) executeQuery(ctx context.Context, helper *legacysql.LegacyDatabaseHelper, query string, args ...any) (*sql.Rows, error) {
|
||||
// Use transaction if available in context.
|
||||
// This allows us to run migrations in a transaction which is specifically required for SQLite.
|
||||
if tx := resource.TransactionFromContext(ctx); tx != nil {
|
||||
return tx.QueryContext(ctx, query, args...)
|
||||
}
|
||||
return helper.DB.GetSqlxSession().Query(ctx, query, args...)
|
||||
}
|
||||
|
||||
func (a *dashboardSqlAccess) getRows(ctx context.Context, helper *legacysql.LegacyDatabaseHelper, query *DashboardQuery) (*rowsWrapper, error) {
|
||||
ctx, span := tracer.Start(ctx, "legacy.dashboardSqlAccess.getRows")
|
||||
defer span.End()
|
||||
|
||||
@@ -150,7 +160,7 @@ func (a *dashboardSqlAccess) getRows(ctx context.Context, sql *legacysql.LegacyD
|
||||
// }
|
||||
}
|
||||
|
||||
req := newQueryReq(sql, query)
|
||||
req := newQueryReq(helper, query)
|
||||
|
||||
tmpl := sqlQueryDashboards
|
||||
if query.UseHistoryTable() && query.GetTrash {
|
||||
@@ -167,7 +177,7 @@ func (a *dashboardSqlAccess) getRows(ctx context.Context, sql *legacysql.LegacyD
|
||||
// fmt.Printf("DASHBOARD QUERY: %s [%+v] // %+v\n", pretty, req.GetArgs(), query)
|
||||
// }
|
||||
|
||||
rows, err := sql.DB.GetSqlxSession().Query(ctx, q, req.GetArgs()...)
|
||||
rows, err := a.executeQuery(ctx, helper, q, req.GetArgs()...)
|
||||
if err != nil {
|
||||
if rows != nil {
|
||||
_ = rows.Close()
|
||||
@@ -465,6 +475,132 @@ func (a *dashboardSqlAccess) MigrateLibraryPanels(ctx context.Context, orgId int
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// MigratePlaylists handles the playlist migration logic
|
||||
func (a *dashboardSqlAccess) MigratePlaylists(ctx context.Context, orgId int64, opts MigrateOptions, stream resourcepb.BulkStore_BulkProcessClient) (*BlobStoreInfo, error) {
|
||||
opts.Progress(-1, "migrating playlists...")
|
||||
rows, err := a.ListPlaylists(ctx, orgId)
|
||||
if rows != nil {
|
||||
defer func() {
|
||||
_ = rows.Close()
|
||||
}()
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Group playlist items by playlist ID
|
||||
type playlistData struct {
|
||||
id int64
|
||||
uid string
|
||||
name string
|
||||
interval string
|
||||
items []playlistv0.PlaylistItem
|
||||
createdAt int64
|
||||
updatedAt int64
|
||||
}
|
||||
|
||||
playlists := make(map[int64]*playlistData)
|
||||
var currentID int64
|
||||
var orgID int64
|
||||
var uid, name, interval string
|
||||
var createdAt, updatedAt int64
|
||||
var itemType, itemValue sql.NullString
|
||||
|
||||
count := 0
|
||||
for rows.Next() {
|
||||
err = rows.Scan(¤tID, &orgID, &uid, &name, &interval, &createdAt, &updatedAt, &itemType, &itemValue)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Get or create playlist entry
|
||||
pl, exists := playlists[currentID]
|
||||
if !exists {
|
||||
pl = &playlistData{
|
||||
id: currentID,
|
||||
uid: uid,
|
||||
name: name,
|
||||
interval: interval,
|
||||
items: []playlistv0.PlaylistItem{},
|
||||
createdAt: createdAt,
|
||||
updatedAt: updatedAt,
|
||||
}
|
||||
playlists[currentID] = pl
|
||||
}
|
||||
|
||||
// Add item if it exists (LEFT JOIN can return NULL for playlists without items)
|
||||
if itemType.Valid && itemValue.Valid {
|
||||
pl.items = append(pl.items, playlistv0.PlaylistItem{
|
||||
Type: playlistv0.PlaylistItemType(itemType.String),
|
||||
Value: itemValue.String,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
if err = rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Convert to K8s objects and send to stream
|
||||
for _, pl := range playlists {
|
||||
playlist := &playlistv0.Playlist{
|
||||
TypeMeta: metav1.TypeMeta{
|
||||
APIVersion: playlistv0.GroupVersion.String(),
|
||||
Kind: "Playlist",
|
||||
},
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: pl.uid,
|
||||
Namespace: opts.Namespace,
|
||||
CreationTimestamp: metav1.NewTime(time.UnixMilli(pl.createdAt)),
|
||||
},
|
||||
Spec: playlistv0.PlaylistSpec{
|
||||
Title: pl.name,
|
||||
Interval: pl.interval,
|
||||
Items: pl.items,
|
||||
},
|
||||
}
|
||||
|
||||
// Set updated timestamp if different from created
|
||||
if pl.updatedAt != pl.createdAt {
|
||||
meta, err := utils.MetaAccessor(playlist)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
updatedTime := time.UnixMilli(pl.updatedAt)
|
||||
meta.SetUpdatedTimestamp(&updatedTime)
|
||||
}
|
||||
|
||||
body, err := json.Marshal(playlist)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
req := &resourcepb.BulkRequest{
|
||||
Key: &resourcepb.ResourceKey{
|
||||
Namespace: opts.Namespace,
|
||||
Group: "playlist.grafana.app",
|
||||
Resource: "playlists",
|
||||
Name: pl.uid,
|
||||
},
|
||||
Value: body,
|
||||
Action: resourcepb.BulkRequest_ADDED,
|
||||
}
|
||||
|
||||
opts.Progress(count, fmt.Sprintf("%s (%d)", pl.name, len(req.Value)))
|
||||
count++
|
||||
|
||||
err = stream.Send(req)
|
||||
if err != nil {
|
||||
if errors.Is(err, io.EOF) {
|
||||
err = nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
opts.Progress(-2, fmt.Sprintf("finished playlists... (%d)", len(playlists)))
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
var _ resource.ListIterator = (*rowsWrapper)(nil)
|
||||
|
||||
type rowsWrapper struct {
|
||||
@@ -903,20 +1039,19 @@ func (a *dashboardSqlAccess) GetLibraryPanels(ctx context.Context, query Library
|
||||
return nil, err
|
||||
}
|
||||
|
||||
sqlx, err := a.sql(ctx)
|
||||
helper, err := a.sql(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
req := newLibraryQueryReq(sqlx, &query)
|
||||
req := newLibraryQueryReq(helper, &query)
|
||||
rawQuery, err := sqltemplate.Execute(sqlQueryPanels, req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("execute template %q: %w", sqlQueryPanels.Name(), err)
|
||||
}
|
||||
q := rawQuery
|
||||
|
||||
res := &dashboardV0.LibraryPanelList{}
|
||||
rows, err := sqlx.DB.GetSqlxSession().Query(ctx, q, req.GetArgs()...)
|
||||
rows, err := a.executeQuery(ctx, helper, rawQuery, req.GetArgs()...)
|
||||
defer func() {
|
||||
if rows != nil {
|
||||
_ = rows.Close()
|
||||
@@ -959,7 +1094,7 @@ func (a *dashboardSqlAccess) GetLibraryPanels(ctx context.Context, query Library
|
||||
}
|
||||
}
|
||||
if query.UID == "" {
|
||||
rv, err := sqlx.GetResourceVersion(ctx, "library_element", "updated")
|
||||
rv, err := helper.GetResourceVersion(ctx, "library_element", "updated")
|
||||
if err == nil {
|
||||
res.ResourceVersion = strconv.FormatInt(rv*1000, 10) // convert to microseconds
|
||||
}
|
||||
@@ -1038,3 +1173,29 @@ func parseLibraryPanelRow(p panel) (dashboardV0.LibraryPanel, error) {
|
||||
func (b *dashboardSqlAccess) RebuildIndexes(ctx context.Context, req *resourcepb.RebuildIndexesRequest) (*resourcepb.RebuildIndexesResponse, error) {
|
||||
return nil, fmt.Errorf("not implemented")
|
||||
}
|
||||
|
||||
func (a *dashboardSqlAccess) ListPlaylists(ctx context.Context, orgID int64) (*sql.Rows, error) {
|
||||
ctx, span := tracer.Start(ctx, "legacy.dashboardSqlAccess.ListPlaylists")
|
||||
defer span.End()
|
||||
|
||||
helper, err := a.sql(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
req := newPlaylistQueryReq(helper, &PlaylistQuery{
|
||||
OrgID: orgID,
|
||||
})
|
||||
|
||||
rawQuery, err := sqltemplate.Execute(sqlQueryPlaylists, req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("execute template %q: %w", sqlQueryPlaylists.Name(), err)
|
||||
}
|
||||
|
||||
rows, err := a.executeQuery(ctx, helper, rawQuery, req.GetArgs()...)
|
||||
if err != nil && rows != nil {
|
||||
_ = rows.Close()
|
||||
return nil, err
|
||||
}
|
||||
return rows, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user