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
2 changes: 1 addition & 1 deletion aggregator/aggregator.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ func runServer() {
}

func subscribeToAdvisoryUpdates(wg *sync.WaitGroup, readerBuilder mqueue.CreateReader) {
handler := mqueue.MakeRetryingHandler(mqueue.MakeAdvisoryUpdateHandler(advisoryUpdateHandler))
handler := mqueue.MakeRetryingHandler(advisoryUpdateHandler)
for i := 0; i < consumerCount; i++ {
mqueue.SpawnReader(base.Context, wg, advisoryUpdateTopic, readerBuilder, handler)
utils.LogDebug("spawned advisory update reader", i)
Expand Down
12 changes: 11 additions & 1 deletion aggregator/events.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,19 @@ package aggregator

import (
"app/base/mqueue"
"app/base/utils"

"github.com/bytedance/sonic"
)

// TODO: stub - will process advisory update events in batches and update account_advisory table
func advisoryUpdateHandler(event mqueue.AdvisoryUpdateEvent) error {
func advisoryUpdateHandler(m mqueue.KafkaMessage) error {
var event mqueue.AdvisoryUpdateEvent
if err := sonic.Unmarshal(m.Value, &event); err != nil {
utils.LogError("err", err, "Could not deserialize advisory update event")
return nil
}
// TODO: advisory update code goes here
_ = event
return nil
}
15 changes: 0 additions & 15 deletions base/mqueue/advisory_update_event.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,28 +2,13 @@ package mqueue

import (
"app/base/types"
"app/base/utils"
"context"

"github.com/bytedance/sonic"
"github.com/google/uuid"
"github.com/pkg/errors"
)

type AdvisoryUpdateEventHandler func(event AdvisoryUpdateEvent) error

func MakeAdvisoryUpdateHandler(handler AdvisoryUpdateEventHandler) MessageHandler {
return func(m KafkaMessage) error {
var event AdvisoryUpdateEvent
err := sonic.Unmarshal(m.Value, &event)
if err != nil {
utils.LogError("err", err, "Could not deserialize advisory update event")
return nil
}
return handler(event)
}
}

type AdvisoryUpdateEvent struct {
RhAccountID int `json:"rh_account_id"`
WorkspaceID uuid.UUID `json:"workspace_id"`
Expand Down
17 changes: 0 additions & 17 deletions base/mqueue/event.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"context"
"time"

"github.com/bytedance/sonic"
"github.com/lestrrat-go/backoff/v2"
)

Expand All @@ -16,26 +15,10 @@ var policy = backoff.Exponential(
backoff.WithMaxRetries(5),
)

type EventHandler func(message PlatformEvent) error

type MessageData interface {
WriteEvents(ctx context.Context, w Writer) error
}

// Performs parsing of kafka message, and then dispatches this message into provided functions
func MakeMessageHandler(eventHandler EventHandler) MessageHandler {
return func(m KafkaMessage) error {
var event PlatformEvent
err := sonic.Unmarshal(m.Value, &event)
// Not a fatal error, invalid data format, log and skip
if err != nil {
utils.LogError("err", err, "Could not deserialize platform event")
return nil
}
return eventHandler(event)
}
}

func SendMessages(ctx context.Context, w Writer, data MessageData) error {
return data.WriteEvents(ctx, w)
}
30 changes: 8 additions & 22 deletions base/mqueue/mqueue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"sync"
"testing"

"github.com/bytedance/sonic"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
)
Expand All @@ -16,30 +17,15 @@ var someid = uuid.MustParse("99c0ffee-0000-0000-0000-0000000050de")

var msg = KafkaMessage{Value: []byte(`{"id": "` + id.String() + `", "type": "delete"}`)}

func TestParseEvents(t *testing.T) {
reached := false

err := MakeMessageHandler(func(event PlatformEvent) error {
assert.Equal(t, event.ID, id)
assert.Equal(t, *event.Type, "delete")
reached = true
return nil
})(msg)

assert.True(t, reached, "Event handler should have been called")
assert.NoError(t, err)
}

func TestRoundTripKafkaGo(t *testing.T) {
utils.SkipWithoutPlatform(t)
reader := NewKafkaReaderFromEnv("test")
defer reader.Close()

var eventOut PlatformEvent
go reader.HandleMessages(t.Context(), MakeMessageHandler(func(event PlatformEvent) error {
eventOut = event
return nil
}))
go reader.HandleMessages(t.Context(), func(m KafkaMessage) error {
return sonic.Unmarshal(m.Value, &eventOut)
})

writer := NewKafkaWriterFromEnv("test")
eventIn := PlatformEvent{ID: someid}
Expand All @@ -53,14 +39,14 @@ func TestSpawnReader(t *testing.T) {
var nReaders int32
wg := sync.WaitGroup{}
SpawnReader(context.Background(), &wg, "", CreateCountedMockReader(&nReaders),
MakeMessageHandler(func(_ PlatformEvent) error { return nil }))
func(_ KafkaMessage) error { return nil })
wg.Wait()
assert.Equal(t, 1, int(nReaders))
}

func TestRetry(t *testing.T) {
i := 0
handler := func(_ PlatformEvent) error {
handler := func(_ KafkaMessage) error {
i++
if i < 2 {
return errors.New("Failed")
Expand All @@ -69,8 +55,8 @@ func TestRetry(t *testing.T) {
}

// Without retry handler should fail
assert.Error(t, MakeMessageHandler(handler)(msg))
assert.Error(t, handler(msg))

// With retry we handler should eventually succeed
assert.NoError(t, MakeRetryingHandler(MakeMessageHandler(handler))(msg))
assert.NoError(t, MakeRetryingHandler(handler)(msg))
}
4 changes: 3 additions & 1 deletion evaluator/advisory_update_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -155,12 +155,14 @@ func TestAdvisoryUpdateKafkaRoundTrip(t *testing.T) {
database.CheckCachesValid(t)

// Run evaluation
err := evaluateHandler(mqueue.PlatformEvent{
data, err := sonic.Marshal(mqueue.PlatformEvent{
SystemIDs: []uuid.UUID{testInventoryID},
RequestIDs: []string{"request-1"},
OrgID: &orgID,
AccountID: rhAccountID})
assert.NoError(t, err)
err = evaluateHandler(mqueue.KafkaMessage{Value: data})
assert.NoError(t, err)

utils.AssertEqualWait(t, 10, func() (exp, act interface{}) {
return true, len(received.Value) > 0
Expand Down
17 changes: 15 additions & 2 deletions evaluator/evaluate.go
Original file line number Diff line number Diff line change
Expand Up @@ -341,6 +341,13 @@ func getUpdatesData(ctx context.Context, system *models.SystemPlatformV2) (*vmaa
}

merged := utils.MergeVMaaSResponses(yumUpdates, vmaasData)
if system.Inventory.Bootc && merged != nil {
// image-mode systems are not patchable via Insights; force all updates applicable
mergedUpdateList := merged.GetUpdateList()
for nevra := range mergedUpdateList {
(*mergedUpdateList[nevra]).SetUpdatesInstallability(APPLICABLE)
}
}
return merged, nil
}

Expand Down Expand Up @@ -733,7 +740,13 @@ func invalidateCaches(orgID string) error {
return err
}

func evaluateHandler(event mqueue.PlatformEvent) error {
func evaluateHandler(m mqueue.KafkaMessage) error {
var event mqueue.PlatformEvent
if err := sonic.Unmarshal(m.Value, &event); err != nil {
utils.LogError("err", err, "Could not deserialize platform event")
return nil
}

var err error
var wg sync.WaitGroup
guard := make(chan struct{}, nEvalGoroutines)
Expand Down Expand Up @@ -800,7 +813,7 @@ func run(wg *sync.WaitGroup, readerBuilder mqueue.CreateReader) {

loadCache()

var handler = mqueue.MakeRetryingHandler(mqueue.MakeMessageHandler(evaluateHandler))
var handler = mqueue.MakeRetryingHandler(evaluateHandler)
// We create multiple consumers, and hope that the partition rebalancing
// algorithm assigns each consumer a single partition
for i := 0; i < consumerCount; i++ {
Expand Down
Loading
Loading