# Implementation architecture

## Runtime boundary

n8n controls scheduling, stage orchestration, retry, recovery and failure routing. Each named HTTP node invokes one real operation in `service/app.py`. The included service isolates algorithms and provider protocols from n8n node-version changes. Code nodes do not make hidden network calls. SQLite keeps raw input, normalized records, snapshots, scoped reports and delivery state outside n8n execution retention.

```mermaid
flowchart TD
  A[Monday 08:00 or manual] --> B[Load configuration and claim period]
  B --> C[HubSpot records, properties, owners and history]
  C --> D[Completeness, validation, deduplication and FX]
  D --> E[Scope audiences, deterministic metrics and prior snapshot]
  E --> F[Persist pre-AI report]
  F --> G[Optional constrained AI with fallback]
  G --> H[Local PNG charts and channel rendering]
  H --> I[Slack root, detail thread and chart uploads]
  I --> J[Accessible MIME email with stable threading]
  J --> K[Audit delivery outcomes]
  L[Hourly recovery] --> I
  K --> M[Admin alert for incomplete delivery]
  N[Production Error Trigger] --> M
```

Slack and email are sequential n8n stages but have independent delivery outcomes. A Slack provider failure is recorded per destination and returns control so email still runs. A service/infrastructure failure stops the chain; after restoration the recovery workflow independently retries both channels. Scale-out requires an external job queue and transactional PostgreSQL leases; do not run multiple copies of this SQLite service.

## Node-by-node specification

All stage calls use POST, authenticated service credentials, a 300-second request timeout, three attempts and five seconds between n8n retries. The service checkpoints completed stages, making HTTP stage retries safe. The service's provider GET retry logic uses exponential backoff plus jitter and respects Retry-After; it fails/defer-retries when a delay exceeds 60 seconds. It does not retry authentication failures as transient errors.

| Workflow/node | Type/version | Input → output | Implementation and failure behavior |
|---|---|---|---|
| 01 / Manual run | Manual Trigger 1 | one trigger item | Uses current previous-week period; no arbitrary backfill |
| 01 / Scheduled run | Schedule Trigger 1.4 | Monday 08:00 | Workflow timezone `Asia/Tbilisi`, replace for workspace |
| Load config and claim period | HTTP Request 4.5 | `{}` → `{run_id,dry_run}` | `/v1/start`; validates server-side config; immutable per-period config hash; one run key per workspace/week |
| HubSpot records and history | HTTP Request 4.5 | `{run_id}` → same run ID | `/v1/collect`; complete paginated active-deal scan plus schemas, pipelines and active/archived owners; raw input stored only after success |
| Validate deduplicate and FX | HTTP Request 4.5 | `{run_id}` → stage result | `/v1/normalize`; fails closed for critical data/FX errors; duplicates choose newest timestamp; conflicting equal-time duplicates block |
| Metrics snapshots and audience scopes | HTTP Request 4.5 | `{run_id}` → stage result | `/v1/analyze`; owner scope first, metrics, immediately preceding-week snapshot, movements and classifiers; persist before AI |
| AI validation fallback and charts | HTTP Request 4.5 | `{run_id}` → stage result | `/v1/prepare`; per-audience evidence ranking; schema/ID validation; deterministic fallback; locally rendered PNG; optional scoped dashboard export |
| Slack digest thread and charts | HTTP Request 4.5 | `{run_id}` → per-destination outcomes | `/v1/slack`; root `ts`, bounded detail replies, upload allocation/bytes/completion persisted separately |
| Threaded accessible email | HTTP Request 4.5 | `{run_id}` → per-destination outcomes | `/v1/email`; one mailbox per audience, CID image, text alternative, stable subject, deterministic Message-ID and persisted References |
| Finalize audit and alert | HTTP Request 4.5 | `{run_id}` → complete/dry_run/incomplete | `/v1/finalize`; verifies expected destination outcomes and component ledger; incomplete → admin alert |
| 02 / Scheduled or manual recovery | Schedule 1.4 / Manual 1 | `{}` → recovery results | `/v1/recover`; up to 20 unfinished prepared runs; failed components only; never re-fetches financial data |
| 03 / Production failure | Error Trigger 1 | n8n error envelope | Link as shared error workflow in workflow settings |
| Record failure and alert administrator | HTTP Request 4.5 | minimal run/stage info → acknowledgment | `/v1/alert`; excludes raw error stack/CRM payload from Slack; six-hour notification suppression |

Every successful HTTP call emits exactly one JSON object containing `run_id` where the next stage needs it. CRM arrays never fan out across the n8n canvas. Production Error Triggers do not cover manual tests. If the service itself is unreachable, its alert endpoint is also unreachable: independent infrastructure/dead-man monitoring is a deployment prerequisite.

## Implemented HubSpot adapter

`service/adapters.py::hubspot` calls actual documented endpoints:

| Endpoint | Purpose |
|---|---|
| `GET /crm/v3/properties/deals` | Verify available property names before querying |
| `GET /crm/v3/pipelines/deals` | Pipeline/stage IDs, display order, status metadata and probability |
| `GET /crm/v3/owners?archived=false/true` | Owner IDs and display names; no email retained in canonical records |
| `GET /crm/v3/objects/deals` | All non-archived deals using cursor pagination and requested property history |

Properties: `dealname`, `amount`, `deal_currency_code`, `pipeline`, `dealstage`, `hubspot_owner_id`, `createdate`, `closedate`, `notes_last_updated`, `hs_next_activity_date`, `closed_lost_reason`. History: `dealstage`, `closedate`, `amount`, `hubspot_owner_id`, `pipeline`. Optional properties are requested only if present. The portal's verified single currency must be configured when the currency property is absent.

The full scan avoids a filtered-query overlap problem and includes old open deals. It does not use dashboard screenshots, OCR, guessed dashboard endpoints or CRM Search's result ceiling. Page cap, repeated cursor, missing array, unsupported schema, timeout and unexpected empty account stop reporting. A known empty account requires `allow_empty:true`.

Boundaries: this is a modest-portal adapter (default 200 pages × 100 records and 240-second collection budget). HubSpot is not transactionally frozen during pagination. Archived/deleted deals are not included, and historical revenue for them is not reconstructed. `notes_last_updated` is HubSpot's activity summary, not a guarantee of meaningful commercial engagement. Engagement bodies/contact/company retrieval are intentionally omitted for minimization. A larger portal should implement an incremental mirror plus tombstones, reconciliation and a collection job queue before using this architecture at scale.

## Data and adapter contracts

The current HubSpot adapter returns `{complete,deals,pipelines,owners,started_at,observed_at,limitations}`. `normalize` outputs the provider-independent canonical records below. To add Salesforce/Pipedrive/etc., implement a separate fetch+normalization adapter that produces this canonical structure; add explicit routing in the service. Merely changing `crm.provider` currently fails validation.

```json
{
  "id": "123", "key": "acme:hubspot:123", "name": "Example deal",
  "owner_id": "101", "owner_name": "Example owner",
  "pipeline": "default", "stage": "proposal", "label": "Proposal",
  "rank": 2, "status": "open", "probability": "0.65",
  "value_original": "42000", "currency_original": "EUR",
  "fx_rate": "1.10", "fx_date": "2026-10-09", "value": "46200.00",
  "created_at": "2026-09-01T10:00:00+00:00",
  "updated_at": "2026-10-10T10:00:00+00:00",
  "close_at": "2026-10-20T12:00:00+00:00",
  "activity_at": null, "next_activity_at": null,
  "lost_reason": "", "flags": ["missing_activity", "missing_next_activity"],
  "history": {}
}
```

Critical: ID, nonnegative finite amount, currency/rate, known pipeline/stage/probability and required timestamps. Invalid critical records stop the digest rather than silently reduce totals. Closed records require a valid close date. Non-critical: missing/unknown owner, activity, future activity, open-deal close date and lost reason remain as visible flags. All sources are read-only.

## Deterministic definitions

| Metric | Definition |
|---|---|
| Window | Previous local Monday 00:00 inclusive through current Monday 00:00 exclusive; ZoneInfo handles DST |
| Open pipeline | Sum current observed open amounts; not period-end reconstruction |
| Weighted | Sum each normalized amount × configured HubSpot stage probability; no forecast guarantee |
| New pipeline | Deals created in window; includes deals subsequently closed; amount observed at collection |
| Won/lost value | Current won/lost status AND current CRM close date in window; CRM deal amount, not recognized accounting revenue |
| Stage/owner breakdown | Scoped open count/value/weighted; owners also include weekly new/won/lost amounts |
| Stalled | Open AND known last activity age ≥ threshold; missing activity remains unknown |
| Closing soon | Open close date between local today and today + configured days, inclusive; overdue kept separate |
| High value | Open normalized value ≥ configured threshold |
| Pushes | Sorted adjacent close-date history changes moving later; includes days moved and repeated pushes within the week |
| Other movement | Deduplicated timed history transitions for stage, owner, amount and pipeline; counts events, not unique deals |
| Snapshot movement | Net current-vs-prior scoped record changes, separately labeled; available when timed history is missing; never added to history counts |
| WoW | Immediately preceding weekly snapshot only; missing week → unavailable; zero prior value → absolute delta with null percentage |
| FX | Base units per original unit; dated table no later than period end and within configured maximum age; preserve originals and rates |
| Rounding | Decimal throughout; normalized values rounded to 2 decimal places, aggregate weighted totals rounded half-up |
| Concentration | Largest open deal / total open pipeline; no target/quota coverage invented |

FX policy uses period-end finance-approved rates for all cohorts. It is not transaction-date accounting FX. Rates must be refreshed every reporting period. Base currency should have two minor-unit decimals in this reference; extend `money()` for zero/three-decimal currencies before deployment. FX effects are included in WoW and labeled; constant-currency variance is an extension. Scope, base currency, analytics policy or stage definitions changing suppress the baseline. Owner transfers can legitimately change a rep's totals.

History is bounded by what the CRM returns. Repeated pushes are counted within the reporting week, not claimed as lifetime push counts. Stage order is compared within a resolved pipeline; an unresolved cross-pipeline transition remains a generic stage change. Newly entered scope is not automatically classified as newly created. Disappeared records are not labeled lost.

## AI contract and fallback

AI is disabled by default. `ai_rank` sends only scoped, deterministic aggregate fact strings and permitted action strings to OpenAI Chat Completions using a deployment-selected structured-output model. No raw CRM names, deal history or contacts are sent. The output is deliberately narrower than the blueprint's free-prose schema:

```json
{"fact_ids":["pipeline","stalled","weekly"],"action_ids":["stalled","hygiene"]}
```

Exact keys, array type, cardinality, uniqueness and membership are validated locally in addition to the provider JSON schema. The model cannot modify a numeric value, add a deal ID or produce arbitrary HTML. Timeout, refusal, malformed JSON, unknown IDs or unsupported model → deterministic fallback. Generated prose is selected from verified templates. This sacrifices free-form interpretation to make validation enforceable; unrestricted prose cannot be made factually safe by a JSON parser alone.

## Audience, chart and delivery contracts

Executive: workspace-wide metrics, commercial priorities and at most three flagged opportunities. Manager: explicit owner allowlist, owner/stage breakdown and up to twenty prioritized records. Rep: exactly one owner, scoped metrics/history/charts. One audience record per mailbox/channel; duplicated destinations across roles are rejected. Route reps to verified private DMs and managers to restricted channels; filtering does not secure a broadly visible Slack channel.

Local charts contain scoped aggregate values. Exact tables/plain text remain available if images are blocked. Slack uses plain-text Block Kit content to prevent CRM names from injecting mentions/markup; a brief parent, ≤2,800-character detail chunks and PNG files share the parent `thread_ts`. File delivery uses `files.getUploadURLExternal` → upload bytes → `files.completeUploadExternal`, not retired `files.upload`.

Email uses STARTTLS SMTP, multipart plain/HTML, semantic headings/table headers and inline CID PNG. Message-ID is a deterministic hash of run + destination. A thread identity includes workspace, audience scope and mailbox. Root/last message IDs persist; later reports use In-Reply-To and References plus a stable subject. Mail clients may still choose their own grouping. No Gmail/Microsoft provider-specific thread guarantee is claimed.

Optional dashboard export configuration specifies an exact authorized HTTPS PNG export URL, host allowlist, token environment name and audience ID. Redirects are not followed, content is capped at 8 MB and PNG magic is checked. It supplements the digest; it never feeds metrics. For exports requiring POST/job polling (e.g. some BI products), implement that provider's documented adapter. Do not paste a dashboard viewing URL and expect it to export. The included generic downloader does not imply a universal dashboard API.

## Storage and delivery state

SQLite tables: `runs` (workspace/period/config hash), `artifacts` (raw/normalized/pre-AI/prepared reports), `deliveries` (destination/component/attempt/status/response), `audit`, `mail_threads`. Primary keys prevent accidental duplicate run/component records. Fully prepared immutable reports support replay without recomputation. Snapshot success is independent of channel success.

Delivery states: absent → sending → sent; known rejection → failed; uncertain acceptance/crash → unknown. Recovery skips sent, retries failed up to five claims, and refuses unknown. Slack upload allocations are persisted so a successful root is not duplicated on later chart failure. Expired allocation URLs need explicit operator reconciliation/reset of the upload component sequence. SMTP network ambiguity needs provider-log investigation by deterministic Message-ID. There is no claim of distributed exactly-once semantics.

## Blueprint coverage and deliberate boundaries

Blueprint sections 1–35: config, adapter, normalize/validate, FX/time, metrics/history/snapshot implemented. Sections 36–43: constrained evidence-ranking AI with deterministic fallback implemented; unrestricted AI prose intentionally replaced. Sections 44–58: role scoping, Slack/email, per-component recovery, persistence and audit implemented. Sections 59–66: empty/zero/missing data behavior handled; other CRM, Teams/webhook output and adjacent lost-reason/follow-up automations remain explicit extension interfaces. Multi-workspace deployments use separate configs, service instances, secrets and volumes; this reference does not promise a multi-tenant credential broker.
