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

178 lines
6.8 KiB
Go

package store
import (
"database/sql"
"encoding/json"
"errors"
"fmt"
"git.ipao.vip/rogee/go-sip/internal/contract"
"git.ipao.vip/rogee/go-sip/internal/tenant"
)
var ErrControlRejected = errors.New("task control rejected by durable task state")
// PrepareControl closes only this task's admission before requesting the
// Agent action. No successful control acknowledgment exists at this point.
// Repeated commands are prepared and dispatched again; this is not control
// deduplication or an expected-revision/CAS API.
func (s *Store) PrepareControl(dispatcherID string, tenantID int64, taskID, action string) error {
if dispatcherID == "" || tenantID <= 0 || taskID == "" {
return errors.New("invalid task control identity")
}
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
var state, status string
var revision int64
var present int
err = tx.QueryRow(`SELECT control_state,status,task_revision,present FROM dispatcher_tasks
WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, dispatcherID, tenantID, taskID).Scan(&state, &status, &revision, &present)
if err != nil {
return fmt.Errorf("load task control state: %w", err)
}
if present != 1 {
return fmt.Errorf("%w: task is not in the assigned discovery list", ErrControlRejected)
}
var prepared string
switch action {
case "pause":
if state == "stopped" || state == "stopping" {
return fmt.Errorf("%w: stopped task cannot be paused", ErrControlRejected)
}
prepared = "pausing"
case "stop":
prepared = "stopping"
case "resume":
if state != "paused" && state != "resuming" {
return fmt.Errorf("%w: resume requires a persistently paused, non-stopped task", ErrControlRejected)
}
if status != "running" {
return fmt.Errorf("%w: fresh task is not running", ErrControlRejected)
}
var configuredRevision int64
err = tx.QueryRow(`SELECT task_revision FROM dispatcher_configs WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, dispatcherID, tenantID, taskID).Scan(&configuredRevision)
if err != nil || configuredRevision != revision {
return errors.New("resume requires the latest verified task configuration")
}
prepared = "resuming"
default:
return fmt.Errorf("invalid task control action %q", action)
}
if _, err := tx.Exec(`UPDATE dispatcher_tasks SET control_state=? WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, prepared, dispatcherID, tenantID, taskID); err != nil {
return fmt.Errorf("persist task control barrier: %w", err)
}
if action == "stop" {
if _, err := tx.Exec(`UPDATE dispatcher_inbox SET status='suppressed' WHERE dispatcher_id=? AND tenant_id=? AND task_id=? AND status='pending'`, dispatcherID, tenantID, taskID); err != nil {
return fmt.Errorf("suppress unstarted stopped-task commands: %w", err)
}
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit task control barrier: %w", err)
}
return nil
}
// CompleteControl atomically records the applied state and durable outbox
// only after an Agent control RPC accepted the action. It does not claim that
// active-call drain or hangup has already finished.
func (s *Store) CompleteControl(dispatcherID string, tenantID int64, taskID, action, eventID string) error {
var expected, final string
switch action {
case "pause":
expected, final = "pausing", "paused"
case "stop":
expected, final = "stopping", "stopped"
case "resume":
expected, final = "resuming", ""
default:
return fmt.Errorf("invalid task control action %q", action)
}
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
var state, status string
err = tx.QueryRow(`SELECT control_state,status FROM dispatcher_tasks WHERE dispatcher_id=? AND tenant_id=? AND task_id=? AND present=1`, dispatcherID, tenantID, taskID).Scan(&state, &status)
if err != nil {
return fmt.Errorf("read prepared task control: %w", err)
}
if state != expected || (action == "resume" && status != "running") {
return fmt.Errorf("task %q control %q is no longer prepared or authorized", taskID, action)
}
if _, err := tx.Exec(`UPDATE dispatcher_tasks SET control_state=? WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, final, dispatcherID, tenantID, taskID); err != nil {
return fmt.Errorf("persist applied task control: %w", err)
}
if err := enqueueControlAck(tx, dispatcherID, tenantID, eventID, "applied"); err != nil {
return err
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit task control and outbox: %w", err)
}
return nil
}
func (s *Store) RejectControl(dispatcherID string, tenantID int64, eventID string) error {
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
if err := enqueueControlAck(tx, dispatcherID, tenantID, eventID, "rejected"); err != nil {
return err
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit rejected task control: %w", err)
}
return nil
}
func enqueueControlAck(tx *sql.Tx, dispatcherID string, tenantID int64, eventID, status string) error {
if eventID == "" || len(eventID) > 255 || tenantID <= 0 || (status != "applied" && status != "rejected") {
return errors.New("invalid task control acknowledgment identity or status")
}
route, err := tenant.ResultRoute(dispatcherID)
if err != nil {
return err
}
body, err := json.Marshal(struct {
EventID string `json:"event_id"`
EventType string `json:"event_type"`
DispatcherID string `json:"dispatcher_id"`
TenantID int64 `json:"tenant_id"`
Payload struct {
Status string `json:"status"`
} `json:"payload"`
}{EventID: eventID, EventType: "task.control", DispatcherID: dispatcherID, TenantID: tenantID, Payload: struct {
Status string `json:"status"`
}{status}})
if err != nil {
return fmt.Errorf("encode task control acknowledgment: %w", err)
}
if err := contract.ValidateCurrent("mq", body); err != nil {
return fmt.Errorf("task control acknowledgment violates MQ contract: %w", err)
}
var oldType string
var oldBody []byte
err = tx.QueryRow(`SELECT event_type,body FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&oldType, &oldBody)
if errors.Is(err, sql.ErrNoRows) {
if _, err := tx.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?)`, dispatcherID, eventID, "task.control", route.BindingKey, body); err != nil {
return fmt.Errorf("persist task control outbox: %w", err)
}
return nil
}
if err != nil {
return fmt.Errorf("inspect existing task control outbox: %w", err)
}
if oldType != "task.control" || string(oldBody) != string(body) {
return fmt.Errorf("control event identity %q conflicts with another outbound result", eventID)
}
if _, err := tx.Exec(`UPDATE dispatcher_outbox SET confirmed=0,confirmed_at=NULL WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID); err != nil {
return fmt.Errorf("requeue repeated control acknowledgment: %w", err)
}
return nil
}