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 != "" }