package application import ( "context" "errors" "reflect" "testing" "time" ) func TestEventRelayUsesBoundedFIFOBackpressure(t *testing.T) { relay, err := NewEventRelay(1) if err != nil { t.Fatal(err) } first := Event{Type: EventCatalogRefreshed, RequestID: "first"} second := Event{Type: EventCatalogRejected, RequestID: "second"} if err := relay.Submit(context.Background(), first); err != nil { t.Fatal(err) } started := make(chan struct{}) secondResult := make(chan error, 1) go func() { close(started) secondResult <- relay.Submit(context.Background(), second) }() <-started select { case err := <-secondResult: t.Fatalf("second Submit() completed while relay was full: %v", err) default: } var received []string if err := relay.Drain(func(event Event) error { received = append(received, event.RequestID) return nil }); err != nil { t.Fatal(err) } if err := waitRelayResult(secondResult); err != nil { t.Fatalf("second Submit() error = %v", err) } if err := relay.Drain(func(event Event) error { received = append(received, event.RequestID) return nil }); err != nil { t.Fatal(err) } if !reflect.DeepEqual(received, []string{"first", "second"}) { t.Fatalf("received order = %v", received) } } func TestEventRelayCloseAndContextUnblockFullSubmit(t *testing.T) { tests := []struct { name string unblock func(*EventRelay, context.CancelFunc) wantErr error }{ { name: "close", unblock: func(relay *EventRelay, _ context.CancelFunc) { relay.Close() }, wantErr: ErrEventRelayClosed, }, { name: "cancel", unblock: func(_ *EventRelay, cancel context.CancelFunc) { cancel() }, wantErr: context.Canceled, }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { relay, err := NewEventRelay(1) if err != nil { t.Fatal(err) } if err := relay.Submit( context.Background(), Event{Type: EventCatalogRefreshed}, ); err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() started := make(chan struct{}) result := make(chan error, 1) go func() { close(started) result <- relay.Submit(ctx, Event{Type: EventCatalogRejected}) }() <-started test.unblock(relay, cancel) if err := waitRelayResult(result); !errors.Is(err, test.wantErr) { t.Fatalf("Submit() error = %v, want %v", err, test.wantErr) } }) } } func TestPumpEventsInvalidatesBeforeUIDrainAndStops(t *testing.T) { relay, err := NewEventRelay(1) if err != nil { t.Fatal(err) } source := make(chan Event, 2) invalidated := make(chan struct{}, 2) ctx, cancel := context.WithCancel(context.Background()) pumpResult := make(chan error, 1) go func() { pumpResult <- PumpEvents(ctx, source, relay, func() { invalidated <- struct{}{} }) }() first := Event{Type: EventCatalogRefreshed, RequestID: "first"} second := Event{Type: EventCatalogRejected, RequestID: "second"} source <- first source <- second waitSignal(t, invalidated) var received []string if len(received) != 0 { t.Fatal("background pump applied an event before UI drain") } if err := relay.Drain(func(event Event) error { received = append(received, event.RequestID) return nil }); err != nil { t.Fatal(err) } waitSignal(t, invalidated) if err := relay.Drain(func(event Event) error { received = append(received, event.RequestID) return nil }); err != nil { t.Fatal(err) } if !reflect.DeepEqual(received, []string{"first", "second"}) { t.Fatalf("received order = %v", received) } cancel() if err := waitRelayResult(pumpResult); !errors.Is(err, context.Canceled) { t.Fatalf("PumpEvents() error = %v", err) } } func TestEventRelayDrainReportsErrorsAndContinues(t *testing.T) { relay, err := NewEventRelay(2) if err != nil { t.Fatal(err) } for _, eventType := range []EventType{EventCatalogRefreshed, EventCatalogRejected} { if err := relay.Submit(context.Background(), Event{Type: eventType}); err != nil { t.Fatal(err) } } wantErr := errors.New("apply failed") applied := 0 err = relay.Drain(func(Event) error { applied++ return wantErr }) if !errors.Is(err, wantErr) || applied != 2 { t.Fatalf("Drain() = (%d applies, %v)", applied, err) } } func TestNewEventRelayRejectsInvalidCapacity(t *testing.T) { if _, err := NewEventRelay(0); !errors.Is(err, ErrEventRelayInvalid) { t.Fatalf("NewEventRelay(0) error = %v", err) } relay, err := NewEventRelay(1) if err != nil { t.Fatal(err) } relay.Close() if err := relay.Submit( context.Background(), Event{Type: EventCatalogRefreshed}, ); !errors.Is(err, ErrEventRelayClosed) { t.Fatalf("Submit(after Close) error = %v", err) } } func waitRelayResult(result <-chan error) error { select { case err := <-result: return err case <-time.After(2 * time.Second): return errors.New("timed out waiting for relay") } } func waitSignal(t *testing.T, signal <-chan struct{}) { t.Helper() select { case <-signal: case <-time.After(2 * time.Second): t.Fatal("timed out waiting for signal") } }