Storage: add support for snapshots, dataframes, and raw json objects (#57934)

This commit is contained in:
Ryan McKinley
2022-11-01 08:28:13 -07:00
committed by GitHub
parent 852d069a3c
commit 5736b46962
12 changed files with 346 additions and 25 deletions
+84 -7
View File
@@ -9,9 +9,11 @@ import (
"github.com/grafana/grafana/pkg/infra/db"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/models"
"github.com/grafana/grafana/pkg/services/dashboardsnapshots"
"github.com/grafana/grafana/pkg/services/playlist"
"github.com/grafana/grafana/pkg/services/sqlstore/session"
"github.com/grafana/grafana/pkg/services/store"
"github.com/grafana/grafana/pkg/services/store/kind/snapshot"
"github.com/grafana/grafana/pkg/services/store/object"
"github.com/grafana/grafana/pkg/services/user"
)
@@ -26,16 +28,26 @@ type objectStoreJob struct {
cfg ExportConfig
broadcaster statusBroadcaster
stopRequested bool
user *user.SignedInUser
sess *session.SessionDB
playlistService playlist.Service
store object.ObjectStoreServer
sess *session.SessionDB
playlistService playlist.Service
store object.ObjectStoreServer
dashboardsnapshots dashboardsnapshots.Service
}
func startObjectStoreJob(cfg ExportConfig, broadcaster statusBroadcaster, db db.DB, playlistService playlist.Service, store object.ObjectStoreServer) (Job, error) {
func startObjectStoreJob(user *user.SignedInUser,
cfg ExportConfig,
broadcaster statusBroadcaster,
db db.DB,
playlistService playlist.Service,
store object.ObjectStoreServer,
dashboardsnapshots dashboardsnapshots.Service,
) (Job, error) {
job := &objectStoreJob{
logger: log.New("export_to_object_store_job"),
cfg: cfg,
user: user,
broadcaster: broadcaster,
status: ExportStatus{
Running: true,
@@ -44,9 +56,10 @@ func startObjectStoreJob(cfg ExportConfig, broadcaster statusBroadcaster, db db.
Count: make(map[string]int, 10),
Index: 0,
},
sess: db.GetSqlxSession(),
playlistService: playlistService,
store: store,
sess: db.GetSqlxSession(),
playlistService: playlistService,
store: store,
dashboardsnapshots: dashboardsnapshots,
}
broadcaster(job.status)
@@ -170,6 +183,70 @@ func (e *objectStoreJob) start() {
e.status.Last = fmt.Sprintf("ITEM: %s", playlist.Uid)
e.broadcaster(e.status)
}
// TODO.. query lookup
orgIDs := []int64{1}
what = "snapshot"
for _, orgId := range orgIDs {
cmd := &dashboardsnapshots.GetDashboardSnapshotsQuery{
OrgId: orgId,
Limit: 500000,
SignedInUser: e.user,
}
err := e.dashboardsnapshots.SearchDashboardSnapshots(ctx, cmd)
if err != nil {
e.status.Status = "error: " + err.Error()
return
}
for _, dto := range cmd.Result {
m := snapshot.Model{
Name: dto.Name,
ExternalURL: dto.ExternalUrl,
Expires: dto.Expires.UnixMilli(),
}
rowUser.OrgID = dto.OrgId
rowUser.UserID = dto.UserId
snapcmd := &dashboardsnapshots.GetDashboardSnapshotQuery{
Key: dto.Key,
}
err = e.dashboardsnapshots.GetDashboardSnapshot(ctx, snapcmd)
if err == nil {
res := snapcmd.Result
m.DeleteKey = res.DeleteKey
m.ExternalURL = res.ExternalUrl
snap := res.Dashboard
m.DashboardUID = snap.Get("uid").MustString("")
snap.Del("uid")
snap.Del("id")
b, _ := snap.MarshalJSON()
m.Snapshot = b
}
_, err = e.store.Write(ctx, &object.WriteObjectRequest{
GRN: &object.GRN{
Scope: models.ObjectStoreScopeEntity,
UID: dto.Key,
Kind: models.StandardKindSnapshot,
},
Body: prettyJSON(m),
Comment: "export from snapshtts",
})
if err != nil {
e.status.Status = "error: " + err.Error()
return
}
e.status.Changed = time.Now().UnixMilli()
e.status.Index++
e.status.Count[what] += 1
e.status.Last = fmt.Sprintf("ITEM: %s", dto.Name)
e.broadcaster(e.status)
}
}
}
type dashInfo struct {
+3 -1
View File
@@ -19,6 +19,7 @@ import (
"github.com/grafana/grafana/pkg/services/live"
"github.com/grafana/grafana/pkg/services/org"
"github.com/grafana/grafana/pkg/services/playlist"
"github.com/grafana/grafana/pkg/services/store"
"github.com/grafana/grafana/pkg/services/store/object"
"github.com/grafana/grafana/pkg/setting"
)
@@ -223,6 +224,7 @@ func (ex *StandardExport) HandleRequestExport(c *models.ReqContext) response.Res
return response.Error(http.StatusLocked, "export already running", nil)
}
user := store.UserFromContext(c.Req.Context())
var job Job
broadcast := func(s ExportStatus) {
ex.broadcastStatus(c.OrgID, s)
@@ -231,7 +233,7 @@ func (ex *StandardExport) HandleRequestExport(c *models.ReqContext) response.Res
case "dummy":
job, err = startDummyExportJob(cfg, broadcast)
case "objectStore":
job, err = startObjectStoreJob(cfg, broadcast, ex.db, ex.playlistService, ex.store)
job, err = startObjectStoreJob(user, cfg, broadcast, ex.db, ex.playlistService, ex.store, ex.dashboardsnapshotsService)
case "git":
dir := filepath.Join(ex.dataDir, "export_git", fmt.Sprintf("git_%d", time.Now().Unix()))
if err := os.MkdirAll(dir, os.ModePerm); err != nil {