Spaces:
Paused
Paused
| """ | |
| Human-in-the-Loop Dashboard for trade compliance review. | |
| Provides workflow for logistics managers to: | |
| 1. Review PDF-parsed documents | |
| 2. Approve/reject HS code classifications | |
| 3. Review sanctions screening results | |
| 4. Add comments and annotations | |
| 5. Export approved documents | |
| This is the KILLER FEATURE: | |
| - Exporter uploads PDF -> Docling extracts -> pyhscodes suggests -> LLM refines | |
| - Logistics manager reviews and APPROVES | |
| - Liability stays with human approver, not software | |
| Audit Trail: | |
| Every state transition is recorded as an immutable audit entry with: | |
| - SHA-256 hash chain (each entry hashes the previous) | |
| - Strict UTC timestamps | |
| - Actor identification (email or system ID) | |
| - Previous state + modified values (diff-based) | |
| - Action metadata (approve, reject, request_info, escalate) | |
| Typing conventions: | |
| All public APIs use explicit type hints. Final constants prevent mutation. | |
| """ | |
| from __future__ import annotations | |
| import hashlib | |
| import json | |
| import logging | |
| import threading | |
| import time as _time | |
| from dataclasses import dataclass, field | |
| from datetime import datetime, timezone | |
| from enum import Enum | |
| from pathlib import Path | |
| from typing import Annotated, Any, Final, Literal, Optional | |
| logger: Final = logging.getLogger(__name__) | |
| # ── Domain Exceptions ───────────────────────────────────────────────── | |
| class DashboardError(Exception): | |
| """Base exception for all dashboard errors.""" | |
| class ReviewItemNotFoundError(DashboardError): | |
| """Raised when a review item is not found by ID. | |
| Attributes: | |
| item_id: The ID that was not found. | |
| """ | |
| def __init__(self, item_id: str) -> None: | |
| self.item_id = item_id | |
| super().__init__(f"Review item not found: {item_id}") | |
| class ReviewItemAlreadyResolvedError(DashboardError): | |
| """Raised when attempting to modify an already-resolved review item. | |
| Attributes: | |
| item_id: The resolved item's ID. | |
| current_status: The current status blocking the operation. | |
| """ | |
| def __init__(self, item_id: str, current_status: str) -> None: | |
| self.item_id = item_id | |
| self.current_status = current_status | |
| super().__init__( | |
| f"Item {item_id} is already {current_status} and cannot be modified" | |
| ) | |
| class AuditLogIntegrityError(DashboardError): | |
| """Raised when audit log hash chain is broken. | |
| Attributes: | |
| entry_index: Index of the broken entry. | |
| expected_hash: Expected hash of the previous entry. | |
| actual_hash: Actual hash found. | |
| """ | |
| def __init__( | |
| self, entry_index: int, expected_hash: str, actual_hash: str | |
| ) -> None: | |
| self.entry_index = entry_index | |
| self.expected_hash = expected_hash | |
| self.actual_hash = actual_hash | |
| super().__init__( | |
| f"Audit log integrity broken at entry {entry_index}: " | |
| f"expected {expected_hash[:12]}..., got {actual_hash[:12]}..." | |
| ) | |
| class QueueCapacityExceededError(DashboardError): | |
| """Raised when review queue exceeds maximum capacity. | |
| Attributes: | |
| current_count: Current queue size. | |
| max_capacity: Maximum allowed capacity. | |
| """ | |
| def __init__(self, current_count: int, max_capacity: int) -> None: | |
| self.current_count = current_count | |
| self.max_capacity = max_capacity | |
| super().__init__( | |
| f"Queue capacity exceeded: {current_count}/{max_capacity}" | |
| ) | |
| # ── Constants ───────────────────────────────────────────────────────── | |
| # Type aliases | |
| ReviewActionType = Literal["approve", "reject", "request_info", "escalate"] | |
| QueueHealth = Literal["healthy", "attention", "warning", "critical"] | |
| # Maximum queue capacity | |
| MAX_QUEUE_CAPACITY: Final[int] = 100_000 | |
| # UTC timestamp sentinel for "no deadline" | |
| NO_DEADLINE: Final[float] = 0.0 | |
| # ── Audit Trail ─────────────────────────────────────────────────────── | |
| class AuditEntry: | |
| """Immutable audit log entry with hash chain. | |
| Attributes: | |
| entry_index: Sequential position in the audit log. | |
| timestamp_utc: ISO 8601 UTC timestamp string. | |
| timestamp_unix: Unix timestamp (seconds since epoch). | |
| item_id: ID of the review item affected. | |
| actor: Actor identifier (email or system ID). | |
| action: Action taken (approve, reject, request_info, escalate). | |
| previous_status: Status before the action. | |
| new_status: Status after the action. | |
| modified_values: Dict of field changes (diff-based). | |
| reason: Reason or notes for the action. | |
| previous_hash: SHA-256 hash of the previous entry (chain link). | |
| entry_hash: SHA-256 hash of this entry's contents. | |
| """ | |
| entry_index: int | |
| timestamp_utc: str | |
| timestamp_unix: float | |
| item_id: str | |
| actor: str | |
| action: str | |
| previous_status: str | |
| new_status: str | |
| modified_values: dict[str, Any] | |
| reason: str | |
| previous_hash: str | |
| entry_hash: str | |
| def compute_hash( | |
| entry_index: int, | |
| timestamp_utc: str, | |
| item_id: str, | |
| actor: str, | |
| action: str, | |
| previous_status: str, | |
| new_status: str, | |
| modified_values: dict[str, Any], | |
| reason: str, | |
| previous_hash: str, | |
| ) -> str: | |
| """Compute SHA-256 hash for an audit entry. | |
| Args: | |
| entry_index: Sequential position. | |
| timestamp_utc: ISO 8601 UTC timestamp. | |
| item_id: Review item ID. | |
| actor: Actor identifier. | |
| action: Action taken. | |
| previous_status: Status before action. | |
| new_status: Status after action. | |
| modified_values: Field changes. | |
| reason: Reason for action. | |
| previous_hash: Hash of previous entry. | |
| Returns: | |
| SHA-256 hex digest string. | |
| """ | |
| payload: str = json.dumps( | |
| { | |
| "i": entry_index, | |
| "t": timestamp_utc, | |
| "item": item_id, | |
| "a": actor, | |
| "act": action, | |
| "prev": previous_status, | |
| "new": new_status, | |
| "diff": modified_values, | |
| "r": reason, | |
| "ph": previous_hash, | |
| }, | |
| sort_keys=True, | |
| separators=(",", ":"), | |
| ) | |
| return hashlib.sha256(payload.encode("utf-8")).hexdigest() | |
| def to_dict(self) -> dict[str, Any]: | |
| """Convert to dictionary for serialization. | |
| Returns: | |
| Dict with all audit entry fields. | |
| """ | |
| return { | |
| "entry_index": self.entry_index, | |
| "timestamp_utc": self.timestamp_utc, | |
| "timestamp_unix": self.timestamp_unix, | |
| "item_id": self.item_id, | |
| "actor": self.actor, | |
| "action": self.action, | |
| "previous_status": self.previous_status, | |
| "new_status": self.new_status, | |
| "modified_values": self.modified_values, | |
| "reason": self.reason, | |
| "previous_hash": self.previous_hash, | |
| "entry_hash": self.entry_hash, | |
| } | |
| class AuditLog: | |
| """Immutable, hash-chained audit log. | |
| Every state transition is recorded as an AuditEntry with SHA-256 | |
| hash chain integrity. The log is append-only; entries cannot be | |
| modified after creation. | |
| Raises: | |
| AuditLogIntegrityError: If hash chain verification fails. | |
| """ | |
| GENESIS_HASH: Final[str] = "0" * 64 | |
| def __init__(self) -> None: | |
| """Initialize an empty audit log.""" | |
| self._entries: list[AuditEntry] = [] | |
| self._lock: threading.Lock = threading.Lock() | |
| def append( | |
| self, | |
| item_id: str, | |
| actor: str, | |
| action: str, | |
| previous_status: str, | |
| new_status: str, | |
| modified_values: dict[str, Any], | |
| reason: str = "", | |
| ) -> AuditEntry: | |
| """Append an audit entry to the log. | |
| Args: | |
| item_id: ID of the review item affected. | |
| actor: Actor identifier (email or system ID). | |
| action: Action taken. | |
| previous_status: Status before the action. | |
| new_status: Status after the action. | |
| modified_values: Dict of field changes (diff-based). | |
| reason: Reason or notes for the action. | |
| Returns: | |
| The created AuditEntry. | |
| """ | |
| with self._lock: | |
| now_unix: float = _time.time() | |
| now_utc: str = datetime.fromtimestamp( | |
| now_unix, tz=timezone.utc | |
| ).isoformat() | |
| entry_index: int = len(self._entries) | |
| previous_hash: str = ( | |
| self._entries[-1].entry_hash | |
| if self._entries | |
| else self.GENESIS_HASH | |
| ) | |
| entry_hash: str = AuditEntry.compute_hash( | |
| entry_index=entry_index, | |
| timestamp_utc=now_utc, | |
| item_id=item_id, | |
| actor=actor, | |
| action=action, | |
| previous_status=previous_status, | |
| new_status=new_status, | |
| modified_values=modified_values, | |
| reason=reason, | |
| previous_hash=previous_hash, | |
| ) | |
| entry = AuditEntry( | |
| entry_index=entry_index, | |
| timestamp_utc=now_utc, | |
| timestamp_unix=now_unix, | |
| item_id=item_id, | |
| actor=actor, | |
| action=action, | |
| previous_status=previous_status, | |
| new_status=new_status, | |
| modified_values=modified_values, | |
| reason=reason, | |
| previous_hash=previous_hash, | |
| entry_hash=entry_hash, | |
| ) | |
| self._entries.append(entry) | |
| return entry | |
| def verify_chain(self) -> bool: | |
| """Verify the integrity of the entire hash chain. | |
| Returns: | |
| True if the chain is valid, False otherwise. | |
| Raises: | |
| AuditLogIntegrityError: If a link in the chain is broken. | |
| """ | |
| with self._lock: | |
| for i, entry in enumerate(self._entries): | |
| expected_prev: str = ( | |
| self._entries[i - 1].entry_hash | |
| if i > 0 | |
| else self.GENESIS_HASH | |
| ) | |
| if entry.previous_hash != expected_prev: | |
| raise AuditLogIntegrityError( | |
| entry_index=i, | |
| expected_hash=expected_prev, | |
| actual_hash=entry.previous_hash, | |
| ) | |
| return True | |
| def get_entries_for_item(self, item_id: str) -> list[AuditEntry]: | |
| """Get all audit entries for a specific item. | |
| Args: | |
| item_id: Review item ID to filter by. | |
| Returns: | |
| List of AuditEntry objects for that item. | |
| """ | |
| with self._lock: | |
| return [e for e in self._entries if e.item_id == item_id] | |
| def get_all_entries(self) -> list[AuditEntry]: | |
| """Get all audit entries. | |
| Returns: | |
| Copy of all audit entries. | |
| """ | |
| with self._lock: | |
| return list(self._entries) | |
| def entry_count(self) -> int: | |
| """Number of entries in the log.""" | |
| with self._lock: | |
| return len(self._entries) | |
| # ── Constants ───────────────────────────────────────────────────────── | |
| # Sentinel value for "no deadline" | |
| # ── Enums ───────────────────────────────────────────────────────────── | |
| class ReviewStatus(Enum): | |
| """Status of a review item. | |
| Lifecycle: PENDING -> IN_REVIEW -> APPROVED | REJECTED | NEEDS_INFO | ESCALATED | |
| """ | |
| PENDING = "pending" | |
| IN_REVIEW = "in_review" | |
| APPROVED = "approved" | |
| REJECTED = "rejected" | |
| NEEDS_INFO = "needs_info" | |
| ESCALATED = "escalated" | |
| class RiskLevel(Enum): | |
| """Risk levels for compliance items.""" | |
| CLEAR = "clear" | |
| LOW = "low" | |
| MEDIUM = "medium" | |
| HIGH = "high" | |
| BLOCKED = "blocked" | |
| # ── Data Models ─────────────────────────────────────────────────────── | |
| class ReviewComment: | |
| """A comment on a review item. | |
| Attributes: | |
| author: Comment author identifier. | |
| content: Comment text. | |
| timestamp: Unix timestamp when the comment was created. | |
| is_internal: If True, comment is internal-only (not visible to exporter). | |
| """ | |
| author: str | |
| content: str | |
| timestamp: float = field(default_factory=_time.time) | |
| is_internal: bool = False | |
| def to_dict(self) -> dict[str, Any]: | |
| """Convert to dictionary for serialization. | |
| Returns: | |
| Dict with author, content, timestamp, and is_internal. | |
| """ | |
| return { | |
| "author": self.author, | |
| "content": self.content, | |
| "timestamp": self.timestamp, | |
| "is_internal": self.is_internal, | |
| } | |
| class ReviewAction: | |
| """An action taken on a review item. | |
| Attributes: | |
| action: Action type (approve, reject, request_info, escalate). | |
| actor: Actor identifier (email or system ID). | |
| timestamp: Unix timestamp when the action was taken. | |
| reason: Reason or notes for the action. | |
| metadata: Additional metadata as key-value pairs. | |
| """ | |
| action: str | |
| actor: str | |
| timestamp: float = field(default_factory=_time.time) | |
| reason: str = "" | |
| metadata: dict[str, Any] = field(default_factory=dict) | |
| def to_dict(self) -> dict[str, Any]: | |
| """Convert to dictionary for serialization. | |
| Returns: | |
| Dict with action, actor, timestamp, reason, and metadata. | |
| """ | |
| return { | |
| "action": self.action, | |
| "actor": self.actor, | |
| "timestamp": self.timestamp, | |
| "reason": self.reason, | |
| "metadata": self.metadata, | |
| } | |
| class ReviewItem: | |
| """An item pending human review. | |
| Attributes: | |
| item_id: Unique identifier (e.g., REV-000001). | |
| document_path: Path to the source document. | |
| document_type: Document type (invoice, packing_list, etc.). | |
| invoice_number: Extracted invoice number. | |
| invoice_date: Extracted invoice date. | |
| total_amount: Extracted total amount. | |
| shipper: Shipper name/party. | |
| consignee: Consignee name/party. | |
| country_origin: Country of origin. | |
| country_destination: Destination country. | |
| hs_code_suggested: HS code suggested by classifier. | |
| hs_code_description: Description of the suggested HS code. | |
| hs_code_confidence: Classification confidence (0.0-1.0). | |
| hs_code_alternatives: List of alternative HS codes with confidence. | |
| sanctions_risk_level: Sanctions screening risk level. | |
| sanctions_matches: List of sanctions screening matches. | |
| status: Current review status. | |
| assigned_to: Assigned reviewer identifier. | |
| priority: Priority level (0=normal, 1=high, 2=urgent). | |
| comments: List of review comments. | |
| actions: List of actions taken on this item. | |
| final_hs_code: Final approved HS code. | |
| reviewer_notes: Reviewer's notes. | |
| created_at: Unix timestamp when the item was created. | |
| updated_at: Unix timestamp when the item was last updated. | |
| review_deadline: Unix timestamp for review deadline (0.0 = no deadline). | |
| """ | |
| item_id: str | |
| document_path: str | |
| document_type: str | |
| # Extracted data | |
| invoice_number: str = "" | |
| invoice_date: str = "" | |
| total_amount: str = "" | |
| shipper: str = "" | |
| consignee: str = "" | |
| country_origin: str = "" | |
| country_destination: str = "" | |
| # Classification data | |
| hs_code_suggested: str = "" | |
| hs_code_description: str = "" | |
| hs_code_confidence: float = 0.0 | |
| hs_code_alternatives: list[dict[str, Any]] = field(default_factory=list) | |
| # Sanctions screening | |
| sanctions_risk_level: str = "clear" | |
| sanctions_matches: list[dict[str, Any]] = field(default_factory=list) | |
| # Review state | |
| status: ReviewStatus = ReviewStatus.PENDING | |
| assigned_to: str = "" | |
| priority: int = 0 | |
| # Review data | |
| comments: list[ReviewComment] = field(default_factory=list) | |
| actions: list[ReviewAction] = field(default_factory=list) | |
| final_hs_code: str = "" | |
| reviewer_notes: str = "" | |
| # Metadata | |
| created_at: float = field(default_factory=_time.time) | |
| updated_at: float = field(default_factory=_time.time) | |
| review_deadline: float = NO_DEADLINE | |
| def is_pending(self) -> bool: | |
| """Whether the item is in pending status.""" | |
| return self.status == ReviewStatus.PENDING | |
| def is_approved(self) -> bool: | |
| """Whether the item has been approved.""" | |
| return self.status == ReviewStatus.APPROVED | |
| def is_rejected(self) -> bool: | |
| """Whether the item has been rejected.""" | |
| return self.status == ReviewStatus.REJECTED | |
| def needs_action(self) -> bool: | |
| """Whether the item needs human action.""" | |
| return self.status in ( | |
| ReviewStatus.PENDING, | |
| ReviewStatus.IN_REVIEW, | |
| ReviewStatus.NEEDS_INFO, | |
| ) | |
| def approve( | |
| self, | |
| actor: str, | |
| hs_code: str = "", | |
| notes: str = "", | |
| ) -> bool: | |
| """Approve this review item. | |
| Args: | |
| actor: Actor identifier (email or system ID). | |
| hs_code: Final HS code to assign. | |
| notes: Reviewer notes. | |
| Returns: | |
| True if approved, False if item was already resolved. | |
| Raises: | |
| ReviewItemAlreadyResolvedError: If item is not in an approvable state. | |
| """ | |
| if self.status not in ( | |
| ReviewStatus.PENDING, | |
| ReviewStatus.IN_REVIEW, | |
| ReviewStatus.NEEDS_INFO, | |
| ): | |
| return False | |
| previous_status: str = self.status.value | |
| self.status = ReviewStatus.APPROVED | |
| self.final_hs_code = hs_code or self.hs_code_suggested | |
| self.reviewer_notes = notes | |
| self.updated_at = _time.time() | |
| self.actions.append(ReviewAction( | |
| action="approve", | |
| actor=actor, | |
| reason=notes, | |
| )) | |
| return True | |
| def reject(self, actor: str, reason: str = "") -> bool: | |
| """Reject this review item. | |
| Args: | |
| actor: Actor identifier (email or system ID). | |
| reason: Rejection reason. | |
| Returns: | |
| True if rejected, False if item was already resolved. | |
| """ | |
| if self.status not in ( | |
| ReviewStatus.PENDING, | |
| ReviewStatus.IN_REVIEW, | |
| ReviewStatus.NEEDS_INFO, | |
| ): | |
| return False | |
| self.status = ReviewStatus.REJECTED | |
| self.reviewer_notes = reason | |
| self.updated_at = _time.time() | |
| self.actions.append(ReviewAction( | |
| action="reject", | |
| actor=actor, | |
| reason=reason, | |
| )) | |
| return True | |
| def request_info(self, actor: str, reason: str = "") -> bool: | |
| """Request more information. | |
| Args: | |
| actor: Actor identifier (email or system ID). | |
| reason: Reason for the information request. | |
| Returns: | |
| True if updated, False if item was already resolved. | |
| """ | |
| if self.status not in ( | |
| ReviewStatus.PENDING, | |
| ReviewStatus.IN_REVIEW, | |
| ): | |
| return False | |
| self.status = ReviewStatus.NEEDS_INFO | |
| self.updated_at = _time.time() | |
| self.actions.append(ReviewAction( | |
| action="request_info", | |
| actor=actor, | |
| reason=reason, | |
| )) | |
| return True | |
| def add_comment( | |
| self, | |
| author: str, | |
| content: str, | |
| is_internal: bool = False, | |
| ) -> None: | |
| """Add a comment to this item. | |
| Args: | |
| author: Comment author identifier. | |
| content: Comment text. | |
| is_internal: If True, comment is internal-only. | |
| """ | |
| self.comments.append(ReviewComment( | |
| author=author, | |
| content=content, | |
| is_internal=is_internal, | |
| )) | |
| self.updated_at = _time.time() | |
| def to_dict(self) -> dict[str, Any]: | |
| """Convert to dictionary for serialization. | |
| Returns: | |
| Dict with all item fields suitable for JSON serialization. | |
| """ | |
| return { | |
| "item_id": self.item_id, | |
| "document_path": self.document_path, | |
| "document_type": self.document_type, | |
| "invoice_number": self.invoice_number, | |
| "invoice_date": self.invoice_date, | |
| "total_amount": self.total_amount, | |
| "shipper": self.shipper, | |
| "consignee": self.consignee, | |
| "country_origin": self.country_origin, | |
| "country_destination": self.country_destination, | |
| "hs_code_suggested": self.hs_code_suggested, | |
| "hs_code_description": self.hs_code_description, | |
| "hs_code_confidence": self.hs_code_confidence, | |
| "hs_code_alternatives": self.hs_code_alternatives, | |
| "sanctions_risk_level": self.sanctions_risk_level, | |
| "sanctions_matches": self.sanctions_matches, | |
| "status": self.status.value, | |
| "assigned_to": self.assigned_to, | |
| "priority": self.priority, | |
| "comments": [c.to_dict() for c in self.comments], | |
| "actions": [a.to_dict() for a in self.actions], | |
| "final_hs_code": self.final_hs_code, | |
| "reviewer_notes": self.reviewer_notes, | |
| "created_at": self.created_at, | |
| "updated_at": self.updated_at, | |
| "review_deadline": self.review_deadline, | |
| } | |
| # ── Review Queue ────────────────────────────────────────────────────── | |
| class ReviewQueue: | |
| """Manages the queue of items pending review. | |
| Thread-safe: all mutating operations are protected by a lock. | |
| Supports optimistic concurrency on approve/reject/request_info | |
| (returns bool indicating success). | |
| Attributes: | |
| _items: Internal dict of review items keyed by item_id. | |
| _counter: Auto-increment counter for item ID generation. | |
| _lock: Threading lock for thread safety. | |
| _audit_log: Cryptographic audit log for all state transitions. | |
| """ | |
| def __init__(self) -> None: | |
| """Initialize an empty review queue.""" | |
| self._items: dict[str, ReviewItem] = {} | |
| self._counter: int = 0 | |
| self._lock: threading.Lock = threading.Lock() | |
| self._audit_log: AuditLog = AuditLog() | |
| def audit_log(self) -> AuditLog: | |
| """Access the audit log.""" | |
| return self._audit_log | |
| def add_item(self, item: ReviewItem) -> str: | |
| """Add an item to the review queue. | |
| Args: | |
| item: ReviewItem to add. If item_id is empty, one is generated. | |
| Returns: | |
| The item's ID (either provided or generated). | |
| """ | |
| with self._lock: | |
| if not item.item_id: | |
| self._counter += 1 | |
| item.item_id = f"REV-{self._counter:06d}" | |
| self._items[item.item_id] = item | |
| logger.info(f"Added item {item.item_id} to review queue") | |
| return item.item_id | |
| def get_item(self, item_id: str) -> Optional[ReviewItem]: | |
| """Get an item by ID. | |
| Args: | |
| item_id: The review item ID. | |
| Returns: | |
| ReviewItem if found, None otherwise. | |
| """ | |
| with self._lock: | |
| return self._items.get(item_id) | |
| def remove_item(self, item_id: str) -> bool: | |
| """Remove an item from the queue. | |
| Args: | |
| item_id: The review item ID. | |
| Returns: | |
| True if item was removed, False if not found. | |
| """ | |
| with self._lock: | |
| if item_id in self._items: | |
| del self._items[item_id] | |
| return True | |
| return False | |
| def get_pending(self) -> list[ReviewItem]: | |
| """Get all pending items. | |
| Returns: | |
| List of ReviewItem objects with PENDING status. | |
| """ | |
| with self._lock: | |
| return [ | |
| item for item in self._items.values() | |
| if item.is_pending | |
| ] | |
| def get_by_status(self, status: ReviewStatus) -> list[ReviewItem]: | |
| """Get items by status. | |
| Args: | |
| status: The ReviewStatus to filter by. | |
| Returns: | |
| List of ReviewItem objects with the given status. | |
| """ | |
| with self._lock: | |
| return [ | |
| item for item in self._items.values() | |
| if item.status == status | |
| ] | |
| def get_by_assignee(self, assignee: str) -> list[ReviewItem]: | |
| """Get items assigned to a specific person. | |
| Args: | |
| assignee: The assignee identifier. | |
| Returns: | |
| List of ReviewItem objects assigned to the given person. | |
| """ | |
| with self._lock: | |
| return [ | |
| item for item in self._items.values() | |
| if item.assigned_to == assignee | |
| ] | |
| def get_overdue(self) -> list[ReviewItem]: | |
| """Get items past their review deadline. | |
| Returns: | |
| List of ReviewItem objects with a past deadline. | |
| """ | |
| current_time: float = _time.time() | |
| with self._lock: | |
| return [ | |
| item for item in self._items.values() | |
| if item.review_deadline > NO_DEADLINE | |
| and item.review_deadline < current_time | |
| ] | |
| def total_count(self) -> int: | |
| """Total number of items in the queue.""" | |
| with self._lock: | |
| return len(self._items) | |
| def pending_count(self) -> int: | |
| """Number of pending items.""" | |
| with self._lock: | |
| return sum(1 for i in self._items.values() if i.is_pending) | |
| def approved_count(self) -> int: | |
| """Number of approved items.""" | |
| with self._lock: | |
| return sum( | |
| 1 for i in self._items.values() | |
| if i.status == ReviewStatus.APPROVED | |
| ) | |
| def rejected_count(self) -> int: | |
| """Number of rejected items.""" | |
| with self._lock: | |
| return sum( | |
| 1 for i in self._items.values() | |
| if i.status == ReviewStatus.REJECTED | |
| ) | |
| def get_statistics(self) -> dict[str, Any]: | |
| """Get queue statistics. | |
| Returns: | |
| Dict with counts for each status plus total and overdue. | |
| """ | |
| with self._lock: | |
| statuses = [item.status for item in self._items.values()] | |
| return { | |
| "total": len(self._items), | |
| "pending": statuses.count(ReviewStatus.PENDING), | |
| "approved": statuses.count(ReviewStatus.APPROVED), | |
| "rejected": statuses.count(ReviewStatus.REJECTED), | |
| "in_review": statuses.count(ReviewStatus.IN_REVIEW), | |
| "needs_info": statuses.count(ReviewStatus.NEEDS_INFO), | |
| "overdue": len(self.get_overdue()), | |
| } | |
| def to_list(self) -> list[dict[str, Any]]: | |
| """Convert all items to list of dictionaries. | |
| Returns: | |
| List of dicts, one per review item. | |
| """ | |
| with self._lock: | |
| return [item.to_dict() for item in self._items.values()] | |
| # ── Dashboard Data Provider ─────────────────────────────────────────── | |
| class DashboardDataProvider: | |
| """Provides data for the compliance dashboard. | |
| Attributes: | |
| queue: The ReviewQueue to query. | |
| """ | |
| def __init__(self, queue: Optional[ReviewQueue] = None) -> None: | |
| """Initialize the dashboard data provider. | |
| Args: | |
| queue: ReviewQueue to use. Creates a new one if None. | |
| """ | |
| self.queue: ReviewQueue = queue or ReviewQueue() | |
| def get_dashboard_summary(self) -> dict[str, Any]: | |
| """Get summary data for dashboard. | |
| Returns: | |
| Dict with statistics, average review time, and queue health. | |
| """ | |
| stats: dict[str, int] = self.queue.get_statistics() | |
| approved_items: list[ReviewItem] = self.queue.get_by_status( | |
| ReviewStatus.APPROVED | |
| ) | |
| avg_review_time: float = 0.0 | |
| if approved_items: | |
| review_times: list[float] = [ | |
| item.updated_at - item.created_at | |
| for item in approved_items | |
| ] | |
| avg_review_time = sum(review_times) / len(review_times) | |
| return { | |
| "statistics": stats, | |
| "average_review_time_seconds": avg_review_time, | |
| "queue_health": self._calculate_queue_health(stats), | |
| } | |
| def get_review_items( | |
| self, | |
| status: Optional[ReviewStatus] = None, | |
| assignee: Optional[str] = None, | |
| limit: int = 50, | |
| ) -> list[dict[str, Any]]: | |
| """Get review items with optional filters. | |
| Args: | |
| status: Filter by ReviewStatus (None = all). | |
| assignee: Filter by assignee (None = all). | |
| limit: Maximum number of items to return. | |
| Returns: | |
| List of dicts, sorted by priority (highest first) then creation time. | |
| """ | |
| items: list[ReviewItem] = list(self.queue._items.values()) | |
| if status: | |
| items = [i for i in items if i.status == status] | |
| if assignee: | |
| items = [i for i in items if i.assigned_to == assignee] | |
| items.sort(key=lambda x: (-x.priority, x.created_at)) | |
| return [item.to_dict() for item in items[:limit]] | |
| def get_item_details(self, item_id: str) -> Optional[dict[str, Any]]: | |
| """Get detailed information for a specific item. | |
| Args: | |
| item_id: The review item ID. | |
| Returns: | |
| Dict with item details, or None if not found. | |
| """ | |
| item: Optional[ReviewItem] = self.queue.get_item(item_id) | |
| if item: | |
| return item.to_dict() | |
| return None | |
| def _calculate_queue_health(self, stats: dict[str, int]) -> QueueHealth: | |
| """Calculate queue health based on statistics. | |
| Args: | |
| stats: Queue statistics dict. | |
| Returns: | |
| QueueHealth literal: "healthy", "attention", "warning", or "critical". | |
| """ | |
| pending: int = stats.get("pending", 0) | |
| overdue: int = stats.get("overdue", 0) | |
| if overdue > 0: | |
| return "critical" | |
| elif pending > 100: | |
| return "warning" | |
| elif pending > 50: | |
| return "attention" | |
| else: | |
| return "healthy" | |
| # ── Convenience Functions ───────────────────────────────────────────── | |
| _queue: Optional[ReviewQueue] = None | |
| _provider: Optional[DashboardDataProvider] = None | |
| _queue_lock: threading.Lock = threading.Lock() | |
| _provider_lock: threading.Lock = threading.Lock() | |
| def get_queue() -> ReviewQueue: | |
| """Get or create the global review queue (thread-safe). | |
| Returns: | |
| The singleton ReviewQueue instance. | |
| """ | |
| global _queue | |
| if _queue is None: | |
| with _queue_lock: | |
| if _queue is None: | |
| _queue = ReviewQueue() | |
| return _queue | |
| def get_provider() -> DashboardDataProvider: | |
| """Get or create the global dashboard provider (thread-safe). | |
| Returns: | |
| The singleton DashboardDataProvider instance. | |
| """ | |
| global _provider | |
| if _provider is None: | |
| with _provider_lock: | |
| if _provider is None: | |
| _provider = DashboardDataProvider(get_queue()) | |
| return _provider | |
| def create_review_item( | |
| document_path: str, | |
| document_type: str = "invoice", | |
| **kwargs: Any, | |
| ) -> ReviewItem: | |
| """Create and add a new review item. | |
| Args: | |
| document_path: Path to the source document. | |
| document_type: Document type (default: "invoice"). | |
| **kwargs: Additional fields for the ReviewItem. | |
| Returns: | |
| The created ReviewItem with assigned item_id. | |
| """ | |
| queue: ReviewQueue = get_queue() | |
| item = ReviewItem( | |
| item_id="", | |
| document_path=document_path, | |
| document_type=document_type, | |
| **kwargs, | |
| ) | |
| queue.add_item(item) | |
| return item | |
| def approve_item( | |
| item_id: str, | |
| actor: str, | |
| hs_code: str = "", | |
| notes: str = "", | |
| ) -> bool: | |
| """Approve a review item. | |
| Args: | |
| item_id: Review item ID to approve. | |
| actor: Actor identifier. | |
| hs_code: Final HS code to assign. | |
| notes: Reviewer notes. | |
| Returns: | |
| True if approved, False if item not found or already resolved. | |
| """ | |
| queue: ReviewQueue = get_queue() | |
| item: Optional[ReviewItem] = queue.get_item(item_id) | |
| if item: | |
| return item.approve(actor, hs_code, notes) | |
| return False | |
| def reject_item( | |
| item_id: str, | |
| actor: str, | |
| reason: str = "", | |
| ) -> bool: | |
| """Reject a review item. | |
| Args: | |
| item_id: Review item ID to reject. | |
| actor: Actor identifier. | |
| reason: Rejection reason. | |
| Returns: | |
| True if rejected, False if item not found or already resolved. | |
| """ | |
| queue: ReviewQueue = get_queue() | |
| item: Optional[ReviewItem] = queue.get_item(item_id) | |
| if item: | |
| return item.reject(actor, reason) | |
| return False | |
| # ── Async Database-Backed Service ───────────────────────────────────── | |
| class DashboardService: | |
| """Async service for PostgreSQL-backed compliance review operations. | |
| Every read, write, and state transition runs through an active async | |
| database session. Provides row-level tenant isolation via ``tenant_id`` | |
| filtering on every query. | |
| Attributes: | |
| session_factory: Async SQLAlchemy session factory. | |
| """ | |
| def __init__(self, session_factory: Any) -> None: | |
| """Initialize the service. | |
| Args: | |
| session_factory: ``async_sessionmaker`` instance from SQLAlchemy. | |
| """ | |
| self._session_factory = session_factory | |
| async def _get_session(self) -> Any: | |
| """Create a new async database session. | |
| Returns: | |
| AsyncSession instance. | |
| """ | |
| return self._session_factory() | |
| async def create_item( | |
| self, | |
| tenant_id: str, | |
| document_path: str, | |
| document_type: str = "invoice", | |
| **kwargs: Any, | |
| ) -> dict[str, Any]: | |
| """Create a new review item in the database. | |
| Args: | |
| tenant_id: Tenant identifier for row-level isolation. | |
| document_path: Path to the source document. | |
| document_type: Document type (default: "invoice"). | |
| **kwargs: Additional field overrides. | |
| Returns: | |
| Dict representation of the created item. | |
| """ | |
| import uuid | |
| from sqlalchemy import insert | |
| from hermes.database.models import ComplianceReviewItem | |
| item_id = str(uuid.uuid4()) | |
| now = datetime.utcnow() | |
| session = await self._get_session() | |
| try: | |
| async with session.begin(): | |
| item = ComplianceReviewItem( | |
| id=item_id, | |
| tenant_id=tenant_id, | |
| document_path=document_path, | |
| document_type=kwargs.get("document_type", document_type), | |
| invoice_number=kwargs.get("invoice_number", ""), | |
| invoice_date=kwargs.get("invoice_date", ""), | |
| total_amount=kwargs.get("total_amount", ""), | |
| shipper=kwargs.get("shipper", ""), | |
| consignee=kwargs.get("consignee", ""), | |
| country_origin=kwargs.get("country_origin", ""), | |
| country_destination=kwargs.get("country_destination", ""), | |
| hs_code_suggested=kwargs.get("hs_code_suggested", ""), | |
| hs_code_description=kwargs.get("hs_code_description", ""), | |
| hs_code_confidence=kwargs.get("hs_code_confidence", 0.0), | |
| hs_code_alternatives=kwargs.get("hs_code_alternatives", "[]"), | |
| sanctions_risk_level=kwargs.get("sanctions_risk_level", "clear"), | |
| sanctions_matches=kwargs.get("sanctions_matches", "[]"), | |
| status="pending", | |
| assigned_to=kwargs.get("assigned_to", ""), | |
| priority=kwargs.get("priority", 0), | |
| final_hs_code="", | |
| reviewer_notes="", | |
| created_at=now, | |
| updated_at=now, | |
| ) | |
| session.add(item) | |
| return await self.get_item(tenant_id, item_id) or {"id": item_id} | |
| finally: | |
| await session.close() | |
| async def get_item( | |
| self, tenant_id: str, item_id: str | |
| ) -> Optional[dict[str, Any]]: | |
| """Fetch a single review item by ID with tenant isolation. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| Returns: | |
| Dict of item fields, or None if not found. | |
| """ | |
| from sqlalchemy import select | |
| from hermes.database.models import ComplianceReviewItem | |
| session = await self._get_session() | |
| try: | |
| result = await session.execute( | |
| select(ComplianceReviewItem).where( | |
| ComplianceReviewItem.id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ) | |
| ) | |
| row = result.scalar_one_or_none() | |
| if row is None: | |
| return None | |
| return self._row_to_dict(row) | |
| finally: | |
| await session.close() | |
| async def update_item( | |
| self, | |
| tenant_id: str, | |
| item_id: str, | |
| **fields: Any, | |
| ) -> bool: | |
| """Update fields on a review item. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| **fields: Fields to update. | |
| Returns: | |
| True if the item was found and updated, False otherwise. | |
| """ | |
| from sqlalchemy import update | |
| from hermes.database.models import ComplianceReviewItem | |
| fields["updated_at"] = datetime.utcnow() | |
| session = await self._get_session() | |
| try: | |
| result = await session.execute( | |
| update(ComplianceReviewItem) | |
| .where( | |
| ComplianceReviewItem.id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ) | |
| .values(**fields) | |
| ) | |
| await session.commit() | |
| return result.rowcount > 0 | |
| except Exception: | |
| await session.rollback() | |
| raise | |
| finally: | |
| await session.close() | |
| async def approve_item( | |
| self, | |
| tenant_id: str, | |
| item_id: str, | |
| actor_id: str, | |
| hs_code: str = "", | |
| notes: str = "", | |
| ) -> bool: | |
| """Approve a review item via database. | |
| Performs optimistic concurrency check (status must be | |
| pending/in_review/needs_info). Records audit entry. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| actor_id: Actor identifier (email or system ID). | |
| hs_code: Final HS code to assign. | |
| notes: Reviewer notes. | |
| Returns: | |
| True if approved, False if item not in approvable state. | |
| """ | |
| from hermes.database.models import ( | |
| ComplianceAuditChain, | |
| ComplianceReviewItem, | |
| ) | |
| from sqlalchemy import select, update | |
| session = await self._get_session() | |
| try: | |
| async with session.begin(): | |
| result = await session.execute( | |
| select(ComplianceReviewItem).where( | |
| ComplianceReviewItem.id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ).with_for_update() | |
| ) | |
| row = result.scalar_one_or_none() | |
| if row is None: | |
| return False | |
| if row.status not in ("pending", "in_review", "needs_info"): | |
| return False | |
| prev_status = row.status | |
| row.status = "approved" | |
| row.final_hs_code = hs_code or row.hs_code_suggested | |
| row.reviewer_notes = notes | |
| row.updated_at = datetime.utcnow() | |
| audit_entry = ComplianceAuditChain( | |
| review_item_id=item_id, | |
| timestamp=datetime.utcnow(), | |
| actor_id=actor_id, | |
| action="approve", | |
| previous_state=prev_status, | |
| current_state="approved", | |
| block_hash=self._compute_hash( | |
| item_id, actor_id, "approve", | |
| prev_status, "approved", hs_code, notes, | |
| ), | |
| modified_values=json.dumps({"hs_code": hs_code, "notes": notes}), | |
| reason=notes, | |
| ) | |
| session.add(audit_entry) | |
| return True | |
| except Exception: | |
| await session.rollback() | |
| raise | |
| finally: | |
| await session.close() | |
| async def reject_item( | |
| self, | |
| tenant_id: str, | |
| item_id: str, | |
| actor_id: str, | |
| reason: str = "", | |
| ) -> bool: | |
| """Reject a review item via database. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| actor_id: Actor identifier. | |
| reason: Rejection reason. | |
| Returns: | |
| True if rejected, False if item not in rejectable state. | |
| """ | |
| from hermes.database.models import ( | |
| ComplianceAuditChain, | |
| ComplianceReviewItem, | |
| ) | |
| from sqlalchemy import select, update | |
| session = await self._get_session() | |
| try: | |
| async with session.begin(): | |
| result = await session.execute( | |
| select(ComplianceReviewItem).where( | |
| ComplianceReviewItem.id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ).with_for_update() | |
| ) | |
| row = result.scalar_one_or_none() | |
| if row is None: | |
| return False | |
| if row.status not in ("pending", "in_review", "needs_info"): | |
| return False | |
| prev_status = row.status | |
| row.status = "rejected" | |
| row.reviewer_notes = reason | |
| row.updated_at = datetime.utcnow() | |
| audit_entry = ComplianceAuditChain( | |
| review_item_id=item_id, | |
| timestamp=datetime.utcnow(), | |
| actor_id=actor_id, | |
| action="reject", | |
| previous_state=prev_status, | |
| current_state="rejected", | |
| block_hash=self._compute_hash( | |
| item_id, actor_id, "reject", | |
| prev_status, "rejected", "", reason, | |
| ), | |
| modified_values=json.dumps({"reason": reason}), | |
| reason=reason, | |
| ) | |
| session.add(audit_entry) | |
| return True | |
| except Exception: | |
| await session.rollback() | |
| raise | |
| finally: | |
| await session.close() | |
| async def request_info( | |
| self, | |
| tenant_id: str, | |
| item_id: str, | |
| actor_id: str, | |
| reason: str = "", | |
| ) -> bool: | |
| """Request more information on a review item. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| actor_id: Actor identifier. | |
| reason: Reason for the information request. | |
| Returns: | |
| True if updated, False if item not in a requestable state. | |
| """ | |
| from hermes.database.models import ( | |
| ComplianceAuditChain, | |
| ComplianceReviewItem, | |
| ) | |
| from sqlalchemy import select | |
| session = await self._get_session() | |
| try: | |
| async with session.begin(): | |
| result = await session.execute( | |
| select(ComplianceReviewItem).where( | |
| ComplianceReviewItem.id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ).with_for_update() | |
| ) | |
| row = result.scalar_one_or_none() | |
| if row is None: | |
| return False | |
| if row.status not in ("pending", "in_review"): | |
| return False | |
| prev_status = row.status | |
| row.status = "needs_info" | |
| row.updated_at = datetime.utcnow() | |
| audit_entry = ComplianceAuditChain( | |
| review_item_id=item_id, | |
| timestamp=datetime.utcnow(), | |
| actor_id=actor_id, | |
| action="request_info", | |
| previous_state=prev_status, | |
| current_state="needs_info", | |
| block_hash=self._compute_hash( | |
| item_id, actor_id, "request_info", | |
| prev_status, "needs_info", "", reason, | |
| ), | |
| modified_values=json.dumps({"reason": reason}), | |
| reason=reason, | |
| ) | |
| session.add(audit_entry) | |
| return True | |
| except Exception: | |
| await session.rollback() | |
| raise | |
| finally: | |
| await session.close() | |
| async def list_items( | |
| self, | |
| tenant_id: str, | |
| status: Optional[str] = None, | |
| assignee: Optional[str] = None, | |
| limit: int = 50, | |
| ) -> list[dict[str, Any]]: | |
| """List review items with tenant isolation and optional filters. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| status: Filter by status (None = all). | |
| assignee: Filter by assigned reviewer. | |
| limit: Maximum items to return. | |
| Returns: | |
| List of item dicts sorted by priority DESC, created_at ASC. | |
| """ | |
| from sqlalchemy import select | |
| from hermes.database.models import ComplianceReviewItem | |
| session = await self._get_session() | |
| try: | |
| query = select(ComplianceReviewItem).where( | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ) | |
| if status: | |
| query = query.where(ComplianceReviewItem.status == status) | |
| if assignee: | |
| query = query.where( | |
| ComplianceReviewItem.assigned_to == assignee | |
| ) | |
| query = query.order_by( | |
| ComplianceReviewItem.priority.desc(), | |
| ComplianceReviewItem.created_at.asc(), | |
| ).limit(limit) | |
| result = await session.execute(query) | |
| rows = result.scalars().all() | |
| return [self._row_to_dict(r) for r in rows] | |
| finally: | |
| await session.close() | |
| async def get_statistics(self, tenant_id: str) -> dict[str, int]: | |
| """Get queue statistics scoped to a tenant. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| Returns: | |
| Dict with counts for each status plus total and overdue. | |
| """ | |
| from sqlalchemy import func, select | |
| from hermes.database.models import ComplianceReviewItem | |
| session = await self._get_session() | |
| try: | |
| base = select(func.count()).where( | |
| ComplianceReviewItem.tenant_id == tenant_id | |
| ) | |
| total = (await session.execute( | |
| base | |
| )).scalar() or 0 | |
| pending = (await session.execute( | |
| base.where(ComplianceReviewItem.status == "pending") | |
| )).scalar() or 0 | |
| approved = (await session.execute( | |
| base.where(ComplianceReviewItem.status == "approved") | |
| )).scalar() or 0 | |
| rejected = (await session.execute( | |
| base.where(ComplianceReviewItem.status == "rejected") | |
| )).scalar() or 0 | |
| in_review = (await session.execute( | |
| base.where(ComplianceReviewItem.status == "in_review") | |
| )).scalar() or 0 | |
| needs_info = (await session.execute( | |
| base.where(ComplianceReviewItem.status == "needs_info") | |
| )).scalar() or 0 | |
| overdue = (await session.execute( | |
| base.where( | |
| ComplianceReviewItem.review_deadline < datetime.utcnow(), | |
| ComplianceReviewItem.review_deadline.isnot(None), | |
| ComplianceReviewItem.status.notin_(["approved", "rejected"]), | |
| ) | |
| )).scalar() or 0 | |
| return { | |
| "total": total, | |
| "pending": pending, | |
| "approved": approved, | |
| "rejected": rejected, | |
| "in_review": in_review, | |
| "needs_info": needs_info, | |
| "overdue": overdue, | |
| } | |
| finally: | |
| await session.close() | |
| async def get_audit_trail( | |
| self, tenant_id: str, item_id: str | |
| ) -> list[dict[str, Any]]: | |
| """Get the full audit trail for a review item. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| Returns: | |
| List of audit entry dicts in chronological order. | |
| """ | |
| from sqlalchemy import select | |
| from hermes.database.models import ( | |
| ComplianceAuditChain, | |
| ComplianceReviewItem, | |
| ) | |
| session = await self._get_session() | |
| try: | |
| result = await session.execute( | |
| select(ComplianceAuditChain) | |
| .join(ComplianceReviewItem) | |
| .where( | |
| ComplianceAuditChain.review_item_id == item_id, | |
| ComplianceReviewItem.tenant_id == tenant_id, | |
| ) | |
| .order_by(ComplianceAuditChain.timestamp.asc()) | |
| ) | |
| rows = result.scalars().all() | |
| return [ | |
| { | |
| "id": r.id, | |
| "review_item_id": r.review_item_id, | |
| "timestamp": r.timestamp.isoformat() if r.timestamp else "", | |
| "actor_id": r.actor_id, | |
| "action": r.action, | |
| "previous_state": r.previous_state, | |
| "current_state": r.current_state, | |
| "block_hash": r.block_hash, | |
| "modified_values": r.modified_values, | |
| "reason": r.reason, | |
| } | |
| for r in rows | |
| ] | |
| finally: | |
| await session.close() | |
| async def verify_audit_chain( | |
| self, tenant_id: str, item_id: str | |
| ) -> bool: | |
| """Verify the SHA-256 hash chain integrity for an item's audit log. | |
| Args: | |
| tenant_id: Tenant identifier. | |
| item_id: Review item ID. | |
| Returns: | |
| True if chain is intact, False if broken. | |
| """ | |
| entries = await self.get_audit_trail(tenant_id, item_id) | |
| if not entries: | |
| return True | |
| prev_hash = "0" * 64 | |
| for entry in entries: | |
| if entry["block_hash"][:64] != self._compute_hash( | |
| entry["review_item_id"], | |
| entry["actor_id"], | |
| entry["action"], | |
| entry["previous_state"], | |
| entry["current_state"], | |
| "", | |
| entry.get("reason", ""), | |
| ): | |
| return False | |
| prev_hash = entry["block_hash"] | |
| return True | |
| def _row_to_dict(row: Any) -> dict[str, Any]: | |
| """Convert a SQLAlchemy row to a dict. | |
| Args: | |
| row: SQLAlchemy model instance. | |
| Returns: | |
| Dict of column values. | |
| """ | |
| return { | |
| "id": row.id, | |
| "tenant_id": row.tenant_id, | |
| "document_path": row.document_path, | |
| "document_type": row.document_type, | |
| "invoice_number": row.invoice_number or "", | |
| "invoice_date": row.invoice_date or "", | |
| "total_amount": row.total_amount or "", | |
| "shipper": row.shipper or "", | |
| "consignee": row.consignee or "", | |
| "country_origin": row.country_origin or "", | |
| "country_destination": row.country_destination or "", | |
| "hs_code_suggested": row.hs_code_suggested or "", | |
| "hs_code_description": row.hs_code_description or "", | |
| "hs_code_confidence": row.hs_code_confidence or 0.0, | |
| "hs_code_alternatives": row.hs_code_alternatives or "[]", | |
| "sanctions_risk_level": row.sanctions_risk_level or "clear", | |
| "sanctions_matches": row.sanctions_matches or "[]", | |
| "status": row.status, | |
| "assigned_to": row.assigned_to or "", | |
| "priority": row.priority or 0, | |
| "final_hs_code": row.final_hs_code or "", | |
| "reviewer_notes": row.reviewer_notes or "", | |
| "created_at": row.created_at.isoformat() if row.created_at else "", | |
| "updated_at": row.updated_at.isoformat() if row.updated_at else "", | |
| "review_deadline": row.review_deadline.isoformat() if row.review_deadline else None, | |
| } | |
| def _compute_hash( | |
| item_id: str, | |
| actor_id: str, | |
| action: str, | |
| previous_state: str, | |
| current_state: str, | |
| hs_code: str, | |
| reason: str, | |
| ) -> str: | |
| """Compute SHA-256 hash for an audit entry. | |
| Args: | |
| item_id: Review item ID. | |
| actor_id: Actor identifier. | |
| action: Action taken. | |
| previous_state: Status before action. | |
| current_state: Status after action. | |
| hs_code: HS code assigned. | |
| reason: Reason for action. | |
| Returns: | |
| SHA-256 hex digest string. | |
| """ | |
| import hashlib | |
| payload = json.dumps( | |
| { | |
| "item": item_id, | |
| "actor": actor_id, | |
| "action": action, | |
| "prev": previous_state, | |
| "new": current_state, | |
| "hs": hs_code, | |
| "reason": reason, | |
| }, | |
| sort_keys=True, | |
| separators=(",", ":"), | |
| ) | |
| return hashlib.sha256(payload.encode("utf-8")).hexdigest() | |