Files
Order-Manager/app/workers.py
T

307 lines
11 KiB
Python

"""
Background work that shouldn't run on the GUI thread.
FetchOrdersWorker runs a service's fetch_orders() call off the main
thread and reports back via a signal. As we add more long-running
operations (pushing to Odoo, etc.) they should follow this same
pattern rather than blocking the UI.
"""
from __future__ import annotations
import datetime as dt
from typing import List, Tuple, TypedDict
from PyQt6.QtCore import QThread, pyqtSignal
from sqlalchemy import select
from app.database import get_session
from app.models import Order
from app.schedule import get_cutoff_time, is_past_cutoff
from app.services.base import OrderService, NormalizedOrder
from app.status_rules import get_active_statuses, get_cancelled_statuses, status_in
from app.tracking import get_fulfilled_statuses
class StatusChange(TypedDict):
ticket_number: str
old_status: str
new_status: str
class SaveResult(TypedDict):
new_count: int
updated_count: int
status_changes: List[StatusChange]
enriched_count: int # tickets that got tracking numbers merged in
unmatched_tracking_tickets: List[str] # tracking data with no matching local ticket
class SendToShipStationWorker(QThread):
"""Runs the emergency ShipStation API send off the GUI thread."""
finished_ok = pyqtSignal(dict) # the created shipment's JSON
failed = pyqtSignal(str)
def __init__(self, order, parent=None):
super().__init__(parent)
self.order = order
def run(self) -> None:
from app.services.shipstation_send import send_order_to_shipstation_api, ShipStationSendError
try:
result = send_order_to_shipstation_api(self.order)
except ShipStationSendError as exc:
self.failed.emit(str(exc))
return
except Exception as exc: # noqa: BLE001
self.failed.emit(f"Unexpected error sending to ShipStation: {exc}")
return
self.finished_ok.emit(result)
class FetchOrdersWorker(QThread):
"""Fetches orders from a given service and saves/merges the results."""
finished_ok = pyqtSignal(object) # emits a SaveResult
failed = pyqtSignal(str)
def __init__(self, service: OrderService, parent=None):
super().__init__(parent)
self.service = service
def run(self) -> None:
try:
orders = self.service.fetch_orders()
except Exception as exc: # noqa: BLE001 - surface any failure to the UI
self.failed.emit(str(exc))
return
try:
result = save_orders(orders)
except Exception as exc: # noqa: BLE001
self.failed.emit(f"Fetched {len(orders)} orders but failed to save them: {exc}")
return
self.finished_ok.emit(result)
def save_orders(orders: List[NormalizedOrder]) -> SaveResult:
"""
JIRA-sourced orders are inserted/updated as usual. ShipStation-sourced
"orders" are actually just tracking-number bundles keyed by ticket
number - rather than creating a second row, they get merged onto the
existing JIRA row for that ticket. If no such row exists locally yet,
the ticket number is reported back as unmatched instead of silently
dropped.
Whenever a JIRA ticket's status transitions INTO a fulfilled status
(Waiting For Return / Device Return Not Needed), fulfilled_at is
stamped - that's what lets the dashboard show "fulfilled today" and
what moves it into the Done pile.
"""
session = get_session()
new_count = 0
updated_count = 0
enriched_count = 0
status_changes: List[StatusChange] = []
unmatched_tracking_tickets: List[str] = []
fulfilled_statuses = get_fulfilled_statuses()
try:
for order in orders:
if order["source"] == "shipstation":
jira_row = session.execute(
select(Order).where(
Order.source == "jira",
Order.ticket_number == order["ticket_number"],
)
).scalar_one_or_none()
if jira_row is None:
unmatched_tracking_tickets.append(order["ticket_number"])
continue
jira_row.tracking_numbers = order.get("tracking_numbers", [])
enriched_count += 1
continue
# source == "jira"
existing = session.execute(
select(Order).where(
Order.source == order["source"],
Order.external_id == order["external_id"],
)
).scalar_one_or_none()
new_status = order["status"]
if existing is None:
session.add(
Order(
source=order["source"],
external_id=order["external_id"],
ticket_number=order.get("ticket_number"),
company=order.get("company", "Unknown"),
skus=order.get("skus", []),
line_items=order.get("line_items", []),
shipping_info=order.get("shipping_info", {}),
creator=order.get("creator"),
summary=order["summary"],
status=new_status,
source_created_at=order["source_created_at"],
raw_data=order["raw_data"],
fulfilled_at=(
dt.datetime.utcnow() if status_in(new_status, fulfilled_statuses) else None
),
)
)
new_count += 1
else:
old_status = existing.status
if old_status != new_status:
status_changes.append(
StatusChange(
ticket_number=existing.ticket_number or existing.external_id,
old_status=old_status,
new_status=new_status,
)
)
if status_in(new_status, fulfilled_statuses) and not status_in(
old_status, fulfilled_statuses
):
existing.fulfilled_at = dt.datetime.utcnow()
existing.ticket_number = order.get("ticket_number")
existing.company = order.get("company", "Unknown")
existing.skus = order.get("skus", [])
existing.line_items = order.get("line_items", [])
existing.shipping_info = order.get("shipping_info", {})
existing.creator = order.get("creator")
existing.summary = order["summary"]
existing.status = new_status
existing.source_created_at = order["source_created_at"]
existing.raw_data = order["raw_data"]
updated_count += 1
session.commit()
finally:
session.close()
return SaveResult(
new_count=new_count,
updated_count=updated_count,
status_changes=status_changes,
enriched_count=enriched_count,
unmatched_tracking_tickets=unmatched_tracking_tickets,
)
def load_all_orders() -> List[Order]:
session = get_session()
try:
return list(
session.execute(select(Order).order_by(Order.source_created_at.desc())).scalars()
)
finally:
session.close()
def split_orders_by_view(orders: List[Order]) -> Tuple[List[Order], List[Order], List[Order]]:
"""
Three tabs, allowlist-driven:
- Active: status is in ACTIVE_STATUSES (just "Created" by default) -
the only tickets that represent real work still to do.
- Cancelled: status is in CANCELLED_STATUSES - its own tab so it
doesn't clutter Active, but still reviewable on demand.
- Done: everything else. This deliberately doesn't enumerate every
"finished" status by name - a JIRA-side automation status like
"1st Contact Attempt" falls in here automatically just by not
being Created or Cancelled, with no code change needed when your
JIRA workflow adds another downstream status later.
"""
active_statuses = get_active_statuses()
cancelled_statuses = get_cancelled_statuses()
active, cancelled, done = [], [], []
for order in orders:
if status_in(order.status, active_statuses):
active.append(order)
elif status_in(order.status, cancelled_statuses):
cancelled.append(order)
else:
done.append(order)
return active, cancelled, done
def load_orders_by_view() -> Tuple[List[Order], List[Order], List[Order]]:
return split_orders_by_view(load_all_orders())
def get_dashboard_stats() -> dict:
"""
Counts for the dashboard. "Active Orders" and the company breakdown
reflect only the Active tab (status in ACTIVE_STATUSES) - the actual
at-hand workload. Cancelled counts the Cancelled tab. Fulfilled-today
and past-cutoff-today look across ALL of today's tickets regardless
of which tab they ended up in, since both are about what happened
today specifically.
Carryover counts any Active ticket that's already known to spill
into tomorrow - either it's genuinely left over from a prior day, or
it arrived today but after the cutoff (same effect, just known a day
earlier). Since we're only looking at the Active bucket, every
ticket here is by definition still unresolved - no separate "is it
still open" check needed.
"""
orders = load_all_orders()
active_orders, cancelled_orders, _done_orders = split_orders_by_view(orders)
cutoff = get_cutoff_time()
today = dt.date.today()
by_company: dict[str, int] = {}
carryover_count = 0
tracking_received_count = 0
for order in active_orders:
by_company[order.company] = by_company.get(order.company, 0) + 1
if order.tracking_numbers:
tracking_received_count += 1
created_before_today = bool(
order.source_created_at and order.source_created_at.date() < today
)
# A ticket that arrived today but after the cutoff won't get worked
# today either - it's known carryover before the date even rolls
# over, not just once tomorrow arrives.
arrived_past_cutoff_today = bool(
order.source_created_at
and order.source_created_at.date() == today
and is_past_cutoff(order.source_created_at, cutoff)
)
if created_before_today or arrived_past_cutoff_today:
carryover_count += 1
fulfilled_today_count = sum(
1 for o in orders if o.fulfilled_at and o.fulfilled_at.date() == today
)
past_cutoff_today_count = sum(
1
for o in orders
if o.source_created_at
and o.source_created_at.date() == today
and is_past_cutoff(o.source_created_at, cutoff)
)
return {
"total": len(active_orders),
"by_company": by_company,
"cancelled_count": len(cancelled_orders),
"carryover_count": carryover_count,
"tracking_received_count": tracking_received_count,
"fulfilled_today_count": fulfilled_today_count,
"past_cutoff_today_count": past_cutoff_today_count,
}