From ce5e8da2ee59691ee3535ea1a01d7c2d8fd7cef2 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 28 Sep 2026 07:17:54 +0800 Subject: [PATCH] refactor: remove unused v0.3 task discovery runtime --- internal/configread/client.go | 61 -------- internal/configread/discovery_v03_test.go | 134 ----------------- internal/configread/discovery_v04_test.go | 11 ++ internal/contract/contract.go | 4 - internal/contract/schema_test.go | 13 -- internal/store/local_v01_discovery.go | 47 +----- internal/store/local_v03_discovery.go | 165 --------------------- internal/store/local_v03_discovery_test.go | 118 --------------- internal/store/local_v03_failure_test.go | 78 ---------- internal/store/local_v04_discovery_test.go | 4 + internal/store/local_v04_lookup.go | 49 ++++++ 11 files changed, 71 insertions(+), 613 deletions(-) delete mode 100644 internal/configread/discovery_v03_test.go delete mode 100644 internal/store/local_v03_discovery.go delete mode 100644 internal/store/local_v03_discovery_test.go delete mode 100644 internal/store/local_v03_failure_test.go create mode 100644 internal/store/local_v04_lookup.go diff --git a/internal/configread/client.go b/internal/configread/client.go index 8459124..ede7646 100644 --- a/internal/configread/client.go +++ b/internal/configread/client.go @@ -56,7 +56,6 @@ type Snapshot struct { QuotaRevision int64 TenantMaxConcurrentCalls int64 QuotaValidUntil time.Time - DiscoveryCursor string FetchedAt time.Time ExpiresAt time.Time SIP json.RawMessage @@ -221,49 +220,6 @@ func (c *Client) ReadTaskStatus(ctx context.Context, taskID, tenantID, tenantKey }, nil } -// ReadTaskDiscovery reads one page of latest task states after a durable event ID. -// SaaS guarantees page completeness and ordering; no per-item event ID is provided. -func (c *Client) ReadTaskDiscovery(ctx context.Context, after string) (TaskDiscovery, error) { - from, err := parseEventCursor(after) - if err != nil { - return TaskDiscovery{}, fmt.Errorf("invalid task-discovery after cursor: %w", err) - } - path := configReadPath + "/tasks?" + url.Values{"after": {after}}.Encode() - body, err := c.getTaskDiscovery(ctx, path) - if err != nil { - return TaskDiscovery{}, err - } - var response taskDiscoveryResponse - if err := json.Unmarshal(body, &response); err != nil { - return TaskDiscovery{}, fmt.Errorf("decode task discovery: %w", err) - } - if response.DispatcherID != c.dispatcherID { - return TaskDiscovery{}, errors.New("task-discovery dispatcher identity does not match the request") - } - to, err := parseEventCursor(response.NextCursor) - if err != nil { - return TaskDiscovery{}, fmt.Errorf("invalid task-discovery next cursor: %w", err) - } - if (len(response.Tasks) == 0 && to != from) || (len(response.Tasks) != 0 && to <= from) { - return TaskDiscovery{}, errors.New("task-discovery page cursor does not match returned tasks") - } - seen := make(map[string]struct{}, len(response.Tasks)) - tenantByID, idByTenant := make(map[string]string), make(map[string]string) - for _, task := range response.Tasks { - if _, duplicate := seen[task.TaskID]; duplicate { - return TaskDiscovery{}, fmt.Errorf("task discovery contains duplicate task ID %q", task.TaskID) - } - seen[task.TaskID] = struct{}{} - if err := validateTenantBinding(tenantByID, idByTenant, task.TenantID, task.TenantKey); err != nil { - return TaskDiscovery{}, err - } - } - return TaskDiscovery{ - DispatcherID: response.DispatcherID, FromCursor: after, NextCursor: response.NextCursor, - Tasks: response.Tasks, Body: append(json.RawMessage(nil), body...), - }, nil -} - func parseEventCursor(cursor string) (uint64, error) { value, err := strconv.ParseUint(cursor, 10, 64) if err != nil || strconv.FormatUint(value, 10) != cursor { @@ -298,17 +254,6 @@ func (c *Client) getConfig(ctx context.Context, path string) (json.RawMessage, e return body, nil } -func (c *Client) getTaskDiscovery(ctx context.Context, path string) (json.RawMessage, error) { - body, err := c.get(ctx, path) - if err != nil { - return nil, err - } - if err := contract.ValidateLocalTaskDiscovery(body); err != nil { - return nil, fmt.Errorf("validate task-discovery response: %w", err) - } - return body, nil -} - func (c *Client) get(ctx context.Context, path string) (json.RawMessage, error) { relative, err := url.Parse(path) if err != nil { @@ -373,12 +318,6 @@ type sipConfigResponse struct { } `json:"artifact"` } -type taskDiscoveryResponse struct { - DispatcherID string `json:"dispatcher_id"` - NextCursor string `json:"next_cursor"` - Tasks []DiscoveredTask `json:"tasks"` -} - type taskConfigResponse struct { Resource string `json:"resource"` DispatcherID string `json:"dispatcher_id"` diff --git a/internal/configread/discovery_v03_test.go b/internal/configread/discovery_v03_test.go deleted file mode 100644 index 1abbaf8..0000000 --- a/internal/configread/discovery_v03_test.go +++ /dev/null @@ -1,134 +0,0 @@ -package configread - -import ( - "context" - "encoding/json" - "net/http" - "net/http/httptest" - "reflect" - "strings" - "testing" -) - -const eventTestDispatcherID = "11111111-1111-4111-8111-111111111111" - -func eventTestClient(t *testing.T, server *httptest.Server) *Client { - t.Helper() - client, err := NewClient(server.URL, eventTestDispatcherID, mockTestSecret, server.Client()) - if err != nil { - t.Fatal(err) - } - return client -} - -func TestEventCursorPagesUseOneTaskListAndAllowRepeatedTaskIDs(t *testing.T) { - pages := map[string]string{ - "0": "task-discovery-page-v0.3.json", - "2": "task-discovery-updated-v0.3.json", - "3": "task-discovery-removed-v0.3.json", - "4": "task-discovery-empty-v0.3.json", - } - var got []string - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - got = append(got, r.URL.RequestURI()) - if r.Header.Get(dispatcherIDHeader) != eventTestDispatcherID || r.Header.Get(dispatcherSecretHeader) != mockTestSecret { - t.Error("Dispatcher credentials were not sent") - } - name := pages[r.URL.Query().Get("after")] - if name == "" || r.URL.Query().Has("page_token") { - t.Errorf("unexpected task request: %s", r.URL.RequestURI()) - w.WriteHeader(http.StatusBadRequest) - return - } - w.Header().Set("Content-Type", "application/json") - _, _ = w.Write(readConfigFixture(t, name)) - })) - defer server.Close() - client := eventTestClient(t, server) - cursor := "0" - var statuses []string - for range 4 { - page, err := client.ReadTaskDiscovery(context.Background(), cursor) - if err != nil { - t.Fatal(err) - } - for _, task := range page.Tasks { - statuses = append(statuses, task.TaskID+":"+task.Status) - } - cursor = page.NextCursor - } - if cursor != "4" || !reflect.DeepEqual(statuses, []string{"a01:running", "a02:paused", "a01:paused", "a02:removed"}) { - t.Fatalf("statuses=%v cursor=%q", statuses, cursor) - } - want := []string{ - "/internal/v1/dispatcher/tasks?after=0", - "/internal/v1/dispatcher/tasks?after=2", - "/internal/v1/dispatcher/tasks?after=3", - "/internal/v1/dispatcher/tasks?after=4", - } - if !reflect.DeepEqual(got, want) { - t.Fatalf("queries=%v want=%v", got, want) - } -} - -func TestEventCursorRejectsBrokenPagesWithoutFallback(t *testing.T) { - valid := string(readConfigFixture(t, "task-discovery-page-v0.3.json")) - var oversized map[string]any - if err := json.Unmarshal([]byte(valid), &oversized); err != nil { - t.Fatal(err) - } - rows := make([]any, 257) - for i := range rows { - rows[i] = oversized["tasks"].([]any)[0] - } - oversized["tasks"] = rows - oversizedBody, err := json.Marshal(oversized) - if err != nil { - t.Fatal(err) - } - cases := map[string]string{ - "nonprogress": strings.Replace(valid, `"next_cursor": "2"`, `"next_cursor": "0"`, 1), - "duplicate": strings.Replace(valid, `"task_id": "a02"`, `"task_id": "a01"`, 1), - "wrong dispatcher": strings.Replace(valid, eventTestDispatcherID, "44444444-4444-4444-8444-444444444444", 1), - "oversized page": string(oversizedBody), - "changes": string(readConfigFixture(t, "task-discovery-invalid-changes-v0.3.json")), - "cursor format": string(readConfigFixture(t, "task-discovery-invalid-cursor-v0.3.json")), - "old contract": string(readConfigFixture(t, "task-discovery-snapshot-v0.2.json")), - } - for name, body := range cases { - t.Run(name, func(t *testing.T) { - requests := 0 - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - requests++ - w.Header().Set("Content-Type", "application/json") - _, _ = w.Write([]byte(body)) - })) - defer server.Close() - client := eventTestClient(t, server) - if _, err := client.ReadTaskDiscovery(context.Background(), "0"); err == nil { - t.Fatal("invalid task page was accepted") - } - if requests != 1 { - t.Fatalf("requests=%d, want one without fallback", requests) - } - }) - } -} - -func TestEventCursorRejectsExpiredStatusWithoutReset(t *testing.T) { - requests := 0 - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - requests++ - w.Header().Set("Content-Type", "application/json") - w.WriteHeader(http.StatusGone) - _, _ = w.Write([]byte(`{"schema_version":"task-discovery.v0.3-proposal","resource":"error","error":{"code":"cursor_expired","message":"expired"}}`)) - })) - defer server.Close() - client := eventTestClient(t, server) - if _, err := client.ReadTaskDiscovery(context.Background(), "1"); err == nil { - t.Fatal("HTTP 410 accepted") - } - if requests != 1 { - t.Fatalf("requests=%d, want one without reset to initial cursor", requests) - } -} diff --git a/internal/configread/discovery_v04_test.go b/internal/configread/discovery_v04_test.go index 6c45912..bfa1484 100644 --- a/internal/configread/discovery_v04_test.go +++ b/internal/configread/discovery_v04_test.go @@ -9,6 +9,17 @@ import ( "testing" ) +const eventTestDispatcherID = "11111111-1111-4111-8111-111111111111" + +func eventTestClient(t *testing.T, server *httptest.Server) *Client { + t.Helper() + client, err := NewClient(server.URL, eventTestDispatcherID, mockTestSecret, server.Client()) + if err != nil { + t.Fatal(err) + } + return client +} + func TestV04SnapshotAndLiveChangesUseSeparateRequests(t *testing.T) { var paths []string server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { diff --git a/internal/contract/contract.go b/internal/contract/contract.go index 0aa57a2..e1e0195 100644 --- a/internal/contract/contract.go +++ b/internal/contract/contract.go @@ -72,10 +72,6 @@ func ValidateLocalConfigRead(raw []byte) error { return validateLocalSchema("config-read-v0.1.schema.json", raw) } -func ValidateLocalTaskDiscovery(raw []byte) error { - return validateLocalSchema("task-discovery-v0.3-proposal.schema.json", raw) -} - func ValidateLocalTaskDiscoveryV04(raw []byte) error { return validateLocalSchema("task-discovery-v0.4-proposal.schema.json", raw) } diff --git a/internal/contract/schema_test.go b/internal/contract/schema_test.go index a262055..4b42e9e 100644 --- a/internal/contract/schema_test.go +++ b/internal/contract/schema_test.go @@ -64,19 +64,6 @@ func TestProjectLocalConfigurationSchemasValidatePositivesAndRejectNegatives(t * if err := ValidateLocalConfigRead(read("config-read-invalid-extra-property-v0.1.json")); err == nil { t.Fatal("config-read schema accepted an additional property") } - page := []byte(`{"schema_version":"task-discovery.v0.3-proposal","dispatcher_id":"11111111-1111-4111-8111-111111111111","next_cursor":"12","tasks":[{"task_id":"22222222-2222-4222-8222-222222222222","tenant_id":"33333333-3333-4333-8333-333333333333","tenant_key":"tenant-A","status":"removed","task_revision":4}]}`) - if err := ValidateLocalTaskDiscovery(page); err != nil { - t.Fatalf("valid discovery event page: %v", err) - } - for name, raw := range map[string][]byte{ - "old snapshot": read("task-discovery-snapshot-v0.2.json"), - "changes": []byte(`{"schema_version":"task-discovery.v0.3-proposal","dispatcher_id":"11111111-1111-4111-8111-111111111111","next_cursor":"12","tasks":[],"changes":[]}`), - "leading zero cursor": []byte(`{"schema_version":"task-discovery.v0.3-proposal","dispatcher_id":"11111111-1111-4111-8111-111111111111","next_cursor":"012","tasks":[]}`), - } { - if err := ValidateLocalTaskDiscovery(raw); err == nil { - t.Fatalf("task-discovery schema accepted %s", name) - } - } } func TestV04LocalSchemasValidateExamples(t *testing.T) { diff --git a/internal/store/local_v01_discovery.go b/internal/store/local_v01_discovery.go index 9028c23..dc19985 100644 --- a/internal/store/local_v01_discovery.go +++ b/internal/store/local_v01_discovery.go @@ -10,14 +10,13 @@ import ( ) var ( - ErrLocalDiscoveryCursorMismatch = errors.New("local task-discovery cursor mismatch") - ErrLocalDiscoveryUnavailable = errors.New("task discovery has no fresh complete response") - ErrLocalDiscoveryTaskLimit = errors.New("task discovery exceeds the 256-task per-Dispatcher limit") - ErrLocalDiscoveryTaskMissing = errors.New("local task-discovery update targets an unknown task") - ErrLocalTaskRemoved = errors.New("local task assignment has been removed") - ErrLocalTaskUnassigned = errors.New("task is not assigned to this Dispatcher") - ErrLocalTaskPaused = errors.New("task admission is paused") - ErrLocalTaskStopped = errors.New("task admission is terminal") + ErrLocalDiscoveryUnavailable = errors.New("task discovery has no fresh complete response") + ErrLocalDiscoveryTaskLimit = errors.New("task discovery exceeds the 256-task per-Dispatcher limit") + ErrLocalDiscoveryTaskMissing = errors.New("local task-discovery update targets an unknown task") + ErrLocalTaskRemoved = errors.New("local task assignment has been removed") + ErrLocalTaskUnassigned = errors.New("task is not assigned to this Dispatcher") + ErrLocalTaskPaused = errors.New("task admission is paused") + ErrLocalTaskStopped = errors.New("task admission is terminal") ) type LocalDiscoveredTask struct { @@ -152,38 +151,6 @@ func mergeLocalTaskState(current, incoming string) string { } } -// CloseLocalTaskDiscoveryAdmission rejects new work without losing the cursor -// or changing an authoritative task pause/stop state. -func (s *Store) CloseLocalTaskDiscoveryAdmission(dispatcherID string) error { - s.mu.Lock() - defer s.mu.Unlock() - _, err := s.db.Exec(`UPDATE local_v03_task_discovery_state SET ready=0 WHERE dispatcher_id=?`, dispatcherID) - return err -} - -func (s *Store) LocalTaskDiscoveryCursor(dispatcherID string) (cursor string, exists bool, err error) { - s.mu.Lock() - defer s.mu.Unlock() - err = s.db.QueryRow(`SELECT cursor FROM local_v03_task_discovery_state WHERE dispatcher_id=?`, dispatcherID).Scan(&cursor) - if errors.Is(err, sql.ErrNoRows) { - var legacy, assignments int - if err := s.db.QueryRow(`SELECT COUNT(*) FROM local_v02_task_discovery_state WHERE dispatcher_id=?`, dispatcherID).Scan(&legacy); err != nil { - return "", false, err - } - if err := s.db.QueryRow(`SELECT COUNT(*) FROM local_v01_task_assignments WHERE dispatcher_id=?`, dispatcherID).Scan(&assignments); err != nil { - return "", false, err - } - if legacy != 0 || assignments != 0 { - return "", false, fmt.Errorf("%w: existing discovery state requires controlled drain before v0.3", ErrLocalDiscoveryUnavailable) - } - return "", false, nil - } - if err != nil { - return "", false, err - } - return cursor, true, nil -} - func (s *Store) LocalTaskAssignments(dispatcherID string) ([]LocalTaskAssignment, error) { s.mu.Lock() defer s.mu.Unlock() diff --git a/internal/store/local_v03_discovery.go b/internal/store/local_v03_discovery.go deleted file mode 100644 index 44e5928..0000000 --- a/internal/store/local_v03_discovery.go +++ /dev/null @@ -1,165 +0,0 @@ -package store - -import ( - "crypto/sha256" - "database/sql" - "encoding/hex" - "encoding/json" - "errors" - "fmt" - "strconv" - "time" -) - -const localEventPageLimit = 256 - -// LocalTaskDiscoveryPage is one complete HTTP response, not an entire task list. -// SaaS guarantees response ordering/completeness; its response-level cursor is -// the only event position supplied by the project-local v0.3 proposal. -type LocalTaskDiscoveryPage struct { - DispatcherID string - FromCursor string - NextCursor string - Tasks []LocalDiscoveredTask - Body json.RawMessage - ObservedAt time.Time -} - -func eventCursorValue(cursor string) (uint64, error) { - value, err := strconv.ParseUint(cursor, 10, 64) - if err != nil || strconv.FormatUint(value, 10) != cursor { - return 0, errors.New("event cursor must be a canonical decimal uint64") - } - return value, nil -} - -// ApplyLocalTaskDiscoveryPage commits assignments, the exact page evidence and -// its cursor together. A nonempty page never opens new execution admission; -// only an empty final page marks the initial catch-up/current poll complete. -func (s *Store) ApplyLocalTaskDiscoveryPage(page LocalTaskDiscoveryPage) error { - from, err := eventCursorValue(page.FromCursor) - if err != nil { - return fmt.Errorf("invalid task-discovery after cursor: %w", err) - } - to, err := eventCursorValue(page.NextCursor) - if err != nil { - return fmt.Errorf("invalid task-discovery next cursor: %w", err) - } - if page.DispatcherID == "" || page.ObservedAt.IsZero() || len(page.Tasks) > localEventPageLimit || !json.Valid(page.Body) || - (len(page.Tasks) == 0 && to != from) || (len(page.Tasks) != 0 && to <= from) { - return ErrLocalDiscoveryUnavailable - } - seen := make(map[string]struct{}, len(page.Tasks)) - for _, task := range page.Tasks { - if err := validateLocalDiscoveredTask(page.DispatcherID, task); err != nil { - return err - } - if _, duplicate := seen[task.TaskID]; duplicate { - return fmt.Errorf("task-discovery page has duplicate task ID %q", task.TaskID) - } - seen[task.TaskID] = struct{}{} - } - - s.mu.Lock() - defer s.mu.Unlock() - tx, err := s.db.Begin() - if err != nil { - return err - } - defer tx.Rollback() - var legacy int - if err := tx.QueryRow(`SELECT COUNT(*) FROM local_v02_task_discovery_state WHERE dispatcher_id=?`, page.DispatcherID).Scan(&legacy); err != nil { - return err - } - if legacy != 0 { - return fmt.Errorf("%w: v0.2 discovery checkpoint must be drained before switching", ErrLocalDiscoveryUnavailable) - } - var cursor string - err = tx.QueryRow(`SELECT cursor FROM local_v03_task_discovery_state WHERE dispatcher_id=?`, page.DispatcherID).Scan(&cursor) - if errors.Is(err, sql.ErrNoRows) { - var existing int - if err := tx.QueryRow(`SELECT COUNT(*) FROM local_v01_task_assignments WHERE dispatcher_id=?`, page.DispatcherID).Scan(&existing); err != nil { - return err - } - if existing != 0 { - return fmt.Errorf("%w: existing assignments have no v0.3 checkpoint", ErrLocalDiscoveryUnavailable) - } - if page.FromCursor != "0" { - return ErrLocalDiscoveryCursorMismatch - } - } else if err != nil { - return err - } else if cursor != page.FromCursor { - return ErrLocalDiscoveryCursorMismatch - } - now := page.ObservedAt.UTC().Format(time.RFC3339Nano) - for _, task := range page.Tasks { - if err := bindLocalTenant(tx, page.DispatcherID, task.TenantID, task.TenantKey, page.ObservedAt); err != nil { - return err - } - if err := applyLocalDiscoveredTask(tx, page.DispatcherID, task, now); err != nil { - return err - } - } - var active int - if err := tx.QueryRow(`SELECT COUNT(*) FROM local_v01_task_assignments WHERE dispatcher_id=? AND removed=0`, page.DispatcherID).Scan(&active); err != nil { - return err - } - if active > localEventPageLimit { - return ErrLocalDiscoveryTaskLimit - } - ready := 0 - if len(page.Tasks) == 0 { - ready = 1 - } - sum := sha256.Sum256(page.Body) - _, err = tx.Exec(`INSERT INTO local_v03_task_discovery_state(dispatcher_id,cursor,ready,body,content_sha256,updated_at) - VALUES(?,?,?,?,?,?) ON CONFLICT(dispatcher_id) DO UPDATE SET cursor=excluded.cursor,ready=excluded.ready, - body=excluded.body,content_sha256=excluded.content_sha256,updated_at=excluded.updated_at`, - page.DispatcherID, page.NextCursor, ready, []byte(page.Body), hex.EncodeToString(sum[:]), now) - if err != nil { - return err - } - return tx.Commit() -} - -// LocalDiscoveredTaskForConfig binds a config read to a fully caught-up, -// durable task assignment. A partial page cannot authorize configuration. -func (s *Store) LocalDiscoveredTaskForConfig(dispatcherID, taskID, tenantID string) (LocalTaskAssignment, error) { - s.mu.Lock() - defer s.mu.Unlock() - var ready int - var observedAt string - err := s.db.QueryRow(`SELECT ready,updated_at FROM local_v04_task_discovery_state WHERE dispatcher_id=?`, dispatcherID).Scan(&ready, &observedAt) - if errors.Is(err, sql.ErrNoRows) || (err == nil && ready != 1) { - return LocalTaskAssignment{}, ErrLocalDiscoveryUnavailable - } - if err != nil { - return LocalTaskAssignment{}, err - } - var assignment LocalTaskAssignment - var removed int - err = s.db.QueryRow(`SELECT dispatcher_id,task_id,tenant_id,tenant_key,task_revision,saas_status,admission_state,removed, - queue_exchange,routing_key,binding_key,queue_name,updated_at FROM local_v01_task_assignments WHERE dispatcher_id=? AND task_id=?`, - dispatcherID, taskID).Scan(&assignment.DispatcherID, &assignment.TaskID, &assignment.TenantID, &assignment.TenantKey, - &assignment.TaskRevision, &assignment.Status, &assignment.AdmissionState, &removed, &assignment.Queue.Exchange, - &assignment.Queue.RoutingKey, &assignment.Queue.BindingKey, &assignment.Queue.QueueName, &assignment.UpdatedAt) - if errors.Is(err, sql.ErrNoRows) { - return LocalTaskAssignment{}, ErrLocalDiscoveryTaskMissing - } - if err != nil { - return LocalTaskAssignment{}, err - } - if assignment.TenantID != tenantID { - return LocalTaskAssignment{}, ErrTenantBindingConflict - } - assignment.Removed = removed != 0 - if assignment.Removed { - return LocalTaskAssignment{}, ErrLocalTaskRemoved - } - assignment.DiscoveryObservedAt, err = time.Parse(time.RFC3339Nano, observedAt) - if err != nil { - return LocalTaskAssignment{}, fmt.Errorf("invalid durable discovery observation time: %w", err) - } - return assignment, nil -} diff --git a/internal/store/local_v03_discovery_test.go b/internal/store/local_v03_discovery_test.go deleted file mode 100644 index 1bd2f99..0000000 --- a/internal/store/local_v03_discovery_test.go +++ /dev/null @@ -1,118 +0,0 @@ -package store - -import ( - "crypto/sha256" - "encoding/hex" - "errors" - "testing" - "time" -) - -func eventPage(after, next string, tasks ...LocalDiscoveredTask) LocalTaskDiscoveryPage { - return LocalTaskDiscoveryPage{ - DispatcherID: "d-1", FromCursor: after, NextCursor: next, - Tasks: tasks, Body: []byte(`{"schema_version":"task-discovery.v0.3-proposal","tasks":[]}`), - ObservedAt: time.Now().UTC(), - } -} - -func eventTask(id, status string, revision int64) LocalDiscoveredTask { - return LocalDiscoveredTask{TaskID: id, TenantID: "tenant-a", TenantKey: "tenant-key-a", Status: status, TaskRevision: revision} -} - -func TestEventPagesPersistCursorAssignmentAndBootGateTogether(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - first := eventPage("0", "1", eventTask("task-a", "running", 1)) - if err := st.ApplyLocalTaskDiscoveryPage(first); err != nil { - t.Fatal(err) - } - cursor, exists, err := st.LocalTaskDiscoveryCursor("d-1") - if err != nil || !exists || cursor != "1" { - t.Fatalf("cursor=%q exists=%v err=%v", cursor, exists, err) - } - var ready int - var saved []byte - var digest string - if err := st.db.QueryRow(`SELECT ready,body,content_sha256 FROM local_v03_task_discovery_state WHERE dispatcher_id=?`, "d-1").Scan(&ready, &saved, &digest); err != nil { - t.Fatal(err) - } - sum := sha256.Sum256(first.Body) - if ready != 0 || string(saved) != string(first.Body) || digest != hex.EncodeToString(sum[:]) { - t.Fatalf("partial page admitted or changed evidence: ready=%d body=%s hash=%s", ready, saved, digest) - } - assignments, err := st.LocalTaskAssignments("d-1") - if err != nil || len(assignments) != 1 || assignments[0].Status != "running" || assignments[0].Queue.QueueName != "agent-call.d.d-1.task.task-a.v3" { - t.Fatalf("assignment=%+v err=%v", assignments, err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("1", "1")); err != nil { - t.Fatal(err) - } - if err := st.db.QueryRow(`SELECT ready FROM local_v03_task_discovery_state WHERE dispatcher_id=?`, "d-1").Scan(&ready); err != nil || ready != 1 { - t.Fatalf("completed page ready=%d err=%v", ready, err) - } -} - -func TestEventPageAllowsSameTaskNewStatusAndRemovedTombstone(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - for _, page := range []LocalTaskDiscoveryPage{ - eventPage("0", "1", eventTask("task-a", "running", 1)), - eventPage("1", "2", eventTask("task-a", "paused", 1)), - eventPage("2", "3", eventTask("task-a", "removed", 1)), - eventPage("3", "3"), - } { - if err := st.ApplyLocalTaskDiscoveryPage(page); err != nil { - t.Fatal(err) - } - } - assignments, err := st.LocalTaskAssignments("d-1") - if err != nil || len(assignments) != 1 || !assignments[0].Removed || assignments[0].Status != "removed" || assignments[0].AdmissionState != "removed" { - t.Fatalf("removed task=%+v err=%v", assignments, err) - } - cursor, exists, err := st.LocalTaskDiscoveryCursor("d-1") - if err != nil || !exists || cursor != "3" { - t.Fatalf("cursor=%q exists=%v err=%v", cursor, exists, err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("1", "4", eventTask("task-a", "running", 1))); !errors.Is(err, ErrLocalDiscoveryCursorMismatch) { - t.Fatalf("old cursor unexpectedly accepted: %v", err) - } -} - -func TestEventPageRunningAfterDiscoveryPauseWaitsForMQResume(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - for _, page := range []LocalTaskDiscoveryPage{ - eventPage("0", "1", eventTask("task-a", "paused", 1)), - eventPage("1", "1"), - eventPage("1", "2", eventTask("task-a", "running", 2)), - eventPage("2", "2"), - } { - if err := st.ApplyLocalTaskDiscoveryPage(page); err != nil { - t.Fatal(err) - } - } - assignment, err := st.LocalTaskAssignment("d-1", "task-a") - if err != nil || assignment.Status != "running" || assignment.TaskRevision != 2 || assignment.AdmissionState != "paused" { - t.Fatalf("discovery running incorrectly reopened admission: %+v err=%v", assignment, err) - } - cursor, exists, err := st.LocalTaskDiscoveryCursor("d-1") - if err != nil || !exists || cursor != "2" { - t.Fatalf("latest task status did not advance cursor: %q exists=%v err=%v", cursor, exists, err) - } - assignment, err = st.ResumeLocalTaskAdmission("d-1", "task-a", "tenant-a", "tenant-key-a", "running", 2) - if err != nil || assignment.Status != "running" || assignment.AdmissionState != "running" { - t.Fatalf("fresh MQ resume did not reopen admission: %+v err=%v", assignment, err) - } -} - -func TestEventPageCannotInterpretOldSnapshotCursor(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - if _, err := st.db.Exec(`INSERT INTO local_v02_task_discovery_state(dispatcher_id,cursor,last_mode,ready,updated_at) VALUES('d-1','opaque-old','snapshot',1,'now')`); err != nil { - t.Fatal(err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("0", "1", eventTask("task-a", "running", 1))); !errors.Is(err, ErrLocalDiscoveryUnavailable) { - t.Fatalf("legacy cursor should block, got %v", err) - } - var count int - if err := st.db.QueryRow(`SELECT COUNT(*) FROM local_v03_task_discovery_state WHERE dispatcher_id='d-1'`).Scan(&count); err != nil || count != 0 { - t.Fatalf("legacy cursor created new checkpoint: count=%d err=%v", count, err) - } -} diff --git a/internal/store/local_v03_failure_test.go b/internal/store/local_v03_failure_test.go deleted file mode 100644 index fd5faa9..0000000 --- a/internal/store/local_v03_failure_test.go +++ /dev/null @@ -1,78 +0,0 @@ -package store - -import ( - "errors" - "strings" - "testing" -) - -func TestEventPageWriteFailureRollsBackAssignmentCursorAndEvidence(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("0", "1", eventTask("task-a", "running", 1))); err != nil { - t.Fatal(err) - } - var before []byte - if err := st.db.QueryRow(`SELECT body FROM local_v03_task_discovery_state WHERE dispatcher_id='d-1'`).Scan(&before); err != nil { - t.Fatal(err) - } - if _, err := st.db.Exec(`CREATE TRIGGER reject_discovery_update BEFORE UPDATE ON local_v01_task_assignments - BEGIN SELECT RAISE(ABORT, 'forced assignment write failure'); END`); err != nil { - t.Fatal(err) - } - err := st.ApplyLocalTaskDiscoveryPage(eventPage("1", "2", eventTask("task-a", "paused", 2))) - if err == nil || !strings.Contains(err.Error(), "forced assignment write failure") { - t.Fatalf("assignment write failure was swallowed: %v", err) - } - cursor, exists, err := st.LocalTaskDiscoveryCursor("d-1") - if err != nil || !exists || cursor != "1" { - t.Fatalf("failed page advanced cursor: %q exists=%v err=%v", cursor, exists, err) - } - assignment, err := st.LocalTaskAssignment("d-1", "task-a") - if err != nil || assignment.Status != "running" || assignment.TaskRevision != 1 { - t.Fatalf("failed page changed assignment: %+v err=%v", assignment, err) - } - var after []byte - if err := st.db.QueryRow(`SELECT body FROM local_v03_task_discovery_state WHERE dispatcher_id='d-1'`).Scan(&after); err != nil || string(after) != string(before) { - t.Fatalf("failed page replaced saved evidence: before=%s after=%s err=%v", before, after, err) - } -} - -func TestEventPageCannotUndoMQPauseStopOrRemoval(t *testing.T) { - st := openLocalDiscoveryTestStore(t) - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("0", "1", eventTask("task-a", "running", 1))); err != nil { - t.Fatal(err) - } - if _, err := st.SetLocalTaskAdmissionBarrier("d-1", "task-a", "tenant-a", "tenant-key-a", "paused"); err != nil { - t.Fatal(err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("1", "2", eventTask("task-a", "running", 2))); err != nil { - t.Fatal(err) - } - assignment, err := st.LocalTaskAssignment("d-1", "task-a") - if err != nil || assignment.AdmissionState != "paused" { - t.Fatalf("stale running discovery overrode MQ pause: %+v err=%v", assignment, err) - } - if _, err := st.SetLocalTaskAdmissionBarrier("d-1", "task-a", "tenant-a", "tenant-key-a", "stopped"); err != nil { - t.Fatal(err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("2", "3", eventTask("task-a", "running", 3))); err != nil { - t.Fatal(err) - } - assignment, err = st.LocalTaskAssignment("d-1", "task-a") - if err != nil || assignment.AdmissionState != "stopped" { - t.Fatalf("stale running discovery overrode MQ stop: %+v err=%v", assignment, err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("3", "4", eventTask("task-a", "removed", 3))); err != nil { - t.Fatal(err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("4", "5", eventTask("task-a", "running", 4))); !errors.Is(err, ErrLocalTaskRemoved) { - t.Fatalf("removed task was revived: %v", err) - } - if err := st.ApplyLocalTaskDiscoveryPage(eventPage("4", "4")); err != nil { - t.Fatal(err) - } - assignments, err := st.LocalTaskAssignments("d-1") - if err != nil || len(assignments) != 1 || !assignments[0].Removed || assignments[0].AdmissionState != "removed" { - t.Fatalf("later event revived removed task: %+v err=%v", assignments, err) - } -} diff --git a/internal/store/local_v04_discovery_test.go b/internal/store/local_v04_discovery_test.go index 783d744..1f0977e 100644 --- a/internal/store/local_v04_discovery_test.go +++ b/internal/store/local_v04_discovery_test.go @@ -6,6 +6,10 @@ import ( "time" ) +func eventTask(id, status string, revision int64) LocalDiscoveredTask { + return LocalDiscoveredTask{TaskID: id, TenantID: "tenant-a", TenantKey: "tenant-key-a", Status: status, TaskRevision: revision} +} + func TestV04CompleteSnapshotReplacesMembershipWithoutReopeningPause(t *testing.T) { st := openLocalDiscoveryTestStore(t) now := time.Now().UTC() diff --git a/internal/store/local_v04_lookup.go b/internal/store/local_v04_lookup.go new file mode 100644 index 0000000..f665158 --- /dev/null +++ b/internal/store/local_v04_lookup.go @@ -0,0 +1,49 @@ +package store + +import ( + "database/sql" + "errors" + "fmt" + "time" +) + +// LocalDiscoveredTaskForConfig binds a config read to a fully caught-up, +// durable task assignment. A partial page cannot authorize configuration. +func (s *Store) LocalDiscoveredTaskForConfig(dispatcherID, taskID, tenantID string) (LocalTaskAssignment, error) { + s.mu.Lock() + defer s.mu.Unlock() + var ready int + var observedAt string + err := s.db.QueryRow(`SELECT ready,updated_at FROM local_v04_task_discovery_state WHERE dispatcher_id=?`, dispatcherID).Scan(&ready, &observedAt) + if errors.Is(err, sql.ErrNoRows) || (err == nil && ready != 1) { + return LocalTaskAssignment{}, ErrLocalDiscoveryUnavailable + } + if err != nil { + return LocalTaskAssignment{}, err + } + var assignment LocalTaskAssignment + var removed int + err = s.db.QueryRow(`SELECT dispatcher_id,task_id,tenant_id,tenant_key,task_revision,saas_status,admission_state,removed, + queue_exchange,routing_key,binding_key,queue_name,updated_at FROM local_v01_task_assignments WHERE dispatcher_id=? AND task_id=?`, + dispatcherID, taskID).Scan(&assignment.DispatcherID, &assignment.TaskID, &assignment.TenantID, &assignment.TenantKey, + &assignment.TaskRevision, &assignment.Status, &assignment.AdmissionState, &removed, &assignment.Queue.Exchange, + &assignment.Queue.RoutingKey, &assignment.Queue.BindingKey, &assignment.Queue.QueueName, &assignment.UpdatedAt) + if errors.Is(err, sql.ErrNoRows) { + return LocalTaskAssignment{}, ErrLocalDiscoveryTaskMissing + } + if err != nil { + return LocalTaskAssignment{}, err + } + if assignment.TenantID != tenantID { + return LocalTaskAssignment{}, ErrTenantBindingConflict + } + assignment.Removed = removed != 0 + if assignment.Removed { + return LocalTaskAssignment{}, ErrLocalTaskRemoved + } + assignment.DiscoveryObservedAt, err = time.Parse(time.RFC3339Nano, observedAt) + if err != nil { + return LocalTaskAssignment{}, fmt.Errorf("invalid durable discovery observation time: %w", err) + } + return assignment, nil +}