Files
go-sip/internal/store/active_task_controls.go
T

56 lines
1.8 KiB
Go

package store
import (
"errors"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
"google.golang.org/protobuf/proto"
)
type ActiveTaskExecutionControl struct {
ExecutionID string
Status string
AgentID string
Binding *agentv1.ExecutionBinding
}
// ActiveTaskExecutionControls returns only reserved or in-flight executions
// that already have a durable Agent assignment. The Store is Dispatcher-local.
func (s *Store) ActiveTaskExecutionControls(tenantID, tenantKey, taskID string) ([]ActiveTaskExecutionControl, error) {
if tenantID == "" || tenantKey == "" || taskID == "" {
return nil, errors.New("tenant and task identity are required")
}
s.mu.Lock()
defer s.mu.Unlock()
rows, err := s.db.Query(`SELECT t.execution_id,t.status,a.agent_id,a.binding
FROM tasks t JOIN execution_agents a ON a.execution_id=t.execution_id
WHERE t.tenant_id=? AND t.tenant_key=? AND t.task_id=?
AND t.status IN ('reserved','running','draining','paused','unknown')
ORDER BY t.execution_id`, tenantID, tenantKey, taskID)
if err != nil {
return nil, err
}
defer rows.Close()
targets := make([]ActiveTaskExecutionControl, 0)
for rows.Next() {
var target ActiveTaskExecutionControl
var encoded []byte
if err := rows.Scan(&target.ExecutionID, &target.Status, &target.AgentID, &encoded); err != nil {
return nil, err
}
binding := &agentv1.ExecutionBinding{}
if err := proto.Unmarshal(encoded, binding); err != nil {
return nil, err
}
if target.ExecutionID != binding.ExecutionId || tenantID != binding.TenantId || tenantKey != binding.TenantKey || taskID != binding.TaskId {
return nil, errors.New("persisted Agent binding does not match task execution")
}
target.Binding = binding
targets = append(targets, target)
}
if err := rows.Err(); err != nil {
return nil, err
}
return targets, nil
}