diff --git a/internal/dispatcher/runtime_integration_test.go b/internal/dispatcher/runtime_integration_test.go index 4cb7d0d..d51ba6c 100644 --- a/internal/dispatcher/runtime_integration_test.go +++ b/internal/dispatcher/runtime_integration_test.go @@ -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 {