Add support for iterators on record service

This commit is contained in:
Denis Arh
2020-03-26 19:07:14 +01:00
parent 33721bf870
commit a3b1450e16
2 changed files with 156 additions and 0 deletions
+119
View File
@@ -82,6 +82,8 @@ type (
DeleteByID(namespaceID, moduleID uint64, recordID ...uint64) error
Organize(namespaceID, moduleID, recordID uint64, sortingField, sortingValue, sortingFilter, valueField, value string) error
Iterator(f types.RecordFilter, fn eventbus.HandlerFn, action string) (err error)
}
Encoder interface {
@@ -795,6 +797,123 @@ func (svc record) Organize(namespaceID, moduleID, recordID uint64, posField, pos
})
}
// Iterator loads and iterates through list of records
//
// For each record, RecordOnIteration is generated and passed to fn()
// to be then passed to automation script that invoked the iteration
//
// No other triggers (before/after update/delete/create) are fired when (if)
// records are changed
//
// action arg enables one of the following scenarios:
// - clone: make new record (unless aborted)
// - update: update records (unless aborted)
// - delete: delete records (unless aborted)
// - default: only iterates over records, records are not changed, return value is ignored
//
//
// Iterator can be invoked only when defined in corredor script:
//
// return default {
// iterator (each) {
// return each({
// resourceType: 'compose:record',
// // action: 'update',
// filter: {
// namespace: '122709101053521922',
// module: '122709116471783426',
// query: 'Status = "foo"',
// sort: 'Status DESC',
// limit: 3,
// },
// })
// },
//
// // this is required in case of a deferred iterator
// // security: { runAs: .... } }
//
// // exec gets called for every record found by iterator
// exec () { ... }
// }
func (svc record) Iterator(f types.RecordFilter, fn eventbus.HandlerFn, action string) (err error) {
var (
invokerID = auth.GetIdentityFromContext(svc.ctx).Identity()
ns *types.Namespace
m *types.Module
set types.RecordSet
)
return svc.db.Transaction(func() (err error) {
if ns, m, _, err = svc.loadCombo(f.NamespaceID, f.ModuleID, 0); err != nil {
return
}
if !svc.ac.CanUpdateRecord(svc.ctx, m) {
return ErrNoUpdatePermissions.withStack()
}
// @todo might be good to split set into smaller chunks
set, f, err = svc.recordRepo.Find(m, f)
if err != nil {
return
}
if err = svc.preloadValues(m, set...); err != nil {
return
}
for _, rec := range set {
if err = fn(svc.ctx, event.RecordOnIteration(rec, nil, m, ns, nil)); err != nil {
if err.Error() != "Aborted" {
// When script was softly aborted (return false),
// proceed with iteration but do not clone, update or delete
// current record!
return
}
}
switch action {
case "clone":
var cln *types.Record
// Assign defaults (only on missing values)
rec.Values = svc.setDefaultValues(m, rec.Values)
// Handle payload from automation scripts
if rve := svc.procCreate(invokerID, m, rec); !rve.IsValid() {
return rve
}
if cln, err = svc.recordRepo.Create(rec); err != nil {
return
} else if err = svc.recordRepo.UpdateValues(cln.ID, cln.Values); err != nil {
return
}
case "update":
// Handle input payload
if rve := svc.procUpdate(invokerID, m, rec, rec); !rve.IsValid() {
return rve
}
if rec, err = svc.recordRepo.Update(rec); err != nil {
return
} else if err = svc.recordRepo.UpdateValues(rec.ID, rec.Values); err != nil {
return
}
case "delete":
if err = svc.recordRepo.Delete(rec); err != nil {
return err
} else if err = svc.recordRepo.DeleteValues(rec); err != nil {
return err
}
}
}
return
})
}
// loadCombo Loads everything we need for record manipulation
//
// Loads namespace, module, record and set of triggers.
+37
View File
@@ -2,6 +2,8 @@ package service
import (
"context"
"errors"
"github.com/cortezaproject/corteza-server/pkg/corredor"
"time"
"go.uber.org/zap"
@@ -141,6 +143,8 @@ func Initialize(ctx context.Context, log *zap.Logger, c Config) (err error) {
DefaultNotification = Notification()
DefaultAttachment = Attachment(DefaultStore)
RegisterIteratorProviders()
return nil
}
@@ -159,6 +163,39 @@ func Watchers(ctx context.Context) {
DefaultPermissions.Watch(ctx)
}
func RegisterIteratorProviders() {
// Register resource finders on iterator
corredor.Service().RegisterIteratorProvider(
"compose:record",
func(ctx context.Context, f map[string]string, h eventbus.HandlerFn, action string) error {
rf := types.RecordFilter{
Filter: f["filter"],
Sort: f["sort"],
}
rf.ParsePagination(f)
if nsLookup, has := f["namespace"]; !has {
return errors.New("namespace for record iteration filter not defined")
} else if ns, err := DefaultNamespace.With(ctx).FindByAny(nsLookup); err != nil {
return err
} else {
rf.NamespaceID = ns.ID
}
if mLookup, has := f["module"]; !has {
return errors.New("module for record iteration filter not defined")
} else if m, err := DefaultModule.With(ctx).FindByAny(rf.NamespaceID, mLookup); err != nil {
return err
} else {
rf.ModuleID = m.ID
}
return DefaultRecord.With(ctx).Iterator(rf, h, action)
},
)
}
// Data is stale when new date does not match updatedAt or createdAt (before first update)
func isStale(new *time.Time, updatedAt *time.Time, createdAt time.Time) bool {
if new == nil {