Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 14 additions & 7 deletions pkg/subscriber/processor/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,18 @@ var (
ErrReasonNil = errors.New("eventpersister: reason is nil")
ErrEvaluationsAreEmpty = errors.New("eventpersister: evaluations are empty")
ErrEvaluationEventIssuedAfterExperimentEnded = errors.New("eventpersister: evaluation event issued after experiment ended") //nolint:lll
ErrFailedToEvaluateUser = errors.New("eventpersister: failed to evaluate user")
ErrAutoOpsRuleNotFound = errors.New("eventpersister: auto ops rule not found")
ErrFeatureEmptyList = errors.New("eventpersister: list feature returned empty")
ErrFeatureVersionNotFound = errors.New("eventpersister: feature version not found")
ErrUnknownEvent = errors.New("metricsevent persister: unknown metrics event")
ErrInvalidDuration = errors.New("metricsevent persister: invalid duration")
ErrUnknownApiId = errors.New("metricsevent persister: unknown api id")
// ErrGoalEventOlderThanEvaluation is returned when the goal event timestamp is older
// than the user's latest evaluation timestamp. Because the client SDK sets both
// timestamps when the events are generated, this means the goal event was created
// before the user was evaluated (incorrect SDK implementation), so the event is
// discarded instead of retried.
ErrGoalEventOlderThanEvaluation = errors.New(
"eventpersister: goal event timestamp is older than the user's latest evaluation timestamp")
ErrFailedToEvaluateUser = errors.New("eventpersister: failed to evaluate user")
ErrAutoOpsRuleNotFound = errors.New("eventpersister: auto ops rule not found")
ErrFeatureEmptyList = errors.New("eventpersister: list feature returned empty")
ErrFeatureVersionNotFound = errors.New("eventpersister: feature version not found")
ErrUnknownEvent = errors.New("metricsevent persister: unknown metrics event")
ErrInvalidDuration = errors.New("metricsevent persister: invalid duration")
ErrUnknownApiId = errors.New("metricsevent persister: unknown api id")
)
56 changes: 38 additions & 18 deletions pkg/subscriber/processor/goal_events_dwh.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,12 @@ func (w *goalEvtWriter) Write(
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeExperimentNotFound).Inc()
continue
}
if errors.Is(err, ErrGoalEventOlderThanEvaluation) {
// The goal event was created before the user was evaluated
// (the client SDK sets the timestamps), so it can never be
// linked. Discard it by acking the message.
continue
}
if !retriable {
w.logger.Error(
"Failed to convert to goal event",
Expand Down Expand Up @@ -264,19 +270,23 @@ func (w *goalEvtWriter) linkGoalEvent(
id, environmentID, tag string,
experiments []*exproto.Experiment,
) ([]*ecdwh.UserEvaluation, bool, error) {
evalExp, retriable, err := w.linkGoalEventByExperiment(ctx, event, id, environmentID, tag, experiments)
evalExp, retriable, err := w.linkGoalEventByExperiment(ctx, event, id, environmentID, tag, experiments, true)
if err != nil {
return nil, retriable, err
}
return evalExp, false, nil
}

// Link one or more experiments by goal ID
// Link one or more experiments by goal ID.
// When storeRetryOnMiss is true (fresh events from Pub/Sub), goal events that can't be
// linked yet are stored in Redis so the retry processor can link them later.
// The retry processor itself passes false because it manages its own retry message.
func (w *goalEvtWriter) linkGoalEventByExperiment(
ctx context.Context,
event *eventproto.GoalEvent,
id, environmentID, tag string,
experiments []*exproto.Experiment,
storeRetryOnMiss bool,
) ([]*ecdwh.UserEvaluation, bool, error) {
// Find the experiment by goal ID
// TODO: we must change the console UI not to allow creating
Expand Down Expand Up @@ -324,19 +334,21 @@ func (w *goalEvtWriter) linkGoalEventByExperiment(
zap.Any("goalEvent", event),
)
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeUserEvaluationNotFound).Inc()
if err := w.storeRetryMessage(&retryMessage{
GoalEvent: event,
EnvironmentID: environmentID,
RetryCount: 0,
ID: id,
}); err != nil {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeFailedToStoreRetryMessage).Inc()
w.logger.Error("Failed to store retry message",
zap.Error(err),
zap.String("environmentId", environmentID),
zap.Any("goalEvent", event),
)
return nil, true, err
if storeRetryOnMiss {
if err := w.storeRetryMessage(&retryMessage{
GoalEvent: event,
EnvironmentID: environmentID,
RetryCount: 0,
ID: id,
}); err != nil {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeFailedToStoreRetryMessage).Inc()
w.logger.Error("Failed to store retry message",
zap.Error(err),
zap.String("environmentId", environmentID),
zap.Any("goalEvent", event),
)
return nil, true, err
}
}
return nil, false, err
}
Expand All @@ -353,11 +365,14 @@ func (w *goalEvtWriter) linkGoalEventByExperiment(
zap.Any("goalEvent", event),
zap.Any("evaluation", eval),
)
// Skip goal events that occurred before the evaluation timestamp.
// This is intentional to trigger retry logic for proper event linking.
// Skip goal events older than the evaluation timestamp.
// Because the client SDK sets the timestamps when the events are generated,
// a goal event older than the user's latest evaluation means the client
// sent the goal event before evaluating the user (incorrect implementation),
// so it must be discarded instead of retried.
if event.Timestamp < eval.Timestamp {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeGoalEventIssuedBeforeEvaluation).Inc()
w.logger.Error("Goal event issued before evaluation",
w.logger.Warn("Goal event is older than the user's latest evaluation, discarding it",
zap.String("environmentId", environmentID),
zap.Any("goalEvent", event),
zap.Any("evaluation", eval),
Expand All @@ -366,6 +381,11 @@ func (w *goalEvtWriter) linkGoalEventByExperiment(
}
evals = append(evals, eval)
}
if len(evals) == 0 {
// Every matching evaluation is newer than the goal event,
// so the goal event can never be linked. The callers must discard it.
return nil, false, ErrGoalEventOlderThanEvaluation
}
return evals, false, nil
}

Expand Down
19 changes: 10 additions & 9 deletions pkg/subscriber/processor/goal_events_dwh_retry.go
Original file line number Diff line number Diff line change
Expand Up @@ -188,19 +188,29 @@ func (w *goalEvtWriter) handleNewRetry(ctx context.Context, msg *retryMessage, k
return
}

// Pass storeRetryOnMiss=false: this path owns the retry message and
// re-stores it below with the incremented retry count, so the linker must
// not overwrite it with a fresh message (which would reset the backoff).
evals, _, err := w.linkGoalEventByExperiment(
ctx,
msg.GoalEvent,
msg.ID,
msg.EnvironmentID,
msg.GoalEvent.Tag,
experiments,
false,
)
if err != nil {
if errors.Is(err, ErrExperimentNotFound) {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeExperimentNotFound).Inc()
lg.Error("Experiment not found, deleting retry message")
w.deleteKey(ctx, key)
} else if errors.Is(err, ErrGoalEventOlderThanEvaluation) {
// The goal event was created before the user was evaluated
// (the client SDK sets the timestamps), so retrying can never
// link it. Discard the retry message.
lg.Warn("Goal event is older than the user's latest evaluation, deleting retry message")
w.deleteKey(ctx, key)
} else {
lg.Error("Linking failed", zap.Error(err))
msg.RetryCount++
Expand All @@ -211,15 +221,6 @@ func (w *goalEvtWriter) handleNewRetry(ctx context.Context, msg *retryMessage, k
}
return
}
if len(evals) == 0 {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeRetryMessageNoEvaluations).Inc()
msg.RetryCount++
if err := w.storeRetryMessage(msg); err != nil {
subscriberHandledCounter.WithLabelValues(subscriberGoalEventDWH, codeFailedToStoreRetryMessage).Inc()
lg.Error("Failed to store retry message", zap.Error(err))
}
return
}

var events []*epproto.GoalEvent
for _, ev := range evals {
Expand Down
Loading
Loading