91 lines
2.9 KiB
Go
91 lines
2.9 KiB
Go
package configread
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/url"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/contract"
|
|
)
|
|
|
|
// CurrentTaskPage is one page of the current, cursor-based discovery list.
|
|
// A page ends the list only when Tasks is empty; short pages do not end it.
|
|
type CurrentTaskPage struct {
|
|
DispatcherID string `json:"dispatcher_id"`
|
|
Cursor string `json:"cursor"`
|
|
Tasks []CurrentDiscoveredTask `json:"tasks"`
|
|
}
|
|
|
|
type CurrentDiscoveredTask struct {
|
|
TaskID string `json:"task_id"`
|
|
TenantID int64 `json:"tenant_id"`
|
|
Status string `json:"status"`
|
|
TaskRevision int64 `json:"task_revision"`
|
|
}
|
|
|
|
// ReadCurrentTasks retrieves exactly one page, with no old snapshot/watermark
|
|
// fallback. The caller must persist non-empty pages before advancing after.
|
|
func (c *Client) ReadCurrentTasks(ctx context.Context, after string) (CurrentTaskPage, error) {
|
|
path := configReadPath + "/tasks"
|
|
if after != "" {
|
|
path += "?after=" + url.QueryEscape(after)
|
|
}
|
|
body, err := c.get(ctx, path)
|
|
if err != nil {
|
|
return CurrentTaskPage{}, err
|
|
}
|
|
if err := contract.ValidateCurrent("task-discovery", body); err != nil {
|
|
return CurrentTaskPage{}, err
|
|
}
|
|
var page CurrentTaskPage
|
|
if err := json.Unmarshal(body, &page); err != nil {
|
|
return CurrentTaskPage{}, fmt.Errorf("decode task discovery: %w", err)
|
|
}
|
|
if page.DispatcherID != c.dispatcherID {
|
|
return CurrentTaskPage{}, errors.New("task discovery dispatcher owner mismatch")
|
|
}
|
|
if len(page.Tasks) > 0 && page.Cursor == after {
|
|
return CurrentTaskPage{}, errors.New("non-empty task discovery page did not advance cursor")
|
|
}
|
|
seen := make(map[string]struct{}, len(page.Tasks))
|
|
for _, task := range page.Tasks {
|
|
if _, exists := seen[task.TaskID]; exists {
|
|
return CurrentTaskPage{}, fmt.Errorf("duplicate task %q in discovery page", task.TaskID)
|
|
}
|
|
seen[task.TaskID] = struct{}{}
|
|
}
|
|
return page, nil
|
|
}
|
|
|
|
// ReadAllCurrentTasks collects a cold-start list. Callers must validate and
|
|
// commit it atomically before opening task admission; no page is durable here.
|
|
func (c *Client) ReadAllCurrentTasks(ctx context.Context) ([]CurrentDiscoveredTask, string, error) {
|
|
var tasks []CurrentDiscoveredTask
|
|
seenTasks := make(map[string]struct{})
|
|
seenCursors := make(map[string]struct{})
|
|
cursor := ""
|
|
for {
|
|
page, err := c.ReadCurrentTasks(ctx, cursor)
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
if len(page.Tasks) == 0 {
|
|
return tasks, page.Cursor, nil
|
|
}
|
|
if _, repeated := seenCursors[page.Cursor]; repeated {
|
|
return nil, "", errors.New("repeated non-empty task discovery cursor")
|
|
}
|
|
seenCursors[page.Cursor] = struct{}{}
|
|
for _, task := range page.Tasks {
|
|
if _, repeated := seenTasks[task.TaskID]; repeated {
|
|
return nil, "", fmt.Errorf("task %q occurs in multiple discovery pages", task.TaskID)
|
|
}
|
|
seenTasks[task.TaskID] = struct{}{}
|
|
tasks = append(tasks, task)
|
|
}
|
|
cursor = page.Cursor
|
|
}
|
|
}
|