Initial KDS scaffold — Phase 2 kitchen port
FastAPI backend (Python 3.11, httpx for SignalR/GraphQL — no MSSQL ODBC), shares kitchen_db directly. React/TS/Vite fullscreen board frontend. Backend: auth.py (APP_SLUG=kds, SimpleNamespace), main.py (4 KDS migrations, SignalR start/stop), kds.py router, models (kds/settings/resos — read from kitchen_db), signalr_listener.py (backoff pre-existing), kds_graphql.py, database.py. Requirements stripped to ~9 packages; image ~400 MB lighter than kitchen (no MSSQL ODBC layer). Frontend: AuthGate (app=kds), single fullscreen route, dark board theme. KDS.tsx URL prefix patched (/api/kds/ → /kds/api/kds/), recipe images cross-app (/kitchen/api/recipes/). nginx: 5 blocks with SSE proxy headers on /kds/api/ block. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
commit
b94585084a
35 changed files with 5195 additions and 0 deletions
0
backend/services/__init__.py
Normal file
0
backend/services/__init__.py
Normal file
382
backend/services/kds_graphql.py
Normal file
382
backend/services/kds_graphql.py
Normal file
|
|
@ -0,0 +1,382 @@
|
|||
"""
|
||||
SambaPOS GraphQL Client for KDS
|
||||
|
||||
Connects to SambaPOS Message Server GraphQL API to fetch open tickets
|
||||
with kitchen orders for the Kitchen Display System.
|
||||
"""
|
||||
|
||||
import logging
|
||||
from typing import Optional
|
||||
from datetime import datetime
|
||||
import httpx
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class SambaPOSGraphQLClient:
|
||||
"""Client for SambaPOS Message Server GraphQL API."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
server_url: str,
|
||||
username: str,
|
||||
password: str,
|
||||
client_id: str
|
||||
):
|
||||
self.server_url = server_url.rstrip('/')
|
||||
self.username = username
|
||||
self.password = password
|
||||
self.client_id = client_id
|
||||
self.access_token: Optional[str] = None
|
||||
self.token_expires_at: Optional[datetime] = None
|
||||
|
||||
async def authenticate(self) -> bool:
|
||||
"""
|
||||
Authenticate with SambaPOS and get access token.
|
||||
|
||||
POST /Token with:
|
||||
grant_type=password
|
||||
username=<user>
|
||||
password=<pass>
|
||||
client_id=<app_key>
|
||||
"""
|
||||
token_url = f"{self.server_url}/Token"
|
||||
|
||||
data = {
|
||||
"grant_type": "password",
|
||||
"username": self.username,
|
||||
"password": self.password,
|
||||
"client_id": self.client_id,
|
||||
}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
response = await client.post(
|
||||
token_url,
|
||||
data=data,
|
||||
headers={"Content-Type": "application/x-www-form-urlencoded"}
|
||||
)
|
||||
if response.status_code == 200:
|
||||
result = response.json()
|
||||
self.access_token = result.get("access_token")
|
||||
expires_in = result.get("expires_in", 86400)
|
||||
self.token_expires_at = datetime.utcnow()
|
||||
logger.info(f"KDS: Authenticated with SambaPOS (expires in {expires_in}s)")
|
||||
return True
|
||||
else:
|
||||
logger.error(f"KDS: Authentication failed: {response.status_code} - {response.text}")
|
||||
return False
|
||||
except httpx.ConnectError:
|
||||
logger.error(f"KDS: Could not connect to {token_url}")
|
||||
return False
|
||||
except Exception as e:
|
||||
logger.error(f"KDS: Authentication error: {e}")
|
||||
return False
|
||||
|
||||
async def ensure_authenticated(self) -> bool:
|
||||
"""Ensure we have a valid token, re-authenticating if needed."""
|
||||
if not self.access_token:
|
||||
return await self.authenticate()
|
||||
return True
|
||||
|
||||
async def graphql_query(self, query: str, variables: Optional[dict] = None) -> dict:
|
||||
"""Execute a GraphQL query against SambaPOS."""
|
||||
if not await self.ensure_authenticated():
|
||||
return {"error": "Authentication failed"}
|
||||
|
||||
graphql_url = f"{self.server_url}/api/graphql"
|
||||
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": f"Bearer {self.access_token}"
|
||||
}
|
||||
|
||||
payload = {
|
||||
"query": query,
|
||||
"variables": variables,
|
||||
"operationName": None
|
||||
}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=15.0) as client:
|
||||
response = await client.post(
|
||||
graphql_url,
|
||||
json=payload,
|
||||
headers=headers
|
||||
)
|
||||
if response.status_code == 200:
|
||||
return response.json()
|
||||
elif response.status_code == 401:
|
||||
# Token expired, try re-auth
|
||||
self.access_token = None
|
||||
if await self.authenticate():
|
||||
return await self.graphql_query(query, variables)
|
||||
return {"error": "Re-authentication failed"}
|
||||
else:
|
||||
return {"error": f"HTTP {response.status_code}: {response.text}"}
|
||||
except httpx.ConnectError:
|
||||
return {"error": f"Could not connect to {graphql_url}"}
|
||||
except Exception as e:
|
||||
return {"error": str(e)}
|
||||
|
||||
async def get_open_tickets(self) -> dict:
|
||||
"""
|
||||
Query for open (not closed) tickets with full order details.
|
||||
|
||||
Returns tickets with:
|
||||
- Ticket info (id, number, date, table)
|
||||
- Orders with Kitchen Course and Kitchen Print states
|
||||
"""
|
||||
query = """
|
||||
{
|
||||
getTickets(isClosed: false, orderBy: date) {
|
||||
id
|
||||
uid
|
||||
number
|
||||
date
|
||||
lastUpdateTime
|
||||
totalAmount
|
||||
remainingAmount
|
||||
note
|
||||
tags {
|
||||
tag
|
||||
tagName
|
||||
}
|
||||
states {
|
||||
stateName
|
||||
state
|
||||
}
|
||||
orders {
|
||||
id
|
||||
uid
|
||||
name
|
||||
portion
|
||||
quantity
|
||||
price
|
||||
priceTag
|
||||
date
|
||||
tags {
|
||||
tag
|
||||
tagName
|
||||
quantity
|
||||
}
|
||||
states {
|
||||
stateName
|
||||
state
|
||||
stateValue
|
||||
}
|
||||
}
|
||||
entities {
|
||||
type
|
||||
name
|
||||
}
|
||||
}
|
||||
}
|
||||
"""
|
||||
return await self.graphql_query(query)
|
||||
|
||||
async def get_ticket_by_id(self, ticket_id: int) -> dict:
|
||||
"""Query for a specific ticket by ID.
|
||||
|
||||
Uses inline ID rather than GraphQL variables because the SambaPOS
|
||||
Message Server has a NullReferenceException bug when processing
|
||||
variable bindings for getTicket.
|
||||
"""
|
||||
query = f"""
|
||||
{{
|
||||
getTicket(id: {int(ticket_id)}) {{
|
||||
id
|
||||
uid
|
||||
number
|
||||
date
|
||||
lastUpdateTime
|
||||
totalAmount
|
||||
tags {{
|
||||
tag
|
||||
tagName
|
||||
}}
|
||||
states {{
|
||||
stateName
|
||||
state
|
||||
}}
|
||||
orders {{
|
||||
id
|
||||
uid
|
||||
name
|
||||
portion
|
||||
quantity
|
||||
price
|
||||
date
|
||||
tags {{
|
||||
tag
|
||||
tagName
|
||||
quantity
|
||||
}}
|
||||
states {{
|
||||
stateName
|
||||
state
|
||||
stateValue
|
||||
}}
|
||||
}}
|
||||
entities {{
|
||||
type
|
||||
name
|
||||
}}
|
||||
}}
|
||||
}}
|
||||
"""
|
||||
return await self.graphql_query(query)
|
||||
|
||||
|
||||
def parse_kitchen_course(order: dict) -> Optional[str]:
|
||||
"""Extract Kitchen Course from order states."""
|
||||
states = order.get("states", [])
|
||||
for state in states:
|
||||
if state.get("stateName") == "Kitchen Course":
|
||||
return state.get("state")
|
||||
return None
|
||||
|
||||
|
||||
def parse_order_status(order: dict) -> str:
|
||||
"""Extract Status state from order (e.g., Submitted, New)."""
|
||||
states = order.get("states", [])
|
||||
for state in states:
|
||||
if state.get("stateName") == "Status":
|
||||
return state.get("state", "Unknown")
|
||||
return "Unknown"
|
||||
|
||||
|
||||
def parse_kitchen_print_state(order: dict) -> Optional[str]:
|
||||
"""Extract Kitchen Print state from order."""
|
||||
states = order.get("states", [])
|
||||
for state in states:
|
||||
if state.get("stateName") == "Kitchen Print":
|
||||
return state.get("state")
|
||||
return None
|
||||
|
||||
|
||||
def parse_gstatus(order: dict) -> tuple[Optional[str], Optional[str]]:
|
||||
"""Extract GStatus state and timestamp from order (used for void detection)."""
|
||||
states = order.get("states", [])
|
||||
for state in states:
|
||||
if state.get("stateName") == "GStatus":
|
||||
return state.get("state"), state.get("stateDateTime")
|
||||
return None, None
|
||||
|
||||
|
||||
def get_table_name(ticket: dict) -> Optional[str]:
|
||||
"""Extract table name from ticket entities."""
|
||||
entities = ticket.get("entities", [])
|
||||
for entity in entities:
|
||||
if entity.get("type") == "Tables":
|
||||
return entity.get("name")
|
||||
return None
|
||||
|
||||
|
||||
def transform_ticket_for_kds(ticket: dict) -> dict:
|
||||
"""
|
||||
Transform a SambaPOS ticket into KDS-friendly format.
|
||||
|
||||
Groups orders by Kitchen Course and filters for kitchen-relevant items.
|
||||
"""
|
||||
table_name = get_table_name(ticket)
|
||||
|
||||
# Group orders by kitchen course
|
||||
orders_by_course = {}
|
||||
all_orders = []
|
||||
deferred_voided = []
|
||||
earliest_kitchen_order_date = None # Track earliest kitchen-printable order time
|
||||
|
||||
for order in ticket.get("orders", []):
|
||||
kitchen_course = parse_kitchen_course(order) or "Uncategorized"
|
||||
order_status = parse_order_status(order)
|
||||
kitchen_print = parse_kitchen_print_state(order)
|
||||
gstatus, gstatus_datetime = parse_gstatus(order)
|
||||
|
||||
# Log order states for debugging
|
||||
logger.debug(f"Order '{order.get('name')}': status={order_status}, gstatus={gstatus}, kitchen_print={kitchen_print}")
|
||||
|
||||
# Must have Kitchen Print state set (meaning it's a kitchen item)
|
||||
if not kitchen_print:
|
||||
continue
|
||||
|
||||
# Determine if item is voided (show with strikethrough)
|
||||
is_voided = (
|
||||
kitchen_print in ["Canceled", "Void"] or
|
||||
gstatus in ["Void", "Cancelled", "Canceled"] or
|
||||
order_status in ["Void", "Cancelled", "Canceled"]
|
||||
)
|
||||
|
||||
# Get void timestamp if available
|
||||
voided_at = gstatus_datetime if is_voided and gstatus in ["Void", "Cancelled", "Canceled"] else None
|
||||
|
||||
# Skip non-voided items that are not submitted (e.g., "New" status)
|
||||
if not is_voided and order_status not in ["Submitted"]:
|
||||
logger.debug(f"Skipping order '{order.get('name')}' - status is '{order_status}', not 'Submitted'")
|
||||
continue
|
||||
|
||||
# Extract order tags (modifiers like "Rare", "No sauce", etc.)
|
||||
order_tags = order.get("tags", [])
|
||||
|
||||
order_data = {
|
||||
"id": order.get("id"),
|
||||
"uid": order.get("uid"),
|
||||
"name": order.get("name"),
|
||||
"portion": order.get("portion"),
|
||||
"quantity": order.get("quantity"),
|
||||
"price": order.get("price"),
|
||||
"kitchen_course": kitchen_course,
|
||||
"status": order_status,
|
||||
"kitchen_print": kitchen_print,
|
||||
"is_voided": is_voided,
|
||||
"voided_at": voided_at,
|
||||
"tags": order_tags,
|
||||
}
|
||||
|
||||
if is_voided:
|
||||
# Defer voided orders - only add to existing course groups later
|
||||
deferred_voided.append((kitchen_course, order_data))
|
||||
else:
|
||||
if kitchen_course not in orders_by_course:
|
||||
orders_by_course[kitchen_course] = []
|
||||
orders_by_course[kitchen_course].append(order_data)
|
||||
all_orders.append(order_data)
|
||||
|
||||
# Track earliest kitchen-printable order date
|
||||
order_date = order.get("date")
|
||||
if order_date:
|
||||
if earliest_kitchen_order_date is None or order_date < earliest_kitchen_order_date:
|
||||
earliest_kitchen_order_date = order_date
|
||||
|
||||
# Add voided orders only to course groups that already exist (avoids "Uncategorized" ghost courses)
|
||||
for course, order_data in deferred_voided:
|
||||
if course in orders_by_course:
|
||||
orders_by_course[course].append(order_data)
|
||||
all_orders.append(order_data)
|
||||
|
||||
# Skip tickets with no kitchen orders
|
||||
if not all_orders:
|
||||
return None
|
||||
|
||||
# Get covers from tags
|
||||
covers = None
|
||||
for tag in ticket.get("tags", []):
|
||||
if tag.get("tagName") == "Covers":
|
||||
try:
|
||||
covers = int(tag.get("tag"))
|
||||
except (ValueError, TypeError):
|
||||
pass
|
||||
|
||||
return {
|
||||
"id": ticket.get("id"),
|
||||
"uid": ticket.get("uid"),
|
||||
"number": ticket.get("number"),
|
||||
"date": ticket.get("date"),
|
||||
"last_update": ticket.get("lastUpdateTime"),
|
||||
"table": table_name,
|
||||
"covers": covers,
|
||||
"total_amount": ticket.get("totalAmount"),
|
||||
"orders": all_orders,
|
||||
"orders_by_course": orders_by_course,
|
||||
"submitted_at": earliest_kitchen_order_date or ticket.get("date"), # First kitchen order time, fallback to ticket date
|
||||
}
|
||||
296
backend/services/signalr_listener.py
Normal file
296
backend/services/signalr_listener.py
Normal file
|
|
@ -0,0 +1,296 @@
|
|||
"""
|
||||
SambaPOS SignalR Listener for KDS
|
||||
|
||||
Connects to SambaPOS Message Server via SignalR 2.x WebSocket and listens
|
||||
for TICKET_REFRESH broadcasts. When received, fetches the specific ticket
|
||||
by ID (works even for closed/zero-total tickets) and creates/updates
|
||||
KDS entries for real-time display.
|
||||
|
||||
This solves the problem of tickets that close instantly (e.g. free breakfast
|
||||
for residents, bar orders paid immediately) not appearing in KDS polling.
|
||||
|
||||
SignalR 2.x Protocol:
|
||||
1. GET /signalr/negotiate - get connection token
|
||||
2. WS /signalr/connect?transport=webSockets&connectionToken=... - WebSocket
|
||||
3. Messages arrive as JSON: {"C": "...", "M": [{...}]}
|
||||
- TICKET_REFRESH: {"H": "Default", "M": "update", "A": ["guid:<TICKET_REFRESH>ticketId"]}
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import urllib.parse
|
||||
from typing import Optional
|
||||
from datetime import datetime
|
||||
|
||||
import httpx
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class KDSEventBus:
|
||||
"""Simple pub/sub for notifying SSE subscribers of KDS events."""
|
||||
|
||||
def __init__(self):
|
||||
self._subscribers: list[asyncio.Queue] = []
|
||||
|
||||
def subscribe(self) -> asyncio.Queue:
|
||||
q: asyncio.Queue = asyncio.Queue()
|
||||
self._subscribers.append(q)
|
||||
return q
|
||||
|
||||
def unsubscribe(self, q: asyncio.Queue):
|
||||
if q in self._subscribers:
|
||||
self._subscribers.remove(q)
|
||||
|
||||
async def publish(self, event: dict):
|
||||
for q in list(self._subscribers):
|
||||
try:
|
||||
q.put_nowait(event)
|
||||
except asyncio.QueueFull:
|
||||
pass
|
||||
|
||||
|
||||
# Global event bus - imported by kds.py for the SSE endpoint
|
||||
kds_event_bus = KDSEventBus()
|
||||
|
||||
|
||||
class SignalRListener:
|
||||
"""Background listener for SambaPOS SignalR broadcasts."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
base_url: str,
|
||||
graphql_username: str,
|
||||
graphql_password: str,
|
||||
graphql_client_id: str,
|
||||
kitchen_id: int,
|
||||
course_order: list,
|
||||
):
|
||||
self.base_url = base_url.rstrip('/')
|
||||
self.graphql_username = graphql_username
|
||||
self.graphql_password = graphql_password
|
||||
self.graphql_client_id = graphql_client_id
|
||||
self.kitchen_id = kitchen_id
|
||||
self.course_order = course_order
|
||||
self._running = False
|
||||
self._task: Optional[asyncio.Task] = None
|
||||
|
||||
def _get_ws_base(self) -> str:
|
||||
"""Convert HTTP URL to WS URL."""
|
||||
return self.base_url.replace('http://', 'ws://').replace('https://', 'wss://')
|
||||
|
||||
async def _negotiate(self) -> Optional[str]:
|
||||
"""Negotiate SignalR connection and get token."""
|
||||
url = f"{self.base_url}/signalr/negotiate"
|
||||
params = {
|
||||
"clientProtocol": "1.5",
|
||||
"connectionData": json.dumps([{"name": "default"}])
|
||||
}
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
response = await client.get(url, params=params)
|
||||
if response.status_code == 200:
|
||||
data = response.json()
|
||||
return data.get("ConnectionToken")
|
||||
else:
|
||||
logger.error(f"SignalR negotiate failed: {response.status_code}")
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.error(f"SignalR negotiate error: {e}")
|
||||
return None
|
||||
|
||||
async def _listen_loop(self):
|
||||
"""Main WebSocket listen loop with auto-reconnection."""
|
||||
try:
|
||||
import websockets
|
||||
except ImportError:
|
||||
logger.error("SignalR: websockets package not installed")
|
||||
return
|
||||
|
||||
while self._running:
|
||||
try:
|
||||
token = await self._negotiate()
|
||||
if not token:
|
||||
logger.warning("SignalR: Failed to negotiate, retrying in 10s...")
|
||||
await asyncio.sleep(10)
|
||||
continue
|
||||
|
||||
encoded_token = urllib.parse.quote(token, safe='')
|
||||
conn_data = urllib.parse.quote(json.dumps([{"name": "default"}]), safe='')
|
||||
ws_url = (
|
||||
f"{self._get_ws_base()}/signalr/connect"
|
||||
f"?transport=webSockets"
|
||||
f"&clientProtocol=1.5"
|
||||
f"&connectionToken={encoded_token}"
|
||||
f"&connectionData={conn_data}"
|
||||
)
|
||||
|
||||
logger.info("SignalR: Connecting to WebSocket...")
|
||||
async with websockets.connect(ws_url) as ws:
|
||||
logger.info("SignalR: Connected, listening for broadcasts")
|
||||
|
||||
while self._running:
|
||||
try:
|
||||
msg = await asyncio.wait_for(ws.recv(), timeout=30)
|
||||
if msg:
|
||||
data = json.loads(msg)
|
||||
messages = data.get("M", [])
|
||||
for m in messages:
|
||||
await self._handle_message(m)
|
||||
except asyncio.TimeoutError:
|
||||
# Normal - no messages received, keep listening
|
||||
continue
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
except Exception as e:
|
||||
logger.warning(f"SignalR: WebSocket recv error: {e}")
|
||||
break
|
||||
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
except Exception as e:
|
||||
logger.warning(f"SignalR: Connection failed: {e}, reconnecting in 5s...")
|
||||
await asyncio.sleep(5)
|
||||
|
||||
async def _handle_message(self, message: dict):
|
||||
"""Handle a SignalR broadcast message."""
|
||||
args = message.get("A", [])
|
||||
for arg in args:
|
||||
if "<TICKET_REFRESH>" in arg:
|
||||
try:
|
||||
ticket_id_str = arg.split("<TICKET_REFRESH>")[1]
|
||||
ticket_id = int(ticket_id_str)
|
||||
logger.info(f"SignalR: TICKET_REFRESH for ticket {ticket_id}")
|
||||
await self._process_ticket_refresh(ticket_id)
|
||||
except (ValueError, IndexError) as e:
|
||||
logger.warning(f"SignalR: Failed to parse TICKET_REFRESH: {arg} - {e}")
|
||||
|
||||
async def _process_ticket_refresh(self, sambapos_ticket_id: int):
|
||||
"""Handle a TICKET_REFRESH broadcast.
|
||||
|
||||
Fetches the specific ticket by ID from SambaPOS GraphQL and
|
||||
persists it as a KDS entry. This captures instantly-closed tickets
|
||||
(free breakfast, bar tabs) that never appear in getTickets(isClosed: false).
|
||||
|
||||
Always publishes an SSE event afterwards so the frontend refreshes.
|
||||
"""
|
||||
# Fetch and persist the ticket directly
|
||||
try:
|
||||
await self._fetch_and_persist_ticket(sambapos_ticket_id)
|
||||
except Exception as e:
|
||||
logger.warning(f"SignalR: Failed to fetch/persist ticket {sambapos_ticket_id}: {e}")
|
||||
|
||||
# Always notify SSE subscribers for instant frontend refresh
|
||||
await kds_event_bus.publish({
|
||||
"type": "ticket_refresh",
|
||||
"sambapos_ticket_id": sambapos_ticket_id,
|
||||
"timestamp": datetime.utcnow().isoformat(),
|
||||
})
|
||||
logger.info(f"SignalR: Published SSE event for ticket {sambapos_ticket_id}")
|
||||
|
||||
async def _fetch_and_persist_ticket(self, sambapos_ticket_id: int):
|
||||
"""Fetch a specific ticket from SambaPOS and create/update a KDS entry."""
|
||||
from services.kds_graphql import SambaPOSGraphQLClient, transform_ticket_for_kds
|
||||
from api.kds import get_or_create_kds_ticket
|
||||
from database import AsyncSessionLocal
|
||||
|
||||
client = SambaPOSGraphQLClient(
|
||||
server_url=self.base_url,
|
||||
username=self.graphql_username,
|
||||
password=self.graphql_password,
|
||||
client_id=self.graphql_client_id,
|
||||
)
|
||||
|
||||
result = await client.get_ticket_by_id(sambapos_ticket_id)
|
||||
|
||||
if "error" in result:
|
||||
logger.warning(f"SignalR: get_ticket_by_id({sambapos_ticket_id}) error: {result['error']}")
|
||||
return
|
||||
|
||||
ticket_data = result.get("data", {}).get("getTicket")
|
||||
if not ticket_data:
|
||||
logger.debug(f"SignalR: get_ticket_by_id({sambapos_ticket_id}) returned no data")
|
||||
return
|
||||
|
||||
transformed = transform_ticket_for_kds(ticket_data)
|
||||
if not transformed:
|
||||
logger.debug(f"SignalR: Ticket {sambapos_ticket_id} has no kitchen orders, skipping")
|
||||
return
|
||||
|
||||
async with AsyncSessionLocal() as db:
|
||||
kds_ticket = await get_or_create_kds_ticket(
|
||||
db, self.kitchen_id, transformed, self.course_order
|
||||
)
|
||||
logger.info(
|
||||
f"SignalR: Persisted ticket {sambapos_ticket_id} "
|
||||
f"(KDS #{kds_ticket.id}, number={kds_ticket.ticket_number})"
|
||||
)
|
||||
|
||||
def start(self):
|
||||
"""Start the listener as a background asyncio task."""
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._listen_loop())
|
||||
logger.info("SignalR: Listener started")
|
||||
|
||||
def stop(self):
|
||||
"""Stop the listener."""
|
||||
self._running = False
|
||||
if self._task:
|
||||
self._task.cancel()
|
||||
logger.info("SignalR: Listener stopped")
|
||||
|
||||
|
||||
# Global listener instance
|
||||
_listener: Optional[SignalRListener] = None
|
||||
|
||||
|
||||
async def start_signalr_listener():
|
||||
"""Start the global SignalR listener using KDS settings from DB."""
|
||||
global _listener
|
||||
|
||||
from database import AsyncSessionLocal
|
||||
from sqlalchemy import select
|
||||
from models.settings import KitchenSettings
|
||||
|
||||
try:
|
||||
async with AsyncSessionLocal() as db:
|
||||
# Get first kitchen with KDS GraphQL configured
|
||||
result = await db.execute(
|
||||
select(KitchenSettings).where(
|
||||
KitchenSettings.kds_graphql_url.isnot(None)
|
||||
)
|
||||
)
|
||||
settings = result.scalar_one_or_none()
|
||||
|
||||
if not settings:
|
||||
logger.info("SignalR: No KDS GraphQL URL configured, listener not started")
|
||||
return
|
||||
|
||||
if not all([settings.kds_graphql_url, settings.kds_graphql_username,
|
||||
settings.kds_graphql_password, settings.kds_graphql_client_id]):
|
||||
logger.info("SignalR: KDS GraphQL credentials incomplete, listener not started")
|
||||
return
|
||||
|
||||
course_order = settings.kds_course_order or ["Starters", "Mains", "Desserts"]
|
||||
|
||||
_listener = SignalRListener(
|
||||
base_url=settings.kds_graphql_url,
|
||||
graphql_username=settings.kds_graphql_username,
|
||||
graphql_password=settings.kds_graphql_password,
|
||||
graphql_client_id=settings.kds_graphql_client_id,
|
||||
kitchen_id=settings.kitchen_id,
|
||||
course_order=course_order,
|
||||
)
|
||||
_listener.start()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"SignalR: Failed to start listener: {e}")
|
||||
|
||||
|
||||
async def stop_signalr_listener():
|
||||
"""Stop the global SignalR listener."""
|
||||
global _listener
|
||||
if _listener:
|
||||
_listener.stop()
|
||||
_listener = None
|
||||
Loading…
Add table
Add a link
Reference in a new issue