-
Notifications
You must be signed in to change notification settings - Fork 86
Expand file tree
/
Copy pathtranscriptToEmailFlow.ts
More file actions
171 lines (138 loc) · 4.52 KB
/
Copy pathtranscriptToEmailFlow.ts
File metadata and controls
171 lines (138 loc) · 4.52 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
166
167
168
169
170
'use server';
import { createAI, createVertexAI } from '../genkitFactory';
import { z } from "genkit";
import { createSimpleFirestoreVSRetriever } from '../retriever/simpleSearchRetriever';
import { vertexAI } from '@genkit-ai/vertexai';
// Create AI instance & retriever using the factory
// const ai = await createAI(gemini25FlashPreview0417);
const ai = await createVertexAI();
await createSimpleFirestoreVSRetriever(ai);
// STEP 1: Task Extraction
const TaskSchema = z.object({
description: z.string(),
});
const TaskArraySchema = ai.defineSchema(
'TaskArraySchema',
z.array(TaskSchema)
);
export const taskExtractionFlow = ai.defineFlow(
{
name: "taskExtractionFlow",
inputSchema: z.string(),
outputSchema: TaskArraySchema,
},
async (transcript) => {
console.log("Running Task Extraction Flow on transcript...");
const taskExtractionPrompt = ai.prompt('taskExtraction');
const { output } = await taskExtractionPrompt(
{
transcript: transcript,
},
{
model: vertexAI.model('gemini-2.5-pro'),
output: { schema: TaskArraySchema }
}
);
return output;
}
);
// STEP 2: Task Research
const TaskResearchResponse = ai.defineSchema(
'TaskResearchResponse',
z.object({
answer: z.string(),
caveats: z.array(z.string()),
docReferences: z.array(z.object({
title: z.string(),
url: z.string(),
relevantContent: z.string().optional(),
})),
})
);
const TaskResearchResponseArray = ai.defineSchema(
'TaskResearchResponseArray',
z.array(TaskResearchResponse)
);
export const taskReseachFlow = ai.defineFlow(
{
name: "taskReseachFlow",
inputSchema: TaskArraySchema,
outputSchema: TaskResearchResponseArray,
},
async (tasks) => {
console.log("Running Task Research Flow on transcript...");
// 1. Retrieve relevant documents for each task in parallel
const retrievalPromises = tasks.map(task => {
return ai.retrieve({
retriever: 'simpleFirestoreVSRetriever',
query: task.description, // Use task description as query
options: { limit: 10 }
});
});
const taskDocsArray = await Promise.all(retrievalPromises);
// 2. Parallel task research generation
const taskResearchPrompt = ai.prompt('taskResearch');
const generationPromises = tasks.map(async (task, index) => {
const docs = taskDocsArray[index];
const { output } = await taskResearchPrompt(
{
question: task.description
},
{
docs: docs,
output: { schema: TaskResearchResponse }
}
);
return output;
});
const researchResults = await Promise.all(generationPromises);
return researchResults;
}
);
// STEP 3: Email Generation based on research
export const emailAggregationFlow = ai.defineFlow(
{
name: "emailAggregationFlow",
inputSchema: z.object({
tasks: TaskArraySchema,
researchResults: TaskResearchResponseArray
}),
outputSchema: z.string(),
},
async (input) => {
const emailGenerationPrompt = ai.prompt('emailGeneration');
const { output } = await emailGenerationPrompt({
tasks: JSON.stringify(input.tasks, null, 2),
research: JSON.stringify(input.researchResults, null, 2)
},
{
model: vertexAI.model('gemini-2.5-pro'),
});
return output.email;
}
);
// Aggregated Orchestration Flow
export const transcriptToEmailFlow = ai.defineFlow(
{
name: "transcriptToEmailFlow",
inputSchema: z.string(),
outputSchema: z.object({
tasks: TaskArraySchema,
research: TaskResearchResponseArray,
email: z.string(),
})
},
async (transcript) => {
const extractedTasks = await taskExtractionFlow(transcript);
const taskResearchResults = await taskReseachFlow(extractedTasks);
const email = await emailAggregationFlow({
tasks: extractedTasks,
researchResults: taskResearchResults
});
return {
tasks: extractedTasks,
research: taskResearchResults,
email: email
};
}
);