initial commit for changes to support database as backend for audit logs

This commit is contained in:
Nirav Parikh
2022-11-16 19:05:38 +05:30
parent 5e8bbffb05
commit 6465afed9b
16 changed files with 1980 additions and 62 deletions
+20
View File
@@ -0,0 +1,20 @@
package query
import v1 "github.com/paralus/paralus/proto/rpc/audit"
type QueryFilters interface {
GetType() string
GetUser() string
GetClient() string
GetTimefrom() string
GetPortal() string
GetCluster() string
GetNamespace() string
GetKind() string
GetMethod() string
GetQueryString() string
GetProjects() []string
}
var _ QueryFilters = (*v1.AuditLogQueryFilter)(nil)
var _ QueryFilters = (*v1.RelayAuditQueryFilter)(nil)
+40
View File
@@ -0,0 +1,40 @@
package service
import (
v1 "github.com/paralus/paralus/proto/rpc/audit"
"github.com/uptrace/bun"
)
type AuditLogService interface {
GetAuditLog(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error)
GetAuditLogByProjects(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error)
}
func NewAuditLogElasticSearchService(url string, auditPattern string, logPrefix string) (AuditLogService, error) {
auditQuery, err := NewElasticSearchQuery(url, auditPattern, logPrefix)
if err != nil {
return nil, err
}
return &auditLogElasticSearchService{auditQuery: auditQuery}, nil
}
func NewAuditLogDatabaseService(db *bun.DB, tag string) (AuditLogService, error) {
return &auditLogDatabaseService{db: db, tag: tag}, nil
}
type RelayAuditService interface {
GetRelayAudit(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error)
GetRelayAuditByProjects(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error)
}
func NewRelayAuditDatabaseService(db *bun.DB, tag string) (RelayAuditService, error) {
return &relayAuditDatabaseService{db: db, tag: tag}, nil
}
func NewRelayAuditElasticSearchService(url string, auditPattern string, logPrefix string) (RelayAuditService, error) {
relayQuery, err := NewElasticSearchQuery(url, auditPattern, logPrefix)
if err != nil {
return nil, err
}
return &relayAuditElasticSearchService{relayQuery: relayQuery}, nil
}
+102
View File
@@ -0,0 +1,102 @@
package service
import (
"context"
"encoding/json"
"github.com/paralus/paralus/internal/dao"
"github.com/paralus/paralus/internal/models"
v1 "github.com/paralus/paralus/proto/rpc/audit"
auditv1 "github.com/paralus/paralus/proto/types/audit"
"github.com/uptrace/bun"
"google.golang.org/protobuf/types/known/structpb"
)
type auditLogDatabaseService struct {
db *bun.DB
tag string
}
func (a *auditLogDatabaseService) GetAuditLog(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
project, err := getProjectFromUrlScope(req.GetMetadata().UrlScope)
if err != nil {
return nil, err
}
req.Filter.Projects = []string{project}
return a.GetAuditLogByProjects(req)
}
func buildAggregators(aggr []models.AggregatorData) []*auditv1.GroupByType {
var groups = make([]*auditv1.GroupByType, 0)
for _, agg := range aggr {
groups = append(groups, &auditv1.GroupByType{
DocCount: int32(agg.Count),
Key: agg.Key,
})
}
return groups
}
func buildDataSource(logs []models.AuditLog) (ds []*auditv1.DataSource) {
for _, log := range logs {
data := &auditv1.Data{}
json.Unmarshal(log.Data, data)
ds = append(ds, &auditv1.DataSource{
XSource: &auditv1.DataSourceJSON{
Json: data,
},
},
)
}
return ds
}
func (a *auditLogDatabaseService) GetAuditLogByProjects(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
err = validateQueryString(req.GetFilter().QueryString)
if err != nil {
return nil, err
}
ctx := context.Background()
auditLogs, err := dao.GetAuditLogs(ctx, a.db, a.tag, req.Filter)
if err != nil {
return nil, err
}
// aggregations
_, err = dao.GetAuditLogAggregations(ctx, a.db, a.tag, "project", req.Filter)
if err != nil {
return nil, err
}
usernameAggr, err := dao.GetAuditLogAggregations(ctx, a.db, a.tag, "username", req.Filter)
if err != nil {
return nil, err
}
typeAggr, err := dao.GetAuditLogAggregations(ctx, a.db, a.tag, "type", req.Filter)
if err != nil {
return nil, err
}
response := &auditv1.AuditResponse{
Aggregations: &auditv1.Aggregations{
GroupByType: &auditv1.AggregatorGroup{
Buckets: buildAggregators(typeAggr),
},
GroupByUsername: &auditv1.AggregatorGroup{
Buckets: buildAggregators(usernameAggr),
},
},
Hits: &auditv1.Hits{Hits: buildDataSource(auditLogs)},
}
var resMap map[string]interface{}
data, _ := json.Marshal(response)
json.Unmarshal(data, &resMap)
result, _ := structpb.NewStruct(resMap)
res = &v1.GetAuditLogSearchResponse{
Result: result,
}
return res, nil
}
@@ -10,23 +10,15 @@ import (
"google.golang.org/protobuf/types/known/structpb"
)
type AuditLogService struct {
type auditLogElasticSearchService struct {
auditQuery ElasticSearchQuery
}
func NewAuditLogService(url string, auditPattern string, logPrefix string) (*AuditLogService, error) {
auditQuery, err := NewElasticSearchQuery(url, auditPattern, logPrefix)
func (a *auditLogElasticSearchService) GetAuditLog(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
if err != nil {
return nil, err
}
return &AuditLogService{auditQuery: auditQuery}, nil
}
func (a *AuditLogService) GetAuditLog(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
if err != nil {
return nil, err
}
project, err := getPrjectFromUrlScope(req.GetMetadata().UrlScope)
project, err := getProjectFromUrlScope(req.GetMetadata().UrlScope)
if err != nil {
return nil, err
}
@@ -44,7 +36,7 @@ func validateQueryString(queryString string) error {
return nil
}
func getPrjectFromUrlScope(urlScope string) (string, error) {
func getProjectFromUrlScope(urlScope string) (string, error) {
s := strings.Split(urlScope, "/")
if len(s) != 2 {
_log.Errorw("Unable to retrieve project from urlScope", "urlScope", urlScope)
@@ -53,7 +45,7 @@ func getPrjectFromUrlScope(urlScope string) (string, error) {
return s[1], nil
}
func (a *AuditLogService) GetAuditLogByProjects(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
func (a *auditLogElasticSearchService) GetAuditLogByProjects(req *v1.GetAuditLogSearchRequest) (res *v1.GetAuditLogSearchResponse, err error) {
err = validateQueryString(req.GetFilter().QueryString)
if err != nil {
return nil, err
@@ -86,7 +86,7 @@ func (m *mockElasticSearchQuery) Handle(msg bytes.Buffer) (map[string]interface{
func TestGetAuditLogByProjectsSimple(t *testing.T) {
esq := &mockElasticSearchQuery{}
al := &AuditLogService{auditQuery: esq}
al := &auditLogElasticSearchService{auditQuery: esq}
req := v1.GetAuditLogSearchRequest{
Filter: &v1.AuditLogQueryFilter{
QueryString: "query-string",
@@ -118,7 +118,7 @@ func TestGetAuditLogByProjectsSimple(t *testing.T) {
func TestGetAuditLogByProjectsNoProject(t *testing.T) {
esq := &mockElasticSearchQuery{}
al := &AuditLogService{auditQuery: esq}
al := &auditLogElasticSearchService{auditQuery: esq}
req := v1.GetAuditLogSearchRequest{
Metadata: &v3.Metadata{UrlScope: "url/project"},
Filter: &v1.AuditLogQueryFilter{
+99
View File
@@ -0,0 +1,99 @@
package service
import (
"context"
"encoding/json"
"github.com/paralus/paralus/internal/dao"
v1 "github.com/paralus/paralus/proto/rpc/audit"
auditv1 "github.com/paralus/paralus/proto/types/audit"
"github.com/uptrace/bun"
"google.golang.org/protobuf/types/known/structpb"
)
type relayAuditDatabaseService struct {
db *bun.DB
tag string
}
func (ra *relayAuditDatabaseService) GetRelayAudit(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
if err != nil {
return nil, err
}
project, err := getProjectFromUrlScope(req.GetMetadata().UrlScope)
if err != nil {
return nil, err
}
req.Filter.Projects = []string{project}
return ra.GetRelayAuditByProjects(req)
}
func (ra *relayAuditDatabaseService) GetRelayAuditByProjects(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
err = validateQueryString(req.GetFilter().QueryString)
if err != nil {
return &v1.RelayAuditResponse{}, err
}
ctx := context.Background()
auditLogs, err := dao.GetAuditLogs(ctx, ra.db, ra.tag, req.Filter)
if err != nil {
return nil, err
}
// aggregations
clusterAggr, err := dao.GetAuditLogAggregations(ctx, ra.db, ra.tag, "cluster", req.Filter)
if err != nil {
return nil, err
}
usernameAggr, err := dao.GetAuditLogAggregations(ctx, ra.db, ra.tag, "username", req.Filter)
if err != nil {
return nil, err
}
nsAggr, err := dao.GetAuditLogAggregations(ctx, ra.db, ra.tag, "namespace", req.Filter)
if err != nil {
return nil, err
}
kindAggr, err := dao.GetAuditLogAggregations(ctx, ra.db, ra.tag, "kind", req.Filter)
if err != nil {
return nil, err
}
methodAggr, err := dao.GetAuditLogAggregations(ctx, ra.db, ra.tag, "method", req.Filter)
if err != nil {
return nil, err
}
response := &auditv1.AuditResponse{
Aggregations: &auditv1.Aggregations{
GroupByCluster: &auditv1.AggregatorGroup{
Buckets: buildAggregators(clusterAggr),
},
GroupByUsername: &auditv1.AggregatorGroup{
Buckets: buildAggregators(usernameAggr),
},
GroupByNamespace: &auditv1.AggregatorGroup{
Buckets: buildAggregators(nsAggr),
},
GroupByKind: &auditv1.AggregatorGroup{
Buckets: buildAggregators(kindAggr),
},
GroupByMethod: &auditv1.AggregatorGroup{
Buckets: buildAggregators(methodAggr),
},
},
Hits: &auditv1.Hits{Hits: buildDataSource(auditLogs)},
}
var resMap map[string]interface{}
data, _ := json.Marshal(response)
json.Unmarshal(data, &resMap)
result, _ := structpb.NewStruct(resMap)
res = &v1.RelayAuditResponse{
Result: result,
}
return res, nil
}
@@ -8,23 +8,15 @@ import (
"google.golang.org/protobuf/types/known/structpb"
)
type RelayAuditService struct {
type relayAuditElasticSearchService struct {
relayQuery ElasticSearchQuery
}
func NewRelayAuditService(url string, auditPattern string, logPrefix string) (*RelayAuditService, error) {
relayQuery, err := NewElasticSearchQuery(url, auditPattern, logPrefix)
func (ra *relayAuditElasticSearchService) GetRelayAudit(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
if err != nil {
return nil, err
}
return &RelayAuditService{relayQuery: relayQuery}, nil
}
func (ra *RelayAuditService) GetRelayAudit(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
if err != nil {
return nil, err
}
project, err := getPrjectFromUrlScope(req.GetMetadata().UrlScope)
project, err := getProjectFromUrlScope(req.GetMetadata().UrlScope)
if err != nil {
return nil, err
}
@@ -32,7 +24,7 @@ func (ra *RelayAuditService) GetRelayAudit(req *v1.RelayAuditRequest) (res *v1.R
return ra.GetRelayAuditByProjects(req)
}
func (ra *RelayAuditService) GetRelayAuditByProjects(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
func (ra *relayAuditElasticSearchService) GetRelayAuditByProjects(req *v1.RelayAuditRequest) (res *v1.RelayAuditResponse, err error) {
err = validateQueryString(req.GetFilter().QueryString)
if err != nil {
return &v1.RelayAuditResponse{}, err
@@ -86,7 +86,7 @@ type rmd struct {
func TestGetRelayAuditLogByProjectsSimple(t *testing.T) {
esq := &mockElasticSearchQuery{}
al := &RelayAuditService{relayQuery: esq}
al := &relayAuditElasticSearchService{relayQuery: esq}
req := v1.RelayAuditRequest{
Filter: &v1.RelayAuditQueryFilter{
QueryString: "query-string",
@@ -122,7 +122,7 @@ func TestGetRelayAuditLogByProjectsSimple(t *testing.T) {
func TestGetRelayAuditLogByProjectsNoProject(t *testing.T) {
esq := &mockElasticSearchQuery{}
al := &RelayAuditService{relayQuery: esq}
al := &relayAuditElasticSearchService{relayQuery: esq}
req := v1.RelayAuditRequest{
Metadata: &v3.Metadata{UrlScope: "url/project"},
Filter: &v1.RelayAuditQueryFilter{