/
aggregate.go
158 lines (139 loc) · 5.58 KB
/
aggregate.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
package comment
import (
"context"
"github.com/openline-ai/openline-customer-os/packages/server/customer-os-common-module/utils"
"github.com/openline-ai/openline-customer-os/packages/server/events-processing-platform/constants"
"github.com/openline-ai/openline-customer-os/packages/server/events-processing-platform/domain/common/aggregate"
commonmodel "github.com/openline-ai/openline-customer-os/packages/server/events-processing-platform/domain/common/model"
"github.com/openline-ai/openline-customer-os/packages/server/events-processing-platform/eventstore"
"github.com/openline-ai/openline-customer-os/packages/server/events-processing-platform/tracing"
commentpb "github.com/openline-ai/openline-customer-os/packages/server/events-processing-proto/gen/proto/go/api/grpc/v1/comment"
"github.com/opentracing/opentracing-go"
"github.com/opentracing/opentracing-go/log"
"github.com/pkg/errors"
"strings"
)
const (
CommentAggregateType eventstore.AggregateType = "comment"
)
type CommentAggregate struct {
*aggregate.CommonTenantIdAggregate
Comment *Comment
}
func GetCommentObjectID(aggregateID string, tenant string) string {
return aggregate.GetAggregateObjectID(aggregateID, tenant, CommentAggregateType)
}
func NewCommentAggregateWithTenantAndID(tenant, id string) *CommentAggregate {
commentAggregate := CommentAggregate{}
commentAggregate.CommonTenantIdAggregate = aggregate.NewCommonAggregateWithTenantAndId(CommentAggregateType, tenant, id)
commentAggregate.SetWhen(commentAggregate.When)
commentAggregate.Comment = &Comment{}
commentAggregate.Tenant = tenant
return &commentAggregate
}
func (a *CommentAggregate) HandleGRPCRequest(ctx context.Context, request any, params map[string]any) (any, error) {
span, ctx := opentracing.StartSpanFromContext(ctx, "CommentAggregate.HandleGRPCRequest")
defer span.Finish()
switch r := request.(type) {
case *commentpb.UpsertCommentGrpcRequest:
return a.UpsertCommentGrpcRequest(ctx, r)
default:
return nil, nil
}
}
func (a *CommentAggregate) UpsertCommentGrpcRequest(ctx context.Context, request *commentpb.UpsertCommentGrpcRequest) (string, error) {
span, _ := opentracing.StartSpanFromContext(ctx, "CommentAggregate.UpsertCommentGrpcRequest")
defer span.Finish()
span.SetTag(tracing.SpanTagTenant, a.GetTenant())
span.SetTag(tracing.SpanTagAggregateId, a.GetID())
span.LogFields(log.Int64("AggregateVersion", a.GetVersion()))
var err error
var event eventstore.Event
dataFields := CommentDataFields{
Content: request.Content,
ContentType: request.ContentType,
AuthorUserId: request.AuthorUserId,
CommentedIssueId: request.CommentedIssueId,
}
source := commonmodel.Source{}
source.FromGrpc(request.SourceFields)
externalSystem := commonmodel.ExternalSystem{}
externalSystem.FromGrpc(request.ExternalSystemFields)
createdAtNotNil := utils.IfNotNilTimeWithDefault(utils.TimestampProtoToTimePtr(request.CreatedAt), utils.Now())
updatedAtNotNil := utils.IfNotNilTimeWithDefault(utils.TimestampProtoToTimePtr(request.UpdatedAt), createdAtNotNil)
if eventstore.IsAggregateNotFound(a) {
event, err = NewCommentCreateEvent(a, dataFields, source, externalSystem, createdAtNotNil, updatedAtNotNil)
} else {
event, err = NewCommentUpdateEvent(a, dataFields.Content, dataFields.ContentType, source.Source, externalSystem, updatedAtNotNil)
}
if err != nil {
tracing.TraceErr(span, err)
return "", errors.Wrap(err, "CommentAggregate.UpsertCommentGrpcRequest failed to create event")
}
aggregate.EnrichEventWithMetadataExtended(&event, span, aggregate.EventMetadata{
Tenant: request.Tenant,
UserId: request.UserId,
App: source.AppSource,
})
return request.Id, a.Apply(event)
}
func (a *CommentAggregate) When(evt eventstore.Event) error {
switch evt.GetEventType() {
case CommentCreateV1:
return a.onCommentCreate(evt)
case CommentUpdateV1:
return a.onCommentUpdate(evt)
default:
if strings.HasPrefix(evt.GetEventType(), constants.EsInternalStreamPrefix) {
return nil
}
err := eventstore.ErrInvalidEventType
err.EventType = evt.GetEventType()
return err
}
}
func (a *CommentAggregate) onCommentCreate(evt eventstore.Event) error {
var eventData CommentCreateEvent
if err := evt.GetJsonData(&eventData); err != nil {
return errors.Wrap(err, "GetJsonData")
}
a.Comment.ID = a.ID
a.Comment.Tenant = a.Tenant
a.Comment.Content = eventData.Content
a.Comment.ContentType = eventData.ContentType
a.Comment.AuthorUserId = eventData.AuthorUserId
a.Comment.CommentedIssueId = eventData.CommentedIssueId
a.Comment.Source = commonmodel.Source{
Source: eventData.Source,
SourceOfTruth: eventData.Source,
AppSource: eventData.AppSource,
}
a.Comment.CreatedAt = eventData.CreatedAt
a.Comment.UpdatedAt = eventData.UpdatedAt
if eventData.ExternalSystem.Available() {
a.Comment.ExternalSystems = []commonmodel.ExternalSystem{eventData.ExternalSystem}
}
return nil
}
func (a *CommentAggregate) onCommentUpdate(evt eventstore.Event) error {
var eventData CommentUpdateEvent
if err := evt.GetJsonData(&eventData); err != nil {
return errors.Wrap(err, "GetJsonData")
}
if eventData.Source == constants.SourceOpenline {
a.Comment.Source.SourceOfTruth = eventData.Source
}
if eventData.Source != a.Comment.Source.SourceOfTruth && a.Comment.Source.SourceOfTruth == constants.SourceOpenline {
if a.Comment.Content == "" {
a.Comment.Content = eventData.Content
}
if a.Comment.ContentType == "" {
a.Comment.ContentType = eventData.ContentType
}
} else {
a.Comment.Content = eventData.Content
a.Comment.ContentType = eventData.ContentType
}
a.Comment.UpdatedAt = eventData.UpdatedAt
return nil
}