Repository navigation
Expand file tree
/
Copy pathagent.py
More file actions
165 lines (143 loc) · 6.46 KB
/
Copy pathagent.py
File metadata and controls
165 lines (143 loc) · 6.46 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
"""
Main agent loop.
Design choices that map directly to the spec:
- Tasks are grouped by (platform, order_id) ONLY to decide batch-vs-sequential
and to open one Stagehand/Browserbase session per platform -- write-back
always happens per line item (excel_io.write_result is called once per
row, never once per order).
- Each line item is processed in its own try/except so one failure never
aborts the rest of the order (partial-success handling).
- A session is opened once per platform (logged in once), reused across
that platform's orders in this run.
"""
import logging
import os
from collections import defaultdict
import config
from excel_io import ExcelTaskStore, ReturnTask
from platforms import get_adapter
from platforms.base import ReturnFlowType, ReturnResult
import stagehand_client as sh
REQUIRED_ENV_VARS = ["BROWSERBASE_API_KEY", "BROWSERBASE_PROJECT_ID", "MODEL_API_KEY"]
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
handlers=[
logging.FileHandler(config.LOG_DIR / "agent.log"),
logging.StreamHandler(),
],
)
log = logging.getLogger("returns_agent")
def group_tasks_by_order(tasks: list[ReturnTask]) -> dict:
groups = defaultdict(list)
for t in tasks:
groups[(t.platform, t.order_id)].append(t)
return groups
def process_order(adapter, session_id: str, order_id: str, line_items: list[ReturnTask],
store: ExcelTaskStore):
flow_type = adapter.detect_flow_type(session_id, order_id)
log.info(f"Order {order_id}: detected flow type = {flow_type.value}, "
f"{len(line_items)} line item(s)")
if flow_type == ReturnFlowType.BATCH:
skus = [t.product_sku for t in line_items]
try:
results = adapter.initiate_return_for_batch(session_id, order_id, skus)
except Exception as e:
log.error(f"Order {order_id}: batch return call failed entirely: {e}")
results = {t.product_sku: ReturnResult(success=False, error_note=str(e))
for t in line_items}
else:
results = {}
for t in line_items:
# Each item runs independently -- a failure here must not stop
# the loop from continuing to the next item.
try:
results[t.product_sku] = adapter.initiate_return_for_item(
session_id, order_id, t.product_sku, t.return_window)
except Exception as e:
log.error(f"Order {order_id} / SKU {t.product_sku}: unhandled exception: {e}")
results[t.product_sku] = ReturnResult(success=False, error_note=str(e))
# Write back per line item, regardless of flow type used above.
for t in line_items:
result = results.get(t.product_sku)
if result is None:
store.write_result(t, return_status=config.RETURN_FAILED,
task_status=config.STATUS_NEEDS_REVIEW,
note="No result returned by adapter")
continue
if result.success:
store.write_result(
t, return_id=result.return_id, return_status=config.RETURN_PLACED,
refund_amount=result.refund_amount, task_status=config.STATUS_DONE,
note="",
)
log.info(f"Order {order_id} / SKU {t.product_sku}: return placed "
f"(Return ID: {result.return_id})")
elif result.out_of_window:
store.write_result(
t, return_status=config.RETURN_OUT_OF_WINDOW,
task_status=config.STATUS_NEEDS_REVIEW, note=result.error_note,
)
log.warning(f"Order {order_id} / SKU {t.product_sku}: out of return window")
elif result.support_needed:
store.write_result(
t, return_status=config.RETURN_SUPPORT_NEEDED,
task_status=config.STATUS_NEEDS_REVIEW, note=result.error_note,
)
log.warning(f"Order {order_id} / SKU {t.product_sku}: no direct return option -- "
f"needs chat/human support")
else:
store.write_result(
t, return_status=config.RETURN_FAILED,
task_status=config.STATUS_NEEDS_REVIEW, note=result.error_note,
)
log.warning(f"Order {order_id} / SKU {t.product_sku}: failed -- {result.error_note}")
if store.mark_order_done_if_complete(order_id):
log.info(f"Order {order_id}: all line items reached a terminal state.")
else:
log.info(f"Order {order_id}: some line items still need human review.")
def run():
missing = [v for v in REQUIRED_ENV_VARS if not os.environ.get(v)]
if missing:
log.error(f"Missing required environment variable(s): {', '.join(missing)}. "
f"Set these before running (Browserbase + model API credentials).")
return
store = ExcelTaskStore()
tasks = store.get_pending_tasks()
if not tasks:
log.info("No pending tasks found. Nothing to do.")
return
grouped = group_tasks_by_order(tasks)
log.info(f"Found {len(tasks)} pending line item(s) across {len(grouped)} order(s).")
# One Stagehand session per platform (logged in once), reused across
# that platform's orders in this run.
sessions = {}
try:
for (platform, order_id), line_items in grouped.items():
adapter = get_adapter(platform)
if platform not in sessions:
log.info(f"Logging into {platform}...")
session_id = sh.start_session(platform)
adapter.login(session_id)
sessions[platform] = session_id
else:
session_id = sessions[platform]
try:
process_order(adapter, session_id, order_id, line_items, store)
except Exception as e:
log.error(f"Order {order_id} on {platform}: order-level failure: {e}")
for t in line_items:
store.write_result(
t, return_status=config.RETURN_FAILED,
task_status=config.STATUS_NEEDS_REVIEW,
note=f"Order-level error: {e}",
)
finally:
for session_id in sessions.values():
try:
sh.end_session(session_id)
except Exception as e:
log.warning(f"Failed to cleanly end session {session_id}: {e}")
log.info("Run complete.")
if __name__ == "__main__":
run()