Adding multithreading to the application
This commit is contained in:
@@ -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"`).
|
both hold the ticket number directly (e.g. both were `"AR-160269"`).
|
||||||
"Pull Tracking Numbers":
|
"Pull Tracking Numbers":
|
||||||
|
|
||||||
1. Fetches today's labels via `GET /v2/labels` - which gives
|
1. Bulk-fetches **all** of today's labels and all shipments from the
|
||||||
`tracking_number` and, importantly, `is_return_label` directly, so
|
last `SHIPSTATION_LOOKBACK_DAYS` (a code constant in
|
||||||
telling a return label apart from an outgoing one needs no guessing.
|
`shipstation_service.py`, default 7 days - wide enough to catch a
|
||||||
2. Looks up each label's shipment (`GET /v2/shipments/{id}`) to read
|
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
|
the ticket number off `shipment_number` (checked first) or
|
||||||
`external_shipment_id`, falling back to scanning the whole payload
|
`external_shipment_id`, falling back to scanning the whole payload
|
||||||
with `TICKET_NUMBER_REGEX` if neither is present.
|
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
|
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),
|
(e.g. JIRA hasn't been imported yet, or it's an account outlier),
|
||||||
|
|||||||
@@ -1,20 +1,30 @@
|
|||||||
"""
|
"""
|
||||||
ShipStation tracking-number puller (API V2).
|
ShipStation tracking-number puller (API V2).
|
||||||
|
|
||||||
ShipStation is used for exactly one thing here: generating shipping
|
Performance note: the previous version did one GET /v2/shipments/{id}
|
||||||
labels for JIRA tickets. It's not an independent order source, so this
|
call per unique shipment referenced by today's labels - an N+1 pattern
|
||||||
doesn't produce its own order rows - it produces tracking numbers keyed
|
that meant a busy day (100+ tickets) meant 100+ sequential network
|
||||||
by ticket number, which app.workers merges directly onto the matching
|
round-trips just for shipment lookups. That was almost certainly the
|
||||||
JIRA-sourced Order row. See app/tracking.py for the "what JIRA status
|
actual bottleneck.
|
||||||
should this become" logic.
|
|
||||||
|
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
|
How this maps to the real API (confirmed against a live payload, not
|
||||||
guessed):
|
guessed):
|
||||||
- GET /v2/labels gives tracking_number + is_return_label directly -
|
- GET /v2/labels gives tracking_number + is_return_label directly -
|
||||||
exactly what's needed to tell "Waiting For Return" apart from
|
exactly what's needed to tell "Waiting For Return" apart from
|
||||||
"Device Return Not Needed".
|
"Device Return Not Needed".
|
||||||
- A label only carries a shipment_id, not the ticket number, so for
|
- A label only carries a shipment_id, not the ticket number. A real
|
||||||
each label we fetch its shipment via GET /v2/shipments/{id}. A real
|
|
||||||
shipment payload showed the ticket number in BOTH shipment_number
|
shipment payload showed the ticket number in BOTH shipment_number
|
||||||
and external_shipment_id (e.g. both were "AR-160269") - we check
|
and external_shipment_id (e.g. both were "AR-160269") - we check
|
||||||
both, then fall back to scanning the whole payload with
|
both, then fall back to scanning the whole payload with
|
||||||
@@ -28,9 +38,11 @@ from __future__ import annotations
|
|||||||
import datetime as dt
|
import datetime as dt
|
||||||
import json
|
import json
|
||||||
import re
|
import re
|
||||||
from typing import Dict, List, Optional
|
import threading
|
||||||
|
from typing import Callable, Dict, List, Optional
|
||||||
|
|
||||||
import requests
|
import requests
|
||||||
|
from PyQt6.QtCore import QRunnable, QThreadPool
|
||||||
|
|
||||||
from app import config
|
from app import config
|
||||||
from app.services.base import OrderService, NormalizedOrder
|
from app.services.base import OrderService, NormalizedOrder
|
||||||
@@ -40,11 +52,45 @@ API_BASE = "https://api.shipstation.com/v2"
|
|||||||
PAGE_SIZE = 100
|
PAGE_SIZE = 100
|
||||||
REQUEST_TIMEOUT_SECONDS = 30
|
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):
|
class ShipStationServiceError(Exception):
|
||||||
"""Raised for any ShipStation fetch failure, with a message safe to show in the UI."""
|
"""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):
|
class ShipStationService(OrderService):
|
||||||
name = "shipstation"
|
name = "shipstation"
|
||||||
|
|
||||||
@@ -66,25 +112,19 @@ class ShipStationService(OrderService):
|
|||||||
"ShipStation is not configured yet. Open Settings and fill in the API Key."
|
"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 = [
|
usable_labels = [
|
||||||
label
|
label
|
||||||
for label in labels
|
for label in labels
|
||||||
if not label.get("voided") and label.get("tracking_number")
|
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
|
# --- Sort/correlate now, entirely in memory, after both batches landed ---
|
||||||
# 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}
|
|
||||||
|
|
||||||
tracking_by_ticket: Dict[str, List[dict]] = {}
|
tracking_by_ticket: Dict[str, List[dict]] = {}
|
||||||
raw_by_ticket: Dict[str, 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:
|
for label in usable_labels:
|
||||||
shipment = shipments_by_id.get(label.get("shipment_id"))
|
shipment = shipments_by_id.get(label.get("shipment_id"))
|
||||||
@@ -111,42 +151,84 @@ class ShipStationService(OrderService):
|
|||||||
for ticket_number in tracking_by_ticket
|
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_start = dt.datetime.combine(dt.date.today(), dt.time.min)
|
||||||
today_end = today_start + dt.timedelta(days=1)
|
today_end = today_start + dt.timedelta(days=1)
|
||||||
|
shipment_start = today_start - dt.timedelta(days=SHIPMENT_LOOKBACK_DAYS)
|
||||||
|
|
||||||
all_labels: List[dict] = []
|
label_params = {
|
||||||
page = 1
|
"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:
|
# Step 1: page 1 of each, fetched concurrently, to learn total page counts.
|
||||||
params = {
|
first_pages = self._run_concurrently(
|
||||||
"created_at_start": today_start.isoformat() + "Z",
|
{
|
||||||
"created_at_end": today_end.isoformat() + "Z",
|
"labels": lambda: self._get("/labels", {**label_params, "page": 1}),
|
||||||
"page": page,
|
"shipments": lambda: self._get("/shipments", {**shipment_params, "page": 1}),
|
||||||
"page_size": PAGE_SIZE,
|
|
||||||
"sort_by": "created_at",
|
|
||||||
"sort_dir": "desc",
|
|
||||||
}
|
}
|
||||||
data = self._get("/labels", params)
|
)
|
||||||
batch = data.get("labels", [])
|
labels = list(first_pages["labels"].get("labels", []))
|
||||||
all_labels.extend(batch)
|
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)
|
# Step 2: every remaining page across BOTH endpoints, fetched
|
||||||
if page >= total_pages or not batch:
|
# concurrently together (not one endpoint at a time).
|
||||||
break
|
remaining_jobs: Dict[str, Callable[[], dict]] = {}
|
||||||
page += 1
|
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]:
|
return labels, shipments
|
||||||
try:
|
|
||||||
return self._get(f"/shipments/{shipment_id}")
|
def _run_concurrently(self, jobs: Dict[str, Callable[[], dict]]) -> Dict[str, dict]:
|
||||||
except ShipStationServiceError:
|
"""Runs each zero-arg callable on its own QThreadPool worker thread,
|
||||||
# Don't let one bad lookup fail the whole import - this label's
|
blocks until all finish, and returns {key: result}. If anything
|
||||||
# tracking number just won't get matched to a ticket this run.
|
failed, raises the first error encountered (with its job key)."""
|
||||||
return None
|
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:
|
def _get(self, path: str, params: Optional[dict] = None) -> dict:
|
||||||
try:
|
try:
|
||||||
@@ -157,7 +239,7 @@ class ShipStationService(OrderService):
|
|||||||
timeout=REQUEST_TIMEOUT_SECONDS,
|
timeout=REQUEST_TIMEOUT_SECONDS,
|
||||||
)
|
)
|
||||||
except requests.RequestException as exc:
|
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:
|
if response.status_code == 401:
|
||||||
raise ShipStationServiceError(
|
raise ShipStationServiceError(
|
||||||
@@ -176,6 +258,8 @@ class ShipStationService(OrderService):
|
|||||||
f"ShipStation returned a response that wasn't valid JSON for {path}."
|
f"ShipStation returned a response that wasn't valid JSON for {path}."
|
||||||
) from exc
|
) from exc
|
||||||
|
|
||||||
|
# -- correlation ------------------------------------------------------
|
||||||
|
|
||||||
def _extract_ticket_number(self, shipment: dict) -> Optional[str]:
|
def _extract_ticket_number(self, shipment: dict) -> Optional[str]:
|
||||||
# Confirmed against a real payload: both of these can carry the
|
# Confirmed against a real payload: both of these can carry the
|
||||||
# ticket number directly. Check the more purpose-built field first.
|
# ticket number directly. Check the more purpose-built field first.
|
||||||
|
|||||||
Reference in New Issue
Block a user