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
14 changes: 13 additions & 1 deletion stovepipe/core/requestlog/materializer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,18 @@ func TestMaterializerPersistRequestStateLog(t *testing.T) {
request.BaseURI = testBaseURI
}
if !tt.wantErr {
expectRequestSummaryUpdate(summaries)
summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(seededRequestSummary(), nil)
summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).DoAndReturn(
func(_ context.Context, summary entity.RequestSummary, _, _ int32) error {
assert.Equal(t, tt.outcomeReason, summary.OutcomeReason)
if tt.state == entity.RequestStateAccepted {
assert.Equal(t, testNowMs, summary.AcceptedAtMs)
} else {
assert.Zero(t, summary.AcceptedAtMs)
}
return nil
},
)
store.EXPECT().Create(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, entry entity.RequestLog) error {
assert.NotEmpty(t, entry.ID)
assert.Equal(t, testNowMs, entry.TimestampMs)
Expand Down Expand Up @@ -157,6 +168,7 @@ func TestMaterializerExistingIdenticalOccurrenceIsSuccess(t *testing.T) {
current.State = entity.RequestStateAccepted
current.RequestVersion = 1
current.StateTimestampMs = testNowMs - 1000
current.AcceptedAtMs = testNowMs - 1000
current.Version = 2
summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil)

Expand Down
81 changes: 58 additions & 23 deletions stovepipe/core/requestlog/request_summary.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,40 +81,71 @@ func updateExistingRequestSummary(
current entity.RequestSummary,
log entity.RequestLog,
) error {
if current.RequestID != log.RequestID || current.Queue != log.Queue {
return fmt.Errorf(
"request summary identity conflicts with state log log_request_id=%q summary_request_id=%q log_queue=%q summary_queue=%q",
log.RequestID, current.RequestID, log.Queue, current.Queue,
)
updated, err := projectRequestSummary(current, log)
if err != nil {
return err
}
if current.RequestVersion > log.RequestVersion {
if updated == current {
return nil
}
if current.RequestVersion == log.RequestVersion {
if requestSummaryMatchesLog(current, log) {
return nil
}
return fmt.Errorf("request summary conflicts with state log request_id=%q version=%d", log.RequestID, log.RequestVersion)
}

updated := current
updated.State = log.State
updated.RequestVersion = log.RequestVersion
updated.StateTimestampMs = log.TimestampMs
if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok {
updated.BaseURI = baseURI
}
oldVersion := current.Version
newVersion := oldVersion + 1
updated.Version = oldVersion
if err := summaryStore.Update(ctx, updated, oldVersion, newVersion); err != nil {
return fmt.Errorf("failed to update request summary request_id=%q: %w", log.RequestID, err)
}
return nil
}

func requestSummaryMatchesLog(summary entity.RequestSummary, log entity.RequestLog) bool {
if summary.State != log.State || summary.StateTimestampMs != log.TimestampMs {
func projectRequestSummary(current entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) {
updated, err := projectRequestSummaryFacts(current, log)
if err != nil {
return entity.RequestSummary{}, err
}
return projectRequestSummaryState(updated, log)
}

// Facts may be filled from any retained version; known values must agree.
func projectRequestSummaryFacts(summary entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) {
if summary.RequestID != log.RequestID || summary.Queue != log.Queue {
return entity.RequestSummary{}, fmt.Errorf(
"request summary identity conflicts with state log log_request_id=%q summary_request_id=%q log_queue=%q summary_queue=%q",
log.RequestID, summary.RequestID, log.Queue, summary.Queue,
)
}
if log.State != entity.RequestStateAccepted {
return summary, nil
}
if summary.AcceptedAtMs != 0 && summary.AcceptedAtMs != log.TimestampMs {
return entity.RequestSummary{}, fmt.Errorf("request summary acceptance time conflicts with retained log request_id=%q", log.RequestID)
}
summary.AcceptedAtMs = log.TimestampMs
return summary, nil
}

func projectRequestSummaryState(summary entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) {
if log.RequestVersion < summary.RequestVersion {
return summary, nil
}
if log.RequestVersion == summary.RequestVersion {
if !requestSummaryStateMatchesLog(summary, log) {
return entity.RequestSummary{}, fmt.Errorf("request summary conflicts with state log request_id=%q version=%d", log.RequestID, log.RequestVersion)
}
summary.OutcomeReason = log.OutcomeReason
return summary, nil
}
summary.State = log.State
summary.RequestVersion = log.RequestVersion
summary.StateTimestampMs = log.TimestampMs
summary.OutcomeReason = log.OutcomeReason
if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok {
summary.BaseURI = baseURI
}
return summary, nil
}

func requestSummaryStateMatchesLog(summary entity.RequestSummary, log entity.RequestLog) bool {
if summary.State != log.State || summary.StateTimestampMs != log.TimestampMs ||
(summary.OutcomeReason != entity.RequestOutcomeReasonUnknown && summary.OutcomeReason != log.OutcomeReason) {
return false
}
baseURI, updatesBaseURI := log.Metadata[MetadataKeyBaseURI]
Expand All @@ -129,6 +160,10 @@ func requestSummaryFromRequestAndLog(request entity.Request, log entity.RequestL
State: log.State,
RequestVersion: log.RequestVersion,
StateTimestampMs: log.TimestampMs,
OutcomeReason: log.OutcomeReason,
}
if log.State == entity.RequestStateAccepted {
summary.AcceptedAtMs = log.TimestampMs
}
if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok {
summary.BaseURI = baseURI
Expand Down
Loading
Loading