From 9df9664c38af80681917bae09aac40c50f61f300 Mon Sep 17 00:00:00 2001 From: Chris Miles Date: Mon, 20 Jul 2026 07:05:38 +0000 Subject: [PATCH 1/3] fix: complete unified event contract --- README.md | 61 +- bus/observer_adapter_contract_test.go | 102 ++++ docs/bus-design.md | 2 +- docs/bus-implementation-checklist.md | 2 +- docs/events.md | 53 +- plan.md => docs/plan.md | 4 +- docs/readme/testcounts/integration_count.json | 2 +- .../observability_contract_test.go | 300 ++++++++++ driver/sqlqueuecore/queue_database_impl.go | 36 +- driver/sqsqueue/worker_sqs_impl_test.go | 9 + event_contract_audit_test.go | 232 ++++++++ event_reference_test.go | 533 ++++++++++++++++++ .../root/process_events_integration_test.go | 46 +- observability.go | 25 +- observability_branches_test.go | 57 +- 15 files changed, 1392 insertions(+), 72 deletions(-) create mode 100644 bus/observer_adapter_contract_test.go rename plan.md => docs/plan.md (98%) create mode 100644 driver/sqlqueuecore/observability_contract_test.go create mode 100644 event_contract_audit_test.go create mode 100644 event_reference_test.go diff --git a/README.md b/README.md index f350d9b..19fcbc0 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ Latest tag - Unit tests (executed count) + Unit tests (executed count) Integration tests (executed count)

@@ -471,38 +471,43 @@ _ = q | Layer | EventKind | Meaning | | ---: | --- | --- | -| **queue** | dispatch_started | Public dispatch began. | -| **queue** | dispatch_succeeded | Backend acceptance completed; synchronous execution may still return an application error. | -| **queue** | dispatch_failed | Public dispatch failed before backend acceptance. | -| **queue** | enqueue_accepted | Job accepted by driver for enqueue. | -| **queue** | enqueue_rejected | Job enqueue failed. | -| **queue** | enqueue_duplicate | Duplicate job rejected due to uniqueness key. | -| **queue** | enqueue_canceled | Context cancellation prevented enqueue. | -| **worker** | process_started | Worker began processing job. | -| **worker** | process_succeeded | Handler returned success. | -| **worker** | process_failed | Handler returned an error or panicked. | +| **queue** | dispatch_started | After job validation, the root facade began a public dispatch. | +| **queue** | dispatch_succeeded | The backend accepted the public dispatch; synchronous execution may still return an application error. | +| **queue** | dispatch_failed | Public dispatch ended before backend acceptance. | +| **queue** | enqueue_accepted | The driver confirmed enqueue acceptance. | +| **queue** | enqueue_rejected | The driver rejected enqueue with an error. | +| **queue** | enqueue_duplicate | Uniqueness policy rejected the logical job as a duplicate. | +| **queue** | enqueue_canceled | Context cancellation or deadline expiry prevented enqueue. | +| **worker** | process_started | A physical handler attempt began. | +| **worker** | process_succeeded | Handler success crossed the driver's settlement boundary; runtimes without a settlement hook emit after handler return. | +| **worker** | process_failed | A handler attempt returned an error or panicked. | | **worker** | process_retried | A numbered application retry attempt began; infrastructure redelivery may repeat the fact. | -| **worker** | process_archived | Driver confirmed terminal settlement; unsupported paths omit this fact. | -| **queue** | queue_paused | Queue was paused (driver supports pause). | -| **queue** | queue_resumed | Queue was resumed. | -| **workflow** | job_started | A workflow job handler started execution. | -| **workflow** | job_succeeded | A workflow job handler completed successfully. | +| **worker** | process_archived | A SQL driver confirmed its fenced transition to terminal `dead` state; other built-in runtimes omit it. | +| **worker** | process_recovered | SQL bulk recovery requeued one stale in-flight claim; identity fields are unavailable at that boundary. | +| **worker** | republish_failed | An internal delay or retry replacement could not be published. | +| **worker** | settlement_failed | Delivery finalization, acknowledgement, deletion, or negative settlement failed or was ambiguous; redelivery remains possible. | +| **queue** | queue_paused | A supporting driver confirmed queue consumption was paused. | +| **queue** | queue_resumed | A supporting driver confirmed queue consumption was resumed. | +| **workflow** | job_started | A logical execution attempt began before handler lookup. | +| **workflow** | job_succeeded | Logical job success committed; settlement-aware drivers publish the fact only after settlement. | | **workflow** | job_failed | A logical job reached permanent or exhausted failure. | -| **workflow** | chain_started | A chain workflow was created and started. | -| **workflow** | chain_advanced | Chain progressed from one node to the next node. | -| **workflow** | chain_completed | Chain reached terminal success. | -| **workflow** | chain_failed | Chain reached terminal failure. | -| **workflow** | batch_started | A batch workflow was created and started. | -| **workflow** | batch_progressed | Batch state changed as jobs completed/failed. | +| **workflow** | chain_started | A chain record was created and initial dispatch began. | +| **workflow** | chain_advanced | A committed node outcome advanced the chain to its next node. | +| **workflow** | chain_completed | The final chain node committed terminal success. | +| **workflow** | chain_failed | A chain committed terminal failure. | +| **workflow** | batch_started | A batch record was created and initial member dispatch began. | +| **workflow** | batch_progressed | A batch member committed a terminal outcome. | | **workflow** | batch_completed | Batch reached terminal success (or allowed-failure completion). | -| **workflow** | batch_failed | Batch reached terminal failure. | -| **workflow** | batch_cancelled | Batch was cancelled before normal completion. | -| **workflow** | callback_started | Chain/batch callback execution started. | -| **workflow** | callback_succeeded | Chain/batch callback completed successfully. | -| **workflow** | callback_failed | Chain/batch callback returned an error. | +| **workflow** | batch_failed | A non-allowed member failure or initial member dispatch rejection committed terminal batch failure. | +| **workflow** | batch_cancelled | Remaining batch work was cancelled after terminal failure. | +| **workflow** | callback_started | A claimed terminal Catch, Then, or Finally callback began execution. | +| **workflow** | callback_succeeded | Terminal callback success crossed the applicable settlement boundary. | +| **workflow** | callback_failed | A terminal callback was invalid or unavailable, returned an error, or panicked. | Handler panics now emit `process_failed` before the original panic value is rethrown. This adds truthful failure telemetry without changing backend panic recovery or retry behavior. +Callback lifecycle facts cover terminal Catch, Then, and Finally callbacks. Inline `Batch.Progress` closures do not emit `callback_*` facts. + ## Examples Runnable examples live in the separate `examples` module ([`./examples`](./examples)). @@ -1131,7 +1136,7 @@ fmt.Println(snapshot.Paused("default")) #### Paused -Paused returns paused count for a queue. +Paused returns the observed pause state for a queue as zero or one. ```go collector := queue.NewStatsCollector() diff --git a/bus/observer_adapter_contract_test.go b/bus/observer_adapter_contract_test.go new file mode 100644 index 0000000..df4f914 --- /dev/null +++ b/bus/observer_adapter_contract_test.go @@ -0,0 +1,102 @@ +package bus + +import ( + "context" + "errors" + "reflect" + "testing" + "time" + + "github.com/goforj/queue/internal/workflow" +) + +// TestLegacyObserverAdapterTranslatesEveryField verifies the deprecated bus +// observer shape remains an exact compatibility projection of engine facts. +func TestLegacyObserverAdapterTranslatesEveryField(t *testing.T) { + assertWorkflowEventShape(t) + + type contextKey struct{} + ctx := context.WithValue(context.Background(), contextKey{}, "legacy-observer-context") + wantErr := errors.New("chain failed") + wantTime := time.Date(2026, time.July, 20, 14, 0, 0, 0, time.UTC) + internalEvent := workflow.Event{ + SchemaVersion: 9, + EventID: "evt_legacy_adapter", + Kind: workflow.EventChainFailed, + DispatchID: "dsp_legacy_adapter", + JobID: "job_legacy_adapter", + ChainID: "chn_legacy_adapter", + BatchID: "bat_legacy_adapter", + Attempt: 4, + JobType: "reports:archive", + JobKey: "job-key-legacy-adapter", + Queue: "critical", + Duration: 73 * time.Millisecond, + Time: wantTime, + Err: wantErr, + } + + var ( + gotContext context.Context + gotEvent Event + ) + adapter := legacyObserverAdapter{observer: ObserverFunc(func(observedContext context.Context, event Event) { + gotContext = observedContext + gotEvent = event + })} + adapter.Observe(ctx, internalEvent) + + wantEvent := Event{ + SchemaVersion: internalEvent.SchemaVersion, + EventID: internalEvent.EventID, + Kind: EventChainFailed, + DispatchID: internalEvent.DispatchID, + JobID: internalEvent.JobID, + ChainID: internalEvent.ChainID, + BatchID: internalEvent.BatchID, + Attempt: internalEvent.Attempt, + JobType: internalEvent.JobType, + JobKey: internalEvent.JobKey, + Queue: internalEvent.Queue, + Duration: internalEvent.Duration, + Time: internalEvent.Time, + Err: internalEvent.Err, + } + if !reflect.DeepEqual(gotEvent, wantEvent) { + t.Fatalf("translated legacy event = %+v, want %+v", gotEvent, wantEvent) + } + if gotContext == nil || gotContext.Value(contextKey{}) != "legacy-observer-context" { + t.Fatalf("legacy observer context = %v, want original workflow context", gotContext) + } +} + +// assertWorkflowEventShape pins the internal adapter input so adding a field +// requires an explicit decision about its legacy projection. +func assertWorkflowEventShape(t *testing.T) { + t.Helper() + wantFields := []string{ + "SchemaVersion", + "EventID", + "Kind", + "DispatchID", + "JobID", + "ChainID", + "BatchID", + "Attempt", + "JobType", + "JobKey", + "Queue", + "Duration", + "Time", + "Err", + } + eventType := reflect.TypeOf(workflow.Event{}) + if eventType.NumField() != len(wantFields) { + t.Fatalf("workflow.Event fields = %d, want %d; update the legacy observer adapter contract", eventType.NumField(), len(wantFields)) + } + for index, wantField := range wantFields { + if gotField := eventType.Field(index).Name; gotField != wantField { + t.Fatalf("workflow.Event field %d = %s, want %s; update the legacy observer adapter contract", index, gotField, wantField) + } + } +} diff --git a/docs/bus-design.md b/docs/bus-design.md index cdf68b1..e76388b 100644 --- a/docs/bus-design.md +++ b/docs/bus-design.md @@ -2,7 +2,7 @@ > **Status:** This is the original bus design record, retained to explain the > version-one wire and API constraints. It is not the current architecture or -> implementation roadmap; use `plan.md` for both. +> implementation roadmap; use [`docs/plan.md`](./plan.md) for both. ## Current Ownership diff --git a/docs/bus-implementation-checklist.md b/docs/bus-implementation-checklist.md index 0b784da..875760f 100644 --- a/docs/bus-implementation-checklist.md +++ b/docs/bus-implementation-checklist.md @@ -1,7 +1,7 @@ # Bus Implementation Checklist > **Status:** Historical checklist for the original independent `bus` -> implementation. `plan.md` is the current source of truth. The engine now lives +> implementation. [`docs/plan.md`](./plan.md) is the current source of truth. The engine now lives > in `internal/workflow`, root `queue` owns the application surface, and public > `bus` is a deprecated forwarding/raw-runtime compatibility facade. diff --git a/docs/events.md b/docs/events.md index a4a0c4a..873589b 100644 --- a/docs/events.md +++ b/docs/events.md @@ -12,9 +12,9 @@ This document defines the root application facade's unified observability contra Queue dispatch lifecycle: -- `EventDispatchStarted`: public dispatch began. +- `EventDispatchStarted`: a public dispatch began after job validation. Jobs rejected during validation emit no dispatch lifecycle facts. - `EventDispatchSucceeded`: public dispatch crossed the backend acceptance boundary. A synchronous handler can still return an application error after this fact. -- `EventDispatchFailed`: public dispatch failed before acceptance. +- `EventDispatchFailed`: a validated public dispatch failed before acceptance. - `EventEnqueueAccepted`: job accepted for dispatch. - `EventEnqueueRejected`: dispatch failed with error. - `EventEnqueueDuplicate`: dispatch rejected as duplicate (`UniqueFor`). @@ -26,21 +26,35 @@ Processing lifecycle: - `EventProcessSucceeded`: handler attempt succeeded. SQL, SQS, and RabbitMQ emit this only after durable row finalization, deletion, or acknowledgement respectively; backends without a post-handler settlement hook retain their documented weaker boundary. - `EventProcessFailed`: handler attempt returned an error or panicked. A panic is reported before the original panic value is rethrown so backend recovery and retry semantics remain unchanged. - `EventProcessRetried`: processing began for a numbered application retry attempt. Infrastructure redelivery of that same attempt may repeat the fact. -- `EventProcessArchived`: the driver confirmed terminal settlement for a failed attempt. +- `EventProcessArchived`: a driver confirmed terminal settlement of a failed attempt. Built-in support is listed below; unsupported drivers omit it rather than predicting a later backend transition. +- `EventProcessRecovered`: SQL bulk recovery requeued one stale in-flight claim. The driver emits one countable fact per affected row, but the bulk update does not load identity or correlation fields. - `EventRepublishFailed`: an internal delay or retry replacement could not be published. -- `EventSettlementFailed`: durable SQL finalization, broker acknowledgement, or broker deletion failed after handler or replacement work completed, so redelivery remains possible. +- `EventSettlementFailed`: delivery finalization, acknowledgement, deletion, or negative settlement failed or was ambiguous, so redelivery remains possible. This includes malformed and unregistered deliveries that fail settlement before handler execution. Queue control lifecycle: -- `EventQueuePaused`: queue consumption paused. -- `EventQueueResumed`: queue consumption resumed. +- `EventQueuePaused`: a supporting driver confirmed that queue consumption was paused. +- `EventQueueResumed`: a supporting driver confirmed that queue consumption was resumed. Workflow lifecycle: -- `EventJobStarted`, `EventJobSucceeded`, `EventJobFailed` -- `EventChainStarted`, `EventChainAdvanced`, `EventChainCompleted`, `EventChainFailed` -- `EventBatchStarted`, `EventBatchProgressed`, `EventBatchCompleted`, `EventBatchFailed`, `EventBatchCancelled` -- `EventCallbackStarted`, `EventCallbackSucceeded`, `EventCallbackFailed` +- `EventJobStarted`: a logical execution attempt began after envelope validation and before handler lookup. A missing handler therefore still has a started fact. +- `EventJobSucceeded`: logical job success committed. Settlement-aware drivers defer the fact until physical settlement succeeds. +- `EventJobFailed`: a logical job reached permanent or exhausted failure. +- `EventChainStarted`: a chain record was created and its initial dispatch began. +- `EventChainAdvanced`: a committed node outcome advanced the chain to its next node. +- `EventChainCompleted`: the final chain node committed terminal success. +- `EventChainFailed`: a chain committed terminal failure. +- `EventBatchStarted`: a batch record was created and its initial member dispatch began. +- `EventBatchProgressed`: a batch member committed a terminal outcome. +- `EventBatchCompleted`: the batch reached terminal success, including completion with allowed member failures. +- `EventBatchFailed`: an initial member dispatch rejection or a non-allowed member failure committed terminal batch failure. +- `EventBatchCancelled`: remaining batch work was cancelled after terminal failure. +- `EventCallbackStarted`: a claimed terminal Catch, Then, or Finally callback began execution after state validation. +- `EventCallbackSucceeded`: a terminal callback completed successfully across the applicable settlement boundary. +- `EventCallbackFailed`: a terminal callback was invalid or unavailable, returned an error, or panicked. + +Retryable callback store failures remain uncommitted and emit no terminal callback fact. Inline `Batch.Progress` closures do not emit `callback_*` facts. Positive job, chain, batch, and callback facts use the same SQL/SQS/RabbitMQ settlement boundary as `EventProcessSucceeded`. The SQL queue gives every processing claim an opaque generation ID. Same-attempt infrastructure redelivery normally retains inherited unsettled-generation provenance. When the current generation commits a receipt-backed workflow transition before later infrastructure work requests redelivery, the workflow engine marks application state committed and SQL retains that current generation instead. The signal selects the truthful receipt owner; it does not commit deferred facts or prove observer delivery. An application retry increments the attempt and clears the link, while recovery flags, aggregate state, and application error text do not supply equivalent authority. @@ -76,18 +90,20 @@ Processing events additionally include: - `MaxRetry` - `Duration` (for `Succeeded` and `Failed`) -Failure/cancel/reject events additionally include: +Failure, rejection, and enqueue-cancellation events additionally include: - `Err` +`EventBatchCancelled` omits `Err`; the adjacent `EventBatchFailed` fact carries the terminal cause. `EventProcessRecovered` is intentionally identity-free because SQL proves recovery with one fenced bulk update rather than a pre-update row read that could race ownership. + Every layer includes the applicable `DispatchID`, `JobID`, `ChainID`, and `BatchID` correlation fields when the delivery carries supported metadata. Queue and worker facts read the versioned direct-driver sidecar or decode a retained workflow envelope, so they can be joined to workflow facts without inspecting payloads in application observers. ## Semantics and guarantees -- Events are per-attempt, not aggregated. +- Physical processing and logical job-execution events are per-attempt. Dispatch and queue-control facts are operation-scoped, while chain and batch facts describe aggregate transitions. - Dispatch, enqueue, and queue-control events use `EventLayerQueue`; physical attempt events use `EventLayerWorker`; logical job, chain, batch, and callback transitions use `EventLayerWorkflow`. - `EventProcessRetried` is emitted when processing begins with `Attempt > 0`. It is intentionally not emitted merely because a handler returned an error, and consumers must tolerate a repeated fact when infrastructure redelivers the same numbered attempt. -- `EventProcessArchived` is reserved for a driver-confirmed terminal settlement; drivers that cannot yet confirm that boundary omit it rather than emitting a prediction. +- `EventProcessArchived` is reserved for a driver-confirmed terminal settlement. Drivers that cannot confirm that boundary omit it rather than emitting a prediction. - `JobKey` is a deterministic hash of the logical job type and payload. Volatile dispatch/workflow IDs are excluded, and the value is not guaranteed globally unique. - Correlated recoverable job successes and emitted positive chain or batch transition facts use a deterministic `EventID` for the same logical fact across settlement recovery. Failure EventIDs remain occurrence-based. Deterministic identity supports deduplication; it does not prove that an observer received the fact or make every event exactly-once. - `Queue` is the effective physical backend name carried by the dispatch. With a namespaced default such as `billing_default`, an explicit logical queue such as `critical` is reported as `billing_critical`. Jobs that omit a queue continue to report `default`; changing how `Config.DefaultQueue` routes empty targets is a separate targeting decision. Correlated queue, worker, workflow, aggregate, and callback facts always report the same name. @@ -106,6 +122,17 @@ Driver-specific capabilities: - Pause/resume control: currently supported by Sync, Workerpool, Redis. - Other drivers still emit collector-based events when `Observer` is configured. +Built-in `EventProcessArchived` support: + +| Runtime | Emits | Confirmed boundary | +| --- | --- | --- | +| SQL (SQLite, MySQL, PostgreSQL) | Yes | The fenced transition to `dead` committed. | +| SQS | No | `DeleteMessage` success cannot prove the receipt was still current or prevent rare standard-queue redelivery. | +| RabbitMQ | No | Consumer acknowledgement is one-way; channel loss can requeue a delivery before the broker processes it. | +| Redis | No | Asynq performs archival after the handler returns, outside the observer's confirmed boundary. | +| NATS | No | Core NATS has no durable terminal archive settlement to confirm. | +| Sync, Workerpool, Null | No | These runtimes expose no durable archive boundary. | + ## Observer behavior contract - Observers are best-effort telemetry hooks only; they must not control queue execution or implement workflow continuations. diff --git a/plan.md b/docs/plan.md similarity index 98% rename from plan.md rename to docs/plan.md index 56aafc5..793991c 100644 --- a/plan.md +++ b/docs/plan.md @@ -25,7 +25,7 @@ This is the living execution plan. Keep it current as work lands. A task is comp - Validate every affected Go module independently. Workspace success alone is insufficient. - Preserve intentional sibling `replace` directives used for repository testing. - Use `/tmp` for all test renders and generated application compositions. -- Keep this file focused on decisions, executable work, evidence, and remaining risk. Move lengthy design specifications into dedicated documents and link them here. +- Keep `docs/plan.md` focused on decisions, executable work, evidence, and remaining risk. Move lengthy design specifications into dedicated documents and link them here. ## North-Star Model @@ -533,11 +533,13 @@ Record accepted decisions here using the next stable ID. | DL-032 | 2026-07-20 | Linearize handler registration with worker activation. | Native runtimes install the current logical handler generation before backend activation. External runtimes publish and catch up the constructed worker before activation. A non-nil registration completed while startup is in flight is therefore present on the potentially consuming backend, while stable handler slots and generation-owned ledgers retain one physical registration across replacements and failed-start retries. Nil and empty registrations remain no-ops and cannot erase pending handlers. This corrects runtime behavior without changing source/API, configuration, persisted data, wire formats, operational rollout, or the minimum Go version. | | DL-033 | 2026-07-20 | Cache only successful workflow SQL schema initialization. | Schema DDL and MySQL key-limit discovery remain serialized, but caller cancellation, connectivity, locking, permission, and idempotent partial-DDL failures no longer poison the store instance permanently. A later operation retries the complete initialization sequence and publishes MySQL limits only after full success. Permanent failures are retried until operators repair the underlying condition, so repeated operations may produce repeated DDL or catalog attempts. This corrects runtime recovery behavior without changing source/API, configuration, persisted schema, wire formats, operational rollout, or the minimum Go version. | | DL-034 | 2026-07-20 | Upgrade pgx to the first release that fixes GO-2026-5004 without raising unrelated module baselines. | `driver/postgresqueue` and the integration module pin pgx v5.9.2 or newer through a checked dependency policy. Because pgx v5.9.2 requires Go 1.25, the PostgreSQL driver, examples, integration tooling, and repository workspace now require Go 1.25. The root library and every non-PostgreSQL published driver remain on Go 1.24.4. An exact per-module policy prevents accidental baseline drift, while the workspace must match the highest module version. This is a minimum-Go-version incompatibility for PostgreSQL driver consumers and repository contributors using the workspace. It does not change queue source/API, configuration shape, persisted data, schema, or wire formats. It does incorporate pgx v5.9 runtime changes, including reduced prepared-statement protocol traffic, discarding pooled connections left in transactions during reset, and defaulting an omitted database user to the current operating-system user. PostgreSQL consumers must build with Go 1.25 or newer and should set the database user explicitly when they do not want that upstream default; consumers of the root or other driver modules need no migration. | +| DL-035 | 2026-07-20 | Keep one exhaustive public event catalog and emit archive facts only from confirmed terminal settlement owners. | The canonical `queue.Event` stream retains queue, worker, and workflow layers because backend acceptance, physical execution, and logical workflow state can diverge, while one observer and one envelope carry all three. Package-wide contract tests require every exported `EventKind` and every private workflow kind cast into the public stream to appear exactly once in the README and detailed contract with the runtime layer. `process_archived` now publishes only after a fenced SQL transition to `dead`. SQS deletion and RabbitMQ consumer acknowledgement remain the documented available settlement boundaries for positive processing facts, but neither proves irreversible archival, so SQS, RabbitMQ, Redis, NATS, and local runtimes omit the stronger archive fact. SQL stale-claim recovery remains a fenced bulk update and emits one identity-free `process_recovered` count fact per affected row. Queue-control calls are serialized with their facts, and collector pause state is an idempotent gauge. This adds truthful observer volume for terminal SQL failures without changing source/API, configuration, persisted data, wire formats, minimum Go versions, or settlement behavior. Consumers counting archive facts should expect the newly confirmed SQL events after upgrade. | ## Progress Log ### 2026-07-20 +- Completed the unified event contract rather than leaving its reference as a partial list. README and detailed docs now cover all 32 exported kinds and their exact validation, settlement, callback, and support boundaries; an AST-backed parity test prevents catalog drift. SQL now emits `process_archived` only after a fenced transition to `dead`, while every built-in runtime that cannot prove irreversible archival omits it. Focused contracts cover archive presence and omission, and real SQL backend contracts prove confirmed archive delivery. - Upgraded the PostgreSQL driver and integration suite to pgx v5.9.2, the first release that fixes GO-2026-5004. Added exact mixed-version module policy, overflow-safe semantic dependency floors across every direct owner, replacement rejection, and manifest-derived minimum-toolchain CI that compiles tagged integration, generated examples, and documentation tooling. The root and unrelated driver modules remain on Go 1.24.4, while the PostgreSQL-only Go 1.25 consumer requirement is explicit. - Replaced one-shot workflow SQL schema failure caching with serialized success-only initialization. A deterministic SQLite regression starts the same store with a canceled first use, then proves a healthy retry completes migration and persists state; focused fix-only reviews verified MySQL limit publication, managed-schema continuity, idempotent DDL replay, and connection ownership. - Closed the handler-registration startup gap for native and external runtimes. Deterministic race-enabled tests cover new and replaced types during live startup, multiple late types, context-derived handles, concurrent start callers, external factory delay, failed-start retry, and shutdown latching while proving one physical registration per type. Corrected the production guide to name the actual SQL driver recovery settings and regenerated the executed unit count to 969. diff --git a/docs/readme/testcounts/integration_count.json b/docs/readme/testcounts/integration_count.json index d997d7b..c44d067 100644 --- a/docs/readme/testcounts/integration_count.json +++ b/docs/readme/testcounts/integration_count.json @@ -1,5 +1,5 @@ { "count": 631, - "source_hash": "sha256:886715b5a5184f7770cc7e723aa6cd424cd882478ea37d11f6b756fa696264c0", + "source_hash": "sha256:4b29b5cb51af6e4bb434cc27e03e5fda07775d3a35a371ec5cb4737779183041", "backend_scope": "all" } diff --git a/driver/sqlqueuecore/observability_contract_test.go b/driver/sqlqueuecore/observability_contract_test.go new file mode 100644 index 0000000..3c681f1 --- /dev/null +++ b/driver/sqlqueuecore/observability_contract_test.go @@ -0,0 +1,300 @@ +package sqlqueuecore + +import ( + "context" + "database/sql/driver" + "errors" + "strings" + "testing" + + "github.com/goforj/queue" + "github.com/goforj/queue/busruntime" +) + +// TestRecoverStaleProcessingEmitsCountFacts verifies bulk recovery reports one +// normalized, identity-free fact for every row the fenced update changed. +func TestRecoverStaleProcessingEmitsCountFacts(t *testing.T) { + connection := &databaseConnStub{ + exec: func(context.Context, string, []driver.NamedValue) (driver.Result, error) { + return driver.RowsAffected(2), nil + }, + } + db := newDatabaseStub(connection) + t.Cleanup(func() { + if err := db.Close(); err != nil { + t.Errorf("close database stub: %v", err) + } + }) + + var events []queue.Event + database := &databaseQueue{ + db: db, + cfg: localDatabaseConfig{ + DriverName: "sqlite", + DefaultQueue: "default", + }, + observer: queue.ObserverFunc(func(_ context.Context, event queue.Event) { + events = append(events, event) + }), + } + if err := database.recoverStaleProcessing(context.Background(), 1_000); err != nil { + t.Fatalf("recover stale processing: %v", err) + } + + if len(events) != 2 { + t.Fatalf("recovery events = %+v, want one fact per recovered row", events) + } + seenIDs := make(map[string]struct{}, len(events)) + for index, event := range events { + if event.Kind != queue.EventProcessRecovered || event.Layer != queue.EventLayerWorker || event.Driver != queue.DriverDatabase { + t.Errorf("recovery event %d = %+v, want normalized database worker fact", index, event) + } + if event.SchemaVersion == 0 || event.EventID == "" || event.Time.IsZero() { + t.Errorf("recovery event %d has incomplete envelope metadata: %+v", index, event) + } + if event.Queue != "" || event.JobType != "" || event.JobKey != "" || event.DispatchID != "" || event.JobID != "" || event.ChainID != "" || event.BatchID != "" { + t.Errorf("recovery event %d invented identity unavailable from the bulk update: %+v", index, event) + } + if _, duplicate := seenIDs[event.EventID]; duplicate { + t.Errorf("recovery event %d reused event ID %q", index, event.EventID) + } + seenIDs[event.EventID] = struct{}{} + } +} + +// TestRecoverStaleProcessingRemainsSilentWithoutConfirmedRows verifies the +// recovery observer never invents a fact when execution fails or row evidence +// is absent or unavailable. +func TestRecoverStaleProcessingRemainsSilentWithoutConfirmedRows(t *testing.T) { + execErr := errors.New("recovery update failed") + rowsErr := errors.New("recovery row count unavailable") + tests := []struct { + name string + result driver.Result + execErr error + wantErr error + }{ + {name: "no stale rows", result: driver.RowsAffected(0)}, + {name: "execution failure", execErr: execErr, wantErr: execErr}, + {name: "row count unavailable", result: databaseResultStub{err: rowsErr}}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + connection := &databaseConnStub{ + exec: func(context.Context, string, []driver.NamedValue) (driver.Result, error) { + return test.result, test.execErr + }, + } + db := newDatabaseStub(connection) + t.Cleanup(func() { + if err := db.Close(); err != nil { + t.Errorf("close database stub: %v", err) + } + }) + + var events []queue.Event + database := &databaseQueue{ + db: db, + cfg: localDatabaseConfig{ + DriverName: "sqlite", + DefaultQueue: "default", + }, + observer: queue.ObserverFunc(func(_ context.Context, event queue.Event) { + events = append(events, event) + }), + } + err := database.recoverStaleProcessing(context.Background(), 1_000) + if test.wantErr != nil && !errors.Is(err, test.wantErr) { + t.Fatalf("recover stale processing error = %v, want %v", err, test.wantErr) + } + if test.wantErr == nil && err != nil { + t.Fatalf("recover stale processing error = %v, want nil", err) + } + if len(events) != 0 { + t.Fatalf("unconfirmed recovery emitted facts: %+v", events) + } + }) + } +} + +// TestDatabaseProcessArchivedRequiresConfirmedTerminalState verifies the SQL +// driver emits archive facts only after a fenced dead-state transition succeeds. +func TestDatabaseProcessArchivedRequiresConfirmedTerminalState(t *testing.T) { + tests := []struct { + name string + attempt int + maxRetry int + register bool + wantArchive bool + wantState string + }{ + {name: "exhausted handler failure", attempt: 2, maxRetry: 2, register: true, wantArchive: true, wantState: "state='dead'"}, + {name: "permanent handler failure", maxRetry: 3, register: true, wantArchive: true, wantState: "state='dead'"}, + {name: "missing handler terminal failure", wantArchive: true, wantState: "state='dead'"}, + {name: "retry remains pending", maxRetry: 1, register: true, wantState: "state='pending'"}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var settlementQuery string + connection := &databaseConnStub{ + exec: func(_ context.Context, query string, _ []driver.NamedValue) (driver.Result, error) { + settlementQuery = query + return driver.RowsAffected(1), nil + }, + } + db := newDatabaseStub(connection) + t.Cleanup(func() { + if err := db.Close(); err != nil { + t.Errorf("close database stub: %v", err) + } + }) + + var events []queue.Event + database := &databaseQueue{ + db: db, + cfg: localDatabaseConfig{DriverName: "sqlite", DefaultQueue: "default"}, + handlers: make(map[string]queue.Handler), + continuation: busruntime.NewContinuationScope(), + observer: queue.ObserverFunc(func(_ context.Context, event queue.Event) { + events = append(events, event) + }), + } + if test.register { + database.handlers["bus:job"] = func(context.Context, queue.Job) error { + if test.name == "permanent handler failure" { + return busruntime.Permanent(errors.New("invalid report")) + } + return errors.New("report failed") + } + } + payload := []byte(`{"schema_version":1,"dispatch_id":"dsp_sql_archive","job_id":"job_sql_archive","job":{"type":"reports:build","payload":"eyJpZCI6MX0="}}`) + database.processJob(&dbJob{ + id: 42, + processingToken: "owned-generation", + queueName: "critical", + jobType: "bus:job", + payload: payload, + attempt: test.attempt, + maxRetry: test.maxRetry, + }) + + if !strings.Contains(settlementQuery, test.wantState) { + t.Fatalf("settlement query = %q, want %q transition", settlementQuery, test.wantState) + } + if !test.wantArchive { + if len(events) != 0 { + t.Fatalf("retryable settlement emitted archive facts: %+v", events) + } + return + } + if len(events) != 1 { + t.Fatalf("archive events = %+v, want exactly one", events) + } + archive := events[0] + if archive.Kind != queue.EventProcessArchived || archive.Layer != queue.EventLayerWorker || archive.Driver != queue.DriverDatabase { + t.Fatalf("archive event = %+v, want normalized database worker fact", archive) + } + if archive.Queue != "critical" || archive.JobType != "reports:build" || archive.DispatchID != "dsp_sql_archive" || archive.JobID != "job_sql_archive" { + t.Fatalf("archive correlation = %+v", archive) + } + if archive.Attempt != test.attempt || archive.MaxRetry != test.maxRetry || archive.Err == nil || archive.SchemaVersion == 0 || archive.EventID == "" || archive.Time.IsZero() { + t.Fatalf("archive attempt or envelope metadata = %+v", archive) + } + }) + } +} + +// TestDatabaseTerminalSettlementFailureNeverArchives verifies an archive fact +// requires positive evidence that exactly one fenced row reached dead state. +func TestDatabaseTerminalSettlementFailureNeverArchives(t *testing.T) { + execErr := errors.New("terminal update failed") + rowsErr := errors.New("terminal row count unavailable") + settlements := []struct { + name string + result driver.Result + execErr error + wantErr error + wantDetail string + }{ + {name: "execution failure", execErr: execErr, wantErr: execErr}, + {name: "lost fence", result: driver.RowsAffected(0), wantDetail: "affected 0 rows, want 1"}, + {name: "ambiguous update", result: driver.RowsAffected(2), wantDetail: "affected 2 rows, want 1"}, + {name: "row count unavailable", result: databaseResultStub{err: rowsErr}, wantErr: rowsErr}, + } + deliveries := []struct { + name string + register bool + attempt int + maxRetry int + }{ + {name: "exhausted handler failure", register: true, attempt: 2, maxRetry: 2}, + {name: "missing handler", attempt: 0, maxRetry: 0}, + } + for _, delivery := range deliveries { + for _, settlement := range settlements { + t.Run(delivery.name+"/"+settlement.name, func(t *testing.T) { + var settlementQueries []string + connection := &databaseConnStub{ + exec: func(_ context.Context, query string, _ []driver.NamedValue) (driver.Result, error) { + settlementQueries = append(settlementQueries, query) + return settlement.result, settlement.execErr + }, + } + db := newDatabaseStub(connection) + t.Cleanup(func() { + if err := db.Close(); err != nil { + t.Errorf("close database stub: %v", err) + } + }) + + var events []queue.Event + database := &databaseQueue{ + db: db, + cfg: localDatabaseConfig{DriverName: "sqlite", DefaultQueue: "default"}, + handlers: make(map[string]queue.Handler), + continuation: busruntime.NewContinuationScope(), + observer: queue.ObserverFunc(func(_ context.Context, event queue.Event) { + events = append(events, event) + }), + } + if delivery.register { + database.handlers["bus:job"] = func(context.Context, queue.Job) error { + return errors.New("report failed") + } + } + payload := []byte(`{"schema_version":1,"dispatch_id":"dsp_sql_failed_archive","job_id":"job_sql_failed_archive","job":{"type":"reports:build","payload":"eyJpZCI6MX0="}}`) + database.processJob(&dbJob{ + id: 42, + processingToken: "owned-generation", + queueName: "critical", + jobType: "bus:job", + payload: payload, + attempt: delivery.attempt, + maxRetry: delivery.maxRetry, + }) + + if len(settlementQueries) != databaseFinalizeRetryCount { + t.Fatalf("terminal settlement attempts = %d, want %d", len(settlementQueries), databaseFinalizeRetryCount) + } + for index, query := range settlementQueries { + if !strings.Contains(query, "state='dead'") { + t.Fatalf("settlement query %d = %q, want terminal dead-state update", index, query) + } + } + if len(events) != 1 { + t.Fatalf("terminal settlement failure events = %+v, want exactly one", events) + } + event := events[0] + if event.Kind != queue.EventSettlementFailed { + t.Fatalf("terminal settlement event = %+v, want settlement_failed without process_archived", event) + } + if settlement.wantErr != nil && !errors.Is(event.Err, settlement.wantErr) { + t.Fatalf("terminal settlement error = %v, want wrapped %v", event.Err, settlement.wantErr) + } + if settlement.wantDetail != "" && !strings.Contains(event.Err.Error(), settlement.wantDetail) { + t.Fatalf("terminal settlement error = %v, want %q", event.Err, settlement.wantDetail) + } + }) + } + } +} diff --git a/driver/sqlqueuecore/queue_database_impl.go b/driver/sqlqueuecore/queue_database_impl.go index 2b0a707..67ba2fe 100644 --- a/driver/sqlqueuecore/queue_database_impl.go +++ b/driver/sqlqueuecore/queue_database_impl.go @@ -729,8 +729,12 @@ func (d *databaseQueue) workerLoop() { func (d *databaseQueue) processJob(job *dbJob) { handler, ok := d.lookup(job.jobType) if !ok { - if err := d.markFailedWithRetry(job, fmt.Errorf("no handler registered for job type %q", job.jobType)); err != nil { + runErr := fmt.Errorf("no handler registered for job type %q", job.jobType) + if err := d.markFailedWithRetry(job, runErr); err != nil { d.handleSettlementFailure(context.Background(), job, err) + } else { + // Missing handlers never cross the process wrapper, so the driver owns any terminal archive fact. + d.observeConfirmedProcessArchive(context.Background(), job, runErr) } return } @@ -757,6 +761,9 @@ func (d *databaseQueue) processJob(job *dbJob) { return } settlement.Commit() + if err != nil { + d.observeConfirmedProcessArchive(ctx, job, err) + } } // handleSettlementFailure preserves inherited recovery lineage before reporting @@ -1235,6 +1242,33 @@ func (d *databaseQueue) observeSettlementFailure(ctx context.Context, job *dbJob }) } +// observeConfirmedProcessArchive publishes terminal failure only after the +// caller has durably fenced the owned row into its dead state. +func (d *databaseQueue) observeConfirmedProcessArchive(ctx context.Context, job *dbJob, runErr error) { + if job == nil || busruntime.ClassifyAttempt(busruntime.DeliveryAttempt{ + Number: job.attempt, + MaxRetry: job.maxRetry, + }, runErr) != busruntime.AttemptFailed { + return + } + metadata := queue.ResolveObservedJobMetadataFromJob(databaseDeliveryJob(job)) + queuecore.SafeObserve(ctx, d.observer, queue.Event{ + Kind: queue.EventProcessArchived, + Driver: queue.DriverDatabase, + Queue: queuecore.NormalizeQueueName(job.queueName), + JobType: metadata.JobType, + JobKey: metadata.JobKey, + DispatchID: metadata.DispatchID, + JobID: metadata.JobID, + ChainID: metadata.ChainID, + BatchID: metadata.BatchID, + Attempt: job.attempt, + MaxRetry: job.maxRetry, + Err: runErr, + Time: time.Now(), + }) +} + // classifyDatabaseFailure derives the durable state transition from the physical attempt and handler result. func classifyDatabaseFailure(job *dbJob, runErr error, now int64) (databaseFailureSettlement, error) { decision := busruntime.ClassifyAttempt(busruntime.DeliveryAttempt{ diff --git a/driver/sqsqueue/worker_sqs_impl_test.go b/driver/sqsqueue/worker_sqs_impl_test.go index 6d500b7..d0e7ecb 100644 --- a/driver/sqsqueue/worker_sqs_impl_test.go +++ b/driver/sqsqueue/worker_sqs_impl_test.go @@ -513,6 +513,7 @@ func TestSQSWorker_ProcessFailureRetryAndTerminal(t *testing.T) { func TestSQSWorker_AttemptDecisionSettlement(t *testing.T) { t.Run("permanent failure deletes without republishing", func(t *testing.T) { stub := &sqsWorkerClientStub{} + var events []queue.Event w := &sqsWorker{ handlers: map[string]queue.Handler{ "job:permanent": func(ctx context.Context, _ queue.Job) error { @@ -525,6 +526,9 @@ func TestSQSWorker_AttemptDecisionSettlement(t *testing.T) { }, client: stub, queueURL: "https://example.local/queue/default", + observer: queue.ObserverFunc(func(_ context.Context, event queue.Event) { + events = append(events, event) + }), } body, err := json.Marshal(sqsMessage{Type: "job:permanent", Queue: "default", MaxRetry: 3}) if err != nil { @@ -539,6 +543,11 @@ func TestSQSWorker_AttemptDecisionSettlement(t *testing.T) { if len(stub.deleteInputs) != 1 { t.Fatalf("permanent failure must delete its receipt, got %d deletes", len(stub.deleteInputs)) } + for _, event := range events { + if event.Kind == queue.EventProcessArchived { + t.Fatalf("SQS deletion cannot prove terminal archival, got event %+v", event) + } + } }) t.Run("uncommitted failure leaves the original receipt", func(t *testing.T) { diff --git a/event_contract_audit_test.go b/event_contract_audit_test.go new file mode 100644 index 0000000..197e001 --- /dev/null +++ b/event_contract_audit_test.go @@ -0,0 +1,232 @@ +package queue + +import ( + "context" + "errors" + "path/filepath" + "reflect" + "testing" + "time" + + "github.com/goforj/queue/internal/workflow" +) + +// TestInternalWorkflowKindsMapToPublicCatalog prevents the private engine from +// casting an undocumented kind or an incorrectly layered fact into queue.Event. +func TestInternalWorkflowKindsMapToPublicCatalog(t *testing.T) { + t.Parallel() + + root := eventReferenceRoot(t) + publicDefinitions := parsePackageEventKindDefinitions(t, root) + publicByKind := make(map[EventKind]string, len(publicDefinitions)) + for _, definition := range publicDefinitions { + publicByKind[definition.kind] = definition.name + } + + workflowDefinitions := parsePackageEventKindDefinitions(t, filepath.Join(root, "internal", "workflow")) + for _, definition := range workflowDefinitions { + publicName, ok := publicByKind[definition.kind] + if !ok { + t.Errorf("internal workflow kind %s (%q) has no public EventKind", definition.name, definition.kind) + continue + } + if publicName != definition.name { + t.Errorf("internal workflow kind %s (%q) maps to public identifier %s", definition.name, definition.kind, publicName) + } + wantLayer := EventLayerWorkflow + switch definition.kind { + case EventDispatchStarted, EventDispatchSucceeded, EventDispatchFailed: + wantLayer = EventLayerQueue + } + if gotLayer := eventLayerForKind(definition.kind); gotLayer != wantLayer { + t.Errorf("internal workflow kind %s (%q) maps to layer %q, want %q", definition.name, definition.kind, gotLayer, wantLayer) + } + } +} + +// TestWorkflowObserverAdapterTranslatesEveryField verifies the internal engine +// cannot lose correlation or application identity at the unified public boundary. +func TestWorkflowObserverAdapterTranslatesEveryField(t *testing.T) { + assertWorkflowEventShape(t) + + type contextKey struct{} + ctx := context.WithValue(context.Background(), contextKey{}, "observer-context") + wantErr := errors.New("batch failed") + wantTime := time.Date(2026, time.July, 20, 12, 30, 0, 0, time.UTC) + internalEvent := workflow.Event{ + SchemaVersion: 7, + EventID: "evt_adapter_contract", + Kind: workflow.EventBatchFailed, + DispatchID: "dsp_adapter_contract", + JobID: "job_adapter_contract", + ChainID: "chn_adapter_contract", + BatchID: "bat_adapter_contract", + Attempt: 3, + JobType: "reports:build", + JobKey: "job-key-adapter-contract", + Queue: "critical", + Duration: 42 * time.Millisecond, + Time: wantTime, + Err: wantErr, + } + + var ( + gotContext context.Context + gotEvent Event + ) + adapter := workflowObserverAdapter{ + driver: DriverSQS, + resolveQueueName: func(queueName string) string { + return "billing_" + queueName + }, + observer: ObserverFunc(func(observedContext context.Context, event Event) { + gotContext = observedContext + gotEvent = event + }), + } + adapter.Observe(ctx, internalEvent) + + wantEvent := Event{ + SchemaVersion: internalEvent.SchemaVersion, + EventID: internalEvent.EventID, + Layer: EventLayerWorkflow, + Kind: EventBatchFailed, + Driver: DriverSQS, + Queue: "billing_critical", + JobType: internalEvent.JobType, + JobKey: internalEvent.JobKey, + DispatchID: internalEvent.DispatchID, + JobID: internalEvent.JobID, + ChainID: internalEvent.ChainID, + BatchID: internalEvent.BatchID, + Attempt: internalEvent.Attempt, + Duration: internalEvent.Duration, + Err: internalEvent.Err, + Time: internalEvent.Time, + } + if !reflect.DeepEqual(gotEvent, wantEvent) { + t.Fatalf("translated workflow event = %+v, want %+v", gotEvent, wantEvent) + } + if gotContext == nil || gotContext.Value(contextKey{}) != "observer-context" { + t.Fatalf("observer context = %v, want original workflow context", gotContext) + } +} + +// assertWorkflowEventShape pins the internal adapter input so adding a field +// requires an explicit decision about its public projection. +func assertWorkflowEventShape(t *testing.T) { + t.Helper() + wantFields := []string{ + "SchemaVersion", + "EventID", + "Kind", + "DispatchID", + "JobID", + "ChainID", + "BatchID", + "Attempt", + "JobType", + "JobKey", + "Queue", + "Duration", + "Time", + "Err", + } + eventType := reflect.TypeOf(workflow.Event{}) + if eventType.NumField() != len(wantFields) { + t.Fatalf("workflow.Event fields = %d, want %d; update the public observer adapter contract", eventType.NumField(), len(wantFields)) + } + for index, wantField := range wantFields { + if gotField := eventType.Field(index).Name; gotField != wantField { + t.Fatalf("workflow.Event field %d = %s, want %s; update the public observer adapter contract", index, gotField, wantField) + } + } +} + +// TestEventLayerForKindCoversPublicKinds pins every currently exported event +// kind to one semantic layer so future edits cannot silently conflate scopes. +func TestEventLayerForKindCoversPublicKinds(t *testing.T) { + kindsByLayer := map[EventLayer][]EventKind{ + EventLayerQueue: { + EventDispatchStarted, + EventDispatchSucceeded, + EventDispatchFailed, + EventEnqueueAccepted, + EventEnqueueRejected, + EventEnqueueDuplicate, + EventEnqueueCanceled, + EventQueuePaused, + EventQueueResumed, + }, + EventLayerWorker: { + EventProcessStarted, + EventProcessSucceeded, + EventProcessFailed, + EventProcessRetried, + EventProcessArchived, + EventProcessRecovered, + EventRepublishFailed, + EventSettlementFailed, + }, + EventLayerWorkflow: { + EventJobStarted, + EventJobSucceeded, + EventJobFailed, + EventChainStarted, + EventChainAdvanced, + EventChainCompleted, + EventChainFailed, + EventBatchStarted, + EventBatchProgressed, + EventBatchCompleted, + EventBatchFailed, + EventBatchCancelled, + EventCallbackStarted, + EventCallbackSucceeded, + EventCallbackFailed, + }, + } + + seen := make(map[EventKind]EventLayer) + for wantLayer, kinds := range kindsByLayer { + for _, kind := range kinds { + if priorLayer, exists := seen[kind]; exists { + t.Fatalf("event kind %q appears in both %q and %q", kind, priorLayer, wantLayer) + } + seen[kind] = wantLayer + if gotLayer := eventLayerForKind(kind); gotLayer != wantLayer { + t.Errorf("event kind %q layer = %q, want %q", kind, gotLayer, wantLayer) + } + } + } + if len(seen) != 32 { + t.Fatalf("covered event kinds = %d, want 32", len(seen)) + } + if gotLayer := eventLayerForKind(EventKind("future_queue_fact")); gotLayer != EventLayerQueue { + t.Fatalf("unknown event layer = %q, want queue compatibility default", gotLayer) + } +} + +// TestStatsCollectorIngestsReservedProcessArchived pins compatibility for +// external drivers that already emit a confirmed terminal-settlement fact. +func TestStatsCollectorIngestsReservedProcessArchived(t *testing.T) { + collector := NewStatsCollector() + collector.Observe(context.Background(), Event{ + Kind: EventProcessArchived, + Layer: EventLayerWorker, + Driver: DriverDatabase, + Queue: "critical", + Time: time.Date(2026, time.July, 20, 13, 0, 0, 0, time.UTC), + }) + + counters, ok := collector.Snapshot().Queue("critical") + if !ok { + t.Fatal("reserved archive fact did not create queue counters") + } + if counters.Archived != 1 { + t.Fatalf("archived count = %d, want 1", counters.Archived) + } + if counters.Pending != 0 || counters.Active != 0 || counters.Retry != 0 || counters.Processed != 0 || counters.Failed != 0 { + t.Fatalf("archive fact mutated unrelated counters: %+v", counters) + } +} diff --git a/event_reference_test.go b/event_reference_test.go new file mode 100644 index 0000000..a1db296 --- /dev/null +++ b/event_reference_test.go @@ -0,0 +1,533 @@ +package queue + +import ( + "bufio" + "fmt" + "go/ast" + "go/parser" + "go/token" + "os" + "path/filepath" + "regexp" + "runtime" + "strconv" + "strings" + "testing" +) + +type eventKindDefinition struct { + name string + kind EventKind +} + +type documentedEventKind struct { + layer EventLayer + line int +} + +// TestEventReferenceMatchesExportedKinds prevents the public event catalog, +// layer mapping, and human-facing references from drifting independently. +func TestEventReferenceMatchesExportedKinds(t *testing.T) { + t.Parallel() + + root := eventReferenceRoot(t) + definitions := parsePackageEventKindDefinitions(t, root) + readmeKinds := parseReadmeEventReference(t, filepath.Join(root, "README.md")) + eventsDocKinds := parseEventsDocEventReference(t, filepath.Join(root, "docs", "events.md")) + + if len(readmeKinds) != len(definitions) { + t.Errorf("README event rows = %d, exported EventKind constants = %d", len(readmeKinds), len(definitions)) + } + if len(eventsDocKinds) != len(definitions) { + t.Errorf("docs/events.md event identifiers = %d, exported EventKind constants = %d", len(eventsDocKinds), len(definitions)) + } + for _, definition := range definitions { + documented, ok := readmeKinds[definition.kind] + if !ok { + t.Errorf("README Events reference is missing %s (%q)", definition.name, definition.kind) + } else if got := eventLayerForKind(definition.kind); got != documented.layer { + t.Errorf("README line %d assigns %s (%q) to layer %q, runtime maps it to %q", documented.line, definition.name, definition.kind, documented.layer, got) + } + if _, ok := eventsDocKinds[definition.name]; !ok { + t.Errorf("docs/events.md Event kinds section is missing %s", definition.name) + } + } + + exported := make(map[EventKind]string, len(definitions)) + for _, definition := range definitions { + exported[definition.kind] = definition.name + } + for kind, documented := range readmeKinds { + if _, ok := exported[kind]; !ok { + t.Errorf("README line %d documents unknown EventKind %q", documented.line, kind) + } + } + exportedNames := exportedEventKindNames(definitions) + for name, line := range eventsDocKinds { + if _, ok := exportedNames[name]; !ok { + t.Errorf("docs/events.md line %d documents unknown EventKind identifier %s", line, name) + } + } +} + +// TestParseEventKindDefinitionsFindsInferredConstants prevents syntactically +// valid typed constants from bypassing the documentation parity contract. +func TestParseEventKindDefinitionsFindsInferredConstants(t *testing.T) { + t.Parallel() + + filename := filepath.Join(t.TempDir(), "event_kinds.go") + source := `package fixture + +type EventKind string + +const ( + EventExplicit EventKind = "explicit" + EventConverted = EventKind("converted") + eventInheritedBase EventKind = "inherited" + EventInherited + EventComposed = EventConverted + "_composed" +) +` + if err := os.WriteFile(filename, []byte(source), 0o600); err != nil { + t.Fatalf("write EventKind parser fixture: %v", err) + } + + directory := filepath.Dir(filename) + additionalSource := `package fixture + +type EventAlias = EventKind + +const ( + EventSeparate EventAlias = "separate" + EventCrossFile = EventConverted + "_cross_file" +) +` + if err := os.WriteFile(filepath.Join(directory, "event_kinds_additional.go"), []byte(additionalSource), 0o600); err != nil { + t.Fatalf("write additional EventKind parser fixture: %v", err) + } + + definitions := parsePackageEventKindDefinitions(t, directory) + want := map[string]EventKind{ + "EventExplicit": "explicit", + "EventConverted": "converted", + "EventInherited": "inherited", + "EventComposed": "converted_composed", + "EventSeparate": "separate", + "EventCrossFile": "converted_cross_file", + } + if len(definitions) != len(want) { + t.Fatalf("EventKind definitions = %d, want %d: %+v", len(definitions), len(want), definitions) + } + for _, definition := range definitions { + if wantKind, ok := want[definition.name]; !ok || definition.kind != wantKind { + t.Errorf("EventKind definition %s = %q, want %q", definition.name, definition.kind, wantKind) + } + } +} + +// parsePackageEventKindDefinitions discovers EventKind constants across every +// production file so moving or adding a declaration cannot bypass the catalog. +func parsePackageEventKindDefinitions(t *testing.T, directory string) []eventKindDefinition { + t.Helper() + entries, err := os.ReadDir(directory) + if err != nil { + t.Fatalf("read package directory %s: %v", directory, err) + } + + filenames := make([]string, 0, len(entries)) + for _, entry := range entries { + if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".go") || strings.HasSuffix(entry.Name(), "_test.go") { + continue + } + filenames = append(filenames, filepath.Join(directory, entry.Name())) + } + + kindTypes := parseEventKindTypeAliases(t, filenames) + knownKinds := make(map[string]EventKind) + definitionsByName := make(map[string]eventKindDefinition) + seenKinds := make(map[EventKind]string) + for { + knownBefore := len(knownKinds) + for _, filename := range filenames { + for _, definition := range parseEventKindDefinitionsFromFile(t, filename, kindTypes, knownKinds) { + if previous, duplicate := definitionsByName[definition.name]; duplicate { + if previous.kind != definition.kind { + t.Fatalf("EventKind name %s resolves to both %q and %q", definition.name, previous.kind, definition.kind) + } + continue + } + if previous, duplicate := seenKinds[definition.kind]; duplicate { + t.Fatalf("EventKind value %q is shared by %s and %s across package files", definition.kind, previous, definition.name) + } + definitionsByName[definition.name] = definition + seenKinds[definition.kind] = definition.name + } + } + if len(knownKinds) == knownBefore { + break + } + } + if len(definitionsByName) == 0 { + t.Fatalf("%s contains no typed EventKind constants", directory) + } + definitions := make([]eventKindDefinition, 0, len(definitionsByName)) + for _, definition := range definitionsByName { + definitions = append(definitions, definition) + } + return definitions +} + +// parseEventKindTypeAliases resolves aliases of EventKind across production +// files so a renamed spelling retains the same catalog obligation. +func parseEventKindTypeAliases(t *testing.T, filenames []string) map[string]struct{} { + t.Helper() + types := map[string]struct{}{"EventKind": {}} + for { + countBefore := len(types) + for _, filename := range filenames { + parsed, err := parser.ParseFile(token.NewFileSet(), filename, nil, 0) + if err != nil { + t.Fatalf("parse %s: %v", filename, err) + } + for _, declaration := range parsed.Decls { + group, ok := declaration.(*ast.GenDecl) + if !ok || group.Tok != token.TYPE { + continue + } + for _, specification := range group.Specs { + alias, ok := specification.(*ast.TypeSpec) + if !ok || !alias.Assign.IsValid() || !isEventKindType(alias.Type, types) { + continue + } + types[alias.Name.Name] = struct{}{} + } + } + } + if len(types) == countBefore { + return types + } + } +} + +// eventReferenceRoot resolves documentation relative to this test file so the +// contract remains stable when tests are launched from another directory. +func eventReferenceRoot(t *testing.T) string { + t.Helper() + _, filename, _, ok := runtime.Caller(0) + if !ok { + t.Fatal("resolve event reference test path") + } + return filepath.Dir(filename) +} + +// parseEventKindDefinitionsFromFile resolves the EventKind constants declared +// by one source file and permits files that contain no event declarations. +func parseEventKindDefinitionsFromFile(t *testing.T, filename string, kindTypes map[string]struct{}, knownKinds map[string]EventKind) []eventKindDefinition { + t.Helper() + parsed, err := parser.ParseFile(token.NewFileSet(), filename, nil, 0) + if err != nil { + t.Fatalf("parse %s: %v", filename, err) + } + + definitions := make([]eventKindDefinition, 0) + seenNames := make(map[string]struct{}) + seenKinds := make(map[EventKind]string) + for _, declaration := range parsed.Decls { + group, ok := declaration.(*ast.GenDecl) + if !ok || group.Tok != token.CONST { + continue + } + var inheritedType ast.Expr + var inheritedValues []ast.Expr + for _, specification := range group.Specs { + values, ok := specification.(*ast.ValueSpec) + if !ok { + continue + } + effectiveType := values.Type + effectiveValues := values.Values + if len(effectiveValues) == 0 { + effectiveType = inheritedType + effectiveValues = inheritedValues + } else { + inheritedType = effectiveType + inheritedValues = effectiveValues + } + if len(effectiveValues) != len(values.Names) { + if isEventKindType(effectiveType, kindTypes) || expressionsReferenceEventKind(effectiveValues, kindTypes, knownKinds) { + t.Fatalf("%s EventKind declaration must pair every name with one value", filename) + } + continue + } + for index, nameIdentifier := range values.Names { + kind, inferred, err := eventKindExpressionValue(effectiveValues[index], kindTypes, knownKinds) + if err != nil { + t.Fatalf("resolve %s value: %v", nameIdentifier.Name, err) + } + if isEventKindType(effectiveType, kindTypes) && !inferred { + kind, err = explicitEventKindValue(effectiveValues[index], kindTypes, knownKinds) + if err != nil { + continue + } + inferred = true + } + if !inferred { + continue + } + name := nameIdentifier.Name + knownKinds[name] = kind + if !ast.IsExported(name) { + continue + } + if previous, duplicate := seenKinds[kind]; duplicate { + t.Fatalf("EventKind value %q is shared by %s and %s", kind, previous, name) + } + if _, duplicate := seenNames[name]; duplicate { + t.Fatalf("EventKind name %s is duplicated", name) + } + seenKinds[kind] = name + seenNames[name] = struct{}{} + definitions = append(definitions, eventKindDefinition{name: name, kind: kind}) + } + } + } + return definitions +} + +// expressionsReferenceEventKind reports whether any expression has the public +// event kind type through conversion, inheritance, or composition. +func expressionsReferenceEventKind(expressions []ast.Expr, kindTypes map[string]struct{}, knownKinds map[string]EventKind) bool { + for _, expression := range expressions { + if _, ok, _ := eventKindExpressionValue(expression, kindTypes, knownKinds); ok { + return true + } + } + return false +} + +// eventKindExpressionValue resolves expressions whose inferred constant type is EventKind. +func eventKindExpressionValue(expression ast.Expr, kindTypes map[string]struct{}, knownKinds map[string]EventKind) (EventKind, bool, error) { + switch typed := expression.(type) { + case *ast.CallExpr: + if !isEventKindType(typed.Fun, kindTypes) || len(typed.Args) != 1 { + return "", false, nil + } + value, err := stringConstantValue(typed.Args[0], kindTypes, knownKinds) + return EventKind(value), true, err + case *ast.Ident: + value, ok := knownKinds[typed.Name] + return value, ok, nil + case *ast.ParenExpr: + return eventKindExpressionValue(typed.X, kindTypes, knownKinds) + case *ast.BinaryExpr: + if typed.Op != token.ADD { + return "", false, nil + } + _, leftTyped, err := eventKindExpressionValue(typed.X, kindTypes, knownKinds) + if err != nil { + return "", false, err + } + _, rightTyped, err := eventKindExpressionValue(typed.Y, kindTypes, knownKinds) + if err != nil { + return "", false, err + } + if !leftTyped && !rightTyped { + return "", false, nil + } + left, err := stringConstantValue(typed.X, kindTypes, knownKinds) + if err != nil { + return "", false, err + } + right, err := stringConstantValue(typed.Y, kindTypes, knownKinds) + if err != nil { + return "", false, err + } + return EventKind(left + right), true, nil + default: + return "", false, nil + } +} + +// explicitEventKindValue resolves the string value of an explicitly typed EventKind constant. +func explicitEventKindValue(expression ast.Expr, kindTypes map[string]struct{}, knownKinds map[string]EventKind) (EventKind, error) { + value, err := stringConstantValue(expression, kindTypes, knownKinds) + return EventKind(value), err +} + +// stringConstantValue resolves the string forms permitted by the EventKind catalog. +func stringConstantValue(expression ast.Expr, kindTypes map[string]struct{}, knownKinds map[string]EventKind) (string, error) { + switch typed := expression.(type) { + case *ast.BasicLit: + if typed.Kind != token.STRING { + return "", fmt.Errorf("expression is not a string constant") + } + value, err := strconv.Unquote(typed.Value) + if err != nil { + return "", fmt.Errorf("unquote %s: %w", typed.Value, err) + } + return value, nil + case *ast.Ident: + value, ok := knownKinds[typed.Name] + if !ok { + return "", fmt.Errorf("identifier %s is not a known EventKind constant", typed.Name) + } + return string(value), nil + case *ast.ParenExpr: + return stringConstantValue(typed.X, kindTypes, knownKinds) + case *ast.CallExpr: + if !isEventKindType(typed.Fun, kindTypes) || len(typed.Args) != 1 { + return "", fmt.Errorf("expression is not an EventKind conversion") + } + return stringConstantValue(typed.Args[0], kindTypes, knownKinds) + case *ast.BinaryExpr: + if typed.Op != token.ADD { + return "", fmt.Errorf("EventKind string expression uses unsupported operator %s", typed.Op) + } + left, err := stringConstantValue(typed.X, kindTypes, knownKinds) + if err != nil { + return "", err + } + right, err := stringConstantValue(typed.Y, kindTypes, knownKinds) + if err != nil { + return "", err + } + return left + right, nil + default: + return "", fmt.Errorf("unsupported EventKind expression %T", expression) + } +} + +// isEventKindType reports whether an AST expression names the public event kind type. +func isEventKindType(expression ast.Expr, kindTypes map[string]struct{}) bool { + identifier, ok := expression.(*ast.Ident) + if !ok { + return false + } + _, ok = kindTypes[identifier.Name] + return ok +} + +// parseReadmeEventReference returns the exact kind-to-layer mapping from the +// manual Events reference table while rejecting duplicate rows. +func parseReadmeEventReference(t *testing.T, filename string) map[EventKind]documentedEventKind { + t.Helper() + file, err := os.Open(filename) + if err != nil { + t.Fatalf("open %s: %v", filename, err) + } + defer func() { + if closeErr := file.Close(); closeErr != nil { + t.Errorf("close %s: %v", filename, closeErr) + } + }() + + kinds := make(map[EventKind]documentedEventKind) + scanner := bufio.NewScanner(file) + inReference := false + line := 0 + for scanner.Scan() { + line++ + text := scanner.Text() + if text == "### Events reference" { + inReference = true + continue + } + if !inReference { + continue + } + if strings.HasPrefix(text, "## ") || strings.HasPrefix(text, "### ") { + break + } + if !strings.HasPrefix(text, "|") || strings.HasPrefix(text, "| Layer ") || strings.HasPrefix(text, "| ---") { + continue + } + columns := strings.Split(text, "|") + if len(columns) != 5 { + t.Fatalf("%s:%d event row has %d columns, want 3", filename, line, len(columns)-2) + } + layer := EventLayer(strings.Trim(strings.TrimSpace(columns[1]), "*")) + kind := EventKind(strings.TrimSpace(columns[2])) + meaning := strings.TrimSpace(columns[3]) + if layer == "" || kind == "" || meaning == "" { + t.Fatalf("%s:%d event row has an empty layer, kind, or meaning", filename, line) + } + if previous, duplicate := kinds[kind]; duplicate { + t.Fatalf("%s:%d duplicates EventKind %q from line %d", filename, line, kind, previous.line) + } + kinds[kind] = documentedEventKind{layer: layer, line: line} + } + if err := scanner.Err(); err != nil { + t.Fatalf("scan %s: %v", filename, err) + } + if !inReference { + t.Fatalf("%s does not contain an Events reference heading", filename) + } + return kinds +} + +// parseEventsDocEventReference returns the exact identifiers listed in the +// detailed Event kinds section while rejecting duplicate entries. +func parseEventsDocEventReference(t *testing.T, filename string) map[string]int { + t.Helper() + file, err := os.Open(filename) + if err != nil { + t.Fatalf("open %s: %v", filename, err) + } + defer func() { + if closeErr := file.Close(); closeErr != nil { + t.Errorf("close %s: %v", filename, closeErr) + } + }() + + identifierPattern := regexp.MustCompile("`(Event[A-Za-z0-9]+)`") + kinds := make(map[string]int) + scanner := bufio.NewScanner(file) + inReference := false + line := 0 + for scanner.Scan() { + line++ + text := scanner.Text() + if text == "## Event kinds" { + inReference = true + continue + } + if !inReference { + continue + } + if strings.HasPrefix(text, "## ") { + break + } + if !strings.HasPrefix(text, "- ") { + continue + } + parts := strings.SplitN(strings.TrimPrefix(text, "- "), ":", 2) + if len(parts) != 2 || strings.TrimSpace(parts[1]) == "" { + t.Fatalf("%s:%d event catalog entry must have a nonempty meaning", filename, line) + } + matches := identifierPattern.FindAllStringSubmatch(parts[0], -1) + if len(matches) != 1 { + t.Fatalf("%s:%d event catalog entry has %d EventKind identifiers, want exactly one", filename, line, len(matches)) + } + name := matches[0][1] + if previousLine, duplicate := kinds[name]; duplicate { + t.Fatalf("%s:%d duplicates EventKind identifier %s from line %d", filename, line, name, previousLine) + } + kinds[name] = line + } + if err := scanner.Err(); err != nil { + t.Fatalf("scan %s: %v", filename, err) + } + if !inReference { + t.Fatalf("%s does not contain an Event kinds heading", filename) + } + return kinds +} + +// exportedEventKindNames indexes the source definitions by public identifier. +func exportedEventKindNames(definitions []eventKindDefinition) map[string]struct{} { + names := make(map[string]struct{}, len(definitions)) + for _, definition := range definitions { + names[definition.name] = struct{}{} + } + return names +} diff --git a/integration/root/process_events_integration_test.go b/integration/root/process_events_integration_test.go index fa0116b..b948c05 100644 --- a/integration/root/process_events_integration_test.go +++ b/integration/root/process_events_integration_test.go @@ -22,7 +22,7 @@ type processEventRecorder struct { func (r *processEventRecorder) Observe(_ context.Context, event queue.Event) { switch event.Kind { - case queue.EventProcessStarted, queue.EventProcessSucceeded, queue.EventProcessFailed: + case queue.EventProcessStarted, queue.EventProcessSucceeded, queue.EventProcessFailed, queue.EventProcessArchived: r.mu.Lock() r.events = append(r.events, event) r.mu.Unlock() @@ -54,10 +54,11 @@ func (r *processEventRecorder) has(kind queue.EventKind, predicate func(queue.Ev func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { fixtures := []struct { - name string - queue string - workers int - newQueue func(t *testing.T, observer queue.Observer) QueueRuntime + name string + queue string + workers int + confirmsArchives bool + newQueue func(t *testing.T, observer queue.Observer) QueueRuntime }{ { name: testenv.BackendRedis, @@ -73,9 +74,10 @@ func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { }, }, { - name: testenv.BackendMySQL, - queue: "obs_events_mysql", - workers: 2, + name: testenv.BackendMySQL, + queue: "obs_events_mysql", + workers: 2, + confirmsArchives: true, newQueue: func(t *testing.T, observer queue.Observer) QueueRuntime { ensureMySQLDB(t) q, err := newQueueRuntime(withObserver(withDefaultQueue(mysqlCfg(mysqlDSN(integrationMySQL.addr)), "obs_events_mysql"), observer)) @@ -86,9 +88,10 @@ func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { }, }, { - name: testenv.BackendPostgres, - queue: "obs_events_postgres", - workers: 2, + name: testenv.BackendPostgres, + queue: "obs_events_postgres", + workers: 2, + confirmsArchives: true, newQueue: func(t *testing.T, observer queue.Observer) QueueRuntime { ensurePostgresDB(t) q, err := newQueueRuntime(withObserver(withDefaultQueue(postgresCfg(postgresDSN(integrationPostgres.addr)), "obs_events_postgres"), observer)) @@ -99,9 +102,10 @@ func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { }, }, { - name: testenv.BackendSQLite, - queue: "obs_events_sqlite", - workers: 2, + name: testenv.BackendSQLite, + queue: "obs_events_sqlite", + workers: 2, + confirmsArchives: true, newQueue: func(t *testing.T, observer queue.Observer) QueueRuntime { dsn := fmt.Sprintf("%s/obs-events-%d.db", t.TempDir(), time.Now().UnixNano()) q, err := newQueueRuntime(withObserver(withDefaultQueue(sqliteCfg(dsn), "obs_events_sqlite"), observer)) @@ -209,6 +213,11 @@ func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { recorder.count(queue.EventProcessSucceeded) >= 1 && recorder.count(queue.EventProcessFailed) >= 1 }) + if fx.confirmsArchives { + waitForObservabilityScenario(t, "process_archive", 12*time.Second, func() bool { + return recorder.count(queue.EventProcessArchived) >= 1 + }) + } requireScenarioTrue(t, "process_started_queue", recorder.has(queue.EventProcessStarted, func(event queue.Event) bool { @@ -228,6 +237,15 @@ func TestObservabilityIntegration_ProcessEvents_AllBackends(t *testing.T) { }), "expected process_failed with queue=%q type=%q max_retry=0", fx.queue, failType, ) + if fx.confirmsArchives { + requireScenarioTrue(t, "process_archived_fields", + recorder.has(queue.EventProcessArchived, func(event queue.Event) bool { + return event.Layer == queue.EventLayerWorker && event.Queue == fx.queue && event.JobType == failType && + event.MaxRetry == 0 && event.Err != nil + }), + "expected process_archived with queue=%q type=%q max_retry=0", fx.queue, failType, + ) + } }) } } diff --git a/observability.go b/observability.go index c3282bc..8bb5f3e 100644 --- a/observability.go +++ b/observability.go @@ -62,11 +62,11 @@ const ( EventQueuePaused EventKind = "queue_paused" // EventQueueResumed indicates queue consumption was resumed. EventQueueResumed EventKind = "queue_resumed" - // EventProcessRecovered indicates a stale in-flight job was requeued for recovery. + // EventProcessRecovered indicates SQL bulk recovery requeued one stale in-flight claim without loading its identity fields. EventProcessRecovered EventKind = "process_recovered" // EventRepublishFailed indicates an internal delay/retry republish attempt failed. EventRepublishFailed EventKind = "republish_failed" - // EventSettlementFailed indicates a broker acknowledgement or deletion failed after handler or replacement work completed. + // EventSettlementFailed indicates delivery finalization or broker settlement failed or remained ambiguous. EventSettlementFailed EventKind = "settlement_failed" // EventJobStarted indicates logical job execution began. EventJobStarted EventKind = "job_started" @@ -547,7 +547,7 @@ func (s StatsSnapshot) Failed(name string) int64 { return counters.Failed } -// Paused returns paused count for a queue. +// Paused returns the observed pause state for a queue as zero or one. // @group Observability // // Example: paused count getter @@ -768,11 +768,9 @@ func (c *StatsCollector) Observe(ctx context.Context, event Event) { case EventProcessArchived: state.counters.Archived++ case EventQueuePaused: - state.counters.Paused++ + state.counters.Paused = 1 case EventQueueResumed: - if state.counters.Paused > 0 { - state.counters.Paused-- - } + state.counters.Paused = 0 case EventProcessSucceeded: state.closeActive(ctx, event) state.counters.Processed++ @@ -889,9 +887,10 @@ func countSince(in []time.Time, cutoff time.Time) int64 { } type observedQueue struct { - inner queueBackend - driver Driver - observer Observer + inner queueBackend + driver Driver + observer Observer + controlMu sync.Mutex } func newObservedQueue(inner queueBackend, driver Driver, observer Observer) queueBackend { @@ -940,6 +939,9 @@ func (q *observedQueue) Ready(ctx context.Context) error { } func (q *observedQueue) Pause(ctx context.Context, queueName string) error { + q.controlMu.Lock() + defer q.controlMu.Unlock() + controller, ok := q.inner.(QueueController) if !ok { return ErrPauseUnsupported @@ -958,6 +960,9 @@ func (q *observedQueue) Pause(ctx context.Context, queueName string) error { } func (q *observedQueue) Resume(ctx context.Context, queueName string) error { + q.controlMu.Lock() + defer q.controlMu.Unlock() + controller, ok := q.inner.(QueueController) if !ok { return ErrPauseUnsupported diff --git a/observability_branches_test.go b/observability_branches_test.go index a9d65eb..ca76e36 100644 --- a/observability_branches_test.go +++ b/observability_branches_test.go @@ -291,8 +291,55 @@ func TestObservedQueue_WrapperMethods(t *testing.T) { if oq.Driver() != DriverWorkerpool { t.Fatalf("expected driver %q, got %q", DriverWorkerpool, oq.Driver()) } - if len(recorder.events) < 2 { - t.Fatalf("expected pause/resume events, got %d", len(recorder.events)) + wantControlEvents := []Event{ + {Layer: EventLayerQueue, Kind: EventQueuePaused, Driver: DriverWorkerpool, Queue: "critical"}, + {Layer: EventLayerQueue, Kind: EventQueueResumed, Driver: DriverWorkerpool, Queue: "critical"}, + } + if len(recorder.events) != len(wantControlEvents) { + t.Fatalf("control events = %+v, want exactly pause and resume", recorder.events) + } + for index, want := range wantControlEvents { + got := recorder.events[index] + if got.Layer != want.Layer || got.Kind != want.Kind || got.Driver != want.Driver || got.Queue != want.Queue { + t.Errorf("control event %d = %+v, want layer=%q kind=%q driver=%q queue=%q", index, got, want.Layer, want.Kind, want.Driver, want.Queue) + } + if got.SchemaVersion == 0 || got.EventID == "" || got.Time.IsZero() { + t.Errorf("control event %d is not normalized: %+v", index, got) + } + } +} + +// TestObservedQueueControlFactsTrackState verifies idempotent control calls +// remain a gauge instead of accumulating pause and resume deltas. +func TestObservedQueueControlFactsTrackState(t *testing.T) { + collector := NewStatsCollector() + overlay := &observedQueue{ + inner: &queueBackendStub{}, + driver: DriverWorkerpool, + observer: collector, + } + for range 2 { + if err := overlay.Pause(context.Background(), "critical"); err != nil { + t.Fatalf("pause: %v", err) + } + } + if err := overlay.Resume(context.Background(), "critical"); err != nil { + t.Fatalf("resume: %v", err) + } + if paused := collector.Snapshot().Paused("critical"); paused != 0 { + t.Fatalf("paused gauge after repeated pause and resume = %d, want 0", paused) + } + + for range 2 { + if err := overlay.Resume(context.Background(), "critical"); err != nil { + t.Fatalf("repeat resume: %v", err) + } + } + if err := overlay.Pause(context.Background(), "critical"); err != nil { + t.Fatalf("final pause: %v", err) + } + if paused := collector.Snapshot().Paused("critical"); paused != 1 { + t.Fatalf("paused gauge after repeated resume and pause = %d, want 1", paused) } } @@ -315,6 +362,9 @@ func TestObservedQueue_UnsupportedAndErrorBranches(t *testing.T) { if err := oqNoRuntime.Resume(context.Background(), "default"); !errors.Is(err, ErrPauseUnsupported) { t.Fatalf("expected ErrPauseUnsupported, got %v", err) } + if len(recorder.events) != 0 { + t.Fatalf("unsupported queue controls emitted events: %+v", recorder.events) + } oqErrs := &observedQueue{ inner: &queueBackendStub{ @@ -334,6 +384,9 @@ func TestObservedQueue_UnsupportedAndErrorBranches(t *testing.T) { if _, err := oqErrs.Stats(context.Background()); err == nil { t.Fatal("expected stats error") } + if len(recorder.events) != 0 { + t.Fatalf("failed queue controls emitted events: %+v", recorder.events) + } } func TestObservedQueue_RegisterBranches(t *testing.T) { From 38904bccc01f4e0934d9cecacef91d366c1b6024 Mon Sep 17 00:00:00 2001 From: Chris Miles Date: Mon, 20 Jul 2026 07:14:22 +0000 Subject: [PATCH 2/3] docs: regenerate pause state example --- examples/statssnapshot-paused/main.go | 4 ++-- observability.go | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/examples/statssnapshot-paused/main.go b/examples/statssnapshot-paused/main.go index aa79c5d..fafc058 100644 --- a/examples/statssnapshot-paused/main.go +++ b/examples/statssnapshot-paused/main.go @@ -13,9 +13,9 @@ import ( ) func main() { - // Paused returns paused count for a queue. + // Paused returns the observed pause state for a queue as zero or one. - // Example: paused count getter + // Example: pause state getter collector := queue.NewStatsCollector() collector.Observe(context.Background(), queue.Event{ Kind: queue.EventQueuePaused, diff --git a/observability.go b/observability.go index 8bb5f3e..5755a7e 100644 --- a/observability.go +++ b/observability.go @@ -550,7 +550,7 @@ func (s StatsSnapshot) Failed(name string) int64 { // Paused returns the observed pause state for a queue as zero or one. // @group Observability // -// Example: paused count getter +// Example: pause state getter // // collector := queue.NewStatsCollector() // collector.Observe(context.Background(), queue.Event{ From b6520cbe3b411958f2f8488eb64bb60afbf37c7d Mon Sep 17 00:00:00 2001 From: Chris Miles Date: Mon, 20 Jul 2026 07:29:50 +0000 Subject: [PATCH 3/3] test: stabilize RabbitMQ uniqueness contract --- docs/readme/testcounts/integration_count.json | 2 +- integration/root/contract_integration_test.go | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/docs/readme/testcounts/integration_count.json b/docs/readme/testcounts/integration_count.json index c44d067..412b5e1 100644 --- a/docs/readme/testcounts/integration_count.json +++ b/docs/readme/testcounts/integration_count.json @@ -1,5 +1,5 @@ { "count": 631, - "source_hash": "sha256:4b29b5cb51af6e4bb434cc27e03e5fda07775d3a35a371ec5cb4737779183041", + "source_hash": "sha256:64018d37bee1280db310b9ddd5606b890827ea0709f6141dca706c574450d1a8", "backend_scope": "all" } diff --git a/integration/root/contract_integration_test.go b/integration/root/contract_integration_test.go index b152b44..883b557 100644 --- a/integration/root/contract_integration_test.go +++ b/integration/root/contract_integration_test.go @@ -207,6 +207,8 @@ func TestQueueContract_RabbitMQ(t *testing.T) { assertMissingHandlerErr: false, supportsPause: false, supportsNativeStats: false, + uniqueTTL: time.Second, + uniqueExpiryWait: 1200 * time.Millisecond, } runQueueContractSuite(t, factory) }