226 lines
9.3 KiB
Go
226 lines
9.3 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 V3 control and task consumers plus durable discovery
|
|
// polling. Broker queues remain SaaS-owned; this runtime only passively consumes
|
|
// them and publishes Dispatcher outbox records.
|
|
type LocalV01Runtime struct {
|
|
dispatcher *Dispatcher
|
|
broker *mq.V3Broker
|
|
client *configread.Client
|
|
verifier SIPConfigVerifier
|
|
queues *v3TaskQueueController
|
|
controls *localTaskControlProcessor
|
|
}
|
|
|
|
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.CloseLocalTaskDiscoveryAdmission(r.dispatcher.dispatcherID); err != nil {
|
|
return fmt.Errorf("close admission before initial discovery: %w", err)
|
|
}
|
|
initialDiscoveryErr := r.refreshDiscoveryState(ctx)
|
|
if initialDiscoveryErr != nil && !errors.Is(initialDiscoveryErr, context.Canceled) {
|
|
slog.Error("initial task discovery 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)
|
|
if err := r.refreshDiscoveryState(ctx); err != nil {
|
|
slog.Error("task discovery refresh failed; task admission closed", "dispatcher_id", r.dispatcher.dispatcherID, "error", err)
|
|
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 {
|
|
cursor, exists, err := r.dispatcher.store.LocalTaskDiscoveryCursor(r.dispatcher.dispatcherID)
|
|
if err != nil {
|
|
return r.failDiscovery(fmt.Errorf("read durable task-discovery cursor: %w", err))
|
|
}
|
|
if !exists {
|
|
cursor = ""
|
|
}
|
|
discovery, err := r.client.ReadTaskDiscovery(ctx, cursor)
|
|
if err != nil {
|
|
return r.failDiscovery(fmt.Errorf("read task discovery: %w", err))
|
|
}
|
|
if err := r.dispatcher.store.ApplyLocalTaskDiscovery(localTaskDiscoveryFromRead(discovery)); err != nil {
|
|
return r.failDiscovery(fmt.Errorf("persist complete task-discovery response: %w", err))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (r *LocalV01Runtime) failDiscovery(cause error) error {
|
|
if err := r.dispatcher.store.CloseLocalTaskDiscoveryAdmission(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 != ""
|
|
}
|