From 43aa2f81ebed8c939cf7f93a6af0cfe04c615891 Mon Sep 17 00:00:00 2001 From: Rogee Date: Fri, 5 Jun 2026 21:14:45 +0800 Subject: [PATCH] feat(reports): add analytics timeseries rollups --- docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md | 67 +++-- docs/parity/gochat_routes.txt | 3 +- internal/app/bootstrap.go | 1 + internal/handler/api/v1/analytics_handler.go | 36 +++ .../handler/api/v1/analytics_handler_test.go | 23 +- .../reporting_events_rollup_repo.go | 1 + internal/router/router.go | 3 +- internal/service/analytics_p513_test.go | 43 +++ internal/service/analytics_query_helpers.go | 269 ++++++++++++++++++ internal/service/analytics_service.go | 18 ++ internal/service/reporting_rollup_worker.go | 62 ++++ 11 files changed, 494 insertions(+), 32 deletions(-) create mode 100644 internal/service/reporting_rollup_worker.go diff --git a/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md b/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md index 666389bc..6a21118c 100644 --- a/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md +++ b/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md @@ -16,12 +16,12 @@ Build GoChat as a Go backend that can directly reuse the frontend from `referenc ## Current Baseline -- Current tracking checkpoint: 2026-06-05 after `37bdac4 feat(captain): queue copilot response jobs`, with this implementation checkpoint prepared as `feat(reports): derive analytics aggregates`. -- Latest implementation checkpoint: this checkpoint, prepared as `feat(reports): derive analytics aggregates`. +- Current tracking checkpoint: 2026-06-05 after `a48b9cd feat(reports): derive analytics aggregates`, with this implementation checkpoint prepared as `feat(reports): add analytics timeseries rollups`. +- Latest implementation checkpoint: this checkpoint, prepared as `feat(reports): add analytics timeseries rollups`. - Latest documentation-only checkpoint: `2923aae docs: land parity execution tracker`; this document is now the active follow-up plan and supersedes `.hermes/plans/*`. -- Worktree status at this implementation checkpoint: B11.1a aligns Captain assistant CRUD/tools/inbox bindings; B11.1b aligns Captain scenarios and custom tools; B11.1c aligns Captain documents, assistant responses, bulk actions, and custom-tool test payloads; B11.2 aligns Copilot thread/message create/list/get/delete payloads, account/user scoping, and no-LLM fallback persistence; B11.3a aligns Captain preferences show/update payloads and account-level model/feature storage; B11.3b aligns Captain playground request/response payloads, account scoping, v2 history handling, and no-LLM fallback; B11.3c adds the fakeable Captain document sync backend gate with disabled, failed, and fake-success states; B11.3d aligns Captain task request/response payloads, no-provider disabled states, follow-up context, suggestion persistence, and Copilot message tool-call key validation; B11.3e aligns Captain stream DTOs/disabled SSE fallbacks and Copilot push-event payload shapes; B12.1 adds the reusable GoChat server/seed entrypoint plus a Meilisearch-first reused Chatwoot frontend smoke harness and report; B12.2a adds API smoke assertions for auth/profile, inbox, conversation/messages, contact/company, widget config/message, and public CSAT; B12.2b adds a zero-dependency Chrome DevTools browser smoke that loads the reused Chatwoot login and dashboard entrypoints through Vite and checks browser auth/dashboard API requests; B12.3a adds enterprise API smoke assertions for SLA reports/download, CSAT reports/download, automation/macros, audit/custom roles, capacity, Captain, and Copilot; B12.3b adds reused-frontend enterprise browser route navigation for SLA, CSAT, automation, macros, audit logs, custom roles, capacity, Captain, and Copilot request coverage; P5.1 adds the PostgreSQL-backed durable `background_jobs` model/migration plus WorkerPool enqueue, schedule, retry/backoff, dead-letter, idempotency, stale-lock recovery, and focused tests; P5.2 wires `channel.Dispatcher` and `dispatch.EventDispatcher` async paths into durable event jobs with worker replay tests; P5.3 queues Meilisearch write-side index/delete jobs for conversations, messages, contacts, companies, and articles while keeping search reads Meilisearch-first; P5.4 queues automation webhook and email transcript side effects as durable jobs while preserving fakeable delivery boundaries; P5.5 queues Chatwoot-style macro execute fan-out through durable `automation:macro_execution` jobs; P5.6 queues resolve-triggered CSAT survey sends and WhatsApp/Twilio CSAT template creation through durable jobs; P5.7 queues Chatwoot enterprise SLA account scans and applied-SLA evaluation jobs through the durable worker; P5.8 queues Chatwoot-style contact export artifact generation through durable `contact:export` jobs; P5.9 queues normalized provider inbound message persistence/dispatch through durable `webhook:incoming_message_persist` jobs; P5.10 queues Chatwoot `SendReplyJob`-style outbound message delivery through durable `message:send_reply` jobs and provider delivery-status/read-receipt updates through durable webhook status jobs; P5.11 queues Captain document sync, crawl/parser, schedule-sync, response-builder, embedding-update, Copilot response, and Captain conversation response-builder work through durable jobs; P5.12 queues scheduled item fan-out, one-off campaigns, snoozed conversation reopening, account auto-resolution, widget/public message status updates, and account conversation bulk actions through durable jobs; P5.13a replaces the live report, bot report, conversation summary, inbox-label matrix, first-response distribution, and outgoing-message placeholder responses with persisted conversation/message/reporting-event aggregations. Next active implementation slice is P5.13b scheduled/cached report rollups and timeseries index parity. +- Worktree status at this implementation checkpoint: B11.1a aligns Captain assistant CRUD/tools/inbox bindings; B11.1b aligns Captain scenarios and custom tools; B11.1c aligns Captain documents, assistant responses, bulk actions, and custom-tool test payloads; B11.2 aligns Copilot thread/message create/list/get/delete payloads, account/user scoping, and no-LLM fallback persistence; B11.3a aligns Captain preferences show/update payloads and account-level model/feature storage; B11.3b aligns Captain playground request/response payloads, account scoping, v2 history handling, and no-LLM fallback; B11.3c adds the fakeable Captain document sync backend gate with disabled, failed, and fake-success states; B11.3d aligns Captain task request/response payloads, no-provider disabled states, follow-up context, suggestion persistence, and Copilot message tool-call key validation; B11.3e aligns Captain stream DTOs/disabled SSE fallbacks and Copilot push-event payload shapes; B12.1 adds the reusable GoChat server/seed entrypoint plus a Meilisearch-first reused Chatwoot frontend smoke harness and report; B12.2a adds API smoke assertions for auth/profile, inbox, conversation/messages, contact/company, widget config/message, and public CSAT; B12.2b adds a zero-dependency Chrome DevTools browser smoke that loads the reused Chatwoot login and dashboard entrypoints through Vite and checks browser auth/dashboard API requests; B12.3a adds enterprise API smoke assertions for SLA reports/download, CSAT reports/download, automation/macros, audit/custom roles, capacity, Captain, and Copilot; B12.3b adds reused-frontend enterprise browser route navigation for SLA, CSAT, automation, macros, audit logs, custom roles, capacity, Captain, and Copilot request coverage; P5.1 adds the PostgreSQL-backed durable `background_jobs` model/migration plus WorkerPool enqueue, schedule, retry/backoff, dead-letter, idempotency, stale-lock recovery, and focused tests; P5.2 wires `channel.Dispatcher` and `dispatch.EventDispatcher` async paths into durable event jobs with worker replay tests; P5.3 queues Meilisearch write-side index/delete jobs for conversations, messages, contacts, companies, and articles while keeping search reads Meilisearch-first; P5.4 queues automation webhook and email transcript side effects as durable jobs while preserving fakeable delivery boundaries; P5.5 queues Chatwoot-style macro execute fan-out through durable `automation:macro_execution` jobs; P5.6 queues resolve-triggered CSAT survey sends and WhatsApp/Twilio CSAT template creation through durable jobs; P5.7 queues Chatwoot enterprise SLA account scans and applied-SLA evaluation jobs through the durable worker; P5.8 queues Chatwoot-style contact export artifact generation through durable `contact:export` jobs; P5.9 queues normalized provider inbound message persistence/dispatch through durable `webhook:incoming_message_persist` jobs; P5.10 queues Chatwoot `SendReplyJob`-style outbound message delivery through durable `message:send_reply` jobs and provider delivery-status/read-receipt updates through durable webhook status jobs; P5.11 queues Captain document sync, crawl/parser, schedule-sync, response-builder, embedding-update, Copilot response, and Captain conversation response-builder work through durable jobs; P5.12 queues scheduled item fan-out, one-off campaigns, snoozed conversation reopening, account auto-resolution, widget/public message status updates, and account conversation bulk actions through durable jobs; P5.13a replaces the live report, bot report, conversation summary, inbox-label matrix, first-response distribution, and outgoing-message placeholder responses with persisted conversation/message/reporting-event aggregations; P5.13b routes `GET /reports` to Chatwoot-style metric timeseries, adds lazy rollup freshness/idempotency, registers durable `reporting:rollup_day` jobs, and makes rollup replacement hard-delete soft-deleted rows before recompute. Next active implementation slice is B9.3 delayed automation action verification, followed by Phase 2/3 drift and Phase 6 placeholder audits. - `go test ./...` passes. -- Route dump succeeds with `TOTAL: 832` after adding the Chatwoot-compatible Twilio delivery-status route plus the legacy namespaced alias. +- Route dump succeeds with `TOTAL: 833` after adding `GET /api/v1/accounts/:account_id/reports` for the Chatwoot reports index path. - Route parity artifacts now exist under `docs/parity/` and are generated by `cmd/route_parity`. - Tracked frontend-critical route audit covers 277 Chatwoot routes: 270 exact, 0 method-compatible, 7 parameter-compatible, 0 missing. The 7 parameter-compatible routes are Gin-internal parameter-name differences for nested AgentCapacityPolicy users/inbox limits; the external URL shape is equivalent. - `/api/v1/widget` stubs are burned down and public inbox/contact/conversation/message core flows are backed by real handlers. @@ -45,16 +45,16 @@ Next ordered checkpoints: | Order | Slice | Required outcome | Primary verification | | --- | --- | --- | --- | -| 1 | P5.13b scheduled/cached analytics | Wire scheduled or lazy cached report rollups and the `/reports` timeseries index path after the P5.13a placeholder burn-down. | Report worker/service fixtures for freshness, idempotency, and timeseries values. | -| 2 | B9.3 delayed automation actions | Confirm current-reference delayed automation params and queue any still-synchronous scheduled action execution. | Automation worker fixtures for schedule time, retry, idempotency, and observable failure. | -| 3 | Phase 2/3 audit pass | Route/controller/serializer drift found by B12 or new reference inspection is captured as named slices, not free-form TODOs. | Regenerated route parity artifacts and fixture-backed serializer tests. | -| 4 | Phase 6 placeholder burn-down | Remaining account/contact/conversation/message/inbox placeholder handlers are either real Chatwoot-compatible flows or explicitly tracked as unsupported reference gaps. | `rg` placeholder audit, route smoke, and endpoint-family tests. | +| 1 | B9.3 delayed automation actions | Confirm current-reference delayed automation params and queue any still-synchronous scheduled action execution. | Automation worker fixtures for schedule time, retry, idempotency, and observable failure. | +| 2 | Phase 2/3 audit pass | Route/controller/serializer drift found by B12 or new reference inspection is captured as named slices, not free-form TODOs. | Regenerated route parity artifacts and fixture-backed serializer tests. | +| 3 | Phase 6 placeholder burn-down | Remaining account/contact/conversation/message/inbox placeholder handlers are either real Chatwoot-compatible flows or explicitly tracked as unsupported reference gaps. | `rg` placeholder audit, route smoke, and endpoint-family tests. | +| 4 | B12 optional live smoke | Run the checked smoke harness in a full PostgreSQL/Redis/Meilisearch/Vite/Chrome environment and turn failures into named slices. | `docs/parity/frontend_smoke_report.md` pass/fail entries linked to owners. | ## Handoff Contract This checkpoint is intended to make the development plan complete enough to track without reading Hermes notes first. -- The next active implementation slice is P5.13b scheduled/cached analytics. P5.13a replaced the most visible report placeholders with persisted aggregations; remaining analytics work is report rollup freshness/idempotency and `/reports` timeseries index parity. B12 has repeatable API/browser/enterprise smoke harnesses in Review; optional live failures should be converted into named slices instead of blocking Phase 5 job work. +- P5.13b scheduled/cached analytics is now in Review. P5.13a replaced the most visible report placeholders with persisted aggregations; P5.13b adds report rollup freshness/idempotency, durable day rollup jobs, and `/reports` timeseries index parity. B12 has repeatable API/browser/enterprise smoke harnesses in Review; optional live failures should be converted into named slices instead of blocking Phase 5 job work. - The Hermes search plan is fully represented by Phase 1/B6. Future search changes must be Meilisearch-first and must not reintroduce production DB fallback. - The Hermes automation/macro/CSAT plan is fully represented by B8/B9 and Phase 5. Durable delayed execution remains visible Phase 5/B9.3 work if the current reference exposes explicit delayed action params; channel-specific template delivery is in Review. - Enterprise scope is fixed: SLA, Audit, CustomRole, AgentCapacity, Captain/Copilot, CSAT, InboxLimit, automation, macros, assignment policies, and related limits/workflows are in scope; SSO/SAML/LDAP/OIDC are out of scope. @@ -64,7 +64,7 @@ Open work after the current checkpoint: | Area | Next concrete action | Tracking location | Done boundary | | --- | --- | --- | --- | -| Phase 5 jobs | Finish P5.13b scheduled/cached analytics and the B9.3 delayed automation action check on top of the committed durable worker. | `Phase 5: Background Jobs And Integrations` | Report and automation tests prove real aggregation/scheduling, freshness/idempotency, and no frontend-visible placeholder values. | +| Phase 5 jobs | Finish the B9.3 delayed automation action check on top of the committed durable worker. | `Phase 5: Background Jobs And Integrations` | Automation tests prove real scheduling, retry, idempotency, and observable failures, or reference inspection proves no delayed action params remain. | | B12 smoke | Run optional live API/browser/enterprise smoke in a full PostgreSQL/Redis/Meilisearch/Vite environment and convert failures into named slices. | `B12 reused frontend verification breakdown` | `docs/parity/frontend_smoke_report.md` records checked pass/fail results and maps failures to slices. | | Phase 2/3 drift | Expand tracked route/serializer fixtures when B12 exposes frontend-critical gaps. | `Phase 2`, `Phase 3`, `docs/parity/` | Route parity remains 0 missing for tracked frontend routes; serializers have reference fixtures. | | Phase 6 placeholders | Re-run placeholder audit and burn down any frontend-reachable stub. | `Phase 6: Core Product Placeholder Burn-down` | Stub list has no reused-frontend critical path without a named owner. | @@ -78,7 +78,7 @@ Open work after the current checkpoint: | Phase 2 | Route and controller parity audit | Doing | Ruby/Bundler unavailable, so Chatwoot route extraction currently uses static `routes.rb` fallback | | Phase 3 | Data and serializer parity | Doing | JSON fixture coverage is partial and still endpoint-family based | | Phase 4 | Enterprise feature completion | Doing | B7, B8, B9, B10, and B11 are in Review; B12 reused frontend smoke harnesses exist and optional live runs can expose follow-up slices | -| Phase 5 | Background jobs and integrations | Doing | P5.1/P5.2/P5.3/P5.4/P5.5/P5.6/P5.7/P5.8/P5.9/P5.10/P5.11/P5.12 durable worker, event dispatch, search indexing, automation delivery, macro, CSAT survey/template, SLA scan, contact export, inbound webhook persistence, outbound/provider delivery status, Captain document sync/crawl/response/embedding/Copilot/conversation responses, conversation maintenance, message status update, and account bulk-action cores are in Review; P5.13a placeholder aggregation is in Review, while scheduled/cached rollups, timeseries index parity, and the B9.3 delayed automation action check remain open | +| Phase 5 | Background jobs and integrations | Doing | P5.1/P5.2/P5.3/P5.4/P5.5/P5.6/P5.7/P5.8/P5.9/P5.10/P5.11/P5.12 durable worker, event dispatch, search indexing, automation delivery, macro, CSAT survey/template, SLA scan, contact export, inbound webhook persistence, outbound/provider delivery status, Captain document sync/crawl/response/embedding/Copilot/conversation responses, conversation maintenance, message status update, account bulk-action cores, and P5.13 analytics rollups/timeseries are in Review; the B9.3 delayed automation action check remains open | | Phase 6 | Core placeholder burn-down | Doing | account/contact/conversation/message/inbox placeholder groups remain broad | | Phase 7 | Verification harness | Review | B12.1 boot/readiness, B12.2a API assertions, B12.2b browser smoke harness, B12.3a enterprise API assertions, and B12.3b enterprise browser route navigation exist; optional live Meilisearch/full-browser runs remain environment-dependent | @@ -88,11 +88,11 @@ This table is the shortest authoritative handoff view. If an older lower section | Priority | Workstream | Current state | Next checkpoint | Commit close rule | | --- | --- | --- | --- | --- | -| 1 | P5.13 reports/analytics | P5.13a derives live reports, bot reports, conversation summary, inbox-label matrix, first-response distribution, and outgoing-message counts from persisted rows; scheduled/cached rollups and `/reports` timeseries index parity remain. | Add scheduled or lazy cached rollup freshness/idempotency and route `GET /reports` to metric timeseries instead of summary. | Report fixtures verify timeseries values, cache/freshness behavior, and no hidden placeholder JSON. | -| 2 | B9.3 delayed automation actions | Macro fan-out, webhook/transcript delivery, and CSAT jobs are durable; current-reference delayed automation action params still need a final check. | Queue any still-synchronous delayed automation action execution or close the row with reference evidence if no explicit delayed params exist. | Automation worker fixtures verify schedule time, retry, idempotency, and observable failure. | -| 3 | Phase 2/3 drift | Tracked route parity is 0 missing for the current critical set; serializer fixtures remain partial. | Expand route/serializer fixtures when smoke or reference inspection exposes drift. | Regenerate parity artifacts and add endpoint-family fixture tests. | -| 4 | Phase 6 placeholder audit | Widget/public/webhook critical placeholders are burned down; account/contact/conversation/message/inbox audit remains broad. | Run a fresh placeholder audit and assign every frontend-reachable stub to a tracked owner. | `rg` audit result is recorded and no reused-frontend blocker is ownerless. | -| 5 | B12 optional live smoke | API/browser/enterprise smoke commands are checked in; live runs need PostgreSQL, Redis, Meilisearch, Vite, and Chrome. | Run full live smoke when environment is available and map failures to the board. | `docs/parity/frontend_smoke_report.md` records pass/fail and linked owners. | +| 1 | B9.3 delayed automation actions | Macro fan-out, webhook/transcript delivery, CSAT jobs, and analytics rollups are durable; current-reference delayed automation action params still need a final check. | Queue any still-synchronous delayed automation action execution or close the row with reference evidence if no explicit delayed params exist. | Automation worker fixtures verify schedule time, retry, idempotency, and observable failure. | +| 2 | Phase 2/3 drift | Tracked route parity is 0 missing for the current critical set; serializer fixtures remain partial. | Expand route/serializer fixtures when smoke or reference inspection exposes drift. | Regenerate parity artifacts and add endpoint-family fixture tests. | +| 3 | Phase 6 placeholder audit | Widget/public/webhook critical placeholders are burned down; account/contact/conversation/message/inbox audit remains broad. | Run a fresh placeholder audit and assign every frontend-reachable stub to a tracked owner. | `rg` audit result is recorded and no reused-frontend blocker is ownerless. | +| 4 | B12 optional live smoke | API/browser/enterprise smoke commands are checked in; live runs need PostgreSQL, Redis, Meilisearch, Vite, and Chrome. | Run full live smoke when environment is available and map failures to the board. | `docs/parity/frontend_smoke_report.md` records pass/fail and linked owners. | +| 5 | P5.13 reports/analytics | P5.13a derives visible report aggregates from persisted rows; P5.13b adds lazy rollup freshness, durable day rollup jobs, and `/reports` metric timeseries. | Keep report drift closed as frontend smoke or reference inspection exposes additional metrics. | Report fixtures verify timeseries values, cache/freshness behavior, and no hidden placeholder JSON. | ## Open Checkpoint Contracts @@ -104,14 +104,14 @@ These rows are the executable development plan from this point forward. A checkp | P5.11b Captain response/embedding fan-out | Captain document/assistant-response services and repositories, Meilisearch/embedding boundaries | `response_builder_job.rb`, `enterprise/app/jobs/captain/llm/update_embedding_job.rb`, FAQ generator/embedding services | Queue FAQ response generation after successful document content changes, reset unedited responses, create/update assistant responses, and fan out embedding update work behind fakeable LLM gates. | Review by `feat(captain): queue response embedding jobs`; tests cover response reset/create, embedding-disabled retry, fake embedding success, account scope, and no external network in default tests. | | P5.11c Copilot and conversation response jobs | `internal/service/copilot_service.go`, `internal/service/copilot_response_worker.go`, `internal/service/captain_conversation_service.go`, `internal/service/message_service.go`, `internal/app/bootstrap.go` | `enterprise/app/jobs/captain/copilot/response_job.rb`, `enterprise/app/jobs/captain/conversation/response_builder_job.rb`, `enterprise/app/services/captain/copilot/chat_service.rb`, `enterprise/app/services/enterprise/message_templates/hook_execution_service.rb`, `enterprise/app/models/copilot_message.rb` | Queue assistant replies after Copilot user messages and Captain pending-conversation triggers. Persist assistant messages, enqueue Captain conversation replies/handoff messages, open handoff conversations, and keep fakeable provider disabled/failure states observable through durable retry. | Review by `feat(captain): queue copilot response jobs`; focused tests cover Copilot enqueue/persist/fallback/retry, Captain conversation enqueue/handoff/retry/non-pending skip, and service/worker/app package replay. | | P5.13a analytics placeholder burn-down | `internal/service/analytics_service.go`, `internal/service/analytics_query_helpers.go`, live/report handlers | Chatwoot `live_reports_controller.rb`, `reports_controller.rb`, `BotMetricsBuilder`, `InboxLabelMatrixBuilder`, `FirstResponseTimeDistributionBuilder`, `OutgoingMessagesCountBuilder` | Replace frontend-visible zero/empty placeholder responses for live conversations, grouped live conversations, bot summary/metrics, conversation summary, inbox-label matrix, first-response distribution, and outgoing-message counts with persisted conversation/message/reporting-event queries. | Review by `feat(reports): derive analytics aggregates`; focused service/handler tests prove non-zero values from persisted rows and `rg` finds no placeholder TODOs in these methods. | -| P5.13b scheduled/cached analytics | `internal/service/analytics_service.go`, `internal/service/reporting_rollup_service.go`, report handlers/services, worker bootstrap | Chatwoot report controllers/services used by dashboard analytics, `Reports::DataSource`, reporting rollup/backfill jobs | Wire scheduled or lazy cached rollup freshness/idempotency and route `GET /reports` to metric timeseries instead of the summary handler. Define freshness rules for expensive rollups. | Report fixtures prove timeseries values are derived from persisted conversations/messages/reporting events, rollups refresh idempotently, and no hidden placeholder report JSON remains. | +| P5.13b scheduled/cached analytics | `internal/service/analytics_service.go`, `internal/service/analytics_query_helpers.go`, `internal/service/reporting_rollup_worker.go`, `internal/service/reporting_rollup_service.go`, report handlers/services, worker bootstrap | Chatwoot report controllers/services used by dashboard analytics, `Reports::DataSource`, reporting rollup/backfill jobs | Wire scheduled or lazy cached rollup freshness/idempotency and route `GET /reports` to metric timeseries instead of the summary handler. Define freshness rules for expensive rollups. | Review by `feat(reports): add analytics timeseries rollups`; report fixtures prove timeseries values are derived from persisted conversations/messages/reporting events, rollups refresh idempotently, durable `reporting:rollup_day` jobs replay, and hidden placeholder report JSON does not reappear. | | Phase 2/3 drift audit | `cmd/route_parity`, `docs/parity/*`, serializer tests | `reference/chatwoot/config/routes.rb`, controller Jbuilder views, reused frontend API clients | Convert any smoke/reference mismatch into a named route, controller, or serializer slice. Static route extraction remains acceptable until Ruby/Bundler is available. | Regenerated route parity shows 0 missing tracked frontend routes; new serializer fixtures cover the drift. | | Phase 6 placeholder burn-down | Account/contact/conversation/message/inbox handlers and services | Matching reference controllers/Jbuilder views plus reused frontend screens | Re-run placeholder audit and assign every frontend-reachable stub to a specific owner. Burn down the highest-impact stubs before broad feature expansion. | `rg` placeholder audit is recorded here; no reused-frontend critical path is ownerless. | | B12 live smoke | `scripts/parity_frontend_smoke.sh`, `docs/parity/frontend_smoke_report.md`, `cmd/gochat` | Reused `reference/chatwoot` Vite frontend, dashboard route/API clients | Run optional live API/browser/enterprise smoke with PostgreSQL, Redis, Meilisearch, GoChat, Vite, and Chrome. Convert failures into named rows above. | Smoke report records command, environment, pass/fail, artifacts, and linked follow-up owners. | Checkpoint sequencing: -1. Run P5.13 now that P5.11a-c Captain/Copilot async behavior is in Review, so report aggregation can include final message/job side effects. +1. Run B9.3 delayed automation action verification next, now that macro/CSAT/report job families are durable. 2. Run Phase 2/3 and Phase 6 audits after each smoke failure or route/serializer change, not as a one-time cleanup. 3. Keep B12 live smoke optional until the full external stack is available, but every failed live smoke must become a named row in this table. @@ -140,6 +140,7 @@ This ledger records the committed parity checkpoints that future slices should b | Commit | Scope | Verification summary | Follow-up state | | --- | --- | --- | --- | +| `feat(reports): add analytics timeseries rollups` | Completes P5.13b for scheduled/cached analytics parity. `GET /api/v1/accounts/:account_id/reports` and v2 `/reports` now route to metric timeseries instead of summary, with support for account/inbox/agent/team/label dimensions, day/hour/week/month/year buckets, conversation/message/reporting-event metrics, and business-hours averages. Analytics summary/dimension/traffic reads call `EnsureRollupsForRange` for lazy freshness; `reporting:rollup_day` jobs provide durable day recompute; rollup replacement uses `Unscoped` delete so soft-deleted rows cannot violate uniqueness on recompute. | `go test ./internal/service -run 'Analytics' -count=1`; `go test ./internal/handler/api/v1 -run 'Analytics\|LiveReport' -count=1`; `go test ./internal/service ./internal/handler/api/v1 ./internal/router ./internal/worker ./internal/app -count=1`; `go run ./cmd/dump_routes > docs/parity/gochat_routes.txt`; `go run ./cmd/route_parity`; `go test ./...`; `git diff --check`; full verification recorded in the P5.13 section. | P5.13 moves to Review; continue B9.3 delayed automation action verification, then Phase 2/3 drift and Phase 6 placeholder audits. | | `feat(reports): derive analytics aggregates` | Advances P5.13a by replacing frontend-visible analytics placeholder responses with persisted aggregations. Live report conversation metrics now count open/unattended/unassigned/pending conversations with team filtering; grouped live reports return assignee/team grouped counts; bot summary/metrics, conversation summary, inbox-label matrix, first-response distribution, and outgoing-message counts are derived from conversations, messages, labels, agent-bot bindings, and reporting events instead of fixed zero/empty JSON. | `go test ./internal/service -run 'Analytics' -count=1`; `go test ./internal/handler/api/v1 -run 'Analytics\|LiveReport' -count=1`; `go test ./internal/service ./internal/handler/api/v1 ./internal/worker ./internal/app -count=1`; `go test ./...`; `git diff --check`; full verification recorded in the P5.13 section. | P5.13a moves to Review; continue P5.13b scheduled/cached rollup freshness and `/reports` timeseries index parity, then B9.3 delayed automation action check. | | `feat(captain): queue copilot response jobs` | Completes P5.11c with durable Chatwoot `Captain::Copilot::ResponseJob` and `Captain::Conversation::ResponseBuilderJob` equivalents. Copilot thread/message creation now persists the user message and enqueues `captain:copilot_response` when a WorkerPool is configured, while worker replay reloads the account/user/thread/message scope and persists assistant replies through a fakeable backend or the existing no-provider fallback. Incoming pending conversation messages for Captain-enabled inboxes now enqueue `captain:conversation_response_builder`; replay collects public incoming/outgoing history, creates Captain outgoing replies, enqueues provider send-reply, and opens the conversation with a handoff message when the backend requests handoff. | `go test ./internal/service -run 'CopilotResponse\|CaptainConversation\|MessageService' -count=1`; `go test ./internal/service ./internal/worker ./internal/app -count=1`; `go test ./...`; `git diff --check`; full verification recorded in the P5.11 section. | P5.11 moves to Review; continue P5.13 analytics aggregation, then Phase 2/3 drift and Phase 6 placeholder audits as smoke exposes gaps. | | `feat(captain): queue response embedding jobs` | Advances P5.11b with durable Captain document response building and embedding update fan-out. Successful document sync/parser content updates enqueue `captain:document_response_builder`, worker replay resets only unedited document responses, preserves edited responses, creates approved `Captain::Document` assistant responses from a fakeable FAQ backend, and enqueues `captain:llm_update_embedding` jobs for created responses. Embedding replay reloads account-scoped responses, uses a fakeable embedding backend or configured LLM provider, and surfaces missing provider config as retryable worker failures. | `go test ./internal/service -run 'CaptainDocumentService\|EnqueueCaptainDocumentScheduleSyncs' -count=1`; `go test ./internal/service ./internal/worker ./internal/app -count=1`; `go test ./...`; `git diff --check`; full verification recorded in the P5.11 section. | P5.11b moves to Review; continue P5.11c Copilot/conversation response jobs, then P5.13 analytics aggregation. | @@ -1642,14 +1643,14 @@ Known hotspots: - `internal/service/message_delivery_worker.go` and `internal/handler/webhook/incoming_persister_jobs.go` now cover P5.10 SendReplyJob-style outbound delivery plus provider delivery-status/read-receipt jobs. - `internal/handler/webhook/*` and provider services now keep verification/parsing in the HTTP path while `IncomingPersister` moves normalized inbound persistence/dispatch behind durable jobs when the WorkerPool is configured. - `internal/service/captain_*` and `internal/service/copilot_*` now cover Captain document sync, crawl/parser, schedule-sync, response-builder, embedding-update, Copilot response, and Captain conversation response-builder jobs. -- `internal/service/analytics_service.go` has placeholder analytics/report paths and is the owner for P5.13 aggregation work. +- `internal/service/analytics_service.go`, `internal/service/analytics_query_helpers.go`, and `internal/service/reporting_rollup_worker.go` now own P5.13 analytics aggregation, lazy rollup freshness, and durable day rollup replay. Checklist: - [x] Add `background_jobs` persistence with job type, payload, queue, status, attempt counters, scheduled/locked timestamps, idempotency key, last error, and completion timestamps. - [x] Replace `internal/worker/worker.go` placeholder with enqueue, schedule, perform, retry/backoff, dead-letter, and graceful shutdown behavior. - [x] Map the committed Chatwoot job/listener families to Go worker responsibilities in the tracking table. -- [ ] Finish durable job dispatch for delayed automation actions if current-reference params require it, and finish report aggregation. +- [ ] Finish durable job dispatch for delayed automation actions if current-reference params require it. - [ ] Add retry and failure logging for the remaining external provider calls. - [ ] Add tests for the remaining enqueueing, uniqueness/idempotency, delayed execution, retries, worker restart pickup, and fakeable external side effects. @@ -1674,7 +1675,7 @@ Tracking table: | P5.10 | Queue outbound message delivery and delivery-status updates. | `send_reply_job.rb`, provider delivery/status jobs | message send/channel services, delivery status handler | Outgoing message creation and provider delivery are separated; retries update message/delivery status exactly once. | Review by `feat(messages): queue send replies` and `feat(messages): queue delivery statuses` | | P5.11 | Queue Captain document sync, crawl, response building, embeddings, and Copilot responses. | Captain document/crawl/response/embedding/Copilot/conversation jobs | `internal/service/captain_document_service.go`, `internal/service/copilot_service.go`, `internal/service/captain_conversation_service.go`, Captain/Copilot services | Existing fakeable disabled/failure gates run under durable jobs; document statuses, Copilot message persistence, and Captain conversation replies survive worker restart. | Review by `feat(captain): queue document syncs`, `feat(captain): queue document crawl jobs`, `feat(captain): queue response embedding jobs`, and `feat(captain): queue copilot response jobs` | | P5.12 | Queue conversation maintenance jobs. | `trigger_scheduled_items_job.rb`, `campaigns/trigger_oneoff_campaign_job.rb`, `conversations/resolution_job.rb`, `reopen_snoozed_conversations_job.rb`, `update_message_status_job.rb`, `bulk_actions_job.rb` | conversation service/handlers | Auto-resolution, snooze reopen, status updates, and bulk actions are scheduled/retryable with idempotent tests. | Review by `feat(conversations): queue maintenance jobs`, `feat(conversations): queue message status updates`, and `feat(conversations): queue bulk actions` | -| P5.13 | Replace placeholder analytics/report builders that need background aggregation. | reporting jobs/services and report controllers | `internal/service/analytics_service.go`, `internal/service/analytics_query_helpers.go`, reporting services | P5.13a derives live/bot/conversation summary/matrix/distribution/outgoing-count report values from persisted rows; remaining expensive aggregation should be scheduled or cached with freshness rules. | Doing: placeholder burn-down Review by `feat(reports): derive analytics aggregates`; scheduled/cached rollups and `/reports` timeseries index parity Todo | +| P5.13 | Replace placeholder analytics/report builders that need background aggregation. | reporting jobs/services and report controllers | `internal/service/analytics_service.go`, `internal/service/analytics_query_helpers.go`, `internal/service/reporting_rollup_worker.go`, reporting services | P5.13a derives live/bot/conversation summary/matrix/distribution/outgoing-count report values from persisted rows; P5.13b adds lazy cached rollup freshness, durable `reporting:rollup_day` jobs, and `/reports` metric timeseries index parity. | Review by `feat(reports): derive analytics aggregates` and `feat(reports): add analytics timeseries rollups`; remaining report work should be created as drift slices when frontend/reference exposes additional metrics | P5.1 current checkpoint: @@ -1682,7 +1683,7 @@ P5.1 current checkpoint: - Replaced the `WorkerPool` stub with enqueue/schedule APIs, handler registration, `ProcessOne`, start/stop worker loops, PostgreSQL `FOR UPDATE SKIP LOCKED` claiming, retry/backoff, dead-letter state, and stale-lock recovery for worker restart pickup. - Kept older `NewWorkerPool()` construction as a no-op-compatible path while adding `NewWorkerPoolWithOptions(db, ...)` for durable wiring. - Added focused tests proving enqueue/idempotency, due job completion, retry to dead-letter, queue/schedule filtering, and stale running job requeue. -- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it, and analytics/report aggregation. +- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it. P5.1 verification: @@ -1698,7 +1699,7 @@ P5.2 current checkpoint: - `dispatch.EventDispatcher` can now attach a `WorkerPool`; sync listeners still execute inline, while async listeners and `DispatchAsync` subscribers enqueue durable per-listener `event:listener_dispatch` jobs. - Worker replay resolves the listener by name and invokes `OnEvent` with the serialized event payload. Missing listeners or malformed payloads fail the job, preserving retry/dead-letter visibility. - `WorkerPool` default queue handling was broadened to process all queues unless explicitly filtered, so `events` jobs are processed by the default worker. -- Remaining integration work is feature-specific: delayed automation actions if still required and analytics aggregation must register producers/handlers on this durable path. +- Remaining integration work is feature-specific: delayed automation actions if still required; analytics aggregation now registers durable day rollup handlers on this path. P5.2 verification: @@ -1732,7 +1733,7 @@ P5.4 current checkpoint: - The durable job handlers reuse the existing fakeable HTTP webhook and SMTP transcript delivery interfaces, so timeout/retry behavior remains tested at the delivery boundary while worker attempts, errors, retry schedules, completion, and dead-letter state are visible through `background_jobs`. - `AutomationRuleService`, `AutomationRuleListener`, and `MacroService` can carry the WorkerPool through action execution. `Bootstrap` now creates a durable WorkerPool, registers automation delivery jobs, wires the dispatcher with event jobs, and starts the worker in `App.Run`. - Existing no-worker constructors still preserve the previous synchronous fallback for focused tests and development paths that do not start the durable worker yet. -- Remaining Phase 5 work: delayed automation actions if current-reference params require them and analytics/report aggregation; macro fan-out, CSAT sends/templates, Meilisearch indexing fan-out, SLA scans, contact export email, provider webhook/offline delivery, Captain/Copilot, and conversation maintenance are in Review. +- Remaining Phase 5 work: delayed automation actions if current-reference params require them; macro fan-out, CSAT sends/templates, Meilisearch indexing fan-out, SLA scans, contact export email, provider webhook/offline delivery, Captain/Copilot, conversation maintenance, and analytics rollups/timeseries are in Review. P5.4 verification: @@ -1769,7 +1770,7 @@ P5.7 current checkpoint: - `sla:process_account` mirrors `Sla::ProcessAccountAppliedSlasJob`: it finds active and active_with_misses AppliedSLA rows for the account and queues `sla:process_applied` jobs. - `sla:process_applied` mirrors `Sla::ProcessAppliedSlaJob`: it calls the existing `AppliedSlaService.Evaluate`, preserving idempotent miss events, notification fan-out, retry/backoff, and dead-letter visibility through `background_jobs`. - Bootstrap registers the SLA job handlers and enqueues the initial root scan when the app boots. -- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it, and analytics aggregation. +- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it. P5.7 verification: @@ -1787,7 +1788,7 @@ P5.8 current checkpoint: - Completed export jobs are idempotent: a duplicate replay of an already-completed export is a no-op, so completion notifications and mail are not duplicated. - Failed worker execution records the export error and leaves the background job retry/dead-letter state observable through `background_jobs`. - No-worker construction still uses the synchronous fallback for focused tests and local paths that do not start the durable worker. -- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it, and analytics aggregation. +- Remaining Phase 5 work: delayed automation scheduled-item execution if reference params require it. P5.8 verification: @@ -1860,7 +1861,7 @@ P5.11 current checkpoint: - Conversation response replay collects public incoming/outgoing history, maps incoming to user and outgoing to assistant context, creates Captain outgoing messages with optional `agent_name`, enqueues `message:send_reply`, and handles handoff responses by creating the configured handoff message and opening the pending conversation. - No-worker construction keeps the previous mark-syncing fallback for focused tests and local paths that do not start the durable worker. - Bootstrap registers Captain document, Copilot response, and Captain conversation response jobs on the shared WorkerPool and wires the assistant-response repository so production sync/crawl/schedule/response/embedding/Copilot/conversation requests are replayable after process restart. -- Remaining P5.11 work: none known for the named durable Captain/Copilot job row; continue P5.13 analytics aggregation. +- Remaining P5.11 work: none known for the named durable Captain/Copilot job row; analytics aggregation is now in P5.13 Review, so the next Phase 5 check is B9.3 delayed automation actions. P5.11 verification: @@ -1908,14 +1909,21 @@ P5.13 current checkpoint: - Bot summary and bot metrics now derive bot resolutions/handoffs from `reporting_events`, exclude handoff conversations from bot resolutions, and calculate bot conversation/message counts and rates from active `agent_bot_inboxes`. - Conversation summary now derives conversation, incoming-message, outgoing-message, resolution, first-response, resolution-time, and reply-time values from persisted conversations/messages/reporting events. - Inbox-label matrix now uses `inboxes`, `tags`, and `conversation_labels`; first-response distribution buckets `first_response` events by inbox channel type; outgoing-message counts group by agent, team, inbox, or label. -- Remaining P5.13 work: scheduled or lazy cached rollup freshness/idempotency and `/api/v2/accounts/:account_id/reports` timeseries index parity. +- P5.13b routes `GET /api/v1/accounts/:account_id/reports` and `GET /api/v2/accounts/:account_id/reports` to metric timeseries instead of summary. Supported metrics include `conversations_count`, `resolutions_count`, incoming/outgoing messages, first-response/resolution/reply averages, and bot resolution/handoff counts. +- Timeseries queries support `type=account|inbox|agent|team|label`, optional `id`, `group_by=hour|day|week|month|year`, and `business_hours=true` for reporting-event average metrics. +- `EnsureRollupsForRange` lazily computes missing daily rollups before summary, dimension, traffic, and timeseries reads, giving expensive reports a deterministic freshness boundary without requiring a separate scheduler for local tests. +- `reporting:rollup_day` durable jobs are registered through `AnalyticsService.SetWorkerPool` and wired in bootstrap; enqueue uses account/date idempotency keys and worker replay calls `ReportingRollupService.ComputeDailyRollup`. +- Rollup recompute now hard-deletes matching soft-deleted rollups through `Unscoped` before insert, preventing unique-key conflicts during idempotent refresh. +- Remaining P5.13 work: none known for the named scheduled/cached analytics row; any additional report metric or serializer mismatch found by reused frontend smoke should become a Phase 2/3 drift slice. P5.13 verification: ```bash env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go test ./internal/service -run 'Analytics' -count=1 env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go test ./internal/handler/api/v1 -run 'Analytics\|LiveReport' -count=1 -env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go test ./internal/service ./internal/handler/api/v1 ./internal/worker ./internal/app -count=1 +env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go test ./internal/service ./internal/handler/api/v1 ./internal/router ./internal/worker ./internal/app -count=1 +env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go run ./cmd/dump_routes > docs/parity/gochat_routes.txt +env GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go run ./cmd/route_parity env TMPDIR=/home/rogee/Projects/gochat/.tmp/test-tmp GOCACHE=/tmp/gochat-gocache GOMODCACHE=/tmp/gochat-gomodcache go test ./... git diff --check ``` @@ -2021,6 +2029,7 @@ Verification milestone gates: ## Progress Log +- 2026-06-05: P5.13b scheduled/cached analytics checkpoint prepared as `feat(reports): add analytics timeseries rollups`; `GET /api/v1/accounts/:account_id/reports` and v2 `/reports` now return Chatwoot-style metric timeseries with account/inbox/agent/team/label dimensions, day/hour/week/month/year buckets, and business-hours average support. Analytics reads lazily ensure missing daily rollups, `reporting:rollup_day` durable jobs replay account/date recomputes with idempotency keys, and rollup replacement hard-deletes soft-deleted rows before refresh. Route dump is now `TOTAL: 833`; tracked route parity remains `270 exact, 0 method-compatible, 7 parameter-compatible, 0 missing out of 277 tracked critical routes`. Focused analytics service tests, analytics/live handler tests, service/handler/router/worker/app package tests, route generation/parity, full `go test ./...`, and `git diff --check` passed. P5.13 moves to Review; next slice is B9.3 delayed automation action verification. - 2026-06-05: P5.13a analytics placeholder burn-down checkpoint prepared as `feat(reports): derive analytics aggregates`; live report conversation metrics and grouped metrics now read persisted conversations, bot summary/metrics read reporting events plus active agent-bot inbox bindings, conversation summary reads conversations/messages/reporting events, inbox-label matrix reads inbox/tag/conversation-label rows, first-response distribution buckets reporting events by channel, and outgoing-message counts group by agent/team/inbox/label. Focused analytics service tests, analytics/live handler tests, service/handler/worker/app package tests, full `go test ./...`, and `git diff --check` passed. Remaining P5.13 follow-up is scheduled/cached rollup freshness and `/reports` timeseries index parity. - 2026-06-05: P5.11c Copilot/conversation response checkpoint prepared as `feat(captain): queue copilot response jobs`; Copilot thread/message creation now enqueues `captain:copilot_response` after persisting the user message, worker replay persists assistant replies through a fakeable backend or the existing no-provider fallback, and backend failures surface as retryable jobs. Pending Captain-enabled incoming conversations now enqueue `captain:conversation_response_builder`, replay creates Captain outgoing replies, queues provider send-reply, and opens conversations on handoff. Focused Copilot/Captain conversation worker tests, service/worker/app package tests, full `go test ./...`, and `git diff --check` passed. P5.11 moves to Review; next slice is P5.13 analytics aggregation. - 2026-06-05: P5.11b Captain response/embedding checkpoint prepared as `feat(captain): queue response embedding jobs`; successful document sync/parser content updates now enqueue `captain:document_response_builder`, response-builder replay resets only unedited document responses, preserves edited responses, creates approved `Captain::Document` assistant responses through a fakeable FAQ backend, and fans out `captain:llm_update_embedding` jobs. Embedding replay reloads account-scoped responses and uses a fakeable embedding backend or configured LLM provider, with missing providers surfacing as retryable worker failures. Focused Captain document worker tests, service/worker/app package tests, full `go test ./...`, and `git diff --check` passed. Remaining P5.11 follow-up is Copilot/conversation response jobs. diff --git a/docs/parity/gochat_routes.txt b/docs/parity/gochat_routes.txt index 85aa34ba..d6398fb4 100644 --- a/docs/parity/gochat_routes.txt +++ b/docs/parity/gochat_routes.txt @@ -310,6 +310,7 @@ GET /api/v1/accounts/:account_id/portals/:portal_id/members/ GET /api/v1/accounts/:account_id/portals/:portal_id/members/:member_id GET /api/v1/accounts/:account_id/portals/:portal_id/ssl_status GET /api/v1/accounts/:account_id/reporting_events +GET /api/v1/accounts/:account_id/reports GET /api/v1/accounts/:account_id/reports/agents GET /api/v1/accounts/:account_id/reports/bot_metrics GET /api/v1/accounts/:account_id/reports/bot_summary @@ -830,4 +831,4 @@ PUT /public/api/v1/csat_survey/:id PUT /public/api/v1/inboxes/:inbox_id/contacts/:contact_id PUT /public/api/v1/inboxes/:inbox_id/contacts/:contact_id/conversations/:conversation_id/messages/:message_id PUT /widget/direct_uploads/:upload_uuid -TOTAL: 832 +TOTAL: 833 diff --git a/internal/app/bootstrap.go b/internal/app/bootstrap.go index ece8bbe5..6dd002aa 100644 --- a/internal/app/bootstrap.go +++ b/internal/app/bootstrap.go @@ -604,6 +604,7 @@ func Bootstrap(env string) (*App, error) { // Analytics services (P11 — Reports/Analytics) analyticsService := service.NewAnalyticsService(reportingEventRepo, reportingEventsRollupRepo) + analyticsService.SetWorkerPool(workerPool) summaryReportService := service.NewSummaryReportService(reportingEventsRollupRepo) dashboardAppService := service.NewDashboardAppService(dashboardAppRepo) platformAppService := service.NewPlatformAppService(platformAppRepo, accessTokenRepo, permissibleRepo) diff --git a/internal/handler/api/v1/analytics_handler.go b/internal/handler/api/v1/analytics_handler.go index b847d47d..10bb3fdb 100644 --- a/internal/handler/api/v1/analytics_handler.go +++ b/internal/handler/api/v1/analytics_handler.go @@ -2,6 +2,7 @@ package v1 import ( "net/http" + "strconv" "time" "github.com/gin-gonic/gin" @@ -21,6 +22,41 @@ func NewAnalyticsHandler(svc *service.AnalyticsService) *AnalyticsHandler { return &AnalyticsHandler{svc: svc} } +// Index returns metric time-series for the overview charts. +// GET /api/v2/accounts/:account_id/reports?metric=...&since=...&until=... +func (h *AnalyticsHandler) Index(c *gin.Context) { + accountID, ok := parseAccountID(c) + if !ok { + return + } + since, until, ok := parseDateRange(c) + if !ok { + return + } + metric := c.Query("metric") + if metric == "" { + response.AbortWithStatusError(c, http.StatusUnprocessableEntity, response.ErrBadRequest, "metric parameter is required") + return + } + id := uint(0) + if rawID := c.Query("id"); rawID != "" { + parsed, err := strconv.ParseUint(rawID, 10, 64) + if err != nil { + response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid id") + return + } + id = uint(parsed) + } + businessHours := c.Query("business_hours") == "true" || c.Query("business_hours") == "1" + result, err := h.svc.GetTimeseries(c.Request.Context(), accountID, metric, since, until, c.DefaultQuery("type", "account"), id, c.Query("group_by"), businessHours) + if err != nil { + applogger.L().Errorf("Timeseries report: %v", err) + response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to generate report") + return + } + response.OK(c, result) +} + // parseAccountID extracts account_id from URL params. func parseAccountID(c *gin.Context) (uint, bool) { id := getAccountID(c) diff --git a/internal/handler/api/v1/analytics_handler_test.go b/internal/handler/api/v1/analytics_handler_test.go index 93de0dbb..255202c1 100644 --- a/internal/handler/api/v1/analytics_handler_test.go +++ b/internal/handler/api/v1/analytics_handler_test.go @@ -37,6 +37,7 @@ func (s *AnalyticsHandlerTestSuite) SetupSuite() { s.db = db s.Require().NoError(db.AutoMigrate( &model.Account{}, + &model.Conversation{}, &model.ReportingEvent{}, &model.ReportingEventsRollup{}, )) @@ -54,6 +55,7 @@ func (s *AnalyticsHandlerTestSuite) SetupSuite() { gin.SetMode(gin.TestMode) r := gin.New() accounts := r.Group("/api/v1/accounts/:account_id") + accounts.GET("/reports", s.handler.Index) accounts.GET("/reports/summary", s.handler.Summary) accounts.GET("/reports/agents", s.handler.AgentMetrics) accounts.GET("/reports/inboxes", s.handler.InboxMetrics) @@ -70,6 +72,8 @@ func (s *AnalyticsHandlerTestSuite) TearDownSuite() { func (s *AnalyticsHandlerTestSuite) SetupTest() { s.db.Exec("DELETE FROM reporting_events") + s.db.Exec("DELETE FROM reporting_events_rollups") + s.db.Exec("DELETE FROM conversations") } // ========== parseAccountID / parseDateRange edge cases ========== @@ -149,6 +153,23 @@ func (s *AnalyticsHandlerTestSuite) TestSummary_WithData() { s.Equal(http.StatusOK, w.Code) } +func (s *AnalyticsHandlerTestSuite) TestIndex_TimeseriesWithData() { + conv := model.Conversation{AccountID: s.accountID, InboxID: 1, ContactID: 1, Status: string(model.ConversationStatusOpen), ChannelType: "web_widget", Channel: "web_widget"} + s.Require().NoError(s.db.Create(&conv).Error) + s.Require().NoError(s.db.Model(&conv).Updates(map[string]interface{}{"created_at": parseTime("2025-01-15T10:00:00Z")}).Error) + + w := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodGet, + "/api/v1/accounts/1/reports?metric=conversations_count&since=2025-01-01T00:00:00Z&until=2025-02-01T00:00:00Z&type=account&group_by=day", nil) + s.router.ServeHTTP(w, req) + s.Equal(http.StatusOK, w.Code) + var body map[string]interface{} + s.NoError(json.Unmarshal(w.Body.Bytes(), &body)) + data := body["data"].([]interface{}) + s.Len(data, 1) + s.Equal(float64(1), data[0].(map[string]interface{})["value"]) +} + // ========== AgentMetrics ========== func (s *AnalyticsHandlerTestSuite) TestAgentMetrics_InvalidAccountID() { @@ -251,4 +272,4 @@ func (s *AnalyticsHandlerTestSuite) TestNilService() { func parseTime(s string) time.Time { t, _ := time.Parse(time.RFC3339, s) return t -} \ No newline at end of file +} diff --git a/internal/repository/reporting_events_rollup_repo.go b/internal/repository/reporting_events_rollup_repo.go index b6d98f4e..bfa218a8 100644 --- a/internal/repository/reporting_events_rollup_repo.go +++ b/internal/repository/reporting_events_rollup_repo.go @@ -64,6 +64,7 @@ func (r *ReportingEventsRollupRepo) AggregateSummary(ctx context.Context, accoun // DeleteByAccountAndDate deletes all rollups for a specific account and date. func (r *ReportingEventsRollupRepo) DeleteByAccountAndDate(ctx context.Context, accountID uint, date time.Time) error { return r.db.WithContext(ctx). + Unscoped(). Where("account_id = ? AND date = ?", accountID, date). Delete(&model.ReportingEventsRollup{}).Error } diff --git a/internal/router/router.go b/internal/router/router.go index 35dad412..74cbddd8 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -1323,6 +1323,7 @@ func registerV1Routes(g *gin.RouterGroup, h *Handlers) { // Reference: Chatwoot reports_controller.rb + live_reports_controller.rb + summary_reports_controller.rb reports := accountScoped.Group("/reports") { + reports.GET("", h.Analytics.Index) reports.GET("/summary", h.Analytics.Summary) reports.GET("/bot_summary", h.Analytics.BotSummary) reports.GET("/agents", h.Analytics.AgentMetrics) @@ -1714,7 +1715,7 @@ func registerV2Routes(g *gin.RouterGroup, h *Handlers) { reports := accountScoped.Group("/reports") { - reports.GET("", h.Analytics.Summary) + reports.GET("", h.Analytics.Index) reports.GET("/summary", h.Analytics.Summary) reports.GET("/bot_summary", h.Analytics.BotSummary) reports.GET("/agents", h.Analytics.AgentMetrics) diff --git a/internal/service/analytics_p513_test.go b/internal/service/analytics_p513_test.go index 11d4b6f1..095f0164 100644 --- a/internal/service/analytics_p513_test.go +++ b/internal/service/analytics_p513_test.go @@ -7,6 +7,7 @@ import ( "github.com/gochat/gochat/internal/model" "github.com/gochat/gochat/internal/repository" + "github.com/gochat/gochat/internal/worker" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gorm.io/driver/sqlite" @@ -32,6 +33,7 @@ func setupAnalyticsP513Test(t *testing.T) (*gorm.DB, *AnalyticsService, *model.A &model.AgentBotInbox{}, &model.Tag{}, &model.ConversationLabel{}, + &model.BackgroundJob{}, )) t.Cleanup(func() { sqlDB, _ := db.DB() @@ -147,3 +149,44 @@ func TestAnalyticsReportsUsePersistedConversationMessageAndEventRows(t *testing. matrixMap := matrix.(map[string]interface{}) assert.Equal(t, [][]int64{{1}}, matrixMap["matrix"]) } + +func TestAnalyticsTimeseriesAndRollupWorker(t *testing.T) { + db, svc, account, inbox, contact, user, _ := setupAnalyticsP513Test(t) + since := time.Date(2026, 6, 2, 0, 0, 0, 0, time.UTC) + until := since.Add(48 * time.Hour) + conv := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, AssigneeID: &user.ID, Status: string(model.ConversationStatusOpen), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType} + require.NoError(t, db.Create(conv).Error) + require.NoError(t, db.Model(conv).Updates(map[string]interface{}{"created_at": since.Add(2 * time.Hour)}).Error) + require.NoError(t, db.Create(&model.ReportingEvent{Base: model.Base{CreatedAt: since.Add(3 * time.Hour)}, AccountID: account.ID, Name: model.MetricNameFirstResponse, Value: 120, ConversationID: &conv.ID, InboxID: &inbox.ID, UserID: &user.ID, EventStartTime: since, EventEndTime: since.Add(2 * time.Minute)}).Error) + + points, err := svc.GetTimeseries(context.Background(), account.ID, "conversations_count", since, until, "account", 0, "day", false) + require.NoError(t, err) + require.Len(t, points, 1) + assert.Equal(t, float64(1), points[0].Value) + assert.Equal(t, since.Unix(), points[0].Timestamp) + + avgPoints, err := svc.GetTimeseries(context.Background(), account.ID, "avg_first_response_time", since, until, "agent", user.ID, "day", false) + require.NoError(t, err) + require.Len(t, avgPoints, 1) + assert.Equal(t, 120.0, avgPoints[0].Value) + assert.Equal(t, int64(1), avgPoints[0].Count) + + wp := worker.NewWorkerPool(db) + svc.SetWorkerPool(wp) + job, err := EnqueueReportingRollupDay(context.Background(), wp, account.ID, since) + require.NoError(t, err) + require.NotNil(t, job) + again, err := EnqueueReportingRollupDay(context.Background(), wp, account.ID, since) + require.NoError(t, err) + assert.Equal(t, job.ID, again.ID) + processed, err := wp.ProcessOne(context.Background()) + require.NoError(t, err) + assert.True(t, processed) + processed, err = wp.ProcessOne(context.Background()) + require.NoError(t, err) + assert.False(t, processed) + + var rollupCount int64 + require.NoError(t, db.Model(&model.ReportingEventsRollup{}).Where("account_id = ? AND date = ?", account.ID, since).Count(&rollupCount).Error) + assert.Greater(t, rollupCount, int64(0)) +} diff --git a/internal/service/analytics_query_helpers.go b/internal/service/analytics_query_helpers.go index dd38b838..59f012a0 100644 --- a/internal/service/analytics_query_helpers.go +++ b/internal/service/analytics_query_helpers.go @@ -3,7 +3,9 @@ package service import ( "context" "errors" + "fmt" "sort" + "strconv" "strings" "time" @@ -11,6 +13,12 @@ import ( "gorm.io/gorm" ) +type AnalyticsTimeseriesPoint struct { + Value float64 `json:"value"` + Timestamp int64 `json:"timestamp"` + Count int64 `json:"count,omitempty"` +} + func (s *AnalyticsService) analyticsDB() (*gorm.DB, error) { if s == nil || s.db == nil { return nil, errors.New("analytics database is required") @@ -18,6 +26,64 @@ func (s *AnalyticsService) analyticsDB() (*gorm.DB, error) { return s.db, nil } +func (s *AnalyticsService) EnsureRollupsForRange(ctx context.Context, accountID uint, since, until time.Time) error { + if s == nil || s.rollups == nil || s.db == nil || since.IsZero() || until.IsZero() { + return nil + } + start := dayStart(since) + end := dayStart(until) + if until.After(end) { + end = end.AddDate(0, 0, 1) + } + for d := start; d.Before(end); d = d.AddDate(0, 0, 1) { + var count int64 + if err := s.db.WithContext(ctx).Model(&model.ReportingEventsRollup{}).Where("account_id = ? AND date = ?", accountID, d).Count(&count).Error; err != nil { + return err + } + if count > 0 { + continue + } + if err := s.rollups.ComputeDailyRollup(ctx, accountID, d); err != nil { + return err + } + } + return nil +} + +func (s *AnalyticsService) GetTimeseries(ctx context.Context, accountID uint, metric string, since, until time.Time, reportType string, id uint, groupBy string, businessHours bool) ([]AnalyticsTimeseriesPoint, error) { + if strings.TrimSpace(groupBy) == "" { + groupBy = "day" + } + if reportType == "" { + reportType = "account" + } + if err := s.EnsureRollupsForRange(ctx, accountID, since, until); err != nil { + return nil, err + } + switch metric { + case "conversations_count": + return s.conversationCountTimeseries(ctx, accountID, since, until, reportType, id, groupBy, false) + case "resolutions_count": + return s.conversationCountTimeseries(ctx, accountID, since, until, reportType, id, groupBy, true) + case "incoming_messages_count": + return s.messageCountTimeseries(ctx, accountID, since, until, reportType, id, groupBy, model.MessageTypeIncoming) + case "outgoing_messages_count": + return s.messageCountTimeseries(ctx, accountID, since, until, reportType, id, groupBy, model.MessageTypeOutgoing) + case "avg_first_response_time": + return s.eventTimeseries(ctx, accountID, since, until, reportType, id, groupBy, []string{model.MetricNameFirstResponse}, true, businessHours) + case "avg_resolution_time": + return s.eventTimeseries(ctx, accountID, since, until, reportType, id, groupBy, []string{"conversation_resolved", model.MetricNameResolutionTime}, true, businessHours) + case "reply_time": + return s.eventTimeseries(ctx, accountID, since, until, reportType, id, groupBy, []string{model.MetricNameReplyTime}, true, businessHours) + case "bot_resolutions_count": + return s.eventTimeseries(ctx, accountID, since, until, reportType, id, groupBy, []string{"conversation_bot_resolved", model.MetricNameBotResolutionsCount}, false, businessHours) + case "bot_handoffs_count": + return s.eventTimeseries(ctx, accountID, since, until, reportType, id, groupBy, []string{"conversation_bot_handoff", model.MetricNameBotHandoffsCount}, false, businessHours) + default: + return nil, fmt.Errorf("unsupported report metric %q", metric) + } +} + func (s *AnalyticsService) liveConversationMetrics(ctx context.Context, accountID uint, teamID uint) (*ConversationMetrics, error) { db, err := s.analyticsDB() if err != nil { @@ -46,6 +112,209 @@ func (s *AnalyticsService) liveConversationMetrics(ctx context.Context, accountI return &result, nil } +func (s *AnalyticsService) conversationCountTimeseries(ctx context.Context, accountID uint, since, until time.Time, reportType string, id uint, groupBy string, resolved bool) ([]AnalyticsTimeseriesPoint, error) { + db, err := s.analyticsDB() + if err != nil { + return nil, err + } + var conversations []model.Conversation + q := db.WithContext(ctx).Model(&model.Conversation{}).Where("account_id = ?", accountID) + if resolved { + q = q.Where("resolved_at IS NOT NULL AND resolved_at >= ? AND resolved_at < ?", since, until) + } else { + q = q.Where("created_at >= ? AND created_at < ?", since, until) + } + q = applyConversationDimension(q, reportType, id) + if err := q.Find(&conversations).Error; err != nil { + return nil, err + } + buckets := map[time.Time]*AnalyticsTimeseriesPoint{} + for _, conversation := range conversations { + t := conversation.CreatedAt + if resolved && conversation.ResolvedAt != nil { + t = *conversation.ResolvedAt + } + bucket := bucketStart(t, groupBy) + point := ensureTimeseriesBucket(buckets, bucket) + point.Value++ + } + return sortedTimeseries(buckets), nil +} + +func (s *AnalyticsService) messageCountTimeseries(ctx context.Context, accountID uint, since, until time.Time, reportType string, id uint, groupBy string, messageType model.MessageType) ([]AnalyticsTimeseriesPoint, error) { + db, err := s.analyticsDB() + if err != nil { + return nil, err + } + var messages []model.Message + q := db.WithContext(ctx).Model(&model.Message{}).Where("messages.account_id = ? AND messages.message_type = ? AND messages.created_at >= ? AND messages.created_at < ?", accountID, string(messageType), since, until) + q = applyMessageDimension(q, reportType, id) + if err := q.Find(&messages).Error; err != nil { + return nil, err + } + buckets := map[time.Time]*AnalyticsTimeseriesPoint{} + for _, message := range messages { + point := ensureTimeseriesBucket(buckets, bucketStart(message.CreatedAt, groupBy)) + point.Value++ + } + return sortedTimeseries(buckets), nil +} + +func (s *AnalyticsService) eventTimeseries(ctx context.Context, accountID uint, since, until time.Time, reportType string, id uint, groupBy string, names []string, average bool, businessHours bool) ([]AnalyticsTimeseriesPoint, error) { + db, err := s.analyticsDB() + if err != nil { + return nil, err + } + var events []model.ReportingEvent + q := db.WithContext(ctx).Model(&model.ReportingEvent{}).Where("account_id = ? AND name IN ? AND created_at >= ? AND created_at < ?", accountID, names, since, until) + q = applyEventDimension(q, reportType, id) + if err := q.Find(&events).Error; err != nil { + return nil, err + } + buckets := map[time.Time]*AnalyticsTimeseriesPoint{} + for _, event := range events { + point := ensureTimeseriesBucket(buckets, bucketStart(event.CreatedAt, groupBy)) + value := event.Value + if businessHours { + value = event.ValueInBusinessHours + } + if average { + point.Value += value + point.Count++ + } else { + point.Value++ + } + } + if average { + for _, point := range buckets { + if point.Count > 0 { + point.Value = point.Value / float64(point.Count) + } + } + } + return sortedTimeseries(buckets), nil +} + +func applyConversationDimension(query *gorm.DB, reportType string, id uint) *gorm.DB { + switch reportType { + case "inbox": + if id > 0 { + query = query.Where("inbox_id = ?", id) + } + case "agent": + if id > 0 { + query = query.Where("assignee_id = ?", id) + } + case "team": + if id > 0 { + query = query.Where("team_id = ?", id) + } + case "label": + if id > 0 { + query = query.Joins("INNER JOIN conversation_labels ON conversation_labels.conversation_id = conversations.id").Where("conversation_labels.tag_id = ?", id) + } + } + return query +} + +func applyMessageDimension(query *gorm.DB, reportType string, id uint) *gorm.DB { + switch reportType { + case "inbox": + if id > 0 { + query = query.Where("messages.inbox_id = ?", id) + } + case "agent": + if id > 0 { + query = query.Where("messages.sender_type = ? AND messages.sender_id = ?", "User", id) + } + case "team": + if id > 0 { + query = query.Joins("INNER JOIN conversations ON conversations.id = messages.conversation_id").Where("conversations.team_id = ?", id) + } + case "label": + if id > 0 { + query = query.Joins("INNER JOIN conversation_labels ON conversation_labels.conversation_id = messages.conversation_id").Where("conversation_labels.tag_id = ?", id) + } + } + return query +} + +func applyEventDimension(query *gorm.DB, reportType string, id uint) *gorm.DB { + switch reportType { + case "inbox": + if id > 0 { + query = query.Where("inbox_id = ?", id) + } + case "agent": + if id > 0 { + query = query.Where("user_id = ?", id) + } + case "team": + if id > 0 { + query = query.Joins("INNER JOIN conversations ON conversations.id = reporting_events.conversation_id").Where("conversations.team_id = ?", id) + } + case "label": + if id > 0 { + query = query.Joins("INNER JOIN conversation_labels ON conversation_labels.conversation_id = reporting_events.conversation_id").Where("conversation_labels.tag_id = ?", id) + } + } + return query +} + +func dayStart(t time.Time) time.Time { + return time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, time.UTC) +} + +func bucketStart(t time.Time, groupBy string) time.Time { + t = t.UTC() + switch groupBy { + case "hour": + return time.Date(t.Year(), t.Month(), t.Day(), t.Hour(), 0, 0, 0, time.UTC) + case "week": + start := dayStart(t) + return start.AddDate(0, 0, -int(start.Weekday())) + case "month": + return time.Date(t.Year(), t.Month(), 1, 0, 0, 0, 0, time.UTC) + case "year": + return time.Date(t.Year(), 1, 1, 0, 0, 0, 0, time.UTC) + default: + return dayStart(t) + } +} + +func ensureTimeseriesBucket(buckets map[time.Time]*AnalyticsTimeseriesPoint, bucket time.Time) *AnalyticsTimeseriesPoint { + point := buckets[bucket] + if point == nil { + point = &AnalyticsTimeseriesPoint{Timestamp: bucket.Unix()} + buckets[bucket] = point + } + return point +} + +func sortedTimeseries(buckets map[time.Time]*AnalyticsTimeseriesPoint) []AnalyticsTimeseriesPoint { + keys := make([]time.Time, 0, len(buckets)) + for key := range buckets { + keys = append(keys, key) + } + sort.Slice(keys, func(i, j int) bool { return keys[i].Before(keys[j]) }) + result := make([]AnalyticsTimeseriesPoint, 0, len(keys)) + for _, key := range keys { + result = append(result, *buckets[key]) + } + return result +} + +func parseReportDimensionID(raw string) uint { + if strings.TrimSpace(raw) == "" { + return 0 + } + parsed, err := strconv.ParseUint(raw, 10, 64) + if err != nil { + return 0 + } + return uint(parsed) +} + func (s *AnalyticsService) groupedLiveConversationMetrics(ctx context.Context, accountID uint, groupBy string) ([]map[string]interface{}, error) { if groupBy != "team_id" && groupBy != "assignee_id" { return nil, errors.New("invalid group_by") diff --git a/internal/service/analytics_service.go b/internal/service/analytics_service.go index 888c07a3..1baa33db 100644 --- a/internal/service/analytics_service.go +++ b/internal/service/analytics_service.go @@ -6,6 +6,7 @@ import ( "github.com/gochat/gochat/internal/model" "github.com/gochat/gochat/internal/repository" + "github.com/gochat/gochat/internal/worker" applogger "github.com/gochat/gochat/pkg/logger" "gorm.io/gorm" ) @@ -16,6 +17,8 @@ type AnalyticsService struct { eventRepo *repository.ReportingEventRepo rollupRepo *repository.ReportingEventsRollupRepo db *gorm.DB + rollups *ReportingRollupService + worker *worker.WorkerPool } func NewAnalyticsService( @@ -33,9 +36,15 @@ func NewAnalyticsService( eventRepo: eventRepo, rollupRepo: rollupRepo, db: db, + rollups: NewReportingRollupService(eventRepo, rollupRepo), } } +func (s *AnalyticsService) SetWorkerPool(wp *worker.WorkerPool) { + s.worker = wp + RegisterReportingRollupJobs(wp, s) +} + // --- Response DTOs --- // MetricSummary holds aggregated metric data. @@ -84,6 +93,9 @@ type GroupedConversationMetric struct { // GetSummary returns account-level aggregated metrics for a date range. func (s *AnalyticsService) GetSummary(ctx context.Context, accountID uint, since, until time.Time) (*SummaryResponse, error) { + if err := s.EnsureRollupsForRange(ctx, accountID, since, until); err != nil { + return nil, err + } rollups, err := s.rollupRepo.FindByAccountAndDateRange(ctx, accountID, since, until) if err != nil { applogger.L().Errorf("GetSummary rollup query: %v", err) @@ -143,6 +155,9 @@ func (s *AnalyticsService) GetTeamMetrics(ctx context.Context, accountID uint, s } func (s *AnalyticsService) getDimensionMetrics(ctx context.Context, accountID uint, dimType model.DimensionType, since, until time.Time) ([]DimensionMetrics, error) { + if err := s.EnsureRollupsForRange(ctx, accountID, since, until); err != nil { + return nil, err + } rollups, err := s.rollupRepo.AggregateSummary(ctx, accountID, dimType, since, until) if err != nil { applogger.L().Errorf("getDimensionMetrics query: %v", err) @@ -192,6 +207,9 @@ func (s *AnalyticsService) getDimensionMetrics(ctx context.Context, accountID ui // GetConversationTraffic returns daily conversation count time-series. func (s *AnalyticsService) GetConversationTraffic(ctx context.Context, accountID uint, since, until time.Time) ([]ConversationTrafficPoint, error) { + if err := s.EnsureRollupsForRange(ctx, accountID, since, until); err != nil { + return nil, err + } rollups, err := s.rollupRepo.FindByMetric(ctx, accountID, model.MetricResolutionsCount, since, until) if err != nil { applogger.L().Errorf("GetConversationTraffic query: %v", err) diff --git a/internal/service/reporting_rollup_worker.go b/internal/service/reporting_rollup_worker.go new file mode 100644 index 00000000..cde01803 --- /dev/null +++ b/internal/service/reporting_rollup_worker.go @@ -0,0 +1,62 @@ +package service + +import ( + "context" + "encoding/json" + "fmt" + "sync" + "time" + + "github.com/gochat/gochat/internal/model" + "github.com/gochat/gochat/internal/worker" +) + +const TaskTypeReportingRollupDay = "reporting:rollup_day" + +type reportingRollupDayJob struct { + AccountID uint `json:"account_id"` + Date string `json:"date"` +} + +var reportingRollupRegistrations sync.Map + +func RegisterReportingRollupJobs(wp *worker.WorkerPool, svc *AnalyticsService) { + if wp == nil || svc == nil { + return + } + if _, loaded := reportingRollupRegistrations.LoadOrStore(wp, struct{}{}); loaded { + return + } + wp.Register(TaskTypeReportingRollupDay, svc.performReportingRollupDayJob) +} + +func EnqueueReportingRollupDay(ctx context.Context, wp *worker.WorkerPool, accountID uint, date time.Time) (*model.BackgroundJob, error) { + if wp == nil || accountID == 0 { + return nil, nil + } + day := time.Date(date.Year(), date.Month(), date.Day(), 0, 0, 0, 0, time.UTC) + dateStr := day.Format(time.DateOnly) + return wp.Enqueue(ctx, TaskTypeReportingRollupDay, reportingRollupDayJob{AccountID: accountID, Date: dateStr}, + worker.WithQueue("low"), + worker.WithMaxAttempts(3), + worker.WithIdempotencyKey(fmt.Sprintf("reporting:rollup_day:%d:%s", accountID, dateStr)), + ) +} + +func (s *AnalyticsService) performReportingRollupDayJob(ctx context.Context, job *model.BackgroundJob) error { + var payload reportingRollupDayJob + if err := json.Unmarshal(job.Payload, &payload); err != nil { + return fmt.Errorf("unmarshal reporting rollup job: %w", err) + } + if payload.AccountID == 0 || payload.Date == "" { + return fmt.Errorf("invalid reporting rollup payload: %#v", payload) + } + date, err := time.Parse(time.DateOnly, payload.Date) + if err != nil { + return fmt.Errorf("parse reporting rollup date: %w", err) + } + if s == nil || s.rollups == nil { + return fmt.Errorf("reporting rollup service is required") + } + return s.rollups.ComputeDailyRollup(ctx, payload.AccountID, date) +}