Make isolated MQ backlog fixtures use confirmed delivery and consumer quiescence
This commit is contained in:
@@ -86,11 +86,32 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
|
||||
}
|
||||
defer func(q, key, exchange string) { _ = admin.QueueUnbind(q, key, exchange, nil) }(queue, route.BindingKey, route.Exchange)
|
||||
}
|
||||
// Socket write completion is not broker acceptance. Establish the same
|
||||
// mandatory/confirm boundary required of the simulated SaaS publisher.
|
||||
if err := admin.Confirm(false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
returned := admin.NotifyReturn(make(chan amqp.Return, 1))
|
||||
publish := func(route tenant.Route, body []byte) {
|
||||
t.Helper()
|
||||
if err := admin.PublishWithContext(context.Background(), route.Exchange, route.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
confirmation, err := admin.PublishWithDeferredConfirmWithContext(ctx, route.Exchange, route.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if confirmation == nil {
|
||||
t.Fatal("isolated SaaS publisher did not enter confirm mode")
|
||||
}
|
||||
accepted, err := confirmation.WaitContext(ctx)
|
||||
if err != nil || !accepted {
|
||||
t.Fatalf("isolated SaaS message not confirmed: accepted=%v error=%v", accepted, err)
|
||||
}
|
||||
select {
|
||||
case msg := <-returned:
|
||||
t.Fatalf("isolated SaaS message returned: code=%d", msg.ReplyCode)
|
||||
default:
|
||||
}
|
||||
}
|
||||
publish(controlRoute, controlBody(t, "control-example", "pause", "drain"))
|
||||
publish(taskRoute, executeBody(t, "call-example", "15803300952"))
|
||||
@@ -366,11 +387,35 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
|
||||
t.Fatalf("pending SIP change originated call %s", spec.EventID)
|
||||
case <-time.After(120 * time.Millisecond):
|
||||
}
|
||||
for _, eventID := range []string{"stop-backlog-1", "stop-backlog-2"} {
|
||||
// The preceding instruction probes the durable admission fence and may
|
||||
// legitimately have reached the inbox before asynchronous consumer stop.
|
||||
// Do not count it as ready-queue backlog. Establish consumer quiescence
|
||||
// before publishing an independent, confirmed three-message purge fixture.
|
||||
deadline = time.After(5 * time.Second)
|
||||
for {
|
||||
state, err := admin.QueueInspect(taskRoute.Queue)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if state.Consumers == 0 {
|
||||
break
|
||||
}
|
||||
select {
|
||||
case <-deadline:
|
||||
t.Fatalf("consumer did not stop behind SIP admission fence: %+v", state)
|
||||
case <-time.After(20 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
for _, eventID := range []string{"stop-backlog-0", "stop-backlog-1", "stop-backlog-2"} {
|
||||
publish(taskRoute, executeBody(t, eventID, "15003164745"))
|
||||
}
|
||||
if state, err := admin.QueueInspect(taskRoute.Queue); err != nil || state.Messages < 3 {
|
||||
t.Fatalf("expected isolated stop backlog: %+v %v", state, err)
|
||||
pendingCommands, pendingErr := db.ListPendingExecute(id)
|
||||
var pendingIDs []string
|
||||
for _, command := range pendingCommands {
|
||||
pendingIDs = append(pendingIDs, command.EventID)
|
||||
}
|
||||
t.Fatalf("expected isolated stop backlog: queue=%+v inspect_error=%v pending_event_ids=%v pending_error=%v", state, err, pendingIDs, pendingErr)
|
||||
}
|
||||
publish(controlRoute, controlBody(t, "stop-after-sip", "stop", ""))
|
||||
select {
|
||||
|
||||
Reference in New Issue
Block a user