diff --git a/README.md b/README.md index 1285038..c071dc0 100644 --- a/README.md +++ b/README.md @@ -72,15 +72,23 @@ guessed): a shipment's `shipment_number` and `external_shipment_id` both hold the ticket number directly (e.g. both were `"AR-160269"`). "Pull Tracking Numbers": -1. Fetches today's labels via `GET /v2/labels` - which gives - `tracking_number` and, importantly, `is_return_label` directly, so - telling a return label apart from an outgoing one needs no guessing. -2. Looks up each label's shipment (`GET /v2/shipments/{id}`) to read +1. Bulk-fetches **all** of today's labels and all shipments from the + last `SHIPSTATION_LOOKBACK_DAYS` (a code constant in + `shipstation_service.py`, default 7 days - wide enough to catch a + shipment that sat "pending" for a day or two before its label was + generated). This replaced an earlier version that looked up each + shipment individually (one HTTP call per ticket) - the actual cause + of the slowness, now down to a handful of bulk list calls. +2. Fetches those (and any additional pages either needs) **concurrently** + via a PyQt `QThreadPool` - several requests in flight at once instead + of one after another. Capped at `MAX_CONCURRENT_REQUESTS` (default 5) + to stay well clear of any ShipStation rate limit. +3. Only after both full batches have landed does it sort/correlate + labels to shipments to ticket numbers, entirely in memory - reading + `tracking_number` and `is_return_label` straight off each label, and the ticket number off `shipment_number` (checked first) or `external_shipment_id`, falling back to scanning the whole payload with `TICKET_NUMBER_REGEX` if neither is present. -3. Groups tracking numbers by ticket number and merges them onto the - matching JIRA row. If a label's ticket number doesn't match any ticket you have locally (e.g. JIRA hasn't been imported yet, or it's an account outlier), diff --git a/app/services/shipstation_service.py b/app/services/shipstation_service.py index 3d712a9..a0ad415 100644 --- a/app/services/shipstation_service.py +++ b/app/services/shipstation_service.py @@ -1,20 +1,30 @@ """ ShipStation tracking-number puller (API V2). -ShipStation is used for exactly one thing here: generating shipping -labels for JIRA tickets. It's not an independent order source, so this -doesn't produce its own order rows - it produces tracking numbers keyed -by ticket number, which app.workers merges directly onto the matching -JIRA-sourced Order row. See app/tracking.py for the "what JIRA status -should this become" logic. +Performance note: the previous version did one GET /v2/shipments/{id} +call per unique shipment referenced by today's labels - an N+1 pattern +that meant a busy day (100+ tickets) meant 100+ sequential network +round-trips just for shipment lookups. That was almost certainly the +actual bottleneck. + +This version instead: + 1. Bulk-fetches ALL of today's labels and ALL recent shipments (a + wider window - see SHIPMENT_LOOKBACK_DAYS - since a shipment can + sit "pending" for a day or two before its label gets generated), + each as a small number of paginated list calls. + 2. Fetches both endpoints, and any additional pages either needs, + concurrently using a PyQt QThreadPool - "multiple workers" pulling + pages at once instead of one request waiting on the last. + 3. Only once both full batches are in memory does it sort/correlate + labels to shipments to ticket numbers - no more network calls + interleaved with processing. How this maps to the real API (confirmed against a live payload, not guessed): - GET /v2/labels gives tracking_number + is_return_label directly - exactly what's needed to tell "Waiting For Return" apart from "Device Return Not Needed". - - A label only carries a shipment_id, not the ticket number, so for - each label we fetch its shipment via GET /v2/shipments/{id}. A real + - A label only carries a shipment_id, not the ticket number. A real shipment payload showed the ticket number in BOTH shipment_number and external_shipment_id (e.g. both were "AR-160269") - we check both, then fall back to scanning the whole payload with @@ -28,9 +38,11 @@ from __future__ import annotations import datetime as dt import json import re -from typing import Dict, List, Optional +import threading +from typing import Callable, Dict, List, Optional import requests +from PyQt6.QtCore import QRunnable, QThreadPool from app import config from app.services.base import OrderService, NormalizedOrder @@ -40,11 +52,45 @@ API_BASE = "https://api.shipstation.com/v2" PAGE_SIZE = 100 REQUEST_TIMEOUT_SECONDS = 30 +# How far back to look for shipments, beyond just "today". Labels are +# fetched strictly for today (that's the definition of "today's tracking +# numbers"), but a shipment can be created a day or two before ShipStation's +# automation actually gets around to producing its label - if that window +# is too narrow, otherwise-matchable labels start showing up as unmatched. +SHIPMENT_LOOKBACK_DAYS = 7 + +# Cap on simultaneous ShipStation requests. Kept modest to stay well clear +# of any API rate limit rather than firing dozens of requests at once. +MAX_CONCURRENT_REQUESTS = 5 + class ShipStationServiceError(Exception): """Raised for any ShipStation fetch failure, with a message safe to show in the UI.""" +class _CallableTask(QRunnable): + """Runs a zero-arg callable on a QThreadPool worker thread and stashes + its result (or exception) into a shared, lock-protected dict/list.""" + + def __init__(self, key: str, fn: Callable[[], dict], results: dict, lock: threading.Lock, errors: list): + super().__init__() + self.key = key + self.fn = fn + self.results = results + self.lock = lock + self.errors = errors + + def run(self) -> None: + try: + value = self.fn() + except Exception as exc: # noqa: BLE001 - surfaced to the caller via `errors` + with self.lock: + self.errors.append(exc) + return + with self.lock: + self.results[self.key] = value + + class ShipStationService(OrderService): name = "shipstation" @@ -66,25 +112,19 @@ class ShipStationService(OrderService): "ShipStation is not configured yet. Open Settings and fill in the API Key." ) - labels = self._fetch_todays_labels() + labels, shipments = self._fetch_labels_and_shipments() - # Only labels that actually produced a usable tracking number matter. usable_labels = [ label for label in labels if not label.get("voided") and label.get("tracking_number") ] + shipments_by_id = {s["shipment_id"]: s for s in shipments if s.get("shipment_id")} - # One shipment lookup per unique shipment_id referenced, not per - # label (an outgoing + return label pair share the same shipment). - shipment_ids = { - label["shipment_id"] for label in usable_labels if label.get("shipment_id") - } - shipments_by_id = {sid: self._fetch_shipment(sid) for sid in shipment_ids} - + # --- Sort/correlate now, entirely in memory, after both batches landed --- tracking_by_ticket: Dict[str, List[dict]] = {} raw_by_ticket: Dict[str, dict] = {} - self.unmatched_labels = [] # labels we couldn't tie to a ticket number + self.unmatched_labels = [] for label in usable_labels: shipment = shipments_by_id.get(label.get("shipment_id")) @@ -111,42 +151,84 @@ class ShipStationService(OrderService): for ticket_number in tracking_by_ticket ] - # -- internals ----------------------------------------------------- + # -- batch fetching (concurrent) ------------------------------------ - def _fetch_todays_labels(self) -> List[dict]: + def _fetch_labels_and_shipments(self) -> tuple[List[dict], List[dict]]: today_start = dt.datetime.combine(dt.date.today(), dt.time.min) today_end = today_start + dt.timedelta(days=1) + shipment_start = today_start - dt.timedelta(days=SHIPMENT_LOOKBACK_DAYS) - all_labels: List[dict] = [] - page = 1 + label_params = { + "created_at_start": today_start.isoformat() + "Z", + "created_at_end": today_end.isoformat() + "Z", + "page_size": PAGE_SIZE, + "sort_by": "created_at", + "sort_dir": "desc", + } + shipment_params = { + "created_at_start": shipment_start.isoformat() + "Z", + "created_at_end": today_end.isoformat() + "Z", + "page_size": PAGE_SIZE, + "sort_by": "created_at", + "sort_dir": "desc", + } - while True: - params = { - "created_at_start": today_start.isoformat() + "Z", - "created_at_end": today_end.isoformat() + "Z", - "page": page, - "page_size": PAGE_SIZE, - "sort_by": "created_at", - "sort_dir": "desc", + # Step 1: page 1 of each, fetched concurrently, to learn total page counts. + first_pages = self._run_concurrently( + { + "labels": lambda: self._get("/labels", {**label_params, "page": 1}), + "shipments": lambda: self._get("/shipments", {**shipment_params, "page": 1}), } - data = self._get("/labels", params) - batch = data.get("labels", []) - all_labels.extend(batch) + ) + labels = list(first_pages["labels"].get("labels", [])) + shipments = list(first_pages["shipments"].get("shipments", [])) + labels_total_pages = first_pages["labels"].get("pages", 1) + shipments_total_pages = first_pages["shipments"].get("pages", 1) - total_pages = data.get("pages", 1) - if page >= total_pages or not batch: - break - page += 1 + # Step 2: every remaining page across BOTH endpoints, fetched + # concurrently together (not one endpoint at a time). + remaining_jobs: Dict[str, Callable[[], dict]] = {} + for page in range(2, labels_total_pages + 1): + remaining_jobs[f"labels:{page}"] = ( + lambda page=page: self._get("/labels", {**label_params, "page": page}) + ) + for page in range(2, shipments_total_pages + 1): + remaining_jobs[f"shipments:{page}"] = ( + lambda page=page: self._get("/shipments", {**shipment_params, "page": page}) + ) - return all_labels + if remaining_jobs: + remaining = self._run_concurrently(remaining_jobs) + for key, data in remaining.items(): + if key.startswith("labels:"): + labels.extend(data.get("labels", [])) + else: + shipments.extend(data.get("shipments", [])) - def _fetch_shipment(self, shipment_id: str) -> Optional[dict]: - try: - return self._get(f"/shipments/{shipment_id}") - except ShipStationServiceError: - # Don't let one bad lookup fail the whole import - this label's - # tracking number just won't get matched to a ticket this run. - return None + return labels, shipments + + def _run_concurrently(self, jobs: Dict[str, Callable[[], dict]]) -> Dict[str, dict]: + """Runs each zero-arg callable on its own QThreadPool worker thread, + blocks until all finish, and returns {key: result}. If anything + failed, raises the first error encountered (with its job key).""" + pool = QThreadPool() + pool.setMaxThreadCount(min(MAX_CONCURRENT_REQUESTS, max(1, len(jobs)))) + + results: Dict[str, dict] = {} + errors: List[Exception] = [] + lock = threading.Lock() + + for key, fn in jobs.items(): + pool.start(_CallableTask(key, fn, results, lock, errors)) + + pool.waitForDone() + + if errors: + raise errors[0] + + return results + + # -- single-request helper ------------------------------------------ def _get(self, path: str, params: Optional[dict] = None) -> dict: try: @@ -157,7 +239,7 @@ class ShipStationService(OrderService): timeout=REQUEST_TIMEOUT_SECONDS, ) except requests.RequestException as exc: - raise ShipStationServiceError(f"Could not reach ShipStation: {exc}") from exc + raise ShipStationServiceError(f"Could not reach ShipStation ({path}): {exc}") from exc if response.status_code == 401: raise ShipStationServiceError( @@ -176,6 +258,8 @@ class ShipStationService(OrderService): f"ShipStation returned a response that wasn't valid JSON for {path}." ) from exc + # -- correlation ------------------------------------------------------ + def _extract_ticket_number(self, shipment: dict) -> Optional[str]: # Confirmed against a real payload: both of these can carry the # ticket number directly. Check the more purpose-built field first.