Files
go-sip/internal/dispatcher/task_runtime_v3.go
T

324 lines
14 KiB
Go

package dispatcher
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"git.ipao.vip/rogee/go-sip/internal/configread"
"git.ipao.vip/rogee/go-sip/internal/mq"
"git.ipao.vip/rogee/go-sip/internal/store"
)
const (
taskDiscoveryPollInterval = 30 * time.Second
outboxFlushInterval = 250 * time.Millisecond
outboxFlushBatchSize = 64
)
// LocalV01Runtime owns the control and task consumers plus full startup discovery
// and in-memory live changes. Broker queues remain SaaS-owned; it only consumes
// them and publishes Dispatcher outbox records.
type LocalV01Runtime struct {
dispatcher *Dispatcher
broker *mq.V3Broker
client *configread.Client
verifier SIPConfigVerifier
queues *v3TaskQueueController
controls *localTaskControlProcessor
cursor string // live event ID; never persisted across a process restart
complete bool // a full snapshot of this process's assignments has been committed
}
func NewLocalV01Runtime(d *Dispatcher, broker *mq.V3Broker, client *configread.Client, verifier SIPConfigVerifier, taskController TaskController) (*LocalV01Runtime, error) {
if d == nil || broker == nil || client == nil || verifier == nil {
return nil, errors.New("dispatcher, V3 broker, config-read client, and SIP verifier are required")
}
if d.dispatcherID == "" || client.DispatcherID() != d.dispatcherID || broker.ControlQueue() != mq.V3ControlQueueName(d.dispatcherID) {
return nil, errors.New("V3 runtime Dispatcher identity mismatch")
}
queues := newV3TaskQueueController(d, broker, taskController)
return &LocalV01Runtime{
dispatcher: d, broker: broker, client: client, verifier: verifier, queues: queues,
controls: newLocalTaskControlProcessor(d, client, queues),
}, nil
}
// EnableMockAuthorizedOrigination binds the single active Agent before task
// consumption starts. Other modes have no authorized execution adapter.
func (r *LocalV01Runtime) EnableMockAuthorizedOrigination(agentID string, agents *AgentCoordinator) error {
if r == nil || r.queues == nil {
return errors.New("local task runtime is not configured")
}
return r.queues.EnableMockAuthorizedOrigination(agentID, agents)
}
// Run keeps control processing available when config discovery is temporarily
// unavailable, while withholding task-queue consumption until a fresh discovery
// and valid execution snapshot have been persisted.
func (r *LocalV01Runtime) Run(ctx context.Context) error {
if r == nil || r.dispatcher == nil || r.broker == nil || r.client == nil || r.verifier == nil || r.queues == nil || r.controls == nil {
return errors.New("local v0.1 runtime is not configured")
}
if err := r.dispatcher.store.CloseLocalTaskDiscoveryAdmissionV04(r.dispatcher.dispatcherID); err != nil {
return fmt.Errorf("close admission before initial discovery: %w", err)
}
initialDiscoveryErr := r.refreshDiscoveryState(ctx)
if initialDiscoveryErr == nil {
initialDiscoveryErr = r.catchUpControlBeforeTasks(ctx)
}
if initialDiscoveryErr != nil && !errors.Is(initialDiscoveryErr, context.Canceled) {
slog.Error("initial task discovery/control backlog failed; task admission remains closed", "dispatcher_id", r.dispatcher.dispatcherID, "error", initialDiscoveryErr)
}
controlConsumer, err := r.broker.StartPredeclaredConsumer(ctx, r.broker.ControlQueue(), r.controls.Handle)
if err != nil {
return fmt.Errorf("start SaaS-predeclared control consumer: %w", err)
}
controlDone := make(chan error, 1)
go func() { controlDone <- controlConsumer.Wait(context.Background()) }()
cleanup := func() {
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := r.queues.StopAll(shutdownCtx); err != nil {
slog.Error("stop task consumers during runtime shutdown", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
if err := controlConsumer.Stop(shutdownCtx); err != nil {
slog.Error("stop control consumer during runtime shutdown", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
}
defer cleanup()
// Results and issued-call reconciliation must keep running even when
// discovery is unavailable and new task consumption remains closed.
r.recoverMockResults(ctx)
if initialDiscoveryErr == nil {
if err := r.reconcileTaskQueues(ctx); err != nil {
slog.Error("initial task queue reconciliation failed; admission remains closed", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
}
discoveryTicker := time.NewTicker(taskDiscoveryPollInterval)
defer discoveryTicker.Stop()
resultTicker := time.NewTicker(time.Second)
defer resultTicker.Stop()
outboxTicker := time.NewTicker(outboxFlushInterval)
defer outboxTicker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-controlDone:
if ctx.Err() != nil {
return nil
}
if err == nil {
return errors.New("control queue consumer exited unexpectedly")
}
return fmt.Errorf("control queue consumer stopped: %w", err)
case err := <-r.queues.Errors():
return fmt.Errorf("task queue consumer stopped: %w", err)
case <-discoveryTicker.C:
r.recoverMockResults(ctx)
catchingUp := !r.complete
if catchingUp {
// Never process the same control backlog with two consumers.
if err := controlConsumer.Stop(ctx); err != nil {
return fmt.Errorf("stop control consumer for snapshot catch-up: %w", err)
}
}
discoveryErr := r.refreshDiscoveryState(ctx)
if catchingUp && discoveryErr == nil {
discoveryErr = r.catchUpControlBeforeTasks(ctx)
}
if catchingUp {
controlConsumer, err = r.broker.StartPredeclaredConsumer(ctx, r.broker.ControlQueue(), r.controls.Handle)
if err != nil {
return fmt.Errorf("restart SaaS-predeclared control consumer: %w", err)
}
controlDone = make(chan error, 1)
go func(done chan error, consumer *mq.V3Consumer) { done <- consumer.Wait(context.Background()) }(controlDone, controlConsumer)
}
if discoveryErr != nil {
slog.Error("task discovery/control backlog refresh failed; task admission closed", "dispatcher_id", r.dispatcher.dispatcherID, "error", discoveryErr)
continue
}
if err := r.reconcileTaskQueues(ctx); err != nil {
slog.Error("task queue reconciliation failed; admission remains fail-closed", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
case <-resultTicker.C:
if r.queues.mockAuthorizedAgents != nil {
if err := r.dispatcher.recoverPendingMockResults(ctx); err != nil && ctx.Err() == nil {
slog.Error("restore pending Mock call results", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
}
case <-outboxTicker.C:
if _, err := r.dispatcher.FlushOutbox(ctx, outboxFlushBatchSize); err != nil && ctx.Err() == nil {
slog.Error("Dispatcher outbox flush failed; records remain retryable", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
}
}
}
func (r *LocalV01Runtime) recoverMockResults(ctx context.Context) {
if r.queues.mockAuthorizedAgents == nil {
return
}
if err := r.dispatcher.RecoverAuthorizedMockResults(ctx, r.queues.mockAuthorizedAgentID, r.queues.mockAuthorizedAgents); err != nil && ctx.Err() == nil {
slog.Error("reconcile issued and final Mock call facts; unresolved calls retain quota", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
}
}
func (r *LocalV01Runtime) refreshDiscoveryState(ctx context.Context) error {
if err := r.dispatcher.store.CloseLocalTaskDiscoveryAdmissionV04(r.dispatcher.dispatcherID); err != nil {
return r.failDiscovery(fmt.Errorf("close task admission before discovery: %w", err))
}
pageCtx, cancel := context.WithTimeout(ctx, taskDiscoveryPollInterval)
defer cancel()
if !r.complete {
var snapshotID, watermark, pageToken string
seen, tokens := make(map[string]bool), make(map[string]bool)
var tasks []store.LocalDiscoveredTask
for {
page, err := r.client.ReadTaskSnapshot(pageCtx, snapshotID, pageToken)
if err != nil {
return r.failDiscovery(fmt.Errorf("read complete task snapshot: %w", err))
}
if snapshotID == "" {
snapshotID, watermark = page.SnapshotID, page.Watermark
} else if page.SnapshotID != snapshotID || page.Watermark != watermark {
return r.failDiscovery(errors.New("task snapshot identity or watermark changed between pages"))
}
for _, task := range page.Tasks {
if seen[task.TaskID] {
return r.failDiscovery(fmt.Errorf("task %s appears twice in complete snapshot", task.TaskID))
}
seen[task.TaskID] = true
tasks = append(tasks, discoveredLocalTask(task))
}
if len(tasks) > 256 {
return r.failDiscovery(errors.New("complete task snapshot exceeds per-Dispatcher task limit"))
}
if page.NextPageToken == "" {
break
}
if tokens[page.NextPageToken] {
return r.failDiscovery(errors.New("task snapshot repeats a page token"))
}
tokens[page.NextPageToken] = true
pageToken = page.NextPageToken
}
if err := r.dispatcher.store.ApplyLocalTaskSnapshot(r.dispatcher.dispatcherID, tasks, r.dispatcher.now().UTC()); err != nil {
return r.failDiscovery(fmt.Errorf("persist complete task snapshot: %w", err))
}
slog.Info("complete task snapshot applied", "dispatcher_id", r.dispatcher.dispatcherID, "snapshot_id", snapshotID, "watermark", watermark, "task_count", len(tasks))
r.cursor, r.complete = watermark, true
return nil
}
for {
page, err := r.client.ReadTaskChanges(pageCtx, r.cursor)
if err != nil {
return r.failDiscovery(fmt.Errorf("read live task changes after %s: %w", r.cursor, err))
}
if len(page.Tasks) == 0 {
if err := r.dispatcher.store.MarkLocalTaskDiscoveryReadyV04(r.dispatcher.dispatcherID, r.dispatcher.now().UTC()); err != nil {
return r.failDiscovery(fmt.Errorf("mark live task changes caught up: %w", err))
}
return nil
}
tasks := make([]store.LocalDiscoveredTask, 0, len(page.Tasks))
for _, task := range page.Tasks {
tasks = append(tasks, discoveredLocalTask(task))
}
if err := r.dispatcher.store.ApplyLocalTaskChanges(r.dispatcher.dispatcherID, tasks, r.dispatcher.now().UTC()); err != nil {
return r.failDiscovery(fmt.Errorf("persist live task changes after %s: %w", r.cursor, err))
}
r.cursor = page.NextCursor
}
}
func discoveredLocalTask(task configread.DiscoveredTask) store.LocalDiscoveredTask {
return store.LocalDiscoveredTask{TaskID: task.TaskID, TenantID: task.TenantID,
TenantKey: task.TenantKey, TaskRevision: task.TaskRevision, Status: task.Status}
}
// catchUpControlBeforeTasks processes queued controls while no task may be
// admitted. It must run without a concurrent control consumer.
func (r *LocalV01Runtime) catchUpControlBeforeTasks(ctx context.Context) error {
if err := r.dispatcher.store.CloseLocalTaskDiscoveryAdmissionV04(r.dispatcher.dispatcherID); err != nil {
return r.failDiscovery(fmt.Errorf("close admission before control backlog: %w", err))
}
count, err := r.queues.broker.DrainControlPredeclared(ctx, mq.V3ControlQueueName(r.dispatcher.dispatcherID), r.controls.Handle)
if err != nil {
return r.failDiscovery(fmt.Errorf("process SaaS-owned control backlog: %w", err))
}
if err := r.dispatcher.store.MarkLocalTaskDiscoveryReadyV04(r.dispatcher.dispatcherID, r.dispatcher.now().UTC()); err != nil {
return r.failDiscovery(fmt.Errorf("open admission after control backlog: %w", err))
}
slog.Info("SaaS-owned control backlog processed", "dispatcher_id", r.dispatcher.dispatcherID, "messages", count)
return nil
}
func (r *LocalV01Runtime) failDiscovery(cause error) error {
r.complete = false // failed changes require a new complete snapshot, never a persisted cursor
r.cursor = ""
if err := r.dispatcher.store.CloseLocalTaskDiscoveryAdmissionV04(r.dispatcher.dispatcherID); err != nil {
cause = errors.Join(cause, fmt.Errorf("close admission after task discovery failure: %w", err))
}
if r.queues != nil {
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := r.queues.StopAll(stopCtx); err != nil {
cause = errors.Join(cause, fmt.Errorf("stop task consumers after discovery failure: %w", err))
}
}
return cause
}
func (r *LocalV01Runtime) reconcileTaskQueues(ctx context.Context) error {
assignments, err := r.dispatcher.store.LocalTaskAssignments(r.dispatcher.dispatcherID)
if err != nil {
return fmt.Errorf("list durable task assignments: %w", err)
}
var failures []error
for _, assignment := range assignments {
if !taskAssignmentCanConsume(assignment) {
if err := r.queues.StopTask(ctx, assignment); err != nil {
failures = append(failures, fmt.Errorf("stop task %s consumer: %w", assignment.TaskID, err))
}
continue
}
if _, valid := r.dispatcher.ProjectConfigSnapshot(assignment.TaskID, assignment.TenantID); !valid {
if err := r.dispatcher.LoadProjectConfig(ctx, r.client, r.verifier, assignment.TaskID, assignment.TenantID); err != nil {
if stopErr := r.queues.StopConsumer(ctx, assignment.TaskID); stopErr != nil {
failures = append(failures, fmt.Errorf("stop task %s consumer after config-read failure: %w", assignment.TaskID, stopErr))
}
failures = append(failures, fmt.Errorf("refresh task %s execution snapshot: %w", assignment.TaskID, err))
continue
}
assignment, err = r.dispatcher.store.LocalTaskAssignment(assignment.DispatcherID, assignment.TaskID)
if err != nil {
failures = append(failures, fmt.Errorf("reload task %s assignment after config refresh: %w", assignment.TaskID, err))
continue
}
if !taskAssignmentCanConsume(assignment) {
if err := r.queues.StopTask(ctx, assignment); err != nil {
failures = append(failures, fmt.Errorf("stop task %s consumer after config refresh: %w", assignment.TaskID, err))
}
continue
}
}
if err := r.queues.StartTask(ctx, assignment); err != nil {
failures = append(failures, fmt.Errorf("start predeclared task %s consumer: %w", assignment.TaskID, err))
}
}
return errors.Join(failures...)
}
func taskAssignmentCanConsume(assignment store.LocalTaskAssignment) bool {
return !assignment.Removed && assignment.Status == "running" && assignment.AdmissionState == "running" && assignment.Queue.QueueName != ""
}