Skip to content

Commit a4517a4

Browse files
FEAT: add the ability to email action plan (#123)
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> Bypassing review due to low engineering coverage over the holidays.
1 parent fb4cb6c commit a4517a4

10 files changed

Lines changed: 482 additions & 40 deletions

File tree

app/src/common/components.py

Lines changed: 63 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,42 @@
3030
logger = logging.getLogger(__name__)
3131

3232

33+
def format_resources(resources: list[dict]) -> str:
34+
return "\n\n".join([format_resource(resource) for resource in resources])
35+
36+
37+
def format_resource(resource: dict) -> str:
38+
return "\n".join(
39+
[
40+
f"### {resource.get('name', 'Unnamed Resource')}",
41+
f"- Referral Type: {resource.get('referral_type', 'None')}",
42+
f"- Description: {resource.get('description', 'None')}",
43+
f"- Website: {resource.get('website', 'None')}",
44+
f"- Phone: {', '.join(resource.get('phones', ['None']))}",
45+
f"- Email: {', '.join(resource.get('emails', ['None']))}",
46+
f"- Addresses: {', '.join(resource.get('addresses', ['None']))}",
47+
]
48+
)
49+
50+
51+
def format_action_plan(action_plan: dict) -> str:
52+
"""Format the action plan for email display. Returns empty string if no action plan."""
53+
if not action_plan:
54+
return ""
55+
56+
title = action_plan.get("title", "Your Action Plan")
57+
summary = action_plan.get("summary", "")
58+
content = action_plan.get("content", "")
59+
60+
parts = [f"## {title}"]
61+
if summary:
62+
parts.append(f"\n{summary}")
63+
if content:
64+
parts.append(f"\n{content}")
65+
66+
return "\n".join(parts)
67+
68+
3369
@component
3470
class EchoNode:
3571
"""
@@ -209,6 +245,32 @@ def run(
209245
"""
210246

211247

248+
@component
249+
class EmailFullResult:
250+
"""
251+
Formats JSON object (representing a list of resources and action plan) and sends it to email address.
252+
"""
253+
254+
@component.output_types(status=str, email=str, message=str)
255+
def run(self, email: str, resources_dict: dict, action_plan_dict: dict) -> dict:
256+
logger.info("Emailing result to %s", email)
257+
logger.debug("Resources JSON content:\n%s", json.dumps(resources_dict, indent=2))
258+
if action_plan_dict:
259+
logger.debug("Action plan JSON content:\n%s", json.dumps(action_plan_dict, indent=2))
260+
261+
formatted_resources = format_resources(resources_dict.get("resources", []))
262+
formatted_action_plan = format_action_plan(action_plan_dict)
263+
264+
message = f"{EMAIL_INTRO}\n{formatted_resources}\n\n{formatted_action_plan}"
265+
266+
# Send email via AWS SES
267+
subject = "Your Requested Resources and Action Plan"
268+
success = send_email(recipient=email, subject=subject, body=message)
269+
status = "success" if success else "failed"
270+
271+
return {"status": status, "email": email, "message": message}
272+
273+
212274
@component
213275
class EmailResult:
214276
"""
@@ -219,7 +281,7 @@ class EmailResult:
219281
def run(self, email: str, json_dict: dict) -> dict:
220282
logger.info("Emailing result to %s", email)
221283
logger.debug("JSON content:\n%s", json.dumps(json_dict, indent=2))
222-
formatted_resources = self.format_resources(json_dict.get("resources", []))
284+
formatted_resources = format_resources(json_dict.get("resources", []))
223285
message = f"{EMAIL_INTRO}\n{formatted_resources}"
224286

225287
# Send email via AWS SES
@@ -229,22 +291,6 @@ def run(self, email: str, json_dict: dict) -> dict:
229291

230292
return {"status": status, "email": email, "message": message}
231293

232-
def format_resources(self, resources: list[dict]) -> str:
233-
return "\n\n".join([self.format_resource(resource) for resource in resources])
234-
235-
def format_resource(self, resource: dict) -> str:
236-
return "\n".join(
237-
[
238-
f"### {resource.get('name', 'Unnamed Resource')}",
239-
f"- Referral Type: {resource.get('referral_type', 'None')}",
240-
f"- Description: {resource.get('description', 'None')}",
241-
f"- Website: {resource.get('website', 'None')}",
242-
f"- Phone: {', '.join(resource.get('phones', ['None']))}",
243-
f"- Email: {', '.join(resource.get('emails', ['None']))}",
244-
f"- Addresses: {', '.join(resource.get('addresses', ['None']))}",
245-
]
246-
)
247-
248294

249295
BaseModelT = TypeVar("BaseModelT", bound=BaseModel)
250296

app/src/pipelines/email_full_result/__init__.py

Whitespace-only changes.
Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
1+
import logging
2+
from pprint import pformat
3+
4+
from fastapi import HTTPException
5+
from hayhooks import BasePipelineWrapper
6+
from haystack import Pipeline
7+
from haystack.core.errors import PipelineRuntimeError
8+
from openinference.instrumentation import _tracers, using_metadata
9+
from opentelemetry.trace.status import Status, StatusCode
10+
11+
from src.common import components, phoenix_utils
12+
13+
logger = logging.getLogger(__name__)
14+
tracer = phoenix_utils.tracer_provider.get_tracer(__name__)
15+
16+
17+
class PipelineWrapper(BasePipelineWrapper):
18+
name = "email_full_result"
19+
20+
def setup(self) -> None:
21+
pipeline = Pipeline()
22+
pipeline.add_component("load_resources", components.LoadResult())
23+
pipeline.add_component("load_action_plan", components.LoadResult())
24+
pipeline.add_component("email_full_result", components.EmailFullResult())
25+
26+
pipeline.connect("load_resources.result_json", "email_full_result.resources_dict")
27+
pipeline.connect("load_action_plan.result_json", "email_full_result.action_plan_dict")
28+
pipeline.add_component("logger", components.ReadableLogger())
29+
30+
self.pipeline = pipeline
31+
32+
def run_api(self, resources_result_id: str, action_plan_result_id: str, email: str) -> dict:
33+
with using_metadata({"email": email}):
34+
# Must set using_metadata context before calling tracer.start_as_current_span()
35+
assert isinstance(tracer, _tracers.OITracer), f"Got unexpected {type(tracer)}"
36+
with tracer.start_as_current_span( # pylint: disable=not-context-manager,unexpected-keyword-arg
37+
self.name, openinference_span_kind="chain"
38+
) as span:
39+
result = self._run(resources_result_id, action_plan_result_id, email)
40+
span.set_input(
41+
{
42+
"resources_result_id": resources_result_id,
43+
"action_plan_result_id": action_plan_result_id,
44+
}
45+
)
46+
span.set_output(result["email_full_result"]["status"])
47+
span.set_status(Status(StatusCode.OK))
48+
return result
49+
50+
def _run(self, resources_result_id: str, action_plan_result_id: str, email: str) -> dict:
51+
try:
52+
run_data = {
53+
"logger": {
54+
"messages_list": [
55+
{
56+
"resources_result_id": resources_result_id,
57+
"action_plan_result_id": action_plan_result_id,
58+
"email": email,
59+
}
60+
],
61+
},
62+
"load_resources": {
63+
"result_id": resources_result_id,
64+
},
65+
"email_full_result": {
66+
"email": email,
67+
},
68+
}
69+
70+
run_data["load_action_plan"] = {
71+
"result_id": action_plan_result_id,
72+
}
73+
74+
response = self.pipeline.run(
75+
run_data,
76+
include_outputs_from={"email_full_result"},
77+
)
78+
logger.debug("Results: %s", pformat(response, width=160))
79+
return response
80+
except PipelineRuntimeError as re:
81+
error_msg = str(re)
82+
if re.component_type == components.LoadResult:
83+
if "Invalid JSON format in result" in error_msg:
84+
status_code = 500 # Internal error
85+
else:
86+
status_code = 400 # User error
87+
else:
88+
status_code = 500 # Internal error
89+
90+
raise HTTPException(
91+
status_code=status_code,
92+
detail=f"Error occurred: {error_msg}",
93+
) from re

app/src/pipelines/email_result/pipeline_wrapper.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,15 +28,15 @@ def setup(self) -> None:
2828

2929
self.pipeline = pipeline
3030

31-
def run_api(self, result_id: str, email: str) -> dict:
31+
def run_api(self, resources_result_id: str, email: str) -> dict:
3232
with using_metadata({"email": email}):
3333
# Must set using_metadata context before calling tracer.start_as_current_span()
3434
assert isinstance(tracer, _tracers.OITracer), f"Got unexpected {type(tracer)}"
3535
with tracer.start_as_current_span( # pylint: disable=not-context-manager,unexpected-keyword-arg
3636
self.name, openinference_span_kind="chain"
3737
) as span:
38-
result = self._run(result_id, email)
39-
span.set_input(result_id)
38+
result = self._run(resources_result_id, email)
39+
span.set_input(resources_result_id)
4040
span.set_output(result["email_result"]["status"])
4141
span.set_status(Status(StatusCode.OK))
4242
return result

app/src/pipelines/generate_action_plan/pipeline_wrapper.py

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,12 @@
99
from pydantic import BaseModel
1010

1111
from src.common import haystack_utils, phoenix_utils
12-
from src.common.components import OpenAIWebSearchGenerator, ReadableLogger
12+
from src.common.components import (
13+
LlmOutputValidator,
14+
OpenAIWebSearchGenerator,
15+
ReadableLogger,
16+
SaveResult,
17+
)
1318
from src.pipelines.generate_referrals.pipeline_wrapper import Resource
1419

1520
logger = logging.getLogger(__name__)
@@ -50,7 +55,12 @@ def setup(self) -> None:
5055
),
5156
name="prompt_builder",
5257
)
58+
pipeline.add_component("output_validator", LlmOutputValidator(ActionPlan))
59+
pipeline.add_component("save_result", SaveResult())
60+
5361
pipeline.connect("prompt_builder", "llm.messages")
62+
pipeline.connect("llm.replies", "output_validator")
63+
pipeline.connect("output_validator.valid_replies", "save_result.replies")
5464

5565
pipeline.add_component("logger", ReadableLogger())
5666
pipeline.connect("llm", "logger")
@@ -90,10 +100,13 @@ def _run(self, resource_objects: list[Resource], user_email: str, user_query: st
90100
},
91101
"llm": {"model": "gpt-5-mini", "reasoning_effort": "low"},
92102
},
93-
include_outputs_from={"llm"},
103+
include_outputs_from={"llm", "save_result"},
94104
)
95105
logger.debug("Results: %s", pformat(response, width=160))
96-
return {"response": response["llm"]["replies"][0]._content[0].text}
106+
return {
107+
"response": response["llm"]["replies"][0]._content[0].text,
108+
"save_result": response["save_result"],
109+
}
97110

98111

99112
def get_resources(resources: list[Resource] | list[dict]) -> list[Resource]:

0 commit comments

Comments
 (0)