Purge stopped task backlog after stopping its consumer

This commit is contained in:
2026-10-01 18:08:09 +08:00
parent ac8a2d0a1e
commit fabfbf71c6
9 changed files with 209 additions and 56 deletions
+14
View File
@@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"fmt"
"log/slog"
"time"
"git.ipao.vip/rogee/go-sip/internal/configread"
@@ -35,6 +36,7 @@ type ControlController struct {
ApprovedSIP *configread.SIP
ApprovedProviders *map[string]configread.Provider
Now func() time.Time
PurgeTaskQueue func(context.Context, string) (int, error)
}
// ProcessControl applies each delivered control independently: it has no
@@ -125,6 +127,18 @@ func (c *ControlController) ProcessControl(ctx context.Context, body []byte) err
}
return fmt.Errorf("prepare task control %q: %w", event.EventID, err)
}
if event.Payload.Action == "stop" {
if c.PurgeTaskQueue == nil {
return errors.New("stop requires a task queue purger")
}
count, err := c.PurgeTaskQueue(ctx, event.Payload.TaskID)
if err != nil {
return fmt.Errorf("purge stopped task %q backlog: %w", event.Payload.TaskID, err)
}
// A purge covers Ready messages only. The durable stop barrier also
// rejects any later or already accepted deliveries.
slog.Info("stopped task backlog purged", "dispatcher_id", event.DispatcherID, "task_id", event.Payload.TaskID, "ready_messages", count)
}
spec := ControlSpec{
DispatcherID: event.DispatcherID, TenantID: event.TenantID,
TaskID: event.Payload.TaskID, Action: event.Payload.Action,
+48 -1
View File
@@ -82,7 +82,7 @@ func newControlFixture(t *testing.T) (*ControlController, *fakeControlAgent, *st
t.Fatal(err)
}
agent := &fakeControlAgent{}
controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, ApprovedSIP: &snapshot.SIP, ApprovedProviders: &providers, Now: func() time.Time { return monday(9, 30) }}
controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, ApprovedSIP: &snapshot.SIP, ApprovedProviders: &providers, Now: func() time.Time { return monday(9, 30) }, PurgeTaskQueue: func(context.Context, string) (int, error) { return 0, nil }}
return controller, agent, s
}
@@ -162,6 +162,53 @@ func TestControlStopCannotResumeOrDispatchPending(t *testing.T) {
}
}
func TestControlStopPurgesAfterBarrierBeforeAgentAndAck(t *testing.T) {
controller, agent, s := newControlFixture(t)
controller.PurgeTaskQueue = func(_ context.Context, taskID string) (int, error) {
if taskID != "task-asr" {
t.Fatalf("purged wrong task: %s", taskID)
}
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, taskID); err != nil || admitted {
t.Fatalf("stop barrier missing at purge: %v %v", admitted, err)
}
if len(agent.calls) != 0 {
t.Fatal("Agent called before purge")
}
if pending, err := s.ListPendingOutbox(controller.DispatcherID); err != nil || len(pending) != 0 {
t.Fatalf("ack before purge: %v %v", pending, err)
}
return 42, nil
}
if err := controller.ProcessControl(context.Background(), controlBody(t, "stop-purge", "stop", "")); err != nil {
t.Fatal(err)
}
if len(agent.calls) != 1 {
t.Fatalf("Agent not called after purge: %v", agent.calls)
}
}
func TestControlStopPurgeFailureDoesNotAckOrCallAgent(t *testing.T) {
controller, agent, s := newControlFixture(t)
controller.PurgeTaskQueue = func(context.Context, string) (int, error) { return 0, errors.New("rabbitmq purge failed") }
body := controlBody(t, "stop-fail", "stop", "")
if err := controller.ProcessControl(context.Background(), body); err == nil || !strings.Contains(err.Error(), "rabbitmq purge failed") {
t.Fatalf("purge failure hidden: %v", err)
}
if len(agent.calls) != 0 {
t.Fatalf("Agent called after purge failure: %v", agent.calls)
}
if pending, err := s.ListPendingOutbox(controller.DispatcherID); err != nil || len(pending) != 0 {
t.Fatalf("stop falsely acknowledged: %v %v", pending, err)
}
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted {
t.Fatalf("failed purge reopened admission: %v %v", admitted, err)
}
controller.PurgeTaskQueue = func(context.Context, string) (int, error) { return 0, nil }
if err := controller.ProcessControl(context.Background(), body); err != nil {
t.Fatalf("redelivered stop did not recover: %v", err)
}
}
func TestStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox(t *testing.T) {
control, _, s := newControlFixture(t)
if err := control.ProcessControl(context.Background(), controlBody(t, "stop-before-execute", "stop", "")); err != nil {
+59 -28
View File
@@ -16,7 +16,8 @@ import (
)
// Runtime owns task/control consumers for one Dispatcher. SaaS creates
// every queue/binding; this process only checks and consumes predeclared ones.
// every queue/binding; this process checks and consumes predeclared ones and
// purges a stopped task's Ready backlog without changing queue topology.
type Runtime struct {
Broker *mq.Broker
Bootstrap Bootstrap
@@ -25,15 +26,17 @@ type Runtime struct {
PollInterval time.Duration
Logger *slog.Logger
gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction
locksMu sync.Mutex
taskLocks map[string]*sync.Mutex
failures chan error
sipUpdates chan struct{}
gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction
consumersMu sync.Mutex // stop consumer and purge cannot race consumer sync
taskConsumers map[string]*mq.Consumer
locksMu sync.Mutex
taskLocks map[string]*sync.Mutex
failures chan error
sipUpdates chan struct{}
}
// Serve closes admission on every shutdown/failure; it never clears durable
// calls, results, task queues, or the SQLite file.
// calls, results, or the SQLite file. Stop purges only the task's Ready backlog.
func (r *Runtime) Serve(ctx context.Context) (result error) {
if r == nil || r.Broker == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Bootstrap.DispatcherID == "" || r.Execute.Store != r.Bootstrap.Store || r.Control.Store != r.Bootstrap.Store || r.Execute.DispatcherID != r.Bootstrap.DispatcherID || r.Control.DispatcherID != r.Bootstrap.DispatcherID || r.PollInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil {
return errors.New("runtime requires one Dispatcher, durable state, verified SIP and configured pending-work polling; external control drain is forbidden")
@@ -48,20 +51,24 @@ func (r *Runtime) Serve(ctx context.Context) (result error) {
r.sipUpdates = make(chan struct{}, 1)
r.taskLocks = make(map[string]*sync.Mutex)
consumers := make(map[string]*mq.Consumer)
r.taskConsumers = consumers
r.Control.PurgeTaskQueue = r.Broker.PurgeTaskQueue
var controlConsumer *mq.Consumer
defer func() {
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
for queue, consumer := range consumers {
if err := consumer.Stop(stopCtx); err != nil {
result = errors.Join(result, fmt.Errorf("stop task consumer %q: %w", queue, err))
}
}
if controlConsumer != nil {
if err := controlConsumer.Stop(stopCtx); err != nil {
result = errors.Join(result, fmt.Errorf("stop control consumer: %w", err))
}
}
r.consumersMu.Lock()
for queue, consumer := range consumers {
if err := consumer.Stop(stopCtx); err != nil {
result = errors.Join(result, fmt.Errorf("stop task consumer %q: %w", queue, err))
}
}
r.consumersMu.Unlock()
if err := r.Bootstrap.Store.CloseAdmission(r.Bootstrap.DispatcherID); err != nil {
result = errors.Join(result, fmt.Errorf("close Dispatcher admission: %w", err))
}
@@ -173,6 +180,7 @@ func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) erro
EventType string `json:"event_type"`
Payload struct {
TaskID string `json:"task_id"`
Action string `json:"action"`
Revision int64 `json:"revision"`
} `json:"payload"`
}
@@ -183,6 +191,29 @@ func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) erro
}
switch event.EventType {
case "task.control":
if event.Payload.Action == "stop" {
// Stop the consumer before purging: closing its channel requeues any
// unacknowledged deliveries so the purge covers them as Ready.
// Hold this lock through the control transition to prevent sync from
// restarting the consumer before the stop barrier is persisted.
r.consumersMu.Lock()
defer r.consumersMu.Unlock()
route, err := tenant.TaskRoute(r.Bootstrap.DispatcherID, event.Payload.TaskID)
if err != nil {
r.signalFailure(err)
return err
}
if consumer := r.taskConsumers[route.Queue]; consumer != nil {
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
err := consumer.Stop(stopCtx)
cancel()
if err != nil {
r.signalFailure(err)
return fmt.Errorf("stop task consumer before purge %q: %w", route.Queue, err)
}
delete(r.taskConsumers, route.Queue)
}
}
return r.withTask(event.Payload.TaskID, func() error {
if err := r.Control.ProcessControl(ctx, body); err != nil {
r.Logger.Error("task control failed", "dispatcher_id", r.Bootstrap.DispatcherID, "task_id", event.Payload.TaskID, "error", err)
@@ -241,6 +272,8 @@ func (r *Runtime) processPending(ctx context.Context) error {
}
func (r *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error {
r.consumersMu.Lock()
defer r.consumersMu.Unlock()
assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID)
if err != nil {
return err
@@ -251,24 +284,22 @@ func (r *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*m
if err != nil {
return err
}
if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" {
if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" || task.ControlState == "stopping" || task.ControlState == "stopped" || task.Status == "stopped" {
continue
}
if task.ControlState != "stopped" && task.ControlState != "stopping" && task.Status != "stopped" {
admitted, err := r.Bootstrap.Store.CanAdmit(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
if err != nil {
return err
}
if !admitted {
continue
}
waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
if err != nil {
return err
}
if waiting != 0 {
continue
}
admitted, err := r.Bootstrap.Store.CanAdmit(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
if err != nil {
return err
}
if !admitted {
continue
}
waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID)
if err != nil {
return err
}
if waiting != 0 {
continue
}
wanted[route.Queue] = task
}
@@ -366,6 +366,12 @@ 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"} {
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)
}
publish(controlRoute, controlBody(t, "stop-after-sip", "stop", ""))
select {
case spec := <-agent.controls:
@@ -408,6 +414,10 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) {
if ack.Payload.Status != "applied" {
t.Fatalf("stop control falsely acknowledged during SIP reload: %+v", ack)
}
state, err := admin.QueueInspect(taskRoute.Queue)
if err != nil || state.Messages != 0 || state.Consumers != 0 {
t.Fatalf("stopped task backlog was not purged: %+v %v", state, err)
}
break
}
}
+18 -23
View File
@@ -275,8 +275,12 @@ func (b *Broker) StartPredeclaredConsumer(ctx context.Context, queue string, han
return consumer, nil
}
func (b *Broker) DrainPredeclared(ctx context.Context, queue string) (int, error) {
if err := b.validateQueue(queue, false); err != nil {
// PurgeTaskQueue removes Ready messages from only this Dispatcher's task queue.
// The caller must first stop its consumer so unacknowledged deliveries are
// requeued before the purge; SaaS still owns queue creation and bindings.
func (b *Broker) PurgeTaskQueue(ctx context.Context, taskID string) (int, error) {
route, err := tenant.TaskRoute(b.dispatcherID, taskID)
if err != nil {
return 0, err
}
if err := ctx.Err(); err != nil {
@@ -291,34 +295,25 @@ func (b *Broker) DrainPredeclared(ctx context.Context, queue string) (int, error
}
channel, err := conn.Channel()
if err != nil {
return 0, fmt.Errorf("open current drain channel: %w", err)
return 0, fmt.Errorf("open current purge channel: %w", err)
}
defer channel.Close()
if _, err := channel.QueueDeclarePassive(queue, true, false, false, false, nil); err != nil {
return 0, fmt.Errorf("required SaaS-owned queue %s unavailable: %w", queue, err)
if _, err := channel.QueueDeclarePassive(route.Queue, true, false, false, false, nil); err != nil {
return 0, fmt.Errorf("required SaaS-owned queue %s unavailable: %w", route.Queue, err)
}
drained := 0
for {
if err := ctx.Err(); err != nil {
return drained, err
}
delivery, ok, err := channel.Get(queue, false)
if err != nil {
return drained, fmt.Errorf("read SaaS-owned queue %s for stop drain: %w", queue, err)
}
if !ok {
return drained, nil
}
if err := delivery.Ack(false); err != nil {
return drained, fmt.Errorf("ack stopped task backlog: %w", err)
}
drained++
if err := ctx.Err(); err != nil {
return 0, err
}
count, err := channel.QueuePurge(route.Queue, false)
if err != nil {
return 0, fmt.Errorf("purge SaaS-owned task queue %s: %w", route.Queue, err)
}
return count, nil
}
// DrainControlPredeclared processes the already queued controls before task admission.
// Unlike stopped task backlogs, control deliveries must pass through the handler
// before ACK; a transient failure is requeued and closes startup admission.
// Control deliveries must pass through the handler before ACK; a transient
// failure is requeued and closes startup admission.
func (b *Broker) DrainControlPredeclared(ctx context.Context, queue string, handler MessageHandler) (int, error) {
if queue == "" || queue != b.controlQueue || handler == nil {
return 0, errors.New("configured SaaS-owned control queue and handler are required")
+56
View File
@@ -149,6 +149,62 @@ func TestBrokerSharedResultQueueAndNoConfigure(t *testing.T) {
t.Fatalf("D2 stole D1 command: %+v %v", state, err)
}
// Purge is permitted with read access and affects only the chosen task;
// it must not remove another Dispatcher's backlog or the control queue.
for i := 0; i < 3; i++ {
if err := admin.PublishWithContext(context.Background(), CommandsExchange, task.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil {
t.Fatal(err)
}
}
if err := admin.PublishWithContext(context.Background(), CommandsExchange, otherTask.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil {
t.Fatal(err)
}
count, err := brokers[0].PurgeTaskQueue(context.Background(), "task-asr")
if err != nil || count != 3 {
t.Fatalf("task-only purge: removed=%d err=%v", count, err)
}
if state, err := admin.QueueInspect(task.Queue); err != nil || state.Messages != 0 {
t.Fatalf("D1 queue not empty: %+v %v", state, err)
}
if state, err := admin.QueueInspect(otherTask.Queue); err != nil || state.Messages != 1 {
t.Fatalf("D2 queue was purged: %+v %v", state, err)
}
if _, err := brokers[0].PurgeTaskQueue(context.Background(), "../foreign"); err == nil {
t.Fatal("invalid task ID accepted for purge")
}
// Stop closes the consumer channel before purge; an in-flight delivery
// must become Ready rather than surviving the purge as Unacked.
inFlight := make(chan struct{}, 1)
consumer, err := brokers[0].StartPredeclaredConsumer(context.Background(), task.Queue, func(ctx context.Context, _ string, _ []byte) error {
select {
case inFlight <- struct{}{}:
default:
}
<-ctx.Done()
return ctx.Err()
})
if err != nil {
t.Fatal(err)
}
if err := admin.PublishWithContext(context.Background(), CommandsExchange, task.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: currentMessage(t, "mq-execute")}); err != nil {
t.Fatal(err)
}
select {
case <-inFlight:
case <-time.After(3 * time.Second):
t.Fatal("task consumer did not receive in-flight delivery")
}
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := consumer.Stop(stopCtx); err != nil {
t.Fatal(err)
}
count, err = brokers[0].PurgeTaskQueue(context.Background(), "task-asr")
if err != nil || count != 1 {
t.Fatalf("in-flight delivery survived stop and purge: removed=%d err=%v", count, err)
}
resultBody := currentMessage(t, "mq-result-no-recording")
for i, broker := range brokers {
message := resultBody